()
| 88 | } |
| 89 | |
| 90 | func (s *SFlow) run() { |
| 91 | var err error |
| 92 | // exit if the sflow is disabled |
| 93 | if !opts.SFlowEnabled { |
| 94 | logger.Println("sflow has been disabled") |
| 95 | return |
| 96 | } |
| 97 | |
| 98 | s.pool = make(chan chan struct{}, maxWorkers) |
| 99 | |
| 100 | hostPort := net.JoinHostPort(s.addr, strconv.Itoa(s.port)) |
| 101 | udpAddr, _ := net.ResolveUDPAddr("udp", hostPort) |
| 102 | |
| 103 | s.conn, err = net.ListenUDP("udp", udpAddr) |
| 104 | if err != nil { |
| 105 | logger.Fatal(err) |
| 106 | } |
| 107 | |
| 108 | atomic.AddInt32(&s.stats.Workers, int32(s.workers)) |
| 109 | for i := 0; i < s.workers; i++ { |
| 110 | go func() { |
| 111 | wQuit := make(chan struct{}) |
| 112 | s.pool <- wQuit |
| 113 | s.sFlowWorker(wQuit) |
| 114 | }() |
| 115 | } |
| 116 | |
| 117 | go mirrorSFlowDispatcher(sFlowMCh) |
| 118 | |
| 119 | logger.Printf("sFlow is running (UDP: listening on [::]:%d workers#: %d)", s.port, s.workers) |
| 120 | |
| 121 | go func() { |
| 122 | if !opts.ProducerEnabled { |
| 123 | return |
| 124 | } |
| 125 | |
| 126 | p := producer.NewProducer(opts.MQName) |
| 127 | p.MQConfigFile = path.Join(opts.VFlowConfigPath, opts.MQConfigFile) |
| 128 | p.MQErrorCount = &s.stats.MQErrorCount |
| 129 | p.Logger = logger |
| 130 | p.Chan = sFlowMQCh |
| 131 | p.Topic = opts.SFlowTopic |
| 132 | |
| 133 | if err := p.Run(); err != nil { |
| 134 | logger.Fatal(err) |
| 135 | } |
| 136 | }() |
| 137 | |
| 138 | go func() { |
| 139 | if !opts.DynWorkers { |
| 140 | logger.Println("sFlow dynamic worker disabled") |
| 141 | return |
| 142 | } |
| 143 | |
| 144 | s.dynWorkers() |
| 145 | }() |
| 146 | |
| 147 | for !s.stop { |
nothing calls this directly
no test coverage detected