Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 33 additions & 2 deletions src/adapters/run-turn-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<object>()): 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;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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;
Comment on lines +137 to 138

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Track only strings retained after delta coalescing

When Kiro emits adjacent phased text_delta events, backlogCodeUnits counts the repeated phase string on every push before coalescing, even though the merged queue item retains that field only once. Dequeueing consequently subtracts only one copy and leaves phantom backlog units, so a long stream of commentary or final_answer deltas can eventually abort despite the actual buffered payload remaining below the limit; with maxBacklogCodeUnits: 12, for example, two one-character commentary deltas should merge into a 12-unit event but the second is rejected. Update the accounting from the post-coalescing retained value, or remove the duplicated phase charge when replacing the tail.

AGENTS.md reference: src/AGENTS.md:L19-L19

Useful? React with 👍 / 👎.

if (queued.length >= maxBacklog) {
opts?.onBacklogExceeded?.();
Expand All @@ -127,6 +157,7 @@ export function createAdapterEventQueue(opts?: {
while (true) {
const next = queued.shift();
if (next) {
backlogCodeUnits -= retainedEventStringCodeUnits(next);
yield next;
continue;
}
Expand Down
52 changes: 51 additions & 1 deletion tests/run-turn-queue.test.ts
Original file line number Diff line number Diff line change
@@ -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 });
Expand Down Expand Up @@ -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);
Expand Down
Loading