diff --git a/src/controllers/anthropic.js b/src/controllers/anthropic.js index d87e7b7d..d1f6344c 100644 --- a/src/controllers/anthropic.js +++ b/src/controllers/anthropic.js @@ -1,5 +1,5 @@ const { isJson, generateUUID } = require('../utils/tools.js'); -const { createUsageObject } = require('../utils/precise-tokenizer.js'); +const { createUsageObject, mergeUpstreamUsage, reportUsage } = require('../utils/precise-tokenizer.js'); const { sendChatRequest, invalidateContextPrefix } = require('../utils/request.js'); const { buildContextPrefixKey } = require('../utils/context-prefix-cache.js'); const accountManager = require('../utils/account.js'); @@ -1252,8 +1252,8 @@ const handleAnthropicStream = async (res, ctx, upstream) => { usage: { input_tokens: 0, output_tokens: 0, - cache_creation_input_tokens: null, - cache_read_input_tokens: null + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0 } } }); @@ -1264,6 +1264,7 @@ const handleAnthropicStream = async (res, ctx, upstream) => { let thinkingSignature = null; let promptTokens = 0; let completionTokens = 0; + let upstreamUsage = null; // 上游逐帧累计的 usage(DashScope 命名已归一化;null = 还没报) let upstreamFinishReason = null; let upstreamCompleted; let upstreamEventCount; @@ -1346,6 +1347,8 @@ const handleAnthropicStream = async (res, ctx, upstream) => { // 但本控制器没有 Agent 回合门禁去解包,标签会原样发给客户端。剥掉它们。 agentTagStripper = createAgentTagStripper(); recoveredBuffer = ''; + // usage 也按轮全新:报的是最后一轮上游给的,没给就估算,绝不继承上一轮的。 + upstreamUsage = null; attemptVisibleText = ''; attemptThinkText = ''; attemptThinkEvidence = false; @@ -1532,10 +1535,8 @@ const handleAnthropicStream = async (res, ctx, upstream) => { const onUpstreamDelta = async (json) => { // 丢弃其余候选回答的帧:上游多路并发会让内容重复 if (!acceptUpstreamFrame(json)) return; - if (json.usage) { - promptTokens = json.usage.prompt_tokens || promptTokens; - completionTokens = json.usage.completion_tokens || completionTokens; - } + // Qwen 的 usage 用 DashScope 命名(input_tokens/output_tokens),每个 typing 帧带累计值 + upstreamUsage = mergeUpstreamUsage(upstreamUsage, json.usage); if (!json.choices || json.choices.length === 0) return; const choice = json.choices[0]; const reportedFinishReason = choice.finish_reason ?? choice.delta?.finish_reason; @@ -1998,11 +1999,10 @@ const handleAnthropicStream = async (res, ctx, upstream) => { return; } - if (promptTokens === 0 && completionTokens === 0) { - const usage = createUsageObject(requestBody?.messages || '', completionContent, null); - promptTokens = usage.prompt_tokens || 0; - completionTokens = usage.completion_tokens || 0; - } + // 只对上游没报的字段补本地估算(早停的回合收不到尾部 usage 帧) + const usage = reportUsage(upstreamUsage, () => createUsageObject(requestBody?.messages || '', completionContent), 'ANTHROPIC'); + promptTokens = usage.prompt_tokens; + completionTokens = usage.completion_tokens; // Daily stats 累计——一次性归属主账户(见模块顶部 attributeChatUsage 注释) attributeChatUsage(ctx.currentAccount, promptTokens, completionTokens); @@ -2013,8 +2013,8 @@ const handleAnthropicStream = async (res, ctx, upstream) => { usage: { input_tokens: promptTokens, output_tokens: completionTokens, - cache_creation_input_tokens: null, - cache_read_input_tokens: null + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0 } }); writeAnthropicEvent(res, 'message_stop', { type: 'message_stop' }); @@ -2043,6 +2043,7 @@ const handleAnthropicNonStream = async (res, ctx, upstream) => { let answerContent = ''; let promptTokens = 0; let completionTokens = 0; + let upstreamUsage = null; // 上游逐帧累计的 usage(DashScope 命名已归一化;null = 还没报) let webSearchInfo = null; let upstreamFinishReason = null; let upstreamCompleted; @@ -2126,10 +2127,8 @@ const handleAnthropicNonStream = async (res, ctx, upstream) => { const onUpstreamDelta = async (json) => { // 丢弃其余候选回答的帧:上游多路并发会让内容重复 if (!acceptUpstreamFrame(json)) return; - if (json.usage) { - promptTokens = json.usage.prompt_tokens || promptTokens; - completionTokens = json.usage.completion_tokens || completionTokens; - } + // Qwen 的 usage 用 DashScope 命名(input_tokens/output_tokens),每个 typing 帧带累计值 + upstreamUsage = mergeUpstreamUsage(upstreamUsage, json.usage); if (!json.choices || json.choices.length === 0) return; const choice = json.choices[0]; const reportedFinishReason = choice.finish_reason ?? choice.delta?.finish_reason; @@ -2446,6 +2445,8 @@ const handleAnthropicNonStream = async (res, ctx, upstream) => { // 判定输入按轮清零(thinkingContent 本身继续累计 —— 响应交付语义不动)。 attemptThinkingContent = ''; upstreamFinishReason = null; + // usage 也按轮全新:报的是最后一轮上游给的,没给就估算,绝不继承上一轮的。 + upstreamUsage = null; const retryResult = await consumeUpstream(retryResp.response, onUpstreamDelta, { shouldStop: () => stopRequested }); upstreamCompleted = retryResult.completed; if (!upstreamCompleted && !upstreamFinishReason) { @@ -2570,13 +2571,14 @@ const handleAnthropicNonStream = async (res, ctx, upstream) => { }); } - if (promptTokens === 0 && completionTokens === 0) { - // 早停的回合收不到上游尾部的 usage 帧:原生调用的参数 JSON 也进本地估算,免得 ~0。 + // 只对上游没报的字段补本地估算。早停的回合收不到上游尾部的 usage 帧: + // 原生调用的参数 JSON 也进本地估算,免得 ~0。 + const usage = reportUsage(upstreamUsage, () => { const nativeArgsText = nativeToolCalls.map(call => call.function.arguments || '').join(''); - const usage = createUsageObject(requestBody?.messages || '', thinkingContent + answerContent + nativeArgsText, null); - promptTokens = usage.prompt_tokens || 0; - completionTokens = usage.completion_tokens || 0; - } + return createUsageObject(requestBody?.messages || '', thinkingContent + answerContent + nativeArgsText); + }, 'ANTHROPIC'); + promptTokens = usage.prompt_tokens; + completionTokens = usage.completion_tokens; const contentBlocks = []; if (thinkingContent && thinkingContent.trim()) { @@ -2618,8 +2620,8 @@ const handleAnthropicNonStream = async (res, ctx, upstream) => { usage: { input_tokens: promptTokens, output_tokens: completionTokens, - cache_creation_input_tokens: null, - cache_read_input_tokens: null + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0 } }); }; diff --git a/src/controllers/chat.js b/src/controllers/chat.js index 6d26fd46..ff475a87 100644 --- a/src/controllers/chat.js +++ b/src/controllers/chat.js @@ -1,5 +1,5 @@ const { isJson, generateUUID } = require('../utils/tools.js') -const { createUsageObject } = require('../utils/precise-tokenizer.js') +const { createUsageObject, mergeUpstreamUsage, reportUsage } = require('../utils/precise-tokenizer.js') const { sendChatRequest } = require('../utils/request.js') const { buildContextPrefixKey } = require('../utils/context-prefix-cache.js') const { @@ -280,14 +280,9 @@ const runWithSSEHeartbeat = async (res, work, intervalMs = 15000) => { } const normalizeAgentUsage = (attempt, requestBody, completionText) => { - let usage = { ...(attempt?.totalTokens || {}) } - if (!usage.prompt_tokens && !usage.completion_tokens) { - usage = createUsageObject(requestBody?.messages || [], completionText, null) - } - usage.prompt_tokens = Math.max(0, Number(usage.prompt_tokens) || 0) - usage.completion_tokens = Math.max(0, Number(usage.completion_tokens) || 0) - usage.total_tokens = usage.prompt_tokens + usage.completion_tokens - return usage + // attempt.upstreamUsage:runtime 逐帧累计的上游 usage(DashScope 命名已归一化;null = 没报)。 + // 只对上游没报的字段补本地估算。 + return reportUsage(attempt?.upstreamUsage ?? null, () => createUsageObject(requestBody?.messages || [], completionText), 'CHAT') } /** @@ -662,6 +657,7 @@ const handleStreamResponse = async (res, response, enable_thinking, enable_web_s completion_tokens: 0, total_tokens: 0 } + let upstreamUsage = null // 上游逐帧累计的 usage(DashScope 命名已归一化;null = 还没报) let completionContent = '' // 收集完整的回复内容用于token估算 let visibleContent = '' @@ -810,13 +806,8 @@ const handleStreamResponse = async (res, response, enable_thinking, enable_web_s // 丢弃其余候选回答的帧:上游多路并发会让内容重复 if (!acceptUpstreamFrame(decodeJson)) return - if (decodeJson.usage) { - totalTokens = { - prompt_tokens: decodeJson.usage.prompt_tokens || totalTokens.prompt_tokens, - completion_tokens: decodeJson.usage.completion_tokens || totalTokens.completion_tokens, - total_tokens: decodeJson.usage.total_tokens || totalTokens.total_tokens - } - } + // Qwen 的 usage 用 DashScope 命名(input_tokens/output_tokens),每个 typing 帧带累计值 + upstreamUsage = mergeUpstreamUsage(upstreamUsage, decodeJson.usage) if (!decodeJson.choices || decodeJson.choices.length === 0) return @@ -1055,17 +1046,8 @@ const handleStreamResponse = async (res, response, enable_thinking, enable_web_s writeContentDelta(`\n\n---\n${webSearchTable}`) } - // 计算最终的token使用量 - if (totalTokens.prompt_tokens === 0 && totalTokens.completion_tokens === 0) { - totalTokens = createUsageObject(requestBody?.messages || promptText, completionContent, null) - logger.info(`流式使用tiktoken计算 - Prompt: ${totalTokens.prompt_tokens}, Completion: ${totalTokens.completion_tokens}, Total: ${totalTokens.total_tokens}`, 'CHAT') - } else { - logger.info(`流式使用上游真实Token - Prompt: ${totalTokens.prompt_tokens}, Completion: ${totalTokens.completion_tokens}, Total: ${totalTokens.total_tokens}`, 'CHAT') - } - - totalTokens.prompt_tokens = Math.max(0, totalTokens.prompt_tokens || 0) - totalTokens.completion_tokens = Math.max(0, totalTokens.completion_tokens || 0) - totalTokens.total_tokens = totalTokens.prompt_tokens + totalTokens.completion_tokens + // 计算最终的token使用量:只对上游没报的字段补本地估算 + totalTokens = reportUsage(upstreamUsage, () => createUsageObject(requestBody?.messages || promptText, completionContent), 'CHAT') // Daily stats 累计——一次性归属到主请求账户 // 注:tool_choice=required retry 走的可能是另一个账户,但 retry 路径罕见, @@ -1181,6 +1163,7 @@ const handleNonStreamResponse = async (res, response, enable_thinking, enable_we completion_tokens: 0, total_tokens: 0 } + let upstreamUsage = null // 上游逐帧累计的 usage(DashScope 命名已归一化;null = 还没报) // 提取prompt文本用于token估算 let promptText = '' @@ -1207,13 +1190,8 @@ const handleNonStreamResponse = async (res, response, enable_thinking, enable_we // 丢弃其余候选回答的帧:上游多路并发会让内容重复 if (!acceptUpstreamFrame(decodeJson)) return - if (decodeJson.usage) { - totalTokens = { - prompt_tokens: decodeJson.usage.prompt_tokens || totalTokens.prompt_tokens, - completion_tokens: decodeJson.usage.completion_tokens || totalTokens.completion_tokens, - total_tokens: decodeJson.usage.total_tokens || totalTokens.total_tokens - } - } + // Qwen 的 usage 用 DashScope 命名(input_tokens/output_tokens),每个 typing 帧带累计值 + upstreamUsage = mergeUpstreamUsage(upstreamUsage, decodeJson.usage) if (!decodeJson.choices || decodeJson.choices.length === 0) return const choice = decodeJson.choices[0] @@ -1424,17 +1402,9 @@ const handleNonStreamResponse = async (res, response, enable_thinking, enable_we assistantContent += `\n\n---\n${webSearchTable}` } - // 计算最终的token使用量(推理内容计入 completion,与 DeepSeek 一致;旧版 fullReasoning 为空) - if (totalTokens.prompt_tokens === 0 && totalTokens.completion_tokens === 0) { - totalTokens = createUsageObject(requestBody?.messages || promptText, fullReasoning + fullContent, null) - logger.info(`非流式使用tiktoken计算 - Prompt: ${totalTokens.prompt_tokens}, Completion: ${totalTokens.completion_tokens}, Total: ${totalTokens.total_tokens}`, 'CHAT') - } else { - logger.info(`非流式使用上游真实Token - Prompt: ${totalTokens.prompt_tokens}, Completion: ${totalTokens.completion_tokens}, Total: ${totalTokens.total_tokens}`, 'CHAT') - } - - totalTokens.prompt_tokens = Math.max(0, totalTokens.prompt_tokens || 0) - totalTokens.completion_tokens = Math.max(0, totalTokens.completion_tokens || 0) - totalTokens.total_tokens = totalTokens.prompt_tokens + totalTokens.completion_tokens + // 计算最终的token使用量:只对上游没报的字段补本地估算 + //(推理内容计入 completion,与 DeepSeek 一致;旧版 fullReasoning 为空) + totalTokens = reportUsage(upstreamUsage, () => createUsageObject(requestBody?.messages || promptText, fullReasoning + fullContent), 'CHAT') // Daily stats 累计——一次性归属到主请求账户(同 stream 分支注释) attributeChatUsage(options.currentAccount, totalTokens) diff --git a/src/utils/openai-agent-runtime.js b/src/utils/openai-agent-runtime.js index 600afad7..c3ea6a75 100644 --- a/src/utils/openai-agent-runtime.js +++ b/src/utils/openai-agent-runtime.js @@ -7,6 +7,7 @@ const { ANSWER_PHASES } = require('./tool-prompt.js') const { consumeSSEStream, createUpstreamResponseFilter } = require('./sse.js') +const { mergeUpstreamUsage } = require('./precise-tokenizer.js') const { createUpstreamDeltaNormalizer, createClientToolNamePredicate } = require('./chat-helpers.js') const { assertNoUpstreamFailure, UpstreamResponseError, isRateLimitError, isWafChallengeError } = require('./upstream-error.js') const { recordFailedAccount, createAccountReplayBody } = require('./agent-account-failover.js') @@ -324,11 +325,7 @@ const collectOpenAIAgentAttempt = async (upstreamResponse, options = {}) => { let lastCreated = null const emittedImages = new Set() const pendingImages = [] - let totalTokens = { - prompt_tokens: 0, - completion_tokens: 0, - total_tokens: 0 - } + let upstreamUsage = null // 上游逐帧累计的 usage(DashScope 命名已归一化;null = 还没报) const streamResult = await consumeSSEStream(upstreamResponse, async (frame) => { if (!frame.data || frame.data.trim() === '[DONE]') return @@ -352,13 +349,8 @@ const collectOpenAIAgentAttempt = async (upstreamResponse, options = {}) => { if (!acceptUpstreamFrame(decoded)) return if (decoded.response_id) acceptedResponseId = decoded.response_id - if (decoded.usage) { - totalTokens = { - prompt_tokens: decoded.usage.prompt_tokens || totalTokens.prompt_tokens, - completion_tokens: decoded.usage.completion_tokens || totalTokens.completion_tokens, - total_tokens: decoded.usage.total_tokens || totalTokens.total_tokens - } - } + // Qwen 的 usage 用 DashScope 命名(input_tokens/output_tokens),每个 typing 帧带累计值 + upstreamUsage = mergeUpstreamUsage(upstreamUsage, decoded.usage) if (!Array.isArray(decoded.choices) || decoded.choices.length === 0) return const choice = decoded.choices[0] @@ -584,7 +576,8 @@ const collectOpenAIAgentAttempt = async (upstreamResponse, options = {}) => { // 门禁靠它识别"原生调用被平台吃掉、只剩叙述"的死亡回合。 interceptedToolNames: normalizeDelta.interceptedToolNames, webSearchInfo, - totalTokens, + // 上游逐帧累计的 usage(null = 没报);chat.js 的 normalizeAgentUsage 只补没报的字段 + upstreamUsage, upstreamFinishReason, upstreamCompleted: streamResult.completed, upstreamEventCount: streamResult.eventCount, diff --git a/src/utils/precise-tokenizer.js b/src/utils/precise-tokenizer.js index e039e990..9cd704a3 100644 --- a/src/utils/precise-tokenizer.js +++ b/src/utils/precise-tokenizer.js @@ -4,6 +4,7 @@ */ const tiktoken = require('tiktoken') +const { logger } = require('./logger.js') /** * 使用tiktoken进行精准token计数 @@ -68,23 +69,13 @@ function countMessagesTokens(messages, model = 'gpt-3.5-turbo') { } /** - * 创建精准的usage对象 + * 本地估算的 usage 对象(只在上游没报时用,见 reportUsage) * @param {Array|string} promptMessages - 提示消息或文本 * @param {string} completionText - 完成文本 - * @param {object} realUsage - 真实的usage数据(如果有) * @param {string} model - 模型名称 * @returns {object} usage对象 */ -function createUsageObject(promptMessages, completionText = '', realUsage = null, model = 'gpt-3.5-turbo') { - // 如果有真实的usage数据,优先使用 - if (realUsage && realUsage.prompt_tokens && realUsage.completion_tokens) { - return { - prompt_tokens: realUsage.prompt_tokens, - completion_tokens: realUsage.completion_tokens, - total_tokens: realUsage.total_tokens || (realUsage.prompt_tokens + realUsage.completion_tokens) - } - } - +function createUsageObject(promptMessages, completionText = '', model = 'gpt-3.5-turbo') { // 计算prompt tokens let promptTokens = 0 if (Array.isArray(promptMessages)) { @@ -103,8 +94,88 @@ function createUsageObject(promptMessages, completionText = '', realUsage = null } } +/** + * 把上游帧里的一个计数字段转成有效数字:负数、NaN、非数字、 + * 以及 0(上游"没数"时也发 0)都当作"没报"→ null。 + */ +function toReportedCount(value) { + return (typeof value === 'number' && Number.isFinite(value) && value > 0) ? value : null +} + +function firstReportedCount(raw, keys) { + for (const key of keys) { + const count = toReportedCount(raw[key]) + if (count !== null) return count + } + return null +} + +/** + * 上游 usage 归一化。Qwen(DashScope 命名)发 input_tokens / output_tokens, + * OpenAI 命名发 prompt_tokens / completion_tokens;统一成 OpenAI 命名。 + * 没报的字段为 null,让调用方只补估算那一个字段。 + * @param {*} raw - 上游帧里的 usage 对象 + * @returns {{prompt_tokens: number|null, completion_tokens: number|null}|null} 一个可用字段都没有时返回 null + */ +function normalizeUpstreamUsage(raw) { + if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return null + const prompt_tokens = firstReportedCount(raw, ['input_tokens', 'prompt_tokens']) + const completion_tokens = firstReportedCount(raw, ['output_tokens', 'completion_tokens']) + if (prompt_tokens === null && completion_tokens === null) return null + return { prompt_tokens, completion_tokens } +} + +/** + * 逐帧累积上游 usage。Qwen 每个 typing 帧都带累计值,最后的 finished 帧不带: + * 报了的字段以最后一次为准,没报的保持已累积的值。 + * @param {{prompt_tokens: number|null, completion_tokens: number|null}|null} acc - 累积值(初始 null) + * @param {*} rawFrameUsage - 当前帧的 usage + */ +function mergeUpstreamUsage(acc, rawFrameUsage) { + const frame = normalizeUpstreamUsage(rawFrameUsage) + if (!frame) return acc + return { + prompt_tokens: frame.prompt_tokens ?? acc?.prompt_tokens ?? null, + completion_tokens: frame.completion_tokens ?? acc?.completion_tokens ?? null + } +} + +/** + * 只对上游没报的字段补本地估算;两项都有时不调用估算(tiktoken 有成本)。 + * @param {{prompt_tokens: number|null, completion_tokens: number|null}|null} acc - 累积的上游 usage + * @param {() => {prompt_tokens: number, completion_tokens: number}} estimate - 惰性本地估算 + * @returns {{prompt_tokens: number, completion_tokens: number, total_tokens: number}} + */ +function resolveUsage(acc, estimate) { + const upstreamPrompt = acc?.prompt_tokens ?? null + const upstreamCompletion = acc?.completion_tokens ?? null + const estimated = (upstreamPrompt === null || upstreamCompletion === null) ? estimate() : null + const prompt_tokens = upstreamPrompt ?? (estimated.prompt_tokens || 0) + const completion_tokens = upstreamCompletion ?? (estimated.completion_tokens || 0) + return { prompt_tokens, completion_tokens, total_tokens: prompt_tokens + completion_tokens } +} + +/** + * resolveUsage + 每个响应一行日志。来源:两项都来自上游是 "upstream"; + * 哪怕只有一项是本地估算的也算 "estimated"。 + * @param {{prompt_tokens: number|null, completion_tokens: number|null}|null} acc - 累积的上游 usage + * @param {() => {prompt_tokens: number, completion_tokens: number}} estimate - 惰性本地估算 + * @param {string} tag - 日志模块标签('ANTHROPIC' / 'CHAT') + * @returns {{prompt_tokens: number, completion_tokens: number, total_tokens: number}} + */ +function reportUsage(acc, estimate, tag) { + let source = 'upstream' + const usage = resolveUsage(acc, () => { source = 'estimated'; return estimate() }) + logger.info(`usage source=${source} input=${usage.prompt_tokens} output=${usage.completion_tokens}`, tag) + return usage +} + module.exports = { countTokens, countMessagesTokens, - createUsageObject + createUsageObject, + normalizeUpstreamUsage, + mergeUpstreamUsage, + resolveUsage, + reportUsage } diff --git a/tests/anthropic-usage-passthrough.test.js b/tests/anthropic-usage-passthrough.test.js new file mode 100644 index 00000000..3739ecfe --- /dev/null +++ b/tests/anthropic-usage-passthrough.test.js @@ -0,0 +1,192 @@ +// Reported usage = upstream usage. Qwen manda `usage` con nombres DashScope +// (input_tokens / output_tokens / total_tokens), acumulado, en cada frame `typing`; +// el frame final no lo trae. El proxy leia prompt_tokens / completion_tokens (OpenAI) +// y por eso TODAS las respuestas caian al estimado local (input_tokens: 8 para "hi", +// cuando Qwen contaba 651). +// +// Seam: los handlers del controller alimentados con frames upstream sinteticos. +// Harness copiado de tests/anthropic-cap-drain.test.js —— cada archivo de test corre +// en su propio proceso. + +const test = require('node:test'); +const { describe, it } = test; +const assert = require('node:assert/strict'); + +// El harness requiere account.js, que en modo file haria login real con data/data.json. +process.env.API_KEY = 'usage-test-key'; +process.env.DATA_SAVE_MODE = 'none'; +process.env.ACCOUNTS = ''; +process.env.ENABLE_CLI = 'false'; + +// Sin red en tests: los parches de require-cache van ANTES de requerir el controller +// (anthropic.js captura sendChatRequest por destructuring en su primer require). +const modelsMap = require('../src/models/models-map.js'); +modelsMap.getLatestModels = async () => { throw new Error('offline test: no model fetch'); }; +const requestModule = require('../src/utils/request.js'); +let upstreamFactory = null; +requestModule.sendChatRequest = async () => (upstreamFactory + ? { status: true, response: upstreamFactory(), currentAccount: null } + : { status: false }); + +const { handleAnthropicStream, handleAnthropicMessages } = require('../src/controllers/anthropic.js'); + +test.after(() => { + require('../src/utils/account.js').destroy(); +}); + +const createMockStreamResponse = () => ({ + output: '', + headers: {}, + writableEnded: false, + set(headers) { Object.assign(this.headers, headers); return this; }, + status() { return this; }, + write(chunk) { this.output += String(chunk); return true; }, + end(chunk = '') { this.output += String(chunk); this.writableEnded = true; } +}); + +const createMockJsonResponse = () => ({ + statusCode: 200, + body: null, + headers: {}, + set(headers) { Object.assign(this.headers, headers); return this; }, + status(code) { this.statusCode = code; return this; }, + json(payload) { this.body = payload; return this; } +}); + +/** Frame de respuesta como lo manda Qwen: `usage` al nivel del frame, junto a `choices`. */ +const frame = (content, usage) => `data: ${JSON.stringify({ + choices: [{ delta: { phase: 'answer', content }, finish_reason: null }], + ...(usage === undefined ? {} : { usage }) +})}\n\n`; + +const STOP = 'data: {"choices":[{"delta":{},"finish_reason":"stop"}]}\n\ndata: [DONE]\n\n'; + +const upstream = (frames) => { + async function* gen() { + for (const f of frames) yield f; + } + return gen(); +}; + +const eventsOf = (output) => output + .split('\n\n') + .filter(Boolean) + .map(chunk => chunk.split('\n').find(line => line.startsWith('data: '))) + .filter(Boolean) + .map(line => JSON.parse(line.slice(6))); + +/** `secondAttemptFrames`: lo que devuelve el sendRequest del segundo attempt; null = un solo attempt. */ +const runStream = async (frames, secondAttemptFrames = null) => { + const res = createMockStreamResponse(); + await handleAnthropicStream(res, { + message_id: 'msg_usage', + model: 'qwen-test', + hasTools: false, + toolChoice: null, + allowedToolNames: [], + toolSchemas: {}, + requestBody: { messages: [{ role: 'user', content: 'hi' }] }, + sendRequest: async () => (secondAttemptFrames ? { status: true, response: upstream(secondAttemptFrames) } : { status: false }) + }, upstream(frames)); + return eventsOf(res.output); +}; + +const deltaUsage = (events) => events.find(e => e.type === 'message_delta').usage; + +describe('reported usage comes from upstream usage (Anthropic /v1/messages)', () => { + it('stream: message_delta carries Qwen input_tokens/output_tokens; cumulative, last typing frame wins', async () => { + const events = await runStream([ + frame('Hel', { input_tokens: 651, output_tokens: 1, total_tokens: 652 }), + frame('lo', { input_tokens: 651, output_tokens: 2, total_tokens: 653 }), + STOP + ]); + const usage = deltaUsage(events); + assert.equal(usage.input_tokens, 651); + assert.equal(usage.output_tokens, 2); + }); + + it('stream: cache fields are 0 (not null) in message_start and message_delta', async () => { + const events = await runStream([frame('Hi', { input_tokens: 651, output_tokens: 1 }), STOP]); + for (const event of [events.find(e => e.type === 'message_start').message, events.find(e => e.type === 'message_delta')]) { + assert.equal(event.usage.cache_creation_input_tokens, 0); + assert.equal(event.usage.cache_read_input_tokens, 0); + } + }); + + it('stream: no upstream usage at all → both counters estimated locally (never 0)', async () => { + const usage = deltaUsage(await runStream([frame('Hello there'), STOP])); + assert.ok(usage.input_tokens > 0, `input_tokens=${usage.input_tokens}`); + assert.ok(usage.output_tokens > 0, `output_tokens=${usage.output_tokens}`); + }); + + it('stream: partial upstream usage → only the missing counter is estimated', async () => { + const usage = deltaUsage(await runStream([frame('Hello there', { input_tokens: 651 }), STOP])); + assert.equal(usage.input_tokens, 651); + assert.ok(usage.output_tokens > 0, `output_tokens=${usage.output_tokens}`); + }); + + it('stream: all-zero upstream usage is treated as absent → estimated', async () => { + const usage = deltaUsage(await runStream([frame('Hello there', { input_tokens: 0, output_tokens: 0, total_tokens: 0 }), STOP])); + assert.ok(usage.input_tokens > 0, `input_tokens=${usage.input_tokens}`); + assert.ok(usage.output_tokens > 0, `output_tokens=${usage.output_tokens}`); + }); + + it('stream: several attempts → the last attempt\'s usage; a later attempt that reports nothing does NOT inherit the first attempt\'s counters', async () => { + // Attempt 1: sin texto visible → el gate reintenta ('empty'); traia 651/9 del upstream. + // Attempt 2: texto valido y SOLO output_tokens. Reportado: output 5 (attempt 2) e input + // estimado —— nunca los 651 del attempt 1. + const usage = deltaUsage(await runStream( + [frame('', { input_tokens: 651, output_tokens: 9, total_tokens: 660 }), STOP], + [frame('Hello there', { output_tokens: 5 }), STOP] + )); + assert.equal(usage.output_tokens, 5); + assert.ok(usage.input_tokens > 0 && usage.input_tokens !== 651, `input_tokens=${usage.input_tokens}`); + }); + + /** Cada argumento es un attempt: el primero es el upstream inicial, los demas, attempts posteriores. */ + const runNonStream = async (...attempts) => { + const queue = [...attempts]; + upstreamFactory = () => { + const frames = queue.shift(); + assert.ok(frames, 'upstream pedido mas veces que attempts preparados'); + return upstream(frames); + }; + try { + const res = createMockJsonResponse(); + await handleAnthropicMessages({ + body: { model: 'qwen3-max', max_tokens: 64, stream: false, messages: [{ role: 'user', content: 'hi' }] } + }, res); + assert.equal(res.statusCode, 200, JSON.stringify(res.body)); + return res.body.usage; + } finally { + upstreamFactory = null; + } + }; + + it('non-stream twin via handleAnthropicMessages: body.usage carries Qwen counts, cache fields 0', async () => { + const usage = await runNonStream([ + frame('Hel', { input_tokens: 651, output_tokens: 1, total_tokens: 652 }), + frame('lo', { input_tokens: 651, output_tokens: 2, total_tokens: 653 }), + STOP + ]); + assert.equal(usage.input_tokens, 651); + assert.equal(usage.output_tokens, 2); + assert.equal(usage.cache_creation_input_tokens, 0); + assert.equal(usage.cache_read_input_tokens, 0); + }); + + it('non-stream: no upstream usage → estimated, never 0', async () => { + const usage = await runNonStream([frame('Hello there'), STOP]); + assert.ok(usage.input_tokens > 0, `input_tokens=${usage.input_tokens}`); + assert.ok(usage.output_tokens > 0, `output_tokens=${usage.output_tokens}`); + }); + + it('non-stream: several attempts → the last attempt\'s usage; a later attempt that reports nothing does NOT inherit the first attempt\'s counters', async () => { + const usage = await runNonStream( + [frame('', { input_tokens: 651, output_tokens: 9, total_tokens: 660 }), STOP], + [frame('Hello there', { output_tokens: 5 }), STOP] + ); + assert.equal(usage.output_tokens, 5); + assert.ok(usage.input_tokens > 0 && usage.input_tokens !== 651, `input_tokens=${usage.input_tokens}`); + }); +}); diff --git a/tests/expected-counts.json b/tests/expected-counts.json index e9f04141..b070841b 100644 --- a/tests/expected-counts.json +++ b/tests/expected-counts.json @@ -1,6 +1,6 @@ { - "tests": 1115, - "suites": 128, + "tests": 1146, + "suites": 134, "note": "Authoritative count. Verify with the per-file sum in AGENTS.md (\"The test gate\"). The -a on that grep is load-bearing: tool-prompt.test.js emits bytes that make grep call the stream binary, and without -a its whole summary line — 133 tests — is silently dropped from the sum.", - "updated": "2026-09-12" + "updated": "2026-09-16" } diff --git a/tests/openai-usage-passthrough.test.js b/tests/openai-usage-passthrough.test.js new file mode 100644 index 00000000..ac26c6fe --- /dev/null +++ b/tests/openai-usage-passthrough.test.js @@ -0,0 +1,175 @@ +// Reported usage = upstream usage en /v1/chat/completions (stream, non-stream y agent runtime). +// Gemelo de tests/anthropic-usage-passthrough.test.js: Qwen manda `usage` con nombres +// DashScope (input_tokens / output_tokens), acumulado por frame; el proxy leia +// prompt_tokens / completion_tokens y caia SIEMPRE al estimado local. +// +// Harness copiado de tests/openai-residue.test.js (cada archivo corre en su propio proceso). + +const test = require('node:test'); +const { describe, it } = test; +const assert = require('node:assert/strict'); + +// El harness requiere account.js, que en modo file haria login real con data/data.json. +process.env.API_KEY = 'usage-test-key'; +process.env.DATA_SAVE_MODE = 'none'; +process.env.ACCOUNTS = ''; +process.env.ENABLE_CLI = 'false'; + +// Sin red en tests: mismos parches de require-cache que el resto de la suite. +const modelsMap = require('../src/models/models-map.js'); +modelsMap.getLatestModels = async () => { throw new Error('offline test: no model fetch'); }; +const requestModule = require('../src/utils/request.js'); +requestModule.sendChatRequest = async () => ({ status: false }); + +const { handleStreamResponse, handleNonStreamResponse } = require('../src/controllers/chat.js'); + +test.after(() => { + require('../src/utils/account.js').destroy(); +}); + +const createMockResponse = () => ({ + output: '', + headers: {}, + headersSent: false, + writableEnded: false, + statusCode: 200, + set(headers) { Object.assign(this.headers, headers); return this; }, + setHeader(name, value) { this.headers[name] = value; }, + write(chunk) { this.headersSent = true; this.output += String(chunk); return true; }, + end(chunk = '') { if (chunk) this.write(chunk); this.writableEnded = true; }, + status(code) { this.statusCode = code; return this; }, + json(value) { + this.headersSent = true; + this.output += JSON.stringify(value); + this.writableEnded = true; + return this; + } +}); + +/** Frame como lo manda Qwen: `usage` al nivel del frame, junto a `choices`. */ +const frame = (content, usage) => `data: ${JSON.stringify({ + choices: [{ delta: { phase: 'answer', content }, finish_reason: null }], + ...(usage === undefined ? {} : { usage }) +})}\n\n`; + +const STOP = 'data: {"choices":[{"delta":{},"finish_reason":"stop"}]}\n\ndata: [DONE]\n\n'; + +const upstreamOf = (frames) => { + async function* gen() { + for (const f of frames) yield f; + } + return gen(); +}; + +const deltasOf = (output) => output + .split('\n\n') + .filter(Boolean) + .map(chunk => chunk.replace(/^data: /, '')) + .filter(payload => payload && payload !== '[DONE]') + .map(payload => JSON.parse(payload)); + +const QWEN_USAGE = [ + frame('Hel', { input_tokens: 651, output_tokens: 1, total_tokens: 652 }), + frame('lo', { input_tokens: 651, output_tokens: 2, total_tokens: 653 }), + STOP +]; +const NO_USAGE = [frame('Hello there'), STOP]; +const PARTIAL_USAGE = [frame('Hello there', { input_tokens: 651 }), STOP]; +const ZERO_USAGE = [frame('Hello there', { input_tokens: 0, output_tokens: 0, total_tokens: 0 }), STOP]; +const REQUEST_BODY = { messages: [{ role: 'user', content: 'hi' }] }; + +const lastStreamUsage = (output) => deltasOf(output).map(e => e.usage).filter(Boolean).pop(); + +const runStream = async (frames, options = { has_tools: false }) => { + const res = createMockResponse(); + await handleStreamResponse(res, upstreamOf(frames), false, false, REQUEST_BODY, options); + return lastStreamUsage(res.output); +}; + +const runNonStream = async (frames, options = { has_tools: false }) => { + const res = createMockResponse(); + await handleNonStreamResponse(res, upstreamOf(frames), false, false, 'qwen-test', REQUEST_BODY, options); + assert.equal(res.statusCode, 200, res.output); + return JSON.parse(res.output).usage; +}; + +const assertEstimated = (usage) => { + assert.ok(usage, 'usage object present'); + assert.ok(usage.prompt_tokens > 0, `prompt_tokens=${usage.prompt_tokens}`); + assert.ok(usage.completion_tokens > 0, `completion_tokens=${usage.completion_tokens}`); + assert.equal(usage.total_tokens, usage.prompt_tokens + usage.completion_tokens); +}; + +describe('reported usage comes from upstream usage (/v1/chat/completions)', () => { + it('stream: the final chunk carries Qwen counts (cumulative, last typing frame wins)', async () => { + const usage = await runStream(QWEN_USAGE); + assert.equal(usage.prompt_tokens, 651); + assert.equal(usage.completion_tokens, 2); + assert.equal(usage.total_tokens, 653); + }); + + it('stream: no upstream usage → estimated, never 0', async () => { + assertEstimated(await runStream(NO_USAGE)); + }); + + it('stream: partial upstream usage → only the missing counter is estimated', async () => { + const usage = await runStream(PARTIAL_USAGE); + assert.equal(usage.prompt_tokens, 651); + assert.ok(usage.completion_tokens > 0, `completion_tokens=${usage.completion_tokens}`); + }); + + it('non-stream: body.usage carries Qwen counts', async () => { + const usage = await runNonStream(QWEN_USAGE); + assert.equal(usage.prompt_tokens, 651); + assert.equal(usage.completion_tokens, 2); + assert.equal(usage.total_tokens, 653); + }); + + it('stream: all-zero upstream usage is treated as absent → estimated', async () => { + assertEstimated(await runStream(ZERO_USAGE)); + }); + + it('non-stream: no upstream usage → estimated, never 0', async () => { + assertEstimated(await runNonStream(NO_USAGE)); + }); + + it('non-stream: all-zero upstream usage is treated as absent → estimated', async () => { + assertEstimated(await runNonStream(ZERO_USAGE)); + }); + + const agentOptions = () => ({ + has_tools: true, + tool_choice: 'auto', + allowed_tool_names: ['Read'], + tool_schemas: { Read: { type: 'object', properties: { file_path: { type: 'string' } }, required: ['file_path'] } }, + agent_turn_max_attempts: 1, + upstream_request_body: REQUEST_BODY, + sendChatRequest: async () => ({ status: false }) + }); + + it('agent runtime (has_tools): the accepted attempt reports Qwen counts, not the estimate', async () => { + const frames = [frame('Hello', { input_tokens: 651, output_tokens: 9, total_tokens: 660 }), STOP]; + const usage = await runNonStream(frames, agentOptions()); + assert.equal(usage.prompt_tokens, 651); + assert.equal(usage.completion_tokens, 9); + }); + + it('agent runtime (has_tools): no upstream usage → estimated, never 0', async () => { + const frames = [frame('Hello there'), STOP]; + assertEstimated(await runNonStream(frames, agentOptions())); + }); + + it('agent runtime (has_tools): several attempts → reported usage is the accepted attempt\'s, not the first nor the sum', async () => { + // Attempt 1: vacío → el gate reintenta ('empty'). Attempt 2: respuesta válida + // con OTROS números. Lo reportado debe ser lo del attempt aceptado (700/5), no 651/9 ni 1351/14. + const first = [frame('', { input_tokens: 651, output_tokens: 9, total_tokens: 660 }), STOP]; + const second = [frame('Hello', { input_tokens: 700, output_tokens: 5, total_tokens: 705 }), STOP]; + let extraAttempts = 0; + const sendChatRequest = async () => { extraAttempts += 1; return { status: true, response: upstreamOf(second) }; }; + const usage = await runNonStream(first, { ...agentOptions(), agent_turn_max_attempts: 2, sendChatRequest }); + assert.equal(extraAttempts, 1, 'hubo exactamente un segundo attempt'); + assert.equal(usage.prompt_tokens, 700); + assert.equal(usage.completion_tokens, 5); + assert.equal(usage.total_tokens, 705); + }); +}); diff --git a/tests/upstream-usage.test.js b/tests/upstream-usage.test.js new file mode 100644 index 00000000..0a582096 --- /dev/null +++ b/tests/upstream-usage.test.js @@ -0,0 +1,114 @@ +// Unidad: normalizacion del `usage` upstream. Qwen (DashScope) manda +// input_tokens / output_tokens; OpenAI manda prompt_tokens / completion_tokens. +// El proxy reporta en formato OpenAI y estima localmente SOLO lo que upstream no dio. + +const test = require('node:test'); +const { describe, it } = test; +const assert = require('node:assert/strict'); + +const { + normalizeUpstreamUsage, + mergeUpstreamUsage, + resolveUsage, + reportUsage +} = require('../src/utils/precise-tokenizer.js'); + +describe('normalizeUpstreamUsage', () => { + it('DashScope naming (what Qwen sends) → OpenAI naming', () => { + assert.deepEqual( + normalizeUpstreamUsage({ input_tokens: 651, output_tokens: 9, total_tokens: 660 }), + { prompt_tokens: 651, completion_tokens: 9 } + ); + }); + + it('OpenAI naming passes through', () => { + assert.deepEqual( + normalizeUpstreamUsage({ prompt_tokens: 12, completion_tokens: 3, total_tokens: 15 }), + { prompt_tokens: 12, completion_tokens: 3 } + ); + }); + + it('partial usage: the present field is kept, the missing one is null', () => { + assert.deepEqual(normalizeUpstreamUsage({ input_tokens: 651 }), { prompt_tokens: 651, completion_tokens: null }); + assert.deepEqual(normalizeUpstreamUsage({ output_tokens: 9 }), { prompt_tokens: null, completion_tokens: 9 }); + }); + + it('negative, NaN and non-numeric fields count as absent', () => { + assert.deepEqual(normalizeUpstreamUsage({ input_tokens: -1, output_tokens: 'nine' }), null); + assert.deepEqual(normalizeUpstreamUsage({ input_tokens: NaN, output_tokens: 4 }), { prompt_tokens: null, completion_tokens: 4 }); + }); + + it('non-object or no usable field → null', () => { + for (const raw of [null, undefined, 'usage', 42, [], {}, { foo: 1 }]) { + assert.equal(normalizeUpstreamUsage(raw), null, `raw=${JSON.stringify(raw)}`); + } + }); + + it('zero is "not reported": all-zero → null, a single zero field → null for that field', () => { + assert.equal(normalizeUpstreamUsage({ input_tokens: 0, output_tokens: 0, total_tokens: 0 }), null); + assert.deepEqual(normalizeUpstreamUsage({ input_tokens: 651, output_tokens: 0 }), { prompt_tokens: 651, completion_tokens: null }); + }); +}); + +describe('reportUsage (resolveUsage + one log line per response)', () => { + it('logs source=upstream only when both counters came from upstream; otherwise source=estimated', () => { + const { logger } = require('../src/utils/logger.js'); + const estimate = () => ({ prompt_tokens: 8, completion_tokens: 3, total_tokens: 11 }); + const lines = []; + const original = logger.info; + logger.info = (message, tag) => { lines.push(`[${tag}] ${message}`); }; + try { + assert.deepEqual( + reportUsage({ prompt_tokens: 651, completion_tokens: 9 }, estimate, 'T'), + { prompt_tokens: 651, completion_tokens: 9, total_tokens: 660 } + ); + reportUsage({ prompt_tokens: 651, completion_tokens: null }, estimate, 'T'); + reportUsage(null, estimate, 'T'); + } finally { + logger.info = original; + } + assert.deepEqual(lines, [ + '[T] usage source=upstream input=651 output=9', + '[T] usage source=estimated input=651 output=3', + '[T] usage source=estimated input=8 output=3' + ]); + }); +}); + +describe('mergeUpstreamUsage (per-frame accumulation)', () => { + it('cumulative counts: the last frame that reports a field wins', () => { + let acc = null; + acc = mergeUpstreamUsage(acc, { input_tokens: 651, output_tokens: 1 }); + acc = mergeUpstreamUsage(acc, { input_tokens: 651, output_tokens: 2 }); + assert.deepEqual(acc, { prompt_tokens: 651, completion_tokens: 2 }); + }); + + it('a frame without usage (or with a partial one) keeps what was already accumulated', () => { + let acc = mergeUpstreamUsage(null, { input_tokens: 651, output_tokens: 5 }); + acc = mergeUpstreamUsage(acc, undefined); + acc = mergeUpstreamUsage(acc, { output_tokens: 7 }); + assert.deepEqual(acc, { prompt_tokens: 651, completion_tokens: 7 }); + }); +}); + +describe('resolveUsage (fill only what upstream never reported)', () => { + const estimate = () => ({ prompt_tokens: 8, completion_tokens: 3, total_tokens: 11 }); + + it('upstream reported both → estimator is not even called', () => { + let calls = 0; + const usage = resolveUsage({ prompt_tokens: 651, completion_tokens: 9 }, () => { calls++; return estimate(); }); + assert.deepEqual(usage, { prompt_tokens: 651, completion_tokens: 9, total_tokens: 660 }); + assert.equal(calls, 0); + }); + + it('only the missing field is estimated; total is recomputed', () => { + assert.deepEqual( + resolveUsage({ prompt_tokens: 651, completion_tokens: null }, estimate), + { prompt_tokens: 651, completion_tokens: 3, total_tokens: 654 } + ); + }); + + it('nothing from upstream → full estimate', () => { + assert.deepEqual(resolveUsage(null, estimate), { prompt_tokens: 8, completion_tokens: 3, total_tokens: 11 }); + }); +});