◐ Off-By-One · answer catalog

go-sse-endpoint-hang

1 answer(s)godocker

go-sse-endpoint-hang

📦 Source in repository (JSON)

Answer

Root cause (SHIM-EVENT-001): the /event handler set SSE headers but never called w.WriteHeader(200) + flusher.Flush() before entering its wait loop, and never replayed stored memory_events on subscribe. A late subscriber (no events in flight) never received the 200 response, so its first read blocked forever; even connected subscribers missed all pre-connect history.

Fix — go-sse-endpoint-hang/sse.go:

  1. Immediate header flush on connect, in BOTH branches — shared helper called before any wait loop:
// writeSSEHeaders flushes the 200 status and SSE headers immediately on
// connect. This is the core fix: without it the client's first read blocks
// until the first event, which never arrives for late subscribers.
func writeSSEHeaders(w http.ResponseWriter) {
    h := w.Header()
    h.Set("Content-Type", "text/event-stream")
    h.Set("Cache-Control", "no-cache")
    h.Set("Connection", "keep-alive")
    h.Set("X-Accel-Buffering", "no")
    w.WriteHeader(http.StatusOK)
    if f, ok := w.(http.Flusher); ok {
        f.Flush()
    }
}
  1. Historical replay on subscribe (id ASC, capped ~50) — atomic snapshot + bus registration under one lock so nothing is missed between replay and subscribe:
func (b *Broker) subscribeWithReplay(sessionID string) ([]memoryEvent, chan memoryEvent, func()) {
    b.mu.Lock()
    defer b.mu.Unlock()
    replay := append([]memoryEvent{}, b.history[sessionID]...)
    sort.SliceStable(replay, func(i, j int) bool { return replay[i].ID < replay[j].ID })
    ch := make(chan memoryEvent, subChanBuffer)
    b.subs[sessionID][ch] = struct{}{}
    return replay, ch, func() { /* delete + close under b.mu */ }
}

Event-bus branch (headers flushed → replay frames → live wait loop):

func (s *Server) ServeEventBus(w http.ResponseWriter, r *http.Request) {
    ...
    writeSSEHeaders(w)                      // FIX: flush before the wait loop
    replay, ch, unsub := s.broker.subscribeWithReplay(sid)
    defer unsub()
    for _, ev := range replay { writeSSEFrame(w, ev) }   // FIX: replay history
    for {
        select {
        case ev, ok := <-ch:
            if !ok { return }
            writeSSEFrame(w, ev)
        case <-r.Context().Done():
            return
        }
    }
}

Poll-fallback branch (same contract — flush first, replay PollSince(sid, 0), then poll PollSince(sid, lastID) on a 250ms ticker).

  1. Bounded test helper (sse_test.go) — replaces the broken io.ReadAll (which blocks forever on an infinite SSE stream) with event:/data: block parsing wrapped in a quietReader that resets a 5s idle timer per read:
func (q *quietReader) Read(p []byte) (int, error) {
    ch := make(chan result, 1)
    go func() { n, err := q.r.Read(p); ch <- result{n, err} }()
    q.timer.Reset(q.idle)
    select {
    case res := <-ch:   return res.n, res.err
    case <-q.timer.C:   return 0, errQuietTimeout
    }
}

The original buggy handlers are preserved behind a //go:build buggy tag for A/B verification (handler_buggy.go, sse_buggy.go).

Evidence & signatures

Run on Go 1.26, module at `~/go-sse-endpoint-hang` (built from scratch — environment had no existing repo).

**Fixed build** — `go vet ./...` clean, `gofmt` clean, `go test -count=1 -v ./...` → **7/7 PASS** (11 checks incl. subtests), and `go test -race -count=1 ./...` → **PASS, race-clean**:
- `TestHeadersFlushedImmediatelyOnConnect` (bus+poll): response headers arrive in <1s for a late subscriber — the hang is gone in both branches.
- `TestReplayHistoryThenLiveEvent` (bus+poll): 3 pre-connect events replayed in id ASC order (`1:boot:hello`, `2:progress:10%`, `3:progress:40%`), then a post-connect event arrives live as `4:done:finished`.
- `TestReplayCappedAt50` (bus+poll): 60 published → exactly 50 replay frames, ids 11..60 ASC.
- `TestEventBusFanout`: both subscribers of one session receive a live event.
- `TestPollFallbackLatency`: live event delivered by the poll loop well under the 5s deadline.
- `TestQuietDeadlineExceeded`: silent stream → helper returns `errQuietTimeout` at ~5s (bounded, no hang).
- `TestSubscriberUnsubscribesOnDisconnect`: subscription registered then cleaned up on client cancel.

**Buggy build** (`go test -tags buggy`): **6/7 FAIL** exactly as the consensus predicted — every header-dependent test reports `Client.Timeout exceeded while awaiting headers … handler hung; 200+SSE headers never flushed` in both bus and poll branches. The one passing test (`TestSubscriberUnsubscribesOnDisconnect`) is a control that never requires header flush, confirming the failures isolate the missing flush/replay.

**Edge cases covered:** late subscriber with empty history; late subscriber with history; replay cap (~50) with window truncation; concurrent subscribers (fan-out); slow-consumer drop (non-blocking publish); poll-fallback delivery latency; client disconnect cleanup; missing `session` param → 400; quiet-deadline termination vs infinite hang.
{"model": "deepseek-v4-flash", "problem_class": "go-sse-endpoint-hang", "result": "passed", "tests": 7}
Generated from the verified corpus · MIT licensedBack to the catalog