Monorepo for Tangled tangled.org
1

Configure Feed

Select the types of activity you want to include in your feed.

eventconsumer: add 90s timeout for conns that somehow become half-open

Signed-off-by: dawn <dawn@tangled.org>

author
dawn
committer
Tangled
date (Jul 21, 2026, 5:03 PM +0300) commit 13d8cece parent 17c82409 change-id qqumuorm
+19
+19
eventconsumer/consumer.go
··· 18 18 19 19 type ProcessFunc func(ctx context.Context, source Source, event eventstream.Event) error 20 20 21 + // server sends a ping every 30s, so any silence longer than this means the 22 + // connection is half-open, dropped without a close frame. 23 + // so without a read deadline, ReadMessage only notices when the kernel's 24 + // tcp keepalive gives up. 25 + const livenessTimeout = 90 * time.Second 26 + 21 27 type ConsumerConfig struct { 22 28 Sources map[Source]struct{} 23 29 ProcessFunc ProcessFunc ··· 325 331 326 332 c.logger.Info("connected", "source", source) 327 333 334 + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) 335 + conn.SetPongHandler(func(string) error { 336 + return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) 337 + }) 338 + conn.SetPingHandler(func(appData string) error { 339 + err := conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(10*time.Second)) 340 + if err != nil { 341 + return err 342 + } 343 + return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) 344 + }) 345 + 328 346 for { 329 347 select { 330 348 case <-ctx.Done(): ··· 337 355 if msgType != websocket.TextMessage { 338 356 continue 339 357 } 358 + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) 340 359 select { 341 360 case c.jobQueue <- job{source: source, message: msg}: 342 361 case <-ctx.Done():