Skip to content
Closed
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
119 changes: 119 additions & 0 deletions apps/server/src/usage/usageTranscriptReader.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
});
49 changes: 35 additions & 14 deletions apps/server/src/usage/usageTranscriptReader.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -194,7 +196,9 @@ export async function readTranscriptRecords(
filePath: string,
provider: UsageProviderKind,
resumeFrom?: TranscriptParsePosition,
options?: { readonly maxLineBytes?: number },
): Promise<TranscriptParseResult | null> {
const maxLineBytes = options?.maxLineBytes ?? DEFAULT_MAX_LINE_BYTES;
let handle: NodeFSP.FileHandle;
try {
handle = await NodeFSP.open(filePath, "r");
Expand Down Expand Up @@ -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<Buffer>;
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
Expand Down
Loading