diff --git a/http.go b/http.go index cae9927..e5b9604 100644 --- a/http.go +++ b/http.go @@ -34,6 +34,15 @@ import ( "go.mitchellh.com/tack/internal/buildkite" ) +// wsWriteWait bounds how long any single websocket write (frame or +// control) is allowed to block before we treat the peer as dead. A +// client that stops reading but keeps the TCP connection open would +// otherwise hang the handler indefinitely on a full kernel send buffer. +// 10s is intentionally generous: real backpressure resolves in +// milliseconds, so anything past that is a stuck peer we'd rather drop +// than serve. +const wsWriteWait = 10 * time.Second + // runHTTP starts the spindle's HTTP server and blocks until ctx is // cancelled or the listener returns a fatal error. On ctx cancellation it // performs a graceful shutdown with a bounded timeout. @@ -285,13 +294,15 @@ func logsHandler(logger *slog.Logger, provider Provider) http.HandlerFunc { defer func() { // Send a close frame on the way out so the appview proxy // sees a clean shutdown. Mirrors upstream - // spindle.(*Spindle).Logs. + // spindle.(*Spindle).Logs. WriteControl honours the + // deadline argument directly, so a stuck peer can't hang + // us here. _ = conn.WriteControl( websocket.CloseMessage, websocket.FormatCloseMessage( websocket.CloseNormalClosure, "log stream complete", ), - time.Now().Add(time.Second), + time.Now().Add(wsWriteWait), ) conn.Close() }() @@ -334,6 +345,20 @@ func logsHandler(logger *slog.Logger, provider Provider) http.HandlerFunc { ) return } + // Bound the write so a client that stopped reading + // but kept the TCP connection open can't hang us on a + // full kernel send buffer. WriteMessage doesn't take a + // deadline argument the way WriteControl does — we + // have to set it on the conn before each frame. + if err := conn.SetWriteDeadline(time.Now().Add(wsWriteWait)); err != nil { + logger.Debug("logs set write deadline failed", + "err", err, + "knot", knot, + "pipeline_rkey", pipelineRkey, + "workflow", workflow, + ) + return + } if err := conn.WriteMessage(websocket.TextMessage, frame); err != nil { logger.Debug("logs frame write failed", "err", err, @@ -440,9 +465,13 @@ func eventsHandler(logger *slog.Logger, br *broker) http.HandlerFunc { return } case <-ticker.C: + // WriteControl takes its own deadline argument, so + // the ping itself can't hang us — but we still want a + // generous-but-bounded ceiling to match the per-frame + // write timeout. if err := conn.WriteControl( websocket.PingMessage, nil, - time.Now().Add(time.Second), + time.Now().Add(wsWriteWait), ); err != nil { logger.Debug("events ping failed", "err", err) return @@ -475,6 +504,13 @@ func streamEvents(ctx context.Context, conn *websocket.Conn, st *store, cursor * if err != nil { return fmt.Errorf("marshal envelope: %w", err) } + // Bound the per-frame write so a client that stopped reading + // (but didn't close the TCP connection) can't hang the + // handler on a full kernel send buffer. WriteMessage has no + // deadline argument of its own — we set it on the conn. + if err := conn.SetWriteDeadline(time.Now().Add(wsWriteWait)); err != nil { + return fmt.Errorf("set write deadline: %w", err) + } if err := conn.WriteMessage(websocket.TextMessage, frame); err != nil { return fmt.Errorf("write frame: %w", err) }