From da9d51b869749d47ef4747978cb03115bf51ffa0 Mon Sep 17 00:00:00 2001 From: Ingwannu Date: Fri, 25 Sep 2026 09:55:11 +0000 Subject: [PATCH 1/3] perf: make fragmented stream work linear --- scripts/test-layout/layout.json | 3 + src/adapters/coding-agent/protocol.ts | 48 +++-- .../serialized-tool-call-content.ts | 70 +++++-- src/usage/log.ts | 7 +- structure/dashboard-and-usage.md | 4 + .../ADR-0102-incremental-stream-accounting.md | 23 +++ .../ADR-5724-serialized-tool-call-content.md | 2 +- structure/providers-and-adapters.md | 6 + ...-chat-serialized-tool-call-scaling.test.ts | 177 ++++++++++++++++++ tests/fixtures/test-layout-expected.json | 3 + .../coding-agent-json-lines-scaling.test.ts | 90 +++++++++ .../usage/usage-snapshot-digest-reuse.test.ts | 112 +++++++++++ 12 files changed, 511 insertions(+), 34 deletions(-) create mode 100644 structure/decisions/ADR-0102-incremental-stream-accounting.md create mode 100644 tests/adapters/openai/openai-chat-serialized-tool-call-scaling.test.ts create mode 100644 tests/providers/coding-agent-json-lines-scaling.test.ts create mode 100644 tests/usage/usage-snapshot-digest-reuse.test.ts diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index b436f658bb5..0788f2a939f 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -168,6 +168,9 @@ } }, "explicit": { + "openai-chat-serialized-tool-call-scaling.test.ts": "adapters/openai", + "coding-agent-json-lines-scaling.test.ts": "providers", + "usage-snapshot-digest-reuse.test.ts": "usage", "release-desktop-scripts.test.ts": "ci-workflows", "installed-gate-drivers.test.ts": "ci-workflows", "gui-desktop-sidecar-script.test.ts": "gui", diff --git a/src/adapters/coding-agent/protocol.ts b/src/adapters/coding-agent/protocol.ts index 032bfba4b22..cccbb05db58 100644 --- a/src/adapters/coding-agent/protocol.ts +++ b/src/adapters/coding-agent/protocol.ts @@ -80,14 +80,11 @@ export async function* readJsonLines( const maxLineBytes = limits.maxLineBytes ?? MAX_STREAM_LINE_BYTES; const maxTotalBytes = limits.maxTotalBytes ?? MAX_STREAM_TOTAL_BYTES; const decoder = new TextDecoder(); - const encoder = new TextEncoder(); - let buffer = ""; + let parts: string[] = []; + let lineBytes = 0; let totalBytes = 0; const flushLine = function* (line: string): Generator { - if (encoder.encode(line).byteLength > maxLineBytes) { - throw new CodingAgentStreamLimitError("Coding-agent stream line exceeded the byte ceiling"); - } const trimmed = line.trim(); if (!trimmed) return; // Blank lines and whitespace-only lines are ignored as padding. let parsed: unknown; @@ -108,26 +105,39 @@ export async function* readJsonLines( yield parsed as StreamMessage; }; + // Decode continuously for split UTF-8/BOM semantics, but search and measure only new + // decoded segments. Joining once per frame avoids repeatedly flattening a growing rope. + const consume = function* (text: string): Generator { + let start = 0; + while (start < text.length) { + const newline = text.indexOf("\n", start); + const end = newline < 0 ? text.length : newline; + const part = text.slice(start, end); + lineBytes += Buffer.byteLength(part); + if (lineBytes > maxLineBytes) { + throw new CodingAgentStreamLimitError("Coding-agent stream line exceeded the byte ceiling"); + } + if (part) parts.push(part); + if (newline < 0) break; + const line = parts.join(""); + parts = []; + lineBytes = 0; + yield* flushLine(line); + start = newline + 1; + } + }; + for await (const chunk of chunks) { totalBytes += chunk.byteLength; if (totalBytes > maxTotalBytes) { throw new CodingAgentStreamLimitError("Coding-agent stream exceeded the total byte ceiling"); } - buffer += decoder.decode(chunk, { stream: true }); - let newline = buffer.indexOf("\n"); - while (newline >= 0) { - const line = buffer.slice(0, newline); - buffer = buffer.slice(newline + 1); - yield* flushLine(line); - newline = buffer.indexOf("\n"); - } - if (encoder.encode(buffer).byteLength > maxLineBytes) { - throw new CodingAgentStreamLimitError("Coding-agent stream line exceeded the byte ceiling"); - } + yield* consume(decoder.decode(chunk, { stream: true })); } - // Flush the decoder's trailing bytes and any final line without a newline terminator. - buffer += decoder.decode(); - if (buffer.trim()) yield* flushLine(buffer); + // Flush incomplete UTF-8 through the same decoded-byte accounting before parsing EOF. + yield* consume(decoder.decode()); + const finalLine = parts.join(""); + if (finalLine.trim()) yield* flushLine(finalLine); } function asRecord(value: unknown): Record | undefined { diff --git a/src/adapters/openai-chat/serialized-tool-call-content.ts b/src/adapters/openai-chat/serialized-tool-call-content.ts index d887a16905c..ff60738068a 100644 --- a/src/adapters/openai-chat/serialized-tool-call-content.ts +++ b/src/adapters/openai-chat/serialized-tool-call-content.ts @@ -192,6 +192,29 @@ function contextAfter(text: string, context: TextContext): TextContext { return splitAtPossibleSerializedToolCall(text.replaceAll(OPEN_TAG, ""), context, true).context; } +/** + * A deferred header/fence can itself be long. While its unresolved portion only grows, a + * chunk-local check suffices; the first resolving character goes back through the full splitter. + * Each unbounded phase is scanned once on entry and once on exit, preserving its Markdown rules. + */ +function deferredContinuationBreak(text: string, context: TextContext): RegExp | undefined { + if (text.startsWith(OPEN_TAG)) { + if (/^\s*$/.test(text)) return /\S/; + if (/^\s*\r\n]*$/.test(text)) return /[>\r\n]/; + } + if (!context.inlineTicks && context.lineStart) { + const fence = /^ {0,3}(`{3,}|~{3,})([^\n]*)$/.exec(text); + if (fence) { + if (!context.fence) return /\n/; + // Closing-fence ticks may grow only before its optional whitespace suffix. + if (!fence[2]) return fence[1]![0] === "`" ? /[^`]/ : /[^~]/; + return /[^ \t\r]/; + } + } + if (/^`+$/.test(text)) return /[^`]/; + return undefined; +} + /** * Holds possible duplicate text within the shared translator budget until the dispatch outcome is * known. While a block candidate is open, any other event (reasoning) is queued at its position in @@ -202,6 +225,11 @@ export class SerializedToolCallContentBuffer { private text = ""; private bytes = 0; private hasOpenTag = false; + private continuationBreak: RegExp | undefined; + private delimiterSuffix = ""; + private sawCloser = false; + private openAfterCloser = false; + private trailingProseChars = 0; private context: TextContext = { fence: null, lineStart: true }; private queued: { offset: number; event: AdapterEvent }[] = []; @@ -217,6 +245,12 @@ export class SerializedToolCallContentBuffer { this.text = next; this.bytes = nextBytes; this.hasOpenTag = hasOpenTag; + this.continuationBreak = undefined; + this.delimiterSuffix = ""; + this.sawCloser = false; + this.openAfterCloser = false; + this.trailingProseChars = 0; + if (hasOpenTag) this.observeDelimiters(next); } catch (error) { reservation.release(); throw error; @@ -225,21 +259,43 @@ export class SerializedToolCallContentBuffer { /** Charges only the appended bytes, so holding an open block never needs twice its retained size. */ private append(delta: string): void { + // Bun counts each isolated surrogate as two bytes, so split pairs sum to the same four-byte + // charge as the joined scalar. Measuring only the delta therefore preserves exact accounting. const deltaBytes = Buffer.byteLength(delta); this.budget.reserveTransient(deltaBytes, { kind: "live_transient" }).commitRetained(); + if (this.hasOpenTag) this.observeDelimiters(delta); this.text += delta; this.bytes += deltaBytes; } + /** Scan new text plus a fixed overlap, never the retained body (which may be a large rope). */ + private observeDelimiters(delta: string): void { + const scan = this.delimiterSuffix + delta; + const closer = scan.lastIndexOf(CLOSE_TAG); + const afterCloser = closer + CLOSE_TAG.length; + if (closer >= 0) { + this.sawCloser = true; + this.openAfterCloser = false; + this.trailingProseChars = 0; + } + if (scan.lastIndexOf(OPEN_TAG) >= (closer >= 0 ? afterCloser : 0)) { + this.openAfterCloser = true; + } + const tail = closer >= 0 ? scan.slice(afterCloser) : delta; + if (this.sawCloser) this.trailingProseChars += tail.replace(/\s+/g, "").length; + this.delimiterSuffix = scan.slice(-(Math.max(OPEN_TAG.length, CLOSE_TAG.length) - 1)); + } + /** Returns immediately safe text and retains only the suffix that still needs reconciliation. */ ingest(delta: string): string { - if (this.hasOpenTag) { + if (this.hasOpenTag || (this.continuationBreak && !this.continuationBreak.test(delta))) { this.append(delta); return ""; } const split = splitAtPossibleSerializedToolCall(this.text + delta, this.context); this.replace(split.defer, split.hasOpenTag); this.context = split.context; + if (!split.hasOpenTag) this.continuationBreak = deferredContinuationBreak(split.defer, split.context); return split.emit; } @@ -262,7 +318,7 @@ export class SerializedToolCallContentBuffer { } const text = this.ingest(delta); // Checked after ingest too: one delta can open a block and already carry more than a bound. - if (this.hasOpenTag && (this.bytes > MAX_HELD_BYTES || proseAfterClosedBlock(this.text) > MAX_TRAILING_CHARS)) { + if (this.hasOpenTag && (this.bytes > MAX_HELD_BYTES || (this.sawCloser && !this.openAfterCloser && this.trailingProseChars > MAX_TRAILING_CHARS))) { return [...textEvents(text), ...this.drain([])]; } return textEvents(text); @@ -334,6 +390,7 @@ export class SerializedToolCallContentBuffer { this.text = ""; this.bytes = 0; this.hasOpenTag = false; + this.continuationBreak = undefined; this.queued = []; } } @@ -595,12 +652,3 @@ export function reconcileSerializedToolCallEvents( function textEvents(text: string): AdapterEvent[] { return text.length > 0 ? [{ type: "text_delta", text }] : []; } - -/** Non-whitespace characters after the last closed block, or 0 while a later block is still open. */ -function proseAfterClosedBlock(text: string): number { - const closer = text.lastIndexOf(CLOSE_TAG); - if (closer < 0) return 0; - const tail = text.slice(closer + CLOSE_TAG.length); - if (tail.includes(OPEN_TAG)) return 0; - return tail.replace(/\s+/g, "").length; -} diff --git a/src/usage/log.ts b/src/usage/log.ts index bc9a5a11e28..d382f202b42 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -1659,9 +1659,10 @@ async function readUsageEntriesIncrementally( // would make a byte-truncated read claim rows were dropped when none were. entriesTruncated: entriesDropped > 0, entriesDropped, - // The digest must describe exactly the region the returned rows came from, which - // is the post-trim window, not the pre-trim one. - prefixDigest: usageRegionDigest(fd, rowsBeginAtBytes, size) ?? "", + // Reuse this read's verified digest only for identical bounds. Growth or trimming + // needs a new digest of the returned region; metadata alone never proves reuse. + prefixDigest: rowsBeginAtBytes === retained.rowsBeginAtBytes && size === retained.coveredThroughBytes + ? covered : usageRegionDigest(fd, rowsBeginAtBytes, size) ?? "", entryLengths: lengths, trailingSkippedBytes: appendedTrailingSkipped, rowsBeginAtBytes, diff --git a/structure/dashboard-and-usage.md b/structure/dashboard-and-usage.md index cea1aa709e5..5e9c5161110 100644 --- a/structure/dashboard-and-usage.md +++ b/structure/dashboard-and-usage.md @@ -346,6 +346,10 @@ batches, so memory stays bounded and unrelated management requests remain servic large existing log. The first read is proportional to ledger size; steady-state refresh work is proportional to newly appended bytes. The Dashboard polls its 30-day usage summary independently once per minute, so usage work cannot delay health/provider/settings state or run every five seconds. +An unchanged retained snapshot verifies its byte region once and reuses that digest only when the +returned region has identical start and end bounds. Appends or window trimming still hash the new +returned region, so digest reuse cannot weaken same-inode rewrite detection. See +[ADR-0102](decisions/ADR-0102-incremental-stream-accounting.md). An oversized row is skipped inside the scanner bound without shortening identities. Accumulators keep normal rows plus `usageIncomplete` / `usageIncompleteReason: "oversized_rows"` on caches and rollups; append ORs the flag and a rebuild recalculates it. Invalid-row counts are not sticky, and absence of the flag is not completeness. GUI caches warn on Usage, Dashboard, provider and key views; CLI warns in human output only; most-used order save refuses an incomplete snapshot. Quota surfaces stay separate. Legacy truncation fields keep their meaning; read/mutation failures still fail closed. diff --git a/structure/decisions/ADR-0102-incremental-stream-accounting.md b/structure/decisions/ADR-0102-incremental-stream-accounting.md new file mode 100644 index 00000000000..00118712fa7 --- /dev/null +++ b/structure/decisions/ADR-0102-incremental-stream-accounting.md @@ -0,0 +1,23 @@ +# ADR-0102 β€” Incremental stream accounting + +- Contract owners: [Chat compatibility](../providers/chat-compat.md#serialized-tool-call-content), + [Providers and adapters](../providers-and-adapters.md), and + [Dashboard and usage](../dashboard-and-usage.md) + +## Decision Log + +- Purpose and intent: Keep proxy CPU and temporary allocation proportional to newly received bytes, + even when providers or child agents fragment a large logical frame into tiny chunks. +- Existing implementation and constraints: Serialized tool-call reconciliation must preserve exact + text and Markdown context; coding-agent JSONL must preserve streaming UTF-8, BOM, newline, and + byte-ceiling behavior; usage snapshots must still detect same-inode rewrites before reuse. +- Alternatives considered: Lower existing size limits; periodically flatten and rescan retained + text; use wall-clock performance assertions; or retain incremental parser/digest state. +- Selected approach: Scan only each new segment plus fixed delimiter overlap, maintain decoded JSONL + line-byte accounting while joining once per completed frame, and reuse a verified usage digest + only when both digest bounds are identical to the returned region. +- Why this approach: Smaller limits reduce functionality without removing quadratic work. Incremental + state preserves the existing wire contracts and gives deterministic work-count regression tests. +- Benefits, tradeoffs, and impact: CPU and cumulative allocation become linear for fragmented + streams and unchanged usage polling avoids a duplicate synchronous hash. The parsers carry small + additional state, and any change in usage-region bounds still requires a fresh digest. diff --git a/structure/decisions/ADR-5724-serialized-tool-call-content.md b/structure/decisions/ADR-5724-serialized-tool-call-content.md index 34bb5bf1079..923b74900d1 100644 --- a/structure/decisions/ADR-5724-serialized-tool-call-content.md +++ b/structure/decisions/ADR-5724-serialized-tool-call-content.md @@ -10,4 +10,4 @@ - Alternatives considered: Keep two regular expressions and try the closed form first (backtracks quadratically on a long unterminated body, and rejecting a closed match on any inner `` hides a body that merely contains that string); drop anything shaped like a tool call (discards real text); restore calls from markup when no structured call exists (a new behaviour this route has no evidence for). - Choice: Read each block by delimiter scan. The first `` preceded by `` closes the block; if none appears before the next block header at the start of a line (`` followed by `` does. A stray `` before the close is markup, and one leading newline in the body is template layout. - Why: It accepts the same grammar on both MiMo routes, keeps a body that contains literal tool-call tags (even a full header) intact, and costs linear time. The agreement rule from ADR-5548 is unchanged, so no new text can disappear without a matching structured call. -- Consequences: The two echo shapes are removed when they duplicate a structured call, streamed and buffered. Markup with no structured call is still shown and still runs nothing. +- Consequences: The two echo shapes are removed when they duplicate a structured call, streamed and buffered. Markup with no structured call is still shown and still runs nothing. Streaming retention scans each new delta with fixed delimiter overlap and carries trailing-prose state; it never searches the whole retained block per delta. See [ADR-0102](ADR-0102-incremental-stream-accounting.md). diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index 74cdcd198ca..391d30acada 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -21,6 +21,12 @@ system prompt through its documented scoped `QODER_APPEND_SYSTEM_PROMPT` or `QODERCN_APPEND_SYSTEM_PROMPT` child environment, never through command-line arguments or inherited vendor variables. +Coding-agent stdout is framed as bounded JSONL directly from decoded stream segments. The framer +tracks the current line's UTF-8 byte count incrementally, searches each decoded segment once, and +joins only when a newline or EOF completes the frame. This preserves split UTF-8, BOM, CRLF, +blank-line, line-limit, and total-limit behavior without re-encoding the growing partial frame on +every child stdout chunk. See [ADR-0102](decisions/ADR-0102-incremental-stream-accounting.md). + Kimi Coding's Chat, API-key, and optional Responses presets consume the same model seeds in `src/providers/registry/model-seeds.ts`, including the native `k3-256k` ID. The Responses preset shares the `kimi` OAuth account and Coding endpoint, keeps Chat as the featured default, and diff --git a/tests/adapters/openai/openai-chat-serialized-tool-call-scaling.test.ts b/tests/adapters/openai/openai-chat-serialized-tool-call-scaling.test.ts new file mode 100644 index 00000000000..c47a6b20b85 --- /dev/null +++ b/tests/adapters/openai/openai-chat-serialized-tool-call-scaling.test.ts @@ -0,0 +1,177 @@ +import { describe, expect, spyOn, test } from "bun:test"; +import { SerializedToolCallContentBuffer, splitAtPossibleSerializedToolCall } from "../../../src/adapters/openai-chat/serialized-tool-call-content"; +import type { AdapterEvent } from "../../../src/types"; +import { createTestTranslatorBudget } from "../../helpers/translator-budget"; + +const OPEN = ""; +const CLOSE = ""; +const BLOCK = OPEN + "ok" + CLOSE; + +function scanWork(count: number, closed: boolean): number { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + buffer.ingestStreaming(closed ? BLOCK : OPEN); + let work = 0; + const original = String.prototype.lastIndexOf; + const scan = spyOn(String.prototype, "lastIndexOf").mockImplementation(function (this: string, ...args) { + work += this.length; + return original.apply(this, args); + }); + try { + // Whitespace after a close cannot trigger the prose limit; the entire prefix stays held. + for (let i = 0; i < count; i++) buffer.ingestStreaming(closed ? " \t\n" : "xyz"); + } finally { + scan.mockRestore(); + buffer.dispose(); + } + return work; +} + +describe("serialized tool content incremental delimiter work", () => { + for (const closed of [false, true]) { + test(`search work is linear for ${closed ? "closed" : "open"} blocks in tiny deltas`, () => { + const small = scanWork(2048, closed); + const large = scanWork(4096, closed); + expect(small).toBeGreaterThan(0); + // Count actual search input, not elapsed time or an implementation-maintained counter. + expect(large).toBeLessThanOrEqual(2 * small + 24); + expect(large).toBeLessThanOrEqual(2 * 4096 * (3 + 11)); + }); + } + + test("every closing delimiter split keeps the inclusive prose limit and event order", () => { + for (let split = 0; split <= CLOSE.length; split++) { + const budget = createTestTranslatorBudget(); + const buffer = new SerializedToolCallContentBuffer(budget); + const reasoning: AdapterEvent = { type: "reasoning_raw_delta", text: "thinking" }; + expect(buffer.ingestStreaming(OPEN + "ok")).toEqual([]); + expect(buffer.hold(reasoning)).toEqual([{ type: "heartbeat" }]); + expect(buffer.ingestStreaming(CLOSE.slice(0, split))).toEqual([]); + expect(buffer.ingestStreaming(CLOSE.slice(split))).toEqual([]); + expect(buffer.ingestStreaming("x".repeat(8192) + " \t\n\u00a0")).toEqual([]); + expect(buffer.ingestStreaming("πŸ˜€")).toEqual([ + { type: "text_delta", text: OPEN + "ok" }, reasoning, + { type: "text_delta", text: CLOSE + "x".repeat(8192) + " \t\n\u00a0πŸ˜€" }, + ]); + expect(budget.snapshot().currentBytes).toBe(0); + buffer.dispose(); + } + }); + + test("split later openers suspend prose release and a later closer resets its count", () => { + for (let split = 0; split <= "".length; split++) { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + expect(buffer.ingestStreaming(BLOCK + "a".repeat(8000))).toEqual([]); + expect(buffer.ingestStreaming("".slice(0, split))).toEqual([]); + expect(buffer.ingestStreaming("".slice(split) + "b".repeat(9000))).toEqual([]); + expect(buffer.ingestStreaming(CLOSE + "x".repeat(8192))).toEqual([]); + const held = buffer.current(); + expect(buffer.ingestStreaming("!")).toEqual([{ type: "text_delta", text: held + "!" }]); + // The next held block has fresh delimiter state, including after a bounded release. + buffer.ingestStreaming("\n"); + expect(buffer.ingestStreaming(OPEN + "y".repeat(9000))).toEqual([]); + buffer.dispose(); + } + }); + + test("opening headers split at every position retain exact reconciliation and queued events", () => { + for (let split = 0; split <= OPEN.length; split++) { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + expect(buffer.ingestStreaming(OPEN.slice(0, split))).toEqual([]); + expect(buffer.ingestStreaming(OPEN.slice(split) + "ok")).toEqual([]); + const reasoning: AdapterEvent = { type: "reasoning_raw_delta", text: "thinking" }; + buffer.hold(reasoning); + expect(buffer.ingestStreaming(CLOSE)).toEqual([]); + expect(buffer.drain([{ names: new Set(["exec"]), argumentsText: '{"input":"ok"}' }])).toEqual([reasoning]); + expect(buffer.ingestStreaming("\n" + BLOCK)).toEqual([{ type: "text_delta", text: "\n" }]); + expect(buffer.drain([])).toEqual([{ type: "text_delta", text: BLOCK }]); + buffer.dispose(); + } + }); +}); + +describe("deferred header and Markdown phases", () => { + const cases = [ + { prefix: "", delta: " \n", end: "ok" + CLOSE, held: true }, + { prefix: "ok" + CLOSE, held: true }, + { prefix: "```", delta: "language", end: "\ncode\n```\n", held: false }, + { prefix: "```\ncode\n```", delta: " ", end: "\n", held: false }, + { prefix: "```\ncode\n```", delta: "`", end: "\n", held: false }, + { prefix: "text `", delta: "`", end: "done", held: false }, + ]; + for (const fixture of cases) { + test(`new chunks alone are examined during the deferred phase ${JSON.stringify(fixture.prefix)}`, () => { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + const initial = buffer.ingest(fixture.prefix); + let examined = 0; + const original = RegExp.prototype.test; + const scan = spyOn(RegExp.prototype, "test").mockImplementation(function (this: RegExp, value) { + examined += value.length; + return original.call(this, value); + }); + try { + for (let i = 0; i < 2048; i++) buffer.ingestStreaming(fixture.delta); + } finally { scan.mockRestore(); } + expect(examined).toBeLessThanOrEqual(2048 * fixture.delta.length + 16); + const emitted = buffer.ingest(fixture.end); + const drained = buffer.flush([]); + expect(initial + emitted + drained).toBe(fixture.prefix + fixture.delta.repeat(2048) + fixture.end); + expect(drained.length > 0).toBe(fixture.held); + buffer.dispose(); + }); + } + + test("a non-whitespace closing fence suffix resolves immediately as literal text", () => { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + expect(buffer.ingest("```\ncode\n``` ")).toBe("```\ncode\n"); + expect(buffer.ingest(" \t")).toBe(""); + expect(buffer.ingest("language")).toBe("``` \tlanguage"); + // The failed closer leaves us inside the fence: tool markup stays literal. + expect(buffer.ingest("\n" + BLOCK)).toBe("\n" + BLOCK); + buffer.dispose(); + }); +}); + +test("incremental deferred phases emit exactly the full splitter's per-delta output", () => { + const fixtures = [ + " \n\tbody" + CLOSE, + " \n\n" + BLOCK, + ]; + for (const text of fixtures) for (let width = 1; width <= text.length; width++) { + const buffer = new SerializedToolCallContentBuffer(createTestTranslatorBudget()); + let retained = ""; + let opened = false; + let context: Parameters[1]; + for (let at = 0; at < text.length; at += width) { + const delta = text.slice(at, at + width); + let expected = ""; + if (opened) retained += delta; + else { + const split = splitAtPossibleSerializedToolCall(retained + delta, context); + retained = split.defer; context = split.context; opened = split.hasOpenTag; + expected = split.emit; + } + expect(buffer.ingest(delta)).toBe(expected); + expect(buffer.current()).toBe(retained); + } + buffer.dispose(); + } +}); + +test("deferred header byte accounting joins split surrogate pairs and releases its reservation", () => { + const budget = createTestTranslatorBudget(); + const buffer = new SerializedToolCallContentBuffer(budget); + const prefix = " { + if (typeof value === "string") measured += value.length; + return originalMeasure(value, encoding); + }); + const encode = spyOn(TextEncoder.prototype, "encode").mockImplementation(function (this: TextEncoder, value) { + encoded += value?.length ?? 0; + return originalEncode.call(this, value); + }); + try { + const result = await collect(bytes); + expect(result).toEqual([{ text: "x".repeat(length) }]); + } finally { + search.mockRestore(); measure.mockRestore(); encode.mockRestore(); + } + return { searched, measured, encoded, length: bytes.length }; +} + +describe("coding-agent incremental JSONL framing", () => { + test("one-byte chunks search and measure linear input without encoding retained frames", async () => { + const small = await framingWork(2048); + const large = await framingWork(4096); + for (const work of [small, large]) { + expect(work.searched).toBeLessThanOrEqual(work.length); + expect(work.measured).toBeLessThanOrEqual(work.length); + expect(work.encoded).toBe(0); + expect(work.searched + work.measured).toBeGreaterThan(0); + } + expect(large.searched + large.measured).toBeLessThanOrEqual(2 * (small.searched + small.measured)); + }); + + test("decoded-byte limits preserve split UTF-8, BOM, CRLF, blank lines and final frames", async () => { + const text = '{"text":"δΈ–η•ŒπŸ˜€"}'; + const maxLineBytes = Buffer.byteLength(text + "\r"); + const bytes = new TextEncoder().encode("\ufeff\r\n" + text + "\r\n \t\n" + text); + for (const width of [1, 2, 3, bytes.length]) { + expect(await collect(bytes, { maxLineBytes, maxTotalBytes: bytes.length }, width)) + .toEqual([{ text: "δΈ–η•ŒπŸ˜€" }, { text: "δΈ–η•ŒπŸ˜€" }]); + await expect(collect(bytes, { maxLineBytes: maxLineBytes - 1 }, width)).rejects.toBeInstanceOf(CodingAgentStreamLimitError); + await expect(collect(bytes, { maxTotalBytes: bytes.length - 1 }, width)).rejects.toThrow("total byte ceiling"); + } + }); + + test("invalid UTF-8 is measured after replacement, including decoder EOF", async () => { + const bytes = Uint8Array.from([...new TextEncoder().encode('{"x":"'), 0xff, ...new TextEncoder().encode('"}\n')]); + expect(await collect(bytes, { maxLineBytes: 11 })).toEqual([{ x: "οΏ½" }]); + await expect(collect(bytes, { maxLineBytes: 10 })).rejects.toBeInstanceOf(CodingAgentStreamLimitError); + // A dangling leading byte expands to a three-byte replacement only at EOF. + await expect(collect(Uint8Array.of(0xe4), { maxLineBytes: 2 })).rejects.toBeInstanceOf(CodingAgentStreamLimitError); + await expect(collect(Uint8Array.of(0xe4), { maxLineBytes: 3 })).rejects.toBeInstanceOf(CodingAgentProtocolError); + }); + + test("limits and error types apply before parsing each line, including whitespace padding", async () => { + for (const text of [" ", " \n", " \r\n"]) { + await expect(collect(new TextEncoder().encode(text), { maxLineBytes: 2 })).rejects.toBeInstanceOf(CodingAgentStreamLimitError); + } + for (const text of ['{"x":', "[]", "null", "false", "42", '"text"']) { + await expect(collect(new TextEncoder().encode(text))).rejects.toBeInstanceOf(CodingAgentProtocolError); + } + const bytes = new TextEncoder().encode('{}\ninvalid\n'); + const iterator = readJsonLines(chunks(bytes, bytes.length)); + expect(await iterator.next()).toMatchObject({ value: {}, done: false }); + await expect(iterator.next()).rejects.toBeInstanceOf(CodingAgentProtocolError); + const totalFirst = readJsonLines(chunks(bytes, bytes.length), { maxTotalBytes: bytes.length - 1 }); + await expect(totalFirst.next()).rejects.toThrow("total byte ceiling"); + }); +}); diff --git a/tests/usage/usage-snapshot-digest-reuse.test.ts b/tests/usage/usage-snapshot-digest-reuse.test.ts new file mode 100644 index 00000000000..e05cf67ff37 --- /dev/null +++ b/tests/usage/usage-snapshot-digest-reuse.test.ts @@ -0,0 +1,112 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { createHash } from "node:crypto"; +import { appendFileSync, closeSync, mkdtempSync, openSync, rmSync, statSync, writeFileSync, writeSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { readUsageSnapshotForManagement, resetUsageReadCacheForTests, usageLogPath, usageReadCacheStatsForTests } from "../../src/usage/log"; + +let directory: string; +let previousHome: string | undefined; +beforeEach(() => { + previousHome = process.env.OPENCODEX_HOME; + directory = mkdtempSync(join(tmpdir(), "ocx-digest-reuse-")); + process.env.OPENCODEX_HOME = directory; + resetUsageReadCacheForTests(); +}); +afterEach(() => { + resetUsageReadCacheForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + rmSync(directory, { recursive: true, force: true }); +}); + +const row = (id: string) => JSON.stringify({ requestId: id, provider: "mock", model: "fixture" }) + "\n"; +const digest = (text: string, from = 0) => `${from}:${from + Buffer.byteLength(text)}:${createHash("sha256").update(text).digest("hex")}`; + +async function measuredRead(maxReadBytes?: number) { + let hashedBytes = 0; + const prototype = Object.getPrototypeOf(createHash("sha256")); + const update = prototype.update; + const spy = spyOn(prototype, "update").mockImplementation(function (this: unknown, data: string | Uint8Array, ...args: unknown[]) { + hashedBytes += typeof data === "string" ? Buffer.byteLength(data) : data.byteLength; + return update.call(this, data, ...args); + }); + try { + return { snapshot: await readUsageSnapshotForManagement(maxReadBytes), get hashedBytes() { return hashedBytes; } }; + } finally { + spy.mockRestore(); + } +} + +describe("retained usage digest reuse", () => { + test("each unchanged poll hashes the retained bytes exactly once, at either input size", async () => { + for (const count of [64, 128]) { + resetUsageReadCacheForTests(); + const text = Array.from({ length: count }, (_, i) => row(String(i).padStart(3, "0"))).join(""); + writeFileSync(usageLogPath(), text); + const initial = await readUsageSnapshotForManagement(); + for (let poll = 0; poll < 3; poll++) { + const { snapshot, hashedBytes } = await measuredRead(); + expect(hashedBytes).toBe(Buffer.byteLength(text)); + expect(snapshot.prefixDigest).toBe(digest(text)); + expect(snapshot.entries).toEqual(initial.entries); + expect(snapshot.entries).not.toBe(initial.entries); + } + expect(usageReadCacheStatsForTests()).toEqual({ fullReads: 1, tailReads: 3, parsedLines: count }); + } + }); + + test("growth hashes the new region and a same-inode fixed-width edit still invalidates retention", async () => { + const original = row("old"); + const appended = row("add"); + const path = usageLogPath(); + writeFileSync(path, original); + await readUsageSnapshotForManagement(); + appendFileSync(path, appended); + const grown = await measuredRead(); + expect(grown.hashedBytes).toBe(Buffer.byteLength(original) + Buffer.byteLength(original + appended)); + expect(grown.snapshot.prefixDigest).toBe(digest(original + appended)); + const inode = statSync(path).ino; + const fd = openSync(path, "r+"); + try { writeSync(fd, Buffer.from(row("new")), 0, Buffer.byteLength(original), 0); } + finally { closeSync(fd); } + expect(statSync(path).ino).toBe(inode); + const rewritten = await readUsageSnapshotForManagement(); + expect(rewritten.entries.map(entry => entry.requestId)).toEqual(["new", "add"]); + expect(rewritten.prefixDigest).toBe(digest(row("new") + appended)); + expect(usageReadCacheStatsForTests().fullReads).toBe(2); + }); + + test("equal-width sliding windows rehash the post-trim region, then reuse only that region", async () => { + const a = row("aaa"), b = row("bbb"), c = row("ccc"), d = row("ddd"); + const width = Buffer.byteLength(b + c); + writeFileSync(usageLogPath(), a + b + c); + const initial = await readUsageSnapshotForManagement(width); + expect(initial.rowsBeginAtBytes).toBe(Buffer.byteLength(a)); + appendFileSync(usageLogPath(), d); + const changed = await measuredRead(width); + expect(changed.hashedBytes).toBe(2 * width); + expect(changed.snapshot.prefixDigest).toBe(digest(c + d, Buffer.byteLength(a + b))); + expect(changed.snapshot.entries.map(entry => entry.requestId)).toEqual(["ccc", "ddd"]); + const unchanged = await measuredRead(width); + expect(unchanged.hashedBytes).toBe(width); + expect(unchanged.snapshot).toEqual(changed.snapshot); + resetUsageReadCacheForTests(); + expect(await readUsageSnapshotForManagement(width)).toEqual(changed.snapshot); + }); + + test("empty files and trailing invalid bytes retain exact range metadata", async () => { + writeFileSync(usageLogPath(), ""); + await readUsageSnapshotForManagement(); + const empty = await measuredRead(); + expect(empty.hashedBytes).toBe(0); + expect(empty.snapshot.prefixDigest).toBe("0:0:empty"); + const text = row("one") + "invalid\n"; + writeFileSync(usageLogPath(), text); + const first = await readUsageSnapshotForManagement(); + const next = await measuredRead(); + expect(next.hashedBytes).toBe(Buffer.byteLength(text)); + expect(next.snapshot).toEqual(first); + expect(next.snapshot.prefixDigest).toBe(digest(text)); + }); +}); From ffc203dad76237d61faf0737ad201cccb6c3f6f4 Mon Sep 17 00:00:00 2001 From: Ingwannu Date: Fri, 25 Sep 2026 10:14:23 +0000 Subject: [PATCH 2/3] fix: normalize split surrogate byte accounting --- .../serialized-tool-call-content.ts | 22 ++++++++++++++----- structure/dashboard-and-usage.md | 9 +++----- .../ADR-0102-incremental-stream-accounting.md | 6 ++--- 3 files changed, 22 insertions(+), 15 deletions(-) diff --git a/src/adapters/openai-chat/serialized-tool-call-content.ts b/src/adapters/openai-chat/serialized-tool-call-content.ts index ff60738068a..d31a339c125 100644 --- a/src/adapters/openai-chat/serialized-tool-call-content.ts +++ b/src/adapters/openai-chat/serialized-tool-call-content.ts @@ -235,6 +235,20 @@ export class SerializedToolCallContentBuffer { constructor(private readonly budget: TranslatorBudget) {} + /** Measures an append exactly even when the runtime prices isolated UTF-16 surrogates differently. */ + private appendedByteLength(delta: string): number { + let bytes = Buffer.byteLength(delta); + if (this.text.length === 0 || delta.length === 0) return bytes; + const tail = this.text[this.text.length - 1]!; + const head = delta[0]!; + const tailCode = tail.charCodeAt(0); + const headCode = head.charCodeAt(0); + if (tailCode >= 0xd800 && tailCode <= 0xdbff && headCode >= 0xdc00 && headCode <= 0xdfff) { + bytes += Buffer.byteLength(tail + head) - Buffer.byteLength(tail) - Buffer.byteLength(head); + } + return bytes; + } + /** Reserves the replacement before releasing the old text, preserving it if the budget rejects growth. */ private replace(next: string, hasOpenTag: boolean): void { const nextBytes = Buffer.byteLength(next); @@ -259,9 +273,7 @@ export class SerializedToolCallContentBuffer { /** Charges only the appended bytes, so holding an open block never needs twice its retained size. */ private append(delta: string): void { - // Bun counts each isolated surrogate as two bytes, so split pairs sum to the same four-byte - // charge as the joined scalar. Measuring only the delta therefore preserves exact accounting. - const deltaBytes = Buffer.byteLength(delta); + const deltaBytes = this.appendedByteLength(delta); this.budget.reserveTransient(deltaBytes, { kind: "live_transient" }).commitRetained(); if (this.hasOpenTag) this.observeDelimiters(delta); this.text += delta; @@ -306,11 +318,11 @@ export class SerializedToolCallContentBuffer { * context. The size bound is checked before the delta is retained. */ ingestStreaming(delta: string): AdapterEvent[] { - const deltaBytes = Buffer.byteLength(delta); + const deltaBytes = this.appendedByteLength(delta); if (this.hasOpenTag && this.bytes + deltaBytes > MAX_HELD_BYTES) { const released = this.drain([]); // A delta that alone passes the bound is delivered as text rather than retained. - if (deltaBytes > MAX_HELD_BYTES) { + if (Buffer.byteLength(delta) > MAX_HELD_BYTES) { this.context = contextAfter(delta, this.context); return [...released, ...textEvents(delta)]; } diff --git a/structure/dashboard-and-usage.md b/structure/dashboard-and-usage.md index 5e9c5161110..cd371cd7e97 100644 --- a/structure/dashboard-and-usage.md +++ b/structure/dashboard-and-usage.md @@ -343,13 +343,10 @@ is treated as an append: the scanner verifies the previous LF and its trailing 6 folds only the suffix into a cloned accumulator and publishes it after validation. Concurrent callers share that work. Cold rebuilds scan the whole ledger in fixed-size chunks and yield between bounded batches, so memory stays bounded and unrelated management requests remain serviceable even for a -large existing log. The first read is proportional to ledger size; steady-state refresh work is -proportional to newly appended bytes. The Dashboard polls its 30-day usage summary independently once +large existing log. The first read is proportional to ledger size; steady-state refresh work is proportional to newly appended bytes. The Dashboard polls its 30-day usage summary independently once per minute, so usage work cannot delay health/provider/settings state or run every five seconds. -An unchanged retained snapshot verifies its byte region once and reuses that digest only when the -returned region has identical start and end bounds. Appends or window trimming still hash the new -returned region, so digest reuse cannot weaken same-inode rewrite detection. See -[ADR-0102](decisions/ADR-0102-incremental-stream-accounting.md). +An unchanged retained snapshot reuses its verified region digest only for identical bounds; appends or trimming hash the returned region, preserving same-inode rewrite detection. +> Decision record: [ADR-0102](decisions/ADR-0102-incremental-stream-accounting.md) An oversized row is skipped inside the scanner bound without shortening identities. Accumulators keep normal rows plus `usageIncomplete` / `usageIncompleteReason: "oversized_rows"` on caches and rollups; append ORs the flag and a rebuild recalculates it. Invalid-row counts are not sticky, and absence of the flag is not completeness. GUI caches warn on Usage, Dashboard, provider and key views; CLI warns in human output only; most-used order save refuses an incomplete snapshot. Quota surfaces stay separate. Legacy truncation fields keep their meaning; read/mutation failures still fail closed. diff --git a/structure/decisions/ADR-0102-incremental-stream-accounting.md b/structure/decisions/ADR-0102-incremental-stream-accounting.md index 00118712fa7..b89ee53f0bc 100644 --- a/structure/decisions/ADR-0102-incremental-stream-accounting.md +++ b/structure/decisions/ADR-0102-incremental-stream-accounting.md @@ -1,8 +1,6 @@ -# ADR-0102 β€” Incremental stream accounting +# ADR-0102 β€” decision recorded under "Usage accounting" -- Contract owners: [Chat compatibility](../providers/chat-compat.md#serialized-tool-call-content), - [Providers and adapters](../providers-and-adapters.md), and - [Dashboard and usage](../dashboard-and-usage.md) +- Contract owner: [Dashboard and usage](../dashboard-and-usage.md#usage-accounting) ## Decision Log From de264791053f343c89cc5292af8a8b6f2f5c400a Mon Sep 17 00:00:00 2001 From: Ingwannu Date: Fri, 25 Sep 2026 10:58:05 +0000 Subject: [PATCH 3/3] docs(usage): state retained-window refresh cost --- structure/dashboard-and-usage.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/structure/dashboard-and-usage.md b/structure/dashboard-and-usage.md index cd371cd7e97..6b7b2a2bf30 100644 --- a/structure/dashboard-and-usage.md +++ b/structure/dashboard-and-usage.md @@ -343,8 +343,8 @@ is treated as an append: the scanner verifies the previous LF and its trailing 6 folds only the suffix into a cloned accumulator and publishes it after validation. Concurrent callers share that work. Cold rebuilds scan the whole ledger in fixed-size chunks and yield between bounded batches, so memory stays bounded and unrelated management requests remain serviceable even for a -large existing log. The first read is proportional to ledger size; steady-state refresh work is proportional to newly appended bytes. The Dashboard polls its 30-day usage summary independently once -per minute, so usage work cannot delay health/provider/settings state or run every five seconds. +large existing log. The first read is proportional to ledger size; later refreshes hash the bounded retained window once and parse only the appended suffix, rather than rescanning the whole ledger. +The Dashboard polls its 30-day usage independently once per minute, separate from five-second state polls. An unchanged retained snapshot reuses its verified region digest only for identical bounds; appends or trimming hash the returned region, preserving same-inode rewrite detection. > Decision record: [ADR-0102](decisions/ADR-0102-incremental-stream-accounting.md)