Merge PR #149 from aakhter/cod-165-ws-resilience: WebSocket durable-delivery resilience

Includes review fixes: real _wsState lifecycle (connecting/connected/disconnected), per-tab supersede identity (multi-tab coexistence), preserved reconnect backoff, connection-dot CSS for connected/fallback states.
This commit is contained in:
Codeman maintainer
2026-07-12 20:07:55 +02:00
13 changed files with 1994 additions and 214 deletions
+204 -43
View File
@@ -455,6 +455,14 @@ class CodemanApp {
this._clientId = (typeof crypto !== 'undefined' && crypto.randomUUID)
? crypto.randomUUID()
: 'c-' + Math.random().toString(36).slice(2) + Date.now().toString(36);
// Per-TAB nonce for the WS registry key (COD-137). _loadReliableState()
// later replaces _clientId with the browser-wide localStorage identity
// (shared by every tab/window of this profile), so the WS upgrade sends
// `clientId:nonce` instead — a same-tab reconnect still supersedes its own
// socket, but two tabs on one session coexist instead of evicting each
// other in a 4010 ping-pong. Input frames keep the bare clientId for seq
// dedup.
this._wsTabNonce = this._clientId;
this.terminal = null;
this.fitAddon = null;
this.activeSessionId = null;
@@ -589,6 +597,8 @@ class CodemanApp {
this._ws = null; // WebSocket instance for active session
this._wsSessionId = null; // Session ID the WS is connected to
this._wsReady = false; // True when WS is open and ready for I/O
this._wsState = 'disconnected'; // connecting | connected | reconnecting | fallback | disconnected
this._wsLastRecvAt = 0; // ms timestamp of the last frame received on the active WS
// Terminal write batching with DEC 2026 sync support
this.pendingWrites = [];
@@ -632,6 +642,9 @@ class CodemanApp {
this._clientId = '';
this._seqCounters = new Map(); // sessionId -> last issued seq
this._pendingDeliveries = new Map(); // sessionId -> [{seq,data,useMux,ts,tries,sentAt}]
// Last rendered connection-indicator tuple; the hot input path skips DOM
// writes when the freshly computed descriptor is identical (COD-136).
this._lastIndicatorDescriptor = null;
this._postDraining = new Set(); // sessionIds with an in-flight POST drainer
this._persistReliableTimer = null;
this._reliableAckTimeoutMs = 4000; // unacked WS frame older than this ⇒ socket likely dead
@@ -2142,9 +2155,21 @@ class CodemanApp {
*/
_connectWs(sessionId) {
this._disconnectWs();
this._wsState = 'connecting';
this._updateConnectionIndicator();
const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal`;
// Pass a per-TAB identity on the upgrade URL so the server's connection
// registry scopes the per-session limit by connection (COD-137): a same-tab
// reconnect supersedes its own socket instead of consuming a new slot and
// tripping a spurious 4008, while two tabs of the same browser (which share
// the localStorage clientId) each keep their own socket. The bare clientId
// still rides the input frames for seq dedup. Omitted if clientId is
// unavailable (server then treats the upgrade as anonymous — still admitted
// up to the limit).
const cid = this._clientId ? `${this._clientId}:${this._wsTabNonce}` : '';
const cidQuery = cid ? `?cid=${encodeURIComponent(cid)}` : '';
const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal${cidQuery}`;
const ws = new WebSocket(url);
this._ws = ws;
this._wsSessionId = sessionId;
@@ -2153,7 +2178,9 @@ class CodemanApp {
// Only mark ready if this is still the intended session
if (this._ws === ws) {
this._wsReady = true;
this._wsState = 'connected';
this._wsReconnectAttempts = 0;
this._updateConnectionIndicator();
// Send a typed resize over the fresh socket: syncs PTY dims after
// (re)connects AND registers the desktop sizing claim server-side —
// selectSession's earlier resizes ran before this WS existed, so they
@@ -2168,6 +2195,9 @@ class CodemanApp {
ws.onmessage = (event) => {
if (this._ws !== ws) return;
// Mark the socket as alive on every received frame (output, ACK, etc.) so
// the redeliver sweep only force-closes a genuinely silent connection.
this._wsLastRecvAt = Date.now();
try {
const msg = JSON.parse(event.data);
if (msg.t === 'o') {
@@ -2194,18 +2224,53 @@ class CodemanApp {
this._wsReady = false;
this._stopMobileResizeRetry();
// Reconnect on unexpected close (server restart, network blip, ping timeout).
// Don't reconnect if we intentionally disconnected (_disconnectWs nulls onclose)
// or if the server rejected the session (4004=not found, 4008=too many, 4009=terminated).
if (event.code < 4004 && this.activeSessionId === sessionId) {
const delay = Math.min(1000 * Math.pow(2, this._wsReconnectAttempts || 0), 10000);
this._wsReconnectAttempts = (this._wsReconnectAttempts || 0) + 1;
this._wsReconnectTimer = setTimeout(() => {
this._wsReconnectTimer = null;
if (this.activeSessionId === sessionId) {
this._connectWs(sessionId);
}
}, delay);
// Decide what to do next from the close code + how many consecutive
// reconnects we've already made (pure policy in constants.js):
// reconnect → transient (server restart, network blip, ping timeout);
// schedule a backoff retry while this session stays active.
// retry-fallback → too-many-connections / unknown rejection; show the HTTP
// fallback but keep retrying so we return to WS when it clears.
// give-up → 4004 (not found) / 4009 (terminated); the session is gone.
// _disconnectWs() nulls onclose for intentional disconnects, so we never land here for those.
const plan = window.CodemanWsReconnect.plan(event.code, this._wsReconnectAttempts || 0);
_crashDiag.log(
`WS CLOSE code=${event.code} reason=${event.reason || ''} action=${plan.action} attempts=${this._wsReconnectAttempts || 0}`
);
const stillActive = this.activeSessionId === sessionId;
if (plan.action === 'give-up') {
this._wsState = stillActive ? 'fallback' : 'disconnected';
this._updateConnectionIndicator();
} else if (plan.action === 'reconnect') {
if (stillActive) {
this._wsState = 'reconnecting';
this._updateConnectionIndicator();
const delay = plan.delayMs + Math.floor(Math.random() * 250); // jitter to de-sync herds
this._wsReconnectAttempts = (this._wsReconnectAttempts || 0) + 1;
this._wsReconnectTimer = setTimeout(() => {
this._wsReconnectTimer = null;
if (this.activeSessionId === sessionId) {
this._connectWs(sessionId);
}
}, delay);
} else {
this._wsState = 'disconnected';
this._updateConnectionIndicator();
}
} else {
// retry-fallback: surface the HTTP fallback, but keep trying on a bounded
// timer so the transport returns to WS once the transient condition clears.
this._wsState = stillActive ? 'fallback' : 'disconnected';
this._updateConnectionIndicator();
if (stillActive) {
this._wsReconnectAttempts = (this._wsReconnectAttempts || 0) + 1;
this._wsReconnectTimer = setTimeout(() => {
this._wsReconnectTimer = null;
if (this.activeSessionId === sessionId) {
this._connectWs(sessionId);
}
}, plan.delayMs);
}
}
};
@@ -2217,7 +2282,11 @@ class CodemanApp {
/** Close the active WebSocket connection (if any). */
_disconnectWs() {
this._clearTimer('_wsReconnectTimer');
this._wsReconnectAttempts = 0;
// Deliberately do NOT reset _wsReconnectAttempts here: _connectWs() calls
// this first, so a reset would restart the exponential backoff ladder at
// attempt 0 on every retry (≈0ms tight reconnect loop during an outage).
// ws.onopen zeroes the counter once a connection actually succeeds.
this._wsState = 'disconnected';
this._stopMobileResizeRetry();
if (this._ws) {
this._ws.onclose = null; // Prevent re-entrant cleanup
@@ -2412,7 +2481,13 @@ class CodemanApp {
this._ws && this._ws.readyState === WebSocket.OPEN && this._wsSessionId === sessionId;
if (isActiveWs) {
const oldest = list[0];
if (oldest && oldest.sentAt && Date.now() - oldest.sentAt > this._reliableAckTimeoutMs) {
// Only tear the socket down when the oldest unacked frame is stale AND the
// socket has been silent for the timeout: a connection still delivering
// output/ACKs is alive (the ACK is just behind), so force-closing it would
// cause needless WS↔HTTP flapping. A truly half-open socket goes quiet.
const stale = oldest && oldest.sentAt && Date.now() - oldest.sentAt > this._reliableAckTimeoutMs;
const silent = Date.now() - this._wsLastRecvAt > this._reliableAckTimeoutMs;
if (stale && silent) {
try {
this._ws.close(); // half-open: never recovers on its own — force reconnect
} catch {
@@ -2420,6 +2495,18 @@ class CodemanApp {
}
continue;
}
if (stale) {
// Stale but the socket is still delivering output: the ACK was lost,
// not the connection. Force-closing isn't warranted (the link is fine),
// but the fast path skips anything with sentAt!==0, so the stranded
// frame would never re-send. Reset sentAt=0 on every stale unacked
// frame so the _drainSession below re-drives them over the live socket
// (server dedups by seq, so a re-sent lost-ACK frame is harmless).
// Frames sent recently (not yet stale) are left untouched.
for (const rec of list) {
if (rec.sentAt && Date.now() - rec.sentAt > this._reliableAckTimeoutMs) rec.sentAt = 0;
}
}
}
this._drainSession(sessionId);
}
@@ -2530,39 +2617,107 @@ class CodemanApp {
}
}
// Pure render of the header connection indicator: reads only `this.*` state,
// touches NO DOM. Returns the exact { display, dotClass, text, title } tuple the
// writer applies. When hidden (display:'none') the other three are normalized to
// '' so the cache compare in _updateConnectionIndicator() is well-defined.
// Every branch/string here must stay byte-identical to what's rendered today.
_computeConnectionDescriptor() {
const { bytes: totalBytes, count } = this._pendingBytes();
const hasQueue = count > 0;
// Only surface a backlog once it's more than a few bytes. A single keystroke
// (1B) ACKs in milliseconds, so without this the label flickered "sending 1B"
// on every key press. Above this threshold means input is genuinely backing up.
const BACKLOG_HINT_BYTES = 4;
const showBacklog = totalBytes > BACKLOG_HINT_BYTES;
const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`);
const queuedSuffix = showBacklog ? ` · ${formatBytes(totalBytes)} queued` : '';
// Hard offline (browser reports no network) dominates everything.
if (!this.isOnline || this._connectionStatus === 'offline') {
return {
display: 'flex',
dotClass: 'connection-dot offline',
text: showBacklog ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline',
title: 'No network connection',
};
}
// With an active terminal, show its transport (WebSocket vs HTTP fallback).
if (this.activeSessionId) {
let cls, label, detail;
switch (this._wsState) {
case 'connected':
cls = 'connected'; label = 'WS'; detail = 'Terminal connected over WebSocket';
break;
case 'fallback':
cls = 'fallback'; label = 'HTTP'; detail = 'WebSocket unavailable — input sent over HTTP';
break;
case 'reconnecting':
cls = 'reconnecting'; label = 'WS…'; detail = 'Reconnecting WebSocket';
break;
case 'connecting':
default:
cls = 'reconnecting'; label = 'WS…'; detail = 'Connecting WebSocket';
break;
}
return {
display: 'flex',
dotClass: `connection-dot ${cls}`,
text: `${label}${queuedSuffix}`,
title: detail,
};
}
// No active terminal — reflect the SSE event stream only when it needs attention.
if (this._connectionStatus === 'reconnecting' || this._connectionStatus === 'disconnected') {
return {
display: 'flex',
dotClass: 'connection-dot reconnecting',
text: showBacklog ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...',
title: 'Reconnecting to server',
};
}
// Idle dashboard, healthy stream — hide unless input is genuinely queued.
if (!hasQueue) {
return { display: 'none', dotClass: '', text: '', title: '' };
}
return {
display: 'flex',
dotClass: 'connection-dot draining',
text: showBacklog ? `Sending ${formatBytes(totalBytes)}...` : 'Sending...',
title: 'Delivering queued input',
};
}
_updateConnectionIndicator() {
const indicator = this.$('connectionIndicator');
const dot = this.$('connectionDot');
const text = this.$('connectionText');
if (!indicator || !dot || !text) return;
const status = this._connectionStatus;
// While the connection is healthy, never surface the input queue. With the
// reliable-delivery layer every keystroke is briefly "pending" until its ACK
// lands a few ms later — showing that flashed "Sending 1B…" on every single
// character. The indicator is only meaningful for an actual connection
// problem (reconnecting / offline), where the queued byte count reassures
// the user their typing is safely buffered and will be sent.
if (status === 'connected' || status === 'connecting') {
indicator.style.display = 'none';
// Called on EVERY keystroke (_reliableSend) and EVERY ACK (_ackDelivery).
// During fast typing the rendered tuple is usually identical, so skip the DOM
// writes when nothing changed (COD-136) — the compute above is DOM-free.
const next = this._computeConnectionDescriptor();
const prev = this._lastIndicatorDescriptor;
if (
prev &&
prev.display === next.display &&
prev.dotClass === next.dotClass &&
prev.text === next.text &&
prev.title === next.title
) {
return;
}
this._lastIndicatorDescriptor = next;
const { bytes: totalBytes, count } = this._pendingBytes();
const hasQueue = count > 0;
indicator.style.display = 'flex';
dot.className = 'connection-dot';
const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`);
if (status === 'reconnecting') {
dot.classList.add('reconnecting');
text.textContent = hasQueue ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...';
} else {
// Offline or disconnected
dot.classList.add('offline');
text.textContent = hasQueue ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline';
indicator.style.display = next.display;
if (next.display !== 'none') {
dot.className = next.dotClass;
text.textContent = next.text;
indicator.title = next.title;
}
}
@@ -3801,6 +3956,9 @@ class CodemanApp {
// the buffer write, causing 70KB+ single-frame flushes that stall WebGL.
// chunkedTerminalWrite also sets this, but we need it before the fetch too.
const bufferLoadOwner = this._beginBufferLoad(selectGen);
// COD-144: track whether the load painted nothing (empty fetch + no cache).
// For that just-created-session case we flush (not discard) queued SSE events.
let bufferWasEmpty = false;
try {
// Fit terminal to container BEFORE writing any buffer data.
// If the browser was resized while viewing another session, the terminal
@@ -3966,13 +4124,16 @@ class CodemanApp {
} else if (!cachedBuffer) {
// No fresh buffer and no cache — clear any stale content
this._resetTerminalForReplay();
bufferWasEmpty = true;
}
// Buffer load complete — unblock live SSE writes (queued events are discarded
// to prevent duplicate content). chunkedTerminalWrite calls _finishBufferLoad
// internally, but if we skipped the write (cache hit or empty), call it here.
// Buffer load complete — unblock live SSE writes. chunkedTerminalWrite calls
// _finishBufferLoad internally (discarding queued events to prevent duplicate
// content); if we skipped the write (cache hit or empty), call it here.
// COD-144: when the load painted nothing, FLUSH the queued events instead of
// discarding — a new session's prompt arrives only as a queued SSE event.
if (this._isLoadingBuffer) {
this._finishBufferLoad(bufferLoadOwner);
this._finishBufferLoad(bufferLoadOwner, { flushQueued: bufferWasEmpty });
}
// Drop the guard so user input clears state normally
this._restoringFlushedState = false;
+27
View File
@@ -156,6 +156,30 @@ function shouldAutoWrapTabs(input) {
return scrollWidth > clientWidth + 1;
}
// COD-134 — Terminal WebSocket reconnect policy.
//
// Decide what to do after a terminal WebSocket closes, given the close `code`
// and `attempt` (0-based count of consecutive reconnects already made):
// - transient closes (code < 4004: 1000/1001/1005/1006/etc.) → 'reconnect'
// with exponential backoff (0 on the first attempt; the caller adds jitter),
// 250ms → 500 → 1000 → ... capped at 10s.
// - 4004 (session not found) / 4009 (session terminated) → 'give-up': the
// session is gone, retrying only wastes connections.
// - 4008 (too many connections) and any other code >= 4004 → 'retry-fallback':
// show the HTTP fallback but keep retrying on a bounded 5s timer so the
// transport returns to WS once the transient condition clears (un-stick).
// Pure: no DOM, no side effects.
function planWsReconnect(code, attempt) {
if (code === 4004 || code === 4009) {
return { action: 'give-up', delayMs: 0 };
}
if (code >= 4004) {
return { action: 'retry-fallback', delayMs: 5000 };
}
const delayMs = attempt <= 0 ? 0 : Math.min(250 * Math.pow(2, attempt - 1), 10000);
return { action: 'reconnect', delayMs };
}
if (typeof window !== 'undefined') {
window.WEBGL_FALLBACK = WEBGL_FALLBACK;
window.evaluateWebGLLongTaskTrip = evaluateWebGLLongTaskTrip;
@@ -163,6 +187,9 @@ if (typeof window !== 'undefined') {
window.CodemanTabOverflow = {
shouldAutoWrapTabs,
};
window.CodemanWsReconnect = {
plan: planWsReconnect,
};
}
// Scheduler API — prioritize terminal writes over background UI updates.
+9
View File
@@ -626,6 +626,15 @@ body {
flex-shrink: 0;
}
.connection-dot.connected {
background: var(--green);
}
.connection-dot.fallback {
background: var(--yellow);
box-shadow: 0 0 6px var(--yellow);
}
.connection-dot.offline {
background: var(--red);
box-shadow: 0 0 6px var(--red);
+27 -5
View File
@@ -2046,11 +2046,25 @@ Object.assign(CodemanApp.prototype, {
* Complete a buffer load: unblock live SSE writes.
* Called when chunkedTerminalWrite finishes (or is skipped for empty buffers).
*
* Queued SSE events are DISCARDED, not flushed. The loaded buffer from the API
* is the source of truth up to the response timestamp. SSE events queued during
* the fetch+write overlap with the buffer — flushing them writes duplicate data
* (especially Ink cursor-up redraws), corrupting the terminal display.
* By default queued SSE events are DISCARDED, not flushed. For an established
* session the loaded buffer from the API is the source of truth up to the
* response timestamp; SSE events queued during the fetch+write overlap already
* appear in that buffer, so flushing them writes duplicate data (especially Ink
* cursor-up redraws), corrupting the terminal display.
*
* COD-144: a brand-new session is the exception. Its terminal fetch can resolve
* BEFORE the PTY emits its first prompt, so the fetched buffer is empty and the
* prompt arrives only as a queued SSE event. Discarding it leaves the terminal
* blank until a tab-switch re-fetches a now-populated buffer. When the caller
* knows the load painted nothing (empty fetch + no cache), it passes
* `{ flushQueued: true }` so the queued events are REPLAYED through
* `batchTerminalWrite()` instead of dropped. Replay runs after `_isLoadingBuffer`
* is cleared, so the events write through normally and are not re-queued.
*
* After unblocking, new SSE/WS events deliver subsequent output normally.
*
* @param {string} [owner] Load token from `_beginBufferLoad`; a stale owner is a no-op.
* @param {{ flushQueued?: boolean }} [opts] When `flushQueued` is true, replay any queued events.
*/
_beginBufferLoad(owner) {
if (this._bufferLoadSeq === undefined) this._bufferLoadSeq = 0;
@@ -2061,13 +2075,21 @@ Object.assign(CodemanApp.prototype, {
return loadOwner;
},
_finishBufferLoad(owner) {
_finishBufferLoad(owner, opts) {
if (owner !== undefined && this._bufferLoadOwner !== owner) {
return false;
}
const queued = this._loadBufferQueue;
this._isLoadingBuffer = false;
this._loadBufferQueue = null;
this._bufferLoadOwner = null;
// COD-144: replay (rather than discard) queued live events when the load
// painted nothing — the queued prompt is the only content a new session has.
if (opts?.flushQueued && queued && queued.length) {
for (const data of queued) {
this.batchTerminalWrite(data);
}
}
return true;
},
+210 -166
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,197 +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;
}
// Enforce per-session connection limit
const currentCount = sessionWsCount.get(id) ?? 0;
if (currentCount >= MAX_WS_PER_SESSION) {
socket.close(4008, 'Too many connections');
return;
}
sessionWsCount.set(id, currentCount + 1);
// 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).
// Swallow socket errors — cleanup happens in 'close'
socket.on('error', () => {});
// 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) });
// Per-connection micro-batch state
let batchChunks: string[] = [];
let batchSize = 0;
let batchTimer: ReturnType<typeof setTimeout> | null = null;
// 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);
});
const flushBatch = () => {
batchTimer = null;
if (batchChunks.length === 0 || socket.readyState !== 1) {
// 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) {
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(() => {
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', () => {
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;
if (count <= 1) {
sessionWsCount.delete(id);
} else {
sessionWsCount.set(id, count - 1);
}
});
});
// 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),
});
});
}
);
}
+123
View File
@@ -0,0 +1,123 @@
/**
* @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 a per-TAB
* connection identity (`cid`, parsed from the upgrade URL query). The browser
* sends `clientId:tabNonce`, NOT the bare localStorage clientId — that one is
* shared by every tab/window of a profile, so keying on it would make two tabs
* on one session evict each other in a 4010 ping-pong. 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). The cid is opaque here; input-frame
* dedup uses the bare clientId separately (`session.shouldApplyInput`).
*
* 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;
}
}