Responses 终态快照同步修正流式回调
902f487ea 引入终态正文快照时,「累加非空」一支只覆盖 accumulation.text、不补
回调。这一支当初只顾着避免正文翻倍,漏了「增量拼接结果本来就可能不等于快照」
——而这恰恰是决定要覆盖的唯一理由,两件事是耦合的。结果是 LlmRunResponse.text
已经是完整值,调用方最后收到的累计正文却停在半截。
实际代价:按累计正文取值的抢救路径(单 Agent 流式在 stream_run 返回 Err 时用
最后一次回调的 accumulated_text 当回复)会拿到半截正文,而「incomplete 且带
工具调用」会被截断拒绝规则判为 Err,恰好走到那里——正是 incomplete 降级回复
该发挥作用的场景。持久化回复流与 UI 实时视图同样停在半截,但 ready 写入用的是
response.text 派生值且直接覆盖 accumulated_text,所以是窗口期问题,不会造成
finalization 冲突或 needs-reconciliation。
改为按「是否等于快照」判断,不一致就补发一次回调。增量字段分三种情形取值:
累加去空白后为空给整个快照(这一支不能并进 strip_prefix,纯空白累加值匹配不上
前缀会退化成空增量,反而让按增量累加的消费者丢内容);快照是增量的延长给后缀;
非前缀关系给空串靠 accumulated_text 纠正。最后一种下按 delta_text 累加的消费者
无法自愈,作为已知残留写进契约文档。
边界:只有「快照延长增量」与「快照与增量分叉」两种情形的回调行为改变,其余
四种(累加为空、完全一致、无快照、纯空白累加)维持字节级不变。不给 Responses
打开 emit_finish_only_delta——那会让每一条流都多一次终态回调,是另一个决定。
LlmStreamDelta 结构、App 侧消费点均未改动。
隔离验证两次:退掉 emit 条件里的 snapshot_corrected,分叉用例转红(前缀延长
那条靠 has_delta 仍绿,说明该标志只对分叉承重);整块回退成旧两支形态,两条
新增用例都转红、其余 101 条全绿,印证边界表里「不变的四行」确实没动。
platform-llm 103 passed(原 101)。新增前缀延长、非前缀分叉两条,并收紧既有的
不重复用例——改为断言回调次数与 (delta_text, accumulated_text) 二元组,原先只
断言累计正文序列,补不补这次回调都可能是绿的。
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -271,7 +271,9 @@ arguments 是否必须是完整 JSON **按流式与非流式区分,两者的
|
||||
|
||||
Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复源,两者必须对称。只恢复工具会造成三种后果:纯文本的 completed-only 响应退化成 `EmptyResponse`(上层白跑一轮重试或降级);`response.incomplete` 携带的截断正文本来是可用的降级结果,同样拿不回来;“正文 + 工具调用”的响应不报错,但模型的前置说明被静默丢掉,最隐蔽。正文提取必须复用非流式那条路径(`output_text` 优先、`output[].content[]` 回退、过滤 `reasoning` / `reasoning_content` / `analysis` / `thinking` 等隐藏 part),不得另写裸 JSON 提取器——漏掉过滤层会把思维链当正文吐给调用方。终态载荷反序列化失败时按“没有快照”静默降级、不报错:这是兜底恢复路径,网关发出未建模的形状时应当退回增量累加结果;这与槽位缺失必须失败关闭的口径不同,那里放过会造成静默的身份与参数错配,这里放过只是回到没有该恢复路径时的行为。
|
||||
|
||||
终态正文按**快照覆盖**而非追加合并,且要按累加状态分两条路:累加为空时(只发终态事件的网关)必须把快照当作一次增量发出去,只覆盖累加值会让调用方的流式通道全程收不到任何文本——Responses 的 finish-only 回调开关是关闭的,只有 Chat 打开,指望终态回调兜底并不成立;累加非空时按权威值覆盖但不补发回调,否则正文在调用方侧翻倍。覆盖语义与工具参数的 `arguments_complete` 一致。
|
||||
终态正文按**快照覆盖**而非追加合并,且要按累加状态分两条路:累加为空时(只发终态事件的网关)必须把快照当作一次增量发出去,只覆盖累加值会让调用方的流式通道全程收不到任何文本——Responses 的 finish-only 回调开关是关闭的,只有 Chat 打开,指望终态回调兜底并不成立;累加非空且与快照一致时不补发回调,否则正文在调用方侧翻倍。覆盖语义与工具参数的 `arguments_complete` 一致。
|
||||
|
||||
快照与增量拼接结果**不一致**时必须补发一次回调,只改累加值不够:调用方最后收到的累计正文会停在增量结果上,而 `LlmRunResponse.text` 已经是完整值,两者在同一次调用里分叉。按累计正文取值的消费者会因此拿到半截回复——流式请求返回 `Err` 时的抢救路径正是这样取值的,而“`incomplete` 且带工具调用”会被截断拒绝规则判为 `Err`,恰好走到那里。补发时的增量字段按三种情形取值:累加去空白后为空时给整个快照(这一支不能并进前缀相减,纯空白累加值匹配不上前缀会退化成空增量,反而让按增量累加的消费者丢内容);快照是增量的延长时给后缀;两者非前缀关系时无法表达成增量,只能给空增量、靠累计正文纠正。最后一种情形下按增量累加的消费者无法自愈,是**已知残留**,只能等调用方在最终回复落地时整体覆盖。该规则不改变“终态事件是否总是触发回调”——只有快照确实纠正了内容才补发,给 Responses 打开 finish-only 回调是另一个决定。
|
||||
|
||||
工具事件的协议槽位缺失时必须失败关闭,不得跳过也不得按事件内位置猜测:槽位是并行分片唯一的归并依据。跳过会静默丢掉整个调用——只剩一个调用时才可能被 `StreamUnavailable` 断言兜住,丢一半毫无察觉;Responses 的整体终态原因是 `completed` / `incomplete`,也不会触发只识别 `tool_use` / `tool_calls` 的那道断言。猜测则会把两个不同调用合并成一个混合体(后者的 id / name 覆盖前者,arguments 被拼接)。判定字段为 Chat 的 `delta.tool_calls[].index`、Responses 的 `output_index`、Anthropic 的 content block `index`。该约束只覆盖工具事件,纯文本增量不依赖槽位,不受影响。
|
||||
|
||||
|
||||
@@ -2024,25 +2024,41 @@ where
|
||||
accumulation.text.push_str(delta_text.as_str());
|
||||
}
|
||||
|
||||
// 终态快照是上游给出的权威完整正文,按累加状态分两条路:
|
||||
// - 累加为空(只发终态事件的网关):当成一次增量发出去。只覆盖 accumulation.text
|
||||
// 的话 response.text 是对了,但调用方的流式通道全程收不到任何文本——Responses
|
||||
// 的 emit_finish_only_delta 是 false,指望终态回调兜底并不成立。
|
||||
// - 累加非空:按权威值覆盖拼接结果,但不补发回调,否则正文在调用方侧翻倍。
|
||||
// 覆盖语义与工具参数的 arguments_complete 一致。
|
||||
// 终态快照是上游给出的权威完整正文,覆盖语义与工具参数的 arguments_complete 一致。
|
||||
// 但只改累加值不够:调用方最后收到的累计正文会停在增量拼接结果上,而
|
||||
// LlmRunResponse.text 已经是完整值,两者在同一次调用里分叉。按 accumulated_text
|
||||
// 取值的消费者(单 Agent 流式在 stream_run 返回 Err 时的抢救路径)会因此拿到半截
|
||||
// 回复——「incomplete + 有工具调用」正好会走到那里。所以不一致时必须补一次回调。
|
||||
//
|
||||
// delta_text 尽量给成「新增的那一截」,让按 delta_text 累加的消费者也能自愈:
|
||||
// - 累加去空白后为空:整个快照就是增量。这一支不能并进 strip_prefix——纯空白累加
|
||||
// 值匹配不上前缀会退化成空 delta,反而让那类消费者丢内容。
|
||||
// - 快照是增量的延长(绝大多数情况):给后缀。
|
||||
// - 两者非前缀关系(增量与终态载荷不同源):无法表达成增量,只能给空串靠
|
||||
// accumulated_text 纠正,按 delta_text 累加的那份副本修不了,是已知残留。
|
||||
//
|
||||
// 相等时不补回调,避免正文在调用方侧翻倍;这也意味着本改动不给 Responses 打开
|
||||
// emit_finish_only_delta——那会让每一条流都多一次终态回调,是另一个决定。
|
||||
let mut snapshot_corrected = false;
|
||||
if let Some(snapshot) = text_snapshot.filter(|text| !text.trim().is_empty()) {
|
||||
if accumulation.text.trim().is_empty() {
|
||||
accumulation.text = snapshot.clone();
|
||||
delta_text = snapshot;
|
||||
has_delta = true;
|
||||
} else {
|
||||
if snapshot != accumulation.text {
|
||||
delta_text = if accumulation.text.trim().is_empty() {
|
||||
snapshot.clone()
|
||||
} else {
|
||||
snapshot
|
||||
.strip_prefix(accumulation.text.as_str())
|
||||
.unwrap_or_default()
|
||||
.to_string()
|
||||
};
|
||||
accumulation.text = snapshot;
|
||||
has_delta = !delta_text.is_empty();
|
||||
snapshot_corrected = true;
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(event_finish_reason) = event_finish_reason {
|
||||
accumulation.finish_reason = Some(event_finish_reason.clone());
|
||||
if has_delta || emit_finish_only_delta {
|
||||
if has_delta || emit_finish_only_delta || snapshot_corrected {
|
||||
let update = LlmStreamDelta {
|
||||
accumulated_text: accumulation.text.clone(),
|
||||
delta_text,
|
||||
@@ -5020,19 +5036,92 @@ mod tests {
|
||||
extra_headers: Vec::new(),
|
||||
}]);
|
||||
|
||||
let mut streamed: Vec<String> = Vec::new();
|
||||
let mut streamed: Vec<(String, String)> = Vec::new();
|
||||
let response = build_test_client(server_url, 0)
|
||||
.stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |delta| {
|
||||
streamed.push(delta.accumulated_text.clone());
|
||||
streamed.push((delta.delta_text.clone(), delta.accumulated_text.clone()));
|
||||
})
|
||||
.await
|
||||
.expect("增量与终态并存时不应重复正文");
|
||||
|
||||
assert_eq!(response.text, "杭州今天多云。");
|
||||
// 已有增量时不补发回调,调用方侧同样不能翻倍。
|
||||
// 快照与增量拼接结果一致时不补发回调,调用方侧同样不能翻倍。这里必须连回调次数
|
||||
// 一起断言:只断言内容序列的话,补不补这一次都可能是绿的。
|
||||
assert_eq!(streamed.len(), 2);
|
||||
assert_eq!(
|
||||
streamed,
|
||||
vec!["杭州今天".to_string(), "杭州今天多云。".to_string()]
|
||||
vec![
|
||||
("杭州今天".to_string(), "杭州今天".to_string()),
|
||||
("多云。".to_string(), "杭州今天多云。".to_string()),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_corrects_streamed_text_when_terminal_snapshot_extends_deltas() {
|
||||
// 增量只流出半截、终态载荷才是完整正文。只改累加值的话 LlmRunResponse.text 对了,
|
||||
// 但调用方最后收到的累计正文停在半截,两者在同一次调用里分叉。
|
||||
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_text.delta","delta":"杭州今天"}"#, "\n\n",
|
||||
r#"data: {"type":"response.completed","response":{"output":[{"id":"msg_0","type":"message","content":[{"type":"output_text","text":"杭州今天多云。"}]}]}}"#, "\n\n"
|
||||
)
|
||||
.to_string(),
|
||||
extra_headers: Vec::new(),
|
||||
}]);
|
||||
|
||||
let mut streamed: Vec<(String, String)> = Vec::new();
|
||||
let response = build_test_client(server_url, 0)
|
||||
.stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |delta| {
|
||||
streamed.push((delta.delta_text.clone(), delta.accumulated_text.clone()));
|
||||
})
|
||||
.await
|
||||
.expect("终态快照延长增量时应补发回调");
|
||||
|
||||
assert_eq!(response.text, "杭州今天多云。");
|
||||
// 快照是增量的延长时给出后缀,按 delta_text 累加的消费者也能自愈。
|
||||
assert_eq!(
|
||||
streamed,
|
||||
vec![
|
||||
("杭州今天".to_string(), "杭州今天".to_string()),
|
||||
("多云。".to_string(), "杭州今天多云。".to_string()),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stream_run_corrects_streamed_text_when_terminal_snapshot_diverges_from_deltas() {
|
||||
// 增量与终态载荷不同源时无法表达成增量:delta_text 只能给空串,靠 accumulated_text
|
||||
// 纠正。按 accumulated_text 取值的消费者自愈,按 delta_text 累加的那份修不了,
|
||||
// 是契约文档里记着的已知残留。
|
||||
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_text.delta","delta":"临安"}"#, "\n\n",
|
||||
r#"data: {"type":"response.completed","response":{"output":[{"id":"msg_0","type":"message","content":[{"type":"output_text","text":"杭州今天多云。"}]}]}}"#, "\n\n"
|
||||
)
|
||||
.to_string(),
|
||||
extra_headers: Vec::new(),
|
||||
}]);
|
||||
|
||||
let mut streamed: Vec<(String, String)> = Vec::new();
|
||||
let response = build_test_client(server_url, 0)
|
||||
.stream_run(weather_tool_request(LlmApiKind::OpenAiResponses), |delta| {
|
||||
streamed.push((delta.delta_text.clone(), delta.accumulated_text.clone()));
|
||||
})
|
||||
.await
|
||||
.expect("终态快照与增量分叉时应补发回调");
|
||||
|
||||
assert_eq!(response.text, "杭州今天多云。");
|
||||
assert_eq!(
|
||||
streamed,
|
||||
vec![
|
||||
("临安".to_string(), "临安".to_string()),
|
||||
(String::new(), "杭州今天多云。".to_string()),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user