diff --git a/docs/architecture.md b/docs/architecture.md index 0b8afa495..c46e463c7 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -209,7 +209,9 @@ There is no dedicated welcome message; the server emits a `status` session messa **Top-level WS envelopes** are `hello`, `recording_state`, `ping`/`pong`, and `session` (which wraps the rich union of session messages). -Client liveness checks use the top-level JSON `ping`/`pong` envelope, not a session RPC and not RFC6455 protocol ping. The app runs through browser and React Native WebSocket APIs, which do not expose protocol ping, so this envelope is the portable way to test the direct or relay data path. Session RPC timeouts are operation failures and must not be treated as proof that the socket is dead. +Client liveness checks use the top-level JSON `ping`/`pong` envelope, not a session RPC or RFC6455 control ping. Current clients ping every 10 seconds, beginning one interval after connecting. The first ping claims an application-ownership lease for that physical socket, all later inbound activity renews it, and the daemon forcibly terminates the socket if the lease expires. A legacy or raw socket that never sends an application ping never enters this lease and is not closed for omitting one. Session RPC timeouts are operation failures and must not be treated as proof that the socket is dead. + +Every physical send path enforces an 8 MiB outbound high-water mark, including JSON broadcasts, binary terminal frames, and the encrypted relay adapter's asynchronous queue. This sits above the terminal stream's 4 MiB soft backpressure threshold, leaving room for snapshot catch-up before the hard cutoff. JSON is serialized once per broadcast after sockets already at the limit are removed, then its exact byte length is checked for every remaining socket. A frame that would cross the limit is not sent; that physical socket is forcibly terminated without disturbing other sockets attached to the same logical session. Multiple tabs and simultaneous direct and relay paths may legitimately share a client id. Client session RPC waits default to 60s so slow relay or mobile networks do not turn a live but delayed daemon response into a false operation failure. Keep connect timeouts, app-level grace windows, explicit diagnostic latency probes, liveness ping timers, and genuinely long-running RPCs separate from this default. diff --git a/packages/server/src/server/relay-transport.ts b/packages/server/src/server/relay-transport.ts index 4aeca54f3..9eb8cc859 100644 --- a/packages/server/src/server/relay-transport.ts +++ b/packages/server/src/server/relay-transport.ts @@ -4,12 +4,12 @@ import { WebSocket } from "ws"; import type pino from "pino"; import { createDaemonChannel, - type EncryptedChannel, type Transport as RelayTransport, type KeyPair, } from "@getpaseo/relay/e2ee"; import { buildRelayWebSocketUrl } from "@getpaseo/protocol/daemon-endpoints"; import type { ExternalSocketMetadata } from "./websocket-server.js"; +import { createEncryptedRelaySocket } from "./websocket/encrypted-relay-socket.js"; interface RelayTransportOptions { logger: pino.Logger; @@ -27,8 +27,10 @@ export interface RelayTransportController { interface RelaySocketLike { readyState: number; + bufferedAmount?: number; send: (data: string | Uint8Array | ArrayBuffer) => void; close: (code?: number, reason?: string) => void; + terminate?: () => void; on: (event: "message" | "close" | "error", listener: (...args: unknown[]) => void) => void; once: (event: "close" | "error", listener: (...args: unknown[]) => void) => void; } @@ -60,18 +62,6 @@ function createDefaultRelayWebSocket(url: string): RelayWebSocketLike { return new WebSocket(url, RELAY_WEBSOCKET_OPTIONS); } -function normalizeRelaySendPayload(data: string | Uint8Array | ArrayBuffer): string | ArrayBuffer { - if (typeof data === "string") return data; - if (data instanceof ArrayBuffer) return data; - if (ArrayBuffer.isView(data)) { - const view = new Uint8Array(data.buffer, data.byteOffset, data.byteLength); - const out = new Uint8Array(view.byteLength); - out.set(view); - return out.buffer; - } - return String(data); -} - function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null; } @@ -450,7 +440,12 @@ async function attachEncryptedSocket( emitter.emit("error", error); }, }); - const encryptedSocket = createEncryptedSocket(channel, emitter); + const encryptedSocket = createEncryptedRelaySocket({ + channel, + emitter, + getTransportBufferedAmount: () => socket.bufferedAmount, + terminateTransport: () => socket.terminate(), + }); await attachSocket(encryptedSocket, metadata); attached = true; for (const message of pendingMessages) { @@ -502,42 +497,6 @@ function createRelayTransportAdapter( return relayTransport; } -function createEncryptedSocket(channel: EncryptedChannel, emitter: EventEmitter): RelaySocketLike { - let readyState = 1; - - channel.setState("open"); - - const close = (code?: number, reason?: string) => { - if (readyState === 3) return; - readyState = 3; - channel.close(code, reason); - }; - - emitter.on("close", () => { - if (readyState === 3) return; - readyState = 3; - }); - - return { - get readyState() { - return readyState; - }, - send: (data) => { - const outbound = normalizeRelaySendPayload(data); - void channel.send(outbound).catch((error) => { - emitter.emit("error", error); - }); - }, - close, - on: (event, listener) => { - emitter.on(event, listener); - }, - once: (event, listener) => { - emitter.once(event, listener); - }, - }; -} - function normalizeMessageData(data: unknown, isBinary: boolean): string | ArrayBuffer { if (!isBinary) { if (typeof data === "string") return data; diff --git a/packages/server/src/server/websocket-server.liveness.e2e.test.ts b/packages/server/src/server/websocket-server.liveness.e2e.test.ts new file mode 100644 index 000000000..cbbc0f705 --- /dev/null +++ b/packages/server/src/server/websocket-server.liveness.e2e.test.ts @@ -0,0 +1,196 @@ +import { expect, test } from "vitest"; +import { WebSocket, type RawData } from "ws"; +import { createTestPaseoDaemon, type TestPaseoDaemon } from "./test-utils/index.js"; +import { WSOutboundMessageSchema, type WSOutboundMessage } from "./messages.js"; + +const LARGE_REQUEST_BYTES = 512 * 1024; +const BURST_MESSAGE_COUNT = 32; +const TEST_TIMEOUT_MS = 30_000; + +interface SocketClose { + code: number; + reason: string; +} + +class ResumedPhysicalSocketSession { + private replacement: WebSocket | null = null; + + private constructor( + private readonly daemon: TestPaseoDaemon, + private readonly original: WebSocket, + ) {} + + static async launch(): Promise { + const daemon = await createTestPaseoDaemon(); + const original = await connectSocket(daemon.port, "stale-physical-socket"); + return new ResumedPhysicalSocketSession(daemon, original); + } + + async abandonOriginal(): Promise { + this.original.pause(); + } + + async resumeSameClient(): Promise { + this.replacement = await connectSocket(this.daemon.port, "stale-physical-socket"); + } + + async broadcastUntilOriginalCloses(): Promise { + const replacement = this.requireReplacement(); + const originalClose = waitForClose(this.original); + const finalRequestId = largeRequestId(BURST_MESSAGE_COUNT - 1); + const finalResponse = waitForMessage(replacement, (message) => { + return ( + message.type === "session" && + message.message.type === "pong" && + message.message.payload.requestId === finalRequestId + ); + }); + + for (let index = 0; index < BURST_MESSAGE_COUNT; index += 1) { + replacement.send( + JSON.stringify({ + type: "session", + message: { + type: "ping", + requestId: largeRequestId(index), + clientSentAt: index, + }, + }), + ); + } + + await finalResponse; + this.original.resume(); + return originalClose; + } + + async replacementRoundTrip(): Promise { + const replacement = this.requireReplacement(); + const requestId = "replacement-still-active"; + await sendAndWait( + replacement, + { + type: "session", + message: { type: "ping", requestId, clientSentAt: 1 }, + }, + (message) => + message.type === "session" && + message.message.type === "pong" && + message.message.payload.requestId === requestId, + ); + } + + async close(): Promise { + this.original.terminate(); + this.replacement?.terminate(); + await this.daemon.close(); + } + + private requireReplacement(): WebSocket { + if (!this.replacement) throw new Error("Replacement socket is not connected"); + return this.replacement; + } +} + +test( + "a resumed stale socket is bounded and removed without disrupting its replacement", + async () => { + const session = await ResumedPhysicalSocketSession.launch(); + try { + await session.abandonOriginal(); + await session.resumeSameClient(); + + const originalClose = await session.broadcastUntilOriginalCloses(); + + expect(originalClose).toEqual({ code: 1006, reason: "" }); + await session.replacementRoundTrip(); + } finally { + await session.close(); + } + }, + TEST_TIMEOUT_MS, +); + +async function connectSocket(port: number, clientId: string): Promise { + const socket = new WebSocket(`ws://127.0.0.1:${port}/ws`); + await waitForOpen(socket); + await sendAndWait( + socket, + { + type: "hello", + clientId, + clientType: "browser", + protocolVersion: 1, + }, + (message) => + message.type === "session" && + message.message.type === "status" && + message.message.payload.status === "server_info", + ); + await sendAndWait(socket, { type: "ping" }, (message) => message.type === "pong"); + return socket; +} + +function largeRequestId(index: number): string { + return `${index}:`.padEnd(LARGE_REQUEST_BYTES, "x"); +} + +function sendAndWait( + socket: WebSocket, + message: unknown, + matches: (message: WSOutboundMessage) => boolean, +): Promise { + const response = waitForMessage(socket, matches); + socket.send(JSON.stringify(message)); + return response; +} + +function waitForOpen(socket: WebSocket): Promise { + return new Promise((resolve, reject) => { + socket.once("open", resolve); + socket.once("error", reject); + }); +} + +function waitForClose(socket: WebSocket): Promise { + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + socket.off("close", onClose); + reject(new Error("Timed out waiting for WebSocket to close")); + }, TEST_TIMEOUT_MS); + const onClose = (code: number, reason: Buffer) => { + clearTimeout(timeout); + resolve({ code, reason: reason.toString() }); + }; + socket.once("close", onClose); + }); +} + +function waitForMessage( + socket: WebSocket, + matches: (message: WSOutboundMessage) => boolean, +): Promise { + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + cleanup(); + reject(new Error("Timed out waiting for WebSocket message")); + }, TEST_TIMEOUT_MS); + const onMessage = (data: RawData) => { + const parsed = WSOutboundMessageSchema.safeParse(JSON.parse(data.toString())); + if (!parsed.success || !matches(parsed.data)) return; + cleanup(); + resolve(parsed.data); + }; + const onClose = () => { + cleanup(); + reject(new Error("WebSocket closed before the expected message arrived")); + }; + const cleanup = () => { + clearTimeout(timeout); + socket.off("message", onMessage); + socket.off("close", onClose); + }; + socket.on("message", onMessage); + socket.on("close", onClose); + }); +} diff --git a/packages/server/src/server/websocket-server.ts b/packages/server/src/server/websocket-server.ts index 0959c0349..104eb08e9 100644 --- a/packages/server/src/server/websocket-server.ts +++ b/packages/server/src/server/websocket-server.ts @@ -89,6 +89,14 @@ import { } from "@getpaseo/protocol/browser-automation/capabilities"; import type { BrowserToolsBroker } from "./browser-tools/broker.js"; import type { DaemonRuntimeConfig } from "./session/daemon/daemon-session.js"; +import { + APPLICATION_SOCKET_LEASE_CHECK_INTERVAL_MS, + ApplicationSocketLease, + MAX_PHYSICAL_SOCKET_BUFFERED_BYTES, + outboundFrameByteLength, + physicalSocketHasCapacity, + sendBoundedPhysicalFrame, +} from "./websocket/physical-socket.js"; const WS_CLOSE_DAEMON_AUTH_FAILED = 4401; @@ -365,6 +373,7 @@ export interface WebSocketLike { bufferedAmount?: number; send: (data: string | Uint8Array | ArrayBuffer) => void; close: (code?: number, reason?: string) => void; + terminate?: () => void; on: (event: "message" | "close" | "error", listener: (...args: unknown[]) => void) => void; once: (event: "close" | "error", listener: (...args: unknown[]) => void) => void; } @@ -420,6 +429,12 @@ interface SocketSessionOptions { hubRelationships?: HubRelationshipManagement; } +interface ClosePhysicalSocketParams { + ws: WebSocketLike; + logMessage: string; + logFields?: Record; +} + const SLOW_REQUEST_THRESHOLD_MS = 500; const EXTERNAL_SESSION_DISCONNECT_GRACE_MS = 90_000; const HELLO_TIMEOUT_MS = 15_000; @@ -520,6 +535,8 @@ export class VoiceAssistantWebSocketServer { private readonly runtimeMetrics = new WebSocketRuntimeMetricsWindow(); private lastRuntimeMetricsSnapshot: WebSocketRuntimeDiagnosticPayload | null = null; private runtimeMetricsInterval: ReturnType | null = null; + private applicationSocketLeaseInterval: ReturnType | null = null; + private readonly applicationSocketLease = new ApplicationSocketLease(); private eventLoopDelayMonitor: ReturnType | null = null; private unsubscribeSpeechReadiness: (() => void) | null = null; private unsubscribeDaemonConfigChange: (() => void) | null = null; @@ -656,6 +673,7 @@ export class VoiceAssistantWebSocketServer { this.wss = this.createWebSocketServer(server, wsConfig, auth); this.startRuntimeMetricsInterval(); + this.startApplicationSocketLeaseInterval(); this.logger.info("WebSocket server initialized on /ws"); } @@ -743,6 +761,19 @@ export class VoiceAssistantWebSocketServer { (runtimeMetricsInterval as unknown as { unref?: () => void }).unref?.(); } + private startApplicationSocketLeaseInterval(): void { + const interval = setInterval(() => { + for (const ws of this.applicationSocketLease.listExpired()) { + this.closePhysicalSocket({ + ws, + logMessage: "Closing physical WebSocket with expired application lease", + }); + } + }, APPLICATION_SOCKET_LEASE_CHECK_INTERVAL_MS); + this.applicationSocketLeaseInterval = interval; + (interval as unknown as { unref?: () => void }).unref?.(); + } + // Main-loop stall visibility: terminal frames and agent traffic share one event // loop, so delay percentiles here are the ground truth for "the daemon is busy". private snapshotEventLoopDelay(): { p50Ms: number; p99Ms: number; maxMs: number } | null { @@ -819,17 +850,10 @@ export class VoiceAssistantWebSocketServer { } public broadcast(message: WSOutboundMessage): void { - const payload = JSON.stringify(message); - for (const [ws, connection] of this.sessions) { - if (connection.kind !== "trusted") { - continue; - } - // WebSocket.OPEN = 1 - if (ws.readyState === 1) { - ws.send(payload); - this.runtimeMetrics.recordOutboundMessage(message, ws.bufferedAmount); - } - } + const trustedSockets = [...this.sessions] + .filter(([, connection]) => connection.kind === "trusted") + .map(([ws]) => ws); + this.sendMessageToSockets(trustedSockets, message); } public listTrustedSessions(): Session[] { @@ -925,6 +949,11 @@ export class VoiceAssistantWebSocketServer { clearInterval(this.runtimeMetricsInterval); this.runtimeMetricsInterval = null; } + if (this.applicationSocketLeaseInterval) { + clearInterval(this.applicationSocketLeaseInterval); + this.applicationSocketLeaseInterval = null; + } + this.applicationSocketLease.clear(); this.flushRuntimeMetrics({ final: true }); this.eventLoopDelayMonitor?.disable(); this.eventLoopDelayMonitor = null; @@ -994,37 +1023,108 @@ export class VoiceAssistantWebSocketServer { } private sendToClient(ws: WebSocketLike, message: WSOutboundMessage): void { - // WebSocket.OPEN = 1. The check is a fast path; the socket can still - // transition to closed between here and ws.send(), so guard the send too — - // a synchronous throw here would propagate as an uncaughtException. - if (ws.readyState !== 1) { + this.sendMessageToSockets([ws], message); + } + + private sendMessageToSockets(sockets: Iterable, message: WSOutboundMessage): void { + const writableSockets = [...sockets].filter((ws) => this.ensureOutboundCapacity(ws, 0)); + if (writableSockets.length === 0) { return; } + + let payload: string; try { - ws.send(JSON.stringify(message)); - this.runtimeMetrics.recordOutboundMessage(message, ws.bufferedAmount); + payload = JSON.stringify(message); + } catch (err) { + this.logger.warn({ err }, "ws_serialize_failed"); + return; + } + + const payloadBytes = outboundFrameByteLength(payload); + for (const ws of writableSockets) { + this.sendFrameToClient(ws, payload, payloadBytes, () => { + this.runtimeMetrics.recordOutboundMessage(message, ws.bufferedAmount); + }); + } + } + + private sendBinaryToClient(ws: WebSocketLike, frame: Uint8Array): void { + this.sendFrameToClient(ws, frame, outboundFrameByteLength(frame), () => { + this.runtimeMetrics.recordOutboundBinaryFrame(ws.bufferedAmount); + }); + } + + private sendFrameToClient( + ws: WebSocketLike, + frame: string | Uint8Array, + frameBytes: number, + recordSent: () => void, + ): void { + try { + const sent = sendBoundedPhysicalFrame({ + socket: ws, + frame, + frameBytes, + onHighWater: () => this.closeAtOutboundHighWater(ws), + }); + if (sent) recordSent(); } catch (err) { this.logger.warn({ err }, "ws_send_failed"); } } - private sendBinaryToClient(ws: WebSocketLike, frame: Uint8Array): void { + private ensureOutboundCapacity(ws: WebSocketLike, frameBytes: number): boolean { + if (ws.readyState !== 1) return false; + if (physicalSocketHasCapacity(ws, frameBytes)) return true; + + this.closeAtOutboundHighWater(ws); + return false; + } + + private closeAtOutboundHighWater(ws: WebSocketLike): void { + this.closePhysicalSocket({ + ws, + logMessage: "Closing physical WebSocket at outbound high-water mark", + logFields: { + bufferedAmount: ws.bufferedAmount, + maxBufferedBytes: MAX_PHYSICAL_SOCKET_BUFFERED_BYTES, + }, + }); + } + + private closePhysicalSocket(params: ClosePhysicalSocketParams): void { + const { ws, logMessage, logFields } = params; + this.applicationSocketLease.release(ws); if (ws.readyState !== 1) { return; } + const identity = this.socketIdentities.get(ws); + this.logger.warn( + { + ...(identity ? toConnectionLogFields(identity) : {}), + ...logFields, + }, + logMessage, + ); try { - ws.send(frame); - this.runtimeMetrics.recordOutboundBinaryFrame(ws.bufferedAmount); + // A close frame queues behind application data, so it cannot enforce a + // hard memory cutoff. Production transports expose terminate(). + if (ws.terminate) { + ws.terminate(); + } else { + ws.close(); + } } catch (err) { - this.logger.warn({ err }, "ws_send_binary_failed"); + this.logger.warn( + { err, ...(identity ? toConnectionLogFields(identity) : {}) }, + "ws_close_failed", + ); } } private sendToConnection(connection: SessionConnection, message: WSOutboundMessage): void { const sockets = connection.kind === "trusted" ? connection.sockets : [connection.socket]; - for (const ws of sockets) { - this.sendToClient(ws, message); - } + this.sendMessageToSockets(sockets, message); } private sendBinaryToConnection(connection: SessionConnection, frame: Uint8Array): void { @@ -1521,6 +1621,7 @@ export class VoiceAssistantWebSocketServer { error?: Error; }, ): Promise { + this.applicationSocketLease.release(ws); const identity = this.socketIdentities.get(ws); const identityFields = identity ? toConnectionLogFields(identity) : {}; const pending = this.clearPendingConnection(ws); @@ -1837,6 +1938,8 @@ export class VoiceAssistantWebSocketServer { return; } + this.applicationSocketLease.renew(ws); + const activeConnection = this.sessions.get(ws); const pendingConnection = this.pendingConnections.get(ws); const log = @@ -1872,6 +1975,7 @@ export class VoiceAssistantWebSocketServer { this.recordInboundMessageType(message.type); if (message.type === "ping") { + this.applicationSocketLease.claim(ws); this.sendToClient(ws, { type: "pong" }); return; } diff --git a/packages/server/src/server/websocket/encrypted-relay-socket.test.ts b/packages/server/src/server/websocket/encrypted-relay-socket.test.ts new file mode 100644 index 000000000..3add059fc --- /dev/null +++ b/packages/server/src/server/websocket/encrypted-relay-socket.test.ts @@ -0,0 +1,120 @@ +import { EventEmitter } from "node:events"; +import { expect, test } from "vitest"; +import { MAX_PHYSICAL_SOCKET_BUFFERED_BYTES } from "./physical-socket.js"; +import { + createEncryptedRelaySocket, + type EncryptedRelayChannel, +} from "./encrypted-relay-socket.js"; + +class BlockingChannel implements EncryptedRelayChannel { + readonly sent: Array = []; + readonly closes: Array<{ code?: number; reason?: string }> = []; + private resolveSend: (() => void) | null = null; + + setState(state: "open"): void { + expect(state).toBe("open"); + } + + send(data: string | ArrayBuffer): Promise { + this.sent.push(data); + return new Promise((resolve) => { + this.resolveSend = resolve; + }); + } + + close(code?: number, reason?: string): void { + this.closes.push({ code, reason }); + } + + drain(): void { + this.resolveSend?.(); + } +} + +test("the encrypted send queue terminates its physical transport at the hard bound", async () => { + const channel = new BlockingChannel(); + let terminations = 0; + const socket = createEncryptedRelaySocket({ + channel, + emitter: new EventEmitter(), + getTransportBufferedAmount: () => 0, + terminateTransport: () => { + terminations += 1; + }, + }); + + socket.send(new Uint8Array(5 * 1024 * 1024)); + expect(channel.sent).toHaveLength(1); + expect(socket.bufferedAmount).toBeGreaterThan(5 * 1024 * 1024); + + socket.send(new Uint8Array(2 * 1024 * 1024)); + + expect(channel.sent).toHaveLength(1); + expect(terminations).toBe(1); + expect(channel.closes).toEqual([]); + expect(socket.readyState).toBe(3); + + channel.drain(); + await Promise.resolve(); +}); + +test("underlying relay backpressure rejects binary before encryption and terminates physically", () => { + const channel = new BlockingChannel(); + let terminations = 0; + const socket = createEncryptedRelaySocket({ + channel, + emitter: new EventEmitter(), + getTransportBufferedAmount: () => MAX_PHYSICAL_SOCKET_BUFFERED_BYTES - 1, + terminateTransport: () => { + terminations += 1; + }, + }); + + socket.send(new Uint8Array(1)); + + expect(channel.sent).toEqual([]); + expect(channel.closes).toEqual([]); + expect(terminations).toBe(1); +}); + +test("pending encryption and underlying relay backpressure share one hard bound", () => { + const channel = new BlockingChannel(); + let transportBufferedAmount = 3 * 1024 * 1024; + let terminations = 0; + const socket = createEncryptedRelaySocket({ + channel, + emitter: new EventEmitter(), + getTransportBufferedAmount: () => transportBufferedAmount, + terminateTransport: () => { + terminations += 1; + }, + }); + + socket.send(new Uint8Array(3 * 1024 * 1024)); + expect(channel.sent).toHaveLength(1); + + transportBufferedAmount = 4 * 1024 * 1024; + socket.send(new Uint8Array(1)); + + expect(channel.sent).toHaveLength(1); + expect(terminations).toBe(1); +}); + +test("explicit encrypted-socket termination forcibly terminates the relay transport", () => { + const channel = new BlockingChannel(); + let terminations = 0; + const socket = createEncryptedRelaySocket({ + channel, + emitter: new EventEmitter(), + getTransportBufferedAmount: () => 0, + terminateTransport: () => { + terminations += 1; + }, + }); + + socket.terminate(); + + expect(terminations).toBe(1); + expect(channel.closes).toEqual([]); + expect(socket.readyState).toBe(3); +}); diff --git a/packages/server/src/server/websocket/encrypted-relay-socket.ts b/packages/server/src/server/websocket/encrypted-relay-socket.ts new file mode 100644 index 000000000..4b00f8200 --- /dev/null +++ b/packages/server/src/server/websocket/encrypted-relay-socket.ts @@ -0,0 +1,100 @@ +import { EventEmitter } from "node:events"; +import { MAX_PHYSICAL_SOCKET_BUFFERED_BYTES, outboundFrameByteLength } from "./physical-socket.js"; + +// NaCl adds a 24-byte nonce and 16-byte authenticator before base64 encoding. +const ENCRYPTED_FRAME_OVERHEAD_BYTES = 40; + +export interface EncryptedRelayChannel { + setState: (state: "open") => void; + send: (data: string | ArrayBuffer) => Promise; + close: (code?: number, reason?: string) => void; +} + +export interface EncryptedRelaySocket { + readonly readyState: number; + readonly bufferedAmount: number; + send: (data: string | Uint8Array | ArrayBuffer) => void; + close: (code?: number, reason?: string) => void; + terminate: () => void; + on: (event: "message" | "close" | "error", listener: (...args: unknown[]) => void) => void; + once: (event: "close" | "error", listener: (...args: unknown[]) => void) => void; +} + +export function createEncryptedRelaySocket(params: { + channel: EncryptedRelayChannel; + emitter: EventEmitter; + getTransportBufferedAmount: () => number | undefined; + terminateTransport: () => void; +}): EncryptedRelaySocket { + const { channel, emitter, getTransportBufferedAmount, terminateTransport } = params; + let readyState = 1; + let pendingEncryptedBytes = 0; + + channel.setState("open"); + + const terminate = () => { + if (readyState === 3) return; + readyState = 3; + terminateTransport(); + }; + + const close = (code?: number, reason?: string) => { + if (readyState === 3) return; + readyState = 3; + channel.close(code, reason); + }; + + emitter.on("close", () => { + readyState = 3; + }); + + return { + get readyState() { + return readyState; + }, + get bufferedAmount() { + return pendingEncryptedBytes + (getTransportBufferedAmount() ?? 0); + }, + send: (data) => { + if (readyState !== 1) return; + const outbound = normalizeRelaySendPayload(data); + const outboundBytes = encryptedRelayFrameByteLength(outbound); + const queuedBytes = pendingEncryptedBytes + (getTransportBufferedAmount() ?? 0); + if (queuedBytes + outboundBytes > MAX_PHYSICAL_SOCKET_BUFFERED_BYTES) { + terminate(); + return; + } + pendingEncryptedBytes += outboundBytes; + void channel + .send(outbound) + .catch((error) => { + emitter.emit("error", error); + }) + .finally(() => { + pendingEncryptedBytes -= outboundBytes; + }); + }, + close, + terminate, + on: (event, listener) => { + emitter.on(event, listener); + }, + once: (event, listener) => { + emitter.once(event, listener); + }, + }; +} + +function normalizeRelaySendPayload(data: string | Uint8Array | ArrayBuffer): string | ArrayBuffer { + if (typeof data === "string") return data; + if (data instanceof ArrayBuffer) return data; + const view = new Uint8Array(data.buffer, data.byteOffset, data.byteLength); + const out = new Uint8Array(view.byteLength); + out.set(view); + return out.buffer; +} + +function encryptedRelayFrameByteLength(data: string | ArrayBuffer): number { + const encryptedBytes = outboundFrameByteLength(data) + ENCRYPTED_FRAME_OVERHEAD_BYTES; + return 4 * Math.ceil(encryptedBytes / 3); +} diff --git a/packages/server/src/server/websocket/physical-socket.test.ts b/packages/server/src/server/websocket/physical-socket.test.ts new file mode 100644 index 000000000..9b334cc11 --- /dev/null +++ b/packages/server/src/server/websocket/physical-socket.test.ts @@ -0,0 +1,68 @@ +import { expect, test } from "vitest"; +import { + APPLICATION_SOCKET_LEASE_MS, + ApplicationSocketLease, + MAX_PHYSICAL_SOCKET_BUFFERED_BYTES, + sendBoundedPhysicalFrame, +} from "./physical-socket.js"; + +test("sockets remain exempt until they send an application ping", () => { + let now = 0; + const lease = new ApplicationSocketLease(() => now); + const legacySocket = {}; + now = APPLICATION_SOCKET_LEASE_MS * 10; + + expect(lease.listExpired()).toEqual([]); + lease.renew(legacySocket); + expect(lease.listExpired()).toEqual([]); +}); + +test("inbound activity renews a claimed lease", () => { + let now = 0; + const lease = new ApplicationSocketLease(() => now); + const applicationSocket = {}; + lease.claim(applicationSocket); + + now = APPLICATION_SOCKET_LEASE_MS - 1; + lease.renew(applicationSocket); + now += APPLICATION_SOCKET_LEASE_MS - 1; + expect(lease.listExpired()).toEqual([]); + + now += 1; + expect(lease.listExpired()).toEqual([applicationSocket]); + lease.release(applicationSocket); + expect(lease.listExpired()).toEqual([]); +}); + +test("an application ping claims a socket lease", () => { + let now = 0; + const lease = new ApplicationSocketLease(() => now); + const rawSocket = {}; + + lease.claim(rawSocket); + now = APPLICATION_SOCKET_LEASE_MS; + + expect(lease.listExpired()).toEqual([rawSocket]); +}); + +test("the shared physical send boundary rejects binary above the hard bound", () => { + const sent: Array = []; + let terminated = false; + const socket = { + readyState: 1, + bufferedAmount: MAX_PHYSICAL_SOCKET_BUFFERED_BYTES - 1, + send: (data: string | Uint8Array | ArrayBuffer) => sent.push(data), + }; + + const accepted = sendBoundedPhysicalFrame({ + socket, + frame: new Uint8Array(2), + onHighWater: () => { + terminated = true; + }, + }); + + expect(accepted).toBe(false); + expect(sent).toEqual([]); + expect(terminated).toBe(true); +}); diff --git a/packages/server/src/server/websocket/physical-socket.ts b/packages/server/src/server/websocket/physical-socket.ts new file mode 100644 index 000000000..eed29ae35 --- /dev/null +++ b/packages/server/src/server/websocket/physical-socket.ts @@ -0,0 +1,78 @@ +// Terminal streams begin snapshot catch-up at 4 MiB. The physical socket gets +// another 4 MiB to recover before the daemon enforces the hard memory bound. +export const MAX_PHYSICAL_SOCKET_BUFFERED_BYTES = 8 * 1024 * 1024; +// Current clients ping every 10 seconds. Four delayed cycles fit inside the +// lease without making an abandoned application socket linger for minutes. +export const APPLICATION_SOCKET_LEASE_MS = 45_000; +export const APPLICATION_SOCKET_LEASE_CHECK_INTERVAL_MS = 10_000; + +type Clock = () => number; + +export class ApplicationSocketLease { + private readonly deadlines = new Map(); + + constructor(private readonly clock: Clock = Date.now) {} + + claim(socket: TSocket): void { + this.deadlines.set(socket, this.clock() + APPLICATION_SOCKET_LEASE_MS); + } + + renew(socket: TSocket): void { + if (this.deadlines.has(socket)) { + this.claim(socket); + } + } + + release(socket: TSocket): void { + this.deadlines.delete(socket); + } + + listExpired(): TSocket[] { + const now = this.clock(); + const expired: TSocket[] = []; + for (const [socket, deadline] of this.deadlines) { + if (deadline > now) continue; + expired.push(socket); + } + return expired; + } + + clear(): void { + this.deadlines.clear(); + } +} + +export function outboundFrameByteLength(data: string | Uint8Array | ArrayBuffer): number { + if (typeof data === "string") return Buffer.byteLength(data); + return data.byteLength; +} + +interface BoundedPhysicalSocket { + readyState: number; + bufferedAmount?: number; + send: (data: string | Uint8Array | ArrayBuffer) => void; +} + +export function physicalSocketHasCapacity( + socket: Pick, + frameBytes: number, +): boolean { + if (typeof socket.bufferedAmount !== "number") return true; + return socket.bufferedAmount + frameBytes <= MAX_PHYSICAL_SOCKET_BUFFERED_BYTES; +} + +export function sendBoundedPhysicalFrame(params: { + socket: BoundedPhysicalSocket; + frame: string | Uint8Array | ArrayBuffer; + frameBytes?: number; + onHighWater: () => void; +}): boolean { + const { socket, frame, frameBytes = outboundFrameByteLength(frame), onHighWater } = params; + if (socket.readyState !== 1) return false; + if (!physicalSocketHasCapacity(socket, frameBytes)) { + onHighWater(); + return false; + } + socket.send(frame); + return true; +}