fix(api): bound SSE relay writes with a deadline to sever stalled readers

relayLogStream copied a pod-log follow to the client with a plain
flusher.Flush per event. On a client that stays connected but stops
reading (its TCP receive window shut), net/http buffers the small
"data:" line and only touches the socket at Flush, which then blocks
forever inside the write. The select's <-ctx.Done() branch is never
reached, because r.Context() cancels on an actual disconnect, not on a
stall, so the relay goroutine and its upstream apiserver follow leak for
the life of the process.

Route every event's write+flush through http.ResponseController with a
per-write deadline (writeTimeout, 30s): a stalled flush now returns
os.ErrDeadlineExceeded, the error plain http.Flusher.Flush swallows, and
the relay abandons the stream so the deferred cancel + src.Close release
the follow. SetWriteDeadline and rc.Flush are best-effort: a writer
without deadline support (httptest recorder; some HTTP/2 origins) ignores
the deadline and behaves exactly as before, so the guard degrades
gracefully.

This closes the leak the per-principal stream cap only bounded the blast
radius of. Verified by a deterministic test with a deadline-aware
ResponseWriter whose flush blocks until the deadline; the test times out
(fails closed) if the guard is removed.
This commit is contained in:
flyemoji committed 2026-07-01 22:30:40 +09:00
1 parent 3c1d64749f
commit d6e3189629
2 files changed
+197 -9

No files matched your search

+145
View File
@@ -5,6 +5,7 @@ import (
"io" "io"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"os"
"strings" "strings"
"sync" "sync"
"testing" "testing"
@@ -528,3 +529,147 @@ func TestServerConsoleDisconnectTeardown(t *testing.T) {
t.Fatalf("expected the first event before disconnect, got %q", w.Body.String()) t.Fatalf("expected the first event before disconnect, got %q", w.Body.String())
} }
} }
// lineOnceThenBlockReadCloser yields exactly one line, then blocks every later Read
// until Close is called. It models a live-but-quiet follow stream: the server printed
// one line and has since gone silent, so nothing on the SOURCE side can end the relay
// — the only thing that can is the client side (here, the write deadline firing on a
// stalled reader). Close unblocks the parked Read so the scan goroutine exits cleanly
// when relayLogStream tears down, mirroring how src.Close aborts a real pods/log Read.
type lineOnceThenBlockReadCloser struct {
line []byte
sentFirst bool // touched only by the single scan goroutine's Read
block chan struct{}
once sync.Once
}
func newLineOnceThenBlock(line string) *lineOnceThenBlockReadCloser {
return &lineOnceThenBlockReadCloser{line: []byte(line), block: make(chan struct{})}
}
func (c *lineOnceThenBlockReadCloser) Read(p []byte) (int, error) {
if !c.sentFirst {
c.sentFirst = true
return copy(p, c.line), nil
}
<-c.block
return 0, io.EOF
}
func (c *lineOnceThenBlockReadCloser) Close() error {
c.once.Do(func() { close(c.block) })
return nil
}
// closed reports whether Close ran. Reading a channel's closed-ness is race-free, so
// the test may call this from another goroutine once the relay has returned.
func (c *lineOnceThenBlockReadCloser) closed() bool {
select {
case <-c.block:
return true
default:
return false
}
}
// deadlineStallWriter models a client that connected — the header flush went out — and
// then stopped reading. Its Write buffers and returns at once (like net/http's bufio-
// backed *response, a small SSE line never touches the socket at Write); its plain
// Flush — the one-time header flush — returns immediately; but every deadline-gated
// FlushError blocks until the write deadline relayLogStream set, then reports
// os.ErrDeadlineExceeded, exactly how a real socket surfaces a SetWriteDeadline expiry
// on a stalled reader. http.NewResponseController(w).Flush() prefers FlushError over
// plain Flush, so the relay's per-event flush travels the blocking path while the
// header flush does not — which is why the guard has to route flushes through the
// ResponseController, not the bare http.Flusher whose Flush swallows the error.
type deadlineStallWriter struct {
mu sync.Mutex
hdr http.Header
deadline time.Time
sawDL bool
}
func newDeadlineStallWriter() *deadlineStallWriter {
return &deadlineStallWriter{hdr: http.Header{}}
}
func (s *deadlineStallWriter) Header() http.Header { return s.hdr }
func (s *deadlineStallWriter) WriteHeader(int) {}
func (s *deadlineStallWriter) Write(p []byte) (int, error) { return len(p), nil } // buffered: never blocks
func (s *deadlineStallWriter) Flush() {} // header flush: instant, best-effort
// FlushError is where the stalled socket bites: it blocks until the deadline the relay
// set via SetWriteDeadline, then returns the same error a real write reports when that
// deadline elapses. With no deadline set it returns nil — a healthy, instant flush.
func (s *deadlineStallWriter) FlushError() error {
s.mu.Lock()
d := s.deadline
s.mu.Unlock()
if d.IsZero() {
return nil
}
t := time.NewTimer(time.Until(d))
defer t.Stop()
<-t.C
return os.ErrDeadlineExceeded
}
func (s *deadlineStallWriter) SetWriteDeadline(t time.Time) error {
s.mu.Lock()
s.deadline = t
s.sawDL = true
s.mu.Unlock()
return nil
}
func (s *deadlineStallWriter) deadlineSet() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.sawDL
}
// TestRelayLogStreamWriteDeadlineSeversStalledReader closes the actual §8 relay leak
// (audit #1): on a client that connected but stopped reading — request context still
// live, socket write blocked — the relay must not pin its goroutine and upstream pod-log
// follow forever. The fix sets a per-write deadline (writeTimeout) via
// http.ResponseController before each event and abandons the stream when a flush exceeds
// it. deadlineStallWriter makes the flush block until that deadline then report
// os.ErrDeadlineExceeded, exactly as a stalled socket does; the source stays open and
// silent (never EOFs), so the ONLY thing that can end the relay is the deadline. Without
// the fix (a bare flusher.Flush that swallows the error), the relay would loop forever
// waiting for a line that never comes and this test would time out — it fails closed on
// the exact leak it guards.
func TestRelayLogStreamWriteDeadlineSeversStalledReader(t *testing.T) {
orig := writeTimeout
writeTimeout = 30 * time.Millisecond
defer func() { writeTimeout = orig }()
src := newLineOnceThenBlock("boot\n")
w := newDeadlineStallWriter()
// A LIVE request context: the client has NOT disconnected. r.Context() never fires
// here, which is precisely why the write deadline — not a context cancel — has to be
// what severs the stalled stream.
r := httptest.NewRequest("GET", "/api/v1/servers/survival/console", nil)
done := make(chan struct{})
go func() {
relayLogStream(w, r, src)
close(done)
}()
select {
case <-done:
// The write deadline fired and the relay tore the stalled stream down.
case <-time.After(2 * time.Second):
t.Fatal("relayLogStream did not return on a stalled-but-open reader; the write deadline never severed the stream (leak)")
}
if !w.deadlineSet() {
t.Fatal("relay never set a write deadline — the leak guard is not wired into the write path")
}
// Returning ran the deferred Close: the upstream follow (a real apiserver
// connection) is released rather than leaked.
if !src.closed() {
t.Fatal("relay returned without closing the source — upstream pod-log follow leaked")
}
}
+52 -9
View File
@@ -3,7 +3,6 @@ package api
import ( import (
"bufio" "bufio"
"context" "context"
"fmt"
"io" "io"
"net/http" "net/http"
"time" "time"
@@ -59,6 +58,23 @@ const sseHeartbeat = ": keepalive\n\n"
// heartbeat without waiting; production never reassigns it. // heartbeat without waiting; production never reassigns it.
var heartbeatInterval = 25 * time.Second var heartbeatInterval = 25 * time.Second
// writeTimeout bounds how long a single SSE write+flush to the client may block on
// the socket before the relay abandons the stream. It is the leak guard's teeth: on
// a stalled-but-open reader (client connected, its TCP receive window shut, never
// reading) the flush — where net/http actually drains the socket, since it buffers
// the small "data:" line rather than writing it through — would otherwise block
// forever INSIDE the write, with the request context never firing (r.Context()
// cancels on an actual disconnect, not on a stall). That pins this goroutine and its
// upstream pod-log follow indefinitely. relayLogStream applies this as a per-write
// deadline via http.ResponseController, so an unresponsive client is torn down
// within writeTimeout of a stalled flush instead of leaking. It sits comfortably
// above any transient slow-client write (a data line is bytes-to-KB) yet well under
// the ~100s proxy idle drop. Best-effort: writers without deadline support
// (httptest.ResponseRecorder; some HTTP/2 origins) ignore it and the relay behaves
// exactly as before. It is a var ONLY so a test can shrink it; production never
// reassigns it.
var writeTimeout = 30 * time.Second
// relayLogStream is the shared §8 read-side relay: it copies a line-oriented log // relayLogStream is the shared §8 read-side relay: it copies a line-oriented log
// source to the client as Server-Sent Events (spec §262 SSE, NOT WebSocket). It // source to the client as Server-Sent Events (spec §262 SSE, NOT WebSocket). It
// is the single reusable artifact the server console (handleServerConsole) and, // is the single reusable artifact the server console (handleServerConsole) and,
@@ -88,6 +104,12 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
"streaming is unsupported by this server")) "streaming is unsupported by this server"))
return return
} }
// rc carries the two capabilities plain http.Flusher lacks: SetWriteDeadline (to
// bound a stalled write) and a Flush whose error is observable — http.Flusher.Flush
// swallows the deadline-exceeded error that a stalled socket flush returns. The
// per-write deadline set inside writeChunk is what severs an unresponsive client;
// the initial header flush below stays a plain best-effort flush (no deadline).
rc := http.NewResponseController(w)
h := w.Header() h := w.Header()
h.Set("Content-Type", "text/event-stream") h.Set("Content-Type", "text/event-stream")
@@ -136,20 +158,19 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
// the deferred Close release the source. // the deferred Close release the source.
return return
} }
// One log line → one SSE "data:" event. A write error means the client // One log line → one SSE "data:" event, written under a per-write deadline
// side is gone; stop (the deferred Close tears the upstream down too). // so a stalled reader severs the stream (writeChunk → false) instead of
if _, err := fmt.Fprintf(w, "data: %s\n\n", line); err != nil { // pinning this goroutine; the deferred Close then tears the upstream down.
if !writeChunk(rc, w, "data: "+line+"\n\n") {
return return
} }
flusher.Flush()
case <-ticker.C: case <-ticker.C:
// No line for a whole interval: emit a comment so the connection stays // No line for a whole interval: emit a comment so the connection stays
// warm past the proxy idle timeout. A write error means the client is // warm past the proxy idle timeout — same deadline-guarded write, so a
// gone; stop. // client that has gone silent-but-stalled is torn down here too.
if _, err := io.WriteString(w, sseHeartbeat); err != nil { if !writeChunk(rc, w, sseHeartbeat) {
return return
} }
flusher.Flush()
case <-ctx.Done(): case <-ctx.Done():
// Client disconnected (request context cancelled). Return; defers run. // Client disconnected (request context cancelled). Return; defers run.
return return
@@ -157,6 +178,28 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
} }
} }
// writeChunk writes one framed SSE chunk to the client under a fresh per-write
// deadline and flushes it, returning false when the client socket is gone so the
// caller tears the relay (and its upstream follow) down. The deadline is the leak
// guard: net/http buffers the small write and only touches the socket at Flush, so a
// stalled reader blocks there — without a deadline that block is unbounded and the
// request context never fires. Both errors are honored: the write error (a line
// larger than the buffer can block mid-write) and the flush error — rc.Flush
// surfaces the os.ErrDeadlineExceeded that plain http.Flusher.Flush swallows.
// SetWriteDeadline and rc.Flush are best-effort: on a writer without deadline
// support (httptest.ResponseRecorder; some HTTP/2 origins) the deadline is ignored
// and rc.Flush reduces to a plain, non-erroring flush, so behaviour is unchanged
// where the guard cannot apply.
func writeChunk(rc *http.ResponseController, w io.Writer, chunk string) bool {
// Best-effort: an unsupported writer returns an error we ignore, leaving the
// write unbounded exactly as before the guard existed.
_ = rc.SetWriteDeadline(time.Now().Add(writeTimeout))
if _, err := io.WriteString(w, chunk); err != nil {
return false
}
return rc.Flush() == nil
}
// serverLogContainer is the container whose logs the read-side relay streams. It // serverLogContainer is the container whose logs the read-side relay streams. It
// mirrors the operator's pod container name (internal/operator builders — // mirrors the operator's pod container name (internal/operator builders —
// containerName), duplicated here for the same reason rconEndpoint duplicates the // containerName), duplicated here for the same reason rconEndpoint duplicates the