From eaf83f4fbe9359decd923f0d5e4fb96ad2451605 Mon Sep 17 00:00:00 2001 From: tgbdhs <49755039+Alexlangl@users.noreply.github.com> Date: Wed, 25 Mar 2026 22:06:21 +0800 Subject: [PATCH] fix(proxy): parse SSE fields with optional spaces in streaming handlers (#1664) * fix(proxy): handle SSE fields with or without spaces * refactor(proxy): deduplicate SSE field parsing --- src-tauri/src/proxy/mod.rs | 1 + src-tauri/src/proxy/providers/streaming.rs | 11 +++++-- .../proxy/providers/streaming_responses.rs | 9 ++++-- src-tauri/src/proxy/response_handler.rs | 23 +++++++++++++- src-tauri/src/proxy/response_processor.rs | 24 +++++++++++++- src-tauri/src/proxy/sse.rs | 31 +++++++++++++++++++ 6 files changed, 91 insertions(+), 8 deletions(-) create mode 100644 src-tauri/src/proxy/sse.rs diff --git a/src-tauri/src/proxy/mod.rs b/src-tauri/src/proxy/mod.rs index 5378823a6..4386230a7 100644 --- a/src-tauri/src/proxy/mod.rs +++ b/src-tauri/src/proxy/mod.rs @@ -22,6 +22,7 @@ pub mod response_handler; pub mod response_processor; pub(crate) mod server; pub mod session; +pub(crate) mod sse; pub mod thinking_budget_rectifier; pub mod thinking_optimizer; pub mod thinking_rectifier; diff --git a/src-tauri/src/proxy/providers/streaming.rs b/src-tauri/src/proxy/providers/streaming.rs index 59c338949..599499016 100644 --- a/src-tauri/src/proxy/providers/streaming.rs +++ b/src-tauri/src/proxy/providers/streaming.rs @@ -2,6 +2,7 @@ //! //! 实现 OpenAI SSE → Anthropic SSE 格式转换 +use crate::proxy::sse::strip_sse_field; use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use serde::{Deserialize, Serialize}; @@ -118,7 +119,7 @@ pub fn create_anthropic_sse_stream( } for l in line.lines() { - if let Some(data) = l.strip_prefix("data: ") { + if let Some(data) = strip_sse_field(l, "data") { if data.trim() == "[DONE]" { log::debug!("[Claude/OpenRouter] <<< OpenAI SSE: [DONE]"); let event = json!({"type": "message_stop"}); @@ -609,7 +610,9 @@ mod tests { let events: Vec = merged .split("\n\n") .filter_map(|block| { - let data = block.lines().find_map(|line| line.strip_prefix("data: "))?; + let data = block + .lines() + .find_map(|line| strip_sse_field(line, "data"))?; serde_json::from_str::(data).ok() }) .collect(); @@ -694,7 +697,9 @@ mod tests { let events: Vec = merged .split("\n\n") .filter_map(|block| { - let data = block.lines().find_map(|line| line.strip_prefix("data: "))?; + let data = block + .lines() + .find_map(|line| strip_sse_field(line, "data"))?; serde_json::from_str::(data).ok() }) .collect(); diff --git a/src-tauri/src/proxy/providers/streaming_responses.rs b/src-tauri/src/proxy/providers/streaming_responses.rs index 285616ea5..6f97b7283 100644 --- a/src-tauri/src/proxy/providers/streaming_responses.rs +++ b/src-tauri/src/proxy/providers/streaming_responses.rs @@ -9,6 +9,7 @@ //! 与 Chat Completions 的 delta chunk 模型完全不同,需要独立的状态机处理。 use super::transform_responses::{build_anthropic_usage_from_responses, map_responses_stop_reason}; +use crate::proxy::sse::strip_sse_field; use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use serde_json::{json, Value}; @@ -133,9 +134,9 @@ pub fn create_anthropic_sse_stream_from_responses( let mut data_parts: Vec = Vec::new(); for line in block.lines() { - if let Some(evt) = line.strip_prefix("event: ") { + if let Some(evt) = strip_sse_field(line, "event") { event_type = Some(evt.trim().to_string()); - } else if let Some(d) = line.strip_prefix("data: ") { + } else if let Some(d) = strip_sse_field(line, "data") { data_parts.push(d.to_string()); } } @@ -810,7 +811,9 @@ mod tests { let events: Vec = merged .split("\n\n") .filter_map(|block| { - let data = block.lines().find_map(|line| line.strip_prefix("data: "))?; + let data = block + .lines() + .find_map(|line| strip_sse_field(line, "data"))?; serde_json::from_str::(data).ok() }) .collect(); diff --git a/src-tauri/src/proxy/response_handler.rs b/src-tauri/src/proxy/response_handler.rs index 7045643a0..408e66a86 100644 --- a/src-tauri/src/proxy/response_handler.rs +++ b/src-tauri/src/proxy/response_handler.rs @@ -5,6 +5,7 @@ use super::session::ProxySession; use super::usage::parser::TokenUsage; use super::ProxyError; +use crate::proxy::sse::strip_sse_field; use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use serde_json::Value; @@ -90,7 +91,7 @@ impl StreamHandler { buffer = buffer[pos + 2..].to_string(); for line in event_text.lines() { - if let Some(data) = line.strip_prefix("data: ") { + if let Some(data) = strip_sse_field(line, "data") { if data.trim() != "[DONE]" { if let Ok(json) = serde_json::from_str::(data) { let mut guard = events.lock().await; @@ -211,4 +212,24 @@ mod tests { let handler = StreamHandler::new(30); assert_eq!(handler.idle_timeout, Duration::from_secs(30)); } + + #[test] + fn test_strip_sse_field_accepts_optional_space() { + assert_eq!( + super::strip_sse_field("data: {\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + super::strip_sse_field("data:{\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + super::strip_sse_field("event: message_start", "event"), + Some("message_start") + ); + assert_eq!( + super::strip_sse_field("event:message_start", "event"), + Some("message_start") + ); + } } diff --git a/src-tauri/src/proxy/response_processor.rs b/src-tauri/src/proxy/response_processor.rs index 18ab278e1..9f329b678 100644 --- a/src-tauri/src/proxy/response_processor.rs +++ b/src-tauri/src/proxy/response_processor.rs @@ -6,6 +6,7 @@ use super::{ handler_config::UsageParserConfig, handler_context::{RequestContext, StreamingTimeoutConfig}, server::ProxyState, + sse::strip_sse_field, usage::parser::TokenUsage, ProxyError, }; @@ -527,7 +528,7 @@ pub fn create_logged_passthrough_stream( if !event_text.trim().is_empty() { // 提取 data 部分并尝试解析为 JSON for line in event_text.lines() { - if let Some(data) = line.strip_prefix("data: ") { + if let Some(data) = strip_sse_field(line, "data") { if data.trim() != "[DONE]" { if let Ok(json_value) = serde_json::from_str::(data) { if let Some(c) = &collector { @@ -591,6 +592,27 @@ mod tests { use std::sync::Arc; use tokio::sync::RwLock; + #[test] + fn test_strip_sse_field_accepts_optional_space() { + assert_eq!( + super::strip_sse_field("data: {\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + super::strip_sse_field("data:{\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + super::strip_sse_field("event: message_start", "event"), + Some("message_start") + ); + assert_eq!( + super::strip_sse_field("event:message_start", "event"), + Some("message_start") + ); + assert_eq!(super::strip_sse_field("id:1", "data"), None); + } + fn build_state(db: Arc) -> ProxyState { ProxyState { db: db.clone(), diff --git a/src-tauri/src/proxy/sse.rs b/src-tauri/src/proxy/sse.rs new file mode 100644 index 000000000..235edce8b --- /dev/null +++ b/src-tauri/src/proxy/sse.rs @@ -0,0 +1,31 @@ +#[inline] +pub(crate) fn strip_sse_field<'a>(line: &'a str, field: &str) -> Option<&'a str> { + line.strip_prefix(&format!("{field}: ")) + .or_else(|| line.strip_prefix(&format!("{field}:"))) +} + +#[cfg(test)] +mod tests { + use super::strip_sse_field; + + #[test] + fn strip_sse_field_accepts_optional_space() { + assert_eq!( + strip_sse_field("data: {\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + strip_sse_field("data:{\"ok\":true}", "data"), + Some("{\"ok\":true}") + ); + assert_eq!( + strip_sse_field("event: message_start", "event"), + Some("message_start") + ); + assert_eq!( + strip_sse_field("event:message_start", "event"), + Some("message_start") + ); + assert_eq!(strip_sse_field("id:1", "data"), None); + } +}