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:
2 files changed
+197
-9
No files matched your search
@@ -5,6 +5,7 @@ import (
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -528,3 +529,147 @@ func TestServerConsoleDisconnectTeardown(t *testing.T) {
|
||||
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")
|
||||
}
|
||||
}
|
||||
@@ -3,7 +3,6 @@ package api
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"time"
|
||||
@@ -59,6 +58,23 @@ const sseHeartbeat = ": keepalive\n\n"
|
||||
// heartbeat without waiting; production never reassigns it.
|
||||
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
|
||||
// source to the client as Server-Sent Events (spec §262 SSE, NOT WebSocket). It
|
||||
// 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"))
|
||||
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.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.
|
||||
return
|
||||
}
|
||||
// One log line → one SSE "data:" event. A write error means the client
|
||||
// side is gone; stop (the deferred Close tears the upstream down too).
|
||||
if _, err := fmt.Fprintf(w, "data: %s\n\n", line); err != nil {
|
||||
// One log line → one SSE "data:" event, written under a per-write deadline
|
||||
// so a stalled reader severs the stream (writeChunk → false) instead of
|
||||
// pinning this goroutine; the deferred Close then tears the upstream down.
|
||||
if !writeChunk(rc, w, "data: "+line+"\n\n") {
|
||||
return
|
||||
}
|
||||
flusher.Flush()
|
||||
case <-ticker.C:
|
||||
// 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
|
||||
// gone; stop.
|
||||
if _, err := io.WriteString(w, sseHeartbeat); err != nil {
|
||||
// warm past the proxy idle timeout — same deadline-guarded write, so a
|
||||
// client that has gone silent-but-stalled is torn down here too.
|
||||
if !writeChunk(rc, w, sseHeartbeat) {
|
||||
return
|
||||
}
|
||||
flusher.Flush()
|
||||
case <-ctx.Done():
|
||||
// Client disconnected (request context cancelled). Return; defers run.
|
||||
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
|
||||
// mirrors the operator's pod container name (internal/operator builders —
|
||||
// containerName), duplicated here for the same reason rconEndpoint duplicates the
|
||||
|
||||
Reference in new issue
Block a user