Skip to content
Merged
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
3 changes: 3 additions & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
48 changes: 29 additions & 19 deletions src/adapters/coding-agent/protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<StreamMessage> {
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;
Expand All @@ -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<StreamMessage> {
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<string, unknown> | undefined {
Expand Down
88 changes: 74 additions & 14 deletions src/adapters/openai-chat/serialized-tool-call-content.ts
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,29 @@ function contextAfter(text: string, context: TextContext): TextContext {
return splitAtPossibleSerializedToolCall(text.replaceAll(OPEN_TAG, "<tool-call>"), 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 (/^<tool_call>\s*$/.test(text)) return /\S/;
if (/^<tool_call>\s*<function=[^>\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
Expand All @@ -202,11 +225,30 @@ 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 }[] = [];

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);
Expand All @@ -217,6 +259,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;
Expand All @@ -225,21 +273,41 @@ export class SerializedToolCallContentBuffer {

/** Charges only the appended bytes, so holding an open block never needs twice its retained size. */
private append(delta: string): void {
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;
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;
}

Expand All @@ -250,19 +318,19 @@ 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)];
}
return [...released, ...textEvents(this.ingest(delta))];
}
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);
Expand Down Expand Up @@ -334,6 +402,7 @@ export class SerializedToolCallContentBuffer {
this.text = "";
this.bytes = 0;
this.hasOpenTag = false;
this.continuationBreak = undefined;
this.queued = [];
}
}
Expand Down Expand Up @@ -595,12 +664,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;
}
7 changes: 4 additions & 3 deletions src/usage/log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
7 changes: 4 additions & 3 deletions structure/dashboard-and-usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -343,9 +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
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.
Comment thread
Ingwannu marked this conversation as resolved.
> 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.

Expand Down
21 changes: 21 additions & 0 deletions structure/decisions/ADR-0102-incremental-stream-accounting.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
# ADR-0102 — decision recorded under "Usage accounting"

- Contract owner: [Dashboard and usage](../dashboard-and-usage.md#usage-accounting)
Comment thread
coderabbitai[bot] marked this conversation as resolved.

## 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.
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<tool_call>` 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 `</tool_call>` preceded by `</function>` closes the block; if none appears before the next block header at the start of a line (`<tool_call>` followed by `<function=`, where a separate bare block can begin), the first `</tool_call>` does. A stray `</parameter>` 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).
6 changes: 6 additions & 0 deletions structure/providers-and-adapters.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading