fix(usage): correct inflated input_tokens in Claude stream parsing

Some Anthropic-compatible SSE providers (e.g. qwen, minimax) report the
full context (fresh + cached) as input_tokens in message_start, double
counting the cached portion that is also reported in
cache_read_input_tokens. This inflated the cacheable-input denominator
and pushed the displayed cache hit rate artificially low.

When a message_delta carries a smaller positive input_tokens, prefer it
over the message_start value and adopt the cache counts from the same
usage block to avoid double counting; fall back to the start cache
values when the delta omits them. Native Claude (no input in delta) and
OpenRouter-converted (input only in delta) paths are unchanged.

Refs #3580
This commit is contained in:
Jason
2026-06-05 21:45:34 +08:00
parent bda625a4f1
commit 8e0e9ac319
+147 -17
View File
@@ -85,6 +85,7 @@ impl TokenUsage {
let mut usage = Self::default();
let mut model: Option<String> = None;
let mut message_id: Option<String> = None;
let mut input_from_delta = false;
for event in events {
if let Some(event_type) = event.get("type").and_then(|v| v.as_str()) {
@@ -129,32 +130,52 @@ impl TokenUsage {
{
usage.output_tokens = output as u32;
}
// OpenRouter 转换后的流式响应:input_tokens 也在 message_delta 中
// 如果 message_start 中没有 input_tokens,则从 message_delta 获取
if usage.input_tokens == 0 {
if let Some(input) =
delta_usage.get("input_tokens").and_then(|v| v.as_u64())
{
usage.input_tokens = input as u32;
let delta_input = delta_usage
.get("input_tokens")
.and_then(|v| v.as_u64())
.map(|v| v as u32);
let delta_cache_read = delta_usage
.get("cache_read_input_tokens")
.and_then(|v| v.as_u64())
.map(|v| v as u32);
let delta_cache_creation = delta_usage
.get("cache_creation_input_tokens")
.and_then(|v| v.as_u64())
.map(|v| v as u32);
// 部分 Anthropic-compatible SSE provider 会在 message_start 上报完整上下文,
// 但在 message_delta 上报修正后的 fresh input。遇到更小的正数 delta input
// 时采用 delta;若同一 usage 块带有缓存计数,也同步采用以避免重复计数。
// 若 delta 缺少缓存字段,则保留 start 中已有的缓存值作为 best-effort fallback。
if let Some(input) = delta_input {
let should_use_delta_input = input > 0
&& (usage.input_tokens == 0
|| input < usage.input_tokens
|| (input_from_delta && input <= usage.input_tokens));
if should_use_delta_input {
usage.input_tokens = input;
input_from_delta = true;
if let Some(cache_read) = delta_cache_read {
usage.cache_read_tokens = cache_read;
}
if let Some(cache_creation) = delta_cache_creation {
usage.cache_creation_tokens = cache_creation;
}
}
}
// 从 message_delta 中处理缓存命中(cache_read_input_tokens)
if usage.cache_read_tokens == 0 {
if let Some(cache_read) = delta_usage
.get("cache_read_input_tokens")
.and_then(|v| v.as_u64())
{
usage.cache_read_tokens = cache_read as u32;
if let Some(cache_read) = delta_cache_read {
usage.cache_read_tokens = cache_read;
}
}
// 从 message_delta 中处理缓存创建(cache_creation_input_tokens)
// 注: 现在 zhipu 没有返回 cache_creation_input_tokens 字段
if usage.cache_creation_tokens == 0 {
if let Some(cache_creation) = delta_usage
.get("cache_creation_input_tokens")
.and_then(|v| v.as_u64())
{
usage.cache_creation_tokens = cache_creation as u32;
if let Some(cache_creation) = delta_cache_creation {
usage.cache_creation_tokens = cache_creation;
}
}
}
@@ -798,6 +819,115 @@ mod tests {
assert_eq!(usage.model, Some("claude-sonnet-4-20250514".to_string()));
}
#[test]
fn test_claude_stream_prefers_smaller_delta_input_and_cache_pair() {
// 部分 Anthropic-compatible provider 会在 message_start 给出包含缓存的总上下文,
// 再在 message_delta 给出修正后的 fresh input,需要以 delta usage 为准。
let events = vec![
json!({
"type": "message_start",
"message": {
"model": "qwen-max",
"usage": {
"input_tokens": 200_000,
"cache_read_input_tokens": 180_000,
"cache_creation_input_tokens": 2_000
}
}
}),
json!({
"type": "message_delta",
"usage": {
"input_tokens": 80_000,
"output_tokens": 1_000,
"cache_read_input_tokens": 120_000,
"cache_creation_input_tokens": 500
}
}),
];
let usage = TokenUsage::from_claude_stream_events(&events).unwrap();
assert_eq!(usage.input_tokens, 80_000);
assert_eq!(usage.output_tokens, 1_000);
assert_eq!(usage.cache_read_tokens, 120_000);
assert_eq!(usage.cache_creation_tokens, 500);
assert_eq!(usage.model, Some("qwen-max".to_string()));
}
#[test]
fn test_claude_stream_updates_cache_pair_from_later_delta_input() {
// 有些 provider 会多次发送带 input 的 message_delta;一旦采用过 delta input
// 后续相同/更小 input 的 delta 应继续更新同一块里的缓存计数。
let events = vec![
json!({
"type": "message_start",
"message": {
"model": "qwen-max",
"usage": {
"input_tokens": 200_000,
"cache_read_input_tokens": 180_000,
"cache_creation_input_tokens": 2_000
}
}
}),
json!({
"type": "message_delta",
"usage": {
"input_tokens": 80_000,
"output_tokens": 100,
"cache_read_input_tokens": 110_000,
"cache_creation_input_tokens": 300
}
}),
json!({
"type": "message_delta",
"usage": {
"input_tokens": 80_000,
"output_tokens": 1_000,
"cache_read_input_tokens": 120_000,
"cache_creation_input_tokens": 500
}
}),
];
let usage = TokenUsage::from_claude_stream_events(&events).unwrap();
assert_eq!(usage.input_tokens, 80_000);
assert_eq!(usage.output_tokens, 1_000);
assert_eq!(usage.cache_read_tokens, 120_000);
assert_eq!(usage.cache_creation_tokens, 500);
assert_eq!(usage.model, Some("qwen-max".to_string()));
}
#[test]
fn test_claude_stream_keeps_start_when_delta_input_is_larger() {
// 正常 Anthropic 语义下,message_start 的 input_tokens 已经可信;
// 如果 delta input 变大,不应覆盖 start input/cache。
let events = vec![
json!({
"type": "message_start",
"message": {
"usage": {
"input_tokens": 100,
"cache_read_input_tokens": 20
}
}
}),
json!({
"type": "message_delta",
"usage": {
"input_tokens": 150,
"output_tokens": 75,
"cache_read_input_tokens": 30
}
}),
];
let usage = TokenUsage::from_claude_stream_events(&events).unwrap();
assert_eq!(usage.input_tokens, 100);
assert_eq!(usage.output_tokens, 75);
assert_eq!(usage.cache_read_tokens, 20);
}
#[test]
fn test_native_claude_stream_parsing() {
// 测试原生 Claude API 流式响应解析