diff --git a/src/types/TokenMetrics.ts b/src/types/TokenMetrics.ts index 042385d9..403e9b5f 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 348ea432..a5fb1f0c 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 88f4c47d..d014e979 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);