(wQuit chan struct{})
| 173 | } |
| 174 | |
| 175 | func (s *SFlow) sFlowWorker(wQuit chan struct{}) { |
| 176 | var ( |
| 177 | reader *bytes.Reader |
| 178 | msg SFUDPMsg |
| 179 | mirror SFUDPMsg |
| 180 | ok bool |
| 181 | b []byte |
| 182 | ) |
| 183 | |
| 184 | LOOP: |
| 185 | for { |
| 186 | |
| 187 | select { |
| 188 | case <-wQuit: |
| 189 | break LOOP |
| 190 | case msg, ok = <-sFlowUDPCh: |
| 191 | if !ok { |
| 192 | break LOOP |
| 193 | } |
| 194 | } |
| 195 | |
| 196 | if opts.Verbose { |
| 197 | logger.Printf("rcvd sflow data from: %s, size: %d bytes", |
| 198 | msg.raddr, len(msg.body)) |
| 199 | } |
| 200 | |
| 201 | if sFlowMirrorEnabled { |
| 202 | mirror.raddr = msg.raddr |
| 203 | mirror.body = sFlowBuffer.Get().([]byte) |
| 204 | mirror.body = append(mirror.body[:0], msg.body...) |
| 205 | |
| 206 | select { |
| 207 | case sFlowMCh <- mirror: |
| 208 | default: |
| 209 | } |
| 210 | } |
| 211 | |
| 212 | reader = bytes.NewReader(msg.body) |
| 213 | d := sflow.NewSFDecoder(reader, opts.SFlowTypeFilter) |
| 214 | datagram, err := d.SFDecode() |
| 215 | if err != nil || (len(datagram.Counters) < 1 && len(datagram.Samples) < 1) { |
| 216 | sFlowBuffer.Put(msg.body[:opts.SFlowUDPSize]) |
| 217 | continue |
| 218 | } |
| 219 | |
| 220 | b, err = json.Marshal(datagram) |
| 221 | if err != nil { |
| 222 | sFlowBuffer.Put(msg.body[:opts.SFlowUDPSize]) |
| 223 | logger.Println(err) |
| 224 | continue |
| 225 | } |
| 226 | |
| 227 | atomic.AddUint64(&s.stats.DecodedCount, 1) |
| 228 | |
| 229 | if opts.Verbose { |
| 230 | logger.Println(string(b)) |
| 231 | } |
| 232 |
no test coverage detected