diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index 353c811bd..bad6bcb7b 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -259,7 +259,9 @@ npm run check:server-rs-ddd 流式收尾时,缺少工具 id / name、或非空 arguments 不是完整 JSON,均返回 `Deserialize`;空 arguments 默认归一为 `{}`。这只是 JSON 语法完整性检查,不是按工具 `parameters` 执行 JSON Schema 校验。非流式 Anthropic 缺失 `tool_use.input` 时也归一为 `{}`;其它协议的非流式 arguments 仍按上游字段解析。 -错误边界固定如下:`StreamUnavailable` 只表示流式响应已给出 `tool_use` / `tool_calls` 完成原因但没有聚合出任何工具 slot,供调用方回退非流式;`EmptyResponse` 表示最终文本和工具调用都为空,纯工具响应合法;`Deserialize` 覆盖 JSON / SSE / UTF-8 解析失败、缺少 `choices[0]`、流式工具身份缺失和流式参数不完整。Anthropic 仍不支持 `web_search`、图片内容和纯 system 消息,必须至少有一条非 system 文本消息。 +流式工具调用必须来自已收尾的流:只要聚合出过工具 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,改动前必须先确认所有在用网关的文本收尾行为。流在任何工具分片到达前就断掉时槽位为空,门禁无从触发,这是已知残留缺口。 + +错误边界固定如下:`StreamUnavailable` 只表示流式响应已给出 `tool_use` / `tool_calls` 完成原因但没有聚合出任何工具 slot,供调用方回退非流式,它不承担截断语义;`EmptyResponse` 表示最终文本和工具调用都为空,纯工具响应合法;`Deserialize` 覆盖 JSON / SSE / UTF-8 解析失败、缺少 `choices[0]`、流式工具身份缺失、流式参数不完整,以及上述工具流未收尾截断。Anthropic 仍不支持 `web_search`、图片内容和纯 system 消息,必须至少有一条非 system 文本消息。 - 图片生成:VectorEngine `gpt-image-2` 图片 provider 归属 `platform-image`,密钥只在后端环境变量中;`api-server` 内的 `openai_image_generation.rs` 只是兼容调用面和外部失败审计桥接,不再承载 provider 协议实现。实际外部生成运行记录统一落 `tracking_event`,`event_key = external_generation_run`,metadata 记录开始 / 结束时间、耗时、状态、成功标记、失败原因、provider task id 和结果摘要,不再写回过时的 `ai_task`。DashScope 只按仍在使用的历史能力单独处理,不作为 GPT-image-2 兜底。VectorEngine `/v1/images/generations` 和 `/v1/images/edits` 上游 POST 使用 `libcurl` 发送;`reqwest` 只保留给参考图 URL 下载和响应中图片 URL 下载。`/v1/images/edits` 的 multipart 参考图必须作为 libcurl 文件上传 part 发送,字段名为 `image`,实现上使用 `Form::buffer(file_name, bytes)` 并设置 `Content-Type`;不能只用 `contents(...).filename(...)`,否则上游会把请求转码为缺少图片并返回 `image is required`。`request_send` 阶段的 curl timeout / connect error 按可重试传输错误处理,最多尝试 5 次,并使用指数退避加短抖动;排障时优先看 `attempt`、`max_attempts`、`retry_delay_ms`、`reference_image_bytes_total` 和 `request_params`,不要把 `SendRequest` 当成上游业务错误。 - 抠图输入以私有 OSS 作为内存生命周期边界:生成原图和角色动作抽取帧上传时消费图片字节所有权,上传完成后不保留原图缓冲;手动去背景直接解析并校验已有 OSS object key,不下载原图。BgFilter 必须为 object key 签发 600 秒 GET URL 并通过 multipart `image_url` 提交,不用 `file` 重传;flat 链路进入阿里云 fallback 时由 `platform-matting` URL 接口单独下载并上传 `AuthorizeFileUpload` 临时对象,在推理前释放下载缓冲,继续 fallback 到本地键色时再单独下载一次原图,本地产出后释放本次原图下载缓冲。签名 URL 不得写入日志、审计或持久化。 diff --git a/server-rs/crates/platform-llm/src/lib.rs b/server-rs/crates/platform-llm/src/lib.rs index c1e471020..f90b323a4 100644 --- a/server-rs/crates/platform-llm/src/lib.rs +++ b/server-rs/crates/platform-llm/src/lib.rs @@ -610,6 +610,11 @@ struct ParsedStreamEvent { finish_reason: Option, usage: Option, is_terminal: bool, + // 本事件是协议层的收尾信号。它和 is_terminal 不同:is_terminal 只有 Chat 的 [DONE] + // 会置位,而 MiniMax 这类网关根本不发 [DONE],只能靠各协议自己的完成事件判断。 + // 也和 finish_reason 分开:Anthropic 的 message_stop 是收尾信号但不带 stop_reason, + // 不能借它写 finish_reason,否则会覆盖 message_delta 给出的真实 end_turn。 + is_completion: bool, tool_fragments: Vec, } @@ -640,6 +645,9 @@ struct StreamAccumulation { finish_reason: Option, usage: Option, tool_calls: Vec, + // 是否观察到过协议收尾信号。字节流干净结束不等于协议收尾:代理超时、网关自行掐断 + // 和 HTTP/2 提前 END_STREAM 都表现为干净 EOF,与正常收尾无法区分。 + completion_observed: bool, } impl StreamAccumulation { @@ -1413,6 +1421,37 @@ impl LlmClient { })?; } + // 截断门禁:出现过工具分片就必须已观察到协议收尾信号。字节流干净结束不构成 + // 收尾证明,参数恰好是合法 JSON 同样不构成——顶层花括号闭合只说明这一个参数 + // 对象字节完整,说明不了模型是否还要发下一个工具块,也说明不了上游随后会不会 + // 报 max_tokens 或 error。这里用未固化的槽位判断,使"参数恰好闭合"的截断仍按 + // 截断归因。 + // + // 残留缺口:流在任何工具分片到达前就断掉时槽位为空,本门禁无从触发;堵它需要 + // 同时收严纯文本路径,本轮不做,只在下方留 warn 攒线上口径。 + if !accumulation.completion_observed && !accumulation.tool_calls.is_empty() { + log_llm_raw_failure( + &self.config, + &request, + true, + 1, + "stream_tool_calls_truncated", + parser.raw_text().as_str(), + ); + return Err(LlmError::Deserialize(format!( + "LLM 流式工具调用在协议完成信号前截断:slots={}, api_kind={:?}", + accumulation.tool_calls.len(), + request.api_kind + ))); + } + if !accumulation.completion_observed && !accumulation.text.trim().is_empty() { + warn!( + "platform-llm stream ended without protocol completion signal: api_kind={:?}, text_chars={}", + request.api_kind, + accumulation.text.chars().count() + ); + } + let tool_calls = accumulation.finish_tool_calls().map_err(|error| { log_llm_raw_failure( &self.config, @@ -1760,6 +1799,7 @@ fn retain_completed_stream_after_tail_error( // 工具调用尚未拼完整时不能保留:半截参数比直接失败更危险。 let tool_calls_complete = accumulation.finish_tool_calls().is_ok(); let retain_response = !accumulation.text.trim().is_empty() + && accumulation.completion_observed && accumulation.finish_reason.is_some() && tool_calls_complete && is_tolerable_tail_error; @@ -1789,9 +1829,14 @@ where finish_reason: event_finish_reason, usage: event_usage, is_terminal, + is_completion, tool_fragments, } = event; + if is_completion { + accumulation.completion_observed = true; + } + if let Some(event_usage) = event_usage { accumulation.usage = Some(event_usage); } @@ -2508,6 +2553,7 @@ fn parse_sse_event_block( return if api_kind == LlmApiKind::OpenAiChat { Ok(Some(ParsedStreamEvent { is_terminal: true, + is_completion: true, ..Default::default() })) } else { @@ -2555,6 +2601,12 @@ fn parse_sse_event_block( delta_text: extract_message_text(first_choice), finish_reason: first_choice.finish_reason.clone(), usage: parsed.usage, + // Chat 的收尾信号是非空 finish_reason,不能只认 [DONE]:部分兼容网关(MiniMax) + // 只发前者。真 OpenAI 两者都发,这里任一到达即视为已收尾。 + is_completion: first_choice + .finish_reason + .as_deref() + .is_some_and(|reason| !reason.trim().is_empty()), tool_fragments: extract_chat_tool_fragments(first_choice), ..Default::default() })) @@ -2608,8 +2660,11 @@ fn parse_responses_sse_event(data: &str) -> Result, Ll })), // completed 事件携带完整 output;有的网关只发它而不发增量事件,这里再取一遍, // 槽位沿用 output 数组下标,与 output_index 语义一致,可安全覆盖增量拼接结果。 + // response.completed 是 Responses 唯一的整体收尾信号;单个 item 的 + // function_call_arguments.done 不算,它只说明该 item 的参数发完了。 "response.completed" => Ok(Some(ParsedStreamEvent { finish_reason: Some("completed".to_string()), + is_completion: true, tool_fragments: extract_responses_completed_tool_fragments(&parsed), ..Default::default() })), @@ -2814,17 +2869,27 @@ fn parse_anthropic_sse_event(data: &str) -> Result, Ll ..Default::default() })) } - "message_delta" => Ok(Some(ParsedStreamEvent { - finish_reason: parsed + "message_delta" => { + let stop_reason = parsed .get("delta") .and_then(|value| value.get("stop_reason")) .and_then(serde_json::Value::as_str) - .map(str::to_string), + .map(str::to_string); + Ok(Some(ParsedStreamEvent { + is_completion: stop_reason + .as_deref() + .is_some_and(|reason| !reason.trim().is_empty()), + finish_reason: stop_reason, + ..Default::default() + })) + } + // message_stop 只是流终止信号;真正的 stop_reason 已由 message_delta 提供, + // 这里不要伪造 finish_reason,否则会覆盖掉 end_turn 等真实值。但它确实是协议 + // 收尾信号,所以单独用 is_completion 记录:兼容网关可能只发它而漏 stop_reason。 + "message_stop" => Ok(Some(ParsedStreamEvent { + is_completion: true, ..Default::default() })), - // message_stop 只是流终止信号;真正的 stop_reason 已由 message_delta 提供, - // 这里不要伪造 finish_reason,否则会覆盖掉 end_turn 等真实值。 - "message_stop" => Ok(None), "error" => { let message = parsed .get("error") @@ -4531,8 +4596,9 @@ mod tests { } #[tokio::test] - async fn stream_run_rejects_truncated_tool_arguments() { - // 参数只拼到一半就断流,不能把半截 JSON 交给业务层。 + async fn stream_run_rejects_incomplete_tool_arguments_json() { + // 收尾信号齐全,但参数只拼到一半,半截 JSON 不能交给业务层。 + // 与下面几个"参数完整但没有收尾信号"的用例是两条独立防线。 let server_url = spawn_mock_server(vec![MockResponse { status_line: "200 OK", content_type: "text/event-stream; charset=utf-8", @@ -4554,6 +4620,169 @@ mod tests { assert!(matches!(error, LlmError::Deserialize(_))); } + #[tokio::test] + async fn stream_run_rejects_anthropic_tool_calls_without_completion_signal() { + // 参数字节完整,但没有 message_delta / message_stop:代理超时或网关掐断都长这样, + // 光凭"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":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"get_weather","input":{}}}"#, "\n\n", + r#"data: {"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"city\":\"杭州\"}"}}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let error = client + .stream_run(weather_tool_request(LlmApiKind::Anthropic), |_| {}) + .await + .expect_err("tool stream without completion signal should fail"); + + let LlmError::Deserialize(message) = error else { + panic!("应按截断失败,实际 {error:?}"); + }; + assert!(message.contains("协议完成信号前截断"), "{message}"); + } + + #[tokio::test] + async fn stream_run_accepts_anthropic_tool_calls_with_message_stop_only() { + // 兼容网关可能漏发 stop_reason 但仍发 message_stop;后者是合法收尾信号, + // 且不能借它伪造 finish_reason。 + 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_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"get_weather","input":{}}}"#, "\n\n", + r#"data: {"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"city\":\"杭州\"}"}}"#, "\n\n", + r#"data: {"type":"content_block_stop","index":1}"#, "\n\n", + r#"data: {"type":"message_stop"}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let response = client + .stream_run(weather_tool_request(LlmApiKind::Anthropic), |_| {}) + .await + .expect("message_stop should count as completion"); + + assert_eq!(response.tool_calls.len(), 1); + assert_eq!(response.tool_calls[0].name, "get_weather"); + assert_eq!(response.finish_reason, None); + } + + #[tokio::test] + async fn stream_run_rejects_chat_tool_calls_without_completion_signal() { + // Chat 分片拼出了完整 arguments,但既无 finish_reason 也无 [DONE]。 + let server_url = spawn_mock_server(vec![MockResponse { + status_line: "200 OK", + content_type: "text/event-stream; charset=utf-8", + body: concat!( + r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"get_weather","arguments":""}}]}}]}"#, "\n\n", + r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{\"city\":\"杭州\"}"}}]}}]}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let error = client + .stream_run(weather_tool_request(LlmApiKind::OpenAiChat), |_| {}) + .await + .expect_err("chat tool stream without completion signal should fail"); + + let LlmError::Deserialize(message) = error else { + panic!("应按截断失败,实际 {error:?}"); + }; + assert!(message.contains("协议完成信号前截断"), "{message}"); + } + + #[tokio::test] + async fn stream_run_rejects_responses_tool_calls_without_completion_signal() { + // Responses 的整体收尾只有 response.completed;单 item 的 + // function_call_arguments.done 不能顶替它。 + 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_item.added","output_index":0,"item":{"type":"function_call","call_id":"call_1","name":"get_weather"}}"#, "\n\n", + r#"data: {"type":"response.function_call_arguments.delta","output_index":0,"delta":"{\"city\":\"杭州\"}"}"#, "\n\n", + r#"data: {"type":"response.function_call_arguments.done","output_index":0,"arguments":"{\"city\":\"杭州\"}"}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let error = client + .stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |_| {}) + .await + .expect_err("responses tool stream without completion signal should fail"); + + let LlmError::Deserialize(message) = error else { + panic!("应按截断失败,实际 {error:?}"); + }; + assert!(message.contains("协议完成信号前截断"), "{message}"); + } + + #[tokio::test] + async fn stream_run_rejects_parallel_tool_calls_truncated_between_blocks() { + // 第一个工具块字节完整,流在第二个 content_block_start 到达前断掉。 + // 旧实现会返回"看起来完整"的单调用结果,静默丢掉模型本要发的第二个调用。 + 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_start","index":1,"content_block":{"type":"tool_use","id":"call_a","name":"get_weather","input":{}}}"#, "\n\n", + r#"data: {"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"city\":\"杭州\"}"}}"#, "\n\n", + r#"data: {"type":"content_block_stop","index":1}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let error = client + .stream_run(weather_tool_request(LlmApiKind::Anthropic), |_| {}) + .await + .expect_err("truncation between tool blocks should fail"); + + let LlmError::Deserialize(message) = error else { + panic!("应按截断失败,实际 {error:?}"); + }; + assert!(message.contains("协议完成信号前截断"), "{message}"); + } + + #[tokio::test] + async fn stream_run_keeps_text_only_response_without_completion_signal() { + // 作用域反向守卫:本轮只收严工具路径。纯文本流缺收尾信号仍按成功返回, + // 只打 warn。改这条断言前必须先确认所有在用网关的文本收尾行为。 + 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","index":0,"delta":{"type":"text_delta","text":"杭州今天"}}"#, "\n\n", + r#"data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"多云。"}}"#, "\n\n" + ) + .to_string(), + extra_headers: Vec::new(), + }]); + + let client = build_test_client(server_url, 0); + let response = client + .stream_run(weather_tool_request(LlmApiKind::Anthropic), |_| {}) + .await + .expect("text-only stream should still succeed"); + + assert_eq!(response.text, "杭州今天多云。"); + assert!(response.tool_calls.is_empty()); + assert_eq!(response.finish_reason, None); + } + #[tokio::test] async fn stream_run_falls_back_when_tool_use_yields_no_fragments() { // 上游说了本轮是工具调用,但事件形状不在已支持范围内,一个分片都没解出来。