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) } }