diff --git a/apps/server/src/project/AgentSessionJson.test.ts b/apps/server/src/project/AgentSessionJson.test.ts new file mode 100644 index 0000000000..515aa50b1b --- /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 d31c47833d..4deba0a677 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/provider/Layers/OpenCodeProvider.test.ts b/apps/server/src/provider/Layers/OpenCodeProvider.test.ts index b0fe26a006..71317b3257 100644 --- a/apps/server/src/provider/Layers/OpenCodeProvider.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeProvider.test.ts @@ -1,4 +1,5 @@ import * as NodeAssert from "node:assert/strict"; +import * as NodeCrypto from "node:crypto"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { it } from "@effect/vitest"; @@ -68,6 +69,10 @@ it.effect("reads Go limits with the instance's XDG credentials and preserves res Effect.provide(NodeServices.layer), ); NodeAssert.equal(limits.unavailable, undefined); + NodeAssert.equal( + limits.credentialFingerprint, + NodeCrypto.createHash("sha256").update("opencode-go\0instance-key").digest("hex"), + ); NodeAssert.deepEqual( limits.windows.map(({ kind, usedPercent, resetsAt: reset }) => ({ kind, @@ -136,6 +141,7 @@ it.effect("keeps Go entitlement absence distinct from failed or malformed usage Effect.provide(NodeServices.layer), ); NodeAssert.equal(limits.unavailable?.reason, reason); + NodeAssert.equal(limits.credentialFingerprint, undefined); NodeAssert.deepEqual(limits.windows, []); } }), diff --git a/apps/server/src/provider/Layers/openCodeUsageLimits.ts b/apps/server/src/provider/Layers/openCodeUsageLimits.ts index 2cabef8791..55fedcdf52 100644 --- a/apps/server/src/provider/Layers/openCodeUsageLimits.ts +++ b/apps/server/src/provider/Layers/openCodeUsageLimits.ts @@ -1,4 +1,5 @@ import * as NodeOS from "node:os"; +import * as NodeCrypto from "node:crypto"; import type { ServerProviderUsageWindow } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; @@ -95,7 +96,16 @@ export const readOpenCodeGoUsageLimits = Effect.fn("readOpenCodeGoUsageLimits")( resetsAt: DateTime.formatIso(body.usage.monthly.resetsAt), }, ]; - return makeUsageLimits({ checkedAt, windows }); + return { + ...makeUsageLimits({ checkedAt, windows }), + // Go's usage response has no account ID. An unkeyed hash matches across + // environments without a shared secret. It permits offline guesses, but + // Go keys are randomly generated. + credentialFingerprint: NodeCrypto.createHash("sha256") + .update("opencode-go\0") + .update(apiKey) + .digest("hex"), + }; }).pipe( Effect.timeout("5 seconds"), Effect.orElseSucceed(() => diff --git a/apps/server/src/usage/UsageService.test.ts b/apps/server/src/usage/UsageService.test.ts index d95772d175..32d40956ba 100644 --- a/apps/server/src/usage/UsageService.test.ts +++ b/apps/server/src/usage/UsageService.test.ts @@ -514,6 +514,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 9e5ab6e0c9..faa6869907 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 0000000000..2187fd4aa4 --- /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 e6d7827985..199bcefba6 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -108,6 +108,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; @@ -223,6 +227,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; @@ -380,6 +388,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; diff --git a/packages/contracts/src/providerUsageLimits.ts b/packages/contracts/src/providerUsageLimits.ts index 2c4b67694b..4d5652950e 100644 --- a/packages/contracts/src/providerUsageLimits.ts +++ b/packages/contracts/src/providerUsageLimits.ts @@ -55,6 +55,8 @@ export const ServerProviderUsageLimits = Schema.Struct({ source: Schema.optional(TrimmedNonEmptyString), checkedAt: IsoDateTime, windows: ForwardCompatibleArray(ServerProviderUsageWindow), + /** Opaque credential identity when the provider does not report an account. */ + credentialFingerprint: Schema.optional(TrimmedNonEmptyString), resetCredits: Schema.optional(ServerProviderResetCredits), unavailable: Schema.optional( Schema.Struct({ diff --git a/packages/shared/src/usageLimits.test.ts b/packages/shared/src/usageLimits.test.ts index 4880df3840..583d12a59c 100644 --- a/packages/shared/src/usageLimits.test.ts +++ b/packages/shared/src/usageLimits.test.ts @@ -542,6 +542,74 @@ describe("pools", () => { expect(accounts[0]?.limits.windows[0]?.usedPercent).toBe(55); }); + it("merges OpenCode Go limits from machines with the same API key", () => { + const go = provider({ + driver: ProviderDriverKind.make("opencode"), + instanceId: ProviderInstanceId.make("opencode"), + auth: { status: "authenticated" }, + usageLimits: { + checkedAt, + credentialFingerprint: "shared-go-key", + windows: [{ ...window, id: "go_rolling", usedPercent: 3 }], + }, + }); + const input = new Map([ + [EnvironmentId.make("env-a"), { ...laptop, serverConfig: { providers: [go] } }], + [ + EnvironmentId.make("env-b"), + { + entry: { target: { label: "Desktop" } }, + serverConfig: { + providers: [ + { + ...go, + usageLimits: { + ...go.usageLimits!, + checkedAt: "2026-09-03T11:30:00.000Z", + windows: [{ ...window, id: "go_rolling", usedPercent: 4 }], + }, + }, + ], + }, + }, + ], + ]); + const accounts = collectLimitAccounts(input); + expect(accounts).toHaveLength(1); + expect(accounts[0]?.environments).toEqual([ + { environmentId: "env-a", label: "Laptop" }, + { environmentId: "env-b", label: "Desktop" }, + ]); + expect(collectLimitPools(accounts, now)[0]?.windows[0]?.members).toHaveLength(1); + expect(accounts[0]?.limits.windows[0]?.usedPercent).toBe(4); + + const differentKey = { + ...go, + usageLimits: { ...go.usageLimits!, credentialFingerprint: "other-go-key" }, + }; + input.set(EnvironmentId.make("env-b"), { + entry: { target: { label: "Desktop" } }, + serverConfig: { providers: [differentKey] }, + }); + expect(collectLimitAccounts(input)).toHaveLength(2); + + input.set(EnvironmentId.make("env-a"), { + ...laptop, + serverConfig: { + providers: [{ ...go, auth: { status: "authenticated", email: "same@example.com" } }], + }, + }); + input.set(EnvironmentId.make("env-b"), { + entry: { target: { label: "Desktop" } }, + serverConfig: { + providers: [ + { ...differentKey, auth: { status: "authenticated", email: "SAME@example.com" } }, + ], + }, + }); + expect(collectLimitAccounts(input)).toHaveLength(1); + }); + it("takes windows from a fresher hub read but credits and redeem from the native instance", () => { const native = provider({ driver: claude, diff --git a/packages/shared/src/usageLimits.ts b/packages/shared/src/usageLimits.ts index a02bf89704..e2395e8998 100644 --- a/packages/shared/src/usageLimits.ts +++ b/packages/shared/src/usageLimits.ts @@ -182,7 +182,7 @@ export function collectLimitSources( const nativeAccounts = new Set(); for (const presentation of presentations.values()) { for (const provider of providersWithLimits(presentation.serverConfig?.providers ?? [])) { - const key = accountKey(provider.driver, provider.auth.email); + const key = accountKey(provider.driver, provider.auth.email, provider.usageLimits); if ( key !== null && provider.usageLimits?.windows.length && @@ -210,7 +210,7 @@ export function collectLimitSources( return perEnvironment.flatMap(({ environmentId, environmentLabel, sources }) => sources.map((source) => { const accounts = source.accounts.filter((account) => { - const key = accountKey(account.driver, account.email); + const key = accountKey(account.driver, account.email, account.usageLimits); return key === null || !nativeAccounts.has(key); }); return { @@ -225,9 +225,17 @@ export function collectLimitSources( ); } -function accountKey(driver: ServerProvider["driver"], email: string | undefined): string | null { +/** Prefer the reported email; use an identical credential when no email is available. */ +function accountKey( + driver: ServerProvider["driver"], + email: string | undefined, + limits?: ServerProviderUsageLimits, +): string | null { const normalizedEmail = email?.trim().toLowerCase(); - return normalizedEmail ? JSON.stringify(["account", driver, normalizedEmail]) : null; + if (normalizedEmail) return JSON.stringify(["account", driver, normalizedEmail]); + return limits?.credentialFingerprint + ? JSON.stringify(["credential", driver, limits.credentialFingerprint]) + : null; } /** Window pushes carry their own time; they must not make an old credit balance look fresh. */ @@ -237,9 +245,9 @@ function creditReadAt(limits: ServerProviderUsageLimits): number { /** * One subscription account as the pooled views see it, whichever way it was - * reported. The same email signed in natively on two environments, or reported - * by a hub as well as natively, is one account: its quota is one bucket, so - * counting it twice would misstate what is left. + * reported. Matching emails or credentials across environments name + * one account. Its quota is one bucket, so counting it twice would misstate + * what is left. */ export interface LimitAccount { readonly key: string; @@ -338,7 +346,7 @@ export function collectLimitAccounts( for (const provider of providersWithLimits(presentation.serverConfig?.providers ?? [])) { if (!provider.usageLimits || limitsNotice(provider.usageLimits) !== null) continue; merge( - accountKey(provider.driver, provider.auth.email) ?? + accountKey(provider.driver, provider.auth.email, provider.usageLimits) ?? JSON.stringify(["native", environmentId, provider.instanceId]), { key: JSON.stringify(["native", environmentId, provider.instanceId]), @@ -367,7 +375,7 @@ export function collectLimitAccounts( for (const account of source.accounts) { if (limitsNotice(account.usageLimits) !== null) continue; merge( - accountKey(account.driver, account.email) ?? + accountKey(account.driver, account.email, account.usageLimits) ?? JSON.stringify(["hub", environmentId, source.id, account.id]), { key: JSON.stringify(["hub", environmentId, source.id, account.id]), @@ -765,7 +773,7 @@ export function collectProviderUsageLimits( ); const nativeAccounts = new Set( native.flatMap((provider) => { - const key = accountKey(provider.driver, provider.auth.email); + const key = accountKey(provider.driver, provider.auth.email, provider.usageLimits); return key && provider.usageLimits?.windows.length && !provider.usageLimits.unavailable ? [key] : []; @@ -775,13 +783,13 @@ export function collectProviderUsageLimits( const notices: string[] = []; for (const provider of native) { if (!provider.usageLimits) continue; - const key = accountKey(provider.driver, provider.auth.email); + const key = accountKey(provider.driver, provider.auth.email, provider.usageLimits); const hubCredits = sources .flatMap((source) => source.accounts.map((account) => ({ source, account }))) .filter( ({ account }) => key !== null && - accountKey(account.driver, account.email) === key && + accountKey(account.driver, account.email, account.usageLimits) === key && account.usageLimits.resetCredits && !limitsNotice(account.usageLimits), ) @@ -824,7 +832,7 @@ export function collectProviderUsageLimits( for (const source of sources) { const matching = source.accounts.filter((account) => account.driver === selected.driver); for (const account of matching) { - const key = accountKey(account.driver, account.email); + const key = accountKey(account.driver, account.email, account.usageLimits); if (key && nativeAccounts.has(key)) continue; accounts.push({ id: JSON.stringify(["hub", source.id, account.id]),