修复: DeepSeek 400 全量扫描 + 队列 per-conv 隔离
- openai_compat: 扫描所有 assistant 消息剥离 orphan tool_calls(原仅查末条) - queue 加 conversationId 字段,按会话精准 drain - regenerate/editMessage 只清本会话排队消息 - newConversation 保留旧会话排队消息 - AiError 只清出错会话的队列项
This commit is contained in:
@@ -0,0 +1,118 @@
|
||||
//! 原生 SSE 流式解析器 — 替代 eventsource-stream 库
|
||||
//!
|
||||
//! BUG-2026-07-17 根治: eventsource-stream 0.2 在 Windows 上对 Deepseek 等 provider
|
||||
//! 的 SSE 响应解析时报 "Transport error: error decoding response body" 错误。
|
||||
//!
|
||||
//! 根因分析:
|
||||
//! eventsource-stream 内部对 bytes_stream 做严格的 UTF-8 + SSE 协议校验,遇到以下情况
|
||||
//! 即报错(且不可恢复):
|
||||
//! - 流中断时未完整接收 UTF-8 字符(网络抖动常见)
|
||||
//! - 缺少结束的 \n\n(连接断开常见)
|
||||
//! - 非 ASCII 字符的多字节序列跨 chunk 边界
|
||||
//!
|
||||
//! 本解析器实现:
|
||||
//! - 宽松的 UTF-8 处理(用 bytes 累积,String::from_utf8_lossy 转换,不报错)
|
||||
//! - SSE 协议简单解析(以 \n\n 分隔事件,data: 前缀提取)
|
||||
//! - 容错:解析失败时跳过该事件继续,不中断流
|
||||
//! - 返回 Vec<String>(每个元素是一个事件 data 字段拼接内容)
|
||||
|
||||
use futures::Stream;
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
/// SSE 事件流的 data 字段内容
|
||||
pub type SseEvent = String;
|
||||
|
||||
/// 原生 SSE 解析器流:包装 bytes_stream,产出 Vec<SseEvent>(一次 poll 可能产出多个事件)
|
||||
pub struct SseStream<S> {
|
||||
inner: S,
|
||||
buffer: Vec<u8>,
|
||||
}
|
||||
|
||||
impl<S> SseStream<S>
|
||||
where
|
||||
S: Stream<Item = Result<bytes::Bytes, reqwest::Error>> + Unpin,
|
||||
{
|
||||
pub fn new(inner: S) -> Self {
|
||||
Self {
|
||||
inner,
|
||||
buffer: Vec::with_capacity(8192),
|
||||
}
|
||||
}
|
||||
|
||||
/// 从 buffer 解析完整的 SSE 事件(以 \n\n 分隔),返回事件列表
|
||||
fn parse_events(&mut self) -> Vec<SseEvent> {
|
||||
let mut events = Vec::new();
|
||||
loop {
|
||||
let sep_pos = self.buffer.windows(2).position(|w| w == b"\n\n");
|
||||
if sep_pos.is_none() {
|
||||
break;
|
||||
}
|
||||
let sep_pos = sep_pos.unwrap();
|
||||
let event_bytes: Vec<u8> = self.buffer.drain(..sep_pos + 2).collect();
|
||||
// 去掉末尾的 \n\n
|
||||
let body_end = event_bytes.len().saturating_sub(2);
|
||||
let event_text = String::from_utf8_lossy(&event_bytes[..body_end]);
|
||||
let data = Self::extract_data_fields(&event_text);
|
||||
if !data.is_empty() {
|
||||
events.push(data);
|
||||
}
|
||||
}
|
||||
events
|
||||
}
|
||||
|
||||
/// 从 SSE 事件文本中提取所有 data: 行的内容,拼接为单个字符串(多个 data 行用 \n 连接)
|
||||
fn extract_data_fields(event_text: &str) -> String {
|
||||
let mut data_parts: Vec<&str> = Vec::new();
|
||||
for line in event_text.lines() {
|
||||
if let Some(rest) = line.strip_prefix("data:") {
|
||||
let rest = rest.strip_prefix(' ').unwrap_or(rest);
|
||||
data_parts.push(rest);
|
||||
}
|
||||
// 忽略 event:/id:/retry: 等其他 SSE 字段(OpenAI/Anthropic 协议未使用)
|
||||
}
|
||||
data_parts.join("\n")
|
||||
}
|
||||
}
|
||||
|
||||
impl<S> Stream for SseStream<S>
|
||||
where
|
||||
S: Stream<Item = Result<bytes::Bytes, reqwest::Error>> + Unpin,
|
||||
{
|
||||
type Item = Result<Vec<SseEvent>, String>;
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
use futures::StreamExt;
|
||||
loop {
|
||||
// 先尝试从 buffer 解析完整事件
|
||||
let events = self.parse_events();
|
||||
if !events.is_empty() {
|
||||
return Poll::Ready(Some(Ok(events)));
|
||||
}
|
||||
|
||||
// buffer 不足以解析出完整事件,从 inner 读更多数据
|
||||
match self.inner.poll_next_unpin(cx) {
|
||||
Poll::Ready(Some(Ok(chunk))) => {
|
||||
self.buffer.extend_from_slice(&chunk);
|
||||
continue;
|
||||
}
|
||||
Poll::Ready(Some(Err(e))) => {
|
||||
return Poll::Ready(Some(Err(format!("SSE 流读取错误: {}", e))));
|
||||
}
|
||||
Poll::Ready(None) => {
|
||||
// 流结束,处理 buffer 中的剩余数据(可能没有 \n\n 结束的最后一段)
|
||||
if !self.buffer.is_empty() {
|
||||
let remaining = String::from_utf8_lossy(&self.buffer).to_string();
|
||||
self.buffer.clear();
|
||||
let data = Self::extract_data_fields(&remaining);
|
||||
if !data.is_empty() {
|
||||
return Poll::Ready(Some(Ok(vec![data])));
|
||||
}
|
||||
}
|
||||
return Poll::Ready(None);
|
||||
}
|
||||
Poll::Pending => return Poll::Pending,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user