77//
88// `bridgeWorkbenchStream` is the route's one call into this module: it
99// wires both this registry and the platform's own event stream onto a
10- // live SSE stream, and owns two defects that matter here.
10+ // live SSE stream, and owns the defects that matter here.
1111//
12- // 1. A write that fails (the client disconnected, but the abort
13- // signal hasn't fired yet) must drop that subscriber immediately
14- // rather than leaving it registered until `stream.onAbort`
15- // eventually runs. A dangling subscriber between disconnect and
16- // abort is a zombie: every event published in that window still
17- // attempts (and fails) a write.
12+ // 1. A write that throws must drop that subscriber and close the
13+ // stream immediately rather than leaving it registered until
14+ // `stream.onAbort` eventually runs — a dangling subscriber between
15+ // disconnect and abort is a zombie: every event published in that
16+ // window still attempts a write. In practice Hono's own
17+ // `StreamingApi.write` swallows writer errors internally and never
18+ // rejects, so `stream.onAbort` (wired by the route) is the only
19+ // signal that actually fires for a disconnected client today; this
20+ // path exists for any write that does throw (a custom `stream`
21+ // passed in tests, or a future Hono release that stops swallowing)
22+ // and is otherwise inert rather than load-bearing (CL-7197).
1823// 2. Access must be re-checked on every delivered event, not only at
1924// connect time. The route resolves access once before opening the
2025// stream, but a share or a share member's row can be revoked at
3136// quiet. A truly instant kill would need a poll/heartbeat
3237// independent of traffic; that's a real, disclosed scope cut, not
3338// a hidden gap, and out of scope here.
39+ // 3. Every delivery — the presence snapshot, each event, and the
40+ // keepalive ping — runs through one chained promise so writes reach
41+ // the stream in the order they were enqueued, never in the order
42+ // their `authorize()` calls happen to resolve. This is also what
43+ // keeps the presence snapshot first: it's enqueued before the local
44+ // and platform subscriptions are installed, so nothing can jump
45+ // ahead of it even if an event arrives the instant a subscription is
46+ // wired up.
3447import type { SSEStreamingApi } from "hono/streaming" ;
48+ import { reportError } from "@corbits/error-sink" ;
3549import type { WorkbenchEvents , ChatWorkbenchEvent } from "./platform-port" ;
3650import { ChatPresenceSnapshotEventData } from "./stream-events" ;
3751import type { WorkbenchPresenceRegistry } from "./workbench-presence" ;
3852
53+ // Below nginx's 60s default proxy timeout and most load balancers' idle
54+ // timeouts, so an idle connection never looks abandoned to anything
55+ // sitting between the client and this process.
56+ const DEFAULT_KEEPALIVE_INTERVAL_MS = 25_000 ;
57+
58+ // A stream whose deliveries have backed up this far has already lost
59+ // coherence with the client (whatever it shows is minutes stale by the
60+ // time it drains) — closing and letting the client reconnect is more
61+ // useful than an unbounded queue holding every event since the backup
62+ // started.
63+ const MAX_QUEUED_DELIVERIES = 200 ;
64+
3965export type WorkbenchSubscriber = ( event : ChatWorkbenchEvent ) => void ;
4066
4167export interface WorkbenchSubscriberRegistry {
@@ -121,18 +147,37 @@ export function createWorkbenchSubscriberRegistry(): WorkbenchSubscriberRegistry
121147 } ;
122148}
123149
150+ /** What `bridgeWorkbenchStream` hands back to the route. */
151+ export interface WorkbenchStreamBridge {
152+ /** Unsubscribes both sources and closes the stream; safe to call more
153+ * than once. The route calls this from `stream.onAbort`. */
154+ teardown : ( ) => void ;
155+ /** Resolves once `teardown` has run, from any cause — client abort,
156+ * a revoked `authorize()`, or a write failure. The route awaits this
157+ * instead of a promise that never resolves, so its closure (and
158+ * everything it closes over) is released once the connection ends
159+ * rather than pinned in memory for the life of the process. */
160+ closed : Promise < void > ;
161+ }
162+
124163/**
125164 * Wires a live SSE stream to both the local registry and the
126165 * platform's own per-workbench event stream, and returns the combined
127- * teardown the route calls from `stream.onAbort`. Before every event
166+ * teardown the route calls from `stream.onAbort` plus a `closed`
167+ * promise it awaits instead of parking forever. Before every event
128168 * from either source is written, `authorize()` re-runs the same
129169 * fail-closed access check the route ran at connect time; a `false`
130- * result unsubscribes both sources and closes the stream rather than
131- * writing the event, so a revoked share or share member stops
132- * receiving events on the very next one published. A failed
133- * `writeSSE` (the client is already gone) unsubscribes that source
134- * immediately rather than only logging and waiting for abort — the
135- * fix for the zombie-subscriber defect.
170+ * result (or a rejection — a transient resolver error must not become
171+ * an unhandled rejection) unsubscribes both sources and closes the
172+ * stream rather than writing the event, so a revoked share or share
173+ * member stops receiving events on the very next one published. A
174+ * `writeSSE` that throws (the client is already gone) unsubscribes
175+ * that source immediately, closes the stream, and reports the
176+ * failure — Hono's own `write` swallows writer errors today, so in
177+ * practice `stream.onAbort` is what actually catches a disconnected
178+ * client, but this path covers any write that does throw rather than
179+ * silently degrading. A periodic keepalive keeps idle connections
180+ * alive behind proxies that time out on silence.
136181 */
137182export function bridgeWorkbenchStream ( input : {
138183 registry : WorkbenchSubscriberRegistry ;
@@ -155,11 +200,19 @@ export function bridgeWorkbenchStream(input: {
155200 registry : WorkbenchPresenceRegistry ;
156201 principalId : string ;
157202 } ;
158- } ) : ( ) => void {
203+ /** Overrides `DEFAULT_KEEPALIVE_INTERVAL_MS`; exists for tests. */
204+ keepaliveIntervalMs ?: number ;
205+ } ) : WorkbenchStreamBridge {
159206 let tornDown = false ;
207+ let resolveClosed : ( ) => void = ( ) => undefined ;
208+ const closed = new Promise < void > ( ( resolve ) => {
209+ resolveClosed = resolve ;
210+ } ) ;
160211
161212 let unsubscribeLocal : ( ) => void = ( ) => undefined ;
162213 let unsubscribePlatform : ( ) => void = ( ) => undefined ;
214+ let keepaliveTimer : ReturnType < typeof setInterval > | undefined = undefined ;
215+
163216 const teardownPresence = ( ) => {
164217 if ( input . presence === undefined ) return ;
165218 const wentOffline = input . presence . registry . disconnect (
@@ -182,67 +235,158 @@ export function bridgeWorkbenchStream(input: {
182235 unsubscribeLocal ( ) ;
183236 unsubscribePlatform ( ) ;
184237 teardownPresence ( ) ;
238+ if ( keepaliveTimer !== undefined ) clearInterval ( keepaliveTimer ) ;
239+ resolveClosed ( ) ;
185240 } ;
241+ const closeStream = ( ) => input . stream . close ( ) . catch ( ( ) => undefined ) ;
186242
187- const deliver = async ( event : ChatWorkbenchEvent ) => {
243+ // Every write to this stream — the presence snapshot, each delivered
244+ // event, and the keepalive ping — is chained onto this one promise so
245+ // they land on the wire in the order they were enqueued, never in the
246+ // order their own async work (an `authorize()` call, say) happens to
247+ // settle. `queuedCount` bounds how far a stuck client can make this
248+ // grow.
249+ let deliveryQueue : Promise < void > = Promise . resolve ( ) ;
250+ let queuedCount = 0 ;
251+
252+ const enqueue = ( task : ( ) => Promise < void > ) => {
188253 if ( tornDown ) return ;
189- if ( ! ( await input . authorize ( ) ) ) {
254+ if ( queuedCount >= MAX_QUEUED_DELIVERIES ) {
255+ reportError ( new Error ( "workbench stream delivery queue overflow" ) , {
256+ operation : "chat.workbenchStream.overflow" ,
257+ roomId : input . workbenchId ,
258+ } ) ;
259+ teardown ( ) ;
260+ void closeStream ( ) ;
261+ return ;
262+ }
263+ queuedCount += 1 ;
264+ deliveryQueue = deliveryQueue
265+ . then ( async ( ) => {
266+ queuedCount -= 1 ;
267+ if ( tornDown ) return ;
268+ await task ( ) ;
269+ } )
270+ // A task's own failure is already handled (reported, torn down)
271+ // inside itself; this catch only exists so a task that somehow
272+ // still throws can't poison every delivery queued after it.
273+ . catch ( ( ) => undefined ) ;
274+ } ;
275+
276+ const deliverEvent = async ( event : ChatWorkbenchEvent ) => {
277+ let authorized : boolean ;
278+ try {
279+ authorized = await input . authorize ( ) ;
280+ } catch ( error ) {
281+ reportError ( error , {
282+ operation : "chat.workbenchStream.authorize" ,
283+ roomId : input . workbenchId ,
284+ } ) ;
190285 teardown ( ) ;
191- await input . stream . close ( ) . catch ( ( ) => undefined ) ;
286+ await closeStream ( ) ;
287+ return ;
288+ }
289+ if ( ! authorized ) {
290+ teardown ( ) ;
291+ await closeStream ( ) ;
192292 return ;
193293 }
194294 try {
195295 await input . stream . writeSSE ( {
196296 event : event . type ,
197297 data : JSON . stringify ( event . data ) ,
198298 } ) ;
199- } catch {
299+ } catch ( error ) {
300+ reportError ( error , {
301+ operation : "chat.workbenchStream.write" ,
302+ roomId : input . workbenchId ,
303+ } ) ;
200304 teardown ( ) ;
305+ await closeStream ( ) ;
201306 }
202307 } ;
203308
309+ // Enqueued before either subscription is installed, so nothing
310+ // delivered by either source can be written ahead of it — the fix
311+ // for the snapshot/delta race.
312+ if ( input . presence !== undefined ) {
313+ const { registry : presenceRegistry , principalId } = input . presence ;
314+ presenceRegistry . connect ( input . workbenchId , principalId ) ;
315+ const snapshot = ChatPresenceSnapshotEventData . assert ( {
316+ members : presenceRegistry . snapshot ( input . workbenchId ) ,
317+ } ) ;
318+ enqueue ( async ( ) => {
319+ try {
320+ await input . stream . writeSSE ( {
321+ event : "chat.presence.snapshot" ,
322+ data : JSON . stringify ( snapshot ) ,
323+ } ) ;
324+ } catch ( error ) {
325+ reportError ( error , {
326+ operation : "chat.workbenchStream.presenceSnapshot" ,
327+ roomId : input . workbenchId ,
328+ } ) ;
329+ teardown ( ) ;
330+ await closeStream ( ) ;
331+ }
332+ } ) ;
333+ }
334+
204335 unsubscribeLocal = input . registry . subscribe ( input . workbenchId , ( event ) => {
205- void deliver ( event ) ;
336+ enqueue ( ( ) => deliverEvent ( event ) ) ;
206337 } ) ;
207338
208339 // The platform side resolves a folded run before it can subscribe
209340 // (see `subscribeToWorkbench` in `platform-adapter.ts`); a transient
210341 // failure there (the run isn't back yet after a hub restart, a slow
211342 // DB) must degrade this stream to registry-only rather than take the
212343 // whole SSE connection down — a client still gets typing/settings
213- // events and its own poll fallback covers the rest.
344+ // events and its own poll fallback covers the rest. The degradation
345+ // itself is still a failure worth knowing about, so it's reported
346+ // rather than swallowed.
214347 try {
215348 unsubscribePlatform = input . platform . subscribeToWorkbench (
216349 input . workbenchId ,
217350 ( event ) => {
218- void deliver ( event ) ;
351+ enqueue ( ( ) => deliverEvent ( event ) ) ;
219352 } ,
220353 ) ;
221- } catch {
354+ } catch ( error ) {
355+ reportError ( error , {
356+ operation : "chat.workbenchStream.platformSubscribe" ,
357+ roomId : input . workbenchId ,
358+ } ) ;
222359 unsubscribePlatform = ( ) => undefined ;
223360 }
224361
225362 if ( input . presence !== undefined ) {
226- const { registry : presenceRegistry , principalId } = input . presence ;
227- presenceRegistry . connect ( input . workbenchId , principalId ) ;
228- const snapshot = ChatPresenceSnapshotEventData . assert ( {
229- members : presenceRegistry . snapshot ( input . workbenchId ) ,
230- } ) ;
231- void input . stream
232- . writeSSE ( {
233- event : "chat.presence.snapshot" ,
234- data : JSON . stringify ( snapshot ) ,
235- } )
236- . catch ( ( ) => undefined ) ;
237363 input . registry . publish ( input . workbenchId , {
238364 type : "chat.presence" ,
239365 data : {
240- principalId,
366+ principalId : input . presence . principalId ,
241367 state : "online" ,
242368 lastActiveAt : new Date ( ) . toISOString ( ) ,
243369 } ,
244370 } ) ;
245371 }
246372
247- return teardown ;
373+ keepaliveTimer = setInterval ( ( ) => {
374+ // A backed-up queue is already proof the connection is alive;
375+ // adding to the backlog would only make it worse.
376+ if ( tornDown || queuedCount > 0 ) return ;
377+ enqueue ( async ( ) => {
378+ try {
379+ await input . stream . writeSSE ( { event : "keepalive" , data : "" } ) ;
380+ } catch ( error ) {
381+ reportError ( error , {
382+ operation : "chat.workbenchStream.keepalive" ,
383+ roomId : input . workbenchId ,
384+ } ) ;
385+ teardown ( ) ;
386+ await closeStream ( ) ;
387+ }
388+ } ) ;
389+ } , input . keepaliveIntervalMs ?? DEFAULT_KEEPALIVE_INTERVAL_MS ) ;
390+
391+ return { teardown, closed } ;
248392}
0 commit comments