fix: requestId dedup to prevent ~2x cost overcounting + stale session detection
Manually ported from upstream: - ef2e0868: deduplicate streaming JSONL entries by requestId - 4f21f267: mark stale ongoing sessions as dead after 5min inactivity
This commit is contained in:
parent
e6fa5c6f06
commit
b87082a915
3 changed files with 51 additions and 3 deletions
|
|
@ -984,6 +984,12 @@ export class ProjectScanner {
|
||||||
? firstMessageTimestampMs
|
? firstMessageTimestampMs
|
||||||
: birthtimeMs;
|
: birthtimeMs;
|
||||||
|
|
||||||
|
// If messages suggest ongoing but the file hasn't been written to in 5+ minutes,
|
||||||
|
// the session likely crashed/was killed (upstream fix #94)
|
||||||
|
const STALE_SESSION_THRESHOLD_MS = 5 * 60 * 1000;
|
||||||
|
const isOngoing =
|
||||||
|
metadata.isOngoing && Date.now() - effectiveMtime < STALE_SESSION_THRESHOLD_MS;
|
||||||
|
|
||||||
return {
|
return {
|
||||||
id: sessionId,
|
id: sessionId,
|
||||||
projectId,
|
projectId,
|
||||||
|
|
@ -993,7 +999,7 @@ export class ProjectScanner {
|
||||||
messageTimestamp: metadata.firstUserMessage?.timestamp,
|
messageTimestamp: metadata.firstUserMessage?.timestamp,
|
||||||
hasSubagents,
|
hasSubagents,
|
||||||
messageCount: metadata.messageCount,
|
messageCount: metadata.messageCount,
|
||||||
isOngoing: metadata.isOngoing,
|
isOngoing,
|
||||||
gitBranch: metadata.gitBranch ?? undefined,
|
gitBranch: metadata.gitBranch ?? undefined,
|
||||||
metadataLevel,
|
metadataLevel,
|
||||||
contextConsumption: metadata.contextConsumption,
|
contextConsumption: metadata.contextConsumption,
|
||||||
|
|
|
||||||
|
|
@ -103,6 +103,8 @@ export interface ParsedMessage {
|
||||||
toolUseResult?: ToolUseResultData;
|
toolUseResult?: ToolUseResultData;
|
||||||
/** Whether this is a compact summary boundary message */
|
/** Whether this is a compact summary boundary message */
|
||||||
isCompactSummary?: boolean;
|
isCompactSummary?: boolean;
|
||||||
|
/** API request ID for deduplicating streaming entries */
|
||||||
|
requestId?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
// =============================================================================
|
// =============================================================================
|
||||||
|
|
|
||||||
|
|
@ -125,6 +125,7 @@ function parseChatHistoryEntry(entry: ChatHistoryEntry): ParsedMessage | null {
|
||||||
let role: string | undefined;
|
let role: string | undefined;
|
||||||
let usage: TokenUsage | undefined;
|
let usage: TokenUsage | undefined;
|
||||||
let model: string | undefined;
|
let model: string | undefined;
|
||||||
|
let requestId: string | undefined;
|
||||||
let cwd: string | undefined;
|
let cwd: string | undefined;
|
||||||
let gitBranch: string | undefined;
|
let gitBranch: string | undefined;
|
||||||
let agentId: string | undefined;
|
let agentId: string | undefined;
|
||||||
|
|
@ -163,6 +164,7 @@ function parseChatHistoryEntry(entry: ChatHistoryEntry): ParsedMessage | null {
|
||||||
usage = entry.message.usage;
|
usage = entry.message.usage;
|
||||||
model = entry.message.model;
|
model = entry.message.model;
|
||||||
agentId = entry.agentId;
|
agentId = entry.agentId;
|
||||||
|
requestId = entry.requestId;
|
||||||
} else if (entry.type === 'system') {
|
} else if (entry.type === 'system') {
|
||||||
isMeta = entry.isMeta ?? false;
|
isMeta = entry.isMeta ?? false;
|
||||||
}
|
}
|
||||||
|
|
@ -195,6 +197,7 @@ function parseChatHistoryEntry(entry: ChatHistoryEntry): ParsedMessage | null {
|
||||||
sourceToolUseID,
|
sourceToolUseID,
|
||||||
sourceToolAssistantUUID,
|
sourceToolAssistantUUID,
|
||||||
toolUseResult,
|
toolUseResult,
|
||||||
|
requestId,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -221,24 +224,61 @@ function parseMessageType(type?: string): MessageType | null {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// =============================================================================
|
||||||
|
// Streaming Deduplication
|
||||||
|
// =============================================================================
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Deduplicate streaming assistant entries by requestId.
|
||||||
|
*
|
||||||
|
* Claude Code writes multiple JSONL entries per API response during streaming,
|
||||||
|
* each with the same requestId but incrementally increasing output_tokens.
|
||||||
|
* Only the last entry per requestId has the final, complete token counts.
|
||||||
|
*
|
||||||
|
* Messages without a requestId (user, system, etc.) pass through unchanged.
|
||||||
|
* Returns a new array with only the last entry per requestId kept.
|
||||||
|
*/
|
||||||
|
export function deduplicateByRequestId(messages: ParsedMessage[]): ParsedMessage[] {
|
||||||
|
const lastIndexByRequestId = new Map<string, number>();
|
||||||
|
for (let i = 0; i < messages.length; i++) {
|
||||||
|
const rid = messages[i].requestId;
|
||||||
|
if (rid) {
|
||||||
|
lastIndexByRequestId.set(rid, i);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (lastIndexByRequestId.size === 0) {
|
||||||
|
return messages;
|
||||||
|
}
|
||||||
|
|
||||||
|
return messages.filter((msg, i) => {
|
||||||
|
if (!msg.requestId) return true;
|
||||||
|
return lastIndexByRequestId.get(msg.requestId) === i;
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
// =============================================================================
|
// =============================================================================
|
||||||
// Metrics Calculation
|
// Metrics Calculation
|
||||||
// =============================================================================
|
// =============================================================================
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Calculate session metrics from parsed messages.
|
* Calculate session metrics from parsed messages.
|
||||||
|
* Deduplicates streaming entries by requestId before summing to avoid ~2x cost overcounting.
|
||||||
*/
|
*/
|
||||||
export function calculateMetrics(messages: ParsedMessage[]): SessionMetrics {
|
export function calculateMetrics(messages: ParsedMessage[]): SessionMetrics {
|
||||||
if (messages.length === 0) {
|
if (messages.length === 0) {
|
||||||
return { ...EMPTY_METRICS };
|
return { ...EMPTY_METRICS };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Deduplicate streaming entries: keep only the last entry per requestId
|
||||||
|
const dedupedMessages = deduplicateByRequestId(messages);
|
||||||
|
|
||||||
let inputTokens = 0;
|
let inputTokens = 0;
|
||||||
let outputTokens = 0;
|
let outputTokens = 0;
|
||||||
let cacheReadTokens = 0;
|
let cacheReadTokens = 0;
|
||||||
let cacheCreationTokens = 0;
|
let cacheCreationTokens = 0;
|
||||||
|
|
||||||
// Get timestamps for duration (loop instead of Math.min/max spread to avoid stack overflow on large sessions)
|
// Get timestamps for duration from ALL messages (not deduped) for accurate session length
|
||||||
const timestamps = messages.map((m) => m.timestamp.getTime()).filter((t) => !isNaN(t));
|
const timestamps = messages.map((m) => m.timestamp.getTime()).filter((t) => !isNaN(t));
|
||||||
|
|
||||||
let minTime = 0;
|
let minTime = 0;
|
||||||
|
|
@ -255,7 +295,7 @@ export function calculateMetrics(messages: ParsedMessage[]): SessionMetrics {
|
||||||
// Calculate cost per-message, then sum (tiered pricing applies per-API-call, not to aggregated totals)
|
// Calculate cost per-message, then sum (tiered pricing applies per-API-call, not to aggregated totals)
|
||||||
let costUsd = 0;
|
let costUsd = 0;
|
||||||
|
|
||||||
for (const msg of messages) {
|
for (const msg of dedupedMessages) {
|
||||||
if (msg.usage) {
|
if (msg.usage) {
|
||||||
const msgInputTokens = msg.usage.input_tokens ?? 0;
|
const msgInputTokens = msg.usage.input_tokens ?? 0;
|
||||||
const msgOutputTokens = msg.usage.output_tokens ?? 0;
|
const msgOutputTokens = msg.usage.output_tokens ?? 0;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue