diff --git a/apps/server/src/project/AgentSessionJson.test.ts b/apps/server/src/project/AgentSessionJson.test.ts new file mode 100644 index 000000000000..515aa50b1b97 --- /dev/null +++ b/apps/server/src/project/AgentSessionJson.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, it } from "@effect/vitest"; + +import { createTranscriptJsonReader, TranscriptJsonLimitError } from "./AgentSessionJson.ts"; + +function read( + json: string, + size: number, + select: (path: ReadonlyArray) => boolean = () => true, +) { + const reader = createTranscriptJsonReader(() => {}, select); + for (let offset = 0; offset < json.length; offset += size) + reader.write(json.slice(offset, offset + size)); + return reader.finish(); +} + +describe("transcript JSON projection", () => { + it.each([1, 2, 7, 64, 1024])("matches JSON.parse across %i-character boundaries", (size) => { + for (const text of [ + '{"a":1,"a":2,"b":"before","b":"after"}', + '{"a":{"x":1},"a":{"y":2},"b":[],"c":{}}', + '{"a":[null,true,false,1,-2.3e4,"😀\\u0061\\\\\\\"",{},[],[1,2]]}', + '{"a":"s","a":null,"b":null,"b":"s","c":0,"c":false}', + '{"__proto__":{"polluted":true},"constructor":1,"__proto__":2}', + ]) + expect(read(text, size)).toEqual(JSON.parse(text)); + }); + + it("projects siblings and array elements without merging repeated parent objects", () => { + const text = + '{"message":{"usage":{"input":100},"content":"large"},"message":{"usage":{"output":5},"content":[1,2]},"rows":[{"keep":1,"drop":2},{"keep":3}],"drop":{"keep":4}}'; + const projected = read(text, 1, (path) => { + if (path[0] === "drop") return false; + return !path.includes("content") && !path.includes("drop"); + }); + expect(projected).toEqual({ + message: { usage: { output: 5 } }, + rows: [{ keep: 1 }, { keep: 3 }], + }); + }); + + it.each(['{"a":', '{"a":1} trailing', '{"a":1}{"a":2}', '{"a":"bad\\x"}', '{"a":[1,]}'])( + "rejects malformed input %s", + (text) => { + expect(read(text, 1)).toBeUndefined(); + }, + ); + + it("retains the import allocation and depth limits", () => { + const limited = createTranscriptJsonReader( + () => { + throw new TranscriptJsonLimitError("budget"); + }, + () => true, + ); + expect(() => limited.write('{"a":1}')).toThrow(TranscriptJsonLimitError); + expect(() => read("[".repeat(129) + "0" + "]".repeat(129), 10)).toThrow( + TranscriptJsonLimitError, + ); + }); +}); diff --git a/apps/server/src/project/AgentSessionJson.ts b/apps/server/src/project/AgentSessionJson.ts index d31c47833d70..4deba0a67794 100644 --- a/apps/server/src/project/AgentSessionJson.ts +++ b/apps/server/src/project/AgentSessionJson.ts @@ -1,7 +1,6 @@ import * as SchemaAST from "effect/SchemaAST"; -import { isMany, none, type Many } from "stream-chain/defs.js"; +import { none, type Many } from "stream-chain/defs.js"; import { Assembler } from "stream-json/core/assembler.js"; -import { filter } from "stream-json/core/filters/filter.js"; import * as StreamJson from "stream-json/core/parser.js"; import type { ParserOptions, Token } from "stream-json/core/parser.js"; @@ -55,6 +54,7 @@ export class TranscriptJsonLimitError extends Error {} export function createTranscriptJsonReader( reserve: (bytes: number) => void, selectPath: (path: JsonPath) => boolean, + options?: { readonly maxDepth?: number }, ) { // The synchronous tokenizer is exported at runtime in 3.6.0, but omitted // from its bundled types. Unlike parser(), it does not wrap tokens in an @@ -65,9 +65,6 @@ export function createTranscriptJsonReader( ) => (input: string | typeof none) => Many | typeof none; }; const tokenize = jsonParser({ packValues: false }); - const select = filter({ filter: selectPath, streamKeys: false }) as ( - input: Token | typeof none, - ) => Token | Many | typeof none; const assembler = new Assembler(); let key: string | null = null; let value = ""; @@ -100,13 +97,62 @@ export function createTranscriptJsonReader( assembler.consume(token); } }; + // Forward actual selected keys instead of reconstructing them from path + // changes: adjacent duplicate keys have the same path but JSON.parse keeps + // the last value. Reconstructing paths can silently retain the first value. + const stack: Array<{ path: JsonPath; key: string | number | null; selected: boolean }> = []; + let selectedValue = false; + const startValue = () => { + const parent = stack.at(-1); + const path = parent?.selected ? [...parent.path, parent.key] : []; + const selected = (parent?.selected ?? true) && selectPath(path); + if (selected && typeof parent?.key === "string") { + assemble({ name: "keyValue", value: parent.key }); + } + return { path, selected }; + }; + const endValue = () => { + const parent = stack.at(-1); + if (parent && typeof parent.key === "number") parent.key++; + }; const selectToken = (token: Token | typeof none) => { - const selected = select(token); - if (selected === none) return; - if (isMany(selected)) { - for (const item of selected.values) assemble(item); - } else { - assemble(selected); + if (token === none) return; + switch (token.name) { + case "keyValue": { + const parent = stack.at(-1); + if (parent) parent.key = token.value; + return; + } + case "startObject": + case "startArray": { + const frame = startValue(); + stack.push({ ...frame, key: token.name === "startArray" ? 0 : null }); + if (frame.selected) assemble(token); + return; + } + case "endObject": + case "endArray": + if (stack.pop()?.selected) assemble(token); + endValue(); + return; + case "startString": + case "startNumber": + selectedValue = startValue().selected; + if (selectedValue) assemble(token); + return; + case "endString": + case "endNumber": + if (selectedValue) assemble(token); + endValue(); + return; + case "nullValue": + case "trueValue": + case "falseValue": + if (startValue().selected) assemble(token); + endValue(); + return; + default: + if (selectedValue) assemble(token); } }; const consume = (input: string | typeof none) => { @@ -116,8 +162,8 @@ export function createTranscriptJsonReader( if (tokens === none) return; for (const token of tokens.values) { if (token.name === "startObject" || token.name === "startArray") { - if (++depth > 128) - throw new TranscriptJsonLimitError("Transcript JSON nesting exceeds 128 levels"); + if (++depth > (options?.maxDepth ?? 128)) + throw new TranscriptJsonLimitError("Transcript JSON nesting exceeds the depth limit"); } else if (token.name === "endObject" || token.name === "endArray") { if (--depth === 0) complete = true; } diff --git a/apps/server/src/usage/UsageService.test.ts b/apps/server/src/usage/UsageService.test.ts index 15ea4b673f22..df6d7ba044c8 100644 --- a/apps/server/src/usage/UsageService.test.ts +++ b/apps/server/src/usage/UsageService.test.ts @@ -645,6 +645,59 @@ describe("UsageService", () => { }).pipe(Effect.scoped), ); + it.live( + "keeps large-record totals and costs exact through append, dedupe, restart and cleanup", + () => + Effect.gen(function* () { + const { transcript, settings, home } = yield* setup; + const large = claudeLine(1, 9900).replace( + '"message":', + '"padding":' + encodeUnknownJsonString("x".repeat(9 * 1024 * 1024)) + ',"message":', + ); + yield* Effect.promise(() => NodeFSP.writeFile(transcript, large)); + yield* Effect.gen(function* () { + const service = yield* UsageService.make; + const first = yield* service.readSummary(WINDOW); + assert.strictEqual(totalOutputTokens(first), 9900); + assert.closeTo( + first.buckets.reduce((sum, bucket) => sum + bucket.costUsd, 0), + 0.4951, + 1e-12, + ); + const warm = yield* service.readSummary(WINDOW); + assert.deepStrictEqual(warm.buckets, first.buckets); + // The repeated content block has the same message/request identity. + yield* Effect.promise(() => NodeFSP.appendFile(transcript, large + claudeLine(2, 100))); + const appended = yield* service.readSummary(WINDOW); + assert.strictEqual(totalOutputTokens(appended), 10000); + assert.strictEqual( + appended.buckets.reduce((sum, bucket) => sum + bucket.totals.uncachedInputTokens, 0), + 20, + ); + const restarted = yield* UsageService.make; + const restored = yield* restarted.readSummary(WINDOW); + assert.deepStrictEqual(restored.buckets, appended.buckets); + yield* Effect.promise(() => NodeFSP.rm(transcript)); + const afterCleanup = yield* UsageService.make; + assert.deepStrictEqual( + (yield* afterCleanup.readSummary(WINDOW)).buckets, + appended.buckets, + ); + }).pipe( + Effect.provide( + serviceLayers({ + prefix: "usage-service-large-record-test", + home, + settings, + ratesDocument: { + "claude-fable-5": { input_cost_per_token: 1e-5, output_cost_per_token: 5e-5 }, + }, + }), + ), + ); + }).pipe(Effect.scoped), + ); + it.live("preserves saved tokens, costs and sessions after transcript cleanup and restart", () => Effect.gen(function* () { const { transcript, settings, home } = yield* setup; diff --git a/apps/server/src/usage/usageTranscriptReader.ts b/apps/server/src/usage/usageTranscriptReader.ts index 9e5ab6e0c9e0..faa686990777 100644 --- a/apps/server/src/usage/usageTranscriptReader.ts +++ b/apps/server/src/usage/usageTranscriptReader.ts @@ -17,15 +17,21 @@ */ import * as NodeFSP from "node:fs/promises"; import * as NodePath from "node:path"; +import * as NodeStringDecoder from "node:string_decoder"; import type { UsageProviderKind } from "@t3tools/contracts"; +import { createTranscriptJsonReader } from "../project/AgentSessionJson.ts"; + import { initialCodexScanState, mightCarryUsage, parseClaudeLine, + parseClaudeRecord, parseCodexLine, + parseCodexRecord, parseGrokLine, + parseGrokRecord, type CodexScanState, type UsageRecord, } from "./usageTranscripts.ts"; @@ -74,9 +80,62 @@ export interface TranscriptParseResult { /** 64 bytes of JSONL tail is ample to distinguish a replaced file. */ export const GUARD_LENGTH = 64; +// Native parsing is faster for common 1–4 MiB context/tool records. Above +// 8 MiB, project usage without allocating the whole record. This switches +// readers; it never discards a record because of its size. +const STREAMING_THRESHOLD_BYTES = 8 * 1024 * 1024; const NEWLINE = 0x0a; const CARRIAGE_RETURN = 0x0d; +type SelectedFields = { readonly [key: string]: true | SelectedFields }; + +// Keep the fields consumed by usageTranscripts, including reducer state and +// dedupe/cost metadata. A selected subtree (usage) keeps future token fields. +const USAGE_FIELDS: Record<"claude" | "codex" | "grok", SelectedFields> = { + claude: { + type: true, + timestamp: true, + requestId: true, + sessionId: true, + costUSD: true, + message: { id: true, model: true, usage: true }, + }, + codex: { + type: true, + timestamp: true, + payload: { + type: true, + id: true, + session_id: true, + model: true, + forked_from_id: true, + source: { subagent: { thread_spawn: { parent_thread_id: true } } }, + info: { last_token_usage: true }, + }, + }, + grok: { + timestamp: true, + params: { + sessionId: true, + _meta: { agentTimestampMs: true }, + update: { sessionUpdate: true, prompt_id: true, usage: true }, + }, + }, +}; + +function selectUsageFields(provider: UsageProviderKind) { + const fields = USAGE_FIELDS[provider === "codex" || provider === "grok" ? provider : "claude"]; + return (path: ReadonlyArray): boolean => { + let selected: true | SelectedFields = fields; + for (const key of path) { + if (selected === true) return true; + if (typeof key !== "string" || !Object.hasOwn(selected, key)) return false; + selected = selected[key]!; + } + return true; + }; +} + function fnv1a(buffer: Buffer): number { let hash = 0x811c9dc5; for (let index = 0; index < buffer.length; index += 1) { @@ -194,7 +253,9 @@ export async function readTranscriptRecords( filePath: string, provider: UsageProviderKind, resumeFrom?: TranscriptParsePosition, + options?: { readonly streamingThresholdBytes?: number }, ): Promise { + const streamingThresholdBytes = options?.streamingThresholdBytes ?? STREAMING_THRESHOLD_BYTES; let handle: NodeFSP.FileHandle; try { handle = await NodeFSP.open(filePath, "r"); @@ -248,43 +309,87 @@ 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. + // Byte offsets remain independent of UTF-8 decoding. Only complete lines + // commit the resume point; an unfinished tail is replayed on the next scan. let resumeOffset = start; + let scanOffset = start; let pendingChunks: Buffer[] = []; + let pendingBytes = 0; + let streaming: ReturnType | undefined; + let decoder: NodeStringDecoder.StringDecoder | undefined; + const selectPath = selectUsageFields(provider); + + const append = (segment: Buffer) => { + if (!streaming && pendingBytes + segment.length <= streamingThresholdBytes) { + if (segment.length > 0) pendingChunks.push(segment); + pendingBytes += segment.length; + return; + } + if (!streaming) { + // Usage has no import-history budget: retain all selected usage fields, + // regardless of the size of the surrounding unselected tool content. + streaming = createTranscriptJsonReader(() => {}, selectPath, { maxDepth: Infinity }); + decoder = new NodeStringDecoder.StringDecoder("utf8"); + for (const pending of pendingChunks) streaming.write(decoder.write(pending)); + pendingChunks = []; + pendingBytes = 0; + } + streaming.write(decoder!.write(segment)); + }; + const finish = (state: CodexScanState, out: UsageRecord[]) => { + if (streaming) { + streaming.write(decoder!.end()); + const projected = streaming.finish(); + if (provider === "grok") { + out.push(...parseGrokRecord(projected)); + } else { + const record = + provider === "codex" + ? parseCodexRecord(projected, state) + : parseClaudeRecord(projected); + if (record !== null) out.push(record); + } + } else if (pendingBytes > 0) { + const line = + pendingChunks.length === 1 + ? pendingChunks[0]! + : Buffer.concat(pendingChunks, pendingBytes); + parseLine(toLineString(line), state, out); + } + pendingChunks = []; + pendingBytes = 0; + streaming = undefined; + decoder = undefined; + }; const stream = handle.createReadStream({ start, autoClose: false, + highWaterMark: 256 * 1024, }) 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); - if (newlineIndex === -1) break; - parseLine(toLineString(buffer.subarray(lineStart, newlineIndex)), codexState, records); + while (lineStart < chunk.length) { + const newlineIndex = chunk.indexOf(NEWLINE, lineStart); + if (newlineIndex === -1) { + append(chunk.subarray(lineStart)); + break; + } + // Most lines fit in the current chunk. Avoid buffering/streaming + // machinery on this hot path. + if (!streaming && pendingBytes === 0) { + parseLine(toLineString(chunk.subarray(lineStart, newlineIndex)), codexState, records); + } else { + append(chunk.subarray(lineStart, newlineIndex)); + finish(codexState, records); + } lineStart = newlineIndex + 1; + resumeOffset = scanOffset + lineStart; } - 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 - // consumed: a writer may still be appending to it, and counting a half - // record now and its full form later would double count. const tailRecords: UsageRecord[] = []; - if (pendingChunks.length > 0) { - const pending = pendingChunks.length === 1 ? pendingChunks[0]! : Buffer.concat(pendingChunks); - if (pending.length > 0) parseLine(toLineString(pending), { ...codexState }, tailRecords); - } + finish({ ...codexState }, tailRecords); const guardLength = Math.min(GUARD_LENGTH, resumeOffset); let guardHash = 0; diff --git a/apps/server/src/usage/usageTranscriptStreaming.test.ts b/apps/server/src/usage/usageTranscriptStreaming.test.ts new file mode 100644 index 000000000000..2187fd4aa431 --- /dev/null +++ b/apps/server/src/usage/usageTranscriptStreaming.test.ts @@ -0,0 +1,323 @@ +// @effect-diagnostics nodeBuiltinImport:off - exercise the real byte reader. +import * as NodeFSP from "node:fs/promises"; +import * as NodeOS from "node:os"; +import * as NodePath from "node:path"; + +import { afterEach, beforeEach, describe, expect, it } from "@effect/vitest"; + +import { + readTranscriptRecords as readWithDefaultThreshold, + type TranscriptParsePosition, +} from "./usageTranscriptReader.ts"; + +// Exercise the same transition with compact fixtures. UsageService tests and +// external 65/517 MiB fixtures also exercise the production threshold. +const readTranscriptRecords = ( + path: string, + provider: "claude" | "codex" | "grok", + position?: TranscriptParsePosition, +) => readWithDefaultThreshold(path, provider, position, { streamingThresholdBytes: 256 * 1024 }); + +let dir: string; +beforeEach(async () => { + dir = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "usage-stream-")); +}); +afterEach(async () => { + await NodeFSP.rm(dir, { recursive: true, force: true }); +}); + +const timestamp = "2026-08-01T10:00:00Z"; +const content = '工具 output \\" usage token_count '.repeat(40_000); +const claude = (id = "m1", output = 99) => ({ + type: "assistant", + timestamp, + sessionId: "s1", + requestId: `r-${id}`, + costUSD: 0.25, + message: { + content: [{ type: "tool_use", input: { text: content } }], + id, + model: "claude-fable-5", + usage: { + input_tokens: 100, + output_tokens: output, + cache_read_input_tokens: 20, + cache_creation_input_tokens: 5, + speed: "fast", + }, + }, +}); +const codex = [ + { type: "session_meta", timestamp, payload: { id: "s1" } }, + { type: "turn_context", timestamp, payload: { model: "gpt-5.6-sol" } }, + { + type: "event_msg", + timestamp, + payload: { + type: "token_count", + info: { + last_token_usage: { + input_tokens: 100, + output_tokens: 99, + cached_input_tokens: 20, + cache_write_input_tokens: 5, + reasoning_output_tokens: 10, + }, + }, + }, + }, +]; +const grok = { + timestamp: 1785578400, + params: { + sessionId: "s1", + _meta: { agentTimestampMs: 1785578400123 }, + update: { + sessionUpdate: "turn_completed", + prompt_id: "p1", + usage: { + inputTokens: 100, + outputTokens: 99, + costUsdTicks: 2500000000, + modelUsage: { + "grok-4.5-build": { + inputTokens: 100, + outputTokens: 99, + cachedReadTokens: 20, + cacheCreationTokens: 5, + reasoningTokens: 10, + }, + }, + }, + }, + }, +}; + +async function scan( + lines: readonly unknown[], + provider: "claude" | "codex" | "grok", + name = "history", +) { + const path = NodePath.join(dir, `${name}.jsonl`); + await NodeFSP.writeFile(path, lines.map((line) => JSON.stringify(line)).join("\n") + "\n"); + const result = await readTranscriptRecords(path, provider); + expect(result).not.toBeNull(); + return result!; +} + +describe("large usage records", () => { + it("keeps usage after large Claude tool input, including fast-mode cost and dedupe metadata", async () => { + const result = await scan([claude()], "claude"); + expect(result.records).toEqual([ + { + provider: "claude", + timestampMs: Date.parse(timestamp), + sessionId: "s1", + model: "claude-fable-5", + totals: { + uncachedInputTokens: 100, + outputTokens: 99, + cachedInputTokens: 20, + cacheCreationTokens: 5, + reasoningTokens: 0, + }, + reportedCostUsd: 0.25, + fast: true, + dedupeKey: "m1:r-m1", + }, + ]); + }); + + it.each(["claude", "codex", "grok"] as const)( + "matches ordinary %s records with large irrelevant fields in either order", + async (provider) => { + const small = + provider === "codex" + ? codex + : provider === "grok" + ? [grok] + : [{ ...claude(), message: { ...claude().message, content: [] } }]; + const expected = await scan(small, provider, "small"); + for (const first of [true, false]) { + const large = small.map((record) => + first ? { padding: content, ...record } : { ...record, padding: content }, + ); + const actual = await scan(large, provider); + expect(actual.records).toEqual(expected.records); + expect(actual.position.codexState).toEqual(expected.position.codexState); + } + }, + ); + + it("preserves Grok model allocation, precise timestamp and prompt identity", async () => { + const result = await scan([{ padding: content, ...grok }], "grok"); + expect(result.records[0]).toMatchObject({ + timestampMs: 1785578400123, + model: "grok-4.5-build", + sessionId: "s1", + totals: { + uncachedInputTokens: 75, + cachedInputTokens: 20, + cacheCreationTokens: 5, + outputTokens: 99, + reasoningTokens: 10, + }, + dedupeKey: "s1:p1:grok-4.5-build", + reportedCostUsd: 0.25, + }); + }); + + it("does not count usage-looking text or nested objects inside tool output", async () => { + const result = await scan( + [ + { type: "user", padding: content, toolOutput: claude() }, + { padding: content, message: { content: JSON.stringify(claude()) } }, + { type: "assistant", timestamp, message: { model: "claude-fable-5", content: [claude()] } }, + ], + "claude", + ); + expect(result.records).toEqual([]); + }); + + it("replays partial UTF-8 and JSON tails, commits exact CRLF offsets, and replaces the tail once", async () => { + const path = NodePath.join(dir, "tail.jsonl"); + const first = JSON.stringify(claude("first", 5)) + "\r\n"; + const second = Buffer.from(JSON.stringify(claude("second", 7))); + const split = second.lastIndexOf(Buffer.from("工具")) + 1; + await NodeFSP.writeFile(path, Buffer.concat([Buffer.from(first), second.subarray(0, split)])); + const partial = await readTranscriptRecords(path, "claude"); + expect(partial?.records.map((r) => r.totals.outputTokens)).toEqual([5]); + expect(partial?.tailRecords).toEqual([]); + expect(partial?.position.resumeOffset).toBe(Buffer.byteLength(first)); + await NodeFSP.appendFile(path, second.subarray(split)); + const complete = await readTranscriptRecords(path, "claude", partial!.position); + expect(complete?.resumed).toBe(true); + expect(complete?.records).toEqual([]); + expect(complete?.tailRecords.map((r) => r.totals.outputTokens)).toEqual([7]); + expect(complete?.position.resumeOffset).toBe(Buffer.byteLength(first)); + await NodeFSP.appendFile(path, "\r\n" + JSON.stringify(claude("third", 11)) + "\n"); + const appended = await readTranscriptRecords(path, "claude", complete!.position); + expect(appended?.records.map((r) => r.totals.outputTokens)).toEqual([7, 11]); + expect(appended?.tailRecords).toEqual([]); + expect(appended?.position.resumeOffset).toBe((await NodeFSP.stat(path)).size); + const full = await readTranscriptRecords(path, "claude"); + expect([...partial!.records, ...appended!.records]).toEqual(full?.records); + }); + + it("preserves Codex model switches and duplicate suppression across streaming resumes", async () => { + const path = NodePath.join(dir, "history.jsonl"); + const first = await scan( + codex.map((record) => ({ padding: content, ...record })), + "codex", + ); + await NodeFSP.appendFile( + path, + [ + { padding: content, ...codex[2] }, + { type: "turn_context", payload: { padding: content, model: "gpt-6" } }, + { + type: "event_msg", + timestamp, + payload: { + type: "token_count", + padding: content, + info: { last_token_usage: { input_tokens: 200, output_tokens: 101 } }, + }, + }, + ] + .map((line) => JSON.stringify(line)) + .join("\n") + "\n", + ); + const result = await readTranscriptRecords(path, "codex", first.position); + expect(result?.resumed).toBe(true); + expect(result?.records).toHaveLength(1); + expect(result?.records[0]).toMatchObject({ + model: "gpt-6", + sessionId: "s1", + totals: { outputTokens: 101 }, + }); + const full = await readTranscriptRecords(path, "codex"); + expect([...first.records, ...result!.records]).toEqual(full?.records); + }); + + it.each(["forked_from_id", "source"])( + "preserves Codex fork-copy suppression from %s", + async (field) => { + const fork = + field === "source" + ? { source: { subagent: { thread_spawn: { parent_thread_id: "parent" } } } } + : { forked_from_id: "parent" }; + const lines = [ + { type: "session_meta", timestamp, payload: { padding: content, id: "child", ...fork } }, + codex[1], + codex[2], + { + ...codex[2], + timestamp: "2026-08-01T10:00:05Z", + payload: { + type: "token_count", + info: { last_token_usage: { input_tokens: 100, output_tokens: 101 } }, + }, + }, + ]; + const result = await scan(lines, "codex"); + expect(result.records).toHaveLength(1); + expect(result.records[0]).toMatchObject({ + sessionId: "child", + totals: { outputTokens: 101 }, + }); + }, + ); + + it("rejects malformed large lines without accepting partial usage or losing following lines", async () => { + const valid = JSON.stringify(claude()); + const path = NodePath.join(dir, "broken.jsonl"); + await NodeFSP.writeFile( + path, + [valid + " junk", valid.slice(0, -1), valid + valid, JSON.stringify(claude("good", 7))].join( + "\n", + ) + "\n", + ); + const result = await readTranscriptRecords(path, "claude"); + expect(result?.records.map((r) => r.totals.outputTokens)).toEqual([7]); + }); + + it("matches JSON.parse for reordered and escaped keys, duplicate fields, arrays and unknown key names", async () => { + const path = NodePath.join(dir, "odd.jsonl"); + const small = JSON.stringify({ ...claude(), message: { ...claude().message, content: [] } }); + const variants = [ + small.replace('"message":', '"mess\\u0061ge":'), + small.replace('"costUSD":0.25', '"costUSD":5,"costUSD":0.25'), + small.replace('"type":"assistant"', '"type":"user","type":"assistant"'), + small.replace('"message":', '"__proto__":{"type":"user"},"message":'), + small.replace('"message":', '"message.usage":{"output_tokens":1234},"message":'), + ]; + for (const line of variants) { + await NodeFSP.writeFile(path, line + "\n"); + const expected = await readTranscriptRecords(path, "claude"); + await NodeFSP.writeFile( + path, + '{"padding":' + JSON.stringify(content) + "," + line.slice(1) + "\n", + ); + const actual = await readTranscriptRecords(path, "claude"); + expect(actual?.records).toEqual(expected?.records); + } + }); + it("reads usage after deeply nested discarded tool content", async () => { + const path = NodePath.join(dir, "deep.jsonl"); + const record = JSON.stringify(claude()); + const nested = "[".repeat(300) + "0" + "]".repeat(300); + await NodeFSP.writeFile(path, '{"toolOutput":' + nested + "," + record.slice(1) + "\n"); + const result = await readTranscriptRecords(path, "claude"); + expect(result?.records[0]?.totals.outputTokens).toBe(99); + }); + + it("keeps JSON number semantics for non-finite values without serializing them into null", async () => { + const path = NodePath.join(dir, "numbers.jsonl"); + const record = JSON.stringify(claude()).replace('"costUSD":0.25', '"costUSD":1e400'); + await NodeFSP.writeFile(path, record + "\n"); + const result = await readTranscriptRecords(path, "claude"); + expect(result?.records[0]?.reportedCostUsd).toBeNull(); + expect(result?.records[0]?.totals.outputTokens).toBe(99); + }); +}); diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index 6e01c2c5a8ed..c3903ef47194 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -113,6 +113,10 @@ export function parseClaudeLine(line: string): UsageRecord | null { } catch { return null; } + return parseClaudeRecord(parsed); +} + +export function parseClaudeRecord(parsed: unknown): UsageRecord | null { if (typeof parsed !== "object" || parsed === null) return null; const record = parsed as Record; @@ -228,6 +232,10 @@ export function parseCodexLine(line: string, state: CodexScanState): UsageRecord } catch { return null; } + return parseCodexRecord(parsed, state); +} + +export function parseCodexRecord(parsed: unknown, state: CodexScanState): UsageRecord | null { if (typeof parsed !== "object" || parsed === null) return null; const record = parsed as Record; @@ -385,6 +393,10 @@ export function parseGrokLine(line: string): readonly UsageRecord[] { } catch { return []; } + return parseGrokRecord(parsed); +} + +export function parseGrokRecord(parsed: unknown): readonly UsageRecord[] { if (typeof parsed !== "object" || parsed === null) return []; const record = parsed as Record;