修复: OpenAI SSE error 识别 + 条件引擎 warn + MidStream stop_flag 优先 + 双窗口 listener 时序加固

This commit is contained in:
lxy
2026-08-01 14:53:03 +08:00
parent ba578dfe6c
commit db99e7107c
6 changed files with 161 additions and 8 deletions
+82 -1
View File
@@ -231,6 +231,49 @@ impl OpenAICompatProvider {
}
}
/// 生成 messages 诊断摘要(每条 role + content 形态 + tool 标记),不含敏感数据。
/// 流中途 error 时附摘要定位哪条非法(对齐 `AnthropicCompatProvider::summarize_messages`)。
fn summarize_openai_messages(messages: &[OpenAiMessage]) -> String {
let lines: Vec<String> = messages
.iter()
.enumerate()
.map(|(i, m)| {
let role = m.role.as_str();
let desc = match &m.content {
serde_json::Value::String(s) => format!("text({}B)", s.len()),
serde_json::Value::Array(blocks) => {
let parts: Vec<String> = blocks
.iter()
.map(|b| {
let ty = b.get("type").and_then(|t| t.as_str()).unwrap_or("?");
match ty {
"text" => format!(
"text({}B)",
b.get("text")
.and_then(|t| t.as_str())
.map(|s| s.len())
.unwrap_or(0)
),
"image_url" => "image".to_string(),
_ => ty.to_string(),
}
})
.collect();
format!("[{}]", parts.join(","))
}
_ => "?".to_string(),
};
let tool_mark = match (&m.tool_calls, &m.tool_call_id) {
(Some(tcs), _) => format!(" tool_calls={}", tcs.len()),
(None, Some(tid)) => format!(" tool_result[tid={}]", tid),
(None, None) => String::new(),
};
format!("#{}:{} {}{}", i, role, desc, tool_mark)
})
.collect();
format!("{} msgs: {}", lines.len(), lines.join(" | "))
}
/// 保证 messages 首条为 user/system(OpenAI 协议要求首条非 assistant/tool)。
///
/// 对齐 `AnthropicCompatProvider::ensure_leading_user`。上游绕过 `ContextManager::sanitize_messages`
@@ -419,6 +462,9 @@ impl LlmProvider for OpenAICompatProvider {
// (严格 UTF-8 + SSE 协议校验,跨 chunk 字符/不完整事件均报错且不可恢复)。
// 原生解析器:bytes 累积 + from_utf8_lossy 宽松处理 + \n\n 分隔,容错不中断流。
let mut last_usage: Option<TokenUsage> = None;
// MidStream error(中转站按 OpenAI 协议在流中途发 error 帧)时附 messages 摘要定位哪条非法
// (对齐 anthropic_compat 672)。
let messages_summary = Self::summarize_openai_messages(&openai_req.messages);
let sse = crate::sse_parser::SseStream::new(resp.bytes_stream());
let stream = sse.flat_map(move |result: Result<Vec<String>, String>| {
@@ -426,7 +472,10 @@ impl LlmProvider for OpenAICompatProvider {
match result {
Ok(events) => {
for data in events {
let chunk = apply_openai_sse(&data, &mut last_usage);
let mut chunk = apply_openai_sse(&data, &mut last_usage);
if let Some(err) = chunk.error.as_mut() {
*err = format!("{} | messages 摘要: {}", err, messages_summary);
}
chunks.push(Ok(chunk));
}
}
@@ -616,6 +665,38 @@ mod tests {
assert!(!c.finished);
}
/// 流中途 error 事件 → error 为 Some(msg)finished=false(避免残缺被当正常完成入库),不污染 usage 累加
#[test]
fn openai_sse_midstream_error_event() {
let mut acc: Option<TokenUsage> = None;
// 先累积一段 usage,验证 error 分支不污染累加器
apply_openai_sse(&usage_only_chunk(10, 20), &mut acc);
let data = r#"{"choices":[],"error":{"message":"context length exceeded","type":"invalid_request_error"}}"#;
let c = apply_openai_sse(data, &mut acc);
assert!(!c.finished, "error 帧不应走 finished 完成路径");
assert_eq!(c.delta, "");
assert!(c.tool_calls.is_none());
assert!(c.usage.is_none(), "error 帧不应带出 usage");
let err = c.error.expect("error 帧应映射为 Some(msg)");
assert_eq!(err, "context length exceeded");
// 累加器保持原值(未被覆盖/清空)
let acc = acc.expect("累加器应保留先前 usage 不受 error 影响");
assert_eq!(acc.prompt_tokens, 10);
assert_eq!(acc.completion_tokens, 20);
}
/// error 无 message 字段 → 兜底 "stream error" 字符串
#[test]
fn openai_sse_midstream_error_without_message_falls_back() {
let mut acc: Option<TokenUsage> = None;
// error 形态异常(只有 type,无 message
let data = r#"{"choices":[],"error":{"type":"server_error"}}"#;
let c = apply_openai_sse(data, &mut acc);
assert!(!c.finished);
assert_eq!(c.error.as_deref(), Some("stream error"), "无 message 字段应兜底");
}
// ---------- 多模态 convert_request ----------
/// 含图消息 → content 数组(text + image_url data URI);纯文本 → 字符串简写
+25 -1
View File
@@ -9,7 +9,7 @@
//! 零行为变更(纯搬迁)。结构对齐 `anthropic_helpers.rs`。
use serde::{Deserialize, Serialize};
use tracing::debug;
use tracing::{debug, error};
use crate::provider::{StreamChunk, TokenUsage, ToolCallDelta};
@@ -111,6 +111,11 @@ pub(crate) struct OpenAiStreamChunk {
/// 末 chunkchoices 为空)携带的累计 usage
#[serde(default)]
pub usage: Option<OpenAiUsage>,
/// 流中途 error 事件(OpenAI 兼容协议:`{"error":{"message":..,"type":..}}`)。
/// 部分中转站按 OpenAI 协议在流中途发 error 帧而非走 HTTP 非 200
/// serde default + Value 兜底:旧响应无此字段不受影响,且对 error 载荷形态不敏感。
#[serde(default)]
pub error: Option<serde_json::Value>,
}
#[derive(Debug, Deserialize)]
@@ -170,6 +175,25 @@ pub(crate) fn apply_openai_sse(data: &str, usage_accum: &mut Option<TokenUsage>)
match serde_json::from_str::<OpenAiStreamChunk>(data) {
Ok(chunk) => {
// 流中途 error 事件(中转站按 OpenAI 协议在流中途发 error 帧)。
// 不走 finished 完成路径(避免残缺响应被当正常完成入库),由 stream_llm
// 识别 error 非空 → 发 AiError + 丢弃残缺(对齐 anthropic_helpers 215-219)。
if let Some(err_val) = chunk.error {
let msg = err_val
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("stream error")
.to_string();
error!(%msg, raw = %err_val, "OpenAI 流式错误事件");
return StreamChunk {
delta: String::new(),
finished: false,
tool_calls: None,
usage: None,
error: Some(msg),
reasoning_content: None,
};
}
// 提取 usage(带 include_usage 时末段 chunk 携带,覆盖累积)
if let Some(u) = chunk.usage {
*usage_accum = Some(TokenUsage {
+37 -4
View File
@@ -38,20 +38,53 @@ impl ConditionEngine {
pub fn evaluate(expr: &str, context: &Value) -> anyhow::Result<bool> {
let trimmed = expr.trim();
if trimmed.is_empty() {
// 空表达式:保守 false。配置漏写条件时留可观测线索(非用户笔误,但利于排查
// 「为什么这条边条件总不满足」——可能是上游未填 condition 字段)。
tracing::warn!(
target: "df_workflow::conditions",
"条件表达式为空,求值保守 false(检查 edge.condition 是否漏填)"
);
return Ok(false);
}
let toks = match tokenize(trimmed) {
Ok(t) => t,
// tokenizer 阶段即非法(如未闭合引号):保守 false,不静默放行
Err(_) => return Ok(false),
// tokenizer 阶段即非法(如未闭合引号、单 = 、未识别关键字):保守 false,不静默放行
// 留 warn 记录表达式片段,排查笔误(如 status 拼错、引号未闭合)有据可循。
// 表达式可能较长,仅取前 200 字符防日志膨胀。
Err(_) => {
let preview = if trimmed.len() > 200 {
format!("{}...(截断)", &trimmed[..trimmed.floor_char_boundary(200)])
} else {
trimmed.to_string()
};
tracing::warn!(
target: "df_workflow::conditions",
expr = %preview,
"条件表达式 tokenize 失败(疑似未闭合引号/非法符号/未识别关键字),求值保守 false"
);
return Ok(false);
}
};
let mut parser = Parser { toks, pos: 0, ctx: context };
match parser.parse_or() {
Ok(v) => Ok(v),
// 解析错误保守 false:与历史行为「未识别表达式默认 false」一致,不破坏调用方
Err(_) => Ok(false),
// 解析错误保守 false:与历史行为「未识别表达式默认 false」一致,不破坏调用方
// 留 warn 记录表达式片段,排查结构错误(如括号不闭合、操作数缺失)有据可循。
Err(_) => {
let preview = if trimmed.len() > 200 {
format!("{}...(截断)", &trimmed[..trimmed.floor_char_boundary(200)])
} else {
trimmed.to_string()
};
tracing::warn!(
target: "df_workflow::conditions",
expr = %preview,
"条件表达式解析失败(疑似括号不闭合/操作数缺失/结构错误),求值保守 false"
);
Ok(false)
}
}
}
}