diff --git a/packages/electron/main.mjs b/packages/electron/main.mjs index e9d92718..9cc499b6 100644 --- a/packages/electron/main.mjs +++ b/packages/electron/main.mjs @@ -42,6 +42,14 @@ log.transports.file.maxSize = 5 * 1024 * 1024; log.transports.file.level = 'info'; log.transports.console.level = isDev ? 'debug' : 'warn'; +// The in-process web server runs in this same Node process and uses plain +// `console.log/warn/error`. Without piping console through electron-log, +// that output never lands in ~/Library/Logs/OpenChamber/main.log and we +// can't diagnose issues (e.g. OpenCode lifecycle, SSE disconnects) after +// the fact. Route all console calls through electron-log so server-side +// diagnostics are persisted. +Object.assign(console, log.functions); + const LOG_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000; try { const logPath = log.transports.file.getFile().path; diff --git a/packages/ui/src/components/chat/ChatInput.tsx b/packages/ui/src/components/chat/ChatInput.tsx index a81bef1e..277acd14 100644 --- a/packages/ui/src/components/chat/ChatInput.tsx +++ b/packages/ui/src/components/chat/ChatInput.tsx @@ -1519,11 +1519,7 @@ const ChatInputComponent: React.FC = ({ onOpenSettings, scrollTo } else if (commandName === 'compact' && currentSessionId) { try { - if (!useConfigStore.getState().isConnected) { - const reason = useConfigStore.getState().lastDisconnectReason; - const suffix = reason ? ` (${reason})` : " (never connected)"; - throw new Error(`Connection lost${suffix}. Please wait for reconnection.`); - } + await sessionActions.waitForConnectionOrThrow(); const { opencodeClient } = await import('@/lib/opencode/client'); const sdk = opencodeClient.getSdkClient(); const configState = useConfigStore.getState(); diff --git a/packages/ui/src/stores/useConfigStore.ts b/packages/ui/src/stores/useConfigStore.ts index d8903c35..a94ba23e 100644 --- a/packages/ui/src/stores/useConfigStore.ts +++ b/packages/ui/src/stores/useConfigStore.ts @@ -468,6 +468,8 @@ interface ConfigStore { agentModelSelections: { [agentName: string]: { providerId: string; modelId: string } }; defaultProviders: { [key: string]: string }; isConnected: boolean; + hasEverConnected: boolean; + connectionPhase: "connecting" | "connected" | "reconnecting"; lastDisconnectReason: string | null; isInitialized: boolean; modelsMetadata: Map; @@ -589,6 +591,8 @@ export const useConfigStore = create()( agentModelSelections: {}, defaultProviders: {}, isConnected: false, + hasEverConnected: false, + connectionPhase: "connecting", lastDisconnectReason: null, isInitialized: false, modelsMetadata: new Map(), @@ -1890,7 +1894,14 @@ export const useConfigStore = create()( while (attempt < maxAttempts) { try { const isHealthy = await opencodeClient.checkHealth(); - set({ isConnected: isHealthy }); + const hasEverConnected = get().hasEverConnected; + set(isHealthy + ? { isConnected: true, hasEverConnected: true, connectionPhase: "connected" } + : { + isConnected: false, + connectionPhase: hasEverConnected ? "reconnecting" : "connecting", + lastDisconnectReason: 'health_check_unhealthy', + }); return isHealthy; } catch (error) { lastError = error; @@ -1903,7 +1914,11 @@ export const useConfigStore = create()( if (lastError) { console.warn("[ConfigStore] Failed to reach OpenCode after retrying:", lastError); } - set({ isConnected: false }); + set({ + isConnected: false, + connectionPhase: get().hasEverConnected ? "reconnecting" : "connecting", + lastDisconnectReason: 'health_check_failed', + }); return false; }, @@ -1917,7 +1932,11 @@ export const useConfigStore = create()( if (!isConnected) { if (debug) console.log("Server not connected"); - set({ isConnected: false }); + // checkConnection already set lastDisconnectReason; do not overwrite. + set({ + isConnected: false, + connectionPhase: get().hasEverConnected ? "reconnecting" : "connecting", + }); return; } @@ -1930,11 +1949,16 @@ export const useConfigStore = create()( if (debug) console.log("Loading agents..."); await get().loadAgents(); - set({ isInitialized: true, isConnected: true }); + set({ isInitialized: true, isConnected: true, hasEverConnected: true, connectionPhase: "connected" }); if (debug) console.log("App initialized successfully"); } catch (error) { console.error("Failed to initialize app:", error); - set({ isInitialized: false, isConnected: false }); + set({ + isInitialized: false, + isConnected: false, + connectionPhase: get().hasEverConnected ? "reconnecting" : "connecting", + lastDisconnectReason: 'init_error', + }); } }, diff --git a/packages/ui/src/sync/__tests__/event-pipeline.test.js b/packages/ui/src/sync/__tests__/event-pipeline.test.js index 5914212e..f508e785 100644 --- a/packages/ui/src/sync/__tests__/event-pipeline.test.js +++ b/packages/ui/src/sync/__tests__/event-pipeline.test.js @@ -666,6 +666,62 @@ describe('createEventPipeline', () => { }, ]); }); + + it('marks the pipeline disconnected on heartbeat timeout and recovers on the next websocket connect', async () => { + installDomStubs(); + globalThis.WebSocket = FakeWebSocket; + + const disconnectReasons = []; + let reconnectCount = 0; + + const sdk = { + global: { + event: async () => { + throw new Error('SSE should not be used in ws mode'); + }, + }, + }; + + const recovered = new Promise((resolve) => { + const { cleanup } = createEventPipeline({ + sdk, + transport: 'ws', + heartbeatTimeoutMs: 20, + reconnectDelayMs: 0, + wsReadyTimeoutMs: 20, + onEvent: () => {}, + onDisconnect: (reason) => { + disconnectReasons.push(reason); + }, + onReconnect: () => { + reconnectCount += 1; + if (reconnectCount === 2) { + cleanup(); + resolve(); + } + }, + }); + }); + + await Promise.resolve(); + + const firstSocket = FakeWebSocket.instances[0]; + firstSocket.emitOpen(); + firstSocket.emitMessage({ type: 'ready', scope: 'global' }); + + await new Promise((resolve) => setTimeout(resolve, 35)); + + const secondSocket = FakeWebSocket.instances[1]; + expect(secondSocket).toBeDefined(); + + secondSocket.emitOpen(); + secondSocket.emitMessage({ type: 'ready', scope: 'global' }); + + await recovered; + + expect(disconnectReasons).toEqual(['ws_heartbeat_timeout']); + expect(reconnectCount).toBe(2); + }); }); // --------------------------------------------------------------------------- diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index 8b056003..8d4d39a5 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -21,10 +21,10 @@ export type FlushHandler = (events: QueuedEvent[]) => void const FLUSH_FRAME_MS = 33 const STREAM_YIELD_MS = 8 -const RECONNECT_DELAY_MS = 250 -const HEARTBEAT_TIMEOUT_MS = 15_000 +const DEFAULT_RECONNECT_DELAY_MS = 250 +const DEFAULT_HEARTBEAT_TIMEOUT_MS = 30_000 const WS_FALLBACK_WINDOW_MS = 60_000 -const WS_READY_TIMEOUT_MS = 2_000 +const DEFAULT_WS_READY_TIMEOUT_MS = 2_000 const ABSOLUTE_URL_PATTERN = /^[a-zA-Z][a-zA-Z\d+\-.]*:\/\// export type EventPipelineInput = { @@ -38,6 +38,9 @@ export type EventPipelineInput = { /** Called when transport switches (e.g. WS timeout → SSE fallback) without actual disconnection. */ onTransportSwitch?: () => void transport?: "auto" | "ws" | "sse" + heartbeatTimeoutMs?: number + reconnectDelayMs?: number + wsReadyTimeoutMs?: number } type MessageStreamWsFrame = { @@ -145,8 +148,25 @@ type DirectoryQueue = { last: number } +type AttemptAbortReason = + | "pipeline_stopped" + | "ws_heartbeat_timeout" + | "sse_heartbeat_timeout" + | null + export function createEventPipeline(input: EventPipelineInput) { - const { sdk, onEvent, onReconnect, onDisconnect, onTransportSwitch, routeDirectory, transport = "auto" } = input + const { + sdk, + onEvent, + onReconnect, + onDisconnect, + onTransportSwitch, + routeDirectory, + transport = "auto", + heartbeatTimeoutMs = DEFAULT_HEARTBEAT_TIMEOUT_MS, + reconnectDelayMs = DEFAULT_RECONNECT_DELAY_MS, + wsReadyTimeoutMs = DEFAULT_WS_READY_TIMEOUT_MS, + } = input const abort = new AbortController() let disconnected = false let lastEventId: string | undefined @@ -244,6 +264,16 @@ export function createEventPipeline(input: EventPipelineInput) { let attempt: AbortController | undefined let lastEventAt = Date.now() let heartbeat: ReturnType | undefined + let activeTransport: "ws" | "sse" = transport === "ws" ? "ws" : "sse" + let attemptAbortReason: AttemptAbortReason = null + + const notifyDisconnected = (reason: string) => { + if (disconnected) { + return + } + disconnected = true + onDisconnect?.(reason) + } const markConnected = () => { disconnected = false @@ -295,8 +325,9 @@ export function createEventPipeline(input: EventPipelineInput) { lastEventAt = Date.now() if (heartbeat) clearTimeout(heartbeat) heartbeat = setTimeout(() => { + attemptAbortReason = `${activeTransport}_heartbeat_timeout` attempt?.abort() - }, HEARTBEAT_TIMEOUT_MS) + }, heartbeatTimeoutMs) } const clearHeartbeat = () => { @@ -359,7 +390,7 @@ export function createEventPipeline(input: EventPipelineInput) { } catch { // ignore } - }, WS_READY_TIMEOUT_MS) + }, wsReadyTimeoutMs) const cleanup = () => { if (readyTimer) { @@ -499,9 +530,12 @@ export function createEventPipeline(input: EventPipelineInput) { while (!abort.signal.aborted) { attempt = new AbortController() lastEventAt = Date.now() - let retryDelayMs = RECONNECT_DELAY_MS + attemptAbortReason = null + let retryDelayMs = reconnectDelayMs const currentTransport = resolveTransport() + activeTransport = currentTransport const onAbort = () => { + attemptAbortReason = "pipeline_stopped" attempt?.abort() } abort.signal.addEventListener("abort", onAbort) @@ -531,21 +565,18 @@ export function createEventPipeline(input: EventPipelineInput) { // update connection state (e.g. set isConnected = false). // Guard: only fire once per disconnection cycle to avoid repeated // setState calls on every failed retry attempt. - if (!disconnected) { - disconnected = true - const taggedReason = typeof error === "object" && error !== null - ? (error as { reason?: unknown }).reason - : undefined - const message = typeof error === "object" && error !== null - ? (error as { message?: unknown }).message - : undefined - const reason = typeof taggedReason === "string" && taggedReason.length > 0 - ? taggedReason - : typeof message === "string" && message.length > 0 - ? `${currentTransport}_error:${message.slice(0, 80)}` - : `${currentTransport}_error:unknown` - onDisconnect?.(reason) - } + const taggedReason = typeof error === "object" && error !== null + ? (error as { reason?: unknown }).reason + : undefined + const message = typeof error === "object" && error !== null + ? (error as { message?: unknown }).message + : undefined + const reason = typeof taggedReason === "string" && taggedReason.length > 0 + ? taggedReason + : typeof message === "string" && message.length > 0 + ? `${currentTransport}_error:${message.slice(0, 80)}` + : `${currentTransport}_error:unknown` + notifyDisconnected(reason) } } finally { abort.signal.removeEventListener("abort", onAbort) @@ -554,6 +585,11 @@ export function createEventPipeline(input: EventPipelineInput) { } if (abort.signal.aborted) return + if (attemptAbortReason && attemptAbortReason !== "pipeline_stopped") { + notifyDisconnected(attemptAbortReason) + retryDelayMs = 0 + attemptAbortReason = null + } if (retryDelayMs > 0) { await wait(retryDelayMs) } @@ -563,7 +599,7 @@ export function createEventPipeline(input: EventPipelineInput) { const onVisibility = () => { if (typeof document === "undefined") return if (document.visibilityState !== "visible") return - if (Date.now() - lastEventAt < HEARTBEAT_TIMEOUT_MS) return + if (Date.now() - lastEventAt < heartbeatTimeoutMs) return attempt?.abort() } diff --git a/packages/ui/src/sync/session-actions.ts b/packages/ui/src/sync/session-actions.ts index 25b576c2..4217b81a 100644 --- a/packages/ui/src/sync/session-actions.ts +++ b/packages/ui/src/sync/session-actions.ts @@ -55,11 +55,30 @@ function dir() { } function connectionLostError(): Error { - const reason = useConfigStore.getState().lastDisconnectReason - const suffix = reason ? ` (${reason})` : " (never connected)" + const { hasEverConnected, lastDisconnectReason } = useConfigStore.getState() + const suffix = lastDisconnectReason + ? ` (${lastDisconnectReason})` + : hasEverConnected + ? "" + : " (never connected)" return new Error(`Connection lost${suffix}. Please wait for reconnection.`) } +// Wait briefly for the pipeline to re-establish connection before failing a +// send. Transient reconnects (heartbeat race, WS→SSE fallback, brief network +// blip) otherwise surface as a hard "Connection lost" toast even though the +// pipeline recovers within a second. Poll isConnected at 100ms intervals. +const CONNECTION_GRACE_MS = 2000 +export async function waitForConnectionOrThrow(): Promise { + if (useConfigStore.getState().isConnected) return + const deadline = Date.now() + CONNECTION_GRACE_MS + while (Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 100)) + if (useConfigStore.getState().isConnected) return + } + throw connectionLostError() +} + function getSessionDirectory(sessionId: string): string | undefined { return useSessionUIStore.getState().getDirectoryForSession(sessionId) || dir() } @@ -332,9 +351,7 @@ export async function optimisticSend(input: { throw new Error("Optimistic refs not set — is useSync() mounted?") } - if (!useConfigStore.getState().isConnected) { - throw connectionLostError() - } + await waitForConnectionOrThrow() const store = dirStore() const messageID = ascendingId("msg") @@ -419,9 +436,7 @@ export async function respondToPermission( requestId: string, response: "once" | "always" | "reject", ): Promise { - if (!useConfigStore.getState().isConnected) { - throw connectionLostError() - } + await waitForConnectionOrThrow() const result = await getRequestReplyClient("permission", sessionId, requestId).permission.reply({ requestID: requestId, reply: response, @@ -435,9 +450,7 @@ export async function dismissPermission( sessionId: string, requestId: string, ): Promise { - if (!useConfigStore.getState().isConnected) { - throw connectionLostError() - } + await waitForConnectionOrThrow() const result = await getRequestReplyClient("permission", sessionId, requestId).permission.reply({ requestID: requestId, reply: "reject", @@ -456,9 +469,7 @@ export async function respondToQuestion( requestId: string, answers: string[] | string[][], ): Promise { - if (!useConfigStore.getState().isConnected) { - throw connectionLostError() - } + await waitForConnectionOrThrow() const result = await getRequestReplyClient("question", sessionId, requestId).question.reply({ requestID: requestId, answers: answers as Array>, @@ -472,9 +483,7 @@ export async function rejectQuestion( sessionId: string, requestId: string, ): Promise { - if (!useConfigStore.getState().isConnected) { - throw connectionLostError() - } + await waitForConnectionOrThrow() const result = await getRequestReplyClient("question", sessionId, requestId).question.reject({ requestID: requestId, }) diff --git a/packages/ui/src/sync/sync-context.tsx b/packages/ui/src/sync/sync-context.tsx index 49254f46..50531cf3 100644 --- a/packages/ui/src/sync/sync-context.tsx +++ b/packages/ui/src/sync/sync-context.tsx @@ -1376,7 +1376,11 @@ export function SyncProvider(props: { handleEvent(directory, payload, childStores, routingIndex) }, onReconnect: () => { - useConfigStore.setState({ isConnected: true, lastDisconnectReason: null }) + useConfigStore.setState({ + isConnected: true, + hasEverConnected: true, + connectionPhase: "connected", + }) for (const [dir, store] of childStores.children) { if (reconnectResyncing.has(dir)) continue if (getReconnectCandidateSessionIds(store.getState()).length === 0) continue @@ -1392,13 +1396,22 @@ export function SyncProvider(props: { } }, onDisconnect: (reason) => { - useConfigStore.setState({ isConnected: false, lastDisconnectReason: reason }) + const { hasEverConnected } = useConfigStore.getState() + useConfigStore.setState({ + isConnected: false, + connectionPhase: hasEverConnected ? "reconnecting" : "connecting", + lastDisconnectReason: reason, + }) }, onTransportSwitch: () => { // Transport switched (e.g. WS timeout → SSE fallback) without // actual disconnection. No events lost — just update connection // state without triggering a full directory resync. - useConfigStore.setState({ isConnected: true }) + useConfigStore.setState({ + isConnected: true, + hasEverConnected: true, + connectionPhase: "connected", + }) }, }) return cleanup diff --git a/packages/web/server/lib/event-stream/DOCUMENTATION.md b/packages/web/server/lib/event-stream/DOCUMENTATION.md index f237c83b..e95a2099 100644 --- a/packages/web/server/lib/event-stream/DOCUMENTATION.md +++ b/packages/web/server/lib/event-stream/DOCUMENTATION.md @@ -28,7 +28,9 @@ This module contains the OpenChamber message-stream WebSocket protocol and runti - Browser clients connect to the WS endpoints above. - OpenChamber still fetches OpenCode upstream event streams over SSE. - Each WS connection proxies one upstream SSE stream. -- Global synthetic events such as `openchamber:session-status`, `openchamber:session-activity`, `openchamber:notification`, and `openchamber:heartbeat` are preserved on the WS path. +- If an upstream SSE stream stalls after the browser WS is already ready, the runtime aborts that upstream fetch and reconnects upstream with `Last-Event-ID`, keeping the browser WS alive when recovery is fast. +- Health checks are reserved for initial upstream connect failures and explicit upstream-unavailable responses, not for ordinary stall recovery on an already-established stream. +- 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. ## Notes for contributors diff --git a/packages/web/server/lib/event-stream/runtime.js b/packages/web/server/lib/event-stream/runtime.js index 701a2b19..cee8f6f1 100644 --- a/packages/web/server/lib/event-stream/runtime.js +++ b/packages/web/server/lib/event-stream/runtime.js @@ -10,6 +10,9 @@ import { sendMessageStreamWsFrame, } from './protocol.js'; +const MESSAGE_STREAM_UPSTREAM_STALL_TIMEOUT_MS = 20_000; +const MESSAGE_STREAM_UPSTREAM_RECONNECT_DELAY_MS = 250; + function shouldTriggerUpstreamHealthCheck(upstream) { if (!upstream) { return true; @@ -67,6 +70,9 @@ export function createMessageStreamWsRuntime({ processForwardedEventPayload, wsClients, triggerHealthCheck, + heartbeatIntervalMs = MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS, + upstreamStallTimeoutMs = MESSAGE_STREAM_UPSTREAM_STALL_TIMEOUT_MS, + upstreamReconnectDelayMs = MESSAGE_STREAM_UPSTREAM_RECONNECT_DELAY_MS, fetchImpl = fetch, }) { const wsServer = new WebSocketServer({ @@ -82,13 +88,22 @@ export function createMessageStreamWsRuntime({ const requestedDirectory = requestUrl.searchParams.get('directory')?.trim() || ''; const controller = new AbortController(); + let currentUpstreamAbortController = null; + let upstreamConnected = false; + let streamReady = false; + let lastEventId = requestedLastEventId; const cleanup = () => { if (!controller.signal.aborted) { controller.abort(); } + if (currentUpstreamAbortController && !currentUpstreamAbortController.signal.aborted) { + currentUpstreamAbortController.abort(); + } wsClients.delete(socket); }; + const wait = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + const pingInterval = setInterval(() => { if (socket.readyState !== 1) { return; @@ -98,15 +113,20 @@ export function createMessageStreamWsRuntime({ socket.ping(); } catch { } - }, MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS); + }, heartbeatIntervalMs); const heartbeatInterval = setInterval(() => { + if (!upstreamConnected) { + return; + } + sendMessageStreamWsEvent(socket, { type: 'openchamber:heartbeat', timestamp: Date.now() }, { directory: 'global' }); - }, MESSAGE_STREAM_WS_HEARTBEAT_INTERVAL_MS); + }, heartbeatIntervalMs); socket.on('close', () => { clearInterval(pingInterval); clearInterval(heartbeatInterval); + upstreamConnected = false; cleanup(); }); @@ -115,73 +135,6 @@ export function createMessageStreamWsRuntime({ }); const run = async () => { - let targetUrl; - try { - targetUrl = new URL(buildOpenCodeUrl(isGlobalStream ? '/global/event' : '/event', '')); - } catch { - sendMessageStreamWsFrame(socket, { type: 'error', message: 'OpenCode service unavailable' }); - socket.close(1011, 'OpenCode service unavailable'); - return; - } - - if (!isGlobalStream && requestedDirectory) { - targetUrl.searchParams.set('directory', requestedDirectory); - } - - const headers = { - Accept: 'text/event-stream', - 'Cache-Control': 'no-cache', - Connection: 'keep-alive', - ...getOpenCodeAuthHeaders(), - }; - - if (requestedLastEventId) { - headers['Last-Event-ID'] = requestedLastEventId; - } - - let upstream; - try { - upstream = await fetchImpl(targetUrl.toString(), { - headers, - signal: controller.signal, - }); - } catch { - if (!controller.signal.aborted) { - sendMessageStreamWsFrame(socket, { type: 'error', message: 'Failed to connect to OpenCode event stream' }); - socket.close(1011, 'Failed to connect to OpenCode event stream'); - // Trigger immediate health check so the server detects and - // restarts a dead OpenCode process without waiting for the next - // periodic interval (up to 15 s). - triggerHealthCheck?.(); - } - return; - } - - if (!upstream.ok || !upstream.body) { - sendMessageStreamWsFrame(socket, { - type: 'error', - message: `OpenCode event stream unavailable (${upstream.status})`, - }); - socket.close(1011, 'OpenCode event stream unavailable'); - if (shouldTriggerUpstreamHealthCheck(upstream)) { - triggerHealthCheck?.(); - } - return; - } - - sendMessageStreamWsFrame(socket, { - type: 'ready', - scope: isGlobalStream ? 'global' : 'directory', - }); - - if (isGlobalStream) { - wsClients.add(socket); - } - - const decoder = new TextDecoder(); - const reader = upstream.body.getReader(); - let buffer = ''; - const forwardBlock = (block) => { if (!block) { return; @@ -197,6 +150,10 @@ export function createMessageStreamWsRuntime({ ? (typeof envelope?.directory === 'string' && envelope.directory.length > 0 ? envelope.directory : 'global') : (requestedDirectory || envelope?.directory || 'global'); + if (typeof envelope?.eventId === 'string' && envelope.eventId.length > 0) { + lastEventId = envelope.eventId; + } + sendMessageStreamWsEvent(socket, payload, { directory, eventId: typeof envelope?.eventId === 'string' && envelope.eventId.length > 0 ? envelope.eventId : undefined, @@ -208,25 +165,148 @@ export function createMessageStreamWsRuntime({ }; try { - while (true) { - const { value, done } = await reader.read(); - if (done) { - break; + while (!controller.signal.aborted) { + let targetUrl; + try { + targetUrl = new URL(buildOpenCodeUrl(isGlobalStream ? '/global/event' : '/event', '')); + } catch { + sendMessageStreamWsFrame(socket, { type: 'error', message: 'OpenCode service unavailable' }); + socket.close(1011, 'OpenCode service unavailable'); + return; } - buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, '\n'); - - let separatorIndex = buffer.indexOf('\n\n'); - while (separatorIndex !== -1) { - const block = buffer.slice(0, separatorIndex); - buffer = buffer.slice(separatorIndex + 2); - forwardBlock(block); - separatorIndex = buffer.indexOf('\n\n'); + if (!isGlobalStream && requestedDirectory) { + targetUrl.searchParams.set('directory', requestedDirectory); } - } - if (buffer.trim().length > 0) { - forwardBlock(buffer.trim()); + const headers = { + Accept: 'text/event-stream', + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + ...getOpenCodeAuthHeaders(), + }; + if (lastEventId) { + headers['Last-Event-ID'] = lastEventId; + } + + const upstreamController = new AbortController(); + currentUpstreamAbortController = upstreamController; + const abortUpstream = () => upstreamController.abort(); + controller.signal.addEventListener('abort', abortUpstream, { once: true }); + + let upstream; + try { + upstream = await fetchImpl(targetUrl.toString(), { + headers, + signal: upstreamController.signal, + }); + } catch { + controller.signal.removeEventListener('abort', abortUpstream); + currentUpstreamAbortController = null; + upstreamConnected = false; + if (controller.signal.aborted) { + return; + } + if (!streamReady) { + sendMessageStreamWsFrame(socket, { type: 'error', message: 'Failed to connect to OpenCode event stream' }); + socket.close(1011, 'Failed to connect to OpenCode event stream'); + triggerHealthCheck?.(); + return; + } + await wait(upstreamReconnectDelayMs); + continue; + } + + if (!upstream.ok || !upstream.body) { + controller.signal.removeEventListener('abort', abortUpstream); + currentUpstreamAbortController = null; + upstreamConnected = false; + if (!streamReady) { + sendMessageStreamWsFrame(socket, { + type: 'error', + message: `OpenCode event stream unavailable (${upstream.status})`, + }); + socket.close(1011, 'OpenCode event stream unavailable'); + if (shouldTriggerUpstreamHealthCheck(upstream)) { + triggerHealthCheck?.(); + } + return; + } + await wait(upstreamReconnectDelayMs); + continue; + } + + if (!streamReady) { + sendMessageStreamWsFrame(socket, { + type: 'ready', + scope: isGlobalStream ? 'global' : 'directory', + }); + streamReady = true; + + if (isGlobalStream) { + wsClients.add(socket); + } + } + + upstreamConnected = true; + const decoder = new TextDecoder(); + const reader = upstream.body.getReader(); + let buffer = ''; + let upstreamAbortReason = null; + let stallTimer = null; + const resetStallTimer = () => { + if (stallTimer) { + clearTimeout(stallTimer); + } + stallTimer = setTimeout(() => { + upstreamAbortReason = 'upstream_stalled'; + upstreamConnected = false; + upstreamController.abort(); + }, upstreamStallTimeoutMs); + }; + + resetStallTimer(); + + try { + while (true) { + const { value, done } = await reader.read(); + if (done) { + break; + } + + resetStallTimer(); + buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, '\n'); + + let separatorIndex = buffer.indexOf('\n\n'); + while (separatorIndex !== -1) { + const block = buffer.slice(0, separatorIndex); + buffer = buffer.slice(separatorIndex + 2); + forwardBlock(block); + separatorIndex = buffer.indexOf('\n\n'); + } + } + + if (buffer.trim().length > 0) { + forwardBlock(buffer.trim()); + } + } catch (error) { + if (!controller.signal.aborted && upstreamAbortReason !== 'upstream_stalled') { + console.warn('Message stream WS proxy error:', error); + } + } finally { + if (stallTimer) { + clearTimeout(stallTimer); + } + upstreamConnected = false; + currentUpstreamAbortController = null; + controller.signal.removeEventListener('abort', abortUpstream); + } + + if (controller.signal.aborted) { + return; + } + + await wait(upstreamReconnectDelayMs); } } catch (error) { if (!controller.signal.aborted) { diff --git a/packages/web/server/lib/event-stream/runtime.test.js b/packages/web/server/lib/event-stream/runtime.test.js index 30d6d21b..5d3c366d 100644 --- a/packages/web/server/lib/event-stream/runtime.test.js +++ b/packages/web/server/lib/event-stream/runtime.test.js @@ -1,6 +1,68 @@ +import { EventEmitter } from 'node:events'; import { describe, expect, it } from 'vitest'; -import { createGlobalUiEventBroadcaster } from './runtime.js'; +import { createGlobalUiEventBroadcaster, createMessageStreamWsRuntime } from './runtime.js'; + +class FakeSocket extends EventEmitter { + constructor() { + super(); + this.readyState = 1; + this.sent = []; + this.closeCalls = []; + } + + send(payload) { + this.sent.push(JSON.parse(payload)); + } + + ping() { + void 0; + } + + close(code, reason) { + if (this.readyState === 3) { + return; + } + this.readyState = 3; + this.closeCalls.push({ code, reason }); + this.emit('close'); + } +} + +function createSseResponse({ blocks = [], signal, holdOpen = false }) { + const encoder = new TextEncoder(); + let index = 0; + + return { + ok: true, + body: { + getReader() { + return { + async read() { + if (index < blocks.length) { + const next = blocks[index++]; + return { value: encoder.encode(next), done: false }; + } + + if (!holdOpen) { + return { value: undefined, done: true }; + } + + return new Promise((resolve, reject) => { + const onAbort = () => { + signal.removeEventListener('abort', onAbort); + const error = new Error('Aborted'); + error.name = 'AbortError'; + reject(error); + }; + signal.addEventListener('abort', onAbort, { once: true }); + }); + }, + }; + }, + }, + }; +} describe('event stream broadcaster', () => { it('fans out synthetic events to SSE and WS clients', () => { @@ -63,3 +125,71 @@ describe('event stream broadcaster', () => { expect(wsClients.size).toBe(0); }); }); + +describe('message stream websocket runtime', () => { + it('reconnects a stalled upstream SSE stream and resumes from the last event id', async () => { + const server = new EventEmitter(); + const wsClients = new Set(); + let triggerHealthCheckCalls = 0; + const fetchCalls = []; + let upstreamAttempt = 0; + + const runtime = createMessageStreamWsRuntime({ + server, + uiAuthController: null, + isRequestOriginAllowed: async () => true, + rejectWebSocketUpgrade() { + throw new Error('upgrade should not be used in this test'); + }, + buildOpenCodeUrl: (path) => `http://127.0.0.1:4096${path}`, + getOpenCodeAuthHeaders: () => ({}), + processForwardedEventPayload() {}, + wsClients, + triggerHealthCheck: () => { + triggerHealthCheckCalls += 1; + }, + heartbeatIntervalMs: 50, + upstreamStallTimeoutMs: 20, + upstreamReconnectDelayMs: 0, + fetchImpl: async (_url, options) => { + const lastEventId = options?.headers?.['Last-Event-ID'] ?? null; + fetchCalls.push(lastEventId); + upstreamAttempt += 1; + + if (upstreamAttempt === 1) { + return createSseResponse({ + signal: options.signal, + holdOpen: true, + blocks: [ + 'id: evt-1\ndata: {"type":"server.connected","properties":{}}\n\n', + ], + }); + } + + return createSseResponse({ + signal: options.signal, + holdOpen: false, + blocks: [ + 'id: evt-2\ndata: {"type":"server.connected","properties":{}}\n\n', + ], + }); + }, + }); + + const socket = new FakeSocket(); + runtime.wsServer.emit('connection', socket, { url: '/api/global/event/ws' }); + + await new Promise((resolve) => setTimeout(resolve, 35)); + + const readyFrames = socket.sent.filter((frame) => frame.type === 'ready'); + const eventFrames = socket.sent.filter((frame) => frame.type === 'event' && frame.payload?.type === 'server.connected'); + + expect(readyFrames).toHaveLength(1); + expect(eventFrames.length).toBeGreaterThanOrEqual(2); + expect(fetchCalls.slice(0, 2)).toEqual([null, 'evt-1']); + expect(triggerHealthCheckCalls).toBe(0); + + socket.close(); + await runtime.close(); + }); +});