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
2 changes: 1 addition & 1 deletion src/types/TokenMetrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ export interface TokenUsage {
}

export interface TranscriptLine {
message?: { usage?: TokenUsage; stop_reason?: string | null };
message?: { id?: string; usage?: TokenUsage; stop_reason?: string | null };
isSidechain?: boolean;
timestamp?: string;
isApiErrorMessage?: boolean;
Expand Down
90 changes: 90 additions & 0 deletions src/utils/__tests__/jsonl-metrics.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,12 +69,14 @@ function makeUsageLine(params: {
isSidechain?: boolean;
isApiErrorMessage?: boolean;
stopReason?: string | null;
messageId?: string;
}): string {
return JSON.stringify({
timestamp: params.timestamp,
isSidechain: params.isSidechain,
isApiErrorMessage: params.isApiErrorMessage,
message: {
id: params.messageId,
stop_reason: params.stopReason,
usage: {
input_tokens: params.input,
Expand All @@ -93,6 +95,7 @@ function makeTranscriptLine(params: {
output?: number;
isSidechain?: boolean;
isApiErrorMessage?: boolean;
messageId?: string;
}): string {
return JSON.stringify({
timestamp: params.timestamp,
Expand All @@ -101,6 +104,7 @@ function makeTranscriptLine(params: {
isApiErrorMessage: params.isApiErrorMessage,
message: typeof params.input === 'number' || typeof params.output === 'number'
? {
id: params.messageId,
usage: {
input_tokens: params.input ?? 0,
output_tokens: params.output ?? 0
Expand Down Expand Up @@ -327,6 +331,68 @@ describe('jsonl transcript metrics', () => {
});
});

it('counts one API response once when its content blocks are written as separate rows', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-metrics-'));
tempRoots.push(root);
const transcriptPath = path.join(root, 'content-blocks.jsonl');

// Claude Code writes one row per content block (thinking, text, tool_use)
// of the same response: same message.id, same finalized stop_reason,
// identical prompt-side usage, non-decreasing output_tokens.
const call1 = { input: 3, cacheRead: 20000, cacheCreate: 5000, stopReason: 'tool_use', messageId: 'msg_1' };
const call2 = { input: 4, cacheRead: 25000, cacheCreate: 0, stopReason: 'end_turn', messageId: 'msg_2' };
fs.writeFileSync(transcriptPath, [
makeUsageLine({ timestamp: '2026-01-01T10:00:00.000Z', output: 40, ...call1 }),
makeUsageLine({ timestamp: '2026-01-01T10:00:01.000Z', output: 90, ...call1 }),
makeUsageLine({ timestamp: '2026-01-01T10:00:02.000Z', output: 120, ...call1 }),
makeUsageLine({ timestamp: '2026-01-01T10:00:10.000Z', output: 15, ...call2 }),
makeUsageLine({ timestamp: '2026-01-01T10:00:11.000Z', output: 60, ...call2 })
].join('\n'));

const metrics = await getTokenMetrics(transcriptPath);

expect(metrics.inputTokens).toBe(7);
expect(metrics.outputTokens).toBe(180);
expect(metrics.cacheReadTokens).toBe(45000);
expect(metrics.cacheCreationTokens).toBe(5000);
expect(metrics.cachedTokens).toBe(50000);
expect(metrics.totalTokens).toBe(50187);
expect(metrics.contextLength).toBe(25004);
});

it('still counts every row when messages carry no id', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-metrics-'));
tempRoots.push(root);
const transcriptPath = path.join(root, 'no-ids.jsonl');

fs.writeFileSync(transcriptPath, [
makeUsageLine({ timestamp: '2026-01-01T10:00:00.000Z', input: 1, output: 10, stopReason: 'end_turn' }),
makeUsageLine({ timestamp: '2026-01-01T10:00:01.000Z', input: 1, output: 10, stopReason: 'end_turn' })
].join('\n'));

const metrics = await getTokenMetrics(transcriptPath);

expect(metrics.inputTokens).toBe(2);
expect(metrics.outputTokens).toBe(20);
});

it('keeps counting a still-streaming response once while its rows share an id', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-metrics-'));
tempRoots.push(root);
const transcriptPath = path.join(root, 'live-blocks.jsonl');

fs.writeFileSync(transcriptPath, [
makeUsageLine({ timestamp: '2026-01-01T10:00:00.000Z', input: 2, output: 50, stopReason: 'end_turn', messageId: 'msg_done' }),
makeUsageLine({ timestamp: '2026-01-01T10:00:05.000Z', input: 3, output: 10, stopReason: null, messageId: 'msg_live' }),
makeUsageLine({ timestamp: '2026-01-01T10:00:06.000Z', input: 3, output: 30, stopReason: null, messageId: 'msg_live' })
].join('\n'));

const metrics = await getTokenMetrics(transcriptPath);

expect(metrics.inputTokens).toBe(5);
expect(metrics.outputTokens).toBe(80);
});

it('counts the latest in-progress streaming entry once when no finalized row exists yet', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-metrics-'));
tempRoots.push(root);
Expand Down Expand Up @@ -831,6 +897,30 @@ describe('jsonl transcript metrics', () => {
});
});

it('counts one speed request per API response, not per content-block row', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-speed-'));
tempRoots.push(root);
const transcriptPath = path.join(root, 'speed-blocks.jsonl');

fs.writeFileSync(transcriptPath, [
makeTranscriptLine({ timestamp: '2026-01-01T10:00:00.000Z', type: 'user' }),
makeTranscriptLine({ timestamp: '2026-01-01T10:00:03.000Z', type: 'assistant', input: 200, output: 40, messageId: 'msg_a' }),
makeTranscriptLine({ timestamp: '2026-01-01T10:00:05.000Z', type: 'assistant', input: 200, output: 100, messageId: 'msg_a' }),
makeTranscriptLine({ timestamp: '2026-01-01T10:00:20.000Z', type: 'user' }),
makeTranscriptLine({ timestamp: '2026-01-01T10:00:24.000Z', type: 'assistant', input: 300, output: 150, messageId: 'msg_b' })
].join('\n'));

const metrics = await getSpeedMetrics(transcriptPath);

expect(metrics).toEqual({
totalDurationMs: 9000,
inputTokens: 500,
outputTokens: 250,
totalTokens: 750,
requestCount: 2
});
});

it('calculates windowed speed metrics from recent requests only', async () => {
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ccstatusline-jsonl-speed-'));
tempRoots.push(root);
Expand Down
92 changes: 77 additions & 15 deletions src/utils/jsonl-metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,12 +99,18 @@ interface TokenMetricAccumulator {
mostRecentPostCompactionTimestampMs: number | null;
}

interface PendingTokenMetricEntry {
messageId: string | null;
entry: TokenMetricEntry;
countable: boolean;
includePostCompactionUsage: boolean;
}

interface TokenMetricState {
metrics: TokenMetricAccumulator;
hasStopReasonField: boolean;
lastUsageEntry: TokenMetricEntry | null;
pending: PendingTokenMetricEntry | null;
sawCompactBoundary: boolean;
boundaryAfterLastUsage: boolean;
lastCompactBoundaryPostTokens: number | null;
}

Expand Down Expand Up @@ -169,18 +175,40 @@ function createTokenMetricState(): TokenMetricState {
return {
metrics: createTokenMetricAccumulator(),
hasStopReasonField: false,
lastUsageEntry: null,
pending: null,
sawCompactBoundary: false,
boundaryAfterLastUsage: false,
lastCompactBoundaryPostTokens: null
};
}

/**
* Claude Code writes one transcript row per content block (thinking, text,
* each tool_use) of a single API response. Those rows share `message.id`,
* carry identical prompt-side usage and a non-decreasing `output_tokens`, so
* only the last row of a message may be counted. Rows are held here until the
* next message starts (or the scan ends), which keeps the single pass over
* the file.
*/
function getMessageId(data: TranscriptLine | null): string | null {
const id = data?.message?.id;
return typeof id === 'string' && id.length > 0 ? id : null;
}

function flushPendingTokenMetricEntry(state: TokenMetricState): void {
const pending = state.pending;
state.pending = null;
if (pending?.countable) {
accumulateTokenMetricEntry(state.metrics, pending.entry, pending.includePostCompactionUsage);
}
}

function collectTokenMetricRecord(state: TokenMetricState, data: TranscriptLine | null, timestampMs: number | null): void {
const compactBoundary = isCompactBoundary(data);
if (compactBoundary) {
state.sawCompactBoundary = true;
state.boundaryAfterLastUsage = true;
if (state.pending) {
state.pending.includePostCompactionUsage = false;
}
state.lastCompactBoundaryPostTokens = getCompactBoundaryPostTokens(data);
resetPostCompactionUsage(state.metrics);
}
Expand All @@ -199,19 +227,30 @@ function collectTokenMetricRecord(state: TokenMetricState, data: TranscriptLine
if (hasStopReason && !state.hasStopReasonField) {
state.hasStopReasonField = true;
state.metrics = createTokenMetricAccumulator();
state.pending = null;
}
if (!state.hasStopReasonField || entry.stopReason) {
accumulateTokenMetricEntry(state.metrics, entry, !compactBoundary);

const messageId = getMessageId(data);
if (messageId === null || messageId !== state.pending?.messageId) {
flushPendingTokenMetricEntry(state);
}
state.lastUsageEntry = entry;
state.boundaryAfterLastUsage = compactBoundary;
// A later row of the same message supersedes the earlier one.
state.pending = {
messageId,
entry,
countable: !state.hasStopReasonField || Boolean(entry.stopReason),
includePostCompactionUsage: !compactBoundary
};
}
}

function finishTokenMetrics(state: TokenMetricState): TokenMetrics {
if (state.hasStopReasonField && state.lastUsageEntry?.stopReason === null) {
accumulateTokenMetricEntry(state.metrics, state.lastUsageEntry, !state.boundaryAfterLastUsage);
// The last message may still be streaming (stop_reason: null); count its
// latest row once so live updates reflect it.
if (state.pending && state.hasStopReasonField && state.pending.entry.stopReason === null) {
state.pending.countable = true;
}
flushPendingTokenMetricEntry(state);

const contextLengthFromUsage = (usage: UsageTokens | null): number | null => usage
? contextLengthFromUsageTokens(usage)
Expand Down Expand Up @@ -317,16 +356,30 @@ function normalizeWindowSeconds(value: number | undefined): number | null {
return normalized > 0 ? normalized : null;
}

interface SpeedMetricCollectorState extends CollectedSpeedMetrics { lastUserTimestampMs: number | null }
interface SpeedMetricCollectorState extends CollectedSpeedMetrics {
lastUserTimestampMs: number | null;
pendingMessageId: string | null;
pendingRequest: SpeedRequest | null;
}

function createSpeedMetricCollector(): SpeedMetricCollectorState {
return {
requests: [],
latestTimestampMs: null,
lastUserTimestampMs: null
lastUserTimestampMs: null,
pendingMessageId: null,
pendingRequest: null
};
}

function flushPendingSpeedRequest(state: SpeedMetricCollectorState): void {
if (state.pendingRequest) {
state.requests.push(state.pendingRequest);
}
state.pendingRequest = null;
state.pendingMessageId = null;
}

function collectSpeedMetricRecord(
state: SpeedMetricCollectorState,
data: TranscriptLine | null,
Expand Down Expand Up @@ -356,12 +409,19 @@ function collectSpeedMetricRecord(
}

const usage = parseUsageTokens(data.message.usage);
state.requests.push({
const messageId = getMessageId(data);
if (messageId === null || messageId !== state.pendingMessageId) {
flushPendingSpeedRequest(state);
}
// One request per API response: rows sharing a message id are content
// blocks of the same response, and the last one carries the full usage.
state.pendingMessageId = messageId;
state.pendingRequest = {
inputTokens: usage.input,
outputTokens: usage.output,
assistantTimestampMs: timestampMs,
interval
});
};
}

async function collectSpeedMetricsFromFile(filePath: string, ignoreSidechain: boolean): Promise<CollectedSpeedMetrics> {
Expand All @@ -370,6 +430,7 @@ async function collectSpeedMetricsFromFile(filePath: string, ignoreSidechain: bo
const data = parseJsonlLine(line) as TranscriptLine | null;
collectSpeedMetricRecord(state, data, parseTimestampMs(data?.timestamp), ignoreSidechain);
}
flushPendingSpeedRequest(state);

return state;
}
Expand Down Expand Up @@ -577,6 +638,7 @@ async function scanTranscript(transcriptPath: string, options: TranscriptScanOpt

let speedMetricsCollection: SpeedMetricsCollection | null = null;
if (speedState) {
flushPendingSpeedRequest(speedState);
const collected: CollectedSpeedMetrics[] = [speedState];
if (referencedAgentIds) {
const subagentPaths = getSubagentTranscriptPaths(transcriptPath, referencedAgentIds);
Expand Down