diff --git a/internal/api/agent.go b/internal/api/agent.go index 9b45989..3eaee6d 100644 --- a/internal/api/agent.go +++ b/internal/api/agent.go @@ -108,6 +108,21 @@ func (h *handlers) agentChat(c *gin.Context) { send(chatEvent{Done: turn}) } +// sseWriteTimeout bounds ONE write to the stream, not the stream itself. +// +// It is refreshed per frame, which is the only shape that satisfies both ends: +// the server's absolute WriteTimeout would cut a long turn (#78), while removing +// the deadline entirely would let a client that stops reading block a write +// forever once the socket buffer fills — pinning the run goroutine and this +// stream's mutex with it, and taking the keep-alive down too since it needs the +// same lock. Generous, because it is a backstop against a stuck peer and not a +// pacing mechanism. +// +// A var, not a const, ONLY so the test can shrink it to prove the deadline is +// refreshed per frame rather than set once — a set-once 30s deadline would pass +// a test whose whole run is under a second. Production never reassigns it. +var sseWriteTimeout = 30 * time.Second + // eventStream serializes writes to one SSE response. // // The mutex is load-bearing, not decoration: step events are sent from the @@ -116,6 +131,7 @@ func (h *handlers) agentChat(c *gin.Context) { // frames long before it crashes anything. type eventStream struct { c *gin.Context + rc *http.ResponseController mu sync.Mutex } @@ -124,12 +140,34 @@ type eventStream struct { // Headers go out before the first write and the stream is flushed immediately, // so a proxy holding the response until it looks complete can't reintroduce // exactly the silence streaming exists to remove. +// +// Taking the write deadline off the server's absolute WriteTimeout and onto a +// per-write one is what makes a turn longer than 30s possible at all (#78). +// WriteTimeout is an ABSOLUTE deadline from when the request header was read, +// not an idle timeout, so a streaming response is cut mid-turn however recently +// it wrote. Without this the 4-minute runTimeout is unreachable and the +// keep-alive below tops out at one tick — pacing a connection that is destroyed +// underneath it. +// +// That failure is INVISIBLE from in here: writes past the deadline return +// err == nil and their bytes are dropped, so there is nothing to detect on the +// write path. Only the client sees it, as a truncated stream it reports as a +// dropped connection. Hence a deadline set up front and refreshed per frame, +// rather than anything checked after the fact. func openEventStream(c *gin.Context) *eventStream { c.Header("Content-Type", "text/event-stream") c.Header("Cache-Control", "no-cache") c.Header("X-Accel-Buffering", "no") + s := &eventStream{c: c, rc: http.NewResponseController(c.Writer)} + // Probe once here rather than reporting per frame: a writer that can't take + // deadlines will fail identically on every write, and the operator needs to + // hear it once. If this fails the stream still works — it is just back to + // being cut at WriteTimeout, which is worth saying out loud. + if err := s.rc.SetWriteDeadline(time.Now().Add(sseWriteTimeout)); err != nil { + slog.Error("api: SSE write deadlines unavailable; long turns will be truncated at the server WriteTimeout", "error", err) + } c.Writer.Flush() - return &eventStream{c: c} + return s } func (s *eventStream) send(ev chatEvent) { @@ -144,6 +182,11 @@ func (s *eventStream) send(ev chatEvent) { func (s *eventStream) write(frame string) { s.mu.Lock() defer s.mu.Unlock() + // Refresh for THIS write, so the stream as a whole is unbounded but no single + // write is. Error deliberately unchecked: openEventStream already reported + // whether deadlines work at all, and this call can only fail the same way, so + // checking here would log once per frame to say the same thing. + _ = s.rc.SetWriteDeadline(time.Now().Add(sseWriteTimeout)) _, _ = io.WriteString(s.c.Writer, frame) s.c.Writer.Flush() } diff --git a/internal/api/sse_deadline_test.go b/internal/api/sse_deadline_test.go new file mode 100644 index 0000000..181bf51 --- /dev/null +++ b/internal/api/sse_deadline_test.go @@ -0,0 +1,101 @@ +package api + +import ( + "bufio" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/gin-gonic/gin" +) + +// streamFrames spins up a real http.Server with the given WriteTimeout and an +// SSE handler that emits `frames` data frames, one every `tick`, then returns. +// It reports how many frames the client actually received and any read error — +// the only vantage point from which the deadline failures in #78/#87 are +// visible, since the writes themselves return nil when the bytes are dropped. +func streamFrames(t *testing.T, serverWriteTimeout, tick time.Duration, frames int) (int, error) { + t.Helper() + gin.SetMode(gin.TestMode) + r := gin.New() + r.GET("/stream", func(c *gin.Context) { + s := openEventStream(c) + for i := 0; i < frames; i++ { + time.Sleep(tick) + s.send(chatEvent{Error: "frame"}) + } + }) + + srv := httptest.NewUnstartedServer(r) + srv.Config.WriteTimeout = serverWriteTimeout + srv.Start() + defer srv.Close() + + resp, err := srv.Client().Get(srv.URL + "/stream") + if err != nil { + t.Fatalf("get: %v", err) + } + defer resp.Body.Close() + + got := 0 + sc := bufio.NewScanner(resp.Body) + for sc.Scan() { + if strings.HasPrefix(sc.Text(), "data: ") { + got++ + } + } + return got, sc.Err() +} + +// TestEventStreamOutlivesServerWriteTimeout is the regression test for #78. +// +// http.Server.WriteTimeout is an ABSOLUTE deadline measured from when the +// request header was read — not an idle timeout — so a streaming response is cut +// once it passes, however recently the handler wrote. pansy sets it to 30s while +// an agent turn may run for minutes. openEventStream must override it. +// +// This has to be asserted from the CLIENT side because the failure cannot be +// observed from the handler: writes made after the deadline return err == nil +// and their bytes are silently discarded. A test that checked the return of +// io.WriteString would pass against the bug. +func TestEventStreamOutlivesServerWriteTimeout(t *testing.T) { + // sseWriteTimeout stays at its 30s default here, so the per-frame refresh + // keeps the stream alive with a huge margin — CI slowness only ever makes + // this pass more surely. The server's 300ms WriteTimeout is the thing being + // overridden; frames straddle it (300ms/600ms/900ms). + got, err := streamFrames(t, 300*time.Millisecond, 300*time.Millisecond, 3) + if err != nil { + t.Errorf("client read error after %d/3 frames: %v", got, err) + } + if got != 3 { + t.Errorf("client received %d frames, want 3 — the stream was cut at the server WriteTimeout", got) + } +} + +// TestEventStreamRefreshesDeadlinePerFrame guards the #87 fix specifically: the +// write deadline is refreshed on EVERY frame, not set once. +// +// A set-once deadline is a plausible "simplification" and it reintroduces the +// unbounded-block risk the per-frame refresh exists to prevent — yet it would +// sail through the test above, whose whole run is far under sseWriteTimeout. So +// shrink sseWriteTimeout below the stream's total duration and send frames whose +// gap stays comfortably under it: per-frame refresh delivers them all, while a +// deadline set once at open would expire mid-stream and cut it short. +func TestEventStreamRefreshesDeadlinePerFrame(t *testing.T) { + orig := sseWriteTimeout + sseWriteTimeout = 400 * time.Millisecond + t.Cleanup(func() { sseWriteTimeout = orig }) + + // The server WriteTimeout is generous (5s), so it isn't the limiter — the + // per-frame sseWriteTimeout is. 8 frames at a 100ms tick span 800ms, well past + // the 400ms deadline, but each 100ms gap is a 4× margin under it. + got, err := streamFrames(t, 5*time.Second, 100*time.Millisecond, 8) + if err != nil { + t.Errorf("client read error after %d/8 frames: %v", got, err) + } + if got != 8 { + t.Errorf("client received %d frames, want 8 — a set-once deadline would cut the stream at ~%v; the refresh must be per-frame", + got, sseWriteTimeout) + } +}