流式工具调用要求协议收尾信号
字节流干净结束不等于协议收尾:代理超时、网关掐断和 HTTP/2 提前 END_STREAM 都表现为干净 EOF。此前只要工具参数恰好是合法 JSON 就会 返回并执行该调用,且并行槽位在断流后会被静默丢弃。 新增 ParsedStreamEvent::is_completion 与 StreamAccumulation:: completion_observed,按协议判定收尾:Chat 用非空 finish_reason 或 [DONE],Responses 用 response.completed,Anthropic 用带 stop_reason 的 message_delta 或 message_stop。不能统一用 [DONE],MiniMax 兼容层 不发该标记。message_stop 改为产生只带 is_completion 的事件,保持不 伪造 finish_reason 覆盖真实 end_turn 的既有约束。 聚合出过工具 slot 但未观察到收尾信号时返回 Deserialize;判定使用未 固化槽位,使参数恰好闭合的截断仍按截断归因。StreamUnavailable 语义 不变,仍只表示事件形状不受支持。本轮只收严工具路径,纯文本缺收尾信 号仍返回成功并打 warn。 真实端点验证:MiniMax Anthropic 兼容层 stop_reason=tool_use、 MiniMax OpenAI 兼容层 finish_reason=tool_calls,均通过门禁。
This commit is contained in:
@@ -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 不得写入日志、审计或持久化。
|
||||
|
||||
@@ -610,6 +610,11 @@ struct ParsedStreamEvent {
|
||||
finish_reason: Option<String>,
|
||||
usage: Option<LlmTokenUsage>,
|
||||
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<ToolCallFragment>,
|
||||
}
|
||||
|
||||
@@ -640,6 +645,9 @@ struct StreamAccumulation {
|
||||
finish_reason: Option<String>,
|
||||
usage: Option<LlmTokenUsage>,
|
||||
tool_calls: Vec<PendingToolCall>,
|
||||
// 是否观察到过协议收尾信号。字节流干净结束不等于协议收尾:代理超时、网关自行掐断
|
||||
// 和 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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() {
|
||||
// 上游说了本轮是工具调用,但事件形状不在已支持范围内,一个分片都没解出来。
|
||||
|
||||
Reference in New Issue
Block a user