diff --git a/apps/server/src/usage/UsageService.test.ts b/apps/server/src/usage/UsageService.test.ts index 226b81b73d11..c3d2e808b28a 100644 --- a/apps/server/src/usage/UsageService.test.ts +++ b/apps/server/src/usage/UsageService.test.ts @@ -102,6 +102,87 @@ function totalOutputTokens(summary: { buckets: readonly { totals: { outputTokens } describe("UsageService", () => { + const usage = (input: number, output: number) => ({ + input_tokens: input, + output_tokens: output, + total_tokens: input + output, + }); + const line = (last: unknown, total: unknown) => + JSON.stringify({ + type: "event_msg", + timestamp: "2026-08-01T10:00:05Z", + payload: { + type: "token_count", + info: { last_token_usage: last, total_token_usage: total }, + }, + }) + "\n"; + const model = + JSON.stringify({ + type: "turn_context", + payload: { model: "gpt-5.4" }, + }) + "\n"; + + it.live("reports inconsistent Codex counters through warm scans, restart, and recovery", () => + Effect.gen(function* () { + const { settings, home } = yield* setup; + const dir = NodePath.join(home, "codex", "sessions"); + const path = NodePath.join(dir, "rollout.jsonl"); + yield* Effect.promise(() => NodeFSP.mkdir(dir, { recursive: true })); + yield* Effect.promise(() => + NodeFSP.writeFile( + path, + model + + line(usage(100, 20), usage(100, 20)) + + line(usage(20, 5), usage(100, 20)) + + line(usage(10, 2), usage(130, 28)), + ), + ); + yield* Effect.gen(function* () { + const service = yield* UsageService.make; + for (let scan = 0; scan < 2; scan++) { + const result = yield* service.readSummary(WINDOW); + assert.strictEqual(totalOutputTokens(result), 20); + const source = result.sources.find((source) => source.fingerprint.provider === "codex"); + assert.strictEqual(source?.status, "partial"); + assert.strictEqual(source?.malformedRecords, 1); + assert.include(source?.message ?? "", "Inconsistent usage records: 1"); + } + const restarted = yield* UsageService.make; + // A warm cache restores diagnostics without touching the transcript. + const open = vi + .spyOn(NodeFSP, "open") + .mockRejectedValue(new Error("unexpected transcript read")); + const restored = yield* restarted + .readSummary(WINDOW) + .pipe(Effect.ensuring(Effect.sync(() => open.mockRestore()))); + assert.strictEqual(totalOutputTokens(restored), 20); + assert.strictEqual( + restored.sources.find((source) => source.fingerprint.provider === "codex") + ?.malformedRecords, + 1, + ); + yield* Effect.promise(() => NodeFSP.appendFile(path, line(usage(10, 2), usage(140, 30)))); + const appended = yield* restarted.readSummary(WINDOW); + assert.strictEqual(totalOutputTokens(appended), 22); + assert.strictEqual( + appended.sources.find((source) => source.fingerprint.provider === "codex") + ?.malformedRecords, + 1, + ); + yield* Effect.promise(() => + NodeFSP.writeFile(path, model + line(usage(100, 20), usage(100, 20))), + ); + const repaired = yield* restarted.readSummary(WINDOW); + const source = repaired.sources.find((source) => source.fingerprint.provider === "codex"); + assert.strictEqual(totalOutputTokens(repaired), 20); + assert.strictEqual(source?.status, "ok"); + assert.strictEqual(source?.malformedRecords, 0); + }).pipe( + Effect.provide(serviceLayers({ prefix: "usage-counter-reconciliation", home, settings })), + ); + }).pipe(Effect.scoped), + ); + it.live("reprices unchanged transcripts when custom prices are added, edited, or removed", () => Effect.gen(function* () { const { transcript, settings, home } = yield* setup; diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index b84adc580193..14015b89e061 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -317,7 +317,10 @@ export const make = Effect.gen(function* () { size: number, mtimeMs: number, provider: UsageProviderKind, - ): Effect.Effect => + ): Effect.Effect<{ + readonly records: readonly UsageRecord[]; + readonly malformedRecords: number; + } | null> => Effect.gen(function* () { const cached = fileCache.get(filePath); // Provider is part of the identity: if both providers were ever pointed @@ -328,9 +331,13 @@ export const make = Effect.gen(function* () { cached.mtimeMs === mtimeMs && cached.provider === provider ) { - return cached.tailRecords.length === 0 - ? cached.records - : [...cached.records, ...cached.tailRecords]; + return { + records: + cached.tailRecords.length === 0 + ? cached.records + : [...cached.records, ...cached.tailRecords], + malformedRecords: cached.malformedRecords, + }; } // Only a strictly grown file may resume. Same size with a new mtime, or @@ -362,10 +369,14 @@ export const make = Effect.gen(function* () { provider, records, tailRecords, + malformedRecords: parsed.malformedRecords, position: parsed.position, }); cacheDirty = true; - return tailRecords.length === 0 ? records : [...records, ...tailRecords]; + return { + records: tailRecords.length === 0 ? records : [...records, ...tailRecords], + malformedRecords: parsed.malformedRecords, + }; }); /** One provider directory's walk and parse, before rates are involved. */ @@ -375,6 +386,7 @@ export const make = Effect.gen(function* () { readonly volumeId: string; readonly status: "ok" | "partial" | "missing" | "failed"; readonly failedEntries: number; + readonly malformedRecords: number; readonly files: readonly { readonly path: string; readonly records: readonly UsageRecord[] | null; @@ -396,10 +408,12 @@ export const make = Effect.gen(function* () { ); const parsedFiles: { path: string; records: readonly UsageRecord[] | null }[] = []; let failedEntries = listing.failedEntries; + let malformedRecords = 0; for (const file of listing.files) { - const records = yield* readFileRecords(file.path, file.size, file.mtimeMs, provider); - if (records === null) failedEntries += 1; - parsedFiles.push({ path: file.path, records }); + const parsed = yield* readFileRecords(file.path, file.size, file.mtimeMs, provider); + if (parsed === null) failedEntries += 1; + malformedRecords += parsed?.malformedRecords ?? 0; + parsedFiles.push({ path: file.path, records: parsed?.records ?? null }); } scanned.push({ provider, @@ -407,6 +421,7 @@ export const make = Effect.gen(function* () { volumeId, files: parsedFiles, failedEntries, + malformedRecords, status: listing.status === "ok" && failedEntries > 0 ? "partial" : listing.status, }); } @@ -484,7 +499,15 @@ export const make = Effect.gen(function* () { const livePaths = new Set(); const walkedRoots: string[] = []; - for (const { provider, dir, volumeId, files, status, failedEntries } of scannedDirs) { + for (const { + provider, + dir, + volumeId, + files, + status, + failedEntries, + malformedRecords, + } of scannedDirs) { if (status === "missing" || status === "failed") { sources.push({ fingerprint: { hostId, provider, resolvedHomePath: dir, volumeId }, @@ -527,14 +550,14 @@ export const make = Effect.gen(function* () { sources.push({ fingerprint: { hostId, provider, resolvedHomePath: dir, volumeId }, - status, + status: malformedRecords > 0 ? "partial" : status, scannedFiles, skippedFiles, - malformedRecords: 0, + malformedRecords, distinctSessions: sessionIds.size, message: - status === "partial" - ? `Usage is incomplete: ${failedEntries} transcript files or directory entries could not be read.` + failedEntries > 0 || malformedRecords > 0 + ? `Usage is incomplete. Unreadable transcript files or directory entries: ${failedEntries}. Inconsistent usage records: ${malformedRecords}.` : null, }); } diff --git a/apps/server/src/usage/usageScanCache.test.ts b/apps/server/src/usage/usageScanCache.test.ts index fdb0aabafa40..8707cf69f865 100644 --- a/apps/server/src/usage/usageScanCache.test.ts +++ b/apps/server/src/usage/usageScanCache.test.ts @@ -8,7 +8,7 @@ import { type CachedFile, type ScanCache, } from "./usageScanCache.ts"; -import type { UsageRecord } from "./usageTranscripts.ts"; +import { initialCodexScanState, type UsageRecord } from "./usageTranscripts.ts"; function record(overrides: Partial = {}): UsageRecord { return { @@ -47,6 +47,7 @@ function cacheWith(entries: readonly [string, number, readonly UsageRecord[]][]) mtimeMs, provider: "claude", records, + malformedRecords: 0, tailRecords: [], position: position(), }); @@ -67,6 +68,7 @@ describe("scan cache round trip", () => { records: [ record({ provider: "grok", model: "grok-4.5-build", dedupeKey: "s:p:grok-4.5-build" }), ], + malformedRecords: 0, tailRecords: [record({ provider: "grok", model: "grok-4.5-build", dedupeKey: null })], position: position({ resumeOffset: 30, guardLength: 30, guardHash: 123 }), }); @@ -75,11 +77,21 @@ describe("scan cache round trip", () => { mtimeMs: 400, provider: "codex", records: [record({ provider: "codex", model: "gpt-5.2-codex", dedupeKey: null })], + malformedRecords: 2, tailRecords: [], position: position({ codexState: { model: "gpt-5.2-codex", sessionId: "session-c", + malformedRecords: 2, + lastCumulativeUsage: { + input_tokens: 10, + cached_input_tokens: 0, + cache_write_input_tokens: 0, + output_tokens: 5, + reasoning_output_tokens: 0, + total_tokens: 15, + }, lastUsageSignature: '{"input_tokens":1}', sawSessionMeta: true, suppressingForkCopies: false, @@ -111,6 +123,33 @@ describe("scan cache round trip", () => { expect(decodeScanCache(JSON.parse(JSON.stringify(poisoned))).has("/a.jsonl")).toBe(false); }); + it.each([ + { malformedRecords: -1 }, + { lastCumulativeUsage: { input_tokens: 10, output_tokens: 5 } }, + { + lastCumulativeUsage: { + input_tokens: 10, + output_tokens: 5, + cached_input_tokens: 0, + cache_write_input_tokens: 0, + reasoning_output_tokens: 6, + total_tokens: 15, + }, + }, + ])("discards corrupt counter state so the transcript is read again %#", (overrides) => { + const encoded = encodeScanCache(cacheWith([["/a.jsonl", 100, [record()]]])); + const poisoned = { + ...encoded, + files: { + "/a.jsonl": { + ...encoded.files["/a.jsonl"]!, + cs: { ...initialCodexScanState(), ...overrides }, + }, + }, + }; + expect(decodeScanCache(JSON.parse(JSON.stringify(poisoned))).size).toBe(0); + }); + it("drops an entry whose guard length is outside the supported range", () => { // The guard length sizes a Buffer in the reader; a bogus value would make // every parse of that file fail and silently drop its usage. @@ -125,7 +164,7 @@ describe("scan cache round trip", () => { it("rejects a document from the previous cache version", () => { const encoded = encodeScanCache(cacheWith([["/a.jsonl", 100, [record()]]])); - const previous = { ...encoded, version: 2 }; + const previous = { ...encoded, version: 3 }; expect(decodeScanCache(JSON.parse(JSON.stringify(previous))).size).toBe(0); }); diff --git a/apps/server/src/usage/usageScanCache.ts b/apps/server/src/usage/usageScanCache.ts index 224f109147e4..6f714183e338 100644 --- a/apps/server/src/usage/usageScanCache.ts +++ b/apps/server/src/usage/usageScanCache.ts @@ -20,13 +20,14 @@ import * as NodePath from "node:path"; import type { UsageProviderKind } from "@t3tools/contracts"; import { GUARD_LENGTH, type TranscriptParsePosition } from "./usageTranscriptReader.ts"; -import type { CodexScanState, UsageRecord } from "./usageTranscripts.ts"; +import { readCodexTokenUsage, type CodexScanState, type UsageRecord } from "./usageTranscripts.ts"; // v2: Codex fork-copy suppression changed what a file parses to, so v1 // entries would keep serving double-counted records forever. // v3: entries carry the parse position and reducer state so a grown file // re-parses only its appended bytes instead of starting over. -const USAGE_SCAN_CACHE_VERSION = 3 as const; +// v4: reconcile Codex cumulative counters and retain malformed-record counts. +const USAGE_SCAN_CACHE_VERSION = 4 as const; export interface CachedFile { readonly size: number; @@ -40,6 +41,7 @@ export interface CachedFile { * re-reads that segment and would otherwise double count it. */ readonly tailRecords: readonly UsageRecord[]; + readonly malformedRecords: number; readonly position: TranscriptParsePosition; } @@ -70,6 +72,7 @@ interface SerializedFile { readonly r: readonly SerializedRecord[]; /** Tail records; see `CachedFile.tailRecords`. */ readonly t: readonly SerializedRecord[]; + readonly mc: number; /** Parse position: resume offset, guard length, guard hash. */ readonly o: number; readonly gl: number; @@ -122,6 +125,7 @@ export function encodeScanCache(cache: ScanCache): SerializedCache { p: entry.provider, r: entry.records.map(serializeRecord), t: entry.tailRecords.map(serializeRecord), + mc: entry.malformedRecords, o: entry.position.resumeOffset, gl: entry.position.guardLength, gh: entry.position.guardHash, @@ -239,6 +243,7 @@ export function decodeScanCache(document: unknown): ScanCache { ) { continue; } + if (typeof entry.mc !== "number" || !Number.isSafeInteger(entry.mc) || entry.mc < 0) continue; const codexState = decodeCodexState(entry.cs); if (codexState === undefined) continue; @@ -253,6 +258,7 @@ export function decodeScanCache(document: unknown): ScanCache { provider, records, tailRecords, + malformedRecords: entry.mc, position: { resumeOffset: entry.o, guardLength: entry.gl, @@ -281,11 +287,25 @@ function decodeCodexState(value: unknown): CodexScanState | null | undefined { typeof state.sawSessionMeta !== "boolean" || typeof state.suppressingForkCopies !== "boolean" || typeof state.forkCopyAnchorMs !== "number" || - !Number.isFinite(state.forkCopyAnchorMs) + !Number.isFinite(state.forkCopyAnchorMs) || + typeof state.malformedRecords !== "number" || + !Number.isSafeInteger(state.malformedRecords) || + state.malformedRecords < 0 ) { return undefined; } + const cumulative = readCodexTokenUsage(state.lastCumulativeUsage); + if ( + state.lastCumulativeUsage !== null && + (cumulative === null || + Object.entries(cumulative).some( + ([key, count]) => (state.lastCumulativeUsage as Record)[key] !== count, + )) + ) + return undefined; return { + lastCumulativeUsage: cumulative, + malformedRecords: state.malformedRecords, model: state.model, sessionId: state.sessionId, lastUsageSignature: state.lastUsageSignature ?? null, diff --git a/apps/server/src/usage/usageTranscriptReader.test.ts b/apps/server/src/usage/usageTranscriptReader.test.ts index 0b36cb5db56f..c761a7f12f94 100644 --- a/apps/server/src/usage/usageTranscriptReader.test.ts +++ b/apps/server/src/usage/usageTranscriptReader.test.ts @@ -66,6 +66,80 @@ function codexUsageLine(outputTokens: number, secondsOffset: number): string { } describe("readTranscriptRecords resume", () => { + it("keeps cumulative counters and malformed counts correct across cache and tentative tails", async () => { + const { encodeScanCache, decodeScanCache } = await import("./usageScanCache.ts"); + const path = NodePath.join(dir, "counter-rollout.jsonl"); + const usage = (input: number, output: number) => ({ + input_tokens: input, + output_tokens: output, + total_tokens: input + output, + }); + const line = (last: unknown, total: unknown) => + JSON.stringify({ + type: "event_msg", + timestamp: "2026-08-01T10:00:05Z", + payload: { + type: "token_count", + info: { last_token_usage: last, total_token_usage: total }, + }, + }); + const firstLine = line(usage(100, 20), usage(100, 20)); + const stale = line(usage(20, 5), usage(100, 20)); + const bad = line(usage(10, 2), usage(130, 28)); + await NodeFSP.writeFile( + path, + codexModelLine("gpt-5.4") + firstLine + "\n" + stale + "\n" + bad, + ); + const first = await readTranscriptRecords(path, "codex"); + assert.isNotNull(first); + assert.strictEqual(first.records.length, 1); + assert.strictEqual(first.malformedRecords, 1); + assert.strictEqual(first.position.codexState?.malformedRecords, 0); + const stats = await NodeFSP.stat(path); + const cache = decodeScanCache( + JSON.parse( + JSON.stringify( + encodeScanCache( + new Map([ + [ + path, + { + ...first, + provider: "codex", + size: stats.size, + mtimeMs: stats.mtimeMs, + }, + ], + ]), + ), + ), + ), + ); + const cached = cache.get(path); + assert.isDefined(cached); + assert.strictEqual(cached.malformedRecords, 1); + const good = line(usage(10, 2), usage(140, 30)); + await NodeFSP.appendFile(path, "\n" + good); + const resumed = await readTranscriptRecords(path, "codex", cached.position); + assert.isNotNull(resumed); + assert.isTrue(resumed.resumed); + assert.strictEqual(resumed.malformedRecords, 1); + assert.strictEqual(resumed.position.codexState?.malformedRecords, 1); + assert.strictEqual(resumed.tailRecords[0]?.totals.outputTokens, 2); + assert.strictEqual(resumed.position.codexState?.lastCumulativeUsage?.output_tokens, 28); + await NodeFSP.appendFile(path, "\n"); + const completed = await readTranscriptRecords(path, "codex", resumed.position); + const full = await readTranscriptRecords(path, "codex"); + assert.isNotNull(completed); + assert.isNotNull(full); + assert.deepStrictEqual( + [...cached.records, ...resumed.records, ...completed.records], + full.records, + ); + assert.strictEqual(completed.malformedRecords, full.malformedRecords); + assert.strictEqual(completed.malformedRecords, 1); + }); + it("parses only appended lines when resuming a grown file", async () => { const path = NodePath.join(dir, "claude.jsonl"); await NodeFSP.writeFile(path, claudeLine(1, 5) + claudeLine(2, 7)); diff --git a/apps/server/src/usage/usageTranscriptReader.ts b/apps/server/src/usage/usageTranscriptReader.ts index 2e64f5c7251f..f1c414cb16d5 100644 --- a/apps/server/src/usage/usageTranscriptReader.ts +++ b/apps/server/src/usage/usageTranscriptReader.ts @@ -67,6 +67,8 @@ export interface TranscriptParseResult { * segment: the next scan re-reads it once the writer finishes the line. */ readonly tailRecords: readonly UsageRecord[]; + /** Invalid Codex counter events, including a complete tentative tail. */ + readonly malformedRecords: number; readonly position: TranscriptParsePosition; /** Whether the parse continued from `resumeFrom` rather than byte 0. */ readonly resumed: boolean; @@ -296,9 +298,10 @@ export async function readTranscriptRecords( // 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[] = []; + const tailState = { ...codexState }; if (pendingChunks.length > 0) { const pending = pendingChunks.length === 1 ? pendingChunks[0]! : Buffer.concat(pendingChunks); - if (pending.length > 0) parseLine(toLineString(pending), { ...codexState }, tailRecords); + if (pending.length > 0) parseLine(toLineString(pending), tailState, tailRecords); } const guardLength = Math.min(GUARD_LENGTH, resumeOffset); @@ -312,6 +315,7 @@ export async function readTranscriptRecords( return { records, tailRecords, + malformedRecords: tailState.malformedRecords, position: { resumeOffset, guardLength, diff --git a/apps/server/src/usage/usageTranscripts.test.ts b/apps/server/src/usage/usageTranscripts.test.ts index b09db613ed85..f5d3cb4d2a3f 100644 --- a/apps/server/src/usage/usageTranscripts.test.ts +++ b/apps/server/src/usage/usageTranscripts.test.ts @@ -137,6 +137,183 @@ describe("parseCodexLine", () => { expect(parseCodexLine(tokenCount(100, 0, 10, 0), state)).not.toBeNull(); }); + const counters = (input: number, output: number) => ({ + input_tokens: input, + cached_input_tokens: 0, + cache_write_input_tokens: 0, + output_tokens: output, + reasoning_output_tokens: 0, + total_tokens: input + output, + }); + const counted = (last: unknown, total: unknown, extra = {}) => + JSON.stringify({ + type: "event_msg", + timestamp: "2026-08-01T05:17:49.919Z", + payload: { + type: "token_count", + info: { last_token_usage: last, total_token_usage: total, ...extra }, + }, + }); + const ready = () => { + const state = initialCodexScanState(); + parseCodexLine(turnContext, state); + return state; + }; + + it("ignores changed request counters when cumulative usage stays unchanged", () => { + const state = ready(); + expect(parseCodexLine(counted(counters(100, 20), counters(100, 20)), state)).not.toBeNull(); + expect(parseCodexLine(counted(counters(20, 5), counters(100, 20)), state)).toBeNull(); + expect(state.malformedRecords).toBe(0); + }); + + it("counts identical requests when cumulative counters advance, including across models", () => { + const state = ready(); + parseCodexLine(counted(counters(100, 20), counters(100, 20)), state); + parseCodexLine(JSON.stringify({ type: "turn_context", payload: { model: "gpt-5.4" } }), state); + const next = parseCodexLine(counted(counters(100, 20), counters(200, 40)), state); + expect(next?.totals.outputTokens).toBe(20); + expect(next?.model).toBe("gpt-5.4"); + expect(state.malformedRecords).toBe(0); + }); + + it("retains cumulative progress across a duplicate legacy event", () => { + const state = ready(); + parseCodexLine(counted(counters(100, 20), counters(100, 20)), state); + expect(parseCodexLine(tokenCount(100, 0, 20, 0), state)).toBeNull(); + expect( + parseCodexLine(counted(counters(100, 20), counters(200, 40)), state)?.totals.outputTokens, + ).toBe(20); + expect(state.malformedRecords).toBe(0); + }); + + it.each([false, true])( + "reconciles counted legacy requests without billing them again, initial cumulative=%s", + (initialCumulative) => { + const state = ready(); + parseCodexLine( + initialCumulative + ? counted(counters(100, 20), counters(100, 20)) + : tokenCount(100, 0, 20, 0), + state, + ); + expect(parseCodexLine(tokenCount(20, 0, 5, 0), state)?.totals.outputTokens).toBe(5); + // The same request is later re-emitted with cumulative counters. + expect(parseCodexLine(counted(counters(20, 5), counters(120, 25)), state)).toBeNull(); + expect( + parseCodexLine(counted(counters(10, 2), counters(130, 27)), state)?.totals.outputTokens, + ).toBe(2); + expect(state.malformedRecords).toBe(0); + }, + ); + + it.each([counters(130, 28), counters(50, 10)])( + "reports inconsistent deltas and recovers from the next valid counter pair %#", + (total) => { + const state = ready(); + parseCodexLine(counted(counters(100, 20), counters(100, 20)), state); + expect(parseCodexLine(counted(counters(20, 5), total), state)).toBeNull(); + expect(state.malformedRecords).toBe(1); + const next = parseCodexLine( + counted(counters(10, 2), counters(total.input_tokens + 10, total.output_tokens + 2)), + state, + ); + expect(next?.totals.outputTokens).toBe(2); + expect(state.malformedRecords).toBe(1); + }, + ); + + it.each([ + { input_tokens: -1 }, + { input_tokens: 1.5 }, + { cached_input_tokens: null }, + { input_tokens: "100" }, + { cached_input_tokens: 101 }, + { reasoning_output_tokens: 21 }, + { total_tokens: 121 }, + { output_tokens: Number.MAX_SAFE_INTEGER + 1 }, + ])("rejects malformed request counters without clamping them %#", (overrides) => { + const state = ready(); + expect( + parseCodexLine(counted({ ...counters(100, 20), ...overrides }, counters(100, 20)), state), + ).toBeNull(); + expect(state.malformedRecords).toBe(1); + }); + + it("reports a malformed cumulative counter rather than falling back to request usage", () => { + const state = ready(); + expect( + parseCodexLine( + counted(counters(100, 20), { ...counters(100, 20), input_tokens: "100" }), + state, + ), + ).toBeNull(); + expect(state.malformedRecords).toBe(1); + expect(parseCodexLine(counted(counters(100, 20), counters(100, 20)), state)).not.toBeNull(); + }); + + it("keeps reasoning inside output while reconciling every component", () => { + const state = ready(); + const usage = { + ...counters(100, 20), + cached_input_tokens: 30, + cache_write_input_tokens: 10, + reasoning_output_tokens: 5, + }; + expect(parseCodexLine(counted(usage, usage), state)?.totals).toEqual({ + uncachedInputTokens: 60, + cachedInputTokens: 30, + cacheCreationTokens: 10, + outputTokens: 20, + reasoningTokens: 5, + }); + const total = { + ...counters(200, 40), + cached_input_tokens: 60, + cache_write_input_tokens: 20, + reasoning_output_tokens: 9, + }; + expect(parseCodexLine(counted(usage, total), state)).toBeNull(); + expect(state.malformedRecords).toBe(1); + }); + + it("ignores context-window estimates and reconciles the following real request", () => { + const state = ready(); + parseCodexLine(counted(counters(100, 20), counters(100, 20)), state); + const estimate = { ...counters(0, 0), total_tokens: 1000 }; + expect( + parseCodexLine( + counted({ ...estimate, total_tokens: 880 }, estimate, { model_context_window: 1000 }), + state, + ), + ).toBeNull(); + expect( + parseCodexLine(counted(counters(10, 2), { ...counters(10, 2), total_tokens: 1012 }), state) + ?.totals.outputTokens, + ).toBe(2); + expect(state.malformedRecords).toBe(0); + }); + + it("uses suppressed fork history as the baseline for real child usage", () => { + const state = initialCodexScanState(); + parseCodexLine( + JSON.stringify({ + type: "session_meta", + timestamp: "2026-08-01T05:17:49.900Z", + payload: { id: "child", forked_from_id: "parent" }, + }), + state, + ); + parseCodexLine(turnContext, state); + expect(parseCodexLine(counted(counters(100, 20), counters(100, 20)), state)).toBeNull(); + const child = JSON.parse(counted(counters(100, 20), counters(200, 40))); + child.timestamp = "2026-08-01T05:17:55.000Z"; + const record = parseCodexLine(JSON.stringify(child), state); + expect(record?.sessionId).toBe("child"); + expect(record?.totals.outputTokens).toBe(20); + expect(state.malformedRecords).toBe(0); + }); + // A forked/subagent rollout opens with the parent's history copied in and // every line re-stamped to the fork instant, then the ancestors' session // metas. Counting those again multiplied usage ~1.85x on real data (#5758). diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index 5d909379eb10..ba8678e2ac0d 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -153,6 +153,54 @@ export function parseClaudeLine(line: string): UsageRecord | null { /* Codex */ /* -------------------------------------------------------------------------- */ +const CODEX_USAGE_FIELDS = [ + "input_tokens", + "cached_input_tokens", + "cache_write_input_tokens", + "output_tokens", + "reasoning_output_tokens", + "total_tokens", +] as const; + +type CodexTokenUsage = Readonly>; + +/** Decode counters without rounding, clamping, or hiding inconsistent subsets. */ +export function readCodexTokenUsage(value: unknown): CodexTokenUsage | null { + if (typeof value !== "object" || value === null || Array.isArray(value)) return null; + const raw = value as Record; + const input = raw["input_tokens"]; + const output = raw["output_tokens"]; + if (typeof input !== "number" || typeof output !== "number") return null; + const usage = { + input_tokens: input, + cached_input_tokens: raw["cached_input_tokens"] === undefined ? 0 : raw["cached_input_tokens"], + cache_write_input_tokens: + raw["cache_write_input_tokens"] === undefined ? 0 : raw["cache_write_input_tokens"], + output_tokens: output, + reasoning_output_tokens: + raw["reasoning_output_tokens"] === undefined ? 0 : raw["reasoning_output_tokens"], + total_tokens: raw["total_tokens"] === undefined ? input + output : raw["total_tokens"], + }; + for (const field of CODEX_USAGE_FIELDS) { + const count = usage[field]; + if (typeof count !== "number" || !Number.isSafeInteger(count) || count < 0) return null; + } + const counters = usage as CodexTokenUsage; + if ( + counters.cached_input_tokens + counters.cache_write_input_tokens > input || + counters.reasoning_output_tokens > output || + counters.total_tokens < input + output + ) + return null; + return counters; +} + +function advanceCodexUsage(previous: CodexTokenUsage | null, last: CodexTokenUsage) { + return readCodexTokenUsage( + Object.fromEntries(CODEX_USAGE_FIELDS.map((key) => [key, (previous?.[key] ?? 0) + last[key]])), + ); +} + /** * Rolling state for a single Codex rollout file. * @@ -164,6 +212,8 @@ export interface CodexScanState { model: string; sessionId: string; lastUsageSignature: string | null; + lastCumulativeUsage: CodexTokenUsage | null; + malformedRecords: number; sawSessionMeta: boolean; /** While true, leading usage events are re-stamped copies of parent history. */ suppressingForkCopies: boolean; @@ -175,6 +225,8 @@ export function initialCodexScanState(): CodexScanState { model: "", sessionId: "", lastUsageSignature: null, + lastCumulativeUsage: null, + malformedRecords: 0, sawSessionMeta: false, suppressingForkCopies: false, forkCopyAnchorMs: 0, @@ -206,9 +258,8 @@ function isForkedSessionMeta(payload: Record): boolean { * Feeds one line of a Codex rollout into `state`, returning a record when the * line was a usage event. * - * Deltas come from `last_token_usage`. Summing those across a session - * reconciles with the session's final `total_token_usage`, provided - * consecutive duplicate events are dropped, which this does. + * Cumulative counters validate each request delta. Legacy events without + * cumulative counters retain consecutive-payload deduplication. */ export function parseCodexLine(line: string, state: CodexScanState): UsageRecord | null { let parsed: unknown; @@ -250,48 +301,98 @@ export function parseCodexLine(line: string, state: CodexScanState): UsageRecord const info = payloadRecord["info"]; if (typeof info !== "object" || info === null) return null; - const last = (info as Record)["last_token_usage"]; - if (typeof last !== "object" || last === null) return null; - const lastRecord = last as Record; - - // Only an event that is otherwise eligible may consume the duplicate - // signature. A token_count arriving before its turn_context (no model yet) - // must not poison it, or the re-emitted copy after the model is known would - // be skipped as a duplicate and those tokens never counted. + const infoRecord = info as Record; + + // An ineligible event must not consume the counters: Codex can re-emit it + // after the model or timestamp becomes available. const timestampMs = parseTimestampMs(record["timestamp"]); - if (timestampMs === null) return null; - if (state.model.length === 0) return null; + if (timestampMs === null || state.model.length === 0) return null; - // Codex re-emits an unchanged token_count on some stream boundaries. Summing - // those would double count, so identical consecutive payloads are skipped. - const signature = JSON.stringify(lastRecord); - if (signature === state.lastUsageSignature) return null; - state.lastUsageSignature = signature; + const last = readCodexTokenUsage(infoRecord["last_token_usage"]); + const rawCumulative = infoRecord["total_token_usage"]; + const hasCumulative = rawCumulative !== undefined && rawCumulative !== null; + const cumulative = hasCumulative ? readCodexTokenUsage(rawCumulative) : null; - // In a forked rollout the copied parent history was already counted from the - // parent's own file. Drop the leading burst; the first usage event separated - // from its predecessor by a real turn's worth of time ends it for good. + // Copied parent history still establishes the child's counter baseline. if (state.suppressingForkCopies) { if (timestampMs - state.forkCopyAnchorMs < FORK_COPY_MAX_GAP_MS) { state.forkCopyAnchorMs = timestampMs; + const signature = last === null ? null : JSON.stringify(last); + if (cumulative !== null) state.lastCumulativeUsage = cumulative; + else if (last !== null && signature !== state.lastUsageSignature) { + state.lastCumulativeUsage = advanceCodexUsage(state.lastCumulativeUsage, last); + } + state.lastUsageSignature = signature; return null; } state.suppressingForkCopies = false; } - const inputTokens = int(lastRecord["input_tokens"]); - const cachedInputTokens = int(lastRecord["cached_input_tokens"]); - const cacheCreationTokens = int(lastRecord["cache_write_input_tokens"]); - const outputTokens = int(lastRecord["output_tokens"]); + if (last === null || (hasCumulative && cumulative === null)) { + state.malformedRecords += 1; + return null; + } + + const previous = state.lastCumulativeUsage; + // A structurally valid counter is the baseline for the next event even if + // this delta is inconsistent, so one bad event does not spoil the whole file. + if (cumulative !== null) state.lastCumulativeUsage = cumulative; + + // Codex fill_to_context_window replaces usage with zero component counters + // and a context estimate in total_tokens. It is not a billable request. + if ( + cumulative !== null && + last.input_tokens + last.output_tokens === 0 && + cumulative.input_tokens + cumulative.output_tokens === 0 && + cumulative.total_tokens === infoRecord["model_context_window"] + ) { + state.lastUsageSignature = null; + return null; + } + + if (last.total_tokens !== last.input_tokens + last.output_tokens) { + state.malformedRecords += 1; + return null; + } + + const signature = JSON.stringify(last); + if (cumulative !== null) { + if (previous !== null && CODEX_USAGE_FIELDS.every((key) => cumulative[key] === previous[key])) { + return null; + } + if ( + !CODEX_USAGE_FIELDS.every((key) => cumulative[key] - (previous?.[key] ?? 0) === last[key]) + ) { + state.malformedRecords += 1; + return null; + } + } else { + if (signature === state.lastUsageSignature) return null; + // Legacy usage already contributes to the records. Advance the expected + // total too, so a later cumulative snapshot neither drops the next delta + // nor bills this legacy request again. + const advanced = advanceCodexUsage(previous, last); + if (advanced === null) { + state.malformedRecords += 1; + return null; + } + state.lastCumulativeUsage = advanced; + } + state.lastUsageSignature = signature; + + const inputTokens = last.input_tokens; + const cachedInputTokens = last.cached_input_tokens; + const cacheCreationTokens = last.cache_write_input_tokens; + const outputTokens = last.output_tokens; const totals: UsageTokenTotals = { // Codex reports `input_tokens` inclusive of the cached portion. - uncachedInputTokens: Math.max(0, inputTokens - cachedInputTokens - cacheCreationTokens), + uncachedInputTokens: inputTokens - cachedInputTokens - cacheCreationTokens, cachedInputTokens, cacheCreationTokens, outputTokens, // Reported inside output_tokens, surfaced separately for the token mix. - reasoningTokens: Math.min(outputTokens, int(lastRecord["reasoning_output_tokens"])), + reasoningTokens: last.reasoning_output_tokens, }; if (totalTokens(totals) === 0) return null;