From d5f7e8e21455c8fc5fe27ae9cda560f91f57edc8 Mon Sep 17 00:00:00 2001 From: ox-alpha Date: Sat, 22 Aug 2026 09:18:57 -0500 Subject: [PATCH] fix(agent): release replay-buffer backlog on terminal consumer/cancel paths MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three memory/lifecycle gaps in ReusableReadableStream and ToolEventBroadcaster, found while profiling fusion OOM (exceededMemory) in openrouter-web: 1. Zero-consumer backlog wipe (active-consumers mode): trimConsumed early-returned when consumers.size === 0, so when the LAST consumer departed, the entire remaining backlog stayed pinned until the stream object itself was GC'd — O(total streamed events) per run, stacking across concurrent runs. Now dropped wholesale once at least one consumer has existed and none remain; guarded on nextConsumerId so a stream nobody has consumed yet keeps its catch-up history. 2. Pending next() hang: return()/throw() deleted the consumer without settling a registered waitingPromise — a concurrently pending next() hung until the next pump tick (or forever on push-driven broadcasters). return() now resolves the waiter; throw() rejects with a normalized error. 3. RRS cancel() now drops the buffered backlog (swept once after source reader cancellation, since the pump can land one in-flight chunk between the sync section and cancellation). 'full' replay mode is untouched: late consumers still replay everything. Tests pin the wipe, slowest-consumer retention, waiter wake, and cancel sweep for both classes. --- packages/agent/src/lib/reusable-stream.ts | 50 ++++- .../agent/src/lib/tool-event-broadcaster.ts | 39 +++- .../agent/tests/unit/reusable-stream.test.ts | 212 ++++++++++++++++++ .../tests/unit/tool-event-broadcaster.test.ts | 85 +++++++ 4 files changed, 384 insertions(+), 2 deletions(-) create mode 100644 packages/agent/tests/unit/reusable-stream.test.ts diff --git a/packages/agent/src/lib/reusable-stream.ts b/packages/agent/src/lib/reusable-stream.ts index f33f5872..fdd4ba41 100644 --- a/packages/agent/src/lib/reusable-stream.ts +++ b/packages/agent/src/lib/reusable-stream.ts @@ -166,6 +166,13 @@ export class ReusableReadableStream { const consumer = self.consumers.get(consumerId); if (consumer) { consumer.cancelled = true; + // Wake a pending next() so it does not hang forever on a consumer + // that has just been removed; the re-run next() observes the + // removal and reports done. + if (consumer.waitingPromise) { + consumer.waitingPromise.resolve(); + consumer.waitingPromise = null; + } self.consumers.delete(consumerId); self.trimConsumed(); } @@ -179,6 +186,12 @@ export class ReusableReadableStream { const consumer = self.consumers.get(consumerId); if (consumer) { consumer.cancelled = true; + // Reject a pending next() with the same error so it does not hang + // on a consumer that has just been removed. + if (consumer.waitingPromise) { + consumer.waitingPromise.reject(e instanceof Error ? e : new Error(String(e))); + consumer.waitingPromise = null; + } self.consumers.delete(consumerId); self.trimConsumed(); } @@ -192,7 +205,11 @@ export class ReusableReadableStream { } private trimConsumed(): void { - if (this.streamReplay === 'full' || this.consumers.size === 0) { + if (this.streamReplay === 'full') { + return; + } + if (this.consumers.size === 0) { + this.dropUnreadBacklog(); return; } @@ -224,6 +241,28 @@ export class ReusableReadableStream { this.bufferHead = nextHead; } + /** + * Drop the retained backlog once no consumer can ever read it again. + * Before the first consumer joins, the backlog IS the catch-up history late + * joiners replay — retain it. Once at least one consumer has existed and + * none remain, any future consumer starts at the watermark anyway, so the + * retained bytes would stay pinned until GC for nothing. `findLastBuffered` + * callers that must observe terminal events after teardown should capture + * them via `onValue` instead of relying on the backlog surviving. + */ + private dropUnreadBacklog(): void { + if (this.nextConsumerId === 0) { + return; + } + const dropped = this.buffer.length - this.bufferHead; + if (dropped === 0) { + return; + } + this.trimOffset += dropped; + this.buffer = []; + this.bufferHead = 0; + } + /** * Start pumping data from the source stream into the buffer */ @@ -313,10 +352,19 @@ export class ReusableReadableStream { } this.consumers.clear(); + // Cancellation is terminal: every consumer is gone, so nobody can read + // the buffered backlog anymore. Drop it instead of pinning it until GC. + this.dropUnreadBacklog(); // Cancel the source stream if (this.sourceReader) { await this.cancelSourceReader(this.sourceReader); } + /* + * The pump may have landed one in-flight chunk between the synchronous + * section above and reader cancellation — sweep once more so the buffer + * is empty regardless of microtask interleaving. + */ + this.dropUnreadBacklog(); } } diff --git a/packages/agent/src/lib/tool-event-broadcaster.ts b/packages/agent/src/lib/tool-event-broadcaster.ts index 6111e090..7f437e09 100644 --- a/packages/agent/src/lib/tool-event-broadcaster.ts +++ b/packages/agent/src/lib/tool-event-broadcaster.ts @@ -142,6 +142,13 @@ export class ToolEventBroadcaster { const consumer = self.consumers.get(consumerId); if (consumer) { consumer.cancelled = true; + // Wake a pending next() so it does not hang forever on a consumer + // that has just been removed; the re-run next() observes the + // removal and reports done. + if (consumer.waitingPromise) { + consumer.waitingPromise.resolve(); + consumer.waitingPromise = null; + } self.consumers.delete(consumerId); self.trimConsumed(); self.cleanup(); @@ -156,6 +163,12 @@ export class ToolEventBroadcaster { const consumer = self.consumers.get(consumerId); if (consumer) { consumer.cancelled = true; + // Reject a pending next() with the same error so it does not hang + // on a consumer that has just been removed. + if (consumer.waitingPromise) { + consumer.waitingPromise.reject(e instanceof Error ? e : new Error(String(e))); + consumer.waitingPromise = null; + } self.consumers.delete(consumerId); self.trimConsumed(); self.cleanup(); @@ -170,7 +183,11 @@ export class ToolEventBroadcaster { } private trimConsumed(): void { - if (this.streamReplay === 'full' || this.consumers.size === 0) { + if (this.streamReplay === 'full') { + return; + } + if (this.consumers.size === 0) { + this.dropUnreadBacklog(); return; } @@ -202,6 +219,26 @@ export class ToolEventBroadcaster { this.bufferHead = nextHead; } + /** + * Drop the retained backlog once no consumer can ever read it again. + * Before the first consumer joins, the buffer IS the catch-up history late + * joiners replay — retain it. Once at least one consumer has existed and + * none remain, any future consumer starts at the watermark anyway, so the + * retained events would stay pinned until GC for nothing. + */ + private dropUnreadBacklog(): void { + if (this.nextConsumerId === 0) { + return; + } + const dropped = this.buffer.length - this.bufferHead; + if (dropped === 0) { + return; + } + this.trimOffset += dropped; + this.buffer = []; + this.bufferHead = 0; + } + /** * Notify all waiting consumers that new data is available or stream completed */ diff --git a/packages/agent/tests/unit/reusable-stream.test.ts b/packages/agent/tests/unit/reusable-stream.test.ts new file mode 100644 index 00000000..bd27dd40 --- /dev/null +++ b/packages/agent/tests/unit/reusable-stream.test.ts @@ -0,0 +1,212 @@ +import { describe, expect, it } from 'vitest'; +import { ReusableReadableStream } from '../../src/lib/reusable-stream.js'; + +/** + * Contract tests for the replay-buffer memory/lifecycle edges: + * - active-consumers mode must not pin O(stream) backlog once every consumer + * has departed (the fusion exceededMemory retainer) + * - full-replay mode (default) keeps replaying everything + * - a pending next() must never hang after return()/throw()/cancel() + * + * A hung next() is caught by vitest's per-test timeout rather than a wall- + * clock guard: the awaited signal is always the promise under test. + */ +function streamOf(values: readonly T[]): ReadableStream { + return new ReadableStream({ + start(controller): void { + for (const value of values) { + controller.enqueue(value); + } + controller.close(); + }, + }); +} + +/** Source whose producer is manually controlled, for mid-stream assertions. */ +function controlledStream(): { + stream: ReadableStream; + push: (value: T) => void; + close: () => void; +} { + let controller!: ReadableStreamDefaultController; + const stream = new ReadableStream({ + start(c): void { + controller = c; + }, + }); + return { + stream, + push: (value) => controller.enqueue(value), + close: () => controller.close(), + }; +} + +describe('ReusableReadableStream', () => { + it('retains full replay for late consumers by default', async () => { + const stream = new ReusableReadableStream( + streamOf([ + 1, + 2, + ]), + ); + const first = stream.createConsumer(); + expect(await Array.fromAsync(first)).toEqual([ + 1, + 2, + ]); + + const second = stream.createConsumer(); + expect(await Array.fromAsync(second)).toEqual([ + 1, + 2, + ]); + }); + + it('active-consumers: backlog follows the slowest remaining consumer', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream, { + streamReplay: 'active-consumers', + }); + const fast = stream.createConsumer(); + const slow = stream.createConsumer(); + + source.push(1); + source.push(2); + expect(await fast.next()).toEqual({ + done: false, + value: 1, + }); + expect(await fast.next()).toEqual({ + done: false, + value: 2, + }); + await fast.return(); + + // The slow consumer still reads from position 0 — trimming follows the + // slowest attached consumer, never the fastest. + expect(await slow.next()).toEqual({ + done: false, + value: 1, + }); + expect(await slow.next()).toEqual({ + done: false, + value: 2, + }); + source.close(); + expect((await slow.next()).done).toBe(true); + }); + + it('active-consumers: backlog is dropped once every consumer departs', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream, { + streamReplay: 'active-consumers', + }); + const only = stream.createConsumer(); + source.push(1); + source.push(2); + expect(await only.next()).toEqual({ + done: false, + value: 1, + }); + expect(await only.next()).toEqual({ + done: false, + value: 2, + }); + await only.return(); + + // Nobody is attached anymore: the retained backlog is dropped instead of + // staying pinned until GC. + expect(stream.findLastBuffered(() => true)).toBeUndefined(); + + // A late consumer starts at the watermark instead of replaying history. + source.push(3); + const late = stream.createConsumer(); + expect(await late.next()).toEqual({ + done: false, + value: 3, + }); + source.close(); + }); + + it('active-consumers: a fully caught-up consumer leaves nothing for later joiners', async () => { + // Trimming follows consumption, not attachment time: once every consumer + // has read everything, the backlog is gone. Replay-after-detach is the + // documented trade-off of active-consumers mode; default ('full') mode + // keeps it. + const stream = new ReusableReadableStream( + streamOf([ + 1, + 2, + ]), + { + streamReplay: 'active-consumers', + }, + ); + const only = stream.createConsumer(); + expect(await Array.fromAsync(only)).toEqual([ + 1, + 2, + ]); + + const late = stream.createConsumer(); + expect(await Array.fromAsync(late)).toEqual([]); + }); + + it('return() wakes a pending next() instead of hanging', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream); + const consumer = stream.createConsumer(); + source.push(1); + expect(await consumer.next()).toEqual({ + done: false, + value: 1, + }); + + const pending = consumer.next(); + await consumer.return(); + await expect(pending).resolves.toEqual({ + done: true, + value: undefined, + }); + source.close(); + }); + + it('throw() rejects a pending next() with a normalized error', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream); + const consumer = stream.createConsumer(); + source.push(1); + await consumer.next(); + + const pending = consumer.next(); + await expect( + consumer.throw(new Error('boom')).catch((error: Error) => error), + ).resolves.toBeInstanceOf(Error); + await expect(pending).rejects.toThrow('boom'); + source.close(); + }); + + it('cancel() settles waiters and leaves no buffered backlog', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream); + const consumer = stream.createConsumer(); + source.push(1); + source.push(2); + + /* + * Whether the parked next() wakes with a buffered value or as done + * depends on microtask interleaving with the pump; the contract is only + * that it settles and that no backlog survives. + */ + const pending = consumer.next(); + await Promise.all([ + pending.catch(() => {}), + stream.cancel(), + ]); + + expect(stream.findLastBuffered(() => true)).toBeUndefined(); + const fresh = stream.createConsumer(); + expect((await fresh.next()).done).toBe(true); + // No source.close(): cancel() already terminated the source stream. + }); +}); diff --git a/packages/agent/tests/unit/tool-event-broadcaster.test.ts b/packages/agent/tests/unit/tool-event-broadcaster.test.ts index 2803e984..d76d8ffc 100644 --- a/packages/agent/tests/unit/tool-event-broadcaster.test.ts +++ b/packages/agent/tests/unit/tool-event-broadcaster.test.ts @@ -398,4 +398,89 @@ describe('ToolEventBroadcaster', () => { }); }); }); + describe('lifecycle compaction (active-consumers)', () => { + it('drops backlog once every consumer departs, late joiner starts at watermark', async () => { + const broadcaster = new ToolEventBroadcaster('active-consumers'); + const only = broadcaster.createConsumer(); + broadcaster.push(1); + broadcaster.push(2); + expect(await only.next()).toEqual({ + done: false, + value: 1, + }); + expect(await only.next()).toEqual({ + done: false, + value: 2, + }); + await only.return(); + + // Nobody is attached: retained history is dropped, so a fresh consumer + // starts at the watermark (retention would replay [1, 2] first). + broadcaster.push(3); + broadcaster.complete(); + const late = broadcaster.createConsumer(); + expect(await Array.fromAsync(late)).toEqual([ + 3, + ]); + }); + + it('backlog survives while at least one consumer remains', async () => { + const broadcaster = new ToolEventBroadcaster('active-consumers'); + const fast = broadcaster.createConsumer(); + const slow = broadcaster.createConsumer(); + broadcaster.push(1); + broadcaster.push(2); + expect(await fast.next()).toEqual({ + done: false, + value: 1, + }); + expect(await fast.next()).toEqual({ + done: false, + value: 2, + }); + await fast.return(); + + // The slow consumer still reads from position 0. + expect(await slow.next()).toEqual({ + done: false, + value: 1, + }); + expect(await slow.next()).toEqual({ + done: false, + value: 2, + }); + broadcaster.complete(); + expect((await slow.next()).done).toBe(true); + }); + + it('return() wakes a pending next() instead of hanging', async () => { + const broadcaster = new ToolEventBroadcaster('active-consumers'); + const consumer = broadcaster.createConsumer(); + broadcaster.push(1); + expect(await consumer.next()).toEqual({ + done: false, + value: 1, + }); + + const pending = consumer.next(); + await consumer.return(); + await expect(pending).resolves.toEqual({ + done: true, + value: undefined, + }); + }); + + it('throw() rejects a pending next() with a normalized error', async () => { + const broadcaster = new ToolEventBroadcaster('active-consumers'); + const consumer = broadcaster.createConsumer(); + broadcaster.push(1); + await consumer.next(); + + const pending = consumer.next(); + await expect( + consumer.throw(new Error('boom')).catch((error: Error) => error), + ).resolves.toBeInstanceOf(Error); + await expect(pending).rejects.toThrow('boom'); + }); + }); });