Area: sdk — streaming. Pre-existing on main; surfaced reviewing #470 and deliberately not fixed there.
StreamController retains every event it has ever received unless something is consuming it through the async iterator. A long-lived .subscribe()-only stream — the pattern our own docs lead with — grows without bound until the process dies or the stream is closed.
Mechanism
clients/ts/src/stream/controller.ts:31-37, in the onEvent handler:
const waiter = this._waiters.shift();
if (waiter) {
waiter.resolve({ value: event, done: false });
} else {
this._buffer.push(event); // :36
}
_buffer has exactly four references in the file (:20 declaration, :36 push, :170-171 read/shift). The only drain is the async iterator's next(). A consumer who calls .subscribe() and never writes for await never pushes a waiter and never calls next(), so the else branch is taken on every single event and nothing ever removes it. There is no cap.
Callback delivery still works — the fan-out at :28-30 happens first — so the stream behaves correctly and leaks silently.
Why it matters
.subscribe() is the primary documented API: it is the first example under Streaming in docs/src/content/docs/sdk/streaming.md, and liveQuery() is built on it. The buffer exists to hand events to a for await consumer that hasn't asked yet; for a callback consumer it is pure retention.
Rough cost: a modest 100 events/sec stream at ~200 bytes/event retains ~70 MB/hour, unbounded. It is the _buffer reference itself that pins the events, so nothing is collectable while the controller lives.
Retention is doubled for a liveQuery() or a .where()-filtered stream. FilteredStreamController extends StreamController (query-builder.ts), so the inner controller buffers every event off the wire while the outer buffers every event that matches the filter. Both layers are subject to the same never-drained path.
Fix sketch
The buffer should only accumulate when someone is actually iterating. Options, roughly in order of preference:
- Only buffer once the iterator has been requested. Set a flag in
[Symbol.asyncIterator]() and make :36 conditional on it. Preserves the current guarantee for for await consumers (including one that attaches late but before events flow) and reduces a subscribe-only stream to zero retention.
- Cap it with a documented bound and drop-oldest, which turns an unbounded leak into a bounded one but silently loses events for a slow iterator.
- Drop the buffer and have
next() only ever resolve from a waiter — simplest, but a for await consumer that awaits something between iterations would miss events, which is a real regression.
(1) is the only one that is purely a fix. Worth checking against #389, which touches this same buffering.
Acceptance
Related: #473 (subscriber-throw isolation, same loops — noticed while verifying that the throw-path buffer growth described there is the smaller of two), #389 (StreamController buffering), #470 (documents the throw-path case; does not change this).
Area: sdk — streaming. Pre-existing on
main; surfaced reviewing #470 and deliberately not fixed there.StreamControllerretains every event it has ever received unless something is consuming it through the async iterator. A long-lived.subscribe()-only stream — the pattern our own docs lead with — grows without bound until the process dies or the stream is closed.Mechanism
clients/ts/src/stream/controller.ts:31-37, in theonEventhandler:_bufferhas exactly four references in the file (:20declaration,:36push,:170-171read/shift). The only drain is the async iterator'snext(). A consumer who calls.subscribe()and never writesfor awaitnever pushes a waiter and never callsnext(), so theelsebranch is taken on every single event and nothing ever removes it. There is no cap.Callback delivery still works — the fan-out at
:28-30happens first — so the stream behaves correctly and leaks silently.Why it matters
.subscribe()is the primary documented API: it is the first example under Streaming indocs/src/content/docs/sdk/streaming.md, andliveQuery()is built on it. The buffer exists to hand events to afor awaitconsumer that hasn't asked yet; for a callback consumer it is pure retention.Rough cost: a modest 100 events/sec stream at ~200 bytes/event retains ~70 MB/hour, unbounded. It is the
_bufferreference itself that pins the events, so nothing is collectable while the controller lives.Retention is doubled for a
liveQuery()or a.where()-filtered stream.FilteredStreamControllerextendsStreamController(query-builder.ts), so the inner controller buffers every event off the wire while the outer buffers every event that matches the filter. Both layers are subject to the same never-drained path.Fix sketch
The buffer should only accumulate when someone is actually iterating. Options, roughly in order of preference:
[Symbol.asyncIterator]()and make:36conditional on it. Preserves the current guarantee forfor awaitconsumers (including one that attaches late but before events flow) and reduces a subscribe-only stream to zero retention.next()only ever resolve from a waiter — simplest, but afor awaitconsumer that awaits something between iterations would miss events, which is a real regression.(1) is the only one that is purely a fix. Worth checking against #389, which touches this same buffering.
Acceptance
.subscribe()-only stream retains no events regardless of how many arrivefor awaitconsumer still receives events that arrived between iterationsRelated: #473 (subscriber-throw isolation, same loops — noticed while verifying that the throw-path buffer growth described there is the smaller of two), #389 (StreamController buffering), #470 (documents the throw-path case; does not change this).