From 8bfbea1e0f5b95e756eeb84e7ccbe59b84845910 Mon Sep 17 00:00:00 2001 From: metaphorics <152830360+metaphorics@users.noreply.github.com> Date: Thu, 2 Jul 2026 17:25:50 +0900 Subject: [PATCH] fix(collab): add WebSocket backpressure with ordered pending-send drain CollabSocket.send previously called ws.send() whenever the socket was OPEN without checking bufferedAmount, allowing the WebSocket send buffer to grow unbounded under a slow relay or guest. The fix adds 64 KiB high / 32 KiB low backpressure thresholds: new frames are queued when bufferedAmount is above the high watermark, and a 25 ms drain timer retries flushing queued frames once the buffer drops below the low watermark, preserving send order via the existing #sendChain. Overflow still drops frames once #pendingSends reaches MAX_PENDING_SENDS, and the drain timer is cleared on close, reconnect, and fatal failure. Verified with `bun test packages/coding-agent/test/collab` (64 pass, 0 fail). Closes #4248 --- packages/coding-agent/CHANGELOG.md | 5 ++ .../coding-agent/src/collab/relay-client.ts | 74 +++++++++++++++++-- 2 files changed, 74 insertions(+), 5 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 11098023c..41d8512a1 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,11 @@ ## [Unreleased] +### Fixed + +- Apply WebSocket send backpressure with ordered drain retry to prevent unbounded `bufferedAmount` growth ([#4248](https://github.com/can1357/oh-my-pi/issues/4248)) + + ## [16.3.0] - 2026-07-02 ### Added diff --git a/packages/coding-agent/src/collab/relay-client.ts b/packages/coding-agent/src/collab/relay-client.ts index 130e92dec..b64e525bc 100644 --- a/packages/coding-agent/src/collab/relay-client.ts +++ b/packages/coding-agent/src/collab/relay-client.ts @@ -21,6 +21,9 @@ const BACKOFF_BASE_MS = 1_000; const BACKOFF_MAX_MS = 30_000; /** Max enveloped frames buffered while a reconnect is pending; overflow is dropped. */ const MAX_PENDING_SENDS = 256; +const WS_BACKPRESSURE_THRESHOLD = 64 * 1024; +const WS_BACKPRESSURE_DRAIN_THRESHOLD = 32 * 1024; +const WS_BACKPRESSURE_DRAIN_RETRY_MS = 25; export interface CollabSocketOptions { /** wss://host[:port]/r/ — no query string. */ @@ -40,6 +43,7 @@ export class CollabSocket { readonly #opts: CollabSocketOptions; #ws: WebSocket | null = null; #retryTimer: NodeJS.Timeout | undefined; + #backpressureDrainTimer: NodeJS.Timeout | undefined; #attempt = 0; /** Terminal state: intentional close or fatal failure. Cleared by connect(). */ #closed = false; @@ -72,28 +76,84 @@ export class CollabSocket { logger.debug("collab: dropping frame, socket closed", { t: frame.t }); return; } + const openWs = this.#ws; + if (openWs && openWs.readyState === WebSocket.OPEN) this.#drainPendingSends(openWs); const sealed = await seal(this.#opts.key, frame); const envelope = packEnvelope(targetPeer, sealed); const ws = this.#ws; if (ws && ws.readyState === WebSocket.OPEN) { + if (this.#pendingSends.length > 0) { + this.#enqueuePendingSend(envelope, frame.t); + if (ws.bufferedAmount < WS_BACKPRESSURE_DRAIN_THRESHOLD) { + this.#drainPendingSends(ws); + } else { + this.#scheduleBackpressureDrain(ws); + } + return; + } + if (ws.bufferedAmount >= WS_BACKPRESSURE_THRESHOLD) { + this.#enqueuePendingSend(envelope, frame.t); + this.#scheduleBackpressureDrain(ws); + return; + } ws.send(envelope); return; } - if (this.#pendingSends.length >= MAX_PENDING_SENDS) { - logger.debug("collab: dropping frame, reconnect buffer full", { t: frame.t }); - return; - } - this.#pendingSends.push(envelope); + this.#enqueuePendingSend(envelope, frame.t); }) .catch((err: unknown) => { logger.debug("collab: send failed", { error: String(err) }); }); } + #enqueuePendingSend(envelope: Uint8Array, frameType: CollabFrame["t"]): void { + if (this.#pendingSends.length >= MAX_PENDING_SENDS) { + logger.debug("collab: dropping frame, reconnect buffer full", { t: frameType }); + return; + } + this.#pendingSends.push(envelope); + } + + #drainPendingSends(ws: WebSocket): void { + while ( + this.#pendingSends.length > 0 && + ws.readyState === WebSocket.OPEN && + ws.bufferedAmount < WS_BACKPRESSURE_DRAIN_THRESHOLD + ) { + const envelope = this.#pendingSends.shift(); + if (!envelope) return; + ws.send(envelope); + } + } + + #scheduleBackpressureDrain(ws: WebSocket): void { + if (this.#backpressureDrainTimer !== undefined) return; + this.#backpressureDrainTimer = setTimeout(() => { + this.#backpressureDrainTimer = undefined; + this.#sendChain = this.#sendChain + .then(async () => { + if (this.#closed || this.#ws !== ws || ws.readyState !== WebSocket.OPEN) return; + this.#drainPendingSends(ws); + if (this.#pendingSends.length > 0) this.#scheduleBackpressureDrain(ws); + }) + .catch((err: unknown) => { + logger.debug("collab: backpressure drain failed", { error: String(err) }); + }); + }, WS_BACKPRESSURE_DRAIN_RETRY_MS); + } + + #clearBackpressureDrain(): void { + if (this.#backpressureDrainTimer !== undefined) { + clearTimeout(this.#backpressureDrainTimer); + this.#backpressureDrainTimer = undefined; + } + } + /** Intentional close: clears any retry timer, suppresses reconnect. A later connect() starts fresh. */ close(): void { const hadActivity = this.#ws !== null || this.#retryTimer !== undefined; this.#clearRetry(); + this.#clearBackpressureDrain(); const wasClosed = this.#closed; this.#closed = true; this.#pendingSends.length = 0; @@ -110,6 +170,7 @@ export class CollabSocket { } #openSocket(): void { + this.#clearBackpressureDrain(); const ws = new WebSocket(`${this.#opts.wsUrl}?role=${this.#opts.role}`); ws.binaryType = "arraybuffer"; this.#ws = ws; @@ -129,6 +190,7 @@ export class CollabSocket { }; ws.onclose = (event: CloseEvent) => { if (this.#ws !== ws) return; + this.#clearBackpressureDrain(); this.#ws = null; this.#handleClose(event.code, event.reason); }; @@ -167,6 +229,7 @@ export class CollabSocket { #handleClose(code: number, reason: string): void { if (this.#closed) return; + this.#clearBackpressureDrain(); const fatalReason = FATAL_CLOSE_REASONS[code]; if (fatalReason !== undefined) { this.#closed = true; @@ -186,6 +249,7 @@ export class CollabSocket { this.#pendingSends.length = 0; const ws = this.#ws; this.#ws = null; + this.#clearBackpressureDrain(); if (ws) { try { ws.close(1000);