diff --git a/internal/api/handlers_logstream_test.go b/internal/api/handlers_logstream_test.go index 6256b72..f62fd30 100644 --- a/internal/api/handlers_logstream_test.go +++ b/internal/api/handlers_logstream_test.go @@ -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") + } +} diff --git a/internal/api/logstream.go b/internal/api/logstream.go index c4cb582..ba1e547 100644 --- a/internal/api/logstream.go +++ b/internal/api/logstream.go @@ -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