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.
This commit is contained in:
flyemoji committed 2026-07-01 22:04:19 +09:00
1 parent c6c0772a7a
commit 3c1d64749f
5 files changed
+175

No files matched your search

+94
View File
@@ -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
+14
View File
@@ -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.
+50
View File
@@ -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: "[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
+13
View File
@@ -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):