From 3c1d64749f6c7ea8092b4953faebaeba232ff3b1 Mon Sep 17 00:00:00 2001 From: Minseong Choi Date: Wed, 1 Jul 2026 22:04:19 +0900 Subject: [PATCH] fix(api): cap concurrent SSE streams per principal Console and build-log relays hold a Server-Sent Event connection open for the life of a client's attachment; a stalled reader pins the relay goroutine plus its upstream kube-apiserver follow. Without a bound, one authenticated principal could open these repeatedly and accumulate leaked control-plane connections. Add a per-principal stream cap (streamLimiter) enforced before either relay opens its follow stream, returning 429 too_many_streams past the limit. cmd/felis wires it to 16; zero disables it, matching the "zero disables" idiom of the other levers. This bounds the blast radius of the stalled-stream leak; it does not close the leak itself -- the per-write deadline that severs a stalled stream is a separate change. --- cmd/felis/api.go | 4 ++ internal/api/api.go | 94 +++++++++++++++++++++++++ internal/api/handlers_logstream.go | 14 ++++ internal/api/handlers_logstream_test.go | 50 +++++++++++++ internal/api/images.go | 13 ++++ 5 files changed, 175 insertions(+) diff --git a/cmd/felis/api.go b/cmd/felis/api.go index 24c2160..7eb831d 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -173,6 +173,10 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { RootDomain: cfg.Server.RootDomain, WakeCooldown: 30 * time.Second, MaxConcurrentLogins: loginBcryptCap, + // Bound concurrent console/build-log SSE streams per principal. Generous enough + // for legitimate multi-tab / multi-server watching, while capping how many + // upstream follow connections a single caller can tie up if their streams stall. + MaxStreamsPerPrincipal: 16, } fmt.Fprintln(stderr, "felis api: external face fails closed (Access JWKS key function not configured)") diff --git a/internal/api/api.go b/internal/api/api.go index 9becb18..e351a23 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -110,6 +110,17 @@ type API struct { // one account they need. Enforced via loginLimiter in handleLogin. MaxConcurrentLogins int + // MaxStreamsPerPrincipal caps how many concurrent Server-Sent Event streams + // (console + build-log relays, spec §8) a single principal may hold open at once. + // Each relay blocks for the lifetime of a client's attachment and, under a stalled + // reader, pins a goroutine plus a kube-apiserver follow connection (relayLogStream). + // The cap does NOT fix that leak — the per-write deadline that severs a stalled + // stream is a separate slice — but it bounds the blast radius so one principal + // cannot accumulate unbounded leaked control-plane connections. Zero — the default + // — disables it (same "zero disables" idiom as the levers above); cmd/felis wires a + // positive value. Enforced via streamGate in the two relay handlers. + MaxStreamsPerPrincipal int + // Now is the clock, injectable for tests. Defaults to time.Now. Now func() time.Time @@ -121,6 +132,9 @@ type API struct { loginCapOnce sync.Once loginCap *concurrencyLimiter + + streamCapOnce sync.Once + streamCap *streamLimiter } // now returns the current time using the injected clock. @@ -160,6 +174,31 @@ func (a *API) loginLimiter() *concurrencyLimiter { return a.loginCap } +// streamGate lazily builds the per-principal SSE stream cap bound to +// MaxStreamsPerPrincipal. A zero cap yields a disabled limiter that admits every +// stream, so a deployment (or test) that leaves it unset pays nothing. +func (a *API) streamGate() *streamLimiter { + a.streamCapOnce.Do(func() { + a.streamCap = newStreamLimiter(a.MaxStreamsPerPrincipal) + }) + return a.streamCap +} + +// streamKey identifies the principal a stream slot is charged to. It prefers the +// stable user id and falls back to the email so a JWT principal without a user id is +// still bucketed by identity; an empty key (no authenticated identity, which the +// external face's auth guard already precludes) shares one bucket, which is safe +// because it is more restrictive, never less. +func streamKey(p *Principal) string { + if p == nil { + return "" + } + if p.UserID != "" { + return p.UserID + } + return p.Email +} + // apiRoute is one served HTTP route. Each face exposes its routes as a single // table (internalAPIRoutes / externalAPIRoutes) so that one declaration drives // BOTH handler construction here AND the OpenAPI parity test (openapi_test.go): @@ -567,6 +606,61 @@ func (l *concurrencyLimiter) acquire() (release func(), ok bool) { } } +// ---- per-principal stream cap ---- + +// streamLimiter bounds how many concurrent guarded sections a single KEY may hold at +// once. It backs the per-principal SSE stream cap (console + build-log relays): each +// relay blocks for the life of a client's attachment and, under a stalled reader, +// pins a goroutine plus a kube-apiserver follow connection, so an unbounded number of +// them from one principal is a control-plane connection-exhaustion vector. Unlike the +// login concurrencyLimiter (a single global semaphore), this counts per key. A +// non-positive max disables it (acquire always admits, release is a no-op), the same +// "zero disables" idiom as the other levers. +type streamLimiter struct { + mu sync.Mutex + n map[string]int + max int +} + +// newStreamLimiter builds a per-key stream cap admitting at most max concurrent +// holders per key. A non-positive max yields a disabled limiter that admits everyone. +func newStreamLimiter(max int) *streamLimiter { + return &streamLimiter{n: map[string]int{}, max: max} +} + +// acquire reserves a slot for key. It returns a release func and true on success, or +// nil and false when key already holds max slots. A non-positive max disables the cap +// (always admits, no-op release). The returned release is guarded by a sync.Once, so +// a defer that runs it exactly once — or even twice on some paths — never +// over-decrements the counter. +func (l *streamLimiter) acquire(key string) (release func(), ok bool) { + if l.max <= 0 { + return func() {}, true + } + l.mu.Lock() + if l.n[key] >= l.max { + l.mu.Unlock() + return nil, false + } + l.n[key]++ + l.mu.Unlock() + + var once sync.Once + return func() { once.Do(func() { l.release(key) }) }, true +} + +// release returns one of key's slots. The counter entry is deleted when it reaches +// zero so the map does not accumulate a permanent entry per principal ever seen. +func (l *streamLimiter) release(key string) { + l.mu.Lock() + defer l.mu.Unlock() + if l.n[key] <= 1 { + delete(l.n, key) + return + } + l.n[key]-- +} + // ---- running-server cap ---- // withinRunningCap reports whether waking info's server is allowed under the diff --git a/internal/api/handlers_logstream.go b/internal/api/handlers_logstream.go index 52d2d95..5cef7d0 100644 --- a/internal/api/handlers_logstream.go +++ b/internal/api/handlers_logstream.go @@ -63,6 +63,20 @@ func (a *API) handleServerConsole(w http.ResponseWriter, r *http.Request) { return } + // Bound concurrent SSE streams per principal BEFORE opening the follow stream, so + // an over-cap caller never even ties up a kube-apiserver connection. A stalled + // reader keeps this relay (and its upstream follow) alive indefinitely — the write + // deadline that actually severs it is a separate slice — so this cap is what stops + // one principal from accumulating unbounded leaked control-plane connections. The + // slot is held for the whole relay and released on every return path. + release, ok := a.streamGate().acquire(streamKey(p)) + if !ok { + writeError(w, r, newError(http.StatusTooManyRequests, "too_many_streams", + "too many open console streams; close one and retry")) + return + } + defer release() + // Open the follow stream. Every error must be resolved HERE, into a normal JSON // envelope, because relayLogStream commits the 200 + SSE headers and no error // body can follow it. diff --git a/internal/api/handlers_logstream_test.go b/internal/api/handlers_logstream_test.go index 269a5b4..6256b72 100644 --- a/internal/api/handlers_logstream_test.go +++ b/internal/api/handlers_logstream_test.go @@ -210,6 +210,56 @@ func TestServerConsoleStream(t *testing.T) { }) } +// 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: "owner1@example.net", 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 diff --git a/internal/api/images.go b/internal/api/images.go index fa31571..42b4a32 100644 --- a/internal/api/images.go +++ b/internal/api/images.go @@ -118,6 +118,19 @@ func (a *API) handleBuildLogs(w http.ResponseWriter, r *http.Request) { return } p := principalFromContext(r.Context()) + + // Bound concurrent SSE streams per principal (shared with the console relay): a + // stalled reader pins this relay and its upstream build-pod follow, so cap how many + // one principal may hold at once. Acquired before opening the stream and released on + // every return path. See streamGate — this bounds blast radius, not the leak itself. + release, ok := a.streamGate().acquire(streamKey(p)) + if !ok { + writeError(w, r, newError(http.StatusTooManyRequests, "too_many_streams", + "too many open build-log streams; close one and retry")) + return + } + defer release() + src, err := a.BuildLogs.StreamLogs(r.Context(), id) switch { case errors.Is(err, ErrNotFound):