From cac919eb57c87e2d9ea106dd6355764e51cc64a0 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:21:00 -0700 Subject: [PATCH 1/4] Add tests for SSE stream ordering, teardown, and keepalive Covers the six workbench-events defects: unhandled rejections from a rejecting authorize(), a dead stream after a failed write, the bare catch swallowing the platform-subscribe failure, out-of-order deliveries under a slow authorize(), the presence-snapshot/delta race, and the periodic keepalive. --- packages/chat/test/workbench-events.test.ts | 273 +++++++++++++++++--- 1 file changed, 232 insertions(+), 41 deletions(-) diff --git a/packages/chat/test/workbench-events.test.ts b/packages/chat/test/workbench-events.test.ts index 9f2f02446..6a9a5e951 100644 --- a/packages/chat/test/workbench-events.test.ts +++ b/packages/chat/test/workbench-events.test.ts @@ -4,7 +4,8 @@ // mid-stream) is covered end to end over real HTTP in // `workbench-share-routes.test.ts`; these tests only pin the bridge's // own unit-level contract with a stub `authorize`. -import { describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, spyOn, test } from "bun:test"; +import * as errorSink from "@corbits/error-sink"; import { bridgeWorkbenchStream, createWorkbenchSubscriberRegistry, @@ -15,6 +16,17 @@ import type { ChatWorkbenchEvent, WorkbenchEvents } from "../src/platform-port"; const alwaysAuthorized = () => Promise.resolve(true); +// Every `deliverEvent`/`writeSSE` in the bridge runs through a chained +// promise queue now (the fix for out-of-order writes), which adds a +// handful of microtask hops between a publish and its write landing on +// the stream. Draining generously here is cheaper than pinning an exact +// hop count that breaks the moment the queue's internals change shape. +async function flush(times = 20): Promise { + for (let i = 0; i < times; i++) { + await Promise.resolve(); + } +} + function fakeStream( writeSSE: (message: unknown) => Promise, close: () => Promise = () => Promise.resolve(), @@ -146,6 +158,14 @@ describe("createPlatformWorkbenchFanout", () => { }); describe("bridgeWorkbenchStream", () => { + afterEach(() => { + // Some tests below stub `@corbits/error-sink`'s `reportError`; always + // restore it so a mock from one test can't leak into the next. + ( + errorSink.reportError as unknown as { mockRestore?: () => void } + ).mockRestore?.(); + }); + test("forwards a registry publish onto the stream as an SSE write", async () => { const registry = createWorkbenchSubscriberRegistry(); const writes: unknown[] = []; @@ -162,19 +182,25 @@ describe("bridgeWorkbenchStream", () => { authorize: alwaysAuthorized, }); registry.publish("chan_1", { type: "chat.typing", data: { a: 1 } }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(writes).toHaveLength(1); }); - test("a subscriber whose write rejects is removed rather than left dangling", async () => { + test("a subscriber whose write rejects is removed and the stream is closed", async () => { const registry = createWorkbenchSubscriberRegistry(); let writeCount = 0; - const stream = fakeStream(() => { - writeCount += 1; - return Promise.reject(new Error("client disconnected")); - }); + let closeCount = 0; + const stream = fakeStream( + () => { + writeCount += 1; + return Promise.reject(new Error("client disconnected")); + }, + () => { + closeCount += 1; + return Promise.resolve(); + }, + ); bridgeWorkbenchStream({ registry, @@ -185,23 +211,22 @@ describe("bridgeWorkbenchStream", () => { }); registry.publish("chan_1", { type: "chat.typing", data: {} }); - // Let the rejected write's `.catch` run before publishing again. - await Promise.resolve(); - await Promise.resolve(); - await Promise.resolve(); + await flush(); registry.publish("chan_1", { type: "chat.typing", data: {} }); - await Promise.resolve(); - await Promise.resolve(); - - // The first publish attempted a write that failed and unsubscribed - // the local subscriber; the second publish must not reach it again - // — a zombie subscriber would keep attempting (and failing) writes - // forever. + await flush(); + + // The first publish attempted a write that failed, unsubscribed the + // local subscriber, and closed the stream so the client's + // `EventSource` actually sees the connection end and reconnects + // (CL-7197) — a zombie subscriber would keep attempting (and + // failing) writes forever, and a stream left open would look live + // while being dead. expect(writeCount).toBe(1); + expect(closeCount).toBe(1); }); - test("a platform subscribe that throws still opens the stream, registry-only", async () => { + test("a platform subscribe that throws still opens the stream, registry-only, and reports the failure", async () => { const registry = createWorkbenchSubscriberRegistry(); const writes: unknown[] = []; const stream = fakeStream((message) => { @@ -213,6 +238,7 @@ describe("bridgeWorkbenchStream", () => { throw new Error("folded run not resolved yet"); }, }; + const report = spyOn(errorSink, "reportError").mockReturnValue("ref_test"); expect(() => bridgeWorkbenchStream({ @@ -225,10 +251,17 @@ describe("bridgeWorkbenchStream", () => { ).not.toThrow(); registry.publish("chan_1", { type: "chat.typing", data: {} }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(writes).toHaveLength(1); + // The bare `catch {}` this degrades through must not swallow the + // failure silently — it reports through `reportError` with context + // a person can act on. + expect(report).toHaveBeenCalledTimes(1); + expect(report.mock.calls[0]?.[1]).toMatchObject({ + operation: "chat.workbenchStream.platformSubscribe", + roomId: "chan_1", + }); }); test("authorize going false unsubscribes both sources and closes the stream, without writing the event", async () => { @@ -263,9 +296,7 @@ describe("bridgeWorkbenchStream", () => { }); registry.publish("chan_1", { type: "chat.typing", data: {} }); - await Promise.resolve(); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(writes).toHaveLength(0); expect(closeCount).toBe(1); @@ -273,13 +304,179 @@ describe("bridgeWorkbenchStream", () => { // A further publish after revocation must not write, or close again. registry.publish("chan_1", { type: "chat.typing", data: {} }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(writes).toHaveLength(0); expect(closeCount).toBe(1); }); + test("an authorize call that rejects is caught, reported, and closes the stream", async () => { + const registry = createWorkbenchSubscriberRegistry(); + let closeCount = 0; + const stream = fakeStream( + () => Promise.resolve(), + () => { + closeCount += 1; + return Promise.resolve(); + }, + ); + const report = spyOn(errorSink, "reportError").mockReturnValue("ref_test"); + const unhandledRejections: unknown[] = []; + const onUnhandledRejection = (reason: unknown) => { + unhandledRejections.push(reason); + }; + process.on("unhandledRejection", onUnhandledRejection); + + try { + bridgeWorkbenchStream({ + registry, + platform: noopPlatformEvents(), + workbenchId: "chan_1", + stream, + authorize: () => Promise.reject(new Error("db blip")), + }); + + registry.publish("chan_1", { type: "chat.typing", data: {} }); + await flush(); + } finally { + process.off("unhandledRejection", onUnhandledRejection); + } + + expect(closeCount).toBe(1); + expect(report).toHaveBeenCalledTimes(1); + expect(report.mock.calls[0]?.[1]).toMatchObject({ + operation: "chat.workbenchStream.authorize", + roomId: "chan_1", + }); + expect(unhandledRejections).toHaveLength(0); + }); + + test("deliveries are serialized: a slow authorize on the first event cannot let a later event write first", async () => { + const registry = createWorkbenchSubscriberRegistry(); + const writes: string[] = []; + const stream = fakeStream((message) => { + const { data } = message as { data: string }; + writes.push(data); + return Promise.resolve(); + }); + let authorizeCalls = 0; + const authorize = (): Promise => { + authorizeCalls += 1; + // The first call (event A's) resolves slowly; every later call + // resolves immediately. Without serialized delivery, B's write + // could land before A's. + if (authorizeCalls === 1) { + return new Promise((resolve) => setTimeout(() => resolve(true), 20)); + } + return Promise.resolve(true); + }; + + bridgeWorkbenchStream({ + registry, + platform: noopPlatformEvents(), + workbenchId: "chan_1", + stream, + authorize, + }); + + registry.publish("chan_1", { type: "chat.typing", data: "A" }); + registry.publish("chan_1", { type: "chat.typing", data: "B" }); + + await new Promise((resolve) => setTimeout(resolve, 50)); + + expect(writes.map((data) => JSON.parse(data) as string)).toEqual([ + "A", + "B", + ]); + }); + + test("the presence snapshot is written before any event racing the setup window", async () => { + const registry = createWorkbenchSubscriberRegistry(); + const presenceRegistry = createWorkbenchPresenceRegistry(); + const writes: { event?: string }[] = []; + const stream = fakeStream((message) => { + writes.push(message as { event?: string }); + return Promise.resolve(); + }); + // A platform whose `subscribeToWorkbench` delivers an event the + // instant it's installed — the exact window in which the old + // implementation's floating snapshot write could still be pending. + const racingPlatform: WorkbenchEvents = { + subscribeToWorkbench(_workbenchId, onEvent) { + onEvent({ type: "chat.agent", data: {} }); + return () => undefined; + }, + }; + + bridgeWorkbenchStream({ + registry, + platform: racingPlatform, + workbenchId: "chan_1", + stream, + authorize: alwaysAuthorized, + presence: { registry: presenceRegistry, principalId: "prn_ada" }, + }); + + await flush(); + + expect(writes[0]?.event).toBe("chat.presence.snapshot"); + }); + + test("the route no longer parks forever: `closed` resolves once teardown runs", async () => { + const registry = createWorkbenchSubscriberRegistry(); + const stream = fakeStream(() => Promise.resolve()); + + const { teardown, closed } = bridgeWorkbenchStream({ + registry, + platform: noopPlatformEvents(), + workbenchId: "chan_1", + stream, + authorize: alwaysAuthorized, + }); + + let settled = false; + void closed.then(() => { + settled = true; + }); + await flush(); + expect(settled).toBe(false); + + teardown(); + await flush(); + + expect(settled).toBe(true); + }); + + test("sends a periodic keepalive and stops once torn down", async () => { + const registry = createWorkbenchSubscriberRegistry(); + const events: (string | undefined)[] = []; + const stream = fakeStream((message) => { + events.push((message as { event?: string }).event); + return Promise.resolve(); + }); + + const { teardown } = bridgeWorkbenchStream({ + registry, + platform: noopPlatformEvents(), + workbenchId: "chan_1", + stream, + authorize: alwaysAuthorized, + keepaliveIntervalMs: 5, + }); + + await new Promise((resolve) => setTimeout(resolve, 30)); + const keepalivesBeforeTeardown = events.filter( + (event) => event === "keepalive", + ).length; + expect(keepalivesBeforeTeardown).toBeGreaterThan(0); + + teardown(); + events.length = 0; + await new Promise((resolve) => setTimeout(resolve, 30)); + + expect(events).toHaveLength(0); + }); + describe("presence", () => { test("connecting delivers a snapshot directly to this stream and broadcasts an online delta", async () => { const registry = createWorkbenchSubscriberRegistry(); @@ -299,8 +496,7 @@ describe("bridgeWorkbenchStream", () => { authorize: alwaysAuthorized, presence: { registry: presenceRegistry, principalId: "prn_ada" }, }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); const snapshotWrite = writes.find( (write) => write.event === "chat.presence.snapshot", @@ -335,7 +531,7 @@ describe("bridgeWorkbenchStream", () => { const observed: ChatWorkbenchEvent[] = []; registry.subscribe("chan_1", (event) => observed.push(event)); - const teardown = bridgeWorkbenchStream({ + const { teardown } = bridgeWorkbenchStream({ registry, platform: noopPlatformEvents(), workbenchId: "chan_1", @@ -343,13 +539,11 @@ describe("bridgeWorkbenchStream", () => { authorize: alwaysAuthorized, presence: { registry: presenceRegistry, principalId: "prn_ada" }, }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); observed.length = 0; teardown(); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(presenceRegistry.snapshot("chan_1")).toEqual([]); const offlineEvent = observed.find( @@ -376,7 +570,7 @@ describe("bridgeWorkbenchStream", () => { return Promise.resolve(); }); - const teardownA = bridgeWorkbenchStream({ + const { teardown: teardownA } = bridgeWorkbenchStream({ registry, platform: noopPlatformEvents(), workbenchId: "chan_1", @@ -392,13 +586,11 @@ describe("bridgeWorkbenchStream", () => { authorize: alwaysAuthorized, presence: { registry: presenceRegistry, principalId: "prn_ada" }, }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); writesB.length = 0; teardownA(); - await Promise.resolve(); - await Promise.resolve(); + await flush(); // Still connected via the second stream — no offline delta. expect( @@ -424,8 +616,7 @@ describe("bridgeWorkbenchStream", () => { stream, authorize: alwaysAuthorized, }); - await Promise.resolve(); - await Promise.resolve(); + await flush(); expect(writes).toHaveLength(0); }); From e6a13e0bbe316f7fd2e2719f6045c402cc3a2b0b Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:21:09 -0700 Subject: [PATCH 2/4] Harden the workbench SSE stream: serialize deliveries, close on teardown, add keepalive Six defects lived in bridgeWorkbenchStream's delivery path: authorize() sat outside the try block and deliveries were floating `void deliver()` calls, so a transient resolver error became an unhandled rejection per event; the write-failure teardown never closed the stream, so the client's EventSource never saw an error and never reconnected, and the route itself parked on a promise that never resolved; a bare `catch {}` around the platform subscribe swallowed failures with no reportError; deliveries weren't sequenced, so two events could land out of publication order depending on when their authorize() calls resolved; and the presence snapshot fired as a floating promise after both subscriptions were installed, letting a racing event overwrite it on the client. Every write (snapshot, event, keepalive) now funnels through a single chained-promise queue so writes land in enqueue order regardless of authorize() timing, bounded by MAX_QUEUED_DELIVERIES to close a stream that has fallen too far behind rather than buffer it unboundedly. The presence snapshot is enqueued before either subscription installs, so nothing can write ahead of it. authorize() and every write path report through reportError and close the stream on failure. bridgeWorkbenchStream now returns { teardown, closed }; the route awaits closed instead of a promise that never resolves. A periodic keepalive keeps idle connections alive behind proxies. --- packages/chat/src/routes.ts | 6 +- packages/chat/src/workbench-events.ts | 201 ++++++++++++++++++++++---- 2 files changed, 173 insertions(+), 34 deletions(-) diff --git a/packages/chat/src/routes.ts b/packages/chat/src/routes.ts index 94656ae36..fc7b5af35 100644 --- a/packages/chat/src/routes.ts +++ b/packages/chat/src/routes.ts @@ -3610,7 +3610,7 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { } return streamSSE(c, async (stream) => { - const unbridge = bridgeWorkbenchStream({ + const { teardown, closed } = bridgeWorkbenchStream({ registry, platform: platformEvents, workbenchId, @@ -3625,8 +3625,8 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { ).then((access) => access !== undefined), presence: { registry: presence, principalId: principal.id }, }); - stream.onAbort(unbridge); - await new Promise(() => undefined); + stream.onAbort(teardown); + await closed; }); }, ); diff --git a/packages/chat/src/workbench-events.ts b/packages/chat/src/workbench-events.ts index ae33e9530..8b4ff8acd 100644 --- a/packages/chat/src/workbench-events.ts +++ b/packages/chat/src/workbench-events.ts @@ -7,14 +7,17 @@ // // `bridgeWorkbenchStream` is the route's one call into this module: it // wires both this registry and the platform's own event stream onto a -// live SSE stream, and owns two defects that matter here. +// live SSE stream, and owns the defects that matter here. // // 1. A write that fails (the client disconnected, but the abort // signal hasn't fired yet) must drop that subscriber immediately // rather than leaving it registered until `stream.onAbort` // eventually runs. A dangling subscriber between disconnect and // abort is a zombie: every event published in that window still -// attempts (and fails) a write. +// attempts (and fails) a write. The same teardown also closes the +// underlying stream so the client's `EventSource` actually sees the +// connection end and reconnects, rather than a stream that looks +// live while nothing is coming through it anymore (CL-7197). // 2. Access must be re-checked on every delivered event, not only at // connect time. The route resolves access once before opening the // stream, but a share or a share member's row can be revoked at @@ -31,11 +34,32 @@ // quiet. A truly instant kill would need a poll/heartbeat // independent of traffic; that's a real, disclosed scope cut, not // a hidden gap, and out of scope here. +// 3. Every delivery — the presence snapshot, each event, and the +// keepalive ping — runs through one chained promise so writes reach +// the stream in the order they were enqueued, never in the order +// their `authorize()` calls happen to resolve. This is also what +// keeps the presence snapshot first: it's enqueued before the local +// and platform subscriptions are installed, so nothing can jump +// ahead of it even if an event arrives the instant a subscription is +// wired up. import type { SSEStreamingApi } from "hono/streaming"; +import { reportError } from "@corbits/error-sink"; import type { WorkbenchEvents, ChatWorkbenchEvent } from "./platform-port"; import { ChatPresenceSnapshotEventData } from "./stream-events"; import type { WorkbenchPresenceRegistry } from "./workbench-presence"; +// Below nginx's 60s default proxy timeout and most load balancers' idle +// timeouts, so an idle connection never looks abandoned to anything +// sitting between the client and this process. +const DEFAULT_KEEPALIVE_INTERVAL_MS = 25_000; + +// A stream whose deliveries have backed up this far has already lost +// coherence with the client (whatever it shows is minutes stale by the +// time it drains) — closing and letting the client reconnect is more +// useful than an unbounded queue holding every event since the backup +// started. +const MAX_QUEUED_DELIVERIES = 200; + export type WorkbenchSubscriber = (event: ChatWorkbenchEvent) => void; export interface WorkbenchSubscriberRegistry { @@ -121,18 +145,34 @@ export function createWorkbenchSubscriberRegistry(): WorkbenchSubscriberRegistry }; } +/** What `bridgeWorkbenchStream` hands back to the route. */ +export interface WorkbenchStreamBridge { + /** Unsubscribes both sources and closes the stream; safe to call more + * than once. The route calls this from `stream.onAbort`. */ + teardown: () => void; + /** Resolves once `teardown` has run, from any cause — client abort, + * a revoked `authorize()`, or a write failure. The route awaits this + * instead of a promise that never resolves, so its closure (and + * everything it closes over) is released once the connection ends + * rather than pinned in memory for the life of the process. */ + closed: Promise; +} + /** * Wires a live SSE stream to both the local registry and the * platform's own per-workbench event stream, and returns the combined - * teardown the route calls from `stream.onAbort`. Before every event + * teardown the route calls from `stream.onAbort` plus a `closed` + * promise it awaits instead of parking forever. Before every event * from either source is written, `authorize()` re-runs the same * fail-closed access check the route ran at connect time; a `false` - * result unsubscribes both sources and closes the stream rather than - * writing the event, so a revoked share or share member stops - * receiving events on the very next one published. A failed - * `writeSSE` (the client is already gone) unsubscribes that source - * immediately rather than only logging and waiting for abort — the - * fix for the zombie-subscriber defect. + * result (or a rejection — a transient resolver error must not become + * an unhandled rejection) unsubscribes both sources and closes the + * stream rather than writing the event, so a revoked share or share + * member stops receiving events on the very next one published. A + * failed `writeSSE` (the client is already gone) unsubscribes that + * source immediately, closes the stream so the client's `EventSource` + * actually reconnects, and reports the failure — the fix for the + * zombie-subscriber and dead-stream defects. */ export function bridgeWorkbenchStream(input: { registry: WorkbenchSubscriberRegistry; @@ -155,11 +195,19 @@ export function bridgeWorkbenchStream(input: { registry: WorkbenchPresenceRegistry; principalId: string; }; -}): () => void { + /** Overrides `DEFAULT_KEEPALIVE_INTERVAL_MS`; exists for tests. */ + keepaliveIntervalMs?: number; +}): WorkbenchStreamBridge { let tornDown = false; + let resolveClosed: () => void = () => undefined; + const closed = new Promise((resolve) => { + resolveClosed = resolve; + }); let unsubscribeLocal: () => void = () => undefined; let unsubscribePlatform: () => void = () => undefined; + let keepaliveTimer: ReturnType | undefined; + const teardownPresence = () => { if (input.presence === undefined) return; const wentOffline = input.presence.registry.disconnect( @@ -182,13 +230,60 @@ export function bridgeWorkbenchStream(input: { unsubscribeLocal(); unsubscribePlatform(); teardownPresence(); + if (keepaliveTimer !== undefined) clearInterval(keepaliveTimer); + resolveClosed(); }; + const closeStream = () => input.stream.close().catch(() => undefined); - const deliver = async (event: ChatWorkbenchEvent) => { + // Every write to this stream — the presence snapshot, each delivered + // event, and the keepalive ping — is chained onto this one promise so + // they land on the wire in the order they were enqueued, never in the + // order their own async work (an `authorize()` call, say) happens to + // settle. `queuedCount` bounds how far a stuck client can make this + // grow. + let deliveryQueue: Promise = Promise.resolve(); + let queuedCount = 0; + + const enqueue = (task: () => Promise) => { if (tornDown) return; - if (!(await input.authorize())) { + if (queuedCount >= MAX_QUEUED_DELIVERIES) { + reportError(new Error("workbench stream delivery queue overflow"), { + operation: "chat.workbenchStream.overflow", + roomId: input.workbenchId, + }); + teardown(); + void closeStream(); + return; + } + queuedCount += 1; + deliveryQueue = deliveryQueue + .then(async () => { + queuedCount -= 1; + if (tornDown) return; + await task(); + }) + // A task's own failure is already handled (reported, torn down) + // inside itself; this catch only exists so a task that somehow + // still throws can't poison every delivery queued after it. + .catch(() => undefined); + }; + + const deliverEvent = async (event: ChatWorkbenchEvent) => { + let authorized: boolean; + try { + authorized = await input.authorize(); + } catch (error) { + reportError(error, { + operation: "chat.workbenchStream.authorize", + roomId: input.workbenchId, + }); teardown(); - await input.stream.close().catch(() => undefined); + await closeStream(); + return; + } + if (!authorized) { + teardown(); + await closeStream(); return; } try { @@ -196,13 +291,44 @@ export function bridgeWorkbenchStream(input: { event: event.type, data: JSON.stringify(event.data), }); - } catch { + } catch (error) { + reportError(error, { + operation: "chat.workbenchStream.write", + roomId: input.workbenchId, + }); teardown(); + await closeStream(); } }; + // Enqueued before either subscription is installed, so nothing + // delivered by either source can be written ahead of it — the fix + // for the snapshot/delta race. + if (input.presence !== undefined) { + const { registry: presenceRegistry, principalId } = input.presence; + presenceRegistry.connect(input.workbenchId, principalId); + const snapshot = ChatPresenceSnapshotEventData.assert({ + members: presenceRegistry.snapshot(input.workbenchId), + }); + enqueue(async () => { + try { + await input.stream.writeSSE({ + event: "chat.presence.snapshot", + data: JSON.stringify(snapshot), + }); + } catch (error) { + reportError(error, { + operation: "chat.workbenchStream.presenceSnapshot", + roomId: input.workbenchId, + }); + teardown(); + await closeStream(); + } + }); + } + unsubscribeLocal = input.registry.subscribe(input.workbenchId, (event) => { - void deliver(event); + enqueue(() => deliverEvent(event)); }); // The platform side resolves a folded run before it can subscribe @@ -210,39 +336,52 @@ export function bridgeWorkbenchStream(input: { // failure there (the run isn't back yet after a hub restart, a slow // DB) must degrade this stream to registry-only rather than take the // whole SSE connection down — a client still gets typing/settings - // events and its own poll fallback covers the rest. + // events and its own poll fallback covers the rest. The degradation + // itself is still a failure worth knowing about, so it's reported + // rather than swallowed. try { unsubscribePlatform = input.platform.subscribeToWorkbench( input.workbenchId, (event) => { - void deliver(event); + enqueue(() => deliverEvent(event)); }, ); - } catch { + } catch (error) { + reportError(error, { + operation: "chat.workbenchStream.platformSubscribe", + roomId: input.workbenchId, + }); unsubscribePlatform = () => undefined; } if (input.presence !== undefined) { - const { registry: presenceRegistry, principalId } = input.presence; - presenceRegistry.connect(input.workbenchId, principalId); - const snapshot = ChatPresenceSnapshotEventData.assert({ - members: presenceRegistry.snapshot(input.workbenchId), - }); - void input.stream - .writeSSE({ - event: "chat.presence.snapshot", - data: JSON.stringify(snapshot), - }) - .catch(() => undefined); input.registry.publish(input.workbenchId, { type: "chat.presence", data: { - principalId, + principalId: input.presence.principalId, state: "online", lastActiveAt: new Date().toISOString(), }, }); } - return teardown; + keepaliveTimer = setInterval(() => { + // A backed-up queue is already proof the connection is alive; + // adding to the backlog would only make it worse. + if (tornDown || queuedCount > 0) return; + enqueue(async () => { + try { + await input.stream.writeSSE({ event: "keepalive", data: "" }); + } catch (error) { + reportError(error, { + operation: "chat.workbenchStream.keepalive", + roomId: input.workbenchId, + }); + teardown(); + await closeStream(); + } + }); + }, input.keepaliveIntervalMs ?? DEFAULT_KEEPALIVE_INTERVAL_MS); + + return { teardown, closed }; } From ddbd2aa726d00b30af19e73c0982243dcc10f629 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:27:04 -0700 Subject: [PATCH 3/4] Update docs: write-failure teardown is inert against Hono's real writeSSE Hono's StreamingApi.write swallows writer errors internally and never rejects, so writeSSE can never throw in production; the write-failure branch in deliverEvent only fires for a stream implementation that does throw (tests today, conceivably a future Hono release). The disconnect signal that actually fires for a real client is stream.onAbort, wired by the route. Correct the comments and the test name that implied write-failure detection was the live mechanism catching a disconnected client. --- packages/chat/src/workbench-events.ts | 31 ++++++++++++--------- packages/chat/test/workbench-events.test.ts | 2 +- 2 files changed, 19 insertions(+), 14 deletions(-) diff --git a/packages/chat/src/workbench-events.ts b/packages/chat/src/workbench-events.ts index 8b4ff8acd..6d5207f84 100644 --- a/packages/chat/src/workbench-events.ts +++ b/packages/chat/src/workbench-events.ts @@ -9,15 +9,17 @@ // wires both this registry and the platform's own event stream onto a // live SSE stream, and owns the defects that matter here. // -// 1. A write that fails (the client disconnected, but the abort -// signal hasn't fired yet) must drop that subscriber immediately -// rather than leaving it registered until `stream.onAbort` -// eventually runs. A dangling subscriber between disconnect and -// abort is a zombie: every event published in that window still -// attempts (and fails) a write. The same teardown also closes the -// underlying stream so the client's `EventSource` actually sees the -// connection end and reconnects, rather than a stream that looks -// live while nothing is coming through it anymore (CL-7197). +// 1. A write that throws must drop that subscriber and close the +// stream immediately rather than leaving it registered until +// `stream.onAbort` eventually runs — a dangling subscriber between +// disconnect and abort is a zombie: every event published in that +// window still attempts a write. In practice Hono's own +// `StreamingApi.write` swallows writer errors internally and never +// rejects, so `stream.onAbort` (wired by the route) is the only +// signal that actually fires for a disconnected client today; this +// path exists for any write that does throw (a custom `stream` +// passed in tests, or a future Hono release that stops swallowing) +// and is otherwise inert rather than load-bearing (CL-7197). // 2. Access must be re-checked on every delivered event, not only at // connect time. The route resolves access once before opening the // stream, but a share or a share member's row can be revoked at @@ -169,10 +171,13 @@ export interface WorkbenchStreamBridge { * an unhandled rejection) unsubscribes both sources and closes the * stream rather than writing the event, so a revoked share or share * member stops receiving events on the very next one published. A - * failed `writeSSE` (the client is already gone) unsubscribes that - * source immediately, closes the stream so the client's `EventSource` - * actually reconnects, and reports the failure — the fix for the - * zombie-subscriber and dead-stream defects. + * `writeSSE` that throws (the client is already gone) unsubscribes + * that source immediately, closes the stream, and reports the + * failure — Hono's own `write` swallows writer errors today, so in + * practice `stream.onAbort` is what actually catches a disconnected + * client, but this path covers any write that does throw rather than + * silently degrading. A periodic keepalive keeps idle connections + * alive behind proxies that time out on silence. */ export function bridgeWorkbenchStream(input: { registry: WorkbenchSubscriberRegistry; diff --git a/packages/chat/test/workbench-events.test.ts b/packages/chat/test/workbench-events.test.ts index 6a9a5e951..574aee2e3 100644 --- a/packages/chat/test/workbench-events.test.ts +++ b/packages/chat/test/workbench-events.test.ts @@ -187,7 +187,7 @@ describe("bridgeWorkbenchStream", () => { expect(writes).toHaveLength(1); }); - test("a subscriber whose write rejects is removed and the stream is closed", async () => { + test("a write that throws removes that subscriber and closes the stream (Hono's own writeSSE never rejects; this covers a stream implementation that does)", async () => { const registry = createWorkbenchSubscriberRegistry(); let writeCount = 0; let closeCount = 0; From f9f8455ebcbd18ed7399eeeb1ff2a01b141e33f5 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:45:29 -0700 Subject: [PATCH 4/4] Initialize the keepalive timer handle at its declaration --- packages/chat/src/workbench-events.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/chat/src/workbench-events.ts b/packages/chat/src/workbench-events.ts index 6d5207f84..eb838adb6 100644 --- a/packages/chat/src/workbench-events.ts +++ b/packages/chat/src/workbench-events.ts @@ -211,7 +211,7 @@ export function bridgeWorkbenchStream(input: { let unsubscribeLocal: () => void = () => undefined; let unsubscribePlatform: () => void = () => undefined; - let keepaliveTimer: ReturnType | undefined; + let keepaliveTimer: ReturnType | undefined = undefined; const teardownPresence = () => { if (input.presence === undefined) return;