Unverified Commit d6e31896 authored by Minseong Choi's avatar Minseong Choi 💬
Browse files

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.
parent 3c1d6474
Loading
Loading
Loading
Loading
+145 −0
Changes for internal/api/handlers_logstream_test.go: 145 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -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")
	}
}
+52 −9
Changes for internal/api/logstream.go: 52 added lines, 9 removed lines.
Original line number Diff line number Diff line
@@ -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