diff --git a/src/adapters/run-turn-queue.ts b/src/adapters/run-turn-queue.ts index 4b63e387de..6cd310fcb6 100644 --- a/src/adapters/run-turn-queue.ts +++ b/src/adapters/run-turn-queue.ts @@ -6,10 +6,29 @@ export const PREFLIGHT_HEARTBEAT_RETAIN_LIMIT = 16; /** * Coalescing threshold for adjacent text/thinking deltas buffered with no - * waiting reader (UTF-16 code units). This is a merge-size ceiling, not a - * byte-memory cap: a single oversized incoming event stays one item. + * waiting reader (UTF-16 code units). The aggregate backlog budget below is + * enforced separately, including for a single oversized incoming event. */ export const COALESCE_MAX_CHUNK_LENGTH = 64 * 1024; +/** Maximum retained string payload across one queue, measured as UTF-16 code units. */ +export const DEFAULT_MAX_BACKLOG_CODE_UNITS = 1024 * 1024; + +function retainedStringCodeUnits(value: unknown, seen = new Set()): number { + if (typeof value === "string") return value.length; + if (!value || typeof value !== "object" || seen.has(value)) return 0; + seen.add(value); + let total = 0; + for (const nested of Object.values(value)) total += retainedStringCodeUnits(nested, seen); + return total; +} + +function retainedEventStringCodeUnits(event: AdapterEvent): number { + let total = 0; + for (const [key, value] of Object.entries(event)) { + if (key !== "type") total += retainedStringCodeUnits(value); + } + return total; +} export interface AdapterEventQueue { push(event: AdapterEvent): void; @@ -64,11 +83,14 @@ export async function preflightAdapterEvents( export function createAdapterEventQueue(opts?: { maxBacklog?: number; + maxBacklogCodeUnits?: number; onBacklogExceeded?: () => void; }): AdapterEventQueue { const queued: AdapterEvent[] = []; const readers: QueueReader[] = []; const maxBacklog = opts?.maxBacklog ?? 1_024; + const maxBacklogCodeUnits = opts?.maxBacklogCodeUnits ?? DEFAULT_MAX_BACKLOG_CODE_UNITS; + let backlogCodeUnits = 0; let closed = false; // Merge an incoming delta into the buffered tail when no reader is waiting. @@ -105,6 +127,14 @@ export function createAdapterEventQueue(opts?: { reader({ done: false, value: event }); return; } + const eventCodeUnits = retainedEventStringCodeUnits(event); + if (eventCodeUnits > maxBacklogCodeUnits - backlogCodeUnits) { + opts?.onBacklogExceeded?.(); + queued.push({ type: "error", message: "consumer stalled: adapter event backlog exceeded — turn aborted" }); + close(); + return; + } + backlogCodeUnits += eventCodeUnits; if (coalesceIntoTail(event)) return; if (queued.length >= maxBacklog) { opts?.onBacklogExceeded?.(); @@ -127,6 +157,7 @@ export function createAdapterEventQueue(opts?: { while (true) { const next = queued.shift(); if (next) { + backlogCodeUnits -= retainedEventStringCodeUnits(next); yield next; continue; } diff --git a/tests/run-turn-queue.test.ts b/tests/run-turn-queue.test.ts index 57406a2e1f..9bfcb1c2b9 100644 --- a/tests/run-turn-queue.test.ts +++ b/tests/run-turn-queue.test.ts @@ -1,5 +1,5 @@ import { describe, expect, test } from "bun:test"; -import { COALESCE_MAX_CHUNK_LENGTH, createAdapterEventQueue, PREFLIGHT_HEARTBEAT_RETAIN_LIMIT, preflightAdapterEvents } from "../src/adapters/run-turn-queue"; +import { COALESCE_MAX_CHUNK_LENGTH, createAdapterEventQueue, DEFAULT_MAX_BACKLOG_CODE_UNITS, PREFLIGHT_HEARTBEAT_RETAIN_LIMIT, preflightAdapterEvents } from "../src/adapters/run-turn-queue"; import type { AdapterEvent } from "../src/types"; const text = (value: string): AdapterEvent => ({ type: "text_delta", text: value }); @@ -113,6 +113,56 @@ describe("run-turn adapter event queue", () => { expect(collected[0]).toEqual(thinking(Array.from({ length: 5_000 }, (_, i) => String(i % 10)).join(""))); }); + test("coalesced deltas cannot exceed the aggregate string backlog budget", async () => { + let backlogExceeded = 0; + const queue = createAdapterEventQueue({ + maxBacklogCodeUnits: 4, + onBacklogExceeded: () => { backlogExceeded += 1; }, + }); + + queue.push(text("ab")); + queue.push(text("cd")); + queue.push(text("e")); + + expect(backlogExceeded).toBe(1); + expect(await queue.collect()).toEqual([ + text("abcd"), + { type: "error", message: "consumer stalled: adapter event backlog exceeded — turn aborted" }, + ]); + }); + + test("an oversized individual event is rejected before it is retained", async () => { + let backlogExceeded = 0; + const queue = createAdapterEventQueue({ + onBacklogExceeded: () => { backlogExceeded += 1; }, + }); + + queue.push(text("x".repeat(DEFAULT_MAX_BACKLOG_CODE_UNITS + 1))); + + expect(backlogExceeded).toBe(1); + expect(await queue.collect()).toEqual([ + { type: "error", message: "consumer stalled: adapter event backlog exceeded — turn aborted" }, + ]); + }); + + test("consuming buffered events releases aggregate string budget", async () => { + let backlogExceeded = 0; + const queue = createAdapterEventQueue({ + maxBacklogCodeUnits: 4, + onBacklogExceeded: () => { backlogExceeded += 1; }, + }); + const iterator = queue.stream()[Symbol.asyncIterator](); + + queue.push(text("abcd")); + expect(await iterator.next()).toEqual({ done: false, value: text("abcd") }); + queue.push(thinking("wxyz")); + queue.close(); + + expect(await iterator.next()).toEqual({ done: false, value: thinking("wxyz") }); + expect(await iterator.next()).toEqual({ done: true, value: undefined }); + expect(backlogExceeded).toBe(0); + }); + test("coalescing splits past the combined-length threshold and preserves concatenation", async () => { const queue = createAdapterEventQueue(); const chunk = "x".repeat(Math.floor(COALESCE_MAX_CHUNK_LENGTH / 3) + 1);