From 7721c2fe10e49e466feb5bc5489d1d85ee2650ec Mon Sep 17 00:00:00 2001 From: Isaac Sanchez-Hawkins <266845420+isanchez404@users.noreply.github.com> Date: Tue, 12 May 2026 04:02:04 -0400 Subject: [PATCH] fix(event-stream): isolate subscriber failures (#1235) * fix(event-stream): isolate subscriber failures * fix(event-stream): handle async subscriber failures --------- Co-authored-by: Isaac Sanchez --- .../web/server/lib/event-stream/global-hub.js | 17 ++- .../lib/event-stream/global-hub.test.js | 140 ++++++++++++++++++ 2 files changed, 155 insertions(+), 2 deletions(-) create mode 100644 packages/web/server/lib/event-stream/global-hub.test.js diff --git a/packages/web/server/lib/event-stream/global-hub.js b/packages/web/server/lib/event-stream/global-hub.js index 86f28b3f..165dc0b7 100644 --- a/packages/web/server/lib/event-stream/global-hub.js +++ b/packages/web/server/lib/event-stream/global-hub.js @@ -22,9 +22,22 @@ export function createGlobalMessageStreamHub({ let everConnected = false; let buildUrlFailed = false; + const notifySubscriber = (kind, subscriber, payload) => { + try { + const result = subscriber(payload); + if (result && typeof result.catch === 'function') { + result.catch((error) => { + console.warn(`Global message stream ${kind} subscriber failed:`, error); + }); + } + } catch (error) { + console.warn(`Global message stream ${kind} subscriber failed:`, error); + } + }; + const notifyStatus = (status) => { for (const subscriber of Array.from(statusSubscribers)) { - subscriber(status); + notifySubscriber('status', subscriber, status); } }; @@ -81,7 +94,7 @@ export function createGlobalMessageStreamHub({ } for (const subscriber of Array.from(eventSubscribers)) { - subscriber(normalized); + notifySubscriber('event', subscriber, normalized); } }, onError(error) { diff --git a/packages/web/server/lib/event-stream/global-hub.test.js b/packages/web/server/lib/event-stream/global-hub.test.js new file mode 100644 index 00000000..d72b136f --- /dev/null +++ b/packages/web/server/lib/event-stream/global-hub.test.js @@ -0,0 +1,140 @@ +import { describe, expect, it, vi } from 'vitest'; + +import { createGlobalMessageStreamHub } from './global-hub.js'; + +function createSseResponse({ blocks = [] } = {}) { + const encoder = new TextEncoder(); + let index = 0; + + return { + ok: true, + body: { + getReader() { + return { + async read() { + if (index < blocks.length) { + return { value: encoder.encode(blocks[index++]), done: false }; + } + return { value: undefined, done: true }; + }, + }; + }, + }, + }; +} + +async function waitForAssertion(assertion) { + const deadline = Date.now() + 1000; + let lastError; + + while (Date.now() < deadline) { + try { + assertion(); + return; + } catch (error) { + lastError = error; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + } + + throw lastError; +} + +describe('createGlobalMessageStreamHub', () => { + it('continues fanout when an event subscriber throws', async () => { + const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const received = []; + const hub = createGlobalMessageStreamHub({ + buildOpenCodeUrl: (pathname) => `http://127.0.0.1:4096${pathname}`, + getOpenCodeAuthHeaders: () => ({}), + upstreamReconnectDelayMs: 100, + fetchImpl: async () => createSseResponse({ + blocks: [ + 'id: evt-1\ndata: {"type":"session.updated","properties":{}}\n\n', + ], + }), + }); + + hub.subscribeEvent(() => { + throw new Error('subscriber failed'); + }); + hub.subscribeEvent((event) => { + received.push(event.eventId); + }); + + try { + hub.start(); + await waitForAssertion(() => { + expect(received).toEqual(['evt-1']); + }); + expect(warnSpy).toHaveBeenCalled(); + } finally { + hub.stop(); + warnSpy.mockRestore(); + } + }); + + it('continues status fanout when a status subscriber throws', async () => { + const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const received = []; + const hub = createGlobalMessageStreamHub({ + buildOpenCodeUrl: (pathname) => `http://127.0.0.1:4096${pathname}`, + getOpenCodeAuthHeaders: () => ({}), + upstreamReconnectDelayMs: 100, + fetchImpl: async () => createSseResponse(), + }); + + hub.subscribeStatus(() => { + throw new Error('status subscriber failed'); + }); + hub.subscribeStatus((status) => { + received.push(status.type); + }); + + try { + hub.start(); + await waitForAssertion(() => { + expect(received).toContain('connect'); + }); + expect(warnSpy).toHaveBeenCalled(); + } finally { + hub.stop(); + warnSpy.mockRestore(); + } + }); + + it('continues fanout when an async event subscriber rejects', async () => { + const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const received = []; + const hub = createGlobalMessageStreamHub({ + buildOpenCodeUrl: (pathname) => `http://127.0.0.1:4096${pathname}`, + getOpenCodeAuthHeaders: () => ({}), + upstreamReconnectDelayMs: 100, + fetchImpl: async () => createSseResponse({ + blocks: [ + 'id: evt-1\ndata: {"type":"session.updated","properties":{}}\n\n', + ], + }), + }); + + hub.subscribeEvent(async () => { + throw new Error('async subscriber failed'); + }); + hub.subscribeEvent((event) => { + received.push(event.eventId); + }); + + try { + hub.start(); + await waitForAssertion(() => { + expect(received).toEqual(['evt-1']); + }); + await waitForAssertion(() => { + expect(warnSpy).toHaveBeenCalled(); + }); + } finally { + hub.stop(); + warnSpy.mockRestore(); + } + }); +});