From 70ce38068f876de78d2ce465b9d6f590ac75e6c4 Mon Sep 17 00:00:00 2001 From: Bohdan Triapitsyn Date: Fri, 24 Apr 2026 12:29:02 +0300 Subject: [PATCH] Split message stream websocket bridges --- .../server/lib/event-stream/DOCUMENTATION.md | 7 +- .../lib/event-stream/directory-ws-bridge.js | 185 +++++++++ .../lib/event-stream/global-ws-bridge.js | 206 ++++++++++ .../web/server/lib/event-stream/runtime.js | 372 ++---------------- 4 files changed, 421 insertions(+), 349 deletions(-) create mode 100644 packages/web/server/lib/event-stream/directory-ws-bridge.js create mode 100644 packages/web/server/lib/event-stream/global-ws-bridge.js diff --git a/packages/web/server/lib/event-stream/DOCUMENTATION.md b/packages/web/server/lib/event-stream/DOCUMENTATION.md index 921d9c9b..f69938c6 100644 --- a/packages/web/server/lib/event-stream/DOCUMENTATION.md +++ b/packages/web/server/lib/event-stream/DOCUMENTATION.md @@ -6,9 +6,11 @@ This module contains the OpenChamber message-stream WebSocket protocol and runti ## Entrypoints and structure - `packages/web/server/lib/event-stream/index.js`: public entrypoint re-exporting protocol and runtime helpers. - `packages/web/server/lib/event-stream/global-hub.js`: shared global upstream SSE hub for server-side subscribers and browser WS fan-out. +- `packages/web/server/lib/event-stream/global-ws-bridge.js`: browser-facing global WS bridge that subscribes clients to the shared global hub. +- `packages/web/server/lib/event-stream/directory-ws-bridge.js`: browser-facing per-directory WS bridge that owns one scoped upstream reader per connection. - `packages/web/server/lib/event-stream/protocol.js`: path constants, SSE envelope parsing, and WebSocket frame serialization helpers. - `packages/web/server/lib/event-stream/upstream-reader.js`: reusable upstream SSE reader with event-id tracking, stall recovery, and reconnect handling. -- `packages/web/server/lib/event-stream/runtime.js`: WebSocket server runtime, upgrade handling, shared global upstream hub, per-directory upstream reader orchestration, and global event broadcasting. +- `packages/web/server/lib/event-stream/runtime.js`: thin WebSocket server runtime for upgrade handling and path dispatch to the global/directory bridges. - `packages/web/server/lib/event-stream/protocol.test.js`: unit tests for protocol helpers. - `packages/web/server/lib/event-stream/upstream-reader.test.js`: unit tests for upstream SSE reader behavior. - `packages/web/server/lib/event-stream/runtime.test.js`: unit tests for runtime-side broadcaster behavior. @@ -44,10 +46,11 @@ This module contains the OpenChamber message-stream WebSocket protocol and runti - Global synthetic events such as `openchamber:session-status`, `openchamber:session-activity`, `openchamber:notification`, and `openchamber:heartbeat` are preserved on the WS path, but heartbeat frames are emitted only while an upstream SSE stream is actively attached. - Global UI broadcasts are fan-out capable across both SSE and WS clients. - The reusable upstream reader centralizes SSE fetch/parsing/reconnect behavior for the WS runtime and OpenCode watcher. Additional event consumers should move to it only with parity tests for their lifecycle and error semantics. +- Browser transport concerns live in the WS bridge modules; server-side global stream ownership lives in `global-hub.js`. ## Notes for contributors - Keep protocol helpers pure and small so they can be unit tested without spinning up a server. -- Keep runtime wiring in this module instead of `packages/web/server/index.js` unless the logic is strictly route-local glue. +- Keep `runtime.js` focused on WebSocket upgrade and endpoint dispatch. Put global browser-client lifecycle in `global-ws-bridge.js`, directory stream lifecycle in `directory-ws-bridge.js`, and upstream stream sharing in `global-hub.js`. - Do not change upstream OpenCode transport assumptions here; OpenCode remains SSE-based. - Keep global replay bounded; do not turn it into an unbounded event log. diff --git a/packages/web/server/lib/event-stream/directory-ws-bridge.js b/packages/web/server/lib/event-stream/directory-ws-bridge.js new file mode 100644 index 00000000..fd338a50 --- /dev/null +++ b/packages/web/server/lib/event-stream/directory-ws-bridge.js @@ -0,0 +1,185 @@ +import { sendMessageStreamWsEvent, sendMessageStreamWsFrame } from './protocol.js'; +import { createUpstreamSseReader } from './upstream-reader.js'; + +function shouldTriggerUpstreamHealthCheck(upstream) { + if (!upstream) { + return true; + } + + if (!upstream.body) { + return upstream.ok || upstream.status >= 500; + } + + return upstream.status >= 500; +} + +export function acceptDirectoryMessageStreamWsConnection({ + socket, + requestedLastEventId, + requestedDirectory, + buildOpenCodeUrl, + getOpenCodeAuthHeaders, + processForwardedEventPayload, + wsClients, + triggerHealthCheck, + heartbeatIntervalMs, + upstreamStallTimeoutMs, + upstreamReconnectDelayMs, + fetchImpl, +}) { + const controller = new AbortController(); + let upstreamConnected = false; + let streamReady = false; + let reader = null; + + const cleanup = () => { + if (!controller.signal.aborted) { + controller.abort(); + } + reader?.stop(); + wsClients.delete(socket); + }; + + const pingInterval = setInterval(() => { + if (socket.readyState !== 1) { + return; + } + + try { + socket.ping(); + } catch { + } + }, heartbeatIntervalMs); + + const heartbeatInterval = setInterval(() => { + if (!upstreamConnected) { + return; + } + + sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' }); + }, heartbeatIntervalMs); + + socket.on('close', () => { + clearInterval(pingInterval); + clearInterval(heartbeatInterval); + upstreamConnected = false; + cleanup(); + }); + + socket.on('error', () => { + void 0; + }); + + const run = async () => { + const forwardEvent = ({ envelope, payload }) => { + const directory = requestedDirectory || envelope?.directory || 'global'; + + sendMessageStreamWsEvent(socket, payload, { + directory, + eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined, + }); + + processForwardedEventPayload(payload, (syntheticPayload) => { + sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' }); + }); + }; + + try { + let buildUrlFailed = false; + const closeWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => { + sendMessageStreamWsFrame(socket, { type: 'error', message }); + socket.close(1011, closeReason); + if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) { + triggerHealthCheck?.(); + } + reader?.stop(); + cleanup(); + }; + + reader = createUpstreamSseReader({ + initialLastEventId: requestedLastEventId, + signal: controller.signal, + stallTimeoutMs: upstreamStallTimeoutMs, + reconnectDelayMs: upstreamReconnectDelayMs, + fetchImpl, + buildUrl: () => { + buildUrlFailed = false; + let targetUrl; + try { + targetUrl = new URL(buildOpenCodeUrl('/event', '')); + } catch { + buildUrlFailed = true; + throw new Error('OpenCode service unavailable'); + } + + if (requestedDirectory) { + targetUrl.searchParams.set('directory', requestedDirectory); + } + + return targetUrl; + }, + getHeaders: getOpenCodeAuthHeaders, + onConnect() { + if (!streamReady) { + sendMessageStreamWsFrame(socket, { + type: 'ready', + scope: 'directory', + }); + streamReady = true; + } + + upstreamConnected = true; + }, + onDisconnect() { + upstreamConnected = false; + }, + onEvent: forwardEvent, + onError(error) { + if (controller.signal.aborted) { + return; + } + + if (!streamReady) { + if (error?.type === 'upstream_unavailable') { + closeWithInitialError({ + message: `OpenCode event stream unavailable (${error.status})`, + closeReason: 'OpenCode event stream unavailable', + triggerHealthCheckFor: error.response, + }); + return; + } + + closeWithInitialError({ + message: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', + closeReason: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', + triggerHealthCheckFor: !buildUrlFailed, + }); + return; + } + + if (error?.type === 'stream_error') { + console.warn('Message stream WS proxy error:', error.error); + } + }, + }); + + await reader.start(); + } catch (error) { + if (!controller.signal.aborted) { + console.warn('Message stream WS proxy error:', error); + sendMessageStreamWsFrame(socket, { type: 'error', message: 'Message stream proxy error' }); + socket.close(1011, 'Message stream proxy error'); + } + } finally { + cleanup(); + try { + if (socket.readyState === 1 || socket.readyState === 0) { + socket.close(); + } + } catch { + } + } + }; + + void run(); +} diff --git a/packages/web/server/lib/event-stream/global-ws-bridge.js b/packages/web/server/lib/event-stream/global-ws-bridge.js new file mode 100644 index 00000000..9ab65c84 --- /dev/null +++ b/packages/web/server/lib/event-stream/global-ws-bridge.js @@ -0,0 +1,206 @@ +import { sendMessageStreamWsEvent, sendMessageStreamWsFrame } from './protocol.js'; + +function shouldTriggerUpstreamHealthCheck(upstream) { + if (!upstream) { + return true; + } + + if (!upstream.body) { + return upstream.ok || upstream.status >= 500; + } + + return upstream.status >= 500; +} + +export function createGlobalMessageStreamWsBridge({ + globalHub, + ownsGlobalHub, + wsClients, + processForwardedEventPayload, + triggerHealthCheck, + heartbeatIntervalMs, +}) { + const clients = new Set(); + const clientLastEventIds = new Map(); + const readyClients = new Set(); + + const removeClient = (socket) => { + clients.delete(socket); + clientLastEventIds.delete(socket); + readyClients.delete(socket); + wsClients.delete(socket); + }; + + const replayEvents = (socket, requestedLastEventId) => { + for (const entry of globalHub.replayAfter(requestedLastEventId)) { + const sent = sendMessageStreamWsEvent(socket, entry.payload, { + directory: entry.directory, + eventId: entry.eventId, + }); + if (!sent) { + removeClient(socket); + return; + } + } + }; + + const markReady = (socket, requestedLastEventId) => { + if (socket.readyState !== 1) { + return; + } + + const sent = sendMessageStreamWsFrame(socket, { + type: 'ready', + scope: 'global', + }); + if (!sent) { + removeClient(socket); + return; + } + + readyClients.add(socket); + wsClients.add(socket); + replayEvents(socket, requestedLastEventId); + }; + + const stopHubIfUnused = () => { + if (ownsGlobalHub && clients.size === 0) { + globalHub.stop(); + } + }; + + const closeClientsWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => { + for (const socket of Array.from(clients)) { + sendMessageStreamWsFrame(socket, { type: 'error', message }); + try { + socket.close(1011, closeReason); + } catch { + } + removeClient(socket); + } + + if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) { + triggerHealthCheck?.(); + } + + if (ownsGlobalHub) { + globalHub.stop(); + } + }; + + const unsubscribeEvent = globalHub.subscribeEvent(({ payload, directory, eventId }) => { + for (const socket of Array.from(clients)) { + if (!readyClients.has(socket)) { + continue; + } + const sent = sendMessageStreamWsEvent(socket, payload, { + directory, + eventId, + }); + if (!sent) { + removeClient(socket); + } + } + + processForwardedEventPayload(payload, (syntheticPayload) => { + for (const socket of Array.from(clients)) { + if (!readyClients.has(socket)) { + continue; + } + const sent = sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' }); + if (!sent) { + removeClient(socket); + } + } + }); + }); + + const unsubscribeStatus = globalHub.subscribeStatus((status) => { + if (status.type === 'connect') { + for (const socket of Array.from(clients)) { + if (!readyClients.has(socket)) { + markReady(socket, clientLastEventIds.get(socket) ?? ''); + } + } + return; + } + + if (status.type === 'initial-error') { + const error = status.error; + if (error?.type === 'upstream_unavailable') { + closeClientsWithInitialError({ + message: `OpenCode event stream unavailable (${error.status})`, + closeReason: 'OpenCode event stream unavailable', + triggerHealthCheckFor: error.response, + }); + return; + } + + closeClientsWithInitialError({ + message: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', + closeReason: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', + triggerHealthCheckFor: !status.buildUrlFailed, + }); + return; + } + + if (status.type === 'error' && status.error?.type === 'stream_error') { + console.warn('Message stream WS proxy error:', status.error.error); + } + }); + + const accept = (socket, { requestedLastEventId = '' } = {}) => { + const pingInterval = setInterval(() => { + if (socket.readyState !== 1) { + return; + } + + try { + socket.ping(); + } catch { + } + }, heartbeatIntervalMs); + + const heartbeatInterval = setInterval(() => { + if (!globalHub.isConnected()) { + return; + } + + sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' }); + }, heartbeatIntervalMs); + + socket.on('close', () => { + clearInterval(pingInterval); + clearInterval(heartbeatInterval); + removeClient(socket); + stopHubIfUnused(); + }); + + socket.on('error', () => { + void 0; + }); + + clients.add(socket); + clientLastEventIds.set(socket, requestedLastEventId); + globalHub.start(); + if (globalHub.isConnected()) { + markReady(socket, requestedLastEventId); + } + }; + + const close = () => { + unsubscribeEvent(); + unsubscribeStatus(); + if (ownsGlobalHub) { + globalHub.stop(); + } + for (const socket of Array.from(clients)) { + removeClient(socket); + } + }; + + return { + accept, + close, + }; +} diff --git a/packages/web/server/lib/event-stream/runtime.js b/packages/web/server/lib/event-stream/runtime.js index 21df5d26..fc619308 100644 --- a/packages/web/server/lib/event-stream/runtime.js +++ b/packages/web/server/lib/event-stream/runtime.js @@ -6,27 +6,15 @@ import { MESSAGE_STREAM_GLOBAL_WS_PATH, MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS, sendMessageStreamWsEvent, - sendMessageStreamWsFrame, } from './protocol.js'; import { createGlobalMessageStreamHub } from './global-hub.js'; +import { createGlobalMessageStreamWsBridge } from './global-ws-bridge.js'; +import { acceptDirectoryMessageStreamWsConnection } from './directory-ws-bridge.js'; import { DEFAULT_UPSTREAM_RECONNECT_DELAY_MS, DEFAULT_UPSTREAM_STALL_TIMEOUT_MS, - createUpstreamSseReader, } from './upstream-reader.js'; -function shouldTriggerUpstreamHealthCheck(upstream) { - if (!upstream) { - return true; - } - - if (!upstream.body) { - return upstream.ok || upstream.status >= 500; - } - - return upstream.status >= 500; -} - export function createGlobalUiEventBroadcaster({ sseClients, wsClients, @@ -91,143 +79,15 @@ export function createMessageStreamWsRuntime({ upstreamReconnectDelayMs, }); - const globalClients = new Set(); - const globalClientLastEventIds = new Map(); - const globalReadyClients = new Set(); - - const replayGlobalEvents = (socket, requestedLastEventId) => { - for (const entry of globalHub.replayAfter(requestedLastEventId)) { - const sent = sendMessageStreamWsEvent(socket, entry.payload, { - directory: entry.directory, - eventId: entry.eventId, - }); - if (!sent) { - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - return; - } - } - }; - - const markGlobalClientReady = (socket, requestedLastEventId) => { - if (socket.readyState !== 1) { - return; - } - - const sent = sendMessageStreamWsFrame(socket, { - type: 'ready', - scope: 'global', - }); - if (!sent) { - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - return; - } - - globalReadyClients.add(socket); - wsClients.add(socket); - replayGlobalEvents(socket, requestedLastEventId); - }; - - const closeGlobalClientsWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => { - for (const socket of Array.from(globalClients)) { - sendMessageStreamWsFrame(socket, { type: 'error', message }); - try { - socket.close(1011, closeReason); - } catch { - } - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - } - - if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) { - triggerHealthCheck?.(); - } - - if (ownsGlobalHub) { - globalHub.stop(); - } - }; - - const unsubscribeGlobalEvent = globalHub.subscribeEvent(({ envelope, payload, directory, eventId }) => { - for (const socket of Array.from(globalClients)) { - if (!globalReadyClients.has(socket)) { - continue; - } - const sent = sendMessageStreamWsEvent(socket, payload, { - directory, - eventId, - }); - if (!sent) { - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - } - } - - processForwardedEventPayload(payload, (syntheticPayload) => { - for (const socket of Array.from(globalClients)) { - if (!globalReadyClients.has(socket)) { - continue; - } - const sent = sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' }); - if (!sent) { - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - } - } - }); + const globalBridge = createGlobalMessageStreamWsBridge({ + globalHub, + ownsGlobalHub, + wsClients, + processForwardedEventPayload, + triggerHealthCheck, + heartbeatIntervalMs, }); - const unsubscribeGlobalStatus = globalHub.subscribeStatus((status) => { - if (status.type === 'connect') { - for (const socket of Array.from(globalClients)) { - if (!globalReadyClients.has(socket)) { - markGlobalClientReady(socket, globalClientLastEventIds.get(socket) ?? ''); - } - } - return; - } - - if (status.type === 'initial-error') { - const error = status.error; - if (error?.type === 'upstream_unavailable') { - closeGlobalClientsWithInitialError({ - message: `OpenCode event stream unavailable (${error.status})`, - closeReason: 'OpenCode event stream unavailable', - triggerHealthCheckFor: error.response, - }); - return; - } - - closeGlobalClientsWithInitialError({ - message: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', - closeReason: status.buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', - triggerHealthCheckFor: !status.buildUrlFailed, - }); - return; - } - - if (status.type === 'error' && status.error?.type === 'stream_error') { - console.warn('Message stream WS proxy error:', status.error.error); - } - }); - - const stopGlobalHubIfUnused = () => { - if (ownsGlobalHub && globalClients.size === 0) { - globalHub.stop(); - } - }; - wsServer.on('connection', (socket, req) => { const rawUrl = typeof req?.url === 'string' ? req.url : MESSAGE_STREAM_GLOBAL_WS_PATH; const pathname = parseRequestPathname(rawUrl); @@ -237,204 +97,26 @@ export function createMessageStreamWsRuntime({ const requestedDirectory = requestUrl.searchParams.get('directory')?.trim() || ''; if (isGlobalStream) { - const pingInterval = setInterval(() => { - if (socket.readyState !== 1) { - return; - } - - try { - socket.ping(); - } catch { - } - }, heartbeatIntervalMs); - - const heartbeatInterval = setInterval(() => { - if (!globalHub.isConnected()) { - return; - } - - sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' }); - }, heartbeatIntervalMs); - - socket.on('close', () => { - clearInterval(pingInterval); - clearInterval(heartbeatInterval); - globalClients.delete(socket); - globalClientLastEventIds.delete(socket); - globalReadyClients.delete(socket); - wsClients.delete(socket); - stopGlobalHubIfUnused(); + globalBridge.accept(socket, { + requestedLastEventId, }); - - socket.on('error', () => { - void 0; - }); - - globalClients.add(socket); - globalClientLastEventIds.set(socket, requestedLastEventId); - globalHub.start(); - if (globalHub.isConnected()) { - markGlobalClientReady(socket, requestedLastEventId); - } return; } - const controller = new AbortController(); - let upstreamConnected = false; - let streamReady = false; - let reader = null; - const cleanup = () => { - if (!controller.signal.aborted) { - controller.abort(); - } - reader?.stop(); - wsClients.delete(socket); - }; - - const pingInterval = setInterval(() => { - if (socket.readyState !== 1) { - return; - } - - try { - socket.ping(); - } catch { - } - }, heartbeatIntervalMs); - - const heartbeatInterval = setInterval(() => { - if (!upstreamConnected) { - return; - } - - sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' }); - }, heartbeatIntervalMs); - - socket.on('close', () => { - clearInterval(pingInterval); - clearInterval(heartbeatInterval); - upstreamConnected = false; - cleanup(); + acceptDirectoryMessageStreamWsConnection({ + socket, + requestedLastEventId, + requestedDirectory, + buildOpenCodeUrl, + getOpenCodeAuthHeaders, + processForwardedEventPayload, + wsClients, + triggerHealthCheck, + heartbeatIntervalMs, + upstreamStallTimeoutMs, + upstreamReconnectDelayMs, + fetchImpl, }); - - socket.on('error', () => { - void 0; - }); - - const run = async () => { - const forwardEvent = ({ envelope, payload }) => { - const directory = isGlobalStream - ? (typeof envelope?.directory === 'string' && envelope.directory.length > 0 ? envelope.directory : 'global') - : (requestedDirectory || envelope?.directory || 'global'); - - sendMessageStreamWsEvent(socket, payload, { - directory, - eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined, - }); - - processForwardedEventPayload(payload, (syntheticPayload) => { - sendMessageStreamWsEvent(socket, syntheticPayload, { directory: 'global' }); - }); - }; - - try { - let buildUrlFailed = false; - const closeWithInitialError = ({ message, closeReason = message, triggerHealthCheckFor = null }) => { - sendMessageStreamWsFrame(socket, { type: 'error', message }); - socket.close(1011, closeReason); - if (triggerHealthCheckFor === true || (triggerHealthCheckFor && shouldTriggerUpstreamHealthCheck(triggerHealthCheckFor))) { - triggerHealthCheck?.(); - } - reader?.stop(); - cleanup(); - }; - - reader = createUpstreamSseReader({ - initialLastEventId: requestedLastEventId, - signal: controller.signal, - stallTimeoutMs: upstreamStallTimeoutMs, - reconnectDelayMs: upstreamReconnectDelayMs, - fetchImpl, - buildUrl: () => { - buildUrlFailed = false; - let targetUrl; - try { - targetUrl = new URL(buildOpenCodeUrl('/event', '')); - } catch { - buildUrlFailed = true; - throw new Error('OpenCode service unavailable'); - } - - if (requestedDirectory) { - targetUrl.searchParams.set('directory', requestedDirectory); - } - - return targetUrl; - }, - getHeaders: getOpenCodeAuthHeaders, - onConnect() { - if (!streamReady) { - sendMessageStreamWsFrame(socket, { - type: 'ready', - scope: 'directory', - }); - streamReady = true; - } - - upstreamConnected = true; - }, - onDisconnect() { - upstreamConnected = false; - }, - onEvent: forwardEvent, - onError(error) { - if (controller.signal.aborted) { - return; - } - - if (!streamReady) { - if (error?.type === 'upstream_unavailable') { - closeWithInitialError({ - message: `OpenCode event stream unavailable (${error.status})`, - closeReason: 'OpenCode event stream unavailable', - triggerHealthCheckFor: error.response, - }); - return; - } - - closeWithInitialError({ - message: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', - closeReason: buildUrlFailed ? 'OpenCode service unavailable' : 'Failed to connect to OpenCode event stream', - triggerHealthCheckFor: !buildUrlFailed, - }); - return; - } - - if (error?.type === 'stream_error') { - console.warn('Message stream WS proxy error:', error.error); - } - }, - }); - - await reader.start(); - } catch (error) { - if (!controller.signal.aborted) { - console.warn('Message stream WS proxy error:', error); - sendMessageStreamWsFrame(socket, { type: 'error', message: 'Message stream proxy error' }); - socket.close(1011, 'Message stream proxy error'); - } - } finally { - cleanup(); - try { - if (socket.readyState === 1 || socket.readyState === 0) { - socket.close(); - } - } catch { - } - } - }; - - void run(); }); const upgradeHandler = (req, socket, head) => { @@ -476,11 +158,7 @@ export function createMessageStreamWsRuntime({ wsServer, async close() { server.off('upgrade', upgradeHandler); - unsubscribeGlobalEvent(); - unsubscribeGlobalStatus(); - if (ownsGlobalHub) { - globalHub.stop(); - } + globalBridge.close(); try { for (const client of wsServer.clients) {