流式最终事件立即收口,工具槽位缺失失败关闭
两处收口不一致,一并处理。 一、Responses 的 response.completed 与 Anthropic 的 message_stop 此前只标记 完成、不终止读取循环。服务端在最终事件后保持连接时,stream_run 会继续等 EOF, 最坏一路等到调用方超时。既有覆盖只有 Chat 的 [DONE]。三者现在统一置终止位。 二、工具事件的协议槽位缺失时此前要么跳过要么猜测。Chat 退回事件内位置,两个各 含一个无 index 调用的事件都落到槽位 0,后者的 id/name 覆盖前者、arguments 被 拼成混合体;Responses 与 Anthropic 五处直接丢事件。跳过只在调用全丢时才会被 StreamUnavailable 断言兜住,丢一半毫无察觉,而 Responses 的 finish_reason 恒为 completed,那道断言对它永远不触发。槽位是并行分片唯一的归并依据,现在一律失败 关闭。收紧只覆盖工具事件,纯文本增量不依赖槽位。 尾错保留用例原本以 message_stop 结尾,加终止位后会在尾部错误之前收口而变成空转, 改为只发 message_delta,使它继续验证“收尾信号已到但流没干净结束”的保留路径。
This commit is contained in:
@@ -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 文本消息。
|
||||
|
||||
@@ -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<ToolCallFragment> {
|
||||
fn extract_chat_tool_fragments(
|
||||
choice: &ChatCompletionsChoice,
|
||||
) -> Result<Vec<ToolCallFragment>, 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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<u64> {
|
||||
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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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<Option<ParsedStreamEvent>, 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() {
|
||||
// 与上一条成对:流式的参数半截意味着流被截断,是传输层事实,必须报错。
|
||||
|
||||
Reference in New Issue
Block a user