diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index 7cf609977..1bbff68dd 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -269,6 +269,10 @@ arguments 是否必须是完整 JSON **按流式与非流式区分,两者的 各协议的最终事件必须同时终止读取循环,不能只标记完成:Chat 的 `data: [DONE]`、Responses 的 `response.completed` 与 `response.incomplete`、Anthropic 的 `message_stop` 都置终止位。Responses 的整体收尾信号有两个——撞到 `max_output_tokens` 时上游**只发 `response.incomplete`、不发 `response.completed`**(真实端点抓包确认),其载荷与 completed 同构,同样带完整 `output[]`,item 上标 `status=incomplete`,`incomplete_details.reason` 给出原因。漏掉它会同时造成三件事:不终止读取循环、流式 Responses 永远产生不出 `incomplete` 这个 `finish_reason`(上面那条截断拒绝规则对它形同虚设)、只在整体终态事件中携带的工具调用被静默丢掉。服务端在最终事件后保持连接(keep-alive、SSE 网关不主动关流)时,只标记完成会让读取一路等到调用方超时。 +Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复源,两者必须对称。只恢复工具会造成三种后果:纯文本的 completed-only 响应退化成 `EmptyResponse`(上层白跑一轮重试或降级);`response.incomplete` 携带的截断正文本来是可用的降级结果,同样拿不回来;“正文 + 工具调用”的响应不报错,但模型的前置说明被静默丢掉,最隐蔽。正文提取必须复用非流式那条路径(`output_text` 优先、`output[].content[]` 回退、过滤 `reasoning` / `reasoning_content` / `analysis` / `thinking` 等隐藏 part),不得另写裸 JSON 提取器——漏掉过滤层会把思维链当正文吐给调用方。终态载荷反序列化失败时按“没有快照”静默降级、不报错:这是兜底恢复路径,网关发出未建模的形状时应当退回增量累加结果;这与槽位缺失必须失败关闭的口径不同,那里放过会造成静默的身份与参数错配,这里放过只是回到没有该恢复路径时的行为。 + +终态正文按**快照覆盖**而非追加合并,且要按累加状态分两条路:累加为空时(只发终态事件的网关)必须把快照当作一次增量发出去,只覆盖累加值会让调用方的流式通道全程收不到任何文本——Responses 的 finish-only 回调开关是关闭的,只有 Chat 打开,指望终态回调兜底并不成立;累加非空时按权威值覆盖但不补发回调,否则正文在调用方侧翻倍。覆盖语义与工具参数的 `arguments_complete` 一致。 + 工具事件的协议槽位缺失时必须失败关闭,不得跳过也不得按事件内位置猜测:槽位是并行分片唯一的归并依据。跳过会静默丢掉整个调用——只剩一个调用时才可能被 `StreamUnavailable` 断言兜住,丢一半毫无察觉;Responses 的整体终态原因是 `completed` / `incomplete`,也不会触发只识别 `tool_use` / `tool_calls` 的那道断言。猜测则会把两个不同调用合并成一个混合体(后者的 id / name 覆盖前者,arguments 被拼接)。判定字段为 Chat 的 `delta.tool_calls[].index`、Responses 的 `output_index`、Anthropic 的 content block `index`。该约束只覆盖工具事件,纯文本增量不依赖槽位,不受影响。 槽位存在但被两个不同调用共用时同样必须失败关闭:同一槽位的 id 与函数名只允许**从缺失变为已知**或**重复同一个值**,出现互不相同的非空值即返回 `Deserialize`。“覆盖身份、追加参数”并不自洽——前一个调用参数为空时拼接结果就是后一个调用的合法 JSON,参数完整性检查兜不住,调用方只会拿到后一个工具,前一个静默消失;Responses 的权威完整参数还会整段覆盖,产出“前一个调用的身份配后一个调用的参数”。两者都会原样交给 Runtime 执行。已知触发路径有两条:兼容网关把 Chat 的 `index` 恒置 0,以及 Responses 的 `response.completed` 回退按 `output[]` 下标重建槽位时与流式 `output_index` 基准错位(例如 completed 载荷省略 reasoning item)。id 必须与函数名一同参与判定——并行调用同一个工具是最常见的并行场景,此时函数名相同,只有 id 能区分。空白身份按缺失跳过、不算冲突:部分兼容网关在续传分片里回发完整 `function` 对象且 `name` / `id` 为空串,按“不等即冲突”会把它们整批误杀,这也与归一层的空白即缺失约定一致。 diff --git a/server-rs/crates/platform-llm/src/lib.rs b/server-rs/crates/platform-llm/src/lib.rs index 06bc7e72f..6fb70b91b 100644 --- a/server-rs/crates/platform-llm/src/lib.rs +++ b/server-rs/crates/platform-llm/src/lib.rs @@ -612,6 +612,10 @@ struct OpenAiCompatibleSseParser { #[derive(Debug, Default)] struct ParsedStreamEvent { delta_text: Option, + // 终态事件携带的完整正文快照。必须与 delta_text 分开:它不是增量,按增量累加会让 + // 正文翻倍。只有 Responses 的 completed / incomplete 会填——Chat 的 [DONE] 与 + // Anthropic 的 message_stop 都不带载荷,那两条协议恒为 None。 + text_snapshot: Option, finish_reason: Option, usage: Option, is_terminal: bool, @@ -1993,6 +1997,7 @@ where for event in events { let ParsedStreamEvent { delta_text, + text_snapshot, finish_reason: event_finish_reason, usage: event_usage, is_terminal, @@ -2013,12 +2018,28 @@ where accumulation.push_tool_fragment(fragment)?; } - let delta_text = delta_text.unwrap_or_default(); - let has_delta = !delta_text.is_empty(); + let mut delta_text = delta_text.unwrap_or_default(); + let mut has_delta = !delta_text.is_empty(); if has_delta { accumulation.text.push_str(delta_text.as_str()); } + // 终态快照是上游给出的权威完整正文,按累加状态分两条路: + // - 累加为空(只发终态事件的网关):当成一次增量发出去。只覆盖 accumulation.text + // 的话 response.text 是对了,但调用方的流式通道全程收不到任何文本——Responses + // 的 emit_finish_only_delta 是 false,指望终态回调兜底并不成立。 + // - 累加非空:按权威值覆盖拼接结果,但不补发回调,否则正文在调用方侧翻倍。 + // 覆盖语义与工具参数的 arguments_complete 一致。 + if let Some(snapshot) = text_snapshot.filter(|text| !text.trim().is_empty()) { + if accumulation.text.trim().is_empty() { + accumulation.text = snapshot.clone(); + delta_text = snapshot; + has_delta = true; + } else { + accumulation.text = snapshot; + } + } + if let Some(event_finish_reason) = event_finish_reason { accumulation.finish_reason = Some(event_finish_reason.clone()); if has_delta || emit_finish_only_delta { @@ -2876,6 +2897,10 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll ), is_completion: true, is_terminal: true, + // 终态载荷既是工具调用的恢复源,也是正文的恢复源,两者必须对称:只恢复 + // 工具会让纯文本的 completed-only 响应变成 EmptyResponse,让「正文 + 工具」 + // 响应静默丢掉模型的前置说明。 + text_snapshot: extract_responses_terminal_text(&parsed), tool_fragments: extract_responses_completed_tool_fragments(&parsed), ..Default::default() })), @@ -2955,6 +2980,19 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll } } +// 终态事件的 response 字段就是一个完整 Response 对象,直接反序列化后复用非流式的正文 +// 提取:它已经处理了 output_text 优先、output[].content[] 回退,以及 reasoning / +// reasoning_content / analysis / thinking 这些隐藏 part 的过滤。另写裸 JSON 提取器必然 +// 漏掉过滤层,会把思维链当正文吐给调用方。 +fn extract_responses_terminal_text(parsed: &serde_json::Value) -> Option { + let response = parsed.get("response")?; + // 反序列化失败按「没有快照」处理而不是报错:这是兜底恢复路径,网关发出我们没建模 + // 的形状时应当退回增量累加结果。这与槽位缺失必须失败关闭的口径不同——那里放过会 + // 造成静默的身份/参数错配,这里放过只是回到本次修复前的行为。 + let envelope: ResponsesResponseEnvelope = serde_json::from_value(response.clone()).ok()?; + extract_responses_text(&envelope).filter(|text| !text.trim().is_empty()) +} + fn extract_responses_completed_tool_fragments(parsed: &serde_json::Value) -> Vec { let Some(items) = parsed .get("response") @@ -4889,6 +4927,102 @@ mod tests { arguments: r#"{"city":"杭州"}"#.to_string(), }] ); + // 载荷里一直带着这句正文,但过去只断言工具调用,缺陷被自己的测试盖住了: + // 终态事件既是工具调用恢复源也是正文恢复源,两者必须对称。 + assert_eq!(response.text, "我来查询。"); + } + + // 只发终态事件的网关:正文只存在于 response.output[],没有任何增量事件。 + async fn expect_responses_terminal_only_text(event_type: &str, finish_reason: &str) { + let server_url = spawn_mock_server(vec![MockResponse { + status_line: "200 OK", + content_type: "text/event-stream; charset=utf-8", + body: format!( + r#"data: {{"type":"{event_type}","response":{{"output":[{{"id":"msg_0","type":"message","content":[{{"type":"output_text","text":"杭州今天多云。"}}]}}]}}}}"# + ) + "\n\n", + extra_headers: Vec::new(), + }]); + + // 只断言 response.text 抓不到「调用方流式通道收不到文本」这个坑:Responses 的 + // emit_finish_only_delta 是 false,只覆盖累加值的话回调根本不会触发。 + let mut streamed: Vec = Vec::new(); + let response = build_test_client(server_url, 0) + .stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |delta| { + streamed.push(delta.accumulated_text.clone()); + }) + .await + .expect("终态事件携带的正文必须能恢复"); + + assert_eq!(response.text, "杭州今天多云。"); + assert_eq!(response.finish_reason.as_deref(), Some(finish_reason)); + assert!(response.tool_calls.is_empty()); + assert_eq!(streamed, vec!["杭州今天多云。".to_string()]); + } + + #[tokio::test] + async fn stream_run_recovers_responses_text_from_completed_event_only() { + // 过去这里返回 EmptyResponse,上层会白跑一轮重试或降级。 + expect_responses_terminal_only_text("response.completed", "completed").await; + } + + #[tokio::test] + async fn stream_run_recovers_responses_text_from_incomplete_event_only() { + // 撞 max_output_tokens 时上游只发 incomplete。截断正文是可用的降级结果, + // 过去同样退化成 EmptyResponse,连降级回复都给不出来。 + expect_responses_terminal_only_text("response.incomplete", "incomplete").await; + } + + #[tokio::test] + async fn stream_run_does_not_duplicate_text_when_terminal_event_repeats_deltas() { + // 终态快照按覆盖而不是追加处理,否则同时发增量和完整 output 的网关会让正文翻倍。 + let server_url = spawn_mock_server(vec![MockResponse { + status_line: "200 OK", + content_type: "text/event-stream; charset=utf-8", + body: concat!( + r#"data: {"type":"response.output_text.delta","delta":"杭州今天"}"#, "\n\n", + r#"data: {"type":"response.output_text.delta","delta":"多云。"}"#, "\n\n", + r#"data: {"type":"response.completed","response":{"output":[{"id":"msg_0","type":"message","content":[{"type":"output_text","text":"杭州今天多云。"}]}]}}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let mut streamed: Vec = Vec::new(); + let response = build_test_client(server_url, 0) + .stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |delta| { + streamed.push(delta.accumulated_text.clone()); + }) + .await + .expect("增量与终态并存时不应重复正文"); + + assert_eq!(response.text, "杭州今天多云。"); + // 已有增量时不补发回调,调用方侧同样不能翻倍。 + assert_eq!( + streamed, + vec!["杭州今天".to_string(), "杭州今天多云。".to_string()] + ); + } + + #[tokio::test] + async fn stream_run_keeps_responses_terminal_reasoning_out_of_text() { + // 锁住「必须复用 extract_responses_text」这个决定:它带隐藏 part 过滤, + // 换成裸 JSON 提取会把思维链当正文吐给调用方。 + let server_url = spawn_mock_server(vec![MockResponse { + status_line: "200 OK", + content_type: "text/event-stream; charset=utf-8", + body: concat!( + r#"data: {"type":"response.completed","response":{"output":[{"id":"rs_0","type":"reasoning","content":[{"type":"reasoning","text":"先判断用户问的是哪座城市。"}]},{"id":"msg_0","type":"message","content":[{"type":"output_text","text":"杭州今天多云。"}]}]}}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let response = build_test_client(server_url, 0) + .stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |_| {}) + .await + .expect("终态正文恢复必须过滤隐藏推理 part"); + + assert_eq!(response.text, "杭州今天多云。"); } #[tokio::test]