@@ -164,39 +164,7 @@ func (p *Protocol) OpenInbound(ctx context.Context) error {
164164 }()
165165
166166 // start to monitor the server heartbeat
167- go func () {
168- // if serverHealthinessTimeout is disable, ignore the server heartbeat
169- if p .serverHealthinessTimeout == 0 {
170- for {
171- select {
172- case heartbeat := <- heartbeatCh :
173- logger .Debugf ("heartbeat received %v" , heartbeat )
174- case <- subCtx .Done ():
175- // context canceled, exiting
176- return
177- }
178- }
179- }
180-
181- // if no heartbeat was received duration the serverHealthinessTimeout, send the timeout error to
182- // reconnectErrorChan
183- for {
184- select {
185- case heartbeat := <- heartbeatCh :
186- logger .Debugf ("heartbeat received %v" , heartbeat )
187- case <- time .After (p .serverHealthinessTimeout ):
188- select {
189- case p .reconnectErrorChan <- fmt .Errorf ("stream timeout: no message received for %d seconds" , p .serverHealthinessTimeout ):
190- default :
191- // error chan is full, do nothing
192- }
193- return
194- case <- subCtx .Done ():
195- // context canceled, exiting
196- return
197- }
198- }
199- }()
167+ go startHeartbeatWatcher (ctx , heartbeatCh )
200168
201169 // Wait until external or internal context done
202170 select {
@@ -229,3 +197,60 @@ func (p *Protocol) Close(ctx context.Context) error {
229197 close (p .closeChan )
230198 return nil
231199}
200+
201+ func startHeartbeatWatcher (
202+ ctx context.Context ,
203+ heartbeatCh <- chan * pbv1.CloudEvent ,
204+ ) {
205+ logger := cecontext .LoggerFrom (ctx )
206+ // if serverHealthinessTimeout is disable, ignore the server heartbeat timeout
207+ if p .serverHealthinessTimeout == 0 {
208+ for {
209+ select {
210+ case msg , ok := <- heartbeatCh :
211+ if ! ok {
212+ return
213+ }
214+ logger .Debugf ("heartbeat received %v" , msg )
215+ case <- ctx .Done ():
216+ return
217+ }
218+ }
219+ }
220+
221+ // if no heartbeat was received duration the serverHealthinessTimeout, send the timeout error to
222+ // reconnectErrorChan
223+ timeout := p .serverHealthinessTimeout
224+ timer := time .NewTimer (timeout )
225+ defer timer .Stop ()
226+
227+ for {
228+ select {
229+ case msg , ok := <- heartbeatCh :
230+ if ! ok {
231+ return
232+ }
233+ logger .Debugf ("heartbeat received %v" , msg )
234+
235+ // reset timer safely
236+ if ! timer .Stop () {
237+ select {
238+ case <- timer .C :
239+ default :
240+ }
241+ }
242+ timer .Reset (timeout )
243+
244+ case <- timer .C :
245+ // timeout reached: no heartbeat received
246+ select {
247+ case p .reconnectErrorChan <- fmt .Errorf ("stream timeout: no message received for %v" , timeout ):
248+ default :
249+ }
250+ return
251+
252+ case <- ctx .Done ():
253+ return
254+ }
255+ }
256+ }
0 commit comments