@@ -223,13 +223,15 @@ function createSession(deps: SessionDeps): CrawlSession {
223223
224224 // ── HookEvent fan-out ──────────────────────────────────────────────────
225225 //
226- // Internal queue + iterator that resolves either a buffered event or the
227- // next emit. Keeping a single queue means hook subscribers and iter
228- // consumers stay in sync: every stable event is pushed exactly once.
229- const queue : HookEvent [ ] = [ ]
230- let resolveNext : ( ( v : IteratorResult < HookEvent > ) => void ) | null = null
231- let iterDone = false
226+ // Every stable event fans out to all registered `handlers` (multicast). Both
227+ // WS broadcast and each `events` async-iterator consumer register here, so a
228+ // second consumer (a second dashboard tab, a CLI tail) sees the full stream
229+ // instead of stealing events from the first — the iterator is per-consumer,
230+ // not a single shared queue.
232231 const handlers = new Set < ( event : HookEvent ) => void > ( )
232+ // Live `events` iterators, closed when the scan ends so their `for await`
233+ // loops terminate.
234+ const iterClosers = new Set < ( ) => void > ( )
233235
234236 // In-memory ring buffer (cap 10k) for `events.subscribe.replay`.
235237 const RING_CAP = 10_000
@@ -248,39 +250,64 @@ function createSession(deps: SessionDeps): CrawlSession {
248250 logOperationalWarn ( 'core.scan_event_subscriber_failed' , err , { scanId } , deps . logger )
249251 }
250252 }
251- if ( iterDone )
252- return
253- if ( resolveNext ) {
254- const r = resolveNext
255- resolveNext = null
256- r ( { value : event , done : false } )
257- }
258- else {
259- queue . push ( event )
260- }
261253 }
262254
263255 function closeIter ( ) : void {
264- iterDone = true
265- if ( resolveNext ) {
266- const r = resolveNext
267- resolveNext = null
268- r ( { value : undefined , done : true } )
269- }
256+ for ( const close of [ ...iterClosers ] )
257+ close ( )
270258 }
271259
272260 const events : AsyncIterable < HookEvent > = {
273- [ Symbol . asyncIterator ] ( ) {
261+ [ Symbol . asyncIterator ] ( ) : AsyncIterator < HookEvent > {
262+ // Per-consumer buffer + waiter, fed by its own subscription. Cleaned up
263+ // on scan end (iterClosers) or when the consumer stops (`return()` — the
264+ // HTTP layer calls it on client disconnect so an abandoned tail unsubs
265+ // and stops holding a slot).
266+ const localQueue : HookEvent [ ] = [ ]
267+ let localResolve : ( ( v : IteratorResult < HookEvent > ) => void ) | null = null
268+ let done = false
269+
270+ const unsub = subscribe ( ( event ) => {
271+ if ( done )
272+ return
273+ if ( localResolve ) {
274+ const r = localResolve
275+ localResolve = null
276+ r ( { value : event , done : false } )
277+ }
278+ else {
279+ localQueue . push ( event )
280+ }
281+ } )
282+
283+ function finish ( ) : void {
284+ if ( done )
285+ return
286+ done = true
287+ unsub ( )
288+ iterClosers . delete ( finish )
289+ if ( localResolve ) {
290+ const r = localResolve
291+ localResolve = null
292+ r ( { value : undefined , done : true } )
293+ }
294+ }
295+ iterClosers . add ( finish )
296+
274297 return {
275298 next ( ) : Promise < IteratorResult < HookEvent > > {
276- if ( queue . length )
277- return Promise . resolve ( { value : queue . shift ( ) ! , done : false } )
278- if ( iterDone )
299+ if ( localQueue . length )
300+ return Promise . resolve ( { value : localQueue . shift ( ) ! , done : false } )
301+ if ( done )
279302 return Promise . resolve ( { value : undefined , done : true } )
280303 return new Promise ( ( r ) => {
281- resolveNext = r
304+ localResolve = r
282305 } )
283306 } ,
307+ return ( ) : Promise < IteratorResult < HookEvent > > {
308+ finish ( )
309+ return Promise . resolve ( { value : undefined , done : true } )
310+ } ,
284311 }
285312 } ,
286313 }
0 commit comments