Skip to content

fix(sdk): a throwing subscriber silently stops delivery to the others #473

Description

@EricAndrechek

Area: sdk — streaming. Surfaced reviewing #470; the crash it caused is fixed there, this is the residual.

StreamController fans out to subscribers with a bare loop (clients/ts/src/stream/controller.ts:28, :45, :58):

for (const sub of this._subscribers) {
  sub.status?.(status);
}

One subscriber throwing aborts the loop, so every subscriber registered after it silently misses that event, status, or error. Iteration order is insertion order, so which consumers are affected depends on the order they happened to subscribe in.

Why now

#470 hit the severe form of this. The transport called onStatus/onError unguarded, so a throw propagated out of the fan-out, unwound the reconnect loop, and ended as a process-fatal unhandled rejection. That is fixed in the transport (sse.ts now isolates all three callbacks, matching what _dispatch already did), which contains the blast radius to "the stream survives".

What it does not fix is the fan-out itself: with two subscribers on one stream, a throw in the first still means the second never hears about that frame. Silent, order-dependent, and invisible to the throwing consumer.

The full surface

Five behaviours, found across three review rounds on #470 (details and execution traces in the comments below). #470 documents all of these in sdk/reference.md and sdk/streaming.md but changes none of them, so the docs currently point here for every one.

# Where What a throw costs
1 controller.ts:28, :45, :58 — the three fan-out loops Every subscriber registered after the thrower misses that event/status/error
2 Same loop, :27-38 A concurrent for await loses the event entirely — the fan-out precedes both _waiters.shift() and _buffer.push(), so it is neither delivered nor queued
3 Same loop, :40-55 On the terminal closed status the aborted loop skips _done = true and the waiter resolution at :48-54, so a for await never exits. Only a later close() heals it (:158-162 sit outside the status guard); unsub() cannot, because auto-close at :134 requires _waiters.length === 0 and the hung iterator is itself a waiter
4 controller.ts:130 — subscribe()'s initial synchronous status call Unguarded, so it throws back out of the call the consumer made. Out of .subscribe(): subscriber registered, no unsubscribe returned. Out of .liveQuery() (which forwards status at live-query.ts:41): no handle returned at all, so nothing can .close() it; _runBackfill never starts so initial() never fires; and the stream opened a line earlier keeps running — connected, reconnecting on its own — with only opts.signal able to stop it, buffering every event into _buffer because _buffering never flips
5 live-query.ts:87-91 The backfill's blanket catch absorbs a throw from initial() or from next() mid-flush and a rejection from the fetch itself (a rejecting auth, a relative baseURL) and clears _buffer, silently dropping the rest. :58-61 drops it a second way — an error Result returns early with _buffer still populated and never flushed, bypassing the catch entirely

Symptom worth naming because it presents as its own puzzle: a throwing status handler registered before .connected() makes that promise reject with its timeout while .status already reads live — connected()'s watcher (:123) is added after the consumer's subscriber, so the throw aborts the fan-out before it observes live. Same mechanism as 1; no separate fix needed.

Scope

Isolate per subscriber in all three loops, so one bad callback costs only its own delivery:

for (const sub of this._subscribers) {
  try { sub.status?.(status); } catch (e) { console.error("[wavehouse] subscriber status handler threw:", e); }
}

That covers 1 and 2. It is not sufficient on its own — three things need separate decisions:

  • 3 needs the bookkeeping made unreachable-proof. Move _done = true and the waiter resolution ahead of the fan-out, or into a finally. Iterator termination should not depend on subscriber behaviour.
  • 4 needs the guard at controller.ts:130, not at subscribe()'s return. "Return the unsubscribe function anyway" does not help the liveQuery() case, because LiveQuery's constructor discards the return value on the way out (live-query.ts:32 assigns _unsubStream only after subscribe() returns). Guarding the call site fixes both at once.
  • 5 needs the two paths distinguished. A fix aimed only at the catch at :87 misses the error-Result early return at :58-61.

Worth deciding at the same time:

  • Whether to report rather than only log. A consumer whose handler throws currently gets no signal at all. There is no obvious channel for it — routing it to error risks a loop if the error handler is the one throwing.
  • Whether next should differ from status/error. A throw in next is the likeliest in practice (it runs consumer rendering/business logic) and the most costly to swallow silently.
  • Whether a backfill error Result should flush the buffer rather than drop it. Arguably the events are still valid even though the historical query failed.

Acceptance

  • A throwing subscriber does not prevent delivery to other subscribers, on all three of next/status/error
  • A throwing subscriber does not cost a concurrent for await its event
  • A throwing status handler on the terminal closed status still terminates a concurrent for await
  • A throwing status handler does not escape .subscribe() or .liveQuery(); both still return their handle
  • A throwing initial()/next() during the backfill flush does not discard the remaining buffered events, and a backfill error Result is handled deliberately either way
  • Pinned by tests, including two subscribers where the first throws, and a for await concurrent with a throwing subscriber
  • sdk/reference.md ("If your own callback throws") and sdk/streaming.md updated — they currently document all five as live behaviour and link here

Related: #470 (fixed the transport-side crash, documented the rest), #389 (StreamController buffering, same file), #449 (live-query dedup boundary, same file).

Activity

  1. added
    bugSomething isn't working
    area/sdkTypeScript SDK (clients/ts/)
    area/streamingSSE / live-query delivery path (/v1/stream)
    on Aug 13, 2026
  2. EricAndrechek commented on Aug 13, 2026

    @EricAndrechek
    MemberAuthor

    Scope is wider than the three fan-out loops — three more unguarded paths

    Found while documenting the throwing-handler contract on #470. This issue as filed named only the _subscribers loops at controller.ts:28, :45, :58. There are three more places a consumer callback runs unprotected, and the first is the one most likely to fire.

    1. subscribe()'s initial status call (clients/ts/src/stream/controller.ts:130)

    subscribe(subscriber: StreamSubscriber<T>): () => void {
      this._subscribers.add(subscriber);
      subscriber.status?.(this._status);   // synchronous, unguarded
      return () => { … };
    }

    This runs before the transport is involved at all, so the transport's guards can't help. A throw propagates straight back out of .subscribe(), and leaves the caller in a genuinely bad state: the subscriber is already registered (so it keeps receiving), but the unsubscribe function was never returned, so it can never be removed. It also always fires — every subscribe() delivers the current state immediately — which makes it the first thing a throwing status handler does, not an edge case. The docs' own example handler is updateIndicator(state), exactly the shape that throws when a DOM node is missing.

    2. A concurrent for await is starved, not just later subscribers

    controller.ts:28-36 runs the _subscribers loop before resolving _waiters / pushing to _buffer. So a throw in any subscriber means the async-iterator consumer never sees that event either — the iterator isn't a separate delivery path, it's downstream of the loop. The issue text said "later subscribers"; it should say "later subscribers and any for await consumer".

    3. liveQuery()'s backfill flush swallows and discards (clients/ts/src/stream/live-query.ts:56, 84-91)

    _runBackfill wraps the initial() call and the whole buffered-event flush in one try, whose catch sets _buffer = []. So a throw from initial(), or from next() partway through the flush, is absorbed and takes the remaining buffered events with it — silently, with no console.error, no error callback, and no way for the consumer to know events were dropped. That is worse than the fan-out case: not just delayed or skipped delivery, but data loss on the backfill seam the feature exists to provide.

    What this changes about the fix

    Guarding only the three loops would leave the most-likely path (1) and the most-damaging path (3) unfixed, while making the docs' "isolated and logged" claim look complete. Worth deciding together:

    • Path 1 needs a decision the loops don't: if the initial status throws, does subscribe() still return the unsubscribe function (recommended — the caller needs it precisely because their handler is broken), or roll back the registration?
    • Path 3 needs to distinguish "the fetch failed" from "your callback threw" — they currently share one catch, and only the first justifies clearing the buffer.

    #470 documents all three as known carve-outs in sdk/reference.md §Error Handling rather than fixing them, since they live in controller.ts / live-query.ts rather than the transport. Whatever lands here should shorten that section.

    — Claude Opus 5, via Claude Code

  3. EricAndrechek commented on Aug 13, 2026

    @EricAndrechek
    MemberAuthor

    One more consequence of path 1 (subscribe()'s unguarded initial status call), found reviewing #470 and worth having here before anyone picks this up — it changes what a fix has to do.

    Through .liveQuery(), the same line strands an unclosable connection.

    LiveQuery's constructor subscribes with status: (s) => subscriber.status?.(s) (live-query.ts:41), so a caller's status handler is on the far end of controller.ts:130 exactly as it is for a bare .subscribe(). The difference is what the throw unwinds through:

    this._unsubStream = stream.subscribe({ … });   // live-query.ts:32

    The assignment happens after subscribe() returns. A throw from that first synchronous status call therefore unwinds before _unsubStream is ever set, out through the LiveQuery constructor, out through client.liveQuery(). The caller gets no object at all — while the controller has already registered the subscriber (controller.ts:129, before the status call) and the SSE connection is open and re-dialing. There is no handle to .close() it with, and _unsubStream?.() at :99 is unreachable because there is no instance.

    So it is strictly worse than the bare-subscribe() case: that one at least leaves the caller holding a live stream they can close, just without the unsubscribe function for the one subscriber.

    Why this bears on the open design question. "Should subscribe() still return the unsubscribe function when the handler threw?" isn't sufficient on its own — LiveQuery discards the return value on the way out of the constructor, so fixing only subscribe()'s return path leaves the .liveQuery() case exactly as broken. Whatever the fix is, it has to hold for a caller who never receives the object that owns the handle. Guarding the initial status call at controller.ts:130 the way the transport guards its own callbacks would cover both at once, which is an argument for fixing it there rather than at the return.

    Pre-existing on main and outside #470's diff, so not fixed there. Verified against live-query.ts:19,32,41,99 and controller.ts:129-131.

  4. EricAndrechek commented on Aug 13, 2026

    @EricAndrechek
    MemberAuthor

    Three more behaviours found while re-reviewing the #470 docs, all executed rather than read. None are new regressions — they are all pre-existing on main — but they change the scope of a fix, so recording them here rather than letting them be rediscovered.

    1. For a concurrent for await, the event is dropped, not merely delayed.

    controller.ts:27-38 runs the subscriber fan-out at :28-30 before both _waiters.shift() at :32 and _buffer.push() at :36. A throw at :29 exits the whole onEvent arrow function, so the event never reaches a waiting iterator and never lands in the buffer for a later next(). It is gone. "Starved" was the wrong word for this — the iterator isn't waiting on something late, it is missing something that will never arrive.

    2. A throw on the terminal closed status leaves a for await hanging forever.

    Same shape, worse consequence. controller.ts:40-55 sets _status at :44, fans out at :45-47, and only then, at :48-54, sets _done = true and resolves the outstanding waiters. A throwing status handler aborts the loop at :46, so _done stays false and every waiter stays pending — against a stream the transport has already torn down. This is reachable from every transport-initiated terminal close: a 4xx, SSE_REDIRECT, SSE_BAD_CONTENT_TYPE.

    It is recoverable, but only by something calling close(): controller.ts:150-163 sets _done and resolves waiters outside the if (this._status !== "closed") guard, so on a second pass the fan-out is skipped, nothing throws, and the iterator terminates. A caller who is sitting in for await and not watching status has no reason to make that call.

    3. .connected() can reject with its timeout against a stream that is already live.

    connected() appends its internal watcher to _subscribers at controller.ts:123. In the ordering the docs themselves demonstrate — subscribe(...) and then await stream.connected() — the user's subscriber is earlier in the set, so a throw from its status handler aborts the fan-out before the watcher is reached. The watcher never observes live, the timer fires, and the promise rejects with Stream did not connect within Nms while controller.status reads live.

    That makes the rejection mean neither of the two things .connected() documents it as meaning.


    Bearing on the fix. (1) and (3) are both consequences of the fan-out loops being unguarded, so guarding each sub.next(...) / sub.status?.(...) / sub.error?.(...) the way the transport guards its own callbacks resolves them together. (2) needs one thing more: the _done/waiter bookkeeping at :48-54 has to be reachable even when a subscriber throws — moving it before the fan-out, or into a finally, since correctness of the iterator's termination shouldn't depend on subscriber behaviour.

    Also worth noting for whoever picks this up: live-query.ts:58-61 drops the buffer on a backfill error Result without going through the catch at :87-91 — _buffering = false and return, with _buffer still populated and never flushed or cleared. Same user-visible outcome as the throwing-handler case (buffered events silently lost), different path, so a fix aimed only at the catch would miss it.

    Verified against controller.ts:27-55,123,150-163 and live-query.ts:56-91 at a67c0648. #470 documents all of this in sdk/reference.md and sdk/streaming.md but does not change the behaviour.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    area/sdkTypeScript SDK (clients/ts/)area/streamingSSE / live-query delivery path (/v1/stream)bugSomething isn't working

    Type

    No type

    Projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions