feat(usage): track cached tokens + correct input/output/cache cost (#2209)
Normalize every provider to one cache-inclusive convention via canonicalizeUsage() before persist, and price cached + cache_creation as subsets of prompt_tokens in calculateCostFromTokens() to stop double-counting. usageRepo now delegates cost math to a single source. Surface Cached tokens/cost across dashboard (overview, tokens, cost, details). Merge Claude message_start cache with message_delta output so cache counts survive. Compatible LLM nodes now allow multiple API-key connections (key pool). Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -391,6 +391,11 @@ export class KiroExecutor extends BaseExecutor {
|
||||
if (metrics && typeof metrics === 'object') {
|
||||
const inputTokens = metrics.inputTokens || 0;
|
||||
const outputTokens = metrics.outputTokens || 0;
|
||||
// ponytail: Amazon Q upstream does not expose cache fields today,
|
||||
// but pick up cache_read_input_tokens / cache_creation_input_tokens
|
||||
// if the event shape grows them so cost tracking stays accurate.
|
||||
const cachedTokens = metrics.cacheReadInputTokens || metrics.cache_read_input_tokens || 0;
|
||||
const cacheCreationInputTokens = metrics.cacheCreationInputTokens || metrics.cache_creation_input_tokens || 0;
|
||||
|
||||
if (inputTokens > 0 || outputTokens > 0) {
|
||||
state.usage = {
|
||||
@@ -398,6 +403,12 @@ export class KiroExecutor extends BaseExecutor {
|
||||
completion_tokens: outputTokens,
|
||||
total_tokens: inputTokens + outputTokens
|
||||
};
|
||||
// Kiro is Claude-backed: inputTokens EXCLUDES cache (Claude convention),
|
||||
// not inclusive like OpenAI's cached_tokens. Emit cache_read_input_tokens
|
||||
// (not cached_tokens) so canonicalizeUsage takes the Claude fold path and
|
||||
// correctly adds cache back into prompt_tokens instead of undercharging.
|
||||
if (cachedTokens > 0) state.usage.cache_read_input_tokens = cachedTokens;
|
||||
if (cacheCreationInputTokens > 0) state.usage.cache_creation_input_tokens = cacheCreationInputTokens;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { saveRequestUsage, appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js";
|
||||
import { COLORS } from "../../utils/stream.js";
|
||||
import { canonicalizeUsage } from "../../utils/usageTracking.js";
|
||||
|
||||
const OPTIONAL_PARAMS = [
|
||||
"temperature", "top_p", "top_k",
|
||||
@@ -48,7 +49,8 @@ export function extractUsageFromResponse(responseBody) {
|
||||
return {
|
||||
prompt_tokens: responseBody.usageMetadata.promptTokenCount || 0,
|
||||
completion_tokens: responseBody.usageMetadata.candidatesTokenCount || 0,
|
||||
reasoning_tokens: responseBody.usageMetadata.thoughtsTokenCount
|
||||
cached_tokens: responseBody.usageMetadata.cachedContentTokenCount || 0,
|
||||
reasoning_tokens: responseBody.usageMetadata.thoughtsTokenCount || 0
|
||||
};
|
||||
}
|
||||
|
||||
@@ -84,8 +86,9 @@ export function saveUsageStats({ provider, model, tokens, connectionId, apiKey,
|
||||
const accountSuffix = connectionId ? ` | account=${connectionId.slice(0, 8)}...` : "";
|
||||
console.log(`${COLORS.green}[${time}] 📊 [${label}] ${provider.toUpperCase()} | in=${inTokens} | out=${outTokens}${accountSuffix}${COLORS.reset}`);
|
||||
|
||||
// Normalize to OpenAI token shape for storage
|
||||
const normalized = {
|
||||
// Canonicalize to one storage convention (prompt_tokens cache-inclusive) so
|
||||
// cached/cache-creation tokens survive to cost calc + stats. See canonicalizeUsage.
|
||||
const normalized = canonicalizeUsage(tokens) || {
|
||||
prompt_tokens: tokens.prompt_tokens ?? tokens.input_tokens ?? 0,
|
||||
completion_tokens: tokens.completion_tokens ?? tokens.output_tokens ?? 0
|
||||
};
|
||||
|
||||
@@ -279,7 +279,10 @@ export function calculateCostFromTokens(tokens, pricing) {
|
||||
|
||||
const inputTokens = tokens.prompt_tokens || tokens.input_tokens || 0;
|
||||
const cachedTokens = tokens.cached_tokens || tokens.cache_read_input_tokens || 0;
|
||||
const nonCachedInput = Math.max(0, inputTokens - cachedTokens);
|
||||
const cacheCreationTokens = tokens.cache_creation_input_tokens || 0;
|
||||
// prompt_tokens is cache-inclusive (see canonicalizeUsage): cached + cache_creation
|
||||
// are subsets, so subtract both to avoid charging them at the full input rate.
|
||||
const nonCachedInput = Math.max(0, inputTokens - cachedTokens - cacheCreationTokens);
|
||||
|
||||
cost += nonCachedInput * (pricing.input / 1000000);
|
||||
|
||||
@@ -295,7 +298,6 @@ export function calculateCostFromTokens(tokens, pricing) {
|
||||
cost += reasoningTokens * ((pricing.reasoning || pricing.output) / 1000000);
|
||||
}
|
||||
|
||||
const cacheCreationTokens = tokens.cache_creation_input_tokens || 0;
|
||||
if (cacheCreationTokens > 0) {
|
||||
cost += cacheCreationTokens * ((pricing.cache_creation || pricing.input) / 1000000);
|
||||
}
|
||||
|
||||
@@ -39,7 +39,16 @@ const USAGE_EXTRACTORS = {
|
||||
},
|
||||
kiro(raw) {
|
||||
const input = n(raw.inputTokens), output = n(raw.outputTokens);
|
||||
return { promptTokens: input, completionTokens: output, totalTokens: input + output };
|
||||
// ponytail: Amazon Q (Kiro upstream) does not expose cache fields today,
|
||||
// but pass through any cache_read/cache_creation/cached_tokens if the
|
||||
// event shape grows them later so cost tracking keeps working without
|
||||
// a second pass.
|
||||
const cached = n(raw.cache_read_input_tokens) || n(raw.cachedTokens) || n(raw.cached_tokens);
|
||||
const cacheCreation = n(raw.cache_creation_input_tokens);
|
||||
const out = { promptTokens: input, completionTokens: output, totalTokens: input + output };
|
||||
if (cached > 0) out.cachedTokens = cached;
|
||||
if (cacheCreation > 0) out.cacheCreationTokens = cacheCreation;
|
||||
return out;
|
||||
},
|
||||
ollama(raw) {
|
||||
const input = n(raw.prompt_eval_count), output = n(raw.eval_count);
|
||||
|
||||
@@ -27,6 +27,25 @@ export function claudeToOpenAIResponse(chunk, state) {
|
||||
state.messageId = chunk.message?.id || `msg_${Date.now()}`;
|
||||
state.model = chunk.message?.model;
|
||||
state.toolCallIndex = 0;
|
||||
// Claude sends input_tokens + cache_read + cache_creation here; message_delta
|
||||
// later carries only the final output_tokens. Capture cache now so the
|
||||
// delta (output-only) doesn't reset it to zero.
|
||||
const startUsage = chunk.message?.usage;
|
||||
if (startUsage && typeof startUsage === "object") {
|
||||
const inputTokens = typeof startUsage.input_tokens === "number" ? startUsage.input_tokens : 0;
|
||||
const cacheReadTokens = typeof startUsage.cache_read_input_tokens === "number" ? startUsage.cache_read_input_tokens : 0;
|
||||
const cacheCreationTokens = typeof startUsage.cache_creation_input_tokens === "number" ? startUsage.cache_creation_input_tokens : 0;
|
||||
const promptTokens = inputTokens + cacheReadTokens + cacheCreationTokens;
|
||||
state.usage = {
|
||||
prompt_tokens: promptTokens,
|
||||
completion_tokens: 0,
|
||||
total_tokens: promptTokens,
|
||||
input_tokens: inputTokens,
|
||||
output_tokens: 0
|
||||
};
|
||||
if (cacheReadTokens > 0) state.usage.cache_read_input_tokens = cacheReadTokens;
|
||||
if (cacheCreationTokens > 0) state.usage.cache_creation_input_tokens = cacheCreationTokens;
|
||||
}
|
||||
results.push(createChunk(state, { role: ROLE.ASSISTANT }));
|
||||
break;
|
||||
}
|
||||
@@ -103,13 +122,15 @@ export function claudeToOpenAIResponse(chunk, state) {
|
||||
}
|
||||
|
||||
case "message_delta": {
|
||||
// Extract usage from message_delta event (Claude native format)
|
||||
// Normalize to OpenAI format (prompt_tokens/completion_tokens) for consistent logging
|
||||
// Extract usage from message_delta event (Claude native format).
|
||||
// Anthropic sends input/cache in message_start and only output here, so
|
||||
// fall back to cache captured in message_start when the delta omits it.
|
||||
if (chunk.usage && typeof chunk.usage === "object") {
|
||||
const inputTokens = typeof chunk.usage.input_tokens === "number" ? chunk.usage.input_tokens : 0;
|
||||
const prev = state.usage || {};
|
||||
const inputTokens = typeof chunk.usage.input_tokens === "number" ? chunk.usage.input_tokens : (prev.input_tokens || 0);
|
||||
const outputTokens = typeof chunk.usage.output_tokens === "number" ? chunk.usage.output_tokens : 0;
|
||||
const cacheReadTokens = typeof chunk.usage.cache_read_input_tokens === "number" ? chunk.usage.cache_read_input_tokens : 0;
|
||||
const cacheCreationTokens = typeof chunk.usage.cache_creation_input_tokens === "number" ? chunk.usage.cache_creation_input_tokens : 0;
|
||||
const cacheReadTokens = typeof chunk.usage.cache_read_input_tokens === "number" ? chunk.usage.cache_read_input_tokens : (prev.cache_read_input_tokens || 0);
|
||||
const cacheCreationTokens = typeof chunk.usage.cache_creation_input_tokens === "number" ? chunk.usage.cache_creation_input_tokens : (prev.cache_creation_input_tokens || 0);
|
||||
|
||||
// prompt_tokens = input_tokens + cache_read + cache_creation (all prompt-side tokens)
|
||||
const promptTokens = inputTokens + cacheReadTokens + cacheCreationTokens;
|
||||
@@ -131,7 +152,14 @@ export function claudeToOpenAIResponse(chunk, state) {
|
||||
const finalChunk = createChunk(state, {}, state.finishReason);
|
||||
|
||||
if (state.usage) {
|
||||
finalChunk.usage = toOpenAIUsage(chunk.usage, "claude");
|
||||
// Build OpenAI usage from the merged state (cache from message_start +
|
||||
// output from message_delta), not the delta chunk alone.
|
||||
finalChunk.usage = toOpenAIUsage({
|
||||
input_tokens: state.usage.input_tokens || 0,
|
||||
output_tokens: state.usage.output_tokens || 0,
|
||||
cache_read_input_tokens: state.usage.cache_read_input_tokens,
|
||||
cache_creation_input_tokens: state.usage.cache_creation_input_tokens
|
||||
}, "claude");
|
||||
}
|
||||
|
||||
results.push(finalChunk);
|
||||
|
||||
@@ -247,9 +247,24 @@ function mergeChunksToResponse(chunks, sourceFormat) {
|
||||
|
||||
if (messageStart?.message) {
|
||||
finalChunk = messageStart.message;
|
||||
// Merge usage if available
|
||||
if (messageDelta?.usage) {
|
||||
finalChunk.usage = messageDelta.usage;
|
||||
// message_start.usage has input + cache; message_delta.usage has the
|
||||
// final output_tokens. Merge so cache survives (delta omits it).
|
||||
const startUsage = messageStart.message.usage;
|
||||
const deltaUsage = messageDelta?.usage;
|
||||
if (startUsage || deltaUsage) {
|
||||
finalChunk.usage = {
|
||||
...(startUsage || {}),
|
||||
...(deltaUsage || {}),
|
||||
...(startUsage?.cache_read_input_tokens !== undefined
|
||||
? { cache_read_input_tokens: startUsage.cache_read_input_tokens }
|
||||
: {}),
|
||||
...(startUsage?.cache_creation_input_tokens !== undefined
|
||||
? { cache_creation_input_tokens: startUsage.cache_creation_input_tokens }
|
||||
: {}),
|
||||
...(startUsage?.input_tokens !== undefined
|
||||
? { input_tokens: startUsage.input_tokens }
|
||||
: {})
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { translateResponse, initState } from "../translator/index.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { trackPendingRequest, appendRequestLog } from "@/lib/usageDb.js";
|
||||
import { extractUsage, hasValidUsage, estimateUsage, logUsage, addBufferToUsage, filterUsageForFormat, COLORS } from "./usageTracking.js";
|
||||
import { extractUsage, mergeUsage, hasValidUsage, estimateUsage, logUsage, addBufferToUsage, filterUsageForFormat, COLORS } from "./usageTracking.js";
|
||||
import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./streamHelpers.js";
|
||||
import { getOpenAIResponsesEventName, isOpenAIResponsesTerminalEvent, formatIncompleteOpenAIResponsesStreamFailure } from "./responsesStreamHelpers.js";
|
||||
import { dbg, isDebugEnabled } from "./debugLog.js";
|
||||
@@ -162,7 +162,7 @@ export function createSSEStream(options = {}) {
|
||||
|
||||
const extracted = extractUsage(parsed);
|
||||
if (extracted) {
|
||||
usage = extracted;
|
||||
usage = mergeUsage(usage, extracted);
|
||||
}
|
||||
|
||||
const isFinishChunk = parsed.choices?.[0]?.finish_reason;
|
||||
@@ -280,7 +280,7 @@ export function createSSEStream(options = {}) {
|
||||
|
||||
// Extract usage
|
||||
const extracted = extractUsage(parsed);
|
||||
if (extracted) state.usage = extracted; // Keep original usage for logging
|
||||
if (extracted) state.usage = mergeUsage(state.usage, extracted); // Keep original usage for logging
|
||||
|
||||
// Responses same-format passthrough: re-emit with original event framing
|
||||
if (keepsOpenAIResponsesFormat && openAIResponsesEventName) {
|
||||
|
||||
@@ -141,6 +141,68 @@ export function normalizeUsage(usage) {
|
||||
return normalized;
|
||||
}
|
||||
|
||||
/**
|
||||
* Canonicalize usage into ONE storage/cost convention so token counts and cost
|
||||
* are consistent across providers:
|
||||
* prompt_tokens = total input INCLUDING cache read + cache creation
|
||||
* cached_tokens = cache-read portion (subset of prompt_tokens)
|
||||
* cache_creation_input_tokens = cache-write portion (subset of prompt_tokens)
|
||||
* completion_tokens, reasoning_tokens, total_tokens
|
||||
*
|
||||
* Discriminator: Claude reports cache_read_input_tokens with a prompt that
|
||||
* EXCLUDES cache, so we fold cache into prompt. OpenAI/Gemini report
|
||||
* cached_tokens already counted inside prompt, so we pass through. Idempotent:
|
||||
* once folded the output carries cached_tokens (not cache_read_input_tokens),
|
||||
* so re-running takes the passthrough branch and does not double-add.
|
||||
*
|
||||
* @param {object} usage - a normalizeUsage()-shaped object
|
||||
* @returns {object|null} canonical token object, or null for invalid input
|
||||
*/
|
||||
export function canonicalizeUsage(usage) {
|
||||
if (!usage || typeof usage !== "object" || Array.isArray(usage)) return null;
|
||||
|
||||
const num = (v) => (Number.isFinite(Number(v)) ? Number(v) : 0);
|
||||
const completion = num(usage.completion_tokens ?? usage.output_tokens);
|
||||
const reasoning = num(usage.reasoning_tokens);
|
||||
// Fall back to the nested prompt_tokens_details.cache_creation_tokens shape
|
||||
// (buildUsage()'s OpenAI-forwarding format) when the top-level field is
|
||||
// absent, so callers that pass a buildUsage() object through don't silently
|
||||
// drop cache_creation.
|
||||
const cacheCreation = num(usage.cache_creation_input_tokens ?? usage.prompt_tokens_details?.cache_creation_tokens);
|
||||
|
||||
let prompt = num(usage.prompt_tokens ?? usage.input_tokens);
|
||||
let cached;
|
||||
|
||||
// Claude path: prompt excludes cache; cache_read_input_tokens and/or
|
||||
// cache_creation_input_tokens are separate. A cache-miss "first write" only
|
||||
// carries cache_creation_input_tokens (no cache_read_input_tokens yet), so
|
||||
// check both fields — otherwise a first-write request falls through to the
|
||||
// OpenAI passthrough branch below and cache_creation never gets folded in.
|
||||
// Guard on the absence of `cached_tokens`: our own canonical output always
|
||||
// sets that key (even to 0), so re-running canonicalizeUsage on an already-
|
||||
// folded result takes the passthrough branch instead of folding again.
|
||||
if (usage.cached_tokens === undefined &&
|
||||
(usage.cache_read_input_tokens !== undefined || usage.cache_creation_input_tokens !== undefined)) {
|
||||
cached = num(usage.cache_read_input_tokens);
|
||||
prompt = prompt + cached + cacheCreation;
|
||||
} else {
|
||||
// OpenAI/Gemini path (or already-canonical input): prompt already includes cached_tokens.
|
||||
cached = num(usage.cached_tokens);
|
||||
}
|
||||
|
||||
const result = {
|
||||
prompt_tokens: prompt,
|
||||
completion_tokens: completion,
|
||||
// Recompute rather than pass through: when the fold branch ran above,
|
||||
// an upstream total_tokens (cache-exclusive) would otherwise be stale.
|
||||
total_tokens: prompt + completion,
|
||||
cached_tokens: cached,
|
||||
cache_creation_input_tokens: cacheCreation,
|
||||
};
|
||||
if (reasoning > 0) result.reasoning_tokens = reasoning;
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if usage has valid token data
|
||||
* Valid = has at least one token field with value > 0
|
||||
@@ -171,6 +233,19 @@ export function hasValidUsage(usage) {
|
||||
export function extractUsage(chunk) {
|
||||
if (!chunk || typeof chunk !== "object") return null;
|
||||
|
||||
// Claude format (message_start event): carries input_tokens + cache_read +
|
||||
// cache_creation. message_delta later carries only the final output_tokens,
|
||||
// so callers must MERGE (mergeUsage), not overwrite, to keep cache counts.
|
||||
if (chunk.type === "message_start" && chunk.message?.usage && typeof chunk.message.usage === "object") {
|
||||
const u = chunk.message.usage;
|
||||
return normalizeUsage({
|
||||
prompt_tokens: u.input_tokens || 0,
|
||||
completion_tokens: u.output_tokens || 0,
|
||||
cache_read_input_tokens: u.cache_read_input_tokens,
|
||||
cache_creation_input_tokens: u.cache_creation_input_tokens
|
||||
});
|
||||
}
|
||||
|
||||
// Claude format (message_delta event)
|
||||
if (chunk.type === "message_delta" && chunk.usage && typeof chunk.usage === "object") {
|
||||
return normalizeUsage({
|
||||
@@ -232,6 +307,27 @@ export function extractUsage(chunk) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Field-wise max-merge of two usage objects. Anthropic splits usage across
|
||||
// events: message_start has real input+cache (output is a placeholder 1),
|
||||
// message_delta has the real cumulative output (input/cache absent). Max keeps
|
||||
// the meaningful value from each without clobbering. Idempotent for other
|
||||
// providers that emit a single complete usage object.
|
||||
export function mergeUsage(prev, next) {
|
||||
if (!prev) return next || null;
|
||||
if (!next) return prev;
|
||||
const merged = { ...prev };
|
||||
for (const [k, v] of Object.entries(next)) {
|
||||
// typeof NaN === "number" — guard with Number.isFinite so one malformed
|
||||
// chunk can't poison the whole accumulation (Math.max(x, NaN) is NaN).
|
||||
if (typeof v === "number" && Number.isFinite(v)) {
|
||||
merged[k] = Math.max(typeof merged[k] === "number" ? merged[k] : 0, v);
|
||||
} else if (v && typeof v === "object") {
|
||||
merged[k] = v; // nested details objects: take latest
|
||||
}
|
||||
}
|
||||
return merged;
|
||||
}
|
||||
|
||||
/**
|
||||
* Estimate input tokens from request body
|
||||
* Calculate total body size for more accurate estimation
|
||||
|
||||
Reference in New Issue
Block a user