Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 28 additions & 26 deletions src/controllers/anthropic.js
Original file line number Diff line number Diff line change
@@ -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');
Expand Down Expand Up @@ -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
}
}
});
Expand All @@ -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;
Expand Down Expand Up @@ -1346,6 +1347,8 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
// 但本控制器没有 Agent 回合门禁去解包,标签会原样发给客户端。剥掉它们。
agentTagStripper = createAgentTagStripper();
recoveredBuffer = '';
// usage 也按轮全新:报的是最后一轮上游给的,没给就估算,绝不继承上一轮的。
upstreamUsage = null;
attemptVisibleText = '';
attemptThinkText = '';
attemptThinkEvidence = false;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -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' });
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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()) {
Expand Down Expand Up @@ -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
}
});
};
Expand Down
60 changes: 15 additions & 45 deletions src/controllers/chat.js
Original file line number Diff line number Diff line change
@@ -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 {
Expand Down Expand Up @@ -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')
}

/**
Expand Down Expand Up @@ -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 = ''

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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 路径罕见,
Expand Down Expand Up @@ -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 = ''
Expand All @@ -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]
Expand Down Expand Up @@ -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)
Expand Down
19 changes: 6 additions & 13 deletions src/utils/openai-agent-runtime.js
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down Expand Up @@ -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
Expand All @@ -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]
Expand Down Expand Up @@ -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,
Expand Down
Loading