Files
Sakurasan d8257df100 fix: messages 流式缓存场景 token 记账错乱
qwen/dashscope 等上游 messages 流式的 usage 语义:
- message_start.usage.input_tokens = 总输入
- message_delta.usage.input_tokens = 非缓存输入(缓存部分单列
  cache_read/cache_creation 字段),是最终计费口径

原 usageSink 字段级合并中 delta 的 input 覆盖 start 的 input,
总输入丢失(31790 → 8);缓存写也未参与计费。

- push:带 cache_* 字段的 usage 视为最终口径,整体替换 sink
- finishUsage:缓存写按 1.25× 输入价计费(Anthropic 5m 口径);
  落库 input_tokens 存总量(含缓存读/写)便于对账
- 估算兜底条件排除已有缓存计数的请求
- 回归测试:缓存写/缓存命中/chat 末块合并不回归
2026-08-28 18:11:24 +08:00

102 lines
4.4 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package proxy
import (
"encoding/json"
"testing"
)
// 复现线上火山方舟 qwen 流:data:{...} 无空格(省略 data: 后的空格)。
func TestSSEDataPayloadNoSpace(t *testing.T) {
lines := []string{
`data:{"message":{"model":"qwen3.8-flash","id":"msg_1","role":"assistant","type":"message","content":[],"usage":{"input_tokens":31626,"output_tokens":0}},"type":"message_start"}`,
`data:{"delta":{"type":"text_delta","text":"你好"},"type":"content_block_delta","index":0}`,
`data:{"delta":{"type":"text_delta","text":"!"},"type":"content_block_delta","index":0}`,
`data:{"delta":{"stop_reason":"end_turn"},"type":"message_delta","usage":{"cache_creation":{"ephemeral_5m_input_tokens":33065},"output_tokens":8,"cache_creation_input_tokens":33065,"input_tokens":8,"cache_read_input_tokens":0}}`,
}
var out string
for _, l := range lines {
out += sseContentText([]byte(l))
}
if out != "你好!" {
t.Fatalf("outputText=%q, want %q", out, "你好!")
}
// message_delta 的 delta.usage 应能提取(output_tokens=8)
u := scanUsage([]byte(lines[3]))
if u == nil {
t.Fatal("scanUsage returned nil for message_delta with usage")
}
var sh usageShape
if err := json.Unmarshal(u, &sh); err != nil {
t.Fatalf("unmarshal usage: %v", err)
}
if sh.OutputTokens != 8 {
t.Fatalf("output_tokens=%d, want 8", sh.OutputTokens)
}
}
// 兼容带空格的单 data: 行(标准 SSE)与 event:+data: 多行块。
func TestSSEDataPayloadSpacedAndMultiLine(t *testing.T) {
// 标准:data: {...}
if got := sseContentText([]byte(`data: {"delta":{"type":"text_delta","text":"hi"},"type":"content_block_delta","index":0}`)); got != "hi" {
t.Fatalf("spaced single line: got %q, want hi", got)
}
// 多行块:event: message_delta\ndata: {...}
block := []byte("event: message_delta\ndata: {\"delta\":{\"type\":\"text_delta\",\"text\":\"yo\"},\"type\":\"content_block_delta\",\"index\":0}\n")
if got := sseContentText(block); got != "yo" {
t.Fatalf("multiline block: got %q, want yo", got)
}
}
// 复现线上 qwen(dashscope)messages 流式缓存场景:
// message_start.usage.input_tokens 是总输入,message_delta.usage.input_tokens 是非缓存输入
// 且带 cache_read/cache_creation,是最终计费口径。合并后:
// in(落库)=input+cache_read+cache_creation,计价 in 只算非缓存部分。
// 此前 delta 的 input 覆盖 start 的 input 导致总输入丢失(31790 → 8)。
func TestUsageSinkMessageDeltaAuthoritative(t *testing.T) {
sink := &usageSink{}
// message_start:总输入 31790
start := json.RawMessage(`{"input_tokens":31790,"output_tokens":0}`)
sink.push(start)
if got := sink.us.InputTokens; got != 31790 {
t.Fatalf("after start: input=%d, want 31790", got)
}
// message_delta:非缓存输入 8 + 缓存写 33229(最终口径,整体替换)
delta := json.RawMessage(`{"output_tokens":8,"cache_creation_input_tokens":33229,"input_tokens":8,"cache_read_input_tokens":0}`)
sink.push(delta)
s := sink.Shape()
if s.InputTokens != 8 || s.CacheCreationInputTokens != 33229 || s.OutputTokens != 8 {
t.Fatalf("after delta: %+v, want input=8 cache_create=33229 output=8", s)
}
// finishUsage 口径:落库 input = 8 + 0 + 33229 = 33237(总量),计价 in=8、cacheCreate=33229
in := s.InputTokens + s.CacheReadInputTokens + s.CacheCreationInputTokens
if in != 33237 {
t.Fatalf("total input=%d, want 33237", in)
}
}
// 缓存命中场景(id=55):delta input=76 非缓存 + cache_read=33229 + cache_creation=17。
func TestUsageSinkCacheHitMerge(t *testing.T) {
sink := &usageSink{}
sink.push(json.RawMessage(`{"input_tokens":31862,"output_tokens":0}`))
sink.push(json.RawMessage(`{"output_tokens":32,"cache_creation_input_tokens":17,"input_tokens":76,"cache_read_input_tokens":33229}`))
s := sink.Shape()
total := s.InputTokens + s.CacheReadInputTokens + s.CacheCreationInputTokens
if total != 33322 {
t.Fatalf("total input=%d, want 33322 (76+33229+17)", total)
}
if s.OutputTokens != 32 {
t.Fatalf("output=%d, want 32", s.OutputTokens)
}
}
// OpenAI chat 末块(无 cache 字段)仍走零值不覆盖合并,不受整体替换影响。
func TestUsageSinkChatLastChunkStillMerges(t *testing.T) {
sink := &usageSink{}
sink.push(json.RawMessage(`{"prompt_tokens":65,"completion_tokens":0}`))
sink.push(json.RawMessage(`{"prompt_tokens":65,"completion_tokens":82}`))
s := sink.Shape()
if s.PromptTokens != 65 || s.CompletionTokens != 82 {
t.Fatalf("chat merge broken: %+v", s)
}
}