Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 28 additions & 12 deletions packages/integrations/claude-agent-sdk/src/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand All @@ -260,18 +278,8 @@ function sumModelUsage(resultMessage: ClaudeSdkMessage | undefined, key: string)
}

function readNumber(record: Record<string, unknown> | 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 {
Expand Down Expand Up @@ -371,3 +379,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
);
}
27 changes: 27 additions & 0 deletions packages/integrations/claude-agent-sdk/tests/session.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
/* eslint-disable require-yield */
import { describe, expect, it } from "vitest";
import {
extractClaudeCodeTokenUsage,
isClaudeCodeMaxTurnsError,
normalizeClaudeModel,
runClaudeAgentSession,
Expand Down Expand Up @@ -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);
});
});
118 changes: 117 additions & 1 deletion packages/integrations/codex-sdk/src/session.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
import type { Dirent } from "node:fs";
import fsp from "node:fs/promises";
import path from "node:path";
import {
HarnessAdapterError,
harnessEventLogLevel,
Expand Down Expand Up @@ -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;
};

Expand Down Expand Up @@ -131,13 +143,22 @@ export async function runCodexSession(input: {
outputSchema?: Record<string, unknown>;
maxToolSteps?: number;
onToolStep?: () => void | Promise<void>;
/**
* 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<CodexSessionResult> {
const sdk = input.sdk ?? (await loadCodexSdk());
const events: CodexEvent[] = [];
let finalMessage = "";
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);
Expand Down Expand Up @@ -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") {
Expand Down Expand Up @@ -209,16 +233,108 @@ 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
* `<codexHome>/sessions/YYYY/MM/DD/rollout-<timestamp>-<threadId>.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<CodexTokenUsage | undefined> {
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<string | undefined> {
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<string, unknown> | 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) &&
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;
}

export function resolveCodexStatus(
iterationError: unknown,
stopReason?: string,
Expand Down
Loading
Loading