From 9d9af7e2636d769485074a8883e3e29d3e09d3b6 Mon Sep 17 00:00:00 2001 From: Bohdan Triapitsyn Date: Thu, 23 Apr 2026 23:40:02 +0300 Subject: [PATCH] Improve event stream resilience --- packages/ui/src/lib/openchamberEvents.ts | 38 +++++++++ .../src/sync/__tests__/event-pipeline.test.js | 79 +++++++++++++++++++ packages/ui/src/sync/event-pipeline.ts | 1 + .../web/server/lib/event-stream/protocol.js | 18 +++++ .../server/lib/event-stream/protocol.test.js | 25 ++++++ packages/web/server/lib/opencode/proxy.js | 42 +++++++++- .../web/server/lib/scheduled-tasks/routes.js | 15 ++++ packages/web/server/opencode-proxy.test.js | 24 +++++- 8 files changed, 240 insertions(+), 2 deletions(-) diff --git a/packages/ui/src/lib/openchamberEvents.ts b/packages/ui/src/lib/openchamberEvents.ts index 1e6e4c58..5468345b 100644 --- a/packages/ui/src/lib/openchamberEvents.ts +++ b/packages/ui/src/lib/openchamberEvents.ts @@ -12,10 +12,20 @@ type Listener = (event: OpenChamberEvent) => void; let eventSource: EventSource | null = null; let reconnectTimer: ReturnType | null = null; +let heartbeatTimer: ReturnType | null = null; let reconnectAttempt = 0; const listeners = new Set(); const MAX_RECONNECT_DELAY_MS = 30_000; +const HEARTBEAT_TIMEOUT_MS = 45_000; + +const clearHeartbeatTimer = () => { + if (!heartbeatTimer) { + return; + } + clearTimeout(heartbeatTimer); + heartbeatTimer = null; +}; const scheduleReconnect = () => { if (reconnectTimer || listeners.size === 0) { @@ -30,12 +40,24 @@ const scheduleReconnect = () => { }; const cleanupSource = () => { + clearHeartbeatTimer(); if (eventSource) { eventSource.close(); } eventSource = null; }; +const resetHeartbeatTimer = () => { + clearHeartbeatTimer(); + if (listeners.size === 0) { + return; + } + heartbeatTimer = setTimeout(() => { + cleanupSource(); + scheduleReconnect(); + }, HEARTBEAT_TIMEOUT_MS); +}; + const parseEnvelope = (raw: string): { type: string; properties: unknown } | null => { if (!raw || raw.trim().length === 0) { return null; @@ -60,6 +82,10 @@ const dispatchFromEnvelope = (envelope: { type: string; properties: unknown }) = return; } + if (envelope.type === 'openchamber:heartbeat') { + return; + } + if (envelope.type !== 'openchamber:scheduled-task-ran') { return; } @@ -93,10 +119,22 @@ const connect = () => { if (typeof window === 'undefined' || listeners.size === 0) { return; } + if (typeof EventSource !== 'function') { + return; + } + + if (eventSource && eventSource.readyState !== EventSource.CLOSED) { + return; + } + cleanupSource(); const source = new EventSource('/api/openchamber/events'); + source.onopen = () => { + resetHeartbeatTimer(); + }; source.onmessage = (event) => { + resetHeartbeatTimer(); const envelope = parseEnvelope(event.data); if (!envelope) { return; diff --git a/packages/ui/src/sync/__tests__/event-pipeline.test.js b/packages/ui/src/sync/__tests__/event-pipeline.test.js index f508e785..f01c724d 100644 --- a/packages/ui/src/sync/__tests__/event-pipeline.test.js +++ b/packages/ui/src/sync/__tests__/event-pipeline.test.js @@ -667,6 +667,85 @@ describe('createEventPipeline', () => { ]); }); + it('passes the last websocket event id when falling back to SSE', async () => { + installDomStubs(); + globalThis.WebSocket = FakeWebSocket; + const originalConsoleError = console.error; + console.error = () => {}; + + let releaseStream; + const hold = new Promise((resolve) => { + releaseStream = resolve; + }); + + const eventOptions = []; + const received = []; + const sdk = { + global: { + event: async (options) => { + eventOptions.push(options); + return { + stream: (async function* () { + yield { + payload: { + type: 'server.connected', + properties: {}, + }, + }; + await hold; + })(), + }; + }, + }, + }; + + const delivered = new Promise((resolve) => { + const { cleanup } = createEventPipeline({ + sdk, + transport: 'auto', + reconnectDelayMs: 0, + wsReadyTimeoutMs: 20, + onEvent: (directory, payload) => { + received.push({ directory, payload }); + if (payload.type !== 'server.connected') { + return; + } + cleanup(); + releaseStream(); + resolve(); + }, + }); + }); + + try { + await Promise.resolve(); + + const firstSocket = FakeWebSocket.instances[0]; + firstSocket.emitOpen(); + firstSocket.emitMessage({ type: 'ready', scope: 'global' }); + firstSocket.emitMessage({ + type: 'event', + eventId: 'evt-1', + directory: '/tmp/project', + payload: { + type: 'session.status', + properties: { + sessionID: 'session-1', + }, + }, + }); + firstSocket.emitClose(); + + await new Promise((resolve) => setTimeout(resolve, 40)); + await delivered; + + expect(eventOptions[0]?.headers?.['Last-Event-ID']).toBe('evt-1'); + expect(received.some((entry) => entry.payload.type === 'server.connected')).toBe(true); + } finally { + console.error = originalConsoleError; + } + }); + it('marks the pipeline disconnected on heartbeat timeout and recovers on the next websocket connect', async () => { installDomStubs(); globalThis.WebSocket = FakeWebSocket; diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index 8d4d39a5..760eed0d 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -339,6 +339,7 @@ export function createEventPipeline(input: EventPipelineInput) { const runSseAttempt = async (signal: AbortSignal) => { const events = await sdk.global.event({ signal, + ...(lastEventId && lastEventId.length > 0 ? { headers: { "Last-Event-ID": lastEventId } } : {}), onSseError: (error: unknown) => { if (isAbortError(error)) return if (streamErrorLogged) return diff --git a/packages/web/server/lib/event-stream/protocol.js b/packages/web/server/lib/event-stream/protocol.js index 033859db..d92ea485 100644 --- a/packages/web/server/lib/event-stream/protocol.js +++ b/packages/web/server/lib/event-stream/protocol.js @@ -1,6 +1,9 @@ export const MESSAGE_STREAM_GLOBAL_WS_PATH = '/api/global/event/ws'; 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; export function parseSseEventEnvelope(block) { if (!block || typeof block !== 'string') { @@ -64,8 +67,23 @@ export function sendMessageStreamWsFrame(socket, payload) { return false; } + if (typeof socket.bufferedAmount === 'number' && socket.bufferedAmount > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { + try { + socket.close(1013, 'Message stream client is too slow'); + } catch { + } + return false; + } + try { socket.send(JSON.stringify(payload)); + if (typeof socket.bufferedAmount === 'number' && socket.bufferedAmount > MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES) { + try { + socket.close(1013, 'Message stream client is too slow'); + } catch { + } + return 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 ca551e32..a3b1417e 100644 --- a/packages/web/server/lib/event-stream/protocol.test.js +++ b/packages/web/server/lib/event-stream/protocol.test.js @@ -3,6 +3,7 @@ import { describe, expect, it } from 'vitest'; import { MESSAGE_STREAM_DIRECTORY_WS_PATH, MESSAGE_STREAM_GLOBAL_WS_PATH, + MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES, parseSseEventEnvelope, sendMessageStreamWsEvent, sendMessageStreamWsFrame, @@ -63,6 +64,30 @@ describe('event stream protocol helpers', () => { expect(rawPayload).toBe('{"type":"ready"}'); }); + it('closes slow websocket clients instead of adding more buffered data', () => { + let closeCall = null; + let sendCalls = 0; + const socket = { + readyState: 1, + bufferedAmount: MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES + 1, + send() { + sendCalls += 1; + }, + close(code, reason) { + closeCall = { code, reason }; + }, + }; + + const sent = sendMessageStreamWsFrame(socket, { type: 'ready' }); + + expect(sent).toBe(false); + expect(sendCalls).toBe(0); + expect(closeCall).toEqual({ + code: 1013, + reason: 'Message stream client is too slow', + }); + }); + it('serializes event frames with routing metadata', () => { let rawPayload = null; const socket = { diff --git a/packages/web/server/lib/opencode/proxy.js b/packages/web/server/lib/opencode/proxy.js index 83928230..69a6c9ba 100644 --- a/packages/web/server/lib/opencode/proxy.js +++ b/packages/web/server/lib/opencode/proxy.js @@ -6,6 +6,43 @@ import { shouldForwardProxyResponseHeader, } from '../../proxy-headers.js'; +export const waitForSseDrain = (res, signal) => new Promise((resolve) => { + if (signal?.aborted || res.writableEnded || res.destroyed) { + resolve(); + return; + } + + const cleanup = () => { + res.off?.('drain', onDone); + res.off?.('close', onDone); + res.off?.('error', onDone); + signal?.removeEventListener?.('abort', onDone); + }; + const onDone = () => { + cleanup(); + resolve(); + }; + + res.once?.('drain', onDone); + res.once?.('close', onDone); + res.once?.('error', onDone); + signal?.addEventListener?.('abort', onDone, { once: true }); +}); + +export const writeSseChunkWithBackpressure = async (res, value, signal) => { + if (!value || value.length === 0 || signal?.aborted || res.writableEnded || res.destroyed) { + return false; + } + + const flushed = res.write(value); + if (flushed !== false) { + return true; + } + + await waitForSseDrain(res, signal); + return !signal?.aborted && !res.writableEnded && !res.destroyed; +}; + export const registerOpenCodeProxy = (app, deps) => { const { fs, @@ -131,7 +168,10 @@ export const registerOpenCodeProxy = (app, deps) => { break; } if (value && value.length > 0) { - res.write(value); + const canContinue = await writeSseChunkWithBackpressure(res, value, abortController.signal); + if (!canContinue) { + break; + } } } diff --git a/packages/web/server/lib/scheduled-tasks/routes.js b/packages/web/server/lib/scheduled-tasks/routes.js index 98d7d36d..5676f0ea 100644 --- a/packages/web/server/lib/scheduled-tasks/routes.js +++ b/packages/web/server/lib/scheduled-tasks/routes.js @@ -209,7 +209,22 @@ export const registerScheduledTaskRoutes = (app, dependencies) => { } catch { } + const heartbeat = setInterval(() => { + try { + writeSseEvent(res, { + type: 'openchamber:heartbeat', + properties: { + timestamp: Date.now(), + }, + }); + } catch { + clearInterval(heartbeat); + clients.delete(res); + } + }, 25_000); + req.on('close', () => { + clearInterval(heartbeat); clients.delete(res); }); }); diff --git a/packages/web/server/opencode-proxy.test.js b/packages/web/server/opencode-proxy.test.js index 39ae3b97..8b4e8564 100644 --- a/packages/web/server/opencode-proxy.test.js +++ b/packages/web/server/opencode-proxy.test.js @@ -1,8 +1,9 @@ import { afterEach, describe, expect, it } from 'vitest'; +import { EventEmitter } from 'node:events'; import express from 'express'; import path from 'path'; -import { registerOpenCodeProxy } from './lib/opencode/proxy.js'; +import { registerOpenCodeProxy, writeSseChunkWithBackpressure } from './lib/opencode/proxy.js'; const listen = (app, host = '127.0.0.1') => new Promise((resolve, reject) => { const server = app.listen(0, host, () => resolve(server)); @@ -81,6 +82,27 @@ describe('OpenCode proxy SSE forwarding', () => { expect(seenAuthorization).toBe('Bearer test-token'); }); + it('waits for drain when writing to a slow SSE response', async () => { + const writes = []; + const res = new EventEmitter(); + res.writableEnded = false; + res.destroyed = false; + res.write = (value) => { + writes.push(value); + return false; + }; + const controller = new AbortController(); + + const write = writeSseChunkWithBackpressure(res, Buffer.from('data: {"ok":true}\n\n'), controller.signal); + + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(writes).toHaveLength(1); + + res.emit('drain'); + + await expect(write).resolves.toBe(true); + }); + it('routes generic API requests through external OpenCode base URL', async () => { const upstream = express(); upstream.get('/config/providers', (_req, res) => {