From 4ab89f9a4e9cf3afd2b785e2660c3cffbd891ff2 Mon Sep 17 00:00:00 2001 From: Aamer Akhter Date: Sat, 20 Jun 2026 09:53:52 -0400 Subject: [PATCH] COD-137 scope WS per-session limit by clientId (fix spurious 4008 on reconnect) MAX_WS_PER_SESSION was gated by a bare Map counter, incremented on upgrade and decremented only on the old socket's async close. A client that dropped and immediately reconnected could land its new upgrade before the old socket's close fired, briefly over-counting and tripping a spurious 4008 (-> HTTP fallback). The limit also counted raw sockets, so a reconnecting client consumed a new slot instead of its own. Replace the counter with WsConnectionRegistry (new pure, unit-tested module) that tracks live sockets per session keyed by clientId. A same-cid upgrade SUPERSEDES its own socket (evicts the stale one with close 4010, reuses the slot, no net count change) -> a reconnect can never be rejected by the cap. The reliable-input protocol (shouldApplyInput(cid,seq)) already assumes one logical client per cid per session, so same-cid eviction is principled, not a regression of multi-tab (which already collides on seq). Slots are freed EAGERLY on error/terminate, not just async close; close is identity-matched so a superseded socket's late close is a no-op. cid-less upgrades are admitted anonymously up to the cap and never evict (backward-compat). Client sends cid on the WS upgrade URL (?cid=, encoded, omitted if absent). Tests: ws-connection-registry.test.ts (reconnect-reclaim at cap, rejects N+1th distinct, eager-terminate frees slot, cid-less up-to-limit + no-evict, late-close-no-evict, per-session isolation) + route integration in ws-routes.test.ts (real upgrade through the cap). 45/45 across registry + ws-routes + input-send-order + ws-reconnect-plan; tsc 0, build, prettier, frontend-syntax clean. --- src/web/public/app.js | 8 +- src/web/routes/ws-routes.ts | 383 +++++++++++++++------------- src/web/ws-connection-registry.ts | 120 +++++++++ test/routes/ws-routes.test.ts | 23 ++ test/ws-connection-registry.test.ts | 112 ++++++++ 5 files changed, 471 insertions(+), 175 deletions(-) create mode 100644 src/web/ws-connection-registry.ts create mode 100644 test/ws-connection-registry.test.ts diff --git a/src/web/public/app.js b/src/web/public/app.js index ccec5458..608bd900 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -1971,7 +1971,13 @@ class CodemanApp { this._disconnectWs(); const proto = location.protocol === 'https:' ? 'wss:' : 'ws:'; - const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal`; + // Pass the stable per-browser clientId on the upgrade URL so the server's + // connection registry scopes the per-session limit by client (COD-137): a + // reconnect supersedes its own socket instead of consuming a new slot and + // tripping a spurious 4008. Omitted if clientId is unavailable (server then + // treats the upgrade as anonymous — still admitted up to the limit). + const cidQuery = this._clientId ? `?cid=${encodeURIComponent(this._clientId)}` : ''; + const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal${cidQuery}`; const ws = new WebSocket(url); this._ws = ws; this._wsSessionId = sessionId; diff --git a/src/web/routes/ws-routes.ts b/src/web/routes/ws-routes.ts index 17d7ee3e..f71a1ed4 100644 --- a/src/web/routes/ws-routes.ts +++ b/src/web/routes/ws-routes.ts @@ -34,6 +34,7 @@ import type { WebSocket } from 'ws'; import type { SessionPort } from '../ports/session-port.js'; import { MAX_INPUT_LENGTH } from '../../config/terminal-limits.js'; import { isAllowedRequestHost, isAllowedRequestOrigin, type HostPolicy } from '../network-auth-policy.js'; +import { WsConnectionRegistry } from '../ws-connection-registry.js'; /** Micro-batch interval for terminal output (ms). Short enough for low latency, * long enough to group Ink's rapid cursor-up redraw sequences into single frames. */ @@ -59,206 +60,240 @@ const DEC_2026_END = '\x1b[?2026l'; /** Max concurrent WS connections per session. Prevents listener/bandwidth multiplication. */ const MAX_WS_PER_SESSION = 5; -/** Track active WS connections per session for connection limiting. */ -const sessionWsCount = new Map(); +/** + * Track live WS connections per session, keyed by clientId (COD-137). + * Replaces a bare counter that over-counted across the async-close gap on + * reconnect (spurious 4008). A same-`cid` reconnect supersedes its own socket + * (reclaims the slot) instead of consuming a new one; cid-less upgrades are + * admitted anonymously up to the cap. See ws-connection-registry.ts. + */ +const sessionWsRegistry = new WsConnectionRegistry(MAX_WS_PER_SESSION); export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort, getHostPolicy: () => HostPolicy): void { - app.get<{ Params: { id: string } }>('/ws/sessions/:id/terminal', { websocket: true }, (socket: WebSocket, req) => { - // Reject cross-site WebSocket hijacking (CSWSH) and DNS-rebinding before doing - // anything: the upgrade must come from an allowed Host and (when the browser - // sends one — it always does for WS) a same-site Origin. Writing to this socket - // injects keystrokes into a --dangerously-skip-permissions agent, so this gate - // matters even on the default no-password install. See security review H5. - const policy = getHostPolicy(); - if (!isAllowedRequestHost(req.headers.host, policy) || !isAllowedRequestOrigin(req.headers.origin, policy)) { - socket.close(4003, 'Forbidden'); - return; - } + app.get<{ Params: { id: string }; Querystring: { cid?: string } }>( + '/ws/sessions/:id/terminal', + { websocket: true }, + (socket: WebSocket, req) => { + // Reject cross-site WebSocket hijacking (CSWSH) and DNS-rebinding before doing + // anything: the upgrade must come from an allowed Host and (when the browser + // sends one — it always does for WS) a same-site Origin. Writing to this socket + // injects keystrokes into a --dangerously-skip-permissions agent, so this gate + // matters even on the default no-password install. See security review H5. + const policy = getHostPolicy(); + if (!isAllowedRequestHost(req.headers.host, policy) || !isAllowedRequestOrigin(req.headers.origin, policy)) { + socket.close(4003, 'Forbidden'); + return; + } - const { id } = req.params; - const session = ctx.sessions.get(id); + const { id } = req.params; + const session = ctx.sessions.get(id); - if (!session) { - socket.close(4004, 'Session not found'); - return; - } + if (!session) { + socket.close(4004, 'Session not found'); + return; + } - // Structured transport logging — surfaces WS open/close/timeout churn so the - // tunnel-flap behavior (COD-134) is observable in the server logs. Fastify is - // configured logger:false, so we log via console (→ journald under systemd). + // Structured transport logging — surfaces WS open/close/timeout churn so the + // tunnel-flap behavior (COD-134) is observable in the server logs. Fastify is + // configured logger:false, so we log via console (→ journald under systemd). - // Enforce per-session connection limit - const currentCount = sessionWsCount.get(id) ?? 0; - if (currentCount >= MAX_WS_PER_SESSION) { - console.warn('[ws] terminal rejected: too many connections', { sessionId: id, wsCount: currentCount }); - socket.close(4008, 'Too many connections'); - return; - } - sessionWsCount.set(id, currentCount + 1); - console.info('[ws] terminal open', { sessionId: id, wsCount: currentCount + 1 }); + // Enforce per-session connection limit, scoped by clientId. A same-cid + // reconnect reclaims its own slot (registry evicts the stale socket), so a + // drop+reconnect burst can no longer over-count across the async-close gap + // and trip a spurious 4008. cid-less upgrades are admitted anonymously. + const cid = typeof req.query?.cid === 'string' && req.query.cid.length > 0 ? req.query.cid : null; + const { admitted, evictedSocket } = sessionWsRegistry.register(id, cid, socket); + if (!admitted) { + console.warn('[ws] terminal rejected: too many connections', { + sessionId: id, + wsCount: sessionWsRegistry.liveCount(id), + }); + socket.close(4008, 'Too many connections'); + return; + } + if (evictedSocket) { + // Same client reconnected; retire the stale socket so it doesn't linger. + console.info('[ws] terminal superseded by reconnect', { sessionId: id }); + try { + evictedSocket.close(4010, 'Superseded by reconnect'); + } catch { + /* socket may already be closing */ + } + } + console.info('[ws] terminal open', { sessionId: id, wsCount: sessionWsRegistry.liveCount(id) }); - // Swallow socket errors — cleanup happens in 'close' - socket.on('error', () => {}); + // Eagerly free the slot on error/terminate — don't wait for the async + // 'close' (idempotent with the 'close' handler below). This is what kills + // the reconnect over-count: the slot is released the instant the socket dies. + socket.on('error', () => { + sessionWsRegistry.unregister(id, socket); + }); - // Per-connection micro-batch state - let batchChunks: string[] = []; - let batchSize = 0; - let batchTimer: ReturnType | null = null; + // Per-connection micro-batch state + let batchChunks: string[] = []; + let batchSize = 0; + let batchTimer: ReturnType | null = null; - const flushBatch = () => { - batchTimer = null; - if (batchChunks.length === 0 || socket.readyState !== 1) { + const flushBatch = () => { + batchTimer = null; + if (batchChunks.length === 0 || socket.readyState !== 1) { + batchChunks = []; + batchSize = 0; + return; + } + const data = batchChunks.join(''); batchChunks = []; batchSize = 0; - return; - } - const data = batchChunks.join(''); - batchChunks = []; - batchSize = 0; - socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`); - }; + socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`); + }; - // Per-connection desktop sizing claim — registered on the first - // desktop-typed resize and released on socket close, so Session.resize() - // can ignore small-viewport resizes only while a desktop is actually - // connected (see Session._desktopSizeClaims). - const sizingToken = Symbol('ws-desktop-sizing'); - let holdsDesktopClaim = false; + // Per-connection desktop sizing claim — registered on the first + // desktop-typed resize and released on socket close, so Session.resize() + // can ignore small-viewport resizes only while a desktop is actually + // connected (see Session._desktopSizeClaims). + const sizingToken = Symbol('ws-desktop-sizing'); + let holdsDesktopClaim = false; - // Attach message handler synchronously BEFORE any async work - // (@fastify/websocket requirement to avoid dropped messages). - socket.on('message', (raw) => { - try { - const msg = JSON.parse(String(raw)); - if (msg.t === 'i' && typeof msg.d === 'string') { - if (msg.d.length > MAX_INPUT_LENGTH) return; - // Reliable delivery: when the frame carries a clientId + seq, apply it - // exactly once (skip a duplicate redelivery) but ACK it regardless so - // the client can drop it from its durable queue. Frames without seq - // (legacy/other tools) are applied as-is — no behavior change. - const cid = typeof msg.cid === 'string' ? msg.cid : null; - const seq = Number.isInteger(msg.seq) ? (msg.seq as number) : null; - const apply = cid && seq !== null ? session.shouldApplyInput(cid, seq) : true; - if (apply) { - // Typed input from a claim-holding desktop keeps the claim "hot" - // and re-asserts the desktop layout after a mobile override. - if (holdsDesktopClaim) session.noteDesktopActivity(); - session.write(msg.d); + // Attach message handler synchronously BEFORE any async work + // (@fastify/websocket requirement to avoid dropped messages). + socket.on('message', (raw) => { + try { + const msg = JSON.parse(String(raw)); + if (msg.t === 'i' && typeof msg.d === 'string') { + if (msg.d.length > MAX_INPUT_LENGTH) return; + // Reliable delivery: when the frame carries a clientId + seq, apply it + // exactly once (skip a duplicate redelivery) but ACK it regardless so + // the client can drop it from its durable queue. Frames without seq + // (legacy/other tools) are applied as-is — no behavior change. + const cid = typeof msg.cid === 'string' ? msg.cid : null; + const seq = Number.isInteger(msg.seq) ? (msg.seq as number) : null; + const apply = cid && seq !== null ? session.shouldApplyInput(cid, seq) : true; + if (apply) { + // Typed input from a claim-holding desktop keeps the claim "hot" + // and re-asserts the desktop layout after a mobile override. + if (holdsDesktopClaim) session.noteDesktopActivity(); + session.write(msg.d); + } + if (seq !== null && socket.readyState === 1) { + socket.send(`{"t":"ia","seq":${seq}}`); + } + } else if ( + msg.t === 'z' && + Number.isInteger(msg.c) && + Number.isInteger(msg.r) && + msg.c >= 1 && + msg.c <= 500 && + msg.r >= 1 && + msg.r <= 200 + ) { + const viewportType = msg.v === 'mobile' || msg.v === 'tablet' || msg.v === 'desktop' ? msg.v : undefined; + if (viewportType === 'desktop') { + session.claimDesktopSizing(sizingToken); + holdsDesktopClaim = true; + } else if (viewportType) { + // The connection's viewport can change (e.g. browser window + // narrowed past the tablet breakpoint) — drop a stale claim. + session.releaseDesktopSizing(sizingToken); + holdsDesktopClaim = false; + } + const force = msg.f === true; + session.resize(msg.c, msg.r, { viewportType, force }); } - if (seq !== null && socket.readyState === 1) { - socket.send(`{"t":"ia","seq":${seq}}`); - } - } else if ( - msg.t === 'z' && - Number.isInteger(msg.c) && - Number.isInteger(msg.r) && - msg.c >= 1 && - msg.c <= 500 && - msg.r >= 1 && - msg.r <= 200 - ) { - const viewportType = msg.v === 'mobile' || msg.v === 'tablet' || msg.v === 'desktop' ? msg.v : undefined; - if (viewportType === 'desktop') { - session.claimDesktopSizing(sizingToken); - holdsDesktopClaim = true; - } else if (viewportType) { - // The connection's viewport can change (e.g. browser window - // narrowed past the tablet breakpoint) — drop a stale claim. - session.releaseDesktopSizing(sizingToken); - holdsDesktopClaim = false; - } - const force = msg.f === true; - session.resize(msg.c, msg.r, { viewportType, force }); + } catch { + // Ignore malformed messages } - } catch { - // Ignore malformed messages - } - }); + }); - // Terminal output -> micro-batched WS send - const onTerminal = (data: string) => { - if (socket.readyState !== 1) return; - batchChunks.push(data); - batchSize += data.length; + // Terminal output -> micro-batched WS send + const onTerminal = (data: string) => { + if (socket.readyState !== 1) return; + batchChunks.push(data); + batchSize += data.length; - // Flush immediately for large batches (responsiveness during bulk output) - if (batchSize > WS_BATCH_FLUSH_THRESHOLD) { - if (batchTimer) { - clearTimeout(batchTimer); + // Flush immediately for large batches (responsiveness during bulk output) + if (batchSize > WS_BATCH_FLUSH_THRESHOLD) { + if (batchTimer) { + clearTimeout(batchTimer); + } + flushBatch(); + return; } - flushBatch(); - return; - } - // Start timer if not already running - if (!batchTimer) { - batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS); - } - }; + // Start timer if not already running + if (!batchTimer) { + batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS); + } + }; - const onClearTerminal = () => { - if (socket.readyState === 1) { - socket.send('{"t":"c"}'); - } - }; + const onClearTerminal = () => { + if (socket.readyState === 1) { + socket.send('{"t":"c"}'); + } + }; - const onNeedsRefresh = () => { - if (socket.readyState === 1) { - socket.send('{"t":"r"}'); - } - }; + const onNeedsRefresh = () => { + if (socket.readyState === 1) { + socket.send('{"t":"r"}'); + } + }; - // Close WS when session exits (deleted, respawned, or crashed) — prevents - // orphaned listeners and stale writes to a dead PTY. - const onSessionExit = () => { - socket.close(4009, 'Session terminated'); - }; + // Close WS when session exits (deleted, respawned, or crashed) — prevents + // orphaned listeners and stale writes to a dead PTY. + const onSessionExit = () => { + socket.close(4009, 'Session terminated'); + }; - session.on('terminal', onTerminal); - session.on('clearTerminal', onClearTerminal); - session.on('needsRefresh', onNeedsRefresh); - session.on('exit', onSessionExit); + session.on('terminal', onTerminal); + session.on('clearTerminal', onClearTerminal); + session.on('needsRefresh', onNeedsRefresh); + session.on('exit', onSessionExit); - // Heartbeat: detect stale connections (especially through tunnels where - // TCP RST can take minutes to propagate). - let pongTimeout: ReturnType | null = null; + // Heartbeat: detect stale connections (especially through tunnels where + // TCP RST can take minutes to propagate). + let pongTimeout: ReturnType | null = null; - socket.on('pong', () => { - if (pongTimeout) { - clearTimeout(pongTimeout); - pongTimeout = null; - } - }); + socket.on('pong', () => { + if (pongTimeout) { + clearTimeout(pongTimeout); + pongTimeout = null; + } + }); - const pingInterval = setInterval(() => { - if (socket.readyState !== 1) return; - socket.ping(); - pongTimeout = setTimeout(() => { - console.warn('[ws] terminal ping timeout — terminating', { sessionId: id }); - socket.terminate(); - }, WS_PONG_TIMEOUT_MS); - }, WS_PING_INTERVAL_MS); + const pingInterval = setInterval(() => { + if (socket.readyState !== 1) return; + socket.ping(); + pongTimeout = setTimeout(() => { + console.warn('[ws] terminal ping timeout — terminating', { sessionId: id }); + // Free the slot eagerly — terminate()'s 'close' may lag, and a client + // reconnecting after a stale-connection drop must not be over-counted. + sessionWsRegistry.unregister(id, socket); + socket.terminate(); + }, WS_PONG_TIMEOUT_MS); + }, WS_PING_INTERVAL_MS); - socket.on('close', (code: number, reason: Buffer) => { - clearInterval(pingInterval); - if (pongTimeout) clearTimeout(pongTimeout); - if (batchTimer) clearTimeout(batchTimer); - batchChunks = []; - session.off('terminal', onTerminal); - session.off('clearTerminal', onClearTerminal); - session.off('needsRefresh', onNeedsRefresh); - session.off('exit', onSessionExit); - session.releaseDesktopSizing(sizingToken); + socket.on('close', (code: number, reason: Buffer) => { + clearInterval(pingInterval); + if (pongTimeout) clearTimeout(pongTimeout); + if (batchTimer) clearTimeout(batchTimer); + batchChunks = []; + session.off('terminal', onTerminal); + session.off('clearTerminal', onClearTerminal); + session.off('needsRefresh', onNeedsRefresh); + session.off('exit', onSessionExit); + session.releaseDesktopSizing(sizingToken); - // Decrement per-session connection count - const count = sessionWsCount.get(id) ?? 1; - const remaining = count <= 1 ? 0 : count - 1; - if (remaining === 0) { - sessionWsCount.delete(id); - } else { - sessionWsCount.set(id, remaining); - } - console.info('[ws] terminal close', { sessionId: id, code, reason: String(reason), wsCount: remaining }); - }); - }); + // Release this socket's slot. Idempotent and identity-matched: if this + // socket was already superseded (a same-cid reconnect took its slot) or + // eagerly unregistered on terminate/error, this is a no-op and the + // reconnected socket keeps the slot. + sessionWsRegistry.unregister(id, socket); + console.info('[ws] terminal close', { + sessionId: id, + code, + reason: String(reason), + wsCount: sessionWsRegistry.liveCount(id), + }); + }); + } + ); } diff --git a/src/web/ws-connection-registry.ts b/src/web/ws-connection-registry.ts new file mode 100644 index 00000000..d0c64cda --- /dev/null +++ b/src/web/ws-connection-registry.ts @@ -0,0 +1,120 @@ +/** + * @fileoverview Per-session WebSocket connection registry (COD-137). + * + * Replaces the bare `Map` counter that previously gated + * `MAX_WS_PER_SESSION`. That counter had two defects: + * + * 1. Transient over-count on reconnect: a client that drops and immediately + * reconnects could land its new upgrade BEFORE the old socket's async + * `close` fired, so the count briefly exceeded the live connection number. + * A reconnect burst could hit the cap and the next upgrade was rejected + * with 4008 → the client fell back to HTTP. (The real spurious-4008 defect.) + * 2. No clientId scoping: the limit counted raw sockets, so a reconnecting + * client consumed a NEW slot instead of replacing its own. + * + * This registry tracks the live socket(s) per session keyed by clientId (`cid`, + * parsed from the upgrade URL query). The reliable-input protocol + * (`session.shouldApplyInput(cid, seq)`) already assumes ONE logical client per + * `cid` per session, so a new upgrade for a `cid` that already holds a socket is + * a SUPERSEDE — the registry evicts the stale socket and reuses its slot, which + * makes a reconnect reclaim rather than double-count (fixes #1 and #2). + * + * Backward-compat: an upgrade with NO `cid` (legacy clients, other tools) is + * admitted anonymously — it counts toward the limit but never evicts another + * client, and several anonymous sockets can coexist up to the cap. + * + * The class is pure (no `ws`/Fastify imports) and generic over a minimal socket + * shape so it can be unit-tested with plain fakes. The route owns the actual + * socket close/terminate; the registry only decides admit/evict and tracks slots. + */ + +/** Minimal socket shape the registry needs — satisfied by `ws` WebSocket. */ +export interface RegistrableSocket { + /** Identity comparison only; never dereferenced beyond `===`. */ + readonly readyState?: number; +} + +export interface RegisterResult { + /** Whether the new socket was admitted (false → caller should reject with 4008). */ + admitted: boolean; + /** + * A stale socket whose slot the new socket reclaimed (same `cid`). The caller + * should close it. Present only on a keyed supersede; never set for anonymous + * upgrades or fresh slots. + */ + evictedSocket?: S; +} + +/** A single live entry: the socket plus its clientId (null = anonymous). */ +interface Entry { + socket: S; + cid: string | null; +} + +export class WsConnectionRegistry { + /** sessionId → live entries (keyed + anonymous). */ + private readonly bySession = new Map[]>(); + + constructor(private readonly maxPerSession: number) {} + + /** + * Attempt to register a new socket for `(sessionId, cid)`. + * + * - cid present and already holds a socket → SUPERSEDE: evict the old one, + * reuse its slot, always admit. + * - otherwise → admit iff distinct-entry count < maxPerSession. + * + * A null/empty `cid` is anonymous: it never matches an existing entry and so + * never evicts; it just consumes a slot. + */ + register(sessionId: string, cid: string | null, socket: S): RegisterResult { + const entries = this.bySession.get(sessionId) ?? []; + + if (cid) { + const existingIdx = entries.findIndex((e) => e.cid === cid); + if (existingIdx !== -1) { + const evicted = entries[existingIdx].socket; + // Reuse the slot in place — no net change to the live count, so a + // reconnect can never be rejected by the cap. + entries[existingIdx] = { socket, cid }; + this.bySession.set(sessionId, entries); + return { admitted: true, evictedSocket: evicted === socket ? undefined : evicted }; + } + } + + if (entries.length >= this.maxPerSession) { + return { admitted: false }; + } + + entries.push({ socket, cid: cid || null }); + this.bySession.set(sessionId, entries); + return { admitted: true }; + } + + /** + * Remove a socket from its session. Idempotent — safe to call on `close`, + * `error`, AND eagerly on `terminate()` (the over-count fix relies on eager + * removal freeing the slot before the async `close` fires). + * + * Matches by socket identity, so a socket that was already superseded + * (replaced in-slot by a same-cid reconnect) is NOT removed by its late + * `close` — the new socket keeps the slot. + */ + unregister(sessionId: string, socket: S): void { + const entries = this.bySession.get(sessionId); + if (!entries) return; + const idx = entries.findIndex((e) => e.socket === socket); + if (idx === -1) return; + entries.splice(idx, 1); + if (entries.length === 0) { + this.bySession.delete(sessionId); + } else { + this.bySession.set(sessionId, entries); + } + } + + /** Number of live entries for a session (0 if none). */ + liveCount(sessionId: string): number { + return this.bySession.get(sessionId)?.length ?? 0; + } +} diff --git a/test/routes/ws-routes.test.ts b/test/routes/ws-routes.test.ts index 93314dd7..5b722036 100644 --- a/test/routes/ws-routes.test.ts +++ b/test/routes/ws-routes.test.ts @@ -491,6 +491,29 @@ describe('ws-routes', () => { for (const ws of connections) ws.close(); } }); + + it('reconnecting client (same cid) is admitted at the cap instead of 4008 (COD-137)', async () => { + const connections: WebSocket[] = []; + try { + // Fill all 5 slots with DISTINCT clients, one of which is "alice". + for (const c of ['alice', 'b', 'c', 'd', 'e']) { + connections.push(await connectWs(`/ws/sessions/ws-test-session/terminal?cid=${c}`)); + } + + // Alice reconnects WHILE her old socket is still registered (the + // over-count window). This must reclaim her slot, not hit the cap. + const aliceNew = await connectWs('/ws/sessions/ws-test-session/terminal?cid=alice'); + connections.push(aliceNew); + + // Sanity: the reconnected socket is live and usable. + ctx._session.emit('terminal', 'reconnected-ok'); + const msg = (await nextMessage(aliceNew)) as { t: string; d: string }; + expect(msg.t).toBe('o'); + expect(msg.d).toContain('reconnected-ok'); + } finally { + for (const ws of connections) ws.close(); + } + }); }); // ========== Heartbeat ========== diff --git a/test/ws-connection-registry.test.ts b/test/ws-connection-registry.test.ts new file mode 100644 index 00000000..35105878 --- /dev/null +++ b/test/ws-connection-registry.test.ts @@ -0,0 +1,112 @@ +/** + * @fileoverview Unit tests for WsConnectionRegistry (COD-137). + * + * The registry is the pure decision unit extracted out of ws-routes.ts so the + * connection-limit / clientId-eviction logic is testable without driving real + * WebSocket upgrades. Uses plain fake sockets (identity only). + * + * @dependency src/web/ws-connection-registry.ts + */ + +import { describe, it, expect } from 'vitest'; +import { WsConnectionRegistry } from '../src/web/ws-connection-registry.js'; + +/** Fake socket — registry only compares identity, so any object works. */ +const sock = (label: string) => ({ readyState: 1, label }); + +describe('WsConnectionRegistry', () => { + it('reconnecting client (same cid) reclaims its slot instead of being rejected at the limit', () => { + const reg = new WsConnectionRegistry(5); + // Fill all 5 slots with distinct clients, one of which is "alice". + for (const c of ['alice', 'b', 'c', 'd', 'e']) { + expect(reg.register('s1', c, sock(c)).admitted).toBe(true); + } + expect(reg.liveCount('s1')).toBe(5); + + // Alice's new upgrade lands BEFORE her old socket's async close fires. + const aliceNew = sock('alice-new'); + const res = reg.register('s1', 'alice', aliceNew); + + expect(res.admitted).toBe(true); // NOT a spurious 4008 + expect(res.evictedSocket).toBeDefined(); // old alice socket handed back to close + expect(reg.liveCount('s1')).toBe(5); // slot reused, not double-counted + }); + + it('still rejects a genuine (N+1)th DISTINCT client', () => { + const reg = new WsConnectionRegistry(5); + for (const c of ['a', 'b', 'c', 'd', 'e']) { + expect(reg.register('s1', c, sock(c)).admitted).toBe(true); + } + const sixth = reg.register('s1', 'f', sock('f')); + expect(sixth.admitted).toBe(false); + expect(sixth.evictedSocket).toBeUndefined(); + expect(reg.liveCount('s1')).toBe(5); + }); + + it('eager removal on terminate frees a slot immediately', () => { + const reg = new WsConnectionRegistry(5); + const sockets = ['a', 'b', 'c', 'd', 'e'].map((c) => { + const s = sock(c); + reg.register('s1', c, s); + return [c, s] as const; + }); + expect(reg.register('s1', 'f', sock('f')).admitted).toBe(false); + + // Eagerly unregister one (simulating terminate/error, not async close). + reg.unregister('s1', sockets[0][1]); + expect(reg.liveCount('s1')).toBe(4); + + // Now a brand-new distinct client is admitted. + expect(reg.register('s1', 'f', sock('f')).admitted).toBe(true); + expect(reg.liveCount('s1')).toBe(5); + }); + + it('cid-less upgrades are admitted up to the limit and never evict a keyed client', () => { + const reg = new WsConnectionRegistry(5); + const keyed = sock('keyed'); + reg.register('s1', 'keyed', keyed); + + // Four anonymous upgrades fill the rest of the cap. + for (let i = 0; i < 4; i++) { + const res = reg.register('s1', null, sock(`anon${i}`)); + expect(res.admitted).toBe(true); + expect(res.evictedSocket).toBeUndefined(); // never evicts the keyed client + } + expect(reg.liveCount('s1')).toBe(5); + + // 6th anonymous is rejected — anonymous sockets count toward the cap. + expect(reg.register('s1', null, sock('anon-extra')).admitted).toBe(false); + + // The keyed client is untouched: a same-cid reconnect still reclaims. + const keyedNew = sock('keyed-new'); + const res = reg.register('s1', 'keyed', keyedNew); + expect(res.admitted).toBe(true); + expect(res.evictedSocket).toBe(keyed); + }); + + it('late close of a superseded socket does not evict the reconnected one', () => { + const reg = new WsConnectionRegistry(5); + const old = sock('old'); + reg.register('s1', 'alice', old); + const fresh = sock('fresh'); + reg.register('s1', 'alice', fresh); // supersede + + // The stale socket's async close arrives late — must NOT remove fresh. + reg.unregister('s1', old); + expect(reg.liveCount('s1')).toBe(1); + + // Fresh is still the live entry: another reconnect evicts fresh, not old. + const fresher = sock('fresher'); + expect(reg.register('s1', 'alice', fresher).evictedSocket).toBe(fresh); + }); + + it('isolates counts per session', () => { + const reg = new WsConnectionRegistry(2); + reg.register('s1', 'a', sock('a')); + reg.register('s1', 'b', sock('b')); + expect(reg.register('s1', 'c', sock('c')).admitted).toBe(false); + // s2 has its own budget. + expect(reg.register('s2', 'a', sock('a2')).admitted).toBe(true); + expect(reg.liveCount('s2')).toBe(1); + }); +});