From 4b69140f8cc5e2200a55ab0483b5da202e9b1df1 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 12 Aug 2026 21:07:38 +0000 Subject: [PATCH 1/3] feat(agent): compact reusable stream replay Co-Authored-By: matt.apperson --- .changeset/quiet-replay-streams.md | 5 + packages/agent/src/inner-loop/call-model.ts | 2 + packages/agent/src/lib/async-params.ts | 9 + packages/agent/src/lib/model-result.ts | 69 ++++++- packages/agent/src/lib/reusable-stream.ts | 137 +++++++++++--- .../agent/src/lib/tool-event-broadcaster.ts | 88 +++++++-- .../model-result-replay-compaction.test.ts | 162 +++++++++++++++++ .../unit/replay-buffer-compaction.test.ts | 171 ++++++++++++++++++ .../tests/unit/tool-event-broadcaster.test.ts | 36 ++++ 9 files changed, 627 insertions(+), 52 deletions(-) create mode 100644 .changeset/quiet-replay-streams.md create mode 100644 packages/agent/tests/unit/model-result-replay-compaction.test.ts create mode 100644 packages/agent/tests/unit/replay-buffer-compaction.test.ts diff --git a/.changeset/quiet-replay-streams.md b/.changeset/quiet-replay-streams.md new file mode 100644 index 00000000..4065926e --- /dev/null +++ b/.changeset/quiet-replay-streams.md @@ -0,0 +1,5 @@ +--- +'@openrouter/agent': patch +--- + +Reduce replay-stream memory usage with opt-in active-consumer compaction and stop provider streams at terminal response events. diff --git a/packages/agent/src/inner-loop/call-model.ts b/packages/agent/src/inner-loop/call-model.ts index 84a687d3..73d2d8ef 100644 --- a/packages/agent/src/inner-loop/call-model.ts +++ b/packages/agent/src/inner-loop/call-model.ts @@ -108,6 +108,7 @@ export function callModel< sharedContextSchema, onTurnStart, onTurnEnd, + streamReplay, allowFinalResponse, strictFinalResponse, hooks, @@ -189,6 +190,7 @@ export function callModel< sharedContextSchema, onTurnStart, onTurnEnd, + streamReplay, allowFinalResponse, strictFinalResponse, hooks: hooks !== undefined ? resolveHooks(hooks) : undefined, diff --git a/packages/agent/src/lib/async-params.ts b/packages/agent/src/lib/async-params.ts index fcbdf1bf..e49ba806 100644 --- a/packages/agent/src/lib/async-params.ts +++ b/packages/agent/src/lib/async-params.ts @@ -4,6 +4,7 @@ import type { DoomLoopOption } from './doom-loop.js'; import type { HooksManager } from './hooks-manager.js'; import type { InlineHookConfig } from './hooks-types.js'; import type { Item } from './item-types.js'; +import type { StreamReplay } from './reusable-stream.js'; import type { ContextInput } from './tool-context.js'; import type { ParsedToolCall, @@ -114,6 +115,13 @@ type BaseCallModelInput< * Receives the turn context and the completed response for that turn */ onTurnEnd?: (context: TurnContext, response: OpenResponsesResult) => void | Promise; + /** + * Controls replay history retained for stream getters. + * `full` preserves all events for delayed and sequential consumers. + * `active-consumers` compacts events after every attached consumer advances. + * @default 'full' + */ + streamReplay?: StreamReplay; /** * When the loop exits because `stopWhen` was met and the last response * still contained tool calls, execute those pending tool calls (so they @@ -332,6 +340,7 @@ export async function resolveAsyncFunctions void | Promise; /** Callback invoked at the end of each tool execution turn */ onTurnEnd?: (context: TurnContext, response: models.OpenResponsesResult) => void | Promise; + /** Replay history retained for delayed and sequential stream consumers. */ + streamReplay?: StreamReplay; /** * When the loop exits because `stopWhen` was met and the last response * still contained tool calls, make one more model request with no tools so @@ -631,6 +642,9 @@ export class ModelResult< null; private initialStreamPipeStarted = false; private initialPipePromise: Promise | null = null; + private initialResponse: models.OpenResponsesResult | null = null; + private initialResponseError: Error | null = null; + private readonly streamReplay: StreamReplay; // Context store for typed tool context (persists across turns) private contextStore: ToolContextStore | null = null; @@ -746,6 +760,7 @@ export class ModelResult< constructor(options: GetResponseOptions) { this.options = options; + this.streamReplay = options.streamReplay ?? 'full'; this.hooksManager = options.hooks; const doomLoopConfig = resolveDoomLoopOption(options.doomLoop); this.doomLoopMonitor = doomLoopConfig ? new DoomLoopMonitor(doomLoopConfig) : null; @@ -860,7 +875,7 @@ export class ModelResult< */ private ensureTurnBroadcaster(): ToolEventBroadcaster> { if (!this.turnBroadcaster) { - this.turnBroadcaster = new ToolEventBroadcaster(); + this.turnBroadcaster = new ToolEventBroadcaster(this.streamReplay); } return this.turnBroadcaster; } @@ -903,7 +918,32 @@ export class ModelResult< timestamp: Date.now(), } satisfies TurnEndEvent); })().catch((error) => { - broadcaster.complete(error instanceof Error ? error : new Error(String(error))); + const normalizedError = error instanceof Error ? error : new Error(String(error)); + this.initialResponseError = normalizedError; + broadcaster.complete(normalizedError); + }); + } + + private captureInitialStreamEvent(event: models.StreamEvents): void { + if (isResponseCompletedEvent(event) || isResponseIncompleteEvent(event)) { + this.initialResponse = event.response; + return; + } + + if (isResponseFailedEvent(event)) { + this.initialResponseError = new Error( + `Response failed: ${JSON.stringify(event.response.error)}`, + ); + } + } + + private setReusableStream(stream: ReadableStream): void { + this.initialResponse = null; + this.initialResponseError = null; + this.reusableStream = new ReusableReadableStream(stream, { + streamReplay: this.streamReplay, + onValue: (event) => this.captureInitialStreamEvent(event), + isTerminalValue: isTerminalResponseStreamEvent, }); } @@ -1068,7 +1108,10 @@ export class ModelResult< turnNumber: number, ): Promise { if (isEventStream(value)) { - const stream = new ReusableReadableStream(value); + const stream = new ReusableReadableStream(value, { + streamReplay: this.streamReplay, + isTerminalValue: isTerminalResponseStreamEvent, + }); if (this.turnBroadcaster) { return this.pipeAndConsumeStream(stream, turnNumber); } @@ -1095,6 +1138,20 @@ export class ModelResult< if (this.finalResponse) { return this.finalResponse; } + + const initialPipePromise = this.initialPipePromise; + if (initialPipePromise) { + await initialPipePromise; + } + + if (this.initialResponseError) { + throw this.initialResponseError; + } + + if (this.initialResponse) { + return this.initialResponse; + } + if (this.reusableStream) { const response = await consumeStreamForCompletion(this.reusableStream); await this.emitPendingModelCallOnce(response); @@ -5540,7 +5597,7 @@ export class ModelResult< // Handle both streaming and non-streaming responses // The API may return a non-streaming response even when stream: true is requested if (isEventStream(apiResult.value)) { - this.reusableStream = new ReusableReadableStream(apiResult.value); + this.setReusableStream(apiResult.value); } else if (this.isNonStreamingResponse(apiResult.value)) { // API returned a complete response directly - use it as the final response this.finalResponse = apiResult.value; @@ -5827,7 +5884,7 @@ export class ModelResult< // Handle both streaming and non-streaming responses if (isEventStream(apiResult.value)) { - this.reusableStream = new ReusableReadableStream(apiResult.value); + this.setReusableStream(apiResult.value); } else if (this.isNonStreamingResponse(apiResult.value)) { this.finalResponse = apiResult.value; await this.emitPendingModelCallOnce(this.finalResponse); @@ -6837,7 +6894,7 @@ export class ModelResult< throw new Error('Stream not initialized'); } - const completedResponse = await consumeStreamForCompletion(this.reusableStream); + const completedResponse = await this.getInitialResponse(); await this.emitPendingModelCallOnce(completedResponse); return extractToolCallsFromResponse(completedResponse) as ParsedToolCall[]; } diff --git a/packages/agent/src/lib/reusable-stream.ts b/packages/agent/src/lib/reusable-stream.ts index bbef31a6..339cbf17 100644 --- a/packages/agent/src/lib/reusable-stream.ts +++ b/packages/agent/src/lib/reusable-stream.ts @@ -5,19 +5,42 @@ * Key features: * - Multiple concurrent consumers with independent read positions * - New consumers can attach while streaming is active - * - Efficient memory management with automatic cleanup + * - Full replay for delayed and sequential consumers by default + * - Opt-in active-consumer replay compaction for bounded memory * - Each consumer can read at their own pace */ +export type StreamReplay = 'full' | 'active-consumers'; + +export interface ReusableReadableStreamOptions { + streamReplay?: StreamReplay; + onValue?: (value: T) => void; + isTerminalValue?: (value: T) => boolean; +} + export class ReusableReadableStream { - private buffer: T[] = []; + private buffer: (T | undefined)[] = []; + private bufferHead = 0; + // Consumer positions are absolute. Buffer index = bufferHead + position - trimOffset. + private trimOffset = 0; private consumers = new Map(); private nextConsumerId = 0; private sourceReader: ReadableStreamDefaultReader | null = null; private sourceComplete = false; private sourceError: Error | null = null; private pumpStarted = false; + private sourceCancelPromise: Promise | null = null; + private readonly streamReplay: StreamReplay; + private readonly onValue: ((value: T) => void) | undefined; + private readonly isTerminalValue: ((value: T) => boolean) | undefined; - constructor(private sourceStream: ReadableStream) {} + constructor( + private sourceStream: ReadableStream, + options: ReusableReadableStreamOptions = {}, + ) { + this.streamReplay = options.streamReplay ?? 'full'; + this.onValue = options.onValue; + this.isTerminalValue = options.isTerminalValue; + } /** * True once the source stream has been fully read into the buffer. @@ -37,9 +60,9 @@ export class ReusableReadableStream { * buffer asynchronously. */ findLastBuffered(predicate: (item: T) => boolean): T | undefined { - for (let i = this.buffer.length - 1; i >= 0; i--) { - const item = this.buffer[i]!; - if (predicate(item)) { + for (let i = this.buffer.length - 1; i >= this.bufferHead; i--) { + const item = this.buffer[i]; + if (item !== undefined && predicate(item)) { return item; } } @@ -48,12 +71,14 @@ export class ReusableReadableStream { /** * Create a new consumer that can independently iterate over the stream. - * Multiple consumers can be created and will all receive the same data. + * Full-replay consumers start at position 0. Active-consumer replay starts + * at the current trim watermark. Multiple attached consumers advance + * independently in either mode. */ createConsumer(): AsyncIterableIterator { const consumerId = this.nextConsumerId++; const state: ConsumerState = { - position: 0, + position: this.trimOffset, waitingPromise: null, cancelled: false, }; @@ -70,25 +95,24 @@ export class ReusableReadableStream { return { async next(): Promise> { const consumer = self.consumers.get(consumerId); - if (!consumer) { - return { - done: true, - value: undefined, - }; - } - - if (consumer.cancelled) { + if (!consumer || consumer.cancelled) { return { done: true, value: undefined, }; } - // If we have buffered data at this position, return it - if (consumer.position < self.buffer.length) { - const value = self.buffer[consumer.position]!; + const bufferIndex = self.bufferHead + consumer.position - self.trimOffset; + if (bufferIndex < self.buffer.length) { + const value = self.buffer[bufferIndex]; + if (value === undefined) { + return { + done: true, + value: undefined, + }; + } consumer.position++; - // Note: We don't clean up buffer to allow sequential/reusable access + self.trimConsumed(); return { done: false, value, @@ -121,7 +145,11 @@ export class ReusableReadableStream { // Immediately check if we should resolve after setting up the promise // This handles the case where data arrived or source completed // between our initial checks and promise creation - if (self.sourceComplete || self.sourceError || consumer.position < self.buffer.length) { + if ( + self.sourceComplete || + self.sourceError || + self.bufferHead + consumer.position - self.trimOffset < self.buffer.length + ) { resolve(); } }); @@ -140,6 +168,7 @@ export class ReusableReadableStream { if (consumer) { consumer.cancelled = true; self.consumers.delete(consumerId); + self.trimConsumed(); } return { done: true, @@ -152,6 +181,7 @@ export class ReusableReadableStream { if (consumer) { consumer.cancelled = true; self.consumers.delete(consumerId); + self.trimConsumed(); } throw e; }, @@ -162,6 +192,39 @@ export class ReusableReadableStream { }; } + private trimConsumed(): void { + if (this.streamReplay === 'full' || this.consumers.size === 0) { + return; + } + + let min = Number.POSITIVE_INFINITY; + for (const consumer of this.consumers.values()) { + if (consumer.position < min) { + min = consumer.position; + } + } + + const nextHead = this.bufferHead + min - this.trimOffset; + if (nextHead <= this.bufferHead) { + return; + } + + this.trimOffset = min; + if (nextHead === this.buffer.length) { + this.buffer = []; + this.bufferHead = 0; + return; + } + + this.buffer.fill(undefined, this.bufferHead, nextHead); + if (nextHead >= BUFFER_COMPACTION_MIN_HEAD && nextHead * 2 >= this.buffer.length) { + this.buffer = this.buffer.slice(nextHead); + this.bufferHead = 0; + return; + } + this.bufferHead = nextHead; + } + /** * Start pumping data from the source stream into the buffer */ @@ -172,11 +235,12 @@ export class ReusableReadableStream { this.pumpStarted = true; this.sourceReader = this.sourceStream.getReader(); + const sourceReader = this.sourceReader; // biome-ignore lint: IIFE used for fire-and-forget stream pump void (async () => { try { while (true) { - const result = await this.sourceReader!.read(); + const result = await sourceReader.read(); if (result.done) { this.sourceComplete = true; @@ -186,21 +250,41 @@ export class ReusableReadableStream { // Add to buffer this.buffer.push(result.value); + this.onValue?.(result.value); // Notify waiting consumers this.notifyAllConsumers(); + + if (this.isTerminalValue?.(result.value)) { + this.sourceComplete = true; + this.notifyAllConsumers(); + try { + await this.cancelSourceReader(sourceReader); + } catch { + // The terminal event is authoritative, cancellation is cleanup only. + } + break; + } } } catch (error) { this.sourceError = error instanceof Error ? error : new Error(String(error)); this.notifyAllConsumers(); } finally { - if (this.sourceReader) { - this.sourceReader.releaseLock(); + sourceReader.releaseLock(); + if (this.sourceReader === sourceReader) { + this.sourceReader = null; } } })(); } + private cancelSourceReader(sourceReader: ReadableStreamDefaultReader): Promise { + if (!this.sourceCancelPromise) { + this.sourceCancelPromise = sourceReader.cancel(); + } + return this.sourceCancelPromise; + } + /** * Notify all waiting consumers that new data is available */ @@ -232,8 +316,7 @@ export class ReusableReadableStream { // Cancel the source stream if (this.sourceReader) { - await this.sourceReader.cancel(); - this.sourceReader.releaseLock(); + await this.cancelSourceReader(this.sourceReader); } } } @@ -246,3 +329,5 @@ interface ConsumerState { } | null; cancelled: boolean; } + +const BUFFER_COMPACTION_MIN_HEAD = 1024; diff --git a/packages/agent/src/lib/tool-event-broadcaster.ts b/packages/agent/src/lib/tool-event-broadcaster.ts index bc4fc069..fd1aba91 100644 --- a/packages/agent/src/lib/tool-event-broadcaster.ts +++ b/packages/agent/src/lib/tool-event-broadcaster.ts @@ -1,20 +1,26 @@ +import type { StreamReplay } from './reusable-stream.js'; + /** * A push-based event broadcaster that supports multiple concurrent consumers. * Similar to ReusableReadableStream but for push-based events from tool execution. * - * Each consumer gets their own position in the buffer and receives all events - * from their join point onward. This enables real-time streaming of generator - * tool preliminary results to multiple consumers simultaneously. + * Each consumer gets their own position in the buffer. Full replay is the + * default, and active-consumer replay can compact consumed events. * * @template T - The event type being broadcast */ export class ToolEventBroadcaster { - private buffer: T[] = []; + private buffer: (T | undefined)[] = []; + private bufferHead = 0; + // Consumer positions are absolute. Buffer index = bufferHead + position - trimOffset. + private trimOffset = 0; private consumers = new Map(); private nextConsumerId = 0; private isComplete = false; private completionError: Error | null = null; + constructor(private readonly streamReplay: StreamReplay = 'full') {} + /** * Push a new event to all consumers. * Events are buffered so late-joining consumers can catch up. @@ -46,20 +52,21 @@ export class ToolEventBroadcaster { */ private cleanup(): void { // Only cleanup if complete and all consumers are done - if (this.isComplete && this.consumers.size === 0) { + if (this.streamReplay === 'active-consumers' && this.isComplete && this.consumers.size === 0) { this.buffer = []; + this.bufferHead = 0; } } /** * Create a new consumer that can independently iterate over events. - * Consumers can join at any time and will receive events from position 0. - * Multiple consumers can be created and will all receive the same events. + * Full-replay consumers start at position 0. Active-consumer replay starts + * at the current trim watermark. */ createConsumer(): AsyncIterableIterator { const consumerId = this.nextConsumerId++; const state: ConsumerState = { - position: 0, + position: this.trimOffset, waitingPromise: null, cancelled: false, }; @@ -71,24 +78,24 @@ export class ToolEventBroadcaster { return { async next(): Promise> { const consumer = self.consumers.get(consumerId); - if (!consumer) { - return { - done: true, - value: undefined, - }; - } - - if (consumer.cancelled) { + if (!consumer || consumer.cancelled) { return { done: true, value: undefined, }; } - // Return buffered event if available - if (consumer.position < self.buffer.length) { - const value = self.buffer[consumer.position]!; + const bufferIndex = self.bufferHead + consumer.position - self.trimOffset; + if (bufferIndex < self.buffer.length) { + const value = self.buffer[bufferIndex]; + if (value === undefined) { + return { + done: true, + value: undefined, + }; + } consumer.position++; + self.trimConsumed(); return { done: false, value, @@ -116,7 +123,11 @@ export class ToolEventBroadcaster { }; // Immediately check if we should resolve after setting up promise - if (self.isComplete || self.completionError || consumer.position < self.buffer.length) { + if ( + self.isComplete || + self.completionError || + self.bufferHead + consumer.position - self.trimOffset < self.buffer.length + ) { resolve(); } }); @@ -133,6 +144,7 @@ export class ToolEventBroadcaster { if (consumer) { consumer.cancelled = true; self.consumers.delete(consumerId); + self.trimConsumed(); self.cleanup(); } return { @@ -146,6 +158,7 @@ export class ToolEventBroadcaster { if (consumer) { consumer.cancelled = true; self.consumers.delete(consumerId); + self.trimConsumed(); self.cleanup(); } throw e; @@ -157,6 +170,39 @@ export class ToolEventBroadcaster { }; } + private trimConsumed(): void { + if (this.streamReplay === 'full' || this.consumers.size === 0) { + return; + } + + let min = Number.POSITIVE_INFINITY; + for (const consumer of this.consumers.values()) { + if (consumer.position < min) { + min = consumer.position; + } + } + + const nextHead = this.bufferHead + min - this.trimOffset; + if (nextHead <= this.bufferHead) { + return; + } + + this.trimOffset = min; + if (nextHead === this.buffer.length) { + this.buffer = []; + this.bufferHead = 0; + return; + } + + this.buffer.fill(undefined, this.bufferHead, nextHead); + if (nextHead >= BUFFER_COMPACTION_MIN_HEAD && nextHead * 2 >= this.buffer.length) { + this.buffer = this.buffer.slice(nextHead); + this.bufferHead = 0; + return; + } + this.bufferHead = nextHead; + } + /** * Notify all waiting consumers that new data is available or stream completed */ @@ -182,3 +228,5 @@ interface ConsumerState { } | null; cancelled: boolean; } + +const BUFFER_COMPACTION_MIN_HEAD = 1024; diff --git a/packages/agent/tests/unit/model-result-replay-compaction.test.ts b/packages/agent/tests/unit/model-result-replay-compaction.test.ts new file mode 100644 index 00000000..5b81aef0 --- /dev/null +++ b/packages/agent/tests/unit/model-result-replay-compaction.test.ts @@ -0,0 +1,162 @@ +import type { OpenRouterCore } from '@openrouter/sdk/core'; +import type * as models from '@openrouter/sdk/models'; +import { describe, expect, it, vi } from 'vitest'; +import { callModel } from '../../src/inner-loop/call-model.js'; + +const mockBetaResponsesSend = vi.hoisted(() => vi.fn()); + +vi.mock('@openrouter/sdk/funcs/betaResponsesSend', () => ({ + betaResponsesSend: mockBetaResponsesSend, +})); + +function response(id: string, output: models.OutputItemsUnion[] = []): models.OpenResponsesResult { + return { + id, + object: 'response', + createdAt: 0, + model: 'test-model', + status: 'completed', + completedAt: 0, + output, + error: null, + incompleteDetails: null, + temperature: null, + topP: null, + presencePenalty: null, + frequencyPenalty: null, + metadata: null, + instructions: null, + tools: [], + toolChoice: 'auto', + parallelToolCalls: false, + } as models.OpenResponsesResult; +} + +const message: models.OutputMessage = { + id: 'message', + type: 'message', + role: 'assistant', + status: 'completed', + content: [ + { + type: 'output_text', + text: 'hello', + annotations: [], + }, + ], +}; + +function completedStream(result: models.OpenResponsesResult): ReadableStream { + return new ReadableStream({ + start(controller) { + controller.enqueue({ + type: 'response.completed', + response: result, + sequenceNumber: 0, + } as models.StreamEvents); + controller.close(); + }, + }); +} + +function failedStream(): ReadableStream { + return new ReadableStream({ + start(controller) { + controller.enqueue({ + type: 'response.failed', + response: response('failed'), + sequenceNumber: 0, + } as models.StreamEvents); + controller.close(); + }, + }); +} + +function client(): OpenRouterCore { + return {} as OpenRouterCore; +} + +describe('ModelResult replay compaction', () => { + it('keeps the default full replay for sequential consumers', async () => { + const result = response('full', [ + message, + ]); + mockBetaResponsesSend.mockResolvedValue({ + ok: true, + value: completedStream(result), + }); + const modelResult = callModel(client(), { + model: 'test-model', + input: 'hello', + }); + + await expect(modelResult.getTextStream().next()).resolves.toEqual({ + done: true, + value: undefined, + }); + await expect(modelResult.getResponse()).resolves.toEqual(result); + }); + + it('returns the terminal response after active-consumer replay trims history', async () => { + const result = response('active', [ + message, + ]); + mockBetaResponsesSend.mockResolvedValue({ + ok: true, + value: completedStream(result), + }); + const modelResult = callModel(client(), { + model: 'test-model', + input: 'hello', + streamReplay: 'active-consumers', + }); + + await expect(modelResult.getFullResponsesStream().next()).resolves.toMatchObject({ + value: { + type: 'response.completed', + }, + }); + await expect(modelResult.getResponse()).resolves.toEqual(result); + }); + + it('strips streamReplay from the resolved request', async () => { + const result = response('request', [ + message, + ]); + const testClient = client(); + mockBetaResponsesSend.mockResolvedValue({ + ok: true, + value: completedStream(result), + }); + const modelResult = callModel(testClient, { + model: 'test-model', + input: 'hello', + streamReplay: 'active-consumers', + }); + + await modelResult.getResponse(); + expect(mockBetaResponsesSend).toHaveBeenCalledWith( + testClient, + expect.objectContaining({ + responsesRequest: expect.not.objectContaining({ + streamReplay: expect.anything(), + }), + }), + expect.anything(), + ); + }); + + it('does not replace the provider terminal result when cleanup fails', async () => { + mockBetaResponsesSend.mockResolvedValue({ + ok: true, + value: failedStream(), + }); + const modelResult = callModel(client(), { + model: 'test-model', + input: 'hello', + streamReplay: 'active-consumers', + }); + + await expect(modelResult.getResponse()).rejects.toThrow(); + }); +}); diff --git a/packages/agent/tests/unit/replay-buffer-compaction.test.ts b/packages/agent/tests/unit/replay-buffer-compaction.test.ts new file mode 100644 index 00000000..9ad8713a --- /dev/null +++ b/packages/agent/tests/unit/replay-buffer-compaction.test.ts @@ -0,0 +1,171 @@ +import { describe, expect, it } from 'vitest'; +import { ReusableReadableStream } from '../../src/lib/reusable-stream.js'; + +function source(values: number[]): ReadableStream { + return new ReadableStream({ + start(controller) { + for (const value of values) { + controller.enqueue(value); + } + controller.close(); + }, + }); +} + +describe('ReusableReadableStream replay policy', () => { + it('replays the complete history to sequential and post-completion consumers by default', async () => { + const stream = new ReusableReadableStream( + source([ + 1, + 2, + 3, + ]), + ); + + expect(await Array.fromAsync(stream.createConsumer())).toEqual([ + 1, + 2, + 3, + ]); + expect(await Array.fromAsync(stream.createConsumer())).toEqual([ + 1, + 2, + 3, + ]); + }); + + it('starts new active-consumer consumers at the current watermark', async () => { + const stream = new ReusableReadableStream( + source([ + 1, + 2, + 3, + ]), + { + streamReplay: 'active-consumers', + }, + ); + const first = stream.createConsumer(); + + expect(await first.next()).toEqual({ + done: false, + value: 1, + }); + const second = stream.createConsumer(); + expect(await Array.fromAsync(first)).toEqual([ + 2, + 3, + ]); + expect(await Array.fromAsync(second)).toEqual([ + 2, + 3, + ]); + }); + + it('continues to replay active-consumer events through repeated compaction', async () => { + const values = Array.from( + { + length: 2500, + }, + (_, index) => index, + ); + const stream = new ReusableReadableStream(source(values), { + streamReplay: 'active-consumers', + }); + const consumer = stream.createConsumer(); + + expect(await Array.fromAsync(consumer)).toEqual(values); + const lateConsumer = stream.createConsumer(); + expect(await Array.fromAsync(lateConsumer)).toEqual([]); + }); + + it('stops at a terminal value and cancels the source once', async () => { + let cancellationCount = 0; + let releaseCount = 0; + const stream = new ReusableReadableStream( + new ReadableStream({ + pull(controller) { + controller.enqueue(1); + controller.enqueue(2); + }, + cancel() { + cancellationCount++; + }, + }), + { + isTerminalValue: (value) => value === 1, + }, + ); + const consumer = stream.createConsumer(); + const originalReleaseLock = ReadableStreamDefaultReader.prototype.releaseLock; + ReadableStreamDefaultReader.prototype.releaseLock = function releaseLock() { + releaseCount++; + originalReleaseLock.call(this); + }; + + try { + expect(await Array.fromAsync(consumer)).toEqual([ + 1, + ]); + } finally { + ReadableStreamDefaultReader.prototype.releaseLock = originalReleaseLock; + } + + expect(cancellationCount).toBe(1); + expect(releaseCount).toBe(1); + }); + + it('retains the terminal value when source cancellation fails', async () => { + const stream = new ReusableReadableStream( + new ReadableStream({ + pull(controller) { + controller.enqueue(1); + }, + cancel() { + return Promise.reject(new Error('cleanup failed')); + }, + }), + { + isTerminalValue: (value) => value === 1, + }, + ); + + await expect(Array.fromAsync(stream.createConsumer())).resolves.toEqual([ + 1, + ]); + }); + + it('treats a failed terminal value as the final buffered event', async () => { + const stream = new ReusableReadableStream( + source([ + 1, + 2, + ]), + { + isTerminalValue: (value) => value === 2, + }, + ); + + await expect(Array.fromAsync(stream.createConsumer())).resolves.toEqual([ + 1, + 2, + ]); + }); + + it('treats an incomplete terminal value as the final buffered event', async () => { + const stream = new ReusableReadableStream( + source([ + 1, + 2, + ]), + { + isTerminalValue: (value) => value === 2, + }, + ); + + await expect(Array.fromAsync(stream.createConsumer())).resolves.toEqual([ + 1, + 2, + ]); + }); +}); diff --git a/packages/agent/tests/unit/tool-event-broadcaster.test.ts b/packages/agent/tests/unit/tool-event-broadcaster.test.ts index 290f13f6..2803e984 100644 --- a/packages/agent/tests/unit/tool-event-broadcaster.test.ts +++ b/packages/agent/tests/unit/tool-event-broadcaster.test.ts @@ -2,6 +2,42 @@ import { describe, expect, it } from 'vitest'; import { ToolEventBroadcaster } from '../../src/lib/tool-event-broadcaster.js'; describe('ToolEventBroadcaster', () => { + it('retains full replay after completion by default', async () => { + const broadcaster = new ToolEventBroadcaster(); + const first = broadcaster.createConsumer(); + + broadcaster.push(1); + broadcaster.complete(); + expect(await Array.fromAsync(first)).toEqual([ + 1, + ]); + + const second = broadcaster.createConsumer(); + expect(await Array.fromAsync(second)).toEqual([ + 1, + ]); + }); + + it('compacts active-consumer history at the watermark', async () => { + const broadcaster = new ToolEventBroadcaster('active-consumers'); + const first = broadcaster.createConsumer(); + + broadcaster.push(1); + broadcaster.push(2); + broadcaster.complete(); + expect(await first.next()).toEqual({ + done: false, + value: 1, + }); + const second = broadcaster.createConsumer(); + expect(await Array.fromAsync(first)).toEqual([ + 2, + ]); + expect(await Array.fromAsync(second)).toEqual([ + 2, + ]); + }); + describe('single consumer', () => { it('should deliver events to a single consumer', async () => { const broadcaster = new ToolEventBroadcaster(); From 56a6a142d38006324ee97e28031fbfb5a2341940 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 12 Aug 2026 21:12:18 +0000 Subject: [PATCH 2/3] fix(agent): preserve replay terminal invariants Co-Authored-By: matt.apperson --- .changeset/quiet-replay-streams.md | 15 +++- packages/agent/README.md | 12 +++ packages/agent/src/index.ts | 1 + packages/agent/src/lib/model-result.ts | 28 +++++- packages/agent/src/lib/reusable-stream.ts | 7 +- .../agent/src/lib/tool-event-broadcaster.ts | 7 +- .../model-result-replay-compaction.test.ts | 88 ++++++++++++++++++- 7 files changed, 144 insertions(+), 14 deletions(-) diff --git a/.changeset/quiet-replay-streams.md b/.changeset/quiet-replay-streams.md index 4065926e..d2f89382 100644 --- a/.changeset/quiet-replay-streams.md +++ b/.changeset/quiet-replay-streams.md @@ -1,5 +1,16 @@ --- -'@openrouter/agent': patch +'@openrouter/agent': minor --- -Reduce replay-stream memory usage with opt-in active-consumer compaction and stop provider streams at terminal response events. +Add opt-in replay compaction and terminal response-event handling for streamed model calls. + +```ts +import { callModel } from '@openrouter/agent'; + +const result = callModel(client, { + model: 'openai/gpt-4o', + input: 'Summarize this document.', + // Retain only the history needed by currently attached consumers. + streamReplay: 'active-consumers', +}); +``` diff --git a/packages/agent/README.md b/packages/agent/README.md index 810bcc56..a44a773c 100644 --- a/packages/agent/README.md +++ b/packages/agent/README.md @@ -87,6 +87,18 @@ console.log(response.usage); // { inputTokens, outputTokens, cost, ... } const usage = await result.getUsage(); console.log(usage); // { modelCalls, inputTokens, outputTokens, totalTokens, cachedTokens, reasoningTokens, cost? } +For long-lived consumers, `streamReplay: 'active-consumers'` releases buffered +events after every attached consumer advances past them. The default +`streamReplay: 'full'` retains complete replay history: + +```typescript +const result = callModel(client, { + model, + input, + streamReplay: 'active-consumers', +}); +``` + // Stream text deltas for await (const delta of result.getTextStream()) { process.stdout.write(delta); diff --git a/packages/agent/src/index.ts b/packages/agent/src/index.ts index 1849332f..c4dfa8e4 100644 --- a/packages/agent/src/index.ts +++ b/packages/agent/src/index.ts @@ -205,6 +205,7 @@ export { buildNextTurnParamsContext, executeNextTurnParamsFunctions, } from './lib/next-turn-params.js'; +export type { StreamReplay } from './lib/reusable-stream.js'; // Stop condition helpers export { finishReasonIs, diff --git a/packages/agent/src/lib/model-result.ts b/packages/agent/src/lib/model-result.ts index 5c13b296..59f974bd 100644 --- a/packages/agent/src/lib/model-result.ts +++ b/packages/agent/src/lib/model-result.ts @@ -1149,6 +1149,7 @@ export class ModelResult< } if (this.initialResponse) { + await this.emitPendingModelCallOnce(this.initialResponse); return this.initialResponse; } @@ -1160,6 +1161,26 @@ export class ModelResult< throw new Error('Neither stream nor response initialized'); } + private extractCachedCompletion(): models.OpenResponsesResult { + if (this.initialResponseError) { + throw this.initialResponseError; + } + if (this.initialResponse) { + return this.initialResponse; + } + if (!this.reusableStream) { + throw new Error('Stream not initialized'); + } + return extractCompletionFromBuffer(this.reusableStream); + } + + private tryExtractCachedCompletion(): models.OpenResponsesResult | undefined { + if (this.initialResponse) { + return this.initialResponse; + } + return this.reusableStream ? tryExtractCompletionFromBuffer(this.reusableStream) : undefined; + } + /** * Save response output to state. * Appends the response output to the message history and records the response ID. @@ -2538,7 +2559,7 @@ export class ModelResult< // Sync backward scan of the retained buffer — not a consumer // replay, which would cost one microtask hop per buffered event // on every hook-less streaming teardown. - await this.emitPendingModelCallOnce(extractCompletionFromBuffer(this.reusableStream)); + await this.emitPendingModelCallOnce(this.extractCachedCompletion()); } else if (this.reusableStream) { // Consumers stop at the terminal event (streamTerminationEvents), // usually before the pump reads the source close that flips @@ -2548,7 +2569,7 @@ export class ModelResult< // dropping the parked telemetry. Stays silent (no emit, no // throw) when nothing terminal was buffered — e.g. an errored // mid-flight stream, where no materialized response exists. - const buffered = tryExtractCompletionFromBuffer(this.reusableStream); + const buffered = this.tryExtractCachedCompletion(); if (buffered) { await this.emitPendingModelCallOnce(buffered); } @@ -4986,6 +5007,7 @@ export class ModelResult< strictFinalResponse: _sfr, hooks: _h, doomLoop: _dl, + streamReplay: _sr, signal: _sig, toolTimeoutMs: _ttm, toolConcurrency: _tc, @@ -7200,7 +7222,7 @@ export class ModelResult< // reusable stream is a passive observation: it buffers events // without executing tools or mutating conversation state, so the // resume generation is counted without advancing the loop. - await this.emitPendingModelCallOnce(await consumeStreamForCompletion(this.reusableStream)); + await this.emitPendingModelCallOnce(await this.getInitialResponse()); } } catch (error) { // Intentionally swallowed — see the "never rejects" note above. The diff --git a/packages/agent/src/lib/reusable-stream.ts b/packages/agent/src/lib/reusable-stream.ts index 339cbf17..f33f5872 100644 --- a/packages/agent/src/lib/reusable-stream.ts +++ b/packages/agent/src/lib/reusable-stream.ts @@ -106,10 +106,9 @@ export class ReusableReadableStream { if (bufferIndex < self.buffer.length) { const value = self.buffer[bufferIndex]; if (value === undefined) { - return { - done: true, - value: undefined, - }; + throw new Error( + 'ReusableReadableStream buffer invariant violated: consumed slot was cleared', + ); } consumer.position++; self.trimConsumed(); diff --git a/packages/agent/src/lib/tool-event-broadcaster.ts b/packages/agent/src/lib/tool-event-broadcaster.ts index fd1aba91..6111e090 100644 --- a/packages/agent/src/lib/tool-event-broadcaster.ts +++ b/packages/agent/src/lib/tool-event-broadcaster.ts @@ -89,10 +89,9 @@ export class ToolEventBroadcaster { if (bufferIndex < self.buffer.length) { const value = self.buffer[bufferIndex]; if (value === undefined) { - return { - done: true, - value: undefined, - }; + throw new Error( + 'ToolEventBroadcaster buffer invariant violated: consumed slot was cleared', + ); } consumer.position++; self.trimConsumed(); diff --git a/packages/agent/tests/unit/model-result-replay-compaction.test.ts b/packages/agent/tests/unit/model-result-replay-compaction.test.ts index 5b81aef0..27e7ec6f 100644 --- a/packages/agent/tests/unit/model-result-replay-compaction.test.ts +++ b/packages/agent/tests/unit/model-result-replay-compaction.test.ts @@ -1,7 +1,9 @@ import type { OpenRouterCore } from '@openrouter/sdk/core'; import type * as models from '@openrouter/sdk/models'; -import { describe, expect, it, vi } from 'vitest'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { z } from 'zod/v4'; import { callModel } from '../../src/inner-loop/call-model.js'; +import { HooksManager } from '../../src/lib/hooks-manager.js'; const mockBetaResponsesSend = vi.hoisted(() => vi.fn()); @@ -77,6 +79,10 @@ function client(): OpenRouterCore { } describe('ModelResult replay compaction', () => { + beforeEach(() => { + mockBetaResponsesSend.mockReset(); + }); + it('keeps the default full replay for sequential consumers', async () => { const result = response('full', [ message, @@ -101,6 +107,13 @@ describe('ModelResult replay compaction', () => { const result = response('active', [ message, ]); + const hooks = new HooksManager(); + const postModelCalls: unknown[] = []; + hooks.on('PostModelCall', { + handler: (payload) => { + postModelCalls.push(payload); + }, + }); mockBetaResponsesSend.mockResolvedValue({ ok: true, value: completedStream(result), @@ -109,6 +122,7 @@ describe('ModelResult replay compaction', () => { model: 'test-model', input: 'hello', streamReplay: 'active-consumers', + hooks, }); await expect(modelResult.getFullResponsesStream().next()).resolves.toMatchObject({ @@ -117,6 +131,78 @@ describe('ModelResult replay compaction', () => { }, }); await expect(modelResult.getResponse()).resolves.toEqual(result); + expect(postModelCalls).toHaveLength(1); + }); + + it('omits streamReplay from the follow-up request body', async () => { + const initial = response('tool', [ + { + type: 'function_call', + id: 'call-item', + callId: 'call-id', + name: 'echo', + arguments: '{}', + status: 'completed', + }, + ]); + const followUp = response('follow-up', [ + message, + ]); + mockBetaResponsesSend + .mockResolvedValueOnce({ + ok: true, + value: initial, + }) + .mockResolvedValueOnce({ + ok: true, + value: followUp, + }); + + const modelResult = callModel(client(), { + model: 'test-model', + input: 'hello', + streamReplay: 'active-consumers', + tools: [ + { + type: 'function', + function: { + name: 'echo', + description: 'Echo input.', + inputSchema: z.object({}), + outputSchema: z.string(), + execute: async () => 'ok', + }, + }, + ], + }); + + await expect(modelResult.getResponse()).resolves.toEqual(followUp); + expect(mockBetaResponsesSend).toHaveBeenCalledTimes(2); + expect(mockBetaResponsesSend.mock.calls[1]?.[1]?.responsesRequest).not.toHaveProperty( + 'streamReplay', + ); + }); + + it('uses the cached terminal response when active replay released the buffer', async () => { + const result = response('usage', [ + message, + ]); + mockBetaResponsesSend.mockResolvedValue({ + ok: true, + value: completedStream(result), + }); + const modelResult = callModel(client(), { + model: 'test-model', + input: 'hello', + streamReplay: 'active-consumers', + }); + + for await (const _event of modelResult.getFullResponsesStream()) { + // Drain the consumer so active replay can release its terminal event. + } + await expect(modelResult.getUsage()).resolves.toMatchObject({ + modelCalls: 1, + }); }); it('strips streamReplay from the resolved request', async () => { From 1e368fd6611ec5bfca108daae9da2e887e864958 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 12 Aug 2026 21:13:19 +0000 Subject: [PATCH 3/3] docs(agent): move replay note out of the consumption code block Co-Authored-By: matt.apperson --- packages/agent/README.md | 30 ++++++++++++++++++------------ 1 file changed, 18 insertions(+), 12 deletions(-) diff --git a/packages/agent/README.md b/packages/agent/README.md index a44a773c..a9047b76 100644 --- a/packages/agent/README.md +++ b/packages/agent/README.md @@ -87,18 +87,6 @@ console.log(response.usage); // { inputTokens, outputTokens, cost, ... } const usage = await result.getUsage(); console.log(usage); // { modelCalls, inputTokens, outputTokens, totalTokens, cachedTokens, reasoningTokens, cost? } -For long-lived consumers, `streamReplay: 'active-consumers'` releases buffered -events after every attached consumer advances past them. The default -`streamReplay: 'full'` retains complete replay history: - -```typescript -const result = callModel(client, { - model, - input, - streamReplay: 'active-consumers', -}); -``` - // Stream text deltas for await (const delta of result.getTextStream()) { process.stdout.write(delta); @@ -139,6 +127,24 @@ What each stream emits: | `getItemsStream()` | all output items (messages, function calls, …) — output items **only**, no usage/response metadata | | `getFullResponsesStream()` | every response event, including `tool.result` / `tool.call_output` execution events, and each round's `response.completed` (with that round's usage block) | +#### Replay history for stream consumers + +Every stream getter above can start from event zero, so by default a result +retains its full event history for the lifetime of the call. Long generator-tool +streams make that history expensive in a constrained runtime. When all consumers +attach before draining, `streamReplay: 'active-consumers'` releases buffered +events once every attached consumer has advanced past them: + +```typescript +const result = callModel(client, { + model, + input, + // Default 'full' retains complete replay history for delayed and + // sequential consumers; 'active-consumers' trades that for bounded memory. + streamReplay: 'active-consumers', +}); +``` + #### Usage across a multi-round tool loop `getResponse()` resolves to the **final** round's response, so in a