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).
This commit is contained in:
2 files changed
+81
No files matched your search
@@ -673,3 +673,71 @@ func TestRelayLogStreamWriteDeadlineSeversStalledReader(t *testing.T) {
|
|||||||
t.Fatal("relay returned without closing the source — upstream pod-log follow leaked")
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -110,6 +110,13 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
|
|||||||
// per-write deadline set inside writeChunk is what severs an unresponsive client;
|
// 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).
|
// the initial header flush below stays a plain best-effort flush (no deadline).
|
||||||
rc := http.NewResponseController(w)
|
rc := http.NewResponseController(w)
|
||||||
|
// Clear any per-write deadline on return. Server.WriteTimeout is deliberately UNSET
|
||||||
|
// (cmd/felis api.go — a WriteTimeout would sever a healthy long SSE stream), and
|
||||||
|
// with it unset net/http never resets the write deadline between keep-alive
|
||||||
|
// requests. So a deadline left set by the last writeChunk would leak onto the NEXT
|
||||||
|
// request that reuses this pooled connection and fail its first write for no reason.
|
||||||
|
// The zero time clears it; best-effort, a no-op on writers without deadline support.
|
||||||
|
defer func() { _ = rc.SetWriteDeadline(time.Time{}) }()
|
||||||
|
|
||||||
h := w.Header()
|
h := w.Header()
|
||||||
h.Set("Content-Type", "text/event-stream")
|
h.Set("Content-Type", "text/event-stream")
|
||||||
@@ -118,6 +125,12 @@ func relayLogStream(w http.ResponseWriter, r *http.Request, src io.ReadCloser) {
|
|||||||
// Defeat proxy buffering (nginx / ingress) so events arrive promptly.
|
// Defeat proxy buffering (nginx / ingress) so events arrive promptly.
|
||||||
h.Set("X-Accel-Buffering", "no")
|
h.Set("X-Accel-Buffering", "no")
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
|
// Best-effort header flush, deliberately WITHOUT a write deadline. A client that
|
||||||
|
// stalls its receive window BEFORE these headers drain is therefore NOT severed by
|
||||||
|
// writeTimeout at connect time — routing this flush through the deadline guard would
|
||||||
|
// break the guard's test specificity, and the connect-time stall is already bounded
|
||||||
|
// by the per-principal stream cap (#44). Only the mid-stream stall (every writeChunk
|
||||||
|
// below) is CLOSED by the deadline guard, not merely bounded.
|
||||||
flusher.Flush()
|
flusher.Flush()
|
||||||
|
|
||||||
// bufio.Scanner.Scan blocks until a line arrives, so to interleave a periodic
|
// bufio.Scanner.Scan blocks until a line arrives, so to interleave a periodic
|
||||||
|
|||||||
Reference in new issue
Block a user