diff --git a/apps/dashboard/lib/activity/client-stream.ts b/apps/dashboard/lib/activity/client-stream.ts index 442897b..a2fc27e 100644 --- a/apps/dashboard/lib/activity/client-stream.ts +++ b/apps/dashboard/lib/activity/client-stream.ts @@ -13,11 +13,17 @@ export interface ActivityStreamConnectionOptions { onEvent: (event: ActivityEvent) => void; onFallback: () => void; onReady?: () => void; + reconnectBaseDelayMs?: number; + reconnectJitterRatio?: number; + reconnectMaxDelayMs?: number; url?: string; } const DEFAULT_CONNECTION_TIMEOUT_MS = 10_000; const DEFAULT_HEARTBEAT_TIMEOUT_MS = 45_000; +const DEFAULT_RECONNECT_BASE_DELAY_MS = 1_000; +const DEFAULT_RECONNECT_JITTER_RATIO = 0.25; +const DEFAULT_RECONNECT_MAX_DELAY_MS = 30_000; export function connectActivityStream({ connectionTimeoutMs = DEFAULT_CONNECTION_TIMEOUT_MS, @@ -26,12 +32,18 @@ export function connectActivityStream({ onEvent, onFallback, onReady = () => {}, + reconnectBaseDelayMs = DEFAULT_RECONNECT_BASE_DELAY_MS, + reconnectJitterRatio = DEFAULT_RECONNECT_JITTER_RATIO, + reconnectMaxDelayMs = DEFAULT_RECONNECT_MAX_DELAY_MS, url = "/api/activity/stream", }: ActivityStreamConnectionOptions): () => void { let source: ActivityEventSourceLike | null = null; let stopped = false; let fallbackStarted = false; - let ready = false; + let connectionReady = false; + let lastEventId: string | null = null; + let reconnectAttempt = 0; + let reconnectTimer: ReturnType | null = null; let watchdog: ReturnType | null = null; const clearWatchdog = () => { @@ -40,20 +52,28 @@ export function connectActivityStream({ watchdog = null; }; - const armWatchdog = (timeoutMs: number) => { + const clearReconnectTimer = () => { + if (reconnectTimer === null) return; + clearTimeout(reconnectTimer); + reconnectTimer = null; + }; + + const armWatchdog = (timeoutMs: number, onTimeout: () => void) => { clearWatchdog(); - watchdog = setTimeout(startFallback, timeoutMs); + watchdog = setTimeout(onTimeout, timeoutMs); }; const markAlive = () => { if (stopped || fallbackStarted) return; - armWatchdog(heartbeatTimeoutMs); + armWatchdog(heartbeatTimeoutMs, scheduleReconnect); }; const onActivity = ((rawEvent: Event) => { markAlive(); const event = parseActivityEvent(rawEvent); - if (event) onEvent(event); + if (!event) return; + lastEventId = event.id; + onEvent(event); }) as EventListener; const onHeartbeat = (() => { @@ -61,8 +81,9 @@ export function connectActivityStream({ }) as EventListener; const onReadyEvent = (() => { - const firstReady = !ready; - ready = true; + const firstReady = !connectionReady; + connectionReady = true; + reconnectAttempt = 0; markAlive(); if (firstReady) onReady(); }) as EventListener; @@ -78,30 +99,69 @@ export function connectActivityStream({ source = null; }; + const streamUrl = () => { + if (!lastEventId) return url; + const separator = url.includes("?") ? "&" : "?"; + return `${url}${separator}lastEventId=${encodeURIComponent(lastEventId)}`; + }; + const startFallback = () => { if (stopped || fallbackStarted) return; fallbackStarted = true; + clearReconnectTimer(); detach(); onFallback(); }; + const reconnectDelay = () => { + const exponentialDelay = Math.min( + reconnectMaxDelayMs, + reconnectBaseDelayMs * 2 ** reconnectAttempt + ); + const jitter = exponentialDelay * reconnectJitterRatio * Math.random(); + reconnectAttempt += 1; + return Math.round(exponentialDelay + jitter); + }; + + const scheduleReconnect = () => { + if (stopped || fallbackStarted) return; + if (!connectionReady) { + startFallback(); + return; + } + + detach(); + clearReconnectTimer(); + reconnectTimer = setTimeout(connect, reconnectDelay()); + }; + const onError = (() => { - startFallback(); + scheduleReconnect(); }) as EventListener; - try { - source = createEventSource(url); + const connect = () => { + if (stopped || fallbackStarted) return; + connectionReady = false; + + try { + source = createEventSource(streamUrl()); + } catch { + startFallback(); + return; + } + source.addEventListener("activity", onActivity); source.addEventListener("error", onError); source.addEventListener("heartbeat", onHeartbeat); source.addEventListener("ready", onReadyEvent); - armWatchdog(connectionTimeoutMs); - } catch { - startFallback(); - } + armWatchdog(connectionTimeoutMs, startFallback); + }; + + connect(); return () => { stopped = true; + clearReconnectTimer(); detach(); }; } diff --git a/apps/dashboard/test/activity-stream.test.ts b/apps/dashboard/test/activity-stream.test.ts index 2d11da5..013ad68 100644 --- a/apps/dashboard/test/activity-stream.test.ts +++ b/apps/dashboard/test/activity-stream.test.ts @@ -75,32 +75,50 @@ describe("activity SSE delivery", () => { assert.equal(getActivitySubscriberCount(), initialSubscribers); }); - test("client connector accepts activity and falls back exactly once on stream error", () => { - const source = new FakeEventSource(); + test("client connector accepts activity and reconnects with backoff on stream error", async () => { + const sources: FakeEventSource[] = []; const received: string[] = []; let fallbackCount = 0; + let readyCount = 0; const event = makeActivityEvent({ id: "evt_client_stream_001" }); const disconnect = connectActivityStream({ - createEventSource: () => source, + createEventSource: (url) => { + const source = new FakeEventSource(); + sources.push(source); + if (sources.length === 2) { + assert.match(url, /lastEventId=evt_client_stream_001/); + } + return source; + }, onEvent: (activity) => received.push(activity.id), onFallback: () => { fallbackCount += 1; }, + onReady: () => { + readyCount += 1; + }, + reconnectBaseDelayMs: 10, + reconnectJitterRatio: 0, + reconnectMaxDelayMs: 10, }); - source.emit("ready", "{}"); - source.emit("activity", JSON.stringify(event)); - source.emit("activity", "not-json"); - source.emit("error"); - source.emit("error"); + sources[0].emit("ready", "{}"); + sources[0].emit("activity", JSON.stringify(event)); + sources[0].emit("activity", "not-json"); + sources[0].emit("error"); + sources[0].emit("error"); + + await delay(30); + sources[1].emit("ready", "{}"); assert.deepEqual(received, [event.id]); - assert.equal(fallbackCount, 1); - assert.equal(source.closeCount, 1); + assert.equal(fallbackCount, 0); + assert.equal(readyCount, 2); + assert.equal(sources[0].closeCount, 1); disconnect(); - assert.equal(source.closeCount, 1); + assert.equal(sources[1].closeCount, 1); }); test("ready handshake reconciles the REST snapshot after subscription", () => { @@ -146,24 +164,32 @@ describe("activity SSE delivery", () => { assert.equal(source.closeCount, 1); }); - test("client falls back when a ready stream stops sending heartbeats", async () => { - const source = new FakeEventSource(); + test("client reconnects when a ready stream stops sending heartbeats", async () => { + const sources: FakeEventSource[] = []; let fallbackCount = 0; connectActivityStream({ connectionTimeoutMs: 100, - createEventSource: () => source, + createEventSource: () => { + const source = new FakeEventSource(); + sources.push(source); + return source; + }, heartbeatTimeoutMs: 10, onEvent: () => {}, onFallback: () => { fallbackCount += 1; }, + reconnectBaseDelayMs: 10, + reconnectJitterRatio: 0, + reconnectMaxDelayMs: 10, }); - source.emit("ready", "{}"); + sources[0].emit("ready", "{}"); await delay(30); - assert.equal(fallbackCount, 1); - assert.equal(source.closeCount, 1); + assert.equal(fallbackCount, 0); + assert.equal(sources[0].closeCount, 1); + assert.equal(sources.length, 2); }); test("server disconnects a stream whose bounded output queue fills", async () => {