From 4ba7b7a18c22813b9e02e77b410138056d47910b Mon Sep 17 00:00:00 2001 From: z1pp090 <286964464+z1pp090@users.noreply.github.com> Date: Sun, 20 Sep 2026 11:23:42 +0200 Subject: [PATCH] fix(metrics): count each API response once when its rows share a message id Claude Code writes one transcript row per content block of an assistant response (thinking, text, each tool_use). The rows share `message.id`, carry identical prompt-side usage and a non-decreasing `output_tokens`, and all of them already hold the finalized `stop_reason`. The stop_reason gate only catches streaming partials (stop_reason: null), so every one of those rows was summed: tokens-input, tokens-output, tokens-cached and tokens-total over-reported by roughly 2x, and the speed widgets counted one request per row. Hold the latest row of the current message id while scanning and count it when the next message starts or the scan ends. Rows without an id keep the old per-row behaviour, the in-progress streaming row is still counted once, and a compaction boundary that lands while a row is pending keeps it out of the post-compaction context the way it did before. Measured on 23 local transcripts: cache read 3.15G -> 1.84G, output 12.1M -> 5.58M, input 55.1k -> 27.8k. Closes #549 Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01StkifHpSLWH4J5TvuP56ur --- src/types/TokenMetrics.ts | 2 +- src/utils/__tests__/jsonl-metrics.test.ts | 90 ++++++++++++++++++++++ src/utils/jsonl-metrics.ts | 92 +++++++++++++++++++---- 3 files changed, 168 insertions(+), 16 deletions(-) diff --git a/src/types/TokenMetrics.ts b/src/types/TokenMetrics.ts index 042385d92..403e9b5f4 100644 --- a/src/types/TokenMetrics.ts +++ b/src/types/TokenMetrics.ts @@ -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; diff --git a/src/utils/__tests__/jsonl-metrics.test.ts b/src/utils/__tests__/jsonl-metrics.test.ts index 348ea432a..a5fb1f0c7 100644 --- a/src/utils/__tests__/jsonl-metrics.test.ts +++ b/src/utils/__tests__/jsonl-metrics.test.ts @@ -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, @@ -93,6 +95,7 @@ function makeTranscriptLine(params: { output?: number; isSidechain?: boolean; isApiErrorMessage?: boolean; + messageId?: string; }): string { return JSON.stringify({ timestamp: params.timestamp, @@ -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 @@ -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); @@ -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); diff --git a/src/utils/jsonl-metrics.ts b/src/utils/jsonl-metrics.ts index 88f4c47d8..d014e9790 100644 --- a/src/utils/jsonl-metrics.ts +++ b/src/utils/jsonl-metrics.ts @@ -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; } @@ -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); } @@ -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) @@ -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, @@ -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 { @@ -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; } @@ -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);