(inch chan<- *baker.Data)
| 71 | } |
| 72 | |
| 73 | func (s *TCP) Run(inch chan<- *baker.Data) error { |
| 74 | s.setOutputChannel(inch) |
| 75 | |
| 76 | ctxLog := log.WithFields(log.Fields{"f": "Run"}) |
| 77 | |
| 78 | addr, err := net.ResolveTCPAddr("tcp", s.Cfg.Listener) |
| 79 | if err != nil { |
| 80 | ctxLog.WithFields(log.Fields{"listener": s.Cfg.Listener}).Error("Can't resolve") |
| 81 | return err |
| 82 | } |
| 83 | |
| 84 | l, err := net.ListenTCP("tcp", addr) |
| 85 | if err != nil { |
| 86 | return err |
| 87 | } |
| 88 | defer l.Close() |
| 89 | |
| 90 | wg := sync.WaitGroup{} |
| 91 | for atomic.LoadInt64(&s.stop) == 0 { |
| 92 | l.SetDeadline(time.Now().Add(1 * time.Second)) |
| 93 | conn, err := l.AcceptTCP() |
| 94 | if err != nil { |
| 95 | if errors.Is(err, os.ErrDeadlineExceeded) { |
| 96 | continue |
| 97 | } |
| 98 | ctxLog.WithFields(log.Fields{"error": err}).Error("Error while accepting") |
| 99 | } |
| 100 | |
| 101 | ctxLog.WithFields(log.Fields{"addr": conn.RemoteAddr()}).Info("Connected") |
| 102 | |
| 103 | wg.Add(1) |
| 104 | go func(conn *net.TCPConn) { |
| 105 | defer func() { |
| 106 | conn.Close() |
| 107 | wg.Done() |
| 108 | }() |
| 109 | |
| 110 | if err := s.handleStream(conn); err != nil { |
| 111 | ctxLog.WithError(err).WithFields(log.Fields{"error": err}).Error("Error when handling stream") |
| 112 | } |
| 113 | }(conn) |
| 114 | } |
| 115 | |
| 116 | wg.Wait() |
| 117 | return nil |
| 118 | } |
| 119 | |
| 120 | func (s *TCP) setOutputChannel(data chan<- *baker.Data) { |
| 121 | s.data = data |
nothing calls this directly
no test coverage detected