MCPcopy Create free account
hub / github.com/AdRoll/baker / Run

Method Run

input/tcp.go:73–118  ·  view source on GitHub ↗
(inch chan<- *baker.Data)

Source from the content-addressed store, hash-verified

71}
72
73func (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
120func (s *TCP) setOutputChannel(data chan<- *baker.Data) {
121 s.data = data

Callers

nothing calls this directly

Calls 6

setOutputChannelMethod · 0.95
handleStreamMethod · 0.95
DoneMethod · 0.80
CloseMethod · 0.65
ErrorMethod · 0.45
WaitMethod · 0.45

Tested by

no test coverage detected