From 9902640a1d69d39d3fcd52903f5a18f29014f4eb Mon Sep 17 00:00:00 2001 From: miguel Date: Mon, 7 Sep 2026 16:49:20 -0700 Subject: [PATCH 1/2] Preserve SDK usage reporting and interrupted Codex usage --- .../claude-agent-sdk/src/session.ts | 40 ++-- .../claude-agent-sdk/tests/session.test.ts | 27 +++ .../integrations/codex-sdk/src/session.ts | 108 +++++++++- .../codex-sdk/tests/session.test.ts | 186 +++++++++++++++++- .../deepagents-sdk/src/session.ts | 12 ++ .../deepagents-sdk/tests/session.test.ts | 23 +++ .../deepagents/runner/run_eval.py | 11 +- .../deepagents/runner/tests/test_run_eval.py | 14 ++ packages/integrations/eve-sdk/src/session.ts | 13 ++ .../eve-sdk/tests/session.test.ts | 22 +++ packages/integrations/fx-sdk/src/session.ts | 21 +- .../integrations/fx-sdk/tests/session.test.ts | 45 +++++ .../integrations/mastra-sdk/src/session.ts | 15 +- .../mastra-sdk/tests/session.test.ts | 55 ++++++ packages/integrations/pi-sdk/src/session.ts | 12 ++ .../integrations/pi-sdk/tests/session.test.ts | 27 +++ 16 files changed, 608 insertions(+), 23 deletions(-) diff --git a/packages/integrations/claude-agent-sdk/src/session.ts b/packages/integrations/claude-agent-sdk/src/session.ts index c1e2effbe..ad8a37485 100644 --- a/packages/integrations/claude-agent-sdk/src/session.ts +++ b/packages/integrations/claude-agent-sdk/src/session.ts @@ -29,6 +29,8 @@ export type ClaudeSessionConfig = { }; export type ClaudeCodeTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; inputTokens: number; outputTokens: number; cacheCreationInputTokens: number; @@ -241,7 +243,23 @@ export function extractClaudeCodeTokenUsage( const cacheReadInputTokens = readNumber(usage, "cache_read_input_tokens") ?? sumModelUsage(resultMessage, "cacheReadInputTokens"); + const reported = + [ + "input_tokens", + "output_tokens", + "cache_creation_input_tokens", + "cache_read_input_tokens", + ].some((key) => isTokenCount(usage?.[key])) || + (isRecord(resultMessage?.modelUsage) && + Object.values(resultMessage.modelUsage).some( + (value) => + isRecord(value) && + ["inputTokens", "outputTokens", "cacheCreationInputTokens", "cacheReadInputTokens"].some( + (key) => isTokenCount(value[key]), + ), + )); return { + reported, inputTokens, outputTokens, cacheCreationInputTokens, @@ -260,18 +278,8 @@ function sumModelUsage(resultMessage: ClaudeSdkMessage | undefined, key: string) } function readNumber(record: Record | undefined, key: string): number | undefined { - if (!record || !(key in record)) return undefined; - return toFiniteNumber(record[key]); -} - -function toFiniteNumber(value: unknown): number { - const parsed = - typeof value === "number" - ? value - : typeof value === "string" && value.trim() - ? Number(value) - : 0; - return Number.isFinite(parsed) ? parsed : 0; + if (!record || !isTokenCount(record[key])) return undefined; + return Number(record[key]); } export function buildClaudeCodeTranscript(messages: ClaudeSdkMessage[]): string { @@ -365,3 +373,11 @@ export function stringifyError(value: unknown): string { export function clip(value: string, maxLength: number): string { return value.length <= maxLength ? value : `${value.slice(0, maxLength - 1)}…`; } + +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); +} diff --git a/packages/integrations/claude-agent-sdk/tests/session.test.ts b/packages/integrations/claude-agent-sdk/tests/session.test.ts index 7ecdf9033..480ad8257 100644 --- a/packages/integrations/claude-agent-sdk/tests/session.test.ts +++ b/packages/integrations/claude-agent-sdk/tests/session.test.ts @@ -1,6 +1,7 @@ /* eslint-disable require-yield */ import { describe, expect, it } from "vitest"; import { + extractClaudeCodeTokenUsage, isClaudeCodeMaxTurnsError, normalizeClaudeModel, runClaudeAgentSession, @@ -198,3 +199,29 @@ describe("Claude Agent SDK session", () => { expect(removed).toEqual(["abort"]); }); }); + +describe("Claude token usage presence", () => { + it.each([ + undefined, + {}, + { usage: {} }, + { usage: { total_tokens: 0 } }, + { usage: { input_tokens: null, output_tokens: -1 } }, + ])("keeps missing or invalid telemetry unreported: %j", (message) => { + expect(extractClaudeCodeTokenUsage(message).reported).toBe(false); + }); + it("preserves observed zero and valid modelUsage fallback", () => { + expect( + extractClaudeCodeTokenUsage({ usage: { input_tokens: 0, output_tokens: "0" } }), + ).toMatchObject({ reported: true, totalTokens: 0 }); + expect( + extractClaudeCodeTokenUsage({ + usage: { input_tokens: "invalid" }, + modelUsage: { model: { inputTokens: 9, outputTokens: 0 } }, + }), + ).toMatchObject({ reported: true, inputTokens: 9, outputTokens: 0 }); + expect( + extractClaudeCodeTokenUsage({ modelUsage: { model: { inputTokens: 0 } } }).reported, + ).toBe(true); + }); +}); diff --git a/packages/integrations/codex-sdk/src/session.ts b/packages/integrations/codex-sdk/src/session.ts index 1c7a77432..9d0e25e9f 100644 --- a/packages/integrations/codex-sdk/src/session.ts +++ b/packages/integrations/codex-sdk/src/session.ts @@ -1,3 +1,6 @@ +import type { Dirent } from "node:fs"; +import fsp from "node:fs/promises"; +import path from "node:path"; import { HarnessAdapterError, harnessEventLogLevel, @@ -36,12 +39,21 @@ export type CodexTokenUsage = { reasoning_output_tokens?: number; }; +/** + * Where the session's token usage came from: the `turn.completed` event, the + * thread's rollout file under CODEX_HOME (when the turn was aborted before it + * completed), or nowhere (all-zero usage that must not be trusted). + */ +export type CodexUsageSource = "turn_completed" | "rollout" | "none"; + export type CodexSessionResult = { events: CodexEvent[]; finalMessage: string; status: "completed" | "max_turns" | "sdk_error"; stopReason?: string; tokenUsage: CodexTokenUsage; + usageSource: CodexUsageSource; + threadId?: string; iterationError?: unknown; }; @@ -131,6 +143,13 @@ export async function runCodexSession(input: { outputSchema?: Record; maxToolSteps?: number; onToolStep?: () => void | Promise; + /** + * CODEX_HOME the binary runs with. `codex exec` only reports usage on + * `turn.completed`, which never arrives when the turn is aborted (step + * budget, caller signal); the cumulative count is then read back from the + * thread's rollout file under this directory. + */ + codexHome?: string; }): Promise { const sdk = input.sdk ?? (await loadCodexSdk()); const events: CodexEvent[] = []; @@ -138,6 +157,8 @@ export async function runCodexSession(input: { let stopReason: string | undefined; let iterationError: unknown; let tokenUsage = emptyTokenUsage(); + let usageSource: CodexUsageSource = "none"; + let threadId: string | undefined; const maxToolSteps = positiveInteger(input.maxToolSteps, 100); const budgetController = new AbortController(); const forwardAbort = () => budgetController.abort(input.signal?.reason); @@ -168,8 +189,11 @@ export async function runCodexSession(input: { for await (const event of streamed.events) { events.push(event); logCodexEvent(input.logger, event); - if (event.type === "turn.completed" && isRecord(event.usage)) { + if (event.type === "thread.started" && typeof event.thread_id === "string") { + threadId = event.thread_id; + } else if (event.type === "turn.completed" && isRecord(event.usage)) { tokenUsage = extractCodexTokenUsage(event.usage); + usageSource = "turn_completed"; } else if (event.type === "turn.failed") { stopReason = readCodexErrorMessage(event.error); } else if (event.type === "error") { @@ -209,16 +233,98 @@ export async function runCodexSession(input: { input.signal?.removeEventListener("abort", forwardAbort); } + if (usageSource === "none" && input.codexHome && threadId) { + const recovered = await readCodexRolloutUsage(input.codexHome, threadId); + if (recovered) { + tokenUsage = recovered; + usageSource = "rollout"; + input.logger.log({ + category: "codex", + level: 1, + message: `token usage recovered from rollout (turn never completed): in=${recovered.input_tokens} out=${recovered.output_tokens}`, + }); + } else { + input.logger.warn({ + category: "codex", + level: 1, + message: "token usage unavailable: turn never completed and no rollout token_count found", + }); + } + } + return { events, finalMessage, status: resolveCodexStatus(iterationError, stopReason, budgetExhausted), ...(stopReason && { stopReason: sanitizeErrorMessage(stopReason) }), tokenUsage, + usageSource, + ...(threadId && { threadId }), ...(iterationError !== undefined && { iterationError }), }; } +/** + * Cumulative usage of a thread from its rollout under + * `/sessions/YYYY/MM/DD/rollout--.jsonl`. + * Codex appends a `token_count` event after every model response, so the last + * one on disk covers everything billed before the process was killed. + */ +export async function readCodexRolloutUsage( + codexHome: string, + threadId: string, +): Promise { + const rollout = await findCodexRollout(codexHome, threadId); + if (!rollout) return undefined; + try { + return parseCodexRolloutUsage(await fsp.readFile(rollout, "utf8")); + } catch { + return undefined; + } +} + +async function findCodexRollout(codexHome: string, threadId: string): Promise { + const sessionsRoot = path.join(codexHome, "sessions"); + const pending = [sessionsRoot]; + while (pending.length > 0) { + const dir = pending.pop()!; + let entries: Dirent[]; + try { + entries = await fsp.readdir(dir, { withFileTypes: true }); + } catch { + continue; + } + for (const entry of entries) { + const full = path.join(dir, entry.name); + if (entry.isDirectory()) pending.push(full); + else if (entry.name.startsWith("rollout-") && entry.name.endsWith(`-${threadId}.jsonl`)) { + return full; + } + } + } + return undefined; +} + +/** Last `token_count` total in a rollout JSONL body; undefined when there is none. */ +export function parseCodexRolloutUsage(body: string): CodexTokenUsage | undefined { + let latest: Record | undefined; + for (const line of body.split(/\r?\n/u)) { + if (!line.includes('"token_count"')) continue; + let record: unknown; + try { + record = JSON.parse(line); + } catch { + continue; + } + if (!isRecord(record)) continue; + const payload = isRecord(record.payload) ? record.payload : record; + if (payload.type !== "token_count" || !isRecord(payload.info)) continue; + const total = payload.info.total_token_usage; + if (isRecord(total)) latest = total; + } + return latest ? extractCodexTokenUsage(latest) : undefined; +} + export function resolveCodexStatus( iterationError: unknown, stopReason?: string, diff --git a/packages/integrations/codex-sdk/tests/session.test.ts b/packages/integrations/codex-sdk/tests/session.test.ts index b26971945..e6c79d875 100644 --- a/packages/integrations/codex-sdk/tests/session.test.ts +++ b/packages/integrations/codex-sdk/tests/session.test.ts @@ -1,6 +1,15 @@ /* eslint-disable require-yield */ +import fsp from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; import { describe, expect, it } from "vitest"; -import { normalizeCodexModel, runCodexSession, type CodexSdk } from "../src/index.js"; +import { + normalizeCodexModel, + parseCodexRolloutUsage, + readCodexRolloutUsage, + runCodexSession, + type CodexSdk, +} from "../src/index.js"; const logger = { log: () => {}, warn: () => {}, error: () => {} }; @@ -105,6 +114,181 @@ describe("Codex SDK session", () => { expect(signal?.aborted).toBe(true); expect(result.status).toBe("max_turns"); expect(result.stopReason).toBe("tool step budget exhausted (1 steps)"); + expect(result.usageSource).toBe("none"); + }); + + it("records where usage came from on a completed turn", async () => { + const sdk: CodexSdk = { + startThread: () => ({ + runStreamed: async () => ({ + events: (async function* () { + yield { type: "thread.started", thread_id: "thread-1" }; + yield { type: "turn.completed", usage: { input_tokens: 5, output_tokens: 2 } }; + })(), + }), + }), + }; + const result = await runCodexSession({ prompt: "task", model: "m", logger, sdk, thread: {} }); + expect(result.usageSource).toBe("turn_completed"); + expect(result.threadId).toBe("thread-1"); + }); + + describe("rollout usage recovery", () => { + const rolloutBody = [ + JSON.stringify({ type: "session_meta", payload: { id: "thread-abc" } }), + JSON.stringify({ + type: "event_msg", + payload: { + type: "token_count", + info: { + total_token_usage: { + input_tokens: 1000, + cached_input_tokens: 800, + cache_write_input_tokens: 0, + output_tokens: 50, + reasoning_output_tokens: 20, + total_tokens: 1050, + }, + last_token_usage: { input_tokens: 1000, output_tokens: 50 }, + }, + }, + }), + JSON.stringify({ type: "response_item", payload: { type: "function_call" } }), + JSON.stringify({ + type: "event_msg", + payload: { + type: "token_count", + info: { + total_token_usage: { + input_tokens: 2_100_000, + cached_input_tokens: 2_000_000, + cache_write_input_tokens: 5000, + output_tokens: 9000, + reasoning_output_tokens: 4000, + total_tokens: 2_109_000, + }, + last_token_usage: { input_tokens: 2_099_000, output_tokens: 8950 }, + }, + }, + }), + "not json {", + JSON.stringify({ type: "event_msg", payload: { type: "token_count", info: null } }), + "", + ].join("\n"); + + async function writeRollout(threadId: string, body = rolloutBody): Promise { + const codexHome = await fsp.mkdtemp(path.join(os.tmpdir(), "codex-home-")); + const dir = path.join(codexHome, "sessions", "2026", "08", "31"); + await fsp.mkdir(dir, { recursive: true }); + await fsp.writeFile(path.join(dir, `rollout-2026-08-31T10-00-00-${threadId}.jsonl`), body); + return codexHome; + } + + it("takes the last cumulative token_count and ignores malformed lines", () => { + expect(parseCodexRolloutUsage(rolloutBody)).toEqual({ + input_tokens: 2_100_000, + cached_input_tokens: 2_000_000, + output_tokens: 9000, + reasoning_output_tokens: 4000, + }); + expect(parseCodexRolloutUsage("")).toBeUndefined(); + expect( + parseCodexRolloutUsage('{"type":"event_msg","payload":{"type":"agent_message"}}'), + ).toBeUndefined(); + }); + + it("finds the thread's rollout under CODEX_HOME/sessions", async () => { + const codexHome = await writeRollout("thread-abc"); + await expect(readCodexRolloutUsage(codexHome, "thread-abc")).resolves.toMatchObject({ + input_tokens: 2_100_000, + }); + await expect(readCodexRolloutUsage(codexHome, "thread-other")).resolves.toBeUndefined(); + await expect( + readCodexRolloutUsage("/nonexistent/codex-home", "thread-abc"), + ).resolves.toBeUndefined(); + }); + + it("recovers usage from the rollout when the budget abort pre-empts turn.completed", async () => { + const codexHome = await writeRollout("thread-abc"); + const logged: string[] = []; + const sdk: CodexSdk = { + startThread: () => ({ + runStreamed: async () => ({ + events: (async function* () { + yield { type: "thread.started", thread_id: "thread-abc" }; + yield { type: "item.completed", item: { type: "mcp_tool_call" } }; + })(), + }), + }), + }; + const result = await runCodexSession({ + prompt: "task", + model: "gpt-5.4-mini", + logger: { ...logger, log: (line: { message: string }) => logged.push(line.message) }, + sdk, + thread: {}, + maxToolSteps: 1, + codexHome, + }); + expect(result.status).toBe("max_turns"); + expect(result.usageSource).toBe("rollout"); + expect(result.tokenUsage).toEqual({ + input_tokens: 2_100_000, + cached_input_tokens: 2_000_000, + output_tokens: 9000, + reasoning_output_tokens: 4000, + }); + expect(logged.some((message) => message.includes("recovered from rollout"))).toBe(true); + }); + + it("does not override usage the turn reported itself", async () => { + const codexHome = await writeRollout("thread-abc"); + const sdk: CodexSdk = { + startThread: () => ({ + runStreamed: async () => ({ + events: (async function* () { + yield { type: "thread.started", thread_id: "thread-abc" }; + yield { type: "turn.completed", usage: { input_tokens: 7, output_tokens: 3 } }; + })(), + }), + }), + }; + const result = await runCodexSession({ + prompt: "task", + model: "m", + logger, + sdk, + thread: {}, + codexHome, + }); + expect(result.usageSource).toBe("turn_completed"); + expect(result.tokenUsage.input_tokens).toBe(7); + }); + + it("leaves usage at none when the rollout is missing", async () => { + const codexHome = await fsp.mkdtemp(path.join(os.tmpdir(), "codex-home-")); + const sdk: CodexSdk = { + startThread: () => ({ + runStreamed: async () => ({ + events: (async function* () { + yield { type: "thread.started", thread_id: "thread-missing" }; + yield { type: "item.completed", item: { type: "command_execution" } }; + })(), + }), + }), + }; + const result = await runCodexSession({ + prompt: "task", + model: "m", + logger, + sdk, + thread: {}, + maxToolSteps: 1, + codexHome, + }); + expect(result.usageSource).toBe("none"); + expect(result.tokenUsage).toEqual({ input_tokens: 0, output_tokens: 0 }); + }); }); it("fails closed to read-only for unknown sandbox modes", async () => { diff --git a/packages/integrations/deepagents-sdk/src/session.ts b/packages/integrations/deepagents-sdk/src/session.ts index 98e39cdd9..c676334d6 100644 --- a/packages/integrations/deepagents-sdk/src/session.ts +++ b/packages/integrations/deepagents-sdk/src/session.ts @@ -31,6 +31,8 @@ export type DeepagentsSessionConfig = { }; export type DeepagentsTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; inputTokens: number; outputTokens: number; cacheReadInputTokens: number; @@ -350,6 +352,8 @@ export function extractDeepagentsTokenUsage( event: Record | undefined, ): DeepagentsTokenUsage { return { + reported: + event?.reported !== false && [event?.input_tokens, event?.output_tokens].some(isTokenCount), inputTokens: toFiniteNumber(event?.input_tokens), outputTokens: toFiniteNumber(event?.output_tokens), cacheReadInputTokens: toFiniteNumber(event?.cache_read_input_tokens), @@ -532,3 +536,11 @@ export function toFiniteNumber(value: unknown): number { : 0; return Number.isFinite(parsed) ? parsed : 0; } + +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); +} diff --git a/packages/integrations/deepagents-sdk/tests/session.test.ts b/packages/integrations/deepagents-sdk/tests/session.test.ts index 3736a0970..f0096dba2 100644 --- a/packages/integrations/deepagents-sdk/tests/session.test.ts +++ b/packages/integrations/deepagents-sdk/tests/session.test.ts @@ -1,6 +1,7 @@ import { PassThrough } from "node:stream"; import { describe, expect, it, vi } from "vitest"; import { + extractDeepagentsTokenUsage, buildDeepagentsRunnerArgs, buildDeepagentsTranscript, normalizeDeepagentsModel, @@ -99,6 +100,7 @@ describe("Deep Agents session", () => { expect(result.finalMessage).toBe("complete"); expect(result.status).toBe("completed"); expect(result.tokenUsage).toEqual({ + reported: true, inputTokens: 10, outputTokens: 4, cacheReadInputTokens: 3, @@ -420,3 +422,24 @@ describe("Deep Agents session", () => { expect(result.stopReason).toContain("SIGKILL"); }); }); + +describe("DeepAgents token usage presence", () => { + it.each([ + undefined, + {}, + { total_tokens: 0 }, + { input_tokens: null, output_tokens: -1 }, + { reported: false, input_tokens: 0, output_tokens: 0 }, + ])("keeps missing or explicitly unavailable usage unreported: %j", (event) => { + expect(extractDeepagentsTokenUsage(event).reported).toBe(false); + }); + it("preserves observed zero and legacy usage events", () => { + expect(extractDeepagentsTokenUsage({ input_tokens: 0, output_tokens: "0" })).toMatchObject({ + reported: true, + totalTokens: 0, + }); + expect( + extractDeepagentsTokenUsage({ reported: true, input_tokens: 2, output_tokens: 1 }), + ).toMatchObject({ reported: true, inputTokens: 2, outputTokens: 1 }); + }); +}); diff --git a/packages/integrations/deepagents/runner/run_eval.py b/packages/integrations/deepagents/runner/run_eval.py index 73e35b522..cb80efa67 100644 --- a/packages/integrations/deepagents/runner/run_eval.py +++ b/packages/integrations/deepagents/runner/run_eval.py @@ -259,7 +259,7 @@ def message_events(message: object, tool_servers: Mapping[str, str]) -> list[Eve return [] -def aggregate_usage(usages: list[Mapping[str, object] | None]) -> dict[str, int]: +def aggregate_usage(usages: list[Mapping[str, object] | None]) -> dict[str, int | bool]: totals = { "input_tokens": 0, "output_tokens": 0, @@ -267,9 +267,16 @@ def aggregate_usage(usages: list[Mapping[str, object] | None]) -> dict[str, int] "reasoning_output_tokens": 0, "total_tokens": 0, } + reported = False for usage in usages: if not usage: continue + reported = reported or any( + isinstance(usage.get(key), int) + and not isinstance(usage.get(key), bool) + and usage[key] >= 0 + for key in ("input_tokens", "output_tokens") + ) totals["input_tokens"] += _integer(usage.get("input_tokens")) totals["output_tokens"] += _integer(usage.get("output_tokens")) totals["total_tokens"] += _integer(usage.get("total_tokens")) @@ -279,7 +286,7 @@ def aggregate_usage(usages: list[Mapping[str, object] | None]) -> dict[str, int] totals["cache_read_input_tokens"] += _integer(input_details.get("cache_read")) if isinstance(output_details, Mapping): totals["reasoning_output_tokens"] += _integer(output_details.get("reasoning")) - return totals + return {**totals, "reported": reported} def _integer(value: object) -> int: diff --git a/packages/integrations/deepagents/runner/tests/test_run_eval.py b/packages/integrations/deepagents/runner/tests/test_run_eval.py index 906b68522..179009f2c 100644 --- a/packages/integrations/deepagents/runner/tests/test_run_eval.py +++ b/packages/integrations/deepagents/runner/tests/test_run_eval.py @@ -113,6 +113,7 @@ async def test_streams_tool_sequence_and_sums_usage() -> None: assert events[-2]["text"] == "done" assert events[-1] == { "type": "usage", + "reported": True, "input_tokens": 15, "output_tokens": 6, "cache_read_input_tokens": 5, @@ -429,6 +430,7 @@ def test_content_blocks_and_usage_helpers() -> None: {"data": "ZGVm", "mime_type": "image/jpeg"}, ] assert aggregate_usage([usage(2, 3, 1, 2), usage(4, 5, 2, 3)]) == { + "reported": True, "input_tokens": 6, "output_tokens": 8, "cache_read_input_tokens": 3, @@ -459,3 +461,15 @@ async def aclose(self) -> None: assert exit_code == 0 assert all(event["type"] != "error" for event in events) + + +@pytest.mark.parametrize("usages", [[], [None], [{}], [{"input_tokens": None, "output_tokens": -1}], [{"total_tokens": 0}]]) +def test_usage_presence_missing(usages: list[dict[str, Any] | None]) -> None: + assert aggregate_usage(usages)["reported"] is False + + +def test_usage_presence_preserves_observed_zero() -> None: + result = aggregate_usage([usage(0, 0, 0, 0), None, {}]) + assert result["reported"] is True + assert result["input_tokens"] == 0 + assert result["output_tokens"] == 0 diff --git a/packages/integrations/eve-sdk/src/session.ts b/packages/integrations/eve-sdk/src/session.ts index 370b89de7..c9d2fb258 100644 --- a/packages/integrations/eve-sdk/src/session.ts +++ b/packages/integrations/eve-sdk/src/session.ts @@ -17,6 +17,8 @@ export type EveEvent = { }; export type EveTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; inputTokens: number; outputTokens: number; cacheReadTokens: number; @@ -339,10 +341,12 @@ export function extractEveTokenUsage(events: EveEvent[]): EveTokenUsage { const usage = { inputTokens: 0, outputTokens: 0, cacheReadTokens: 0, cacheWriteTokens: 0 }; let costUsd = 0; let hasCostUsd = false; + let reported = false; for (const event of events) { if (event.type !== "step.completed") continue; const data = isRecord(event.data) ? event.data : undefined; const stepUsage = isRecord(data?.usage) ? data.usage : undefined; + reported ||= [stepUsage?.inputTokens, stepUsage?.outputTokens].some(isTokenCount); usage.inputTokens += toFiniteNumber(stepUsage?.inputTokens); usage.outputTokens += toFiniteNumber(stepUsage?.outputTokens); usage.cacheReadTokens += toFiniteNumber(stepUsage?.cacheReadTokens); @@ -354,6 +358,7 @@ export function extractEveTokenUsage(events: EveEvent[]): EveTokenUsage { } return { ...usage, + reported, totalTokens: Object.values(usage).reduce((sum, value) => sum + value, 0), ...(hasCostUsd && { costUsd }), }; @@ -489,3 +494,11 @@ export function stringifyError(value: unknown): string { export function clip(value: string, maxLength: number): string { return value.length <= maxLength ? value : `${value.slice(0, maxLength - 1)}…`; } + +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); +} diff --git a/packages/integrations/eve-sdk/tests/session.test.ts b/packages/integrations/eve-sdk/tests/session.test.ts index 9d8ef17da..7a16c6eaa 100644 --- a/packages/integrations/eve-sdk/tests/session.test.ts +++ b/packages/integrations/eve-sdk/tests/session.test.ts @@ -4,6 +4,7 @@ import { PassThrough } from "node:stream"; import { describe, expect, it, vi } from "vitest"; import { HarnessAdapterError } from "@browserbasehq/stagehand-integrations/harness"; import { + extractEveTokenUsage, logEveEvent, parseEveDevServerUrl, runEveSession, @@ -217,6 +218,7 @@ describe("Eve SDK session", () => { sessionId: "session-1", serverUrl: "http://eve", tokenUsage: { + reported: true, inputTokens: 15, outputTokens: 7, cacheReadTokens: 2, @@ -633,3 +635,23 @@ describe("Eve SDK session", () => { expect(JSON.stringify(result)).not.toContain("SUPERSECRET"); }); }); + +describe("Eve token usage presence", () => { + it.each([{}, { costUsd: 0 }, { totalTokens: 0 }, { inputTokens: null, outputTokens: -1 }])( + "does not infer tokens from missing or invalid telemetry: %j", + (usage) => { + expect(extractEveTokenUsage([{ type: "step.completed", data: { usage } }]).reported).toBe( + false, + ); + }, + ); + it("preserves observed zero across missing later steps without inventing cost", () => { + expect(extractEveTokenUsage([]).reported).toBe(false); + const usage = extractEveTokenUsage([ + { type: "step.completed", data: { usage: { inputTokens: 0, outputTokens: "0" } } }, + { type: "step.completed", data: {} }, + ]); + expect(usage).toMatchObject({ reported: true, totalTokens: 0 }); + expect(usage).not.toHaveProperty("costUsd"); + }); +}); diff --git a/packages/integrations/fx-sdk/src/session.ts b/packages/integrations/fx-sdk/src/session.ts index 0aa6d9949..a0d343f1d 100644 --- a/packages/integrations/fx-sdk/src/session.ts +++ b/packages/integrations/fx-sdk/src/session.ts @@ -63,6 +63,8 @@ export type FxEvent = | { type: "turn_committed"; terminal_reason?: string; turn_kind?: string }; export type FxTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; input_tokens: number; cached_input_tokens: number; output_tokens: number; @@ -603,6 +605,9 @@ export function extractFxTokenUsage( const committed = findLastCommittedTurn(events); if (committed) { return { + reported: [committed.payload.total_input_tokens, committed.payload.total_output_tokens].some( + isTokenCount, + ), input_tokens: toFiniteNumber(committed.payload.total_input_tokens), cached_input_tokens: 0, output_tokens: toFiniteNumber(committed.payload.total_output_tokens), @@ -792,13 +797,15 @@ function hasUsageFields(record: Record): boolean { function usageFromRecord(record?: Record): FxTokenUsage { return { + reported: [record?.input_tokens, record?.output_tokens].some(isTokenCount), input_tokens: toFiniteNumber(record?.input_tokens), cached_input_tokens: toFiniteNumber(record?.cache_read_tokens), output_tokens: toFiniteNumber(record?.output_tokens), reasoning_output_tokens: toFiniteNumber(record?.reasoning_tokens), - ...(record && - "total_cost" in record && { - total_cost: toFiniteNumber(record.total_cost), + ...(typeof record?.total_cost === "number" && + Number.isFinite(record.total_cost) && + record.total_cost >= 0 && { + total_cost: record.total_cost, }), }; } @@ -857,3 +864,11 @@ export function stringifyError(value: unknown): string { export function clip(value: string, maxLength: number): string { return value.length <= maxLength ? value : `${value.slice(0, maxLength - 1)}…`; } + +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); +} diff --git a/packages/integrations/fx-sdk/tests/session.test.ts b/packages/integrations/fx-sdk/tests/session.test.ts index b715fa178..564a2787c 100644 --- a/packages/integrations/fx-sdk/tests/session.test.ts +++ b/packages/integrations/fx-sdk/tests/session.test.ts @@ -4,6 +4,7 @@ import type { ChildProcess } from "node:child_process"; import { describe, expect, it, vi } from "vitest"; import { HarnessAdapterError } from "@browserbasehq/stagehand-integrations/harness"; import { + extractFxTokenUsage, buildFxTranscript, createFxProcessRunner, normalizeFxModel, @@ -120,6 +121,7 @@ describe("fx CLI session", () => { "ask_result", ]); expect(result.tokenUsage).toEqual({ + reported: true, input_tokens: 100, cached_input_tokens: 30, output_tokens: 20, @@ -545,3 +547,46 @@ describe("fx CLI session", () => { expect(buildFxTranscript(result.events)).not.toContain("1234567890"); }); }); + +describe("FX token usage presence", () => { + it.each([null, undefined, -1, NaN, Infinity, "invalid"])( + "omits invalid reported bills: %j", + (total_cost) => { + expect(extractFxTokenUsage([], { total_cost })).not.toHaveProperty("total_cost"); + }, + ); + it("preserves an explicitly reported zero-dollar bill", () => { + expect(extractFxTokenUsage([], { total_cost: 0 })).toMatchObject({ + total_cost: 0, + reported: false, + }); + }); + it.each([undefined, {}, { total_cost: 0 }, { input_tokens: null, output_tokens: -1 }])( + "keeps missing or cost-only telemetry unreported: %j", + (snapshot) => { + expect(extractFxTokenUsage([], snapshot).reported).toBe(false); + }, + ); + it("preserves observed zero in snapshots, committed totals and checkpoints", () => { + expect( + extractFxTokenUsage([], { snapshot: { input_tokens: 0, output_tokens: "0" } }), + ).toMatchObject({ reported: true, input_tokens: 0, output_tokens: 0 }); + const committed = committedEvent(); + committed.payload.total_input_tokens = 0; + committed.payload.total_output_tokens = 0; + expect(extractFxTokenUsage([committed]).reported).toBe(true); + expect( + extractFxTokenUsage([ + { kind: "usage_checkpointed", payload: { usage: { input_tokens: 0, output_tokens: 0 } } }, + ]).reported, + ).toBe(true); + expect( + extractFxTokenUsage([ + { + kind: "history_turn_committed", + payload: { turn: { kind: "completed", assistant: "done" } }, + }, + ]).reported, + ).toBe(false); + }); +}); diff --git a/packages/integrations/mastra-sdk/src/session.ts b/packages/integrations/mastra-sdk/src/session.ts index 58935a375..1c9ba9eb4 100644 --- a/packages/integrations/mastra-sdk/src/session.ts +++ b/packages/integrations/mastra-sdk/src/session.ts @@ -71,6 +71,8 @@ export type MastraSessionConfig = { }; export type MastraTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; inputTokens: number; outputTokens: number; reasoningTokens: number; @@ -310,7 +312,6 @@ export async function runMastraSession(input: { if (stream.error && !stopReason) { stopReason = sanitizeErrorMessage(stringifyError(stream.error)); } - if (isEmptyTokenUsage(tokenUsage)) tokenUsage = summedStepUsage; } } catch (error) { iterationError = error; @@ -359,7 +360,7 @@ export async function runMastraSession(input: { ...(stopReason && { stopReason: sanitizeErrorMessage(stopReason) }), ...(finishReason && { finishReason }), stepCount, - tokenUsage, + tokenUsage: tokenUsage.reported ? tokenUsage : summedStepUsage, ...(iterationError !== undefined && { iterationError }), }; } @@ -387,6 +388,7 @@ export function extractMastraTokenUsage( const inputTokens = toFiniteNumber(usage?.inputTokens); const outputTokens = toFiniteNumber(usage?.outputTokens); return { + reported: [usage?.inputTokens, usage?.outputTokens].some(isTokenCount), inputTokens, outputTokens, reasoningTokens: toFiniteNumber(usage?.reasoningTokens), @@ -617,6 +619,7 @@ function flattenMcpResult(value: unknown): string { function addTokenUsage(left: MastraTokenUsage, right: MastraTokenUsage): MastraTokenUsage { return { + reported: left.reported === true || right.reported === true, inputTokens: left.inputTokens + right.inputTokens, outputTokens: left.outputTokens + right.outputTokens, reasoningTokens: left.reasoningTokens + right.reasoningTokens, @@ -625,6 +628,10 @@ function addTokenUsage(left: MastraTokenUsage, right: MastraTokenUsage): MastraT }; } -function isEmptyTokenUsage(usage: MastraTokenUsage): boolean { - return Object.values(usage).every((value) => value === 0); +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); } diff --git a/packages/integrations/mastra-sdk/tests/session.test.ts b/packages/integrations/mastra-sdk/tests/session.test.ts index 095b9d0ed..f7a4160ec 100644 --- a/packages/integrations/mastra-sdk/tests/session.test.ts +++ b/packages/integrations/mastra-sdk/tests/session.test.ts @@ -1,6 +1,7 @@ /* eslint-disable require-yield */ import { describe, expect, it, vi } from "vitest"; import { + extractMastraTokenUsage, buildMastraTranscript, compactMastraEvent, normalizeMastraModel, @@ -141,6 +142,7 @@ describe("Mastra SDK session", () => { expect(result.finishReason).toBe("stop"); expect(result.stepCount).toBe(1); expect(result.tokenUsage).toEqual({ + reported: true, inputTokens: 100, outputTokens: 25, reasoningTokens: 5, @@ -428,3 +430,56 @@ describe("Mastra SDK session", () => { expect(result.finalText).toBe("done https://x.test?apiKey=[redacted]"); }); }); + +describe("Mastra token usage presence", () => { + it.each([undefined, {}, { totalTokens: 0 }, { inputTokens: null, outputTokens: -1 }])( + "keeps absent or invalid telemetry unreported: %j", + (usage) => { + expect(extractMastraTokenUsage(usage).reported).toBe(false); + }, + ); + it("recognizes observed zero", () => { + expect(extractMastraTokenUsage({ inputTokens: 0, outputTokens: "0" })).toMatchObject({ + reported: true, + totalTokens: 0, + }); + }); + it.each([ + { finish: undefined, expectedInput: 7, reported: true }, + { finish: {}, expectedInput: 7, reported: true }, + { finish: { inputTokens: 0, outputTokens: 0 }, expectedInput: 0, reported: true }, + ])( + "uses step usage only when finish telemetry is absent: %j", + async ({ finish, expectedInput, reported }) => { + const result = await runMastraSession({ + prompt: "task", + model: "gpt-5.4-mini", + logger, + session: {}, + sdk: fakeSdk({ + events: [ + { + type: "step-finish", + payload: { output: { usage: { inputTokens: 7, outputTokens: 3 } } }, + }, + { + type: "finish", + payload: { stepResult: { reason: "stop" }, output: { usage: finish } }, + }, + ], + }), + }); + expect(result.tokenUsage).toMatchObject({ reported, inputTokens: expectedInput }); + }, + ); + it("does not invent usage on an empty stream", async () => { + const result = await runMastraSession({ + prompt: "task", + model: "gpt-5.4-mini", + logger, + session: {}, + sdk: fakeSdk({ events: [] }), + }); + expect(result.tokenUsage.reported).toBe(false); + }); +}); diff --git a/packages/integrations/pi-sdk/src/session.ts b/packages/integrations/pi-sdk/src/session.ts index 5336cc074..9dad319d7 100644 --- a/packages/integrations/pi-sdk/src/session.ts +++ b/packages/integrations/pi-sdk/src/session.ts @@ -49,6 +49,8 @@ export type PiSessionConfig = { mcpServers?: Record; }; export type PiTokenUsage = { + /** Whether token counters were observed; false distinguishes missing telemetry from zero. */ + reported?: boolean; inputTokens: number; outputTokens: number; cacheReadTokens: number; @@ -267,6 +269,7 @@ export async function runPiSession(input: { export function extractPiTokenUsage(events: PiEvent[]): PiTokenUsage { const total: PiTokenUsage = { + reported: false, inputTokens: 0, outputTokens: 0, cacheReadTokens: 0, @@ -281,6 +284,7 @@ export function extractPiTokenUsage(events: PiEvent[]): PiTokenUsage { const message = event.message; if (message.role !== "assistant" || !isRecord(message.usage)) continue; const usage = message.usage; + total.reported ||= [usage.input, usage.output].some(isTokenCount); total.inputTokens += toFiniteNumber(usage.input); total.outputTokens += toFiniteNumber(usage.output); total.cacheReadTokens += toFiniteNumber(usage.cacheRead); @@ -495,3 +499,11 @@ export function toFiniteNumber(value: unknown): number { : 0; return Number.isFinite(parsed) ? parsed : 0; } + +function isTokenCount(value: unknown): boolean { + return ( + ((typeof value === "number" && Number.isFinite(value)) || + (typeof value === "string" && value.trim().length > 0 && Number.isFinite(Number(value)))) && + Number(value) >= 0 + ); +} diff --git a/packages/integrations/pi-sdk/tests/session.test.ts b/packages/integrations/pi-sdk/tests/session.test.ts index e6b06657e..549ee8420 100644 --- a/packages/integrations/pi-sdk/tests/session.test.ts +++ b/packages/integrations/pi-sdk/tests/session.test.ts @@ -1,5 +1,6 @@ import { describe, expect, it, vi } from "vitest"; import { + extractPiTokenUsage, buildPiMcpToolName, buildPiTranscript, definePiCodeRunTool, @@ -132,6 +133,7 @@ describe("pi SDK session", () => { expect(result.turns).toBe(2); expect(result.events.some((event) => event.type === "message_update")).toBe(false); expect(result.tokenUsage).toEqual({ + reported: true, inputTokens: 15, outputTokens: 8, cacheReadTokens: 4, @@ -428,3 +430,28 @@ describe("pi SDK session", () => { expect(removeListener).toHaveBeenCalledWith("abort", expect.any(Function)); }); }); + +describe("Pi token usage presence", () => { + const event = (usage: Record) => ({ + type: "message_end", + message: { role: "assistant", usage }, + }); + it.each([{}, { cost: { total: 0 } }, { totalTokens: 0 }, { input: null, output: -1 }])( + "does not infer tokens from absent or invalid telemetry: %j", + (usage) => { + expect(extractPiTokenUsage([event(usage)]).reported).toBe(false); + }, + ); + it("preserves observed zero and presence across later missing usage", () => { + expect(extractPiTokenUsage([]).reported).toBe(false); + expect(extractPiTokenUsage([event({ input: 0, output: "0" }), event({})])).toMatchObject({ + reported: true, + inputTokens: 0, + outputTokens: 0, + }); + expect( + extractPiTokenUsage([{ type: "message_end", message: { role: "user", usage: { input: 1 } } }]) + .reported, + ).toBe(false); + }); +}); From de45faba201044a137654c894617ce8a9d4b95cf Mon Sep 17 00:00:00 2001 From: miguel Date: Mon, 7 Sep 2026 17:42:29 -0700 Subject: [PATCH 2/2] Retain the last valid Codex rollout usage total --- .../integrations/codex-sdk/src/session.ts | 12 ++++++- .../codex-sdk/tests/session.test.ts | 34 +++++++++++++++++++ 2 files changed, 45 insertions(+), 1 deletion(-) diff --git a/packages/integrations/codex-sdk/src/session.ts b/packages/integrations/codex-sdk/src/session.ts index 9d0e25e9f..0eb8f6f68 100644 --- a/packages/integrations/codex-sdk/src/session.ts +++ b/packages/integrations/codex-sdk/src/session.ts @@ -320,7 +320,17 @@ export function parseCodexRolloutUsage(body: string): CodexTokenUsage | undefine const payload = isRecord(record.payload) ? record.payload : record; if (payload.type !== "token_count" || !isRecord(payload.info)) continue; const total = payload.info.total_token_usage; - if (isRecord(total)) latest = total; + if ( + isRecord(total) && + typeof total.input_tokens === "number" && + Number.isFinite(total.input_tokens) && + total.input_tokens >= 0 && + typeof total.output_tokens === "number" && + Number.isFinite(total.output_tokens) && + total.output_tokens >= 0 + ) { + latest = total; + } } return latest ? extractCodexTokenUsage(latest) : undefined; } diff --git a/packages/integrations/codex-sdk/tests/session.test.ts b/packages/integrations/codex-sdk/tests/session.test.ts index e6c79d875..b31cb6658 100644 --- a/packages/integrations/codex-sdk/tests/session.test.ts +++ b/packages/integrations/codex-sdk/tests/session.test.ts @@ -197,6 +197,40 @@ describe("Codex SDK session", () => { ).toBeUndefined(); }); + it.each([ + {}, + { input_tokens: 10 }, + { output_tokens: 5 }, + { input_tokens: "10", output_tokens: 5 }, + { input_tokens: null, output_tokens: 5 }, + { input_tokens: -1, output_tokens: 5 }, + { input_tokens: 10, output_tokens: -1 }, + { input_tokens: 10, output_tokens: Number.POSITIVE_INFINITY }, + ])("ignores incomplete or invalid cumulative usage: %j", (total) => { + const incomplete = JSON.stringify({ + type: "event_msg", + payload: { type: "token_count", info: { total_token_usage: total } }, + }); + expect(parseCodexRolloutUsage(`${rolloutBody}\n${incomplete}`)).toEqual( + parseCodexRolloutUsage(rolloutBody), + ); + expect(parseCodexRolloutUsage(incomplete)).toBeUndefined(); + }); + + it("preserves explicitly reported zero cumulative usage", () => { + const zero = JSON.stringify({ + type: "event_msg", + payload: { + type: "token_count", + info: { total_token_usage: { input_tokens: 0, output_tokens: 0 } }, + }, + }); + expect(parseCodexRolloutUsage(`${rolloutBody}\n${zero}`)).toEqual({ + input_tokens: 0, + output_tokens: 0, + }); + }); + it("finds the thread's rollout under CODEX_HOME/sessions", async () => { const codexHome = await writeRollout("thread-abc"); await expect(readCodexRolloutUsage(codexHome, "thread-abc")).resolves.toMatchObject({