From 8e0e9ac319f47b85949200c626f90917f8697653 Mon Sep 17 00:00:00 2001 From: Jason Date: Fri, 5 Jun 2026 21:45:34 +0800 Subject: [PATCH] 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 --- src-tauri/src/proxy/usage/parser.rs | 164 +++++++++++++++++++++++++--- 1 file changed, 147 insertions(+), 17 deletions(-) diff --git a/src-tauri/src/proxy/usage/parser.rs b/src-tauri/src/proxy/usage/parser.rs index f2b797ee3..8b186efd0 100644 --- a/src-tauri/src/proxy/usage/parser.rs +++ b/src-tauri/src/proxy/usage/parser.rs @@ -85,6 +85,7 @@ impl TokenUsage { let mut usage = Self::default(); let mut model: Option = None; let mut message_id: Option = 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 流式响应解析