diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index ce0fab7d6..b187acf76 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -267,6 +267,10 @@ arguments 是否必须是完整 JSON **按流式与非流式区分,两者的 流式工具调用必须来自已收尾的流:只要聚合出过工具 slot,收尾时就必须已观察到本协议的完成信号,否则按截断返回 `Deserialize`。完成信号按协议判定——Chat 为非空 `choices[].finish_reason` 或 `data: [DONE]`,Responses 为 `response.completed`,Anthropic 为带 `stop_reason` 的 `message_delta` 或 `message_stop`。不能用 `data: [DONE]` 作为统一判据:MiniMax 兼容层不发该标记,只发 `finish_reason`。也不能只用 “参数是合法 JSON” 当完成证明——顶层花括号闭合只说明单个参数对象字节完整,说明不了模型是否还要发下一个工具块,更说明不了上游随后会不会报 `max_tokens` 或 error;代理超时、网关自行掐断和 HTTP/2 提前 `END_STREAM` 都表现为干净 EOF,与正常收尾在字节层无法区分。该门禁当前只覆盖工具路径;纯文本响应缺完成信号仍按成功返回并打 warn,改动前必须先确认所有在用网关的文本收尾行为。流在任何工具分片到达前就断掉时槽位为空,门禁无从触发,这是已知残留缺口。 +各协议的最终事件必须同时终止读取循环,不能只标记完成:Chat 的 `data: [DONE]`、Responses 的 `response.completed`、Anthropic 的 `message_stop` 都置终止位。服务端在最终事件后保持连接(keep-alive、SSE 网关不主动关流)时,只标记完成会让读取一路等到调用方超时。 + +工具事件的协议槽位缺失时必须失败关闭,不得跳过也不得按事件内位置猜测:槽位是并行分片唯一的归并依据。跳过会静默丢掉整个调用——只剩一个调用时才会被 `StreamUnavailable` 断言兜住,丢一半毫无察觉,而 Responses 的 `finish_reason` 恒为 `completed`,那道断言对它永远不触发;猜测则会把两个不同调用合并成一个混合体(后者的 id / name 覆盖前者,arguments 被拼接)。判定字段为 Chat 的 `delta.tool_calls[].index`、Responses 的 `output_index`、Anthropic 的 content block `index`。该约束只覆盖工具事件,纯文本增量不依赖槽位,不受影响。 + 反过来,已经收尾的流遇到尾部传输 / 解析错误时必须保留结果,不能重跑 Provider。判断“有没有值得保留的东西”要看正文或工具调用任一非空,不能只看正文——纯工具调用响应的正文本来就是空的(MiniMax 的 Anthropic 工具流恒定如此),只看正文会让这类响应每次都被丢弃,白白多跑一轮往返。保留的安全性由“协议完成信号已到 + `finish_reason` 已到 + 工具参数完整 + 错误属可容忍尾部错误”共同保证,与正常路径判据一致。Anthropic 的 `message_stop` 与 Responses 的 `response.completed` 目前只标记完成、不标记流终止,因此收尾事件之后仍会读到 EOF,这条尾部路径是常态而非边缘情况。 错误边界固定如下:`StreamUnavailable` 只表示流式响应已给出 `tool_use` / `tool_calls` 完成原因但没有聚合出任何工具 slot,供调用方回退非流式,它不承担截断语义;`EmptyResponse` 表示最终文本和工具调用都为空,纯工具响应合法;`Deserialize` 覆盖 JSON / SSE / UTF-8 解析失败、缺少 `choices[0]`、流式工具身份缺失、流式参数不完整,以及上述工具流未收尾截断。Anthropic 仍不支持 `web_search`、图片内容和纯 system 消息,必须至少有一条非 system 文本消息。 diff --git a/server-rs/crates/platform-llm/src/lib.rs b/server-rs/crates/platform-llm/src/lib.rs index 47a41a246..457bc76db 100644 --- a/server-rs/crates/platform-llm/src/lib.rs +++ b/server-rs/crates/platform-llm/src/lib.rs @@ -2758,36 +2758,44 @@ fn parse_sse_event_block( .finish_reason .as_deref() .is_some_and(|reason| !reason.trim().is_empty()), - tool_fragments: extract_chat_tool_fragments(first_choice), + tool_fragments: extract_chat_tool_fragments(first_choice)?, ..Default::default() })) } // Chat 分片:首片带 index + id + function.name,后续片只有 index + function.arguments。 -fn extract_chat_tool_fragments(choice: &ChatCompletionsChoice) -> Vec { +fn extract_chat_tool_fragments( + choice: &ChatCompletionsChoice, +) -> Result, LlmError> { let Some(tool_calls) = choice .delta .as_ref() .and_then(|delta| delta.tool_calls.as_deref()) else { - return Vec::new(); + return Ok(Vec::new()); }; tool_calls .iter() - .enumerate() - .map(|(position, tool_call)| ToolCallFragment { - slot: tool_call.index.unwrap_or(position as u64), - id: tool_call.id.clone(), - name: tool_call - .function - .as_ref() - .and_then(|function| function.name.clone()), - arguments_delta: tool_call - .function - .as_ref() - .and_then(|function| function.arguments.clone()), - arguments_complete: None, + .map(|tool_call| { + // 不能退回事件内位置:两个各含一个无 index 调用的事件都会落到槽位 0, + // 后者的 id / name 覆盖前者,arguments 还会被拼在一起。 + let slot = tool_call + .index + .ok_or_else(|| missing_tool_slot_error("Chat", "delta.tool_calls[]", "index"))?; + Ok(ToolCallFragment { + slot, + id: tool_call.id.clone(), + name: tool_call + .function + .as_ref() + .and_then(|function| function.name.clone()), + arguments_delta: tool_call + .function + .as_ref() + .and_then(|function| function.arguments.clone()), + arguments_complete: None, + }) }) .collect() } @@ -2816,6 +2824,7 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll "response.completed" => Ok(Some(ParsedStreamEvent { finish_reason: Some("completed".to_string()), is_completion: true, + is_terminal: true, tool_fragments: extract_responses_completed_tool_fragments(&parsed), ..Default::default() })), @@ -2830,9 +2839,8 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll { return Ok(None); } - let Some(slot) = responses_output_slot(&parsed) else { - return Ok(None); - }; + let slot = responses_output_slot(&parsed) + .ok_or_else(|| missing_tool_slot_error("Responses", event_type, "output_index"))?; Ok(Some(ParsedStreamEvent { tool_fragments: vec![ToolCallFragment { slot, @@ -2850,9 +2858,8 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll })) } "response.function_call_arguments.delta" => { - let Some(slot) = responses_output_slot(&parsed) else { - return Ok(None); - }; + let slot = responses_output_slot(&parsed) + .ok_or_else(|| missing_tool_slot_error("Responses", event_type, "output_index"))?; Ok(Some(ParsedStreamEvent { tool_fragments: vec![ToolCallFragment { slot, @@ -2866,9 +2873,8 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll })) } "response.function_call_arguments.done" => { - let Some(slot) = responses_output_slot(&parsed) else { - return Ok(None); - }; + let slot = responses_output_slot(&parsed) + .ok_or_else(|| missing_tool_slot_error("Responses", event_type, "output_index"))?; Ok(Some(ParsedStreamEvent { tool_fragments: vec![ToolCallFragment { slot, @@ -2944,6 +2950,16 @@ fn anthropic_block_slot(parsed: &serde_json::Value) -> Option { parsed.get("index").and_then(serde_json::Value::as_u64) } +// 槽位是并行工具分片唯一的归并依据,缺失时必须失败关闭而不是跳过或猜测:跳过会静默 +// 丢掉整个调用(只剩一个调用时才会被 StreamUnavailable 断言兜住,丢一半就毫无察觉, +// 而 Responses 的 finish_reason 恒为 completed,那道断言对它永远不触发),猜测则会把 +// 两个不同调用合并成一个混合体。 +fn missing_tool_slot_error(protocol: &str, event: &str, field: &str) -> LlmError { + LlmError::Deserialize(format!( + "LLM {protocol} 流式工具事件缺少槽位字段 {field}:event={event}" + )) +} + fn parse_anthropic_sse_event(data: &str) -> Result, LlmError> { let parsed: serde_json::Value = serde_json::from_str(data).map_err(|error| { LlmError::Deserialize(format!("解析 LLM Anthropic SSE 事件失败:{error}")) @@ -2965,9 +2981,8 @@ fn parse_anthropic_sse_event(data: &str) -> Result, Ll { return Ok(None); } - let Some(slot) = anthropic_block_slot(&parsed) else { - return Ok(None); - }; + let slot = anthropic_block_slot(&parsed) + .ok_or_else(|| missing_tool_slot_error("Anthropic", event_type, "index"))?; Ok(Some(ParsedStreamEvent { tool_fragments: vec![ToolCallFragment { slot, @@ -2992,9 +3007,8 @@ fn parse_anthropic_sse_event(data: &str) -> Result, Ll .unwrap_or_default(); if delta_type == "input_json_delta" { - let Some(slot) = anthropic_block_slot(&parsed) else { - return Ok(None); - }; + let slot = anthropic_block_slot(&parsed) + .ok_or_else(|| missing_tool_slot_error("Anthropic", event_type, "index"))?; return Ok(Some(ParsedStreamEvent { tool_fragments: vec![ToolCallFragment { slot, @@ -3039,6 +3053,7 @@ fn parse_anthropic_sse_event(data: &str) -> Result, Ll // 收尾信号,所以单独用 is_completion 记录:兼容网关可能只发它而漏 stop_reason。 "message_stop" => Ok(Some(ParsedStreamEvent { is_completion: true, + is_terminal: true, ..Default::default() })), "error" => { @@ -4272,9 +4287,9 @@ mod tests { "\n\n", r#"data: {"type":"content_block_stop","index":0}"#, "\n\n", + // 刻意不发 message_stop:它是终止事件,会让读取循环在尾部错误之前就收口, + // 本用例要验证的正是“收尾信号已到、流却没干净结束”时的保留行为。 r#"data: {"type":"message_delta","delta":{"stop_reason":"tool_use"}}"#, - "\n\n", - r#"data: {"type":"message_stop"}"#, "\n\n" ); let raw_response = format!( @@ -5177,6 +5192,168 @@ mod tests { expect_tool_call_deserialize_error(error, "流式工具调用来自未完成的响应"); } + // 终止事件:服务端在最终事件之后保持连接不关时,stream_run 必须立即收口, + // 否则会一路等到调用方超时。既有覆盖只有 Chat 的 [DONE]。 + async fn assert_stream_stops_at_terminal_event( + api_kind: LlmApiKind, + sse_body: &'static str, + expected_text: &str, + ) { + let listener = TcpListener::bind("127.0.0.1:0").expect("listener should bind"); + let address = listener.local_addr().expect("listener should have addr"); + let (stream_done_sender, stream_done_receiver) = std::sync::mpsc::channel(); + let server_handle = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("request should connect"); + read_request(&mut stream); + stream + .write_all( + format!( + concat!( + "HTTP/1.1 200 OK\r\n", + "Content-Type: text/event-stream; charset=utf-8\r\n", + "Connection: keep-alive\r\n\r\n", + "{}" + ), + sse_body + ) + .as_bytes(), + ) + .expect("stream response should be written"); + stream.flush().expect("stream response should flush"); + stream_done_receiver + .recv_timeout(StdDuration::from_secs(1)) + .expect("stream_run should complete before upstream closes"); + }); + + let client = build_test_client(format!("http://{address}"), 0); + let response = tokio::time::timeout( + StdDuration::from_secs(1), + client.stream_run( + LlmRunRequest::single_turn("系统", "用户").with_api_kind(api_kind), + |_| {}, + ), + ) + .await + .expect("最终事件应在上游关闭前收口") + .expect("stream_run should succeed"); + + assert_eq!(response.text, expected_text); + stream_done_sender + .send(()) + .expect("server should still keep the stream open"); + server_handle.join().expect("server thread should join"); + } + + #[tokio::test] + async fn stream_run_stops_at_responses_completed_without_waiting_for_eof() { + assert_stream_stops_at_terminal_event( + LlmApiKind::OpenAiResponses, + concat!( + r#"data: {"type":"response.output_text.delta","delta":"你好"}"#, + "\n\n", + r#"data: {"type":"response.completed"}"#, + "\n\n" + ), + "你好", + ) + .await; + } + + #[tokio::test] + async fn stream_run_stops_at_anthropic_message_stop_without_waiting_for_eof() { + assert_stream_stops_at_terminal_event( + LlmApiKind::Anthropic, + concat!( + r#"data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"你好"}}"#, "\n\n", + r#"data: {"type":"message_delta","delta":{"stop_reason":"end_turn"}}"#, "\n\n", + r#"data: {"type":"message_stop"}"#, "\n\n" + ), + "你好", + ) + .await; + } + + // 槽位缺失:并行分片唯一的归并依据没了,跳过会静默丢调用,按事件内位置猜会把两个 + // 不同调用合并成混合体,两者都比直接失败危险。 + async fn expect_stream_missing_slot_error(api_kind: LlmApiKind, body: &str) { + let server_url = spawn_mock_server(vec![MockResponse { + status_line: "200 OK", + content_type: "text/event-stream; charset=utf-8", + body: body.to_string(), + extra_headers: Vec::new(), + }]); + + let error = build_test_client(server_url, 0) + .stream_run(weather_tool_request(api_kind), |_| {}) + .await + .expect_err("缺少协议槽位必须失败关闭"); + + expect_tool_call_deserialize_error(error, "流式工具事件缺少槽位字段"); + } + + #[tokio::test] + async fn stream_run_rejects_chat_tool_fragments_without_index() { + // 两个事件各带一个无 index 调用:旧实现都归到槽位 0,后者覆盖前者的身份。 + expect_stream_missing_slot_error( + LlmApiKind::OpenAiChat, + concat!( + r#"data: {"choices":[{"delta":{"tool_calls":[{"id":"call_a","function":{"name":"get_weather","arguments":"{\"city\":\"杭州\"}"}}]}}]}"#, "\n\n", + r#"data: {"choices":[{"delta":{"tool_calls":[{"id":"call_b","function":{"name":"get_air_quality","arguments":"{\"city\":\"杭州\"}"}}]}}]}"#, "\n\n", + r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#, "\n\n", + "data: [DONE]\n\n" + ), + ) + .await; + } + + #[tokio::test] + async fn stream_run_rejects_responses_function_call_without_output_index() { + expect_stream_missing_slot_error( + LlmApiKind::OpenAiResponses, + concat!( + r#"data: {"type":"response.output_item.added","item":{"type":"function_call","call_id":"call_1","name":"get_weather"}}"#, "\n\n", + r#"data: {"type":"response.completed"}"#, "\n\n" + ), + ) + .await; + } + + #[tokio::test] + async fn stream_run_rejects_anthropic_tool_use_without_block_index() { + expect_stream_missing_slot_error( + LlmApiKind::Anthropic, + concat!( + r#"data: {"type":"content_block_start","content_block":{"type":"tool_use","id":"call_1","name":"get_weather","input":{}}}"#, "\n\n", + r#"data: {"type":"message_delta","delta":{"stop_reason":"tool_use"}}"#, "\n\n" + ), + ) + .await; + } + + #[tokio::test] + async fn stream_run_keeps_text_only_anthropic_events_without_block_index() { + // 作用域守卫:只有工具事件收紧。text_delta 不依赖槽位,缺 index 不应受影响。 + 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":"content_block_delta","delta":{"type":"text_delta","text":"杭州多云。"}}"#, "\n\n", + r#"data: {"type":"message_delta","delta":{"stop_reason":"end_turn"}}"#, "\n\n", + r#"data: {"type":"message_stop"}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let response = build_test_client(server_url, 0) + .stream_run(weather_tool_request(LlmApiKind::Anthropic), |_| {}) + .await + .expect("纯文本事件不依赖槽位"); + + assert_eq!(response.text, "杭州多云。"); + assert!(response.tool_calls.is_empty()); + } + #[tokio::test] async fn stream_chat_tool_call_with_incomplete_arguments_json_still_fails() { // 与上一条成对:流式的参数半截意味着流被截断,是传输层事实,必须报错。