From 79545aa48e3f9f0d6e0b648fa4a7df33ab9884cb Mon Sep 17 00:00:00 2001 From: Sakurasan <26715255+Sakurasan@users.noreply.github.com> Date: Thu, 27 Aug 2026 01:24:30 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E5=85=BC=E5=AE=B9=E4=B8=8A=E6=B8=B8=20S?= =?UTF-8?q?SE=20data:=20=E5=90=8E=E6=97=A0=E7=A9=BA=E6=A0=BC=E6=A0=BC?= =?UTF-8?q?=E5=BC=8F=EF=BC=8C=E6=B5=81=E5=BC=8F=E8=AE=B0=E8=B4=A6=E4=B8=8D?= =?UTF-8?q?=E5=86=8D=200=20=E6=B6=88=E8=80=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 火山方舟等上游流式返回 data:{...}(data: 后无空格),原解析只认 'data: '(带空格),导致 content 文本与 usage 全部漏解析: - sseContentText 抽成 sseDataPayload,统一处理 data: 前缀 (支持带/不带空格 + event:+data: 多行块) - scanUsage 复用 sseDataPayload - convert/stream.go parseLine 也支持 data:{ 无空格 - 新增回归测试覆盖无空格/带空格/多行块三种形态 --- server/internal/proxy/convert/stream.go | 7 ++-- server/internal/proxy/passthrough.go | 43 ++++++++++++---------- server/internal/proxy/sse_fix_test.go | 48 +++++++++++++++++++++++++ 3 files changed, 77 insertions(+), 21 deletions(-) create mode 100644 server/internal/proxy/sse_fix_test.go diff --git a/server/internal/proxy/convert/stream.go b/server/internal/proxy/convert/stream.go index 644a82f..fe2732a 100644 --- a/server/internal/proxy/convert/stream.go +++ b/server/internal/proxy/convert/stream.go @@ -11,16 +11,17 @@ type sseState struct { } // parseLine 解析一行 SSE;返回是否 data 行及其内容、是否 [DONE]。 +// data: 后可跟空格(标准)或紧贴 JSON(上游如火山方舟会省略空格)。 func (s *sseState) parseLine(line []byte) (isData bool, data string, done bool) { str := strings.TrimRight(string(line), "\r\n") switch { case strings.HasPrefix(str, "event: "): s.event = strings.TrimSpace(strings.TrimPrefix(str, "event: ")) return false, "", false - case str == "data: [DONE]": + case str == "data: [DONE]" || str == "data:[DONE]": return true, "[DONE]", true - case strings.HasPrefix(str, "data: "): - return true, strings.TrimPrefix(str, "data: "), false + case strings.HasPrefix(str, "data:"): + return true, strings.TrimLeft(strings.TrimPrefix(str, "data:"), " "), false default: return false, "", false } diff --git a/server/internal/proxy/passthrough.go b/server/internal/proxy/passthrough.go index f711cc8..c56cfc0 100644 --- a/server/internal/proxy/passthrough.go +++ b/server/internal/proxy/passthrough.go @@ -101,18 +101,32 @@ func requestText(body []byte) string { return strings.Join(parts, "\n") } -// sseContentText 提取一条 SSE 中的内容文本(chat delta.content / responses delta / messages delta.text)。 -// 兼容单 data: 行与 event:+data: 多行块(转换器 eventLine 产出的块)。 -func sseContentText(line []byte) string { +// sseDataPayload 提取一条 SSE 的 JSON 载荷(去掉 data: 前缀与空白)。 +// 兼容三种写法: +// - 单 data: 行:data: {...} 或 data:{...}(上游如火山方舟会省略 data: 后的空格) +// - event:+data: 多行块:转换器 eventLine 产出的块(event: xxx\ndata: {...} 拼在一个 []byte) +func sseDataPayload(line []byte) (string, bool) { s := string(line) - // 多行块:取最后一个 data: 行(event: 头 + data: 载荷拼在一个 []byte 里) - if idx := strings.LastIndex(s, "\ndata: "); idx >= 0 { - s = s[idx+len("\ndata: "):] - } else if strings.HasPrefix(s, "data: ") { - s = strings.TrimPrefix(s, "data: ") + idx := strings.LastIndex(s, "\ndata:") + if idx >= 0 { + s = s[idx+len("\ndata:"):] // 跳过 event: 头,落在 data: 之后 + } else if strings.HasPrefix(s, "data:") { + s = strings.TrimPrefix(s, "data:") + } else { + return "", false } + s = strings.TrimLeft(s, " ") // data: 后的可选空格 s = strings.TrimSpace(s) if s == "" || s == "[DONE]" { + return "", false + } + return s, true +} + +// sseContentText 提取一条 SSE 中的内容文本(chat delta.content / responses delta / messages delta.text)。 +func sseContentText(line []byte) string { + s, ok := sseDataPayload(line) + if !ok { return "" } var m map[string]any @@ -394,18 +408,11 @@ func extractUsage(data []byte) json.RawMessage { // scanUsage 从 SSE 一行中提取 usage(OpenAI 末块 / responses completed / messages message_delta 等)。 func scanUsage(line []byte) json.RawMessage { - s := string(line) - if !strings.Contains(s, `"usage"`) { + if !bytes.Contains(line, []byte(`"usage"`)) { return nil } - // 兼容单 data: 行与 event:+data: 多行块(转换器 eventLine 产出的块) - if idx := strings.LastIndex(s, "\ndata: "); idx >= 0 { - s = s[idx+len("\ndata: "):] - } else if strings.HasPrefix(s, "data: ") { - s = strings.TrimPrefix(s, "data: ") - } - s = strings.TrimSpace(s) - if s == "[DONE]" || s == "" { + s, ok := sseDataPayload(line) + if !ok { return nil } var m map[string]json.RawMessage diff --git a/server/internal/proxy/sse_fix_test.go b/server/internal/proxy/sse_fix_test.go new file mode 100644 index 0000000..0135080 --- /dev/null +++ b/server/internal/proxy/sse_fix_test.go @@ -0,0 +1,48 @@ +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) + } +}