From abbf43ea6fa9f36d0cbd49b81083294e64539b44 Mon Sep 17 00:00:00 2001 From: YoVinchen Date: Mon, 6 Apr 2026 11:09:23 +0800 Subject: [PATCH] refactor(proxy): extract take_sse_block helper with CRLF delimiter support Replace inline `buffer.find("\n\n")` SSE splitting logic across streaming, streaming_responses, response_handler, and response_processor with a shared `take_sse_block` function that handles both `\n\n` and `\r\n\r\n` delimiters. --- src-tauri/src/proxy/providers/streaming.rs | 7 +--- .../proxy/providers/streaming_responses.rs | 7 +--- src-tauri/src/proxy/response_handler.rs | 7 +--- src-tauri/src/proxy/response_processor.rs | 7 +--- src-tauri/src/proxy/sse.rs | 42 ++++++++++++++++++- 5 files changed, 49 insertions(+), 21 deletions(-) diff --git a/src-tauri/src/proxy/providers/streaming.rs b/src-tauri/src/proxy/providers/streaming.rs index 65a10260e..f6c70d919 100644 --- a/src-tauri/src/proxy/providers/streaming.rs +++ b/src-tauri/src/proxy/providers/streaming.rs @@ -2,7 +2,7 @@ //! //! 实现 OpenAI SSE → Anthropic SSE 格式转换 -use crate::proxy::sse::strip_sse_field; +use crate::proxy::sse::{strip_sse_field, take_sse_block}; use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use serde::{Deserialize, Serialize}; @@ -110,10 +110,7 @@ pub fn create_anthropic_sse_stream( let text = String::from_utf8_lossy(&bytes); buffer.push_str(&text); - while let Some(pos) = buffer.find("\n\n") { - let line = buffer[..pos].to_string(); - buffer = buffer[pos + 2..].to_string(); - + while let Some(line) = take_sse_block(&mut buffer) { if line.trim().is_empty() { continue; } diff --git a/src-tauri/src/proxy/providers/streaming_responses.rs b/src-tauri/src/proxy/providers/streaming_responses.rs index ea9274ff8..32510abdb 100644 --- a/src-tauri/src/proxy/providers/streaming_responses.rs +++ b/src-tauri/src/proxy/providers/streaming_responses.rs @@ -9,7 +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 crate::proxy::sse::{strip_sse_field, take_sse_block}; use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use serde_json::{json, Value}; @@ -122,10 +122,7 @@ pub fn create_anthropic_sse_stream_from_responses(line: &'a str, field: &str) -> Option<&'a str> .or_else(|| line.strip_prefix(&format!("{field}:"))) } +#[inline] +pub(crate) fn take_sse_block(buffer: &mut String) -> Option { + let mut best: Option<(usize, usize)> = None; + + for (delimiter, len) in [("\r\n\r\n", 4usize), ("\n\n", 2usize)] { + if let Some(pos) = buffer.find(delimiter) { + if best.is_none_or(|(best_pos, _)| pos < best_pos) { + best = Some((pos, len)); + } + } + } + + let (pos, len) = best?; + let block = buffer[..pos].to_string(); + buffer.drain(..pos + len); + Some(block) +} + #[cfg(test)] mod tests { - use super::strip_sse_field; + use super::{strip_sse_field, take_sse_block}; #[test] fn strip_sse_field_accepts_optional_space() { @@ -28,4 +46,26 @@ mod tests { ); assert_eq!(strip_sse_field("id:1", "data"), None); } + + #[test] + fn take_sse_block_supports_lf_delimiters() { + let mut buffer = "data: {\"ok\":true}\n\nrest".to_string(); + + assert_eq!( + take_sse_block(&mut buffer), + Some("data: {\"ok\":true}".to_string()) + ); + assert_eq!(buffer, "rest"); + } + + #[test] + fn take_sse_block_supports_crlf_delimiters() { + let mut buffer = "data: {\"ok\":true}\r\n\r\nrest".to_string(); + + assert_eq!( + take_sse_block(&mut buffer), + Some("data: {\"ok\":true}".to_string()) + ); + assert_eq!(buffer, "rest"); + } }