From 3532362498f19374cc222a5aeabe3409482c164e Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Thu, 9 Jul 2026 09:27:26 -0600 Subject: [PATCH 1/2] fix(agents-runtime): avoid stack overflow on large live batches --- .changeset/large-live-batches.md | 5 +++++ packages/agents-runtime/src/process-wake.ts | 12 +++++++++--- 2 files changed, 14 insertions(+), 3 deletions(-) create mode 100644 .changeset/large-live-batches.md diff --git a/.changeset/large-live-batches.md b/.changeset/large-live-batches.md new file mode 100644 index 0000000000..5a5128f468 --- /dev/null +++ b/.changeset/large-live-batches.md @@ -0,0 +1,5 @@ +--- +"@electric-ax/agents-runtime": patch +--- + +Avoid stack overflows when processing large live StreamDB batches. diff --git a/packages/agents-runtime/src/process-wake.ts b/packages/agents-runtime/src/process-wake.ts index 472cba39b5..95ddc746aa 100644 --- a/packages/agents-runtime/src/process-wake.ts +++ b/packages/agents-runtime/src/process-wake.ts @@ -84,6 +84,12 @@ interface RawClaimCallbackResponse extends Omit { const DEFAULT_IDLE_TIMEOUT = 20_000 const DEFAULT_HEARTBEAT_INTERVAL = 10_000 + +function appendAll(target: Array, items: Array): void { + for (const item of items) { + target.push(item) + } +} type EntityStreamOptions = NonNullable< Parameters[3] > @@ -1075,7 +1081,7 @@ export async function processWake( // ctx.events so no child result is silently acked without being visible. pendingLiveBatches.shift() batches.push(batch) - deltaEvents.push(...changeEvents) + appendAll(deltaEvents, changeEvents) if (freshKind === `inbox`) { selectedKind = `inbox` @@ -1130,7 +1136,7 @@ export async function processWake( if (!preloaded) { const changeEvents = toChangeEvents(batch) if (changeEvents.length > 0) { - catchUpEvents.push(...changeEvents) + appendAll(catchUpEvents, changeEvents) } lastCatchUpOffset = batch.offset return @@ -1145,7 +1151,7 @@ export async function processWake( handleLatestSignalEvents(changeEvents) - catchUpEvents.push(...changeEvents) + appendAll(catchUpEvents, changeEvents) if ( resolveCurrentWakeReady !== null && From 41483590ad25cdaae26f08dc5c60bda9094aca66 Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Thu, 9 Jul 2026 09:27:26 -0600 Subject: [PATCH 2/2] fix(agents-runtime): avoid stack overflow on large live batches --- .changeset/large-live-batches.md | 6 ++++ packages/agents-runtime/src/process-wake.ts | 12 ++++++-- packages/agents-server/src/stream-client.ts | 6 +++- .../agents-server/test/stream-client.test.ts | 29 +++++++++++++++++++ 4 files changed, 49 insertions(+), 4 deletions(-) create mode 100644 .changeset/large-live-batches.md diff --git a/.changeset/large-live-batches.md b/.changeset/large-live-batches.md new file mode 100644 index 0000000000..adf35cf2a7 --- /dev/null +++ b/.changeset/large-live-batches.md @@ -0,0 +1,6 @@ +--- +"@electric-ax/agents-runtime": patch +"@electric-ax/agents-server": patch +--- + +Avoid stack overflows when processing large live StreamDB batches and reading large JSON streams during server rehydration. diff --git a/packages/agents-runtime/src/process-wake.ts b/packages/agents-runtime/src/process-wake.ts index 472cba39b5..95ddc746aa 100644 --- a/packages/agents-runtime/src/process-wake.ts +++ b/packages/agents-runtime/src/process-wake.ts @@ -84,6 +84,12 @@ interface RawClaimCallbackResponse extends Omit { const DEFAULT_IDLE_TIMEOUT = 20_000 const DEFAULT_HEARTBEAT_INTERVAL = 10_000 + +function appendAll(target: Array, items: Array): void { + for (const item of items) { + target.push(item) + } +} type EntityStreamOptions = NonNullable< Parameters[3] > @@ -1075,7 +1081,7 @@ export async function processWake( // ctx.events so no child result is silently acked without being visible. pendingLiveBatches.shift() batches.push(batch) - deltaEvents.push(...changeEvents) + appendAll(deltaEvents, changeEvents) if (freshKind === `inbox`) { selectedKind = `inbox` @@ -1130,7 +1136,7 @@ export async function processWake( if (!preloaded) { const changeEvents = toChangeEvents(batch) if (changeEvents.length > 0) { - catchUpEvents.push(...changeEvents) + appendAll(catchUpEvents, changeEvents) } lastCatchUpOffset = batch.offset return @@ -1145,7 +1151,7 @@ export async function processWake( handleLatestSignalEvents(changeEvents) - catchUpEvents.push(...changeEvents) + appendAll(catchUpEvents, changeEvents) if ( resolveCurrentWakeReady !== null && diff --git a/packages/agents-server/src/stream-client.ts b/packages/agents-server/src/stream-client.ts index 96e92de279..292cf75c26 100644 --- a/packages/agents-server/src/stream-client.ts +++ b/packages/agents-server/src/stream-client.ts @@ -417,7 +417,11 @@ export class StreamClient { offset: fromOffset ?? `-1`, live: false, }) - return await response.json() + const items: Array = [] + for await (const item of response.jsonStream()) { + items.push(item) + } + return items }) } diff --git a/packages/agents-server/test/stream-client.test.ts b/packages/agents-server/test/stream-client.test.ts index 6daa4aa173..46188ec15a 100644 --- a/packages/agents-server/test/stream-client.test.ts +++ b/packages/agents-server/test/stream-client.test.ts @@ -7,6 +7,7 @@ const { flushMock, detachMock, headMock, + streamMock, MockFetchError, MockDurableStreamError, } = vi.hoisted(() => { @@ -29,6 +30,7 @@ const { flushMock: vi.fn().mockResolvedValue(undefined), detachMock: vi.fn().mockResolvedValue(undefined), headMock: vi.fn(), + streamMock: vi.fn(), MockFetchError: HoistedFetchError, MockDurableStreamError: HoistedDurableStreamError, } @@ -38,6 +40,8 @@ vi.mock(`@durable-streams/client`, () => ({ DurableStream: class { constructor(_opts: { url: string; contentType?: string }) {} + stream = streamMock + static create = vi.fn() static delete = vi.fn() static head = headMock @@ -58,6 +62,31 @@ describe(`StreamClient`, () => { flushMock.mockClear() detachMock.mockClear() headMock.mockReset() + streamMock.mockReset() + }) + + it(`readJson reads large batches without using StreamResponse.json`, async () => { + const largeBatch = Array.from({ length: 70_000 }, (_, i) => ({ i })) + const jsonMock = vi.fn(() => { + throw new RangeError(`Maximum call stack size exceeded`) + }) + streamMock.mockResolvedValueOnce({ + json: jsonMock, + async *jsonStream() { + for (const item of largeBatch) { + yield item + } + }, + }) + + const client = new StreamClient(`http://127.0.0.1:4545`) + + await expect(client.readJson(`/large-json`)).resolves.toEqual(largeBatch) + expect(jsonMock).not.toHaveBeenCalled() + expect(streamMock).toHaveBeenCalledWith({ + offset: `-1`, + live: false, + }) }) it(`appendIdempotent uses IdempotentProducer append/flush/detach`, async () => {