diff --git a/apps/server/src/usage/usageTranscriptReader.test.ts b/apps/server/src/usage/usageTranscriptReader.test.ts index 5feb68b2ff58..23f718ace4f8 100644 --- a/apps/server/src/usage/usageTranscriptReader.test.ts +++ b/apps/server/src/usage/usageTranscriptReader.test.ts @@ -208,3 +208,122 @@ describe("readTranscriptRecords resume", () => { assert.isNull(await readTranscriptRecords(NodePath.join(dir, "missing.jsonl"), "claude")); }); }); + +describe("readTranscriptRecords bounded lines", () => { + const options = { maxLineBytes: 1024 }; + const oversized = JSON.stringify({ type: "tool_output", output: "x".repeat(128 * 1024) }); + + it.each(["\n", "\r\n"])( + "skips oversized records before and after a resume with %j line endings", + async (newline) => { + const path = NodePath.join(dir, "rollout.jsonl"); + const prefix = ( + codexMetaLine() + + codexModelLine("gpt-5.2-codex") + + codexUsageLine(9, 5) + + `${oversized}\n` + ).replaceAll("\n", newline); + await NodeFSP.writeFile(path, prefix); + const first = await readTranscriptRecords(path, "codex", undefined, options); + assert.isNotNull(first); + assert.strictEqual(first.position.resumeOffset, Buffer.byteLength(prefix)); + assert.deepStrictEqual(first.tailRecords, []); + + const appended = (`${oversized}\n` + codexUsageLine(9, 5) + codexUsageLine(21, 8)).replaceAll( + "\n", + newline, + ); + await NodeFSP.appendFile(path, appended); + const second = await readTranscriptRecords(path, "codex", first.position, options); + assert.isNotNull(second); + assert.isTrue(second.resumed); + assert.strictEqual(second.position.resumeOffset, Buffer.byteLength(prefix + appended)); + assert.deepStrictEqual( + second.records.map((record) => record.totals.outputTokens), + [21], + ); + assert.strictEqual(second.records[0]?.model, "gpt-5.2-codex"); + assert.strictEqual(second.records[0]?.sessionId, "codex-session-1"); + const full = await readTranscriptRecords(path, "codex", undefined, options); + assert.isNotNull(full); + assert.deepStrictEqual([...first.records, ...second.records], full.records); + assert.deepStrictEqual(second.position, full.position); + }, + ); + + it.each([false, true])( + "leaves an oversized unfinished tail outside the committed position (resumed: %s)", + async (resume) => { + const path = NodePath.join(dir, "rollout.jsonl"); + const prefix = codexMetaLine() + codexModelLine("gpt-5.2-codex"); + await NodeFSP.writeFile(path, prefix); + const initial = await readTranscriptRecords(path, "codex", undefined, options); + assert.isNotNull(initial); + await NodeFSP.appendFile(path, oversized); + const first = await readTranscriptRecords( + path, + "codex", + resume ? initial.position : undefined, + options, + ); + assert.isNotNull(first); + assert.strictEqual(first.resumed, resume); + assert.deepStrictEqual(first.position, initial.position); + assert.deepStrictEqual(first.records, []); + assert.deepStrictEqual(first.tailRecords, []); + + const appended = `\n${codexUsageLine(21, 8)}`; + await NodeFSP.appendFile(path, appended); + const second = await readTranscriptRecords(path, "codex", first.position, options); + assert.isNotNull(second); + assert.isTrue(second.resumed); + assert.strictEqual( + second.position.resumeOffset, + Buffer.byteLength(prefix + oversized + appended), + ); + assert.strictEqual(second.records.length, 1); + assert.strictEqual(second.records[0]?.totals.outputTokens, 21); + assert.strictEqual(second.records[0]?.model, "gpt-5.2-codex"); + assert.strictEqual(second.records[0]?.sessionId, "codex-session-1"); + }, + ); + + it("enforces the limit in bytes, accepting the exact limit and rejecting one extra byte", async () => { + const path = NodePath.join(dir, "claude.jsonl"); + const record = claudeLine(1, 5).trimEnd(); + const unicodeRecord = record.replace("session-1", "session-🌍"); + const maxLineBytes = Buffer.byteLength(unicodeRecord); + await NodeFSP.writeFile(path, `${unicodeRecord}\n${unicodeRecord} \n${claudeLine(2, 7)}`); + const result = await readTranscriptRecords(path, "claude", undefined, { maxLineBytes }); + assert.isNotNull(result); + assert.deepStrictEqual( + result.records.map((entry) => entry.totals.outputTokens), + [5, 7], + ); + assert.strictEqual(result.position.resumeOffset, (await NodeFSP.stat(path)).size); + }); + + it("does not parse an oversized usage tail or commit its Codex reducer state", async () => { + const path = NodePath.join(dir, "rollout.jsonl"); + const prefix = codexMetaLine() + codexModelLine("gpt-5.2-codex"); + await NodeFSP.writeFile(path, prefix); + const initial = await readTranscriptRecords(path, "codex", undefined, options); + assert.isNotNull(initial); + const tail = codexModelLine("ignored-model").trimEnd() + " ".repeat(2048); + await NodeFSP.appendFile(path, tail); + const first = await readTranscriptRecords(path, "codex", initial.position, options); + assert.isNotNull(first); + assert.deepStrictEqual(first.position, initial.position); + assert.deepStrictEqual(first.tailRecords, []); + await NodeFSP.appendFile(path, `\n${codexUsageLine(21, 8)}`); + const second = await readTranscriptRecords(path, "codex", first.position, options); + assert.isNotNull(second); + assert.strictEqual(second.records[0]?.model, "gpt-5.2-codex"); + + await NodeFSP.writeFile(path, claudeLine(1, 5).trimEnd() + " ".repeat(2048)); + const oversizedTail = await readTranscriptRecords(path, "claude", undefined, options); + assert.isNotNull(oversizedTail); + assert.deepStrictEqual(oversizedTail.tailRecords, []); + assert.strictEqual(oversizedTail.position.resumeOffset, 0); + }); +}); diff --git a/apps/server/src/usage/usageTranscriptReader.ts b/apps/server/src/usage/usageTranscriptReader.ts index 9e5ab6e0c9e0..c7fa1d406148 100644 --- a/apps/server/src/usage/usageTranscriptReader.ts +++ b/apps/server/src/usage/usageTranscriptReader.ts @@ -74,6 +74,8 @@ export interface TranscriptParseResult { /** 64 bytes of JSONL tail is ample to distinguish a replaced file. */ export const GUARD_LENGTH = 64; +// Bound records before decoding, well below V8's maximum string length. +const DEFAULT_MAX_LINE_BYTES = 64 * 1024 * 1024; const NEWLINE = 0x0a; const CARRIAGE_RETURN = 0x0d; @@ -194,7 +196,9 @@ export async function readTranscriptRecords( filePath: string, provider: UsageProviderKind, resumeFrom?: TranscriptParsePosition, + options?: { readonly maxLineBytes?: number }, ): Promise { + const maxLineBytes = options?.maxLineBytes ?? DEFAULT_MAX_LINE_BYTES; let handle: NodeFSP.FileHandle; try { handle = await NodeFSP.open(filePath, "r"); @@ -250,31 +254,48 @@ export async function readTranscriptRecords( const records: UsageRecord[] = []; // Buffer-level line splitting rather than `readline`, because resuming // needs byte-exact offsets and decoded strings cannot provide them. - // Newline-free chunks are collected rather than concatenated as they - // arrive, so a single huge line costs one copy instead of one per chunk. + // Oversized lines are drained without decoding or retaining their bytes. + // The scan offset advances through them, but only a newline commits it. let resumeOffset = start; + let scanOffset = start; let pendingChunks: Buffer[] = []; + let pendingBytes = 0; + let discardingLine = false; const stream = handle.createReadStream({ start, autoClose: false, }) as AsyncIterable; for await (const chunk of stream) { - if (!chunk.includes(NEWLINE)) { - pendingChunks.push(chunk); - continue; - } - const buffer: Buffer = - pendingChunks.length === 0 ? chunk : Buffer.concat([...pendingChunks, chunk]); - pendingChunks = []; let lineStart = 0; - for (;;) { - const newlineIndex = buffer.indexOf(NEWLINE, lineStart); + while (lineStart < chunk.length) { + const newlineIndex = chunk.indexOf(NEWLINE, lineStart); + const lineEnd = newlineIndex === -1 ? chunk.length : newlineIndex; + const segment = chunk.subarray(lineStart, lineEnd); + if (!discardingLine) { + if (pendingBytes + segment.length > maxLineBytes) { + pendingChunks = []; + pendingBytes = 0; + discardingLine = true; + } else if (segment.length > 0) { + pendingChunks.push(segment); + pendingBytes += segment.length; + } + } if (newlineIndex === -1) break; - parseLine(toLineString(buffer.subarray(lineStart, newlineIndex)), codexState, records); + if (!discardingLine && pendingBytes > 0) { + const line = + pendingChunks.length === 1 + ? pendingChunks[0]! + : Buffer.concat(pendingChunks, pendingBytes); + parseLine(toLineString(line), codexState, records); + } lineStart = newlineIndex + 1; + resumeOffset = scanOffset + lineStart; + pendingChunks = []; + pendingBytes = 0; + discardingLine = false; } - resumeOffset += lineStart; - if (lineStart < buffer.length) pendingChunks.push(buffer.subarray(lineStart)); + scanOffset += chunk.length; } // A trailing segment without its newline is parsed for this result but not