COD-137 scope WS per-session limit by clientId (fix spurious 4008 on reconnect)

MAX_WS_PER_SESSION was gated by a bare Map<sessionId,number> 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.
This commit is contained in:
Aamer Akhter
2026-07-10 14:39:01 -04:00
parent 20cb42d202
commit 4ab89f9a4e
5 changed files with 471 additions and 175 deletions
+7 -1
View File
@@ -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;
+209 -174
View File
@@ -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<string, number>();
/**
* 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<WebSocket>(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<typeof setTimeout> | null = null;
// Per-connection micro-batch state
let batchChunks: string[] = [];
let batchSize = 0;
let batchTimer: ReturnType<typeof setTimeout> | 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<typeof setTimeout> | null = null;
// Heartbeat: detect stale connections (especially through tunnels where
// TCP RST can take minutes to propagate).
let pongTimeout: ReturnType<typeof setTimeout> | 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),
});
});
}
);
}
+120
View File
@@ -0,0 +1,120 @@
/**
* @fileoverview Per-session WebSocket connection registry (COD-137).
*
* Replaces the bare `Map<sessionId, number>` 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<S> {
/** 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<S> {
socket: S;
cid: string | null;
}
export class WsConnectionRegistry<S extends RegistrableSocket = RegistrableSocket> {
/** sessionId → live entries (keyed + anonymous). */
private readonly bySession = new Map<string, Entry<S>[]>();
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<S> {
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;
}
}
+23
View File
@@ -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 ==========
+112
View File
@@ -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);
});
});