Improve event stream resilience

This commit is contained in:
Bohdan Triapitsyn
2026-04-24 12:30:11 +03:00
parent f0c3d5e326
commit 9d9af7e263
8 changed files with 240 additions and 2 deletions
+38
View File
@@ -12,10 +12,20 @@ type Listener = (event: OpenChamberEvent) => void;
let eventSource: EventSource | null = null; let eventSource: EventSource | null = null;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null; let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let heartbeatTimer: ReturnType<typeof setTimeout> | null = null;
let reconnectAttempt = 0; let reconnectAttempt = 0;
const listeners = new Set<Listener>(); const listeners = new Set<Listener>();
const MAX_RECONNECT_DELAY_MS = 30_000; const MAX_RECONNECT_DELAY_MS = 30_000;
const HEARTBEAT_TIMEOUT_MS = 45_000;
const clearHeartbeatTimer = () => {
if (!heartbeatTimer) {
return;
}
clearTimeout(heartbeatTimer);
heartbeatTimer = null;
};
const scheduleReconnect = () => { const scheduleReconnect = () => {
if (reconnectTimer || listeners.size === 0) { if (reconnectTimer || listeners.size === 0) {
@@ -30,12 +40,24 @@ const scheduleReconnect = () => {
}; };
const cleanupSource = () => { const cleanupSource = () => {
clearHeartbeatTimer();
if (eventSource) { if (eventSource) {
eventSource.close(); eventSource.close();
} }
eventSource = null; 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 => { const parseEnvelope = (raw: string): { type: string; properties: unknown } | null => {
if (!raw || raw.trim().length === 0) { if (!raw || raw.trim().length === 0) {
return null; return null;
@@ -60,6 +82,10 @@ const dispatchFromEnvelope = (envelope: { type: string; properties: unknown }) =
return; return;
} }
if (envelope.type === 'openchamber:heartbeat') {
return;
}
if (envelope.type !== 'openchamber:scheduled-task-ran') { if (envelope.type !== 'openchamber:scheduled-task-ran') {
return; return;
} }
@@ -93,10 +119,22 @@ const connect = () => {
if (typeof window === 'undefined' || listeners.size === 0) { if (typeof window === 'undefined' || listeners.size === 0) {
return; return;
} }
if (typeof EventSource !== 'function') {
return;
}
if (eventSource && eventSource.readyState !== EventSource.CLOSED) {
return;
}
cleanupSource(); cleanupSource();
const source = new EventSource('/api/openchamber/events'); const source = new EventSource('/api/openchamber/events');
source.onopen = () => {
resetHeartbeatTimer();
};
source.onmessage = (event) => { source.onmessage = (event) => {
resetHeartbeatTimer();
const envelope = parseEnvelope(event.data); const envelope = parseEnvelope(event.data);
if (!envelope) { if (!envelope) {
return; return;
@@ -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 () => { it('marks the pipeline disconnected on heartbeat timeout and recovers on the next websocket connect', async () => {
installDomStubs(); installDomStubs();
globalThis.WebSocket = FakeWebSocket; globalThis.WebSocket = FakeWebSocket;
+1
View File
@@ -339,6 +339,7 @@ export function createEventPipeline(input: EventPipelineInput) {
const runSseAttempt = async (signal: AbortSignal) => { const runSseAttempt = async (signal: AbortSignal) => {
const events = await sdk.global.event({ const events = await sdk.global.event({
signal, signal,
...(lastEventId && lastEventId.length > 0 ? { headers: { "Last-Event-ID": lastEventId } } : {}),
onSseError: (error: unknown) => { onSseError: (error: unknown) => {
if (isAbortError(error)) return if (isAbortError(error)) return
if (streamErrorLogged) return if (streamErrorLogged) return
@@ -1,6 +1,9 @@
export const MESSAGE_STREAM_GLOBAL_WS_PATH = '/api/global/event/ws'; 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_DIRECTORY_WS_PATH = '/api/event/ws';
export const MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS = 15 * 1000; 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) { export function parseSseEventEnvelope(block) {
if (!block || typeof block !== 'string') { if (!block || typeof block !== 'string') {
@@ -64,8 +67,23 @@ export function sendMessageStreamWsFrame(socket, payload) {
return false; 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 { try {
socket.send(JSON.stringify(payload)); 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; return true;
} catch { } catch {
return false; return false;
@@ -3,6 +3,7 @@ import { describe, expect, it } from 'vitest';
import { import {
MESSAGE_STREAM_DIRECTORY_WS_PATH, MESSAGE_STREAM_DIRECTORY_WS_PATH,
MESSAGE_STREAM_GLOBAL_WS_PATH, MESSAGE_STREAM_GLOBAL_WS_PATH,
MESSAGE_STREAM_WS_MAX_BUFFERED_BYTES,
parseSseEventEnvelope, parseSseEventEnvelope,
sendMessageStreamWsEvent, sendMessageStreamWsEvent,
sendMessageStreamWsFrame, sendMessageStreamWsFrame,
@@ -63,6 +64,30 @@ describe('event stream protocol helpers', () => {
expect(rawPayload).toBe('{"type":"ready"}'); 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', () => { it('serializes event frames with routing metadata', () => {
let rawPayload = null; let rawPayload = null;
const socket = { const socket = {
+41 -1
View File
@@ -6,6 +6,43 @@ import {
shouldForwardProxyResponseHeader, shouldForwardProxyResponseHeader,
} from '../../proxy-headers.js'; } 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) => { export const registerOpenCodeProxy = (app, deps) => {
const { const {
fs, fs,
@@ -131,7 +168,10 @@ export const registerOpenCodeProxy = (app, deps) => {
break; break;
} }
if (value && value.length > 0) { if (value && value.length > 0) {
res.write(value); const canContinue = await writeSseChunkWithBackpressure(res, value, abortController.signal);
if (!canContinue) {
break;
}
} }
} }
@@ -209,7 +209,22 @@ export const registerScheduledTaskRoutes = (app, dependencies) => {
} catch { } 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', () => { req.on('close', () => {
clearInterval(heartbeat);
clients.delete(res); clients.delete(res);
}); });
}); });
+23 -1
View File
@@ -1,8 +1,9 @@
import { afterEach, describe, expect, it } from 'vitest'; import { afterEach, describe, expect, it } from 'vitest';
import { EventEmitter } from 'node:events';
import express from 'express'; import express from 'express';
import path from 'path'; 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 listen = (app, host = '127.0.0.1') => new Promise((resolve, reject) => {
const server = app.listen(0, host, () => resolve(server)); const server = app.listen(0, host, () => resolve(server));
@@ -81,6 +82,27 @@ describe('OpenCode proxy SSE forwarding', () => {
expect(seenAuthorization).toBe('Bearer test-token'); 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 () => { it('routes generic API requests through external OpenCode base URL', async () => {
const upstream = express(); const upstream = express();
upstream.get('/config/providers', (_req, res) => { upstream.get('/config/providers', (_req, res) => {