diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index 83741854..331ef9e6 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -20,6 +20,8 @@ export type QueuedEvent = { export type FlushHandler = (events: QueuedEvent[]) => void const FLUSH_FRAME_MS = 33 +const BACKPRESSURE_FLUSH_FRAME_MS = 200 +const BACKPRESSURE_MODE_MS = 10_000 const STREAM_YIELD_MS = 8 const DEFAULT_RECONNECT_DELAY_MS = 250 const DEFAULT_HEARTBEAT_TIMEOUT_MS = 30_000 @@ -44,7 +46,7 @@ export type EventPipelineInput = { } type MessageStreamWsFrame = { - type: "ready" | "event" | "error" + type: "ready" | "event" | "error" | "backpressure" payload?: unknown eventId?: string directory?: string @@ -254,7 +256,8 @@ export function createEventPipeline(input: EventPipelineInput) { const d = getOrCreateDir(directory) if (d.timer) return const elapsed = Date.now() - d.last - d.timer = setTimeout(() => flushDir(directory), Math.max(0, FLUSH_FRAME_MS - elapsed)) + const flushFrameMs = Date.now() < backpressureUntil ? BACKPRESSURE_FLUSH_FRAME_MS : FLUSH_FRAME_MS + d.timer = setTimeout(() => flushDir(directory), Math.max(0, flushFrameMs - elapsed)) } const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)) @@ -269,6 +272,7 @@ export function createEventPipeline(input: EventPipelineInput) { let activeTransport: "ws" | "sse" = transport === "ws" ? "ws" : "sse" let attemptAbortReason: AttemptAbortReason = null let consecutiveFailures = 0 + let backpressureUntil = 0 const notifyDisconnected = (reason: string) => { if (disconnected) { @@ -489,6 +493,11 @@ export function createEventPipeline(input: EventPipelineInput) { return } + if (frame.type === "backpressure") { + backpressureUntil = Date.now() + BACKPRESSURE_MODE_MS + return + } + if (frame.type !== "event") { return } diff --git a/packages/web/server/lib/event-stream/global-hub.js b/packages/web/server/lib/event-stream/global-hub.js index 1bc98abe..86f28b3f 100644 --- a/packages/web/server/lib/event-stream/global-hub.js +++ b/packages/web/server/lib/event-stream/global-hub.js @@ -1,6 +1,8 @@ import { createUpstreamSseReader } from './upstream-reader.js'; -export const MESSAGE_STREAM_GLOBAL_REPLAY_LIMIT = 512; +// Raised from 512 → 2048 to improve recovery after brief disconnects during +// long-running agent sessions where many events accumulate quickly. +export const MESSAGE_STREAM_GLOBAL_REPLAY_LIMIT = 2048; export function createGlobalMessageStreamHub({ buildOpenCodeUrl, diff --git a/packages/web/server/lib/event-stream/protocol.js b/packages/web/server/lib/event-stream/protocol.js index d92ea485..394fa502 100644 --- a/packages/web/server/lib/event-stream/protocol.js +++ b/packages/web/server/lib/event-stream/protocol.js @@ -3,7 +3,13 @@ export const MESSAGE_STREAM_DIRECTORY_WS_PATH = '/api/event/ws'; export const MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS = 15 * 1000; // Per-client pending outbound WS buffer, not a payload or stream-size limit. // Healthy clients stay near 0; this only trips when a client is far behind. -export const MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES = 4 * 1024 * 1024; +// Raised from 4 MB → 16 MB to tolerate bursts during long agent sessions +// (e.g. ultrawork / multi-tool loops) where the browser briefly falls behind. +export const MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES = 16 * 1024 * 1024; + +// Threshold at which we emit a backpressure warning frame so the client can +// proactively start shedding low-priority updates before the hard disconnect. +export const MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES = 12 * 1024 * 1024; export function parseSseEventEnvelope(block) { if (!block || typeof block !== 'string') { @@ -67,7 +73,9 @@ export function sendMessageStreamWsFrame(socket, payload) { return false; } - if (typeof socket.bufferedAmount === 'number' && socket.bufferedAmount > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { + const buffered = typeof socket.bufferedAmount === 'number' ? socket.bufferedAmount : 0; + + if (buffered > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { try { socket.close(1013, 'Message stream client is too slow'); } catch { @@ -77,13 +85,36 @@ export function sendMessageStreamWsFrame(socket, payload) { try { socket.send(JSON.stringify(payload)); - if (typeof socket.bufferedAmount === 'number' && socket.bufferedAmount > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { + const bufferedAfter = typeof socket.bufferedAmount === 'number' ? socket.bufferedAmount : 0; + if (bufferedAfter > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { try { socket.close(1013, 'Message stream client is too slow'); } catch { } return false; } + + // Emit a one-shot backpressure warning when the buffer is building up. + // The flag prevents sending repeated warnings that would themselves + // increase the buffer. It resets once the buffer drains below the + // threshold. + if (bufferedAfter > MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES) { + if (!socket._ocBackpressureWarned) { + socket._ocBackpressureWarned = true; + try { + socket.send(JSON.stringify({ + type: 'backpressure', + bufferedBytes: bufferedAfter, + maxBytes: MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES, + })); + } catch { + // Best-effort warning — ignore send failures. + } + } + } else if (socket._ocBackpressureWarned) { + socket._ocBackpressureWarned = false; + } + return true; } catch { return false; diff --git a/packages/web/server/lib/event-stream/protocol.test.js b/packages/web/server/lib/event-stream/protocol.test.js index a3b1417e..ec8c7aaf 100644 --- a/packages/web/server/lib/event-stream/protocol.test.js +++ b/packages/web/server/lib/event-stream/protocol.test.js @@ -4,6 +4,7 @@ import { MESSAGE_STREAM_DIRECTORY_WS_PATH, MESSAGE_STREAM_GLOBAL_WS_PATH, MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES, + MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES, parseSseEventEnvelope, sendMessageStreamWsEvent, sendMessageStreamWsFrame, @@ -88,6 +89,73 @@ describe('event stream protocol helpers', () => { }); }); + it('emits a backpressure warning when buffer exceeds the warn threshold', () => { + const sentPayloads = []; + const socket = { + readyState: 1, + bufferedAmount: MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES + 1, + send(payload) { + sentPayloads.push(payload); + }, + }; + + const sent = sendMessageStreamWsFrame(socket, { type: 'test' }); + + expect(sent).toBe(true); + expect(sentPayloads).toHaveLength(2); + + const warning = JSON.parse(sentPayloads[1]); + expect(warning.type).toBe('backpressure'); + expect(warning.bufferedBytes).toBeGreaterThan(0); + expect(warning.maxBytes).toBe(MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES); + }); + + it('does not repeat backpressure warnings while still above threshold', () => { + const sentPayloads = []; + const socket = { + readyState: 1, + bufferedAmount: MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES + 1, + send(payload) { + sentPayloads.push(payload); + }, + }; + + sendMessageStreamWsFrame(socket, { type: 'test1' }); + sendMessageStreamWsFrame(socket, { type: 'test2' }); + + // First call: data + warning = 2 sends. Second call: data only = 1 send. + expect(sentPayloads).toHaveLength(3); + expect(JSON.parse(sentPayloads[1]).type).toBe('backpressure'); + expect(JSON.parse(sentPayloads[2])).toEqual({ type: 'test2' }); + }); + + it('resets backpressure warning flag when buffer drains', () => { + const sentPayloads = []; + const socket = { + readyState: 1, + bufferedAmount: MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES + 1, + send(payload) { + sentPayloads.push(payload); + }, + }; + + sendMessageStreamWsFrame(socket, { type: 'first' }); + expect(socket._ocBackpressureWarned).toBe(true); + + // Buffer drains + socket.bufferedAmount = 100; + sendMessageStreamWsFrame(socket, { type: 'recovered' }); + expect(socket._ocBackpressureWarned).toBe(false); + + // Buffer spikes again — warning should fire again + socket.bufferedAmount = MESSAGE_STREAM_WS_BACKPRESSURE_WARN_BYTES + 1; + sendMessageStreamWsFrame(socket, { type: 'again' }); + const backpressureFrames = sentPayloads + .map((p) => JSON.parse(p)) + .filter((p) => p.type === 'backpressure'); + expect(backpressureFrames).toHaveLength(2); + }); + it('serializes event frames with routing metadata', () => { let rawPayload = null; const socket = {