This commit is contained in:
2026-07-06 00:01:34 +07:00
parent 2729408ef3
commit e7470e955e
108 changed files with 4131 additions and 529 deletions

View File

@@ -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 }
: {})
};
}
}
}

View File

@@ -5,6 +5,7 @@ import { formatSSE } from "./streamHelpers.js";
// Responses API events that signal the stream has reached a terminal state
const OPENAI_RESPONSES_TERMINAL_EVENTS = new Set([
"response.completed",
"response.done",
"response.failed",
"error"
]);

View File

@@ -1,8 +1,12 @@
import { translateResponse, initState } from "../translator/index.js";
import { FORMATS } from "../translator/formats.js";
import { trackPendingRequest, appendRequestLog } from "@/lib/usageDb.js";
<<<<<<< HEAD
import { extractUsage, hasValidUsage, estimateUsage, addBufferToUsage, filterUsageForFormat, COLORS } from "./usageTracking.js";
import { saveUsageStats } from "../handlers/chatCore/requestDetail.js";
=======
import { extractUsage, mergeUsage, hasValidUsage, estimateUsage, logUsage, addBufferToUsage, filterUsageForFormat, COLORS } from "./usageTracking.js";
>>>>>>> 7f436e2792be4fa5a4d1c4d6b8e9bc85eaaa6a3d
import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./streamHelpers.js";
import { getOpenAIResponsesEventName, isOpenAIResponsesTerminalEvent, formatIncompleteOpenAIResponsesStreamFailure } from "./responsesStreamHelpers.js";
import { dbg, isDebugEnabled } from "./debugLog.js";
@@ -131,6 +135,20 @@ export function createSSEStream(options = {}) {
}
}
// Strip empty tool_calls arrays that break AI SDK reasoning tracking.
// Some providers (e.g. CodeBuddy CN) include `"tool_calls": []` in
// every streaming delta. @ai-sdk/openai-compatible checks
// `delta.tool_calls != null` — an empty array passes this check,
// causing premature `reasoning-end` on every chunk.
if (parsed?.choices) {
for (const choice of parsed.choices) {
if (choice.delta?.tool_calls && Array.isArray(choice.delta.tool_calls) && choice.delta.tool_calls.length === 0) {
delete choice.delta.tool_calls;
fieldsInjected = true;
}
}
}
if (!hasValuableContent(parsed, FORMATS.OPENAI)) {
continue;
}
@@ -149,7 +167,7 @@ export function createSSEStream(options = {}) {
const extracted = extractUsage(parsed);
if (extracted) {
usage = extracted;
usage = mergeUsage(usage, extracted);
}
const isFinishChunk = parsed.choices?.[0]?.finish_reason;
@@ -218,9 +236,11 @@ export function createSSEStream(options = {}) {
sseEmittedCount++;
}
// [DONE] not emitted in translate mode — some clients' SSE decoders
// fail to parse the OpenAI sentinel on Claude-format translated streams.
// message_stop already signals end-of-response; stream close handles it.
if (keepsOpenAIResponsesFormat && !streamDoneSent) {
const doneOutput = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(doneOutput);
controller.enqueue(sharedEncoder.encode(doneOutput));
}
streamDoneSent = true;
if (keepsOpenAIResponsesFormat) openAIResponsesDoneSent = true;
continue;
@@ -265,7 +285,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) {
@@ -351,7 +371,9 @@ export function createSSEStream(options = {}) {
// Some clients (e.g. OpenClaw) expect the OpenAI-style sentinel:
// data: [DONE]\n\n
// Without it they can hang until timeout and trigger failover.
if (!streamDoneSent) {
// Gemini-family clients (Antigravity, Vertex, Gemini) reject this sentinel with 400 syntax errors.
const isGeminiFamily = provider === "antigravity" || provider === "gemini" || provider === "vertex";
if (!streamDoneSent && !isGeminiFamily) {
const doneOutput = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(doneOutput);
controller.enqueue(sharedEncoder.encode(doneOutput));
@@ -416,8 +438,13 @@ export function createSSEStream(options = {}) {
openAIResponsesTerminalSeen = true;
}
// [DONE] not emitted in translate mode — see comment above.
// Passthrough mode still emits it for standard OpenAI clients.
if (keepsOpenAIResponsesFormat && !openAIResponsesDoneSent && !streamDoneSent) {
const doneOutput = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(doneOutput);
controller.enqueue(sharedEncoder.encode(doneOutput));
openAIResponsesDoneSent = true;
streamDoneSent = true;
}
if (!hasValidUsage(state?.usage) && totalContentLength > 0) {
state.usage = estimateUsage(body, totalContentLength, sourceFormat);

View File

@@ -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