From 13d8ceceec456827c37a21e6fb103d7d49890181 Mon Sep 17 00:00:00 2001 From: dawn Date: Tue, 21 Jul 2026 15:54:45 +0300 Subject: [PATCH] eventconsumer: add 90s timeout for conns that somehow become half-open Signed-off-by: dawn --- eventconsumer/consumer.go | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/eventconsumer/consumer.go b/eventconsumer/consumer.go index 08182f5a..7d5d4c4c 100644 --- a/eventconsumer/consumer.go +++ b/eventconsumer/consumer.go @@ -18,6 +18,12 @@ import ( type ProcessFunc func(ctx context.Context, source Source, event eventstream.Event) error +// server sends a ping every 30s, so any silence longer than this means the +// connection is half-open, dropped without a close frame. +// so without a read deadline, ReadMessage only notices when the kernel's +// tcp keepalive gives up. +const livenessTimeout = 90 * time.Second + type ConsumerConfig struct { Sources map[Source]struct{} ProcessFunc ProcessFunc @@ -325,6 +331,18 @@ func (c *Consumer) runConnection(ctx context.Context, source Source) error { c.logger.Info("connected", "source", source) + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + conn.SetPongHandler(func(string) error { + return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + }) + conn.SetPingHandler(func(appData string) error { + err := conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(10*time.Second)) + if err != nil { + return err + } + return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + }) + for { select { case <-ctx.Done(): @@ -337,6 +355,7 @@ func (c *Consumer) runConnection(ctx context.Context, source Source) error { if msgType != websocket.TextMessage { continue } + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) select { case c.jobQueue <- job{source: source, message: msg}: case <-ctx.Done(): -- 2.51.2