流式工具槽位身份冲突改为失败关闭
push_tool_fragment 遇到已有槽位时无条件覆盖 id / 函数名,参数却是追加 (Chat / Anthropic)或整段覆盖(Responses 的 arguments_complete)。两者 不自洽:前一个调用参数为空时,拼接结果正好是后一个调用的合法 JSON, 参数完整性检查兜不住,调用方只拿到后一个工具,前一个静默消失并直接交给 Runtime 执行。 守规协议不会命中——Chat 的 index、Responses 的 output_index、Anthropic 的 content block index 在单条响应内都不重复。已知触发路径两条:兼容网关把 Chat 的 index 恒置 0;Responses 的 response.completed 回退按 output[] 下标 重建槽位时与流式 output_index 基准错位(如 completed 载荷省略 reasoning item),后者会产出「前一个调用的身份配后一个调用的参数」。 改为同槽位身份只允许从缺失变已知或重复同一个值,互不相同的非空值返回 Deserialize。id 与函数名都查:并行调用同一个工具时函数名相同,只有 id 能 区分。空白身份按缺失跳过,不算冲突——部分兼容网关在续传分片里回发完整 function 对象且 name / id 为空串,按「不等即冲突」会整批误杀,这也与归一层 的空白即缺失约定一致。身份冲突不走尾部错误保留路径,累加状态已不可信。 隔离验证:退回无条件覆盖后,四条拒绝用例全红,其中 Chat 那条返回 Ok(tool_calls=[call_b/get_air_quality]),call_a/get_weather 无声消失。 platform-llm 90 passed(原 84),新增 6 条:Chat 换工具 / Chat 同名换 id / Responses completed 与增量矛盾 / Anthropic 复用 index 四条拒绝,外加 网关每片回发同一身份、completed 重复确认同槽位两条兼容性守卫。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -271,9 +271,11 @@ arguments 是否必须是完整 JSON **按流式与非流式区分,两者的
|
||||
|
||||
工具事件的协议槽位缺失时必须失败关闭,不得跳过也不得按事件内位置猜测:槽位是并行分片唯一的归并依据。跳过会静默丢掉整个调用——只剩一个调用时才可能被 `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` 为空串,按“不等即冲突”会把它们整批误杀,这也与归一层的空白即缺失约定一致。
|
||||
|
||||
反过来,已经收尾的流遇到尾部传输 / 解析错误时必须保留结果,不能重跑 Provider。判断“有没有值得保留的东西”要看正文或工具调用任一非空,不能只看正文——纯工具调用响应的正文本来就是空的(MiniMax 的 Anthropic 工具流恒定如此),只看正文会让这类响应每次都被丢弃,白白多跑一轮往返。保留的安全性由“协议完成信号已到 + `finish_reason` 已到 + 工具参数完整 + 错误属可容忍尾部错误”共同保证,与正常路径判据一致。Chat 的非空 `finish_reason` 与 Anthropic 带 `stop_reason` 的 `message_delta` 会标记完成但不直接终止读取,因此在后续终止事件缺失时仍可能进入这条尾部保留路径;Chat 的 `[DONE]`、Anthropic 的 `message_stop` 以及 Responses 的 `response.completed` / `response.incomplete` 已经直接终止读取。
|
||||
|
||||
错误边界固定如下:`StreamUnavailable` 只表示流式响应已给出 `tool_use` / `tool_calls` 完成原因但没有聚合出任何工具 slot,供调用方回退非流式,它不承担截断语义;`EmptyResponse` 表示最终文本和工具调用都为空,纯工具响应合法;`Deserialize` 覆盖 JSON / SSE / UTF-8 解析失败、缺少 `choices[0]`、流式工具身份缺失、流式参数不完整,以及上述工具流未收尾截断。Anthropic 仍不支持 `web_search`、图片内容和纯 system 消息,必须至少有一条非 system 文本消息。
|
||||
错误边界固定如下:`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 不得写入日志、审计或持久化。
|
||||
|
||||
@@ -766,8 +766,39 @@ struct StreamAccumulation {
|
||||
completion_observed: bool,
|
||||
}
|
||||
|
||||
// 身份字段只允许「从缺失变已知」或「重复同一个值」。同槽位换身份说明上游把两个不同调用
|
||||
// 挤进了一个槽,此时覆盖身份却追加参数并不自洽:前一个调用参数为空时拼接结果仍是合法
|
||||
// JSON,截断断言兜不住,调用方只会看到后一个工具,前一个无声消失;Responses 的
|
||||
// arguments_complete 还会整段覆盖,产出「A 的身份配 B 的参数」。两种结果都会直接交给
|
||||
// Runtime 执行,所以必须失败关闭。
|
||||
//
|
||||
// 空白值按缺失处理,不算冲突:部分 OpenAI 兼容网关在续传分片里回发完整 function 对象,
|
||||
// name / id 是空串,按「不等即冲突」会把它们整批误杀。这也与 normalize_tool_calls 的
|
||||
// 空白即缺失约定一致。
|
||||
fn merge_tool_identity(
|
||||
current: &mut Option<String>,
|
||||
incoming: Option<String>,
|
||||
field: &str,
|
||||
slot: u64,
|
||||
) -> Result<(), LlmError> {
|
||||
let Some(incoming) = incoming.filter(|value| !value.trim().is_empty()) else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
match current {
|
||||
Some(existing) if existing.trim() == incoming.trim() => Ok(()),
|
||||
Some(existing) => Err(LlmError::Deserialize(format!(
|
||||
"LLM 流式工具分片槽位 {slot} 的 {field} 冲突:已有 {existing},又收到 {incoming}"
|
||||
))),
|
||||
None => {
|
||||
*current = Some(incoming);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl StreamAccumulation {
|
||||
fn push_tool_fragment(&mut self, fragment: ToolCallFragment) {
|
||||
fn push_tool_fragment(&mut self, fragment: ToolCallFragment) -> Result<(), LlmError> {
|
||||
if !self
|
||||
.tool_calls
|
||||
.iter()
|
||||
@@ -786,12 +817,9 @@ impl StreamAccumulation {
|
||||
.find(|pending| pending.slot == fragment.slot)
|
||||
.expect("slot was just ensured");
|
||||
|
||||
if let Some(id) = fragment.id {
|
||||
entry.id = Some(id);
|
||||
}
|
||||
if let Some(name) = fragment.name {
|
||||
entry.name = Some(name);
|
||||
}
|
||||
// 身份先校验:冲突时连参数都不能并进去,累加状态已经不可信。
|
||||
merge_tool_identity(&mut entry.id, fragment.id, "id", fragment.slot)?;
|
||||
merge_tool_identity(&mut entry.name, fragment.name, "函数名", fragment.slot)?;
|
||||
if let Some(delta) = fragment.arguments_delta {
|
||||
entry.arguments.push_str(delta.as_str());
|
||||
}
|
||||
@@ -799,6 +827,8 @@ impl StreamAccumulation {
|
||||
if let Some(complete) = fragment.arguments_complete {
|
||||
entry.arguments = complete;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// 流结束后固化,走与非流式相同的归一:缺 id / 函数名报错,空参数归一为 {},
|
||||
@@ -1884,8 +1914,9 @@ where
|
||||
Ok(events) => (events, None),
|
||||
Err(error) => (error.parsed_events, Some(error.error)),
|
||||
};
|
||||
// 槽位身份冲突比尾部错误更根本:累加出的工具调用已不可信,不能再走保留路径。
|
||||
let stream_terminated =
|
||||
consume_stream_events(events, accumulation, emit_finish_only_delta, on_delta);
|
||||
consume_stream_events(events, accumulation, emit_finish_only_delta, on_delta)?;
|
||||
|
||||
if stream_terminated {
|
||||
return Ok(true);
|
||||
@@ -1941,7 +1972,7 @@ fn consume_stream_events<F>(
|
||||
accumulation: &mut StreamAccumulation,
|
||||
emit_finish_only_delta: bool,
|
||||
on_delta: &mut F,
|
||||
) -> bool
|
||||
) -> Result<bool, LlmError>
|
||||
where
|
||||
F: FnMut(&LlmStreamDelta),
|
||||
{
|
||||
@@ -1965,7 +1996,7 @@ where
|
||||
|
||||
// 工具调用只累加,不进 on_delta:调用方的流式通道仍然只承载文本。
|
||||
for fragment in tool_fragments {
|
||||
accumulation.push_tool_fragment(fragment);
|
||||
accumulation.push_tool_fragment(fragment)?;
|
||||
}
|
||||
|
||||
let delta_text = delta_text.unwrap_or_default();
|
||||
@@ -1994,11 +2025,11 @@ where
|
||||
}
|
||||
|
||||
if is_terminal {
|
||||
return true;
|
||||
return Ok(true);
|
||||
}
|
||||
}
|
||||
|
||||
false
|
||||
Ok(false)
|
||||
}
|
||||
|
||||
fn normalize_non_empty(value: String, error_message: &str) -> Result<String, LlmError> {
|
||||
@@ -5410,6 +5441,151 @@ mod tests {
|
||||
.await;
|
||||
}
|
||||
|
||||
// 槽位身份冲突:槽位在但被两个不同调用共用。比缺槽位更隐蔽——身份被覆盖、参数却是
|
||||
// 追加/整段覆盖,两者不自洽,产出的调用会直接交给 Runtime 执行。
|
||||
async fn expect_stream_slot_identity_conflict_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_reusing_slot_for_another_call() {
|
||||
// 兼容网关把 index 恒置 0 时两个串行调用挤进同一槽位。这里刻意让前一个调用参数为空:
|
||||
// 拼接结果是后者的合法 JSON,截断断言兜不住,旧实现会静默丢掉 get_weather,
|
||||
// 只把 get_air_quality 交出去。
|
||||
expect_stream_slot_identity_conflict_error(
|
||||
LlmApiKind::OpenAiChat,
|
||||
concat!(
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_a","function":{"name":"get_weather","arguments":""}}]}}]}"#, "\n\n",
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"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_chat_tool_fragments_reusing_slot_for_same_tool_name() {
|
||||
// 并行调用同一个工具是最常见的并行场景:name 相同,只有 id 能区分。
|
||||
// 只查 name 的实现会把这两个调用合并成一个混合参数体。
|
||||
expect_stream_slot_identity_conflict_error(
|
||||
LlmApiKind::OpenAiChat,
|
||||
concat!(
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_a","function":{"name":"get_weather","arguments":""}}]}}]}"#, "\n\n",
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_b","function":{"name":"get_weather","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_completed_event_contradicting_streamed_slot() {
|
||||
// completed 回退按 output[] 下标重建槽位,前提是它与 output_index 同义。网关若在
|
||||
// completed 载荷里省掉 reasoning item,基准就错位:槽位 0 会拿到另一个工具的身份,
|
||||
// arguments_complete 再整段覆盖,产出「A 的身份配 B 的参数」且是合法 JSON。
|
||||
expect_stream_slot_identity_conflict_error(
|
||||
LlmApiKind::OpenAiResponses,
|
||||
concat!(
|
||||
r#"data: {"type":"response.output_item.added","item":{"id":"fc_0","type":"function_call","call_id":"call_a","name":"get_weather"},"output_index":0}"#, "\n\n",
|
||||
r#"data: {"type":"response.function_call_arguments.done","item_id":"fc_0","output_index":0,"arguments":"{\"city\":\"杭州\"}"}"#, "\n\n",
|
||||
r#"data: {"type":"response.completed","response":{"output":[{"id":"fc_1","type":"function_call","call_id":"call_b","name":"get_air_quality","arguments":"{\"city\":\"苏州\"}"}]}}"#, "\n\n"
|
||||
),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_rejects_anthropic_content_block_reusing_index() {
|
||||
expect_stream_slot_identity_conflict_error(
|
||||
LlmApiKind::Anthropic,
|
||||
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_start","index":1,"content_block":{"type":"tool_use","id":"call_b","name":"get_air_quality","input":{}}}"#, "\n\n",
|
||||
r#"data: {"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"city\":\"杭州\"}"}}"#, "\n\n",
|
||||
r#"data: {"type":"message_delta","delta":{"stop_reason":"tool_use"}}"#, "\n\n"
|
||||
),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_accepts_chat_gateway_repeating_tool_identity_every_chunk() {
|
||||
// 兼容性守卫:不少网关每个续传分片都回发完整 function 对象,身份要么是同一个值、
|
||||
// 要么是空串。前者不算冲突,后者按缺失跳过——否则这批网关会被整批误杀。
|
||||
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_7gOveph","function":{"name":"get_weather","arguments":"{\"city\""}}]}}]}"#, "\n\n",
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_7gOveph","function":{"name":"get_weather","arguments":":\"杭州"}}]}}]}"#, "\n\n",
|
||||
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"","function":{"name":"","arguments":"\"}"}}]}}]}"#, "\n\n",
|
||||
r#"data: {"choices":[{"finish_reason":"tool_calls"}]}"#, "\n\n",
|
||||
"data: [DONE]\n\n"
|
||||
)
|
||||
.to_string(),
|
||||
extra_headers: Vec::new(),
|
||||
}]);
|
||||
|
||||
let response = build_test_client(server_url, 0)
|
||||
.stream_run(weather_tool_request(LlmApiKind::OpenAiChat), |_| {})
|
||||
.await
|
||||
.expect("重复回发同一身份不算冲突");
|
||||
|
||||
assert_eq!(
|
||||
response.tool_calls,
|
||||
vec![LlmToolCall {
|
||||
id: "call_7gOveph".to_string(),
|
||||
name: "get_weather".to_string(),
|
||||
arguments: r#"{"city":"杭州"}"#.to_string(),
|
||||
}]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_keeps_responses_completed_event_reconfirming_streamed_slot() {
|
||||
// 作用域守卫:completed 回退与增量事件槽位一致时是正常路径,身份重复不能报错,
|
||||
// arguments_complete 仍要覆盖拼接结果。
|
||||
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","item":{"id":"fc_0","type":"function_call","call_id":"call_a","name":"get_weather"},"output_index":0}"#, "\n\n",
|
||||
r#"data: {"type":"response.function_call_arguments.delta","delta":"{\"city","item_id":"fc_0","output_index":0}"#, "\n\n",
|
||||
r#"data: {"type":"response.completed","response":{"output":[{"id":"fc_0","type":"function_call","call_id":"call_a","name":"get_weather","arguments":"{\"city\":\"杭州\"}"}]}}"#, "\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("同槽位重复确认身份不算冲突");
|
||||
|
||||
assert_eq!(
|
||||
response.tool_calls,
|
||||
vec![LlmToolCall {
|
||||
id: "call_a".to_string(),
|
||||
name: "get_weather".to_string(),
|
||||
arguments: r#"{"city":"杭州"}"#.to_string(),
|
||||
}]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_keeps_text_only_anthropic_events_without_block_index() {
|
||||
// 作用域守卫:只有工具事件收紧。text_delta 不依赖槽位,缺 index 不应受影响。
|
||||
|
||||
Reference in New Issue
Block a user