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:
5 files changed
+175
No files matched your search
@@ -173,6 +173,10 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int {
|
|||||||
RootDomain: cfg.Server.RootDomain,
|
RootDomain: cfg.Server.RootDomain,
|
||||||
WakeCooldown: 30 * time.Second,
|
WakeCooldown: 30 * time.Second,
|
||||||
MaxConcurrentLogins: loginBcryptCap,
|
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)")
|
fmt.Fprintln(stderr, "felis api: external face fails closed (Access JWKS key function not configured)")
|
||||||
|
|
||||||
|
|||||||
@@ -110,6 +110,17 @@ type API struct {
|
|||||||
// one account they need. Enforced via loginLimiter in handleLogin.
|
// one account they need. Enforced via loginLimiter in handleLogin.
|
||||||
MaxConcurrentLogins int
|
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 is the clock, injectable for tests. Defaults to time.Now.
|
||||||
Now func() time.Time
|
Now func() time.Time
|
||||||
|
|
||||||
@@ -121,6 +132,9 @@ type API struct {
|
|||||||
|
|
||||||
loginCapOnce sync.Once
|
loginCapOnce sync.Once
|
||||||
loginCap *concurrencyLimiter
|
loginCap *concurrencyLimiter
|
||||||
|
|
||||||
|
streamCapOnce sync.Once
|
||||||
|
streamCap *streamLimiter
|
||||||
}
|
}
|
||||||
|
|
||||||
// now returns the current time using the injected clock.
|
// now returns the current time using the injected clock.
|
||||||
@@ -160,6 +174,31 @@ func (a *API) loginLimiter() *concurrencyLimiter {
|
|||||||
return a.loginCap
|
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
|
// apiRoute is one served HTTP route. Each face exposes its routes as a single
|
||||||
// table (internalAPIRoutes / externalAPIRoutes) so that one declaration drives
|
// table (internalAPIRoutes / externalAPIRoutes) so that one declaration drives
|
||||||
// BOTH handler construction here AND the OpenAPI parity test (openapi_test.go):
|
// 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 ----
|
// ---- running-server cap ----
|
||||||
|
|
||||||
// withinRunningCap reports whether waking info's server is allowed under the
|
// withinRunningCap reports whether waking info's server is allowed under the
|
||||||
|
|||||||
@@ -63,6 +63,20 @@ func (a *API) handleServerConsole(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
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
|
// 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
|
// envelope, because relayLogStream commits the 200 + SSE headers and no error
|
||||||
// body can follow it.
|
// body can follow it.
|
||||||
|
|||||||
@@ -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
|
// 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
|
// 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
|
// that has booted and gone silent (no chat, no log output), which is exactly the
|
||||||
|
|||||||
@@ -118,6 +118,19 @@ func (a *API) handleBuildLogs(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
p := principalFromContext(r.Context())
|
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)
|
src, err := a.BuildLogs.StreamLogs(r.Context(), id)
|
||||||
switch {
|
switch {
|
||||||
case errors.Is(err, ErrNotFound):
|
case errors.Is(err, ErrNotFound):
|
||||||
|
|||||||
Reference in new issue
Block a user