Files
Felis/internal/api/handlers_logstream_test.go
T
flyemoji 8f41a003b0 fix(api): clear the SSE write deadline on return so it can't leak to a reused connection
The per-write deadline that severs a stalled SSE reader was never cleared on
return. Server.WriteTimeout is deliberately unset -- a WriteTimeout would sever
a healthy long-lived stream -- and with it unset net/http never resets the
connection write deadline between keep-alive requests. So the deadline the last
writeChunk left set leaks onto the next request that reuses the pooled
connection and fails its first write for no reason. Clear it to the zero value
on return via a deferred rc.SetWriteDeadline; best-effort, a no-op on writers
without deadline support.

Also record honestly at the header flush that the connect-time stall stays
bounded only by the per-principal stream cap, not severed by this guard -- only
the mid-stream stall is closed. Adds a test pinning the clear (fails closed:
neutering the deferred clear leaves a +writeTimeout deadline set on return).
2026-07-01 23:06:26 +09:00

744 lines
28 KiB
Go

package api
import (
"context"
"io"
"net/http"
"net/http/httptest"
"os"
"strings"
"sync"
"testing"
"time"
)
// fakeLogStreamer drives the §8 read-side handler without a cluster: it returns a
// canned source (or error) and records what it was asked to stream. The real
// K8sLogStreamer's pod selection and pods/log follow are integration-only, so the
// handler is tested against this fake (spec §8 读=pods/log follow).
type fakeLogStreamer struct {
src io.ReadCloser
err error
calls int
gotName string
gotCtx context.Context
// srcFromCtx, when set, builds the returned source from the context
// StreamLogs actually receives, instead of returning the pre-baked src. The
// teardown test uses it to prove the handler threads r.Context() through:
// the real K8sLogStreamer opens the pods/log follow under the passed ctx, so
// a source whose only cancellation path is that ctx faithfully models the
// leak surface. A handler that passed context.Background() would hand this a
// never-cancelled context and hang.
srcFromCtx func(ctx context.Context) io.ReadCloser
}
func (f *fakeLogStreamer) StreamLogs(ctx context.Context, name string) (io.ReadCloser, error) {
f.calls++
f.gotName = name
f.gotCtx = ctx
if f.err != nil {
return nil, f.err
}
if f.srcFromCtx != nil {
return f.srcFromCtx(ctx), nil
}
return f.src, nil
}
// recordReadCloser is a finite log source that records whether Close ran, so the
// relay's teardown (deferred src.Close) can be asserted.
type recordReadCloser struct {
r *strings.Reader
closed bool
}
func (s *recordReadCloser) Read(p []byte) (int, error) { return s.r.Read(p) }
func (s *recordReadCloser) Close() error { s.closed = true; return nil }
// ctxBlockingReadCloser emits one line, then blocks until its context is
// cancelled — modeling a live `pods/log` follow stream that yields boot output
// and then waits. It is how the disconnect-teardown test proves a client
// disconnect (request-context cancel) unblocks the relay and releases the stream.
type ctxBlockingReadCloser struct {
ctx context.Context
first []byte
firstRead chan struct{}
sentFirst bool
closed chan struct{}
}
func (b *ctxBlockingReadCloser) Read(p []byte) (int, error) {
if !b.sentFirst {
b.sentFirst = true
n := copy(p, b.first)
close(b.firstRead) // signal the relay has begun streaming
return n, nil
}
// No more data until the request context is cancelled (client disconnect).
<-b.ctx.Done()
return 0, b.ctx.Err()
}
func (b *ctxBlockingReadCloser) Close() error {
close(b.closed)
return nil
}
// TestServerConsoleStream exercises the §8 read-side SSE relay end-to-end through
// the external app-tier router: ownership (owner/admin), the nil-Logs and
// open-stream failure surface, and the SSE framing itself. The real pod-log
// follow is integration-only; this drives the handler against a fakeLogStreamer.
func TestServerConsoleStream(t *testing.T) {
owner := &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}
// mk builds an API whose "survival" server is owned by owner1, with a fresh
// fakeLogStreamer wired. Subtests override src/err as needed.
mk := func() (*API, *fakeRepo, *fakeLogStreamer) {
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
streamer := &fakeLogStreamer{}
api := newTestAPI(repo, newFakeCluster())
api.Logs = streamer
return api, repo, streamer
}
t.Run("owner attaches -> 200 SSE framing + flush + audit", func(t *testing.T) {
api, repo, streamer := mk()
streamer.src = &recordReadCloser{r: strings.NewReader("line one\nline two\n")}
api.External = staticExternal{p: owner}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusOK {
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
}
if ct := w.Header().Get("Content-Type"); ct != "text/event-stream" {
t.Fatalf("Content-Type = %q, want text/event-stream", ct)
}
if cc := w.Header().Get("Cache-Control"); cc != "no-cache" {
t.Fatalf("Cache-Control = %q, want no-cache", cc)
}
// Each log line is framed as a single SSE data event.
for _, want := range []string{"data: line one\n\n", "data: line two\n\n"} {
if !strings.Contains(w.Body.String(), want) {
t.Fatalf("body missing %q; got %q", want, w.Body.String())
}
}
if !w.Flushed {
t.Fatal("SSE relay must flush per event (recorder not flushed)")
}
if !streamer.src.(*recordReadCloser).closed {
t.Fatal("relay must Close the log source when the stream ends")
}
if streamer.gotName != "survival" {
t.Fatalf("streamer saw name %q, want survival", streamer.gotName)
}
if len(repo.audits) != 1 || repo.audits[0].Action != "console.attach" || repo.audits[0].Actor != "[email protected]" {
t.Fatalf("audit not written as expected: %+v", repo.audits)
}
})
t.Run("admin attaches to another's server -> 200", func(t *testing.T) {
api, _, streamer := mk()
streamer.src = &recordReadCloser{r: strings.NewReader("boot\n")}
api.External = staticExternal{p: &Principal{UserID: "admin1", Email: "[email protected]",
Role: "admin", ViaAdminAccess: true}}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusOK {
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
}
if streamer.calls != 1 {
t.Fatalf("admin attach reached the streamer %d times, want 1", streamer.calls)
}
})
t.Run("non-owner -> 403, no stream opened", func(t *testing.T) {
api, _, streamer := mk()
api.External = staticExternal{p: &Principal{UserID: "stranger", Role: "user"}}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusForbidden {
t.Fatalf("code = %d, want 403", w.Code)
}
if streamer.calls != 0 {
t.Fatal("a forbidden caller must not open a log stream")
}
})
t.Run("unknown server -> 404", func(t *testing.T) {
api, _, streamer := mk()
api.External = staticExternal{p: owner}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/missing/console", "", nil)
if w.Code != http.StatusNotFound {
t.Fatalf("code = %d, want 404", w.Code)
}
if streamer.calls != 0 {
t.Fatal("an unknown server must not open a log stream")
}
})
t.Run("nil Logs -> 503 console_unavailable", func(t *testing.T) {
api, _, _ := mk()
api.Logs = nil
api.External = staticExternal{p: owner}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusServiceUnavailable || decodeErr(t, w) != "console_unavailable" {
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
}
})
t.Run("no running pod -> 409 not_running, no audit", func(t *testing.T) {
api, repo, streamer := mk()
streamer.err = ErrNotFound
api.External = staticExternal{p: owner}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusConflict || decodeErr(t, w) != "not_running" {
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
}
if len(repo.audits) != 0 {
t.Fatalf("a failed attach must not audit: %+v", repo.audits)
}
})
t.Run("stream unavailable -> 503 console_unavailable", func(t *testing.T) {
api, _, streamer := mk()
streamer.err = ErrConsoleUnavailable
api.External = staticExternal{p: owner}
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusServiceUnavailable || decodeErr(t, w) != "console_unavailable" {
t.Fatalf("code = %d body %s", w.Code, w.Body.String())
}
})
}
// TestServerConsoleStreamPerPrincipalCap pins the audit-hardening bound: a single
// principal may hold at most MaxStreamsPerPrincipal concurrent SSE streams, and an
// attach past that is shed with 429 too_many_streams BEFORE any upstream follow is
// opened. (Honest scope: this bounds the blast radius of the stalled-stream leak, it
// does NOT close the leak — the per-write deadline that severs a stalled stream is a
// separate slice.) With the principal already at its one-stream cap the next attach is
// refused without reaching the streamer; releasing the held slot lets an identical
// attach through, proving the 429 was the cap and not something else.
func TestServerConsoleStreamPerPrincipalCap(t *testing.T) {
owner := &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
api := newTestAPI(repo, newFakeCluster())
api.MaxStreamsPerPrincipal = 1
streamer := &fakeLogStreamer{}
api.Logs = streamer
api.External = staticExternal{p: owner}
// Occupy the principal's one stream slot, mimicking a live attach in flight. This
// lazily builds the same one-slot limiter the handler consults.
release, ok := api.streamGate().acquire(streamKey(owner))
if !ok {
t.Fatal("could not acquire the sole stream slot in test setup")
}
// A second concurrent attach is shed with 429 and never reaches the streamer.
w := do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusTooManyRequests {
t.Fatalf("over-cap attach: code = %d, want 429 (%s)", w.Code, w.Body.String())
}
if code := decodeErr(t, w); code != "too_many_streams" {
t.Fatalf("over-cap attach: error code = %q, want too_many_streams", code)
}
if streamer.calls != 0 {
t.Fatalf("an over-cap attach must not open an upstream stream (streamer.calls = %d)", streamer.calls)
}
// Releasing the held slot lets an identical attach through: the finite source EOFs,
// so the relay returns immediately with the SSE framing.
release()
streamer.src = &recordReadCloser{r: strings.NewReader("boot\n")}
w = do(api.ExternalHandler(), "GET", "/api/v1/servers/survival/console", "", nil)
if w.Code != http.StatusOK {
t.Fatalf("after releasing the slot: code = %d, want 200 (%s)", w.Code, w.Body.String())
}
if streamer.calls != 1 {
t.Fatalf("after release the attach should reach the streamer once, got %d", streamer.calls)
}
}
// idleReadCloser is a perfectly quiet log follow: every Read blocks until the
// context is cancelled, yielding no line at all. It models a Minecraft server
// that has booted and gone silent (no chat, no log output), which is exactly the
// case the SSE heartbeat exists to keep alive.
type idleReadCloser struct {
ctx context.Context
closed chan struct{}
}
func (b *idleReadCloser) Read(p []byte) (int, error) {
<-b.ctx.Done()
return 0, b.ctx.Err()
}
func (b *idleReadCloser) Close() error {
if b.closed != nil {
close(b.closed)
}
return nil
}
// signalWriter wraps the recorder so the test learns the moment the relay makes
// its first body write WITHOUT racing on the recorder's buffer: on an idle stream
// that first write can only be a heartbeat (no data line will ever arrive), so
// closing `fired` there lets the test cancel deterministically instead of sleeping
// a guessed interval. Flush is forwarded because relayLogStream type-asserts
// http.Flusher and flushes every event.
type signalWriter struct {
http.ResponseWriter
fired chan struct{}
once sync.Once
}
func (s *signalWriter) Write(p []byte) (int, error) {
s.once.Do(func() { close(s.fired) })
return s.ResponseWriter.Write(p)
}
func (s *signalWriter) Flush() {
if f, ok := s.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
// TestServerConsoleHeartbeat proves the §8 relay keeps an idle stream warm: with
// no log line ever arriving, it must still emit SSE comment heartbeats (so an
// intermediary's idle timeout — Cloudflare's ~100s — never tears the console
// down) and must emit NO data event. heartbeatInterval is shrunk so a heartbeat
// lands within the test window; the first body write is necessarily that
// heartbeat, so the test waits for it rather than sleeping a fixed duration.
func TestServerConsoleHeartbeat(t *testing.T) {
old := heartbeatInterval
heartbeatInterval = 2 * time.Millisecond
defer func() { heartbeatInterval = old }()
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
api := newTestAPI(repo, newFakeCluster())
api.External = staticExternal{p: &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}}
ctx, cancel := context.WithCancel(context.Background())
streamer := &fakeLogStreamer{srcFromCtx: func(streamCtx context.Context) io.ReadCloser {
return &idleReadCloser{ctx: streamCtx}
}}
api.Logs = streamer
rec := httptest.NewRecorder()
w := &signalWriter{ResponseWriter: rec, fired: make(chan struct{})}
r := httptest.NewRequest("GET", "/api/v1/servers/survival/console", nil).WithContext(ctx)
done := make(chan struct{})
go func() {
api.ExternalHandler().ServeHTTP(w, r)
close(done)
}()
// Wait for the relay's first body write — on an idle stream, a heartbeat — then
// disconnect. No sleep-on-a-guessed-interval.
select {
case <-w.fired:
case <-time.After(2 * time.Second):
t.Fatal("idle relay never emitted a heartbeat")
}
cancel()
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("handler did not return after client disconnect")
}
// Read the body only after <-done, so the relay's writes happen-before this.
body := rec.Body.String()
if !strings.Contains(body, sseHeartbeat) {
t.Fatalf("idle relay emitted no SSE heartbeat comment; body = %q", body)
}
// An idle stream carries keep-alives only — never a data event.
if strings.Contains(body, "data:") {
t.Fatalf("idle relay must emit only heartbeats, got a data event: %q", body)
}
}
// interleaveWriter signals once the relay has written BOTH a data event and a
// heartbeat, so a test can prove the two coexist on one stream without racing on
// the recorder. Each relay event is a single Write (fmt.Fprintf / io.WriteString),
// and the relay is the lone writer, so inspecting p per Write is race-free.
type interleaveWriter struct {
http.ResponseWriter
sawData bool
sawBeat bool
both chan struct{}
once sync.Once
}
func (b *interleaveWriter) Write(p []byte) (int, error) {
s := string(p)
if strings.HasPrefix(s, "data:") {
b.sawData = true
}
if s == sseHeartbeat {
b.sawBeat = true
}
if b.sawData && b.sawBeat {
b.once.Do(func() { close(b.both) })
}
return b.ResponseWriter.Write(p)
}
func (b *interleaveWriter) Flush() {
if f, ok := b.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
// TestServerConsoleHeartbeatInterleavesWithData proves a heartbeat following a
// real log line does not corrupt it: the source emits one line then goes idle, so
// the relay writes a data event and then — past the (shrunk) interval — a
// keep-alive. The select serializes the two writes, so the data event must appear
// intact (no comment spliced into the middle of `data: …\n\n`).
func TestServerConsoleHeartbeatInterleavesWithData(t *testing.T) {
old := heartbeatInterval
heartbeatInterval = 2 * time.Millisecond
defer func() { heartbeatInterval = old }()
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
api := newTestAPI(repo, newFakeCluster())
api.External = staticExternal{p: &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}}
ctx, cancel := context.WithCancel(context.Background())
streamer := &fakeLogStreamer{srcFromCtx: func(streamCtx context.Context) io.ReadCloser {
return &ctxBlockingReadCloser{
ctx: streamCtx,
first: []byte("boot line\n"),
firstRead: make(chan struct{}),
closed: make(chan struct{}),
}
}}
api.Logs = streamer
rec := httptest.NewRecorder()
w := &interleaveWriter{ResponseWriter: rec, both: make(chan struct{})}
r := httptest.NewRequest("GET", "/api/v1/servers/survival/console", nil).WithContext(ctx)
done := make(chan struct{})
go func() {
api.ExternalHandler().ServeHTTP(w, r)
close(done)
}()
// Wait until the relay has emitted the data event AND a heartbeat after it.
select {
case <-w.both:
case <-time.After(2 * time.Second):
t.Fatal("relay did not produce both a data event and a heartbeat")
}
cancel()
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("handler did not return after disconnect")
}
body := rec.Body.String()
// The data event must appear intact — a heartbeat must not splice into it.
if !strings.Contains(body, "data: boot line\n\n") {
t.Fatalf("data event not intact in interleaved stream; body = %q", body)
}
if !strings.Contains(body, sseHeartbeat) {
t.Fatalf("no heartbeat after the idle gap; body = %q", body)
}
}
// TestServerConsoleDisconnectTeardown proves the linchpin of the read-side relay:
// a client disconnect (request-context cancel) unblocks the follow stream and the
// deferred Close releases it — no leaked apiserver connection. The source emits
// one line, then blocks on its context; cancelling the request must make the
// handler return AND Close the source.
func TestServerConsoleDisconnectTeardown(t *testing.T) {
repo := newFakeRepo()
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner1"}
api := newTestAPI(repo, newFakeCluster())
api.External = staticExternal{p: &Principal{UserID: "owner1", Email: "[email protected]", Role: "user"}}
ctx, cancel := context.WithCancel(context.Background())
// The channels are owned by the test, but the source itself is built inside
// StreamLogs from the context the handler passes in — NOT from the test's ctx.
// This is what actually backs the teardown claim: the source's only
// cancellation path is the context the handler threaded through, mirroring the
// real K8sLogStreamer (GetLogs(...).Stream(ctx)). If handleServerConsole
// streamed under context.Background() instead of r.Context(), this source's
// second Read would block on a never-cancelled context, the handler would
// never return, and <-done below would time out — the test fails closed on the
// exact leak it exists to prevent.
firstRead := make(chan struct{})
closed := make(chan struct{})
streamer := &fakeLogStreamer{srcFromCtx: func(streamCtx context.Context) io.ReadCloser {
return &ctxBlockingReadCloser{
ctx: streamCtx,
first: []byte("boot progress 50%\n"),
firstRead: firstRead,
closed: closed,
}
}}
api.Logs = streamer
r := httptest.NewRequest("GET", "/api/v1/servers/survival/console", nil).WithContext(ctx)
w := httptest.NewRecorder()
done := make(chan struct{})
go func() {
api.ExternalHandler().ServeHTTP(w, r)
close(done)
}()
// Wait until the relay has streamed the first line and is blocked on the next
// read, then simulate the client going away.
select {
case <-firstRead:
case <-time.After(2 * time.Second):
t.Fatal("relay never read the first log line")
}
cancel()
// The handler must return promptly...
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("handler did not return after client disconnect (leaked stream)")
}
// ...and the deferred Close must have released the upstream stream.
select {
case <-closed:
case <-time.After(2 * time.Second):
t.Fatal("relay did not Close the log source on disconnect")
}
// Belt-and-suspenders: the context StreamLogs received must itself be Done.
// Read after <-done, so the write inside StreamLogs is happens-before this.
// Together with the source-from-ctx wiring above, this nails the one
// production-critical property — the relay streams under the request context.
select {
case <-streamer.gotCtx.Done():
default:
t.Fatal("StreamLogs did not receive the cancellable request context (gotCtx not Done)")
}
if !strings.Contains(w.Body.String(), "data: boot progress 50%\n\n") {
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")
}
}
// deadlineRecordWriter records the write deadlines the relay sets and never blocks on
// flush — a healthy client whose stream simply ends. It pins the deadline-CLEAR half of
// the leak guard: SetWriteDeadline stores every value, so the test can read back the
// LAST one the relay left behind after it returns. sawPositive proves a real per-write
// deadline was applied during streaming, so a broken fix that never sets a deadline at
// all cannot pass the clear-check by leaving the field zero throughout.
type deadlineRecordWriter struct {
mu sync.Mutex
hdr http.Header
lastDeadline time.Time
sawPositive bool
}
func newDeadlineRecordWriter() *deadlineRecordWriter {
return &deadlineRecordWriter{hdr: http.Header{}}
}
func (s *deadlineRecordWriter) Header() http.Header { return s.hdr }
func (s *deadlineRecordWriter) WriteHeader(int) {}
func (s *deadlineRecordWriter) Write(p []byte) (int, error) { return len(p), nil }
func (s *deadlineRecordWriter) Flush() {}
func (s *deadlineRecordWriter) FlushError() error { return nil } // healthy: never blocks
func (s *deadlineRecordWriter) SetWriteDeadline(t time.Time) error {
s.mu.Lock()
s.lastDeadline = t
if !t.IsZero() {
s.sawPositive = true
}
s.mu.Unlock()
return nil
}
func (s *deadlineRecordWriter) finalDeadline() (last time.Time, sawPositive bool) {
s.mu.Lock()
defer s.mu.Unlock()
return s.lastDeadline, s.sawPositive
}
// TestRelayLogStreamClearsWriteDeadlineOnReturn pins the keep-alive hygiene half of the
// §8 leak guard (audit #1 follow-up). Server.WriteTimeout is deliberately UNSET so a
// healthy long SSE stream is never severed (cmd/felis api.go), and with it unset net/http
// never resets the connection's write deadline between keep-alive requests. So the
// per-write deadline the relay sets must be CLEARED when the relay returns — otherwise it
// leaks onto the NEXT request that reuses this pooled connection and fails that request's
// first write for no reason. Here a finite source EOFs cleanly; after the relay returns
// the writer's final deadline must be the zero value, and a positive deadline must have
// been set first (so a fix that never sets a deadline at all cannot pass by leaving zero
// the whole time).
func TestRelayLogStreamClearsWriteDeadlineOnReturn(t *testing.T) {
src := &recordReadCloser{r: strings.NewReader("boot\n")}
w := newDeadlineRecordWriter()
r := httptest.NewRequest("GET", "/api/v1/servers/survival/console", nil)
relayLogStream(w, r, src)
last, sawPositive := w.finalDeadline()
if !sawPositive {
t.Fatal("relay never set a per-write deadline — the leak guard is not wired into the write path")
}
if !last.IsZero() {
t.Fatalf("relay left a write deadline of %v set on return; it must clear it to the zero value so it cannot leak onto a reused keep-alive connection", last)
}
if !src.closed {
t.Fatal("relay returned without closing the source")
}
}