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..eb838adb6 100644 --- a/packages/chat/src/workbench-events.ts +++ b/packages/chat/src/workbench-events.ts @@ -7,14 +7,19 @@ // // `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. +// 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 @@ -31,11 +36,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 +147,37 @@ 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 + * `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; @@ -155,11 +200,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 = undefined; + const teardownPresence = () => { if (input.presence === undefined) return; const wentOffline = input.presence.registry.disconnect( @@ -182,13 +235,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 +296,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 +341,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 }; } diff --git a/packages/chat/test/workbench-events.test.ts b/packages/chat/test/workbench-events.test.ts index 9f2f02446..574aee2e3 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 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; - 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); });