Clear the SSE write deadline so an agent turn can outlive WriteTimeout (#87)
Build image / build-and-push (push) Successful in 11s
Build image / build-and-push (push) Successful in 11s
Closes #78. The 30s server WriteTimeout is an absolute deadline that was cutting every agent turn over 30s mid-stream — and the keep-alive from #73 could never work because it ticked into a connection destroyed at 30s. Refresh a per-write deadline instead: unbounded stream, bounded writes, so a stuck reader still can't pin the run goroutine. Two client-side regression tests, one per failure mode. Gadfly: no material issues on final review.
This commit was merged in pull request #87.
This commit is contained in:
+44
-1
@@ -108,6 +108,21 @@ func (h *handlers) agentChat(c *gin.Context) {
|
|||||||
send(chatEvent{Done: turn})
|
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.
|
// eventStream serializes writes to one SSE response.
|
||||||
//
|
//
|
||||||
// The mutex is load-bearing, not decoration: step events are sent from the
|
// 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.
|
// frames long before it crashes anything.
|
||||||
type eventStream struct {
|
type eventStream struct {
|
||||||
c *gin.Context
|
c *gin.Context
|
||||||
|
rc *http.ResponseController
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -124,12 +140,34 @@ type eventStream struct {
|
|||||||
// Headers go out before the first write and the stream is flushed immediately,
|
// 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
|
// so a proxy holding the response until it looks complete can't reintroduce
|
||||||
// exactly the silence streaming exists to remove.
|
// 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 {
|
func openEventStream(c *gin.Context) *eventStream {
|
||||||
c.Header("Content-Type", "text/event-stream")
|
c.Header("Content-Type", "text/event-stream")
|
||||||
c.Header("Cache-Control", "no-cache")
|
c.Header("Cache-Control", "no-cache")
|
||||||
c.Header("X-Accel-Buffering", "no")
|
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()
|
c.Writer.Flush()
|
||||||
return &eventStream{c: c}
|
return s
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *eventStream) send(ev chatEvent) {
|
func (s *eventStream) send(ev chatEvent) {
|
||||||
@@ -144,6 +182,11 @@ func (s *eventStream) send(ev chatEvent) {
|
|||||||
func (s *eventStream) write(frame string) {
|
func (s *eventStream) write(frame string) {
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
defer s.mu.Unlock()
|
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)
|
_, _ = io.WriteString(s.c.Writer, frame)
|
||||||
s.c.Writer.Flush()
|
s.c.Writer.Flush()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user