From 987b4761ed69e3395ec26a8568007ba49b12bec1 Mon Sep 17 00:00:00 2001 From: lzhs1995 Date: Sun, 26 Jul 2026 03:26:14 +0530 Subject: [PATCH 1/2] fix(stream): report upstream interruption instead of faking success MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 上游断流时,流式与缓冲两条路径都走正常收尾(generate_final_events / finish_and_get_all_events),客户端因此收到 `message_delta` (stop_reason=end_turn) + `message_stop`——一次失败被伪装成一次完整成功的 回合。Claude Code CLI 据此认为本轮已完成:不重试、继续往下走,于是半截的 工具调用参数被丢弃、多步任务在中途"正常"结束。这是 CLI 侧观察到的 「工具调用不严格执行」的一个真因。 改为在断流时只关闭已打开的内容块,并下发 Anthropic `overloaded_error` 事件,客户端 SDK / CLI 对该类型走重试路径。断流收尾时不做任何残留 flush (thinking / invoke 嗅探 / XML 过滤器 / 工具 JSON 累积器),那些残留本身 就是不完整数据,flush 出去只会把半截内容伪装成有效内容。 同时置位 message_delta_sent / message_ended 幂等位,使断流之后任何路径 再调 generate_final_events 都无法补出正常收尾(防回归)。 trace 侧行为不变:仍记 interrupted + STREAM_INTERRUPTED,用量仍记 error。 新增 6 个测试锁定契约:流式 / 缓冲各自「必须报错且无正常收尾」与「事后 不可回退成伪造成功」、半截工具 JSON 不得作为 tool_use 发出,以及反向对照 「正常收尾不受影响,仍发 message_delta + message_stop」。 CI:把 stream.rs 纳入候选分支的 rustfmt 闸门,并把本分支加入触发列表。 --- .github/workflows/kiro-document-candidate.yml | 2 + src/anthropic/handlers.rs | 19 +- src/anthropic/stream.rs | 267 ++++++++++++++++++ 3 files changed, 283 insertions(+), 5 deletions(-) diff --git a/.github/workflows/kiro-document-candidate.yml b/.github/workflows/kiro-document-candidate.yml index 454fbb03..72e443f1 100644 --- a/.github/workflows/kiro-document-candidate.yml +++ b/.github/workflows/kiro-document-candidate.yml @@ -5,6 +5,7 @@ on: push: branches: - feat/native-document-contract-20260725 + - feat/stream-interrupt-no-fake-success-20260726 concurrency: group: kiro-document-candidate-${{ github.ref }} @@ -53,6 +54,7 @@ jobs: rustfmt --edition 2024 --check \ src/anthropic/converter.rs \ src/anthropic/handlers.rs \ + src/anthropic/stream.rs \ src/anthropic/types.rs \ src/anthropic/websearch_loop.rs \ src/kiro/model/requests/conversation.rs diff --git a/src/anthropic/handlers.rs b/src/anthropic/handlers.rs index a8de6dff..c3cbeaf1 100644 --- a/src/anthropic/handlers.rs +++ b/src/anthropic/handlers.rs @@ -33,7 +33,9 @@ use uuid::Uuid; use super::converter::{ConversionError, convert_request_with_mode}; use super::middleware::{AppState, KeyContext}; -use super::stream::{BufferedStreamContext, SseEvent, StreamContext}; +use super::stream::{ + BufferedStreamContext, STREAM_INTERRUPTED_CLIENT_MESSAGE, SseEvent, StreamContext, +}; use super::types::{ CountTokensRequest, CountTokensResponse, ErrorResponse, MessagesRequest, Model, ModelsResponse, OutputConfig, Thinking, @@ -1232,8 +1234,11 @@ fn create_sse_stream( } Some(Err(e)) => { tracing::error!("读取响应流失败: {}", e); - // 发送最终事件并结束(记为 error) - let final_events = ctx.generate_final_events(); + // 上游断流 ≠ 正常收尾:只关闭未闭合的块并下发 error 事件, + // 绝不补发 message_delta(stop_reason=end_turn) + message_stop, + // 否则客户端会把半截响应当成一次成功完成的回合而不重试。 + let final_events = + ctx.generate_interrupted_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); record_stream_usage(&hook, &ctx, credential_id, "error"); // 已开始返回内容后上游断流:标记为 interrupted,带已发送字节数 tracer.finalize( @@ -2098,8 +2103,12 @@ fn create_buffered_sse_stream( } Some(Err(e)) => { tracing::error!("读取响应流失败: {}", e); - // 发生错误,完成处理并返回所有事件 - let all_events = ctx.finish_and_get_all_events(); + // 上游断流:把已缓冲的事件连同 error 事件一起下发, + // 但不补发 message_delta / message_stop(详见 + // BufferedStreamContext::interrupt_and_get_all_events)。 + let all_events = ctx.interrupt_and_get_all_events( + STREAM_INTERRUPTED_CLIENT_MESSAGE, + ); let (i, o, cc, cr, credits) = ctx.final_usage(); hook.record(credential_id, i, o, cc, cr, credits, "error"); // 缓冲模式 chunk 读取失败:上游中途断流 diff --git a/src/anthropic/stream.rs b/src/anthropic/stream.rs index 39e2caa4..0ce2f080 100644 --- a/src/anthropic/stream.rs +++ b/src/anthropic/stream.rs @@ -22,6 +22,14 @@ use crate::kiro::model::events::Event; /// signature,因此该占位字符串只在客户端 ↔ kiro.rs 之间存在,不会影响转发。 pub(super) const THINKING_SIGNATURE_PLACEHOLDER: &str = "kiro-rs-thinking-signature"; +/// 上游断流时下发给客户端的错误文案 +/// +/// 用 `overloaded_error`(Anthropic 官方错误类型之一)而不是伪造正常收尾:客户端 +/// SDK / Claude Code CLI 对该类型走重试路径,而 `stop_reason: end_turn` 会被当成 +/// 本轮已成功完成,半截的工具调用与多步任务就此静默中断。 +pub(super) const STREAM_INTERRUPTED_CLIENT_MESSAGE: &str = + "Upstream connection was interrupted before the response finished. Please retry."; + const TOOL_USE_XML_PREFIX: &str = " Vec { + let mut events = Vec::new(); + + for (index, block) in self.active_blocks.iter_mut() { + if block.started && !block.stopped { + events.push(SseEvent::new( + "content_block_stop", + json!({ + "type": "content_block_stop", + "index": index + }), + )); + block.stopped = true; + } + } + + // 占位:后续任何 generate_final_events 调用都不得再补发正常收尾。 + self.message_delta_sent = true; + self.message_ended = true; + + events + } + /// 生成最终事件序列 pub fn generate_final_events( &mut self, @@ -2347,6 +2384,33 @@ impl StreamContext { events } + /// 上游断流收尾:关闭未闭合的块,并补发 Anthropic `error` 事件,**不发** + /// `message_delta` / `message_stop`。 + /// + /// 背景(Claude Code CLI 工具调用「不严格执行」的一个真因):原实现在上游 + /// 断流时也走 `generate_final_events()`,客户端收到 `stop_reason: end_turn` + /// + `message_stop`,把截断的半截响应当成一次**成功完成**的回合——工具调用 + /// 参数写到一半就没了、多步任务在中途"正常"结束。伪造成功比明确报错更糟: + /// 客户端不会重试。改为下发 `overloaded_error`,让客户端走重试路径。 + /// + /// 注意:此处**不做**任何 flush(thinking / invoke 嗅探 / XML 过滤器残留、 + /// 工具 JSON 累积器收尾),因为那些残留本身就是不完整数据,flush 出去只会把 + /// 半截内容伪装成有效内容。 + pub fn generate_interrupted_events(&mut self, message: &str) -> Vec { + let mut events = self.state_manager.generate_interrupted_events(); + events.push(SseEvent::new( + "error", + json!({ + "type": "error", + "error": { + "type": "overloaded_error", + "message": message + } + }), + )); + events + } + /// 生成最终事件序列 pub fn generate_final_events(&mut self) -> Vec { let mut events = Vec::new(); @@ -2589,6 +2653,41 @@ impl BufferedStreamContext { std::mem::take(&mut self.event_buffer) } + /// 上游断流收尾(缓冲模式):下发已缓冲的事件 + `error`,**不补** + /// `message_delta` / `message_stop`。 + /// + /// 与 [`Self::finish_and_get_all_events`] 的区别:不做任何残留 flush、不做 + /// `message_start` 的 usage 更正之外的收尾动作。缓冲模式此时一个字节都还没发给 + /// 客户端,理论上可以整体丢弃,但保留已缓冲内容对排查更有用;关键是末尾必须是 + /// `error` 而不是成功收尾,否则客户端把截断响应当成完整回合(详见 + /// `StreamContext::generate_interrupted_events`)。 + pub fn interrupt_and_get_all_events(&mut self, message: &str) -> Vec { + if !self.initial_events_generated { + let initial_events = self.inner.generate_initial_events(); + self.event_buffer.extend(initial_events); + self.initial_events_generated = true; + } + + let (final_input_tokens, cache_creation, cache_read) = self.inner.resolved_usage(); + + let interrupted_events = self.inner.generate_interrupted_events(message); + self.event_buffer.extend(interrupted_events); + + // 与成功路径一致:把 message_start 里的 usage 更正为真实口径。 + for event in &mut self.event_buffer { + if event.event == "message_start" + && let Some(message) = event.data.get_mut("message") + && let Some(usage) = message.get_mut("usage") + { + usage["input_tokens"] = serde_json::json!(final_input_tokens); + usage["cache_creation_input_tokens"] = serde_json::json!(cache_creation); + usage["cache_read_input_tokens"] = serde_json::json!(cache_read); + } + } + + std::mem::take(&mut self.event_buffer) + } + /// 取出最终用量(在 finish_and_get_all_events 之后调用) /// /// 返回顺序:(input_tokens, output_tokens, cache_creation_tokens, cache_read_tokens, credits) @@ -5099,4 +5198,172 @@ mod tests { && e.data["content_block"]["data"] == "encrypted-thinking" })); } + + // ---- 上游断流:必须报错,绝不伪造成功收尾 ---- + // + // 背景:断流时若走正常收尾(message_delta stop_reason=end_turn + message_stop), + // Claude Code CLI 会把截断的半截响应当成一次完整成功的回合,于是不重试、 + // 继续往下走(工具调用参数写到一半就没了、多步任务中途"正常"结束)。 + + fn has_normal_completion(events: &[SseEvent]) -> bool { + events + .iter() + .any(|e| e.event == "message_delta" || e.event == "message_stop") + } + + fn overloaded_errors(events: &[SseEvent]) -> Vec<&SseEvent> { + events + .iter() + .filter(|e| e.event == "error" && e.data["error"]["type"] == "overloaded_error") + .collect() + } + + #[test] + fn interrupted_stream_emits_error_and_no_normal_completion() { + let mut ctx = StreamContext::new_with_thinking( + "test-model", + 10, + false, + HashMap::new(), + test_known_tools(), + ); + let mut all = ctx.generate_initial_events(); + all.extend(ctx.process_assistant_response("half an ans")); + + let interrupted = ctx.generate_interrupted_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); + all.extend(interrupted.clone()); + + assert!( + !has_normal_completion(&all), + "断流不得下发 message_delta / message_stop,否则客户端把截断当成功: {:?}", + all.iter().map(|e| &e.event).collect::>() + ); + assert_eq!( + overloaded_errors(&interrupted).len(), + 1, + "断流必须下发恰好一个 overloaded_error 事件" + ); + assert_eq!( + interrupted.last().map(|e| e.event.as_str()), + Some("error"), + "error 必须是最后一个事件" + ); + // 已打开的 text 块要闭合,客户端才能干净地丢弃这一轮 + assert!( + interrupted.iter().any(|e| e.event == "content_block_stop"), + "断流应关闭已打开的内容块" + ); + } + + #[test] + fn interrupted_stream_cannot_regress_to_fake_success() { + // 防回归:断流收尾之后,任何路径再调 generate_final_events 都不能补出 + // 正常收尾事件(message_delta / message_stop 的幂等位必须已置位)。 + let mut ctx = StreamContext::new_with_thinking( + "test-model", + 10, + false, + HashMap::new(), + test_known_tools(), + ); + let _ = ctx.generate_initial_events(); + let _ = ctx.process_assistant_response("partial"); + let _ = ctx.generate_interrupted_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); + + let after = ctx.generate_final_events(); + assert!( + !has_normal_completion(&after), + "断流后不得再补发正常收尾: {:?}", + after.iter().map(|e| &e.event).collect::>() + ); + } + + #[test] + fn interrupted_stream_does_not_flush_truncated_tool_json() { + // 上游在工具参数写到一半时断流:不得把半截 JSON 当成完整工具调用发出。 + let mut ctx = StreamContext::new_with_thinking( + "test-model", + 10, + false, + HashMap::new(), + test_known_tools(), + ); + let mut all = ctx.generate_initial_events(); + all.extend(ctx.process_tool_use(&tool_evt("t1", "Write", "{\"file_pa", false))); + + let interrupted = ctx.generate_interrupted_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); + all.extend(interrupted); + + assert!( + !all.iter().any(|e| { + e.event == "content_block_start" && e.data["content_block"]["type"] == "tool_use" + }), + "半截工具调用不得作为 tool_use 块发出" + ); + assert!(!has_normal_completion(&all)); + } + + #[test] + fn buffered_interrupted_stream_emits_error_and_no_normal_completion() { + let mut ctx = + BufferedStreamContext::new("test-model", 10, false, HashMap::new(), test_known_tools()); + ctx.process_and_buffer(&Event::AssistantResponse( + serde_json::from_value(serde_json::json!({ "content": "half an ans" })) + .expect("assistantResponseEvent fixture"), + )); + + let all = ctx.interrupt_and_get_all_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); + + assert!( + all.iter().any(|e| e.event == "message_start"), + "缓冲模式应保留已缓冲事件便于排查" + ); + assert!( + !has_normal_completion(&all), + "缓冲模式断流同样不得下发 message_delta / message_stop: {:?}", + all.iter().map(|e| &e.event).collect::>() + ); + assert_eq!(overloaded_errors(&all).len(), 1); + assert_eq!(all.last().map(|e| e.event.as_str()), Some("error")); + } + + #[test] + fn buffered_interrupted_stream_cannot_regress_to_fake_success() { + let mut ctx = + BufferedStreamContext::new("test-model", 10, false, HashMap::new(), test_known_tools()); + ctx.process_and_buffer(&Event::AssistantResponse( + serde_json::from_value(serde_json::json!({ "content": "partial" })) + .expect("assistantResponseEvent fixture"), + )); + let _ = ctx.interrupt_and_get_all_events(STREAM_INTERRUPTED_CLIENT_MESSAGE); + + let after = ctx.finish_and_get_all_events(); + assert!( + !has_normal_completion(&after), + "缓冲模式断流后不得再补发正常收尾: {:?}", + after.iter().map(|e| &e.event).collect::>() + ); + } + + #[test] + fn normal_completion_still_emits_message_stop() { + // 反向对照:正常收尾路径不受影响,仍必须发 message_delta + message_stop。 + let mut ctx = StreamContext::new_with_thinking( + "test-model", + 10, + false, + HashMap::new(), + test_known_tools(), + ); + let mut all = ctx.generate_initial_events(); + all.extend(ctx.process_assistant_response("a complete answer")); + all.extend(ctx.generate_final_events()); + + assert!( + all.iter() + .any(|e| e.event == "message_delta" && e.data["delta"]["stop_reason"] == "end_turn") + ); + assert!(all.iter().any(|e| e.event == "message_stop")); + assert!(overloaded_errors(&all).is_empty()); + } } From 762df8bd96c5d7d38d5f6d1d07caa63a1ae4bf5d Mon Sep 17 00:00:00 2001 From: LZHS Date: Sun, 26 Jul 2026 03:56:45 +0530 Subject: [PATCH 2/2] fix(converter): derive a stable agentContinuationId per client session MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit conversationId 已经从 metadata.user_id 的 session UUID 派生,但 agentContinuationId 仍是每请求 Uuid::new_v4()。上游因此把同一个客户端会话 的每一轮当成一条全新的 agent 任务线,多步任务(工具调用链、长任务续写)在 上游侧没有连续性。 改为用同一个会话锚点派生 continuation(domain 前缀 + SHA-256,保持与 conversationId 一一对应但互不相等)。拿不到锚点时仍退回随机值,保持旧行为 ——宁可丢连续性,不可乱绑任务线。 Tests: 585 passed (新增 5 个派生级 + 3 个端到端) --- src/anthropic/converter.rs | 166 ++++++++++++++++++++++++++++++++++++- 1 file changed, 163 insertions(+), 3 deletions(-) diff --git a/src/anthropic/converter.rs b/src/anthropic/converter.rs index ae660df5..81897278 100644 --- a/src/anthropic/converter.rs +++ b/src/anthropic/converter.rs @@ -637,6 +637,32 @@ impl std::fmt::Display for ConversionError { impl std::error::Error for ConversionError {} +/// 由会话锚点派生稳定的 `agentContinuationId` +/// +/// 原实现每个请求都 `Uuid::new_v4()`,于是同一个客户端会话的每一轮在上游看来都是 +/// 一条全新的 agent 任务线:多步任务(工具调用链、长任务续写)失去连续性。 +/// `conversationId` 已经从 `metadata.user_id` 的 session UUID 派生,这里用同一个 +/// 锚点派生 continuation,让「同一会话 = 同一条任务线」在上游成立。 +/// +/// 取 SHA-256 而不直接复用 conversationId:两个字段在上游是不同维度,直接相等会 +/// 让上游把「会话」和「任务线」当成同一个键;派生值保证一一对应但互不相等。 +/// domain 前缀防止与其它用途的摘要撞用途。 +fn derive_agent_continuation_id(conversation_anchor: &str) -> String { + let mut hasher = Sha256::new(); + hasher.update(b"kiro-rs/agent-continuation/v1\0"); + hasher.update(conversation_anchor.as_bytes()); + let hex = format!("{:x}", hasher.finalize()); + // 沿用 UUID 形状:上游只要求稳定字符串,UUID 形状便于日志辨识。 + format!( + "{}-{}-{}-{}-{}", + &hex[..8], + &hex[8..12], + &hex[12..16], + &hex[16..20], + &hex[20..32] + ) +} + /// 从 metadata.user_id 中提取 session UUID /// /// 支持两种格式: @@ -747,13 +773,20 @@ pub fn convert_request_with_mode( // 3. 生成会话 ID 和代理 ID // 优先从 metadata.user_id 中提取 session UUID 作为 conversationId - let conversation_id = req + let session_anchor = req .metadata .as_ref() .and_then(|m| m.user_id.as_ref()) - .and_then(|user_id| extract_session_id(user_id)) + .and_then(|user_id| extract_session_id(user_id)); + let conversation_id = session_anchor + .clone() .unwrap_or_else(|| Uuid::new_v4().to_string()); - let agent_continuation_id = Uuid::new_v4().to_string(); + // 有会话锚点时派生稳定的 continuation(同一会话 = 同一条 agent 任务线); + // 拿不到锚点则退回随机值,保持旧行为——宁可丢连续性,不可乱绑任务线。 + let agent_continuation_id = match session_anchor.as_deref() { + Some(anchor) => derive_agent_continuation_id(anchor), + None => Uuid::new_v4().to_string(), + }; // 4. 确定触发类型 let chat_trigger_type = determine_chat_trigger_type(req); @@ -4018,4 +4051,131 @@ mod tests { matches!(err, ConversionError::UnsupportedDocument(message) if message.contains("URL source")) ); } + + // ---- 会话连续性:agentContinuationId 必须由会话锚点确定性派生 ---- + // + // 原实现每请求一个随机 UUID,上游因此把同一客户端会话的每一轮当成新的 agent + // 任务线,多步任务(工具调用链)失去连续性。 + + #[test] + fn agent_continuation_id_is_stable_for_the_same_anchor() { + let anchor = "0b4445e1-f5be-49e1-87ce-62bbc28ad705"; + assert_eq!( + derive_agent_continuation_id(anchor), + derive_agent_continuation_id(anchor), + "同一会话锚点必须派生出同一个 continuation" + ); + } + + #[test] + fn agent_continuation_id_differs_across_anchors() { + let a = derive_agent_continuation_id("0b4445e1-f5be-49e1-87ce-62bbc28ad705"); + let b = derive_agent_continuation_id("0b4445e1-f5be-49e1-87ce-62bbc28ad706"); + assert_ne!(a, b, "不同会话不得共用 agent 任务线"); + } + + #[test] + fn agent_continuation_id_is_not_equal_to_the_anchor() { + // conversationId 直接用锚点;continuation 必须是另一个值,否则上游会把 + // 「会话」与「任务线」当成同一个键。 + let anchor = "0b4445e1-f5be-49e1-87ce-62bbc28ad705"; + assert_ne!(derive_agent_continuation_id(anchor), anchor); + } + + #[test] + fn agent_continuation_id_keeps_uuid_shape() { + let id = derive_agent_continuation_id("0b4445e1-f5be-49e1-87ce-62bbc28ad705"); + assert_eq!(id.len(), 36, "{}", id); + assert_eq!(id.chars().filter(|c| *c == '-').count(), 4, "{}", id); + assert!( + id.chars().all(|c| c.is_ascii_hexdigit() || c == '-'), + "{}", + id + ); + } + + #[test] + fn distinct_anchors_do_not_collide_in_bulk() { + let ids: std::collections::HashSet = (0..2000) + .map(|i| derive_agent_continuation_id(&format!("session-{i}"))) + .collect(); + assert_eq!(ids.len(), 2000, "派生值发生碰撞"); + } + + /// 构造一个带 metadata.user_id 的最小请求,用于端到端验证转换结果。 + fn request_with_user_id(user_id: Option<&str>) -> MessagesRequest { + MessagesRequest { + force_web_search_loop: false, + model: "claude-sonnet-4-5-20250929".to_string(), + max_tokens: 1024, + messages: vec![super::super::types::Message { + role: "user".to_string(), + content: serde_json::json!("hi"), + }], + stream: false, + system: None, + tools: None, + tool_choice: None, + thinking: None, + output_config: None, + metadata: Some(super::super::types::Metadata { + user_id: user_id.map(|s| s.to_string()), + }), + } + } + + #[test] + fn same_client_session_reuses_conversation_and_continuation() { + let user_id = "user_abc_account__session_0b4445e1-f5be-49e1-87ce-62bbc28ad705"; + let first = convert_request(&request_with_user_id(Some(user_id))).unwrap(); + let second = convert_request(&request_with_user_id(Some(user_id))).unwrap(); + + let first_state = &first.conversation_state; + let second_state = &second.conversation_state; + + assert_eq!( + first_state.conversation_id, second_state.conversation_id, + "同一会话的两轮必须复用 conversationId" + ); + assert_eq!( + first_state.agent_continuation_id, second_state.agent_continuation_id, + "同一会话的两轮必须复用 agentContinuationId" + ); + } + + #[test] + fn distinct_client_sessions_do_not_share_continuation() { + let a = convert_request(&request_with_user_id(Some( + "user_abc_account__session_0b4445e1-f5be-49e1-87ce-62bbc28ad705", + ))) + .unwrap(); + let b = convert_request(&request_with_user_id(Some( + "user_abc_account__session_11111111-2222-3333-4444-555555555555", + ))) + .unwrap(); + + assert_ne!( + a.conversation_state.conversation_id, + b.conversation_state.conversation_id + ); + assert_ne!( + a.conversation_state.agent_continuation_id, + b.conversation_state.agent_continuation_id + ); + } + + #[test] + fn missing_session_anchor_falls_back_to_random_per_request() { + // 拿不到锚点时保持旧行为:每请求独立随机,宁可丢连续性也不乱绑任务线。 + let first = convert_request(&request_with_user_id(None)).unwrap(); + let second = convert_request(&request_with_user_id(None)).unwrap(); + assert_ne!( + first.conversation_state.conversation_id, + second.conversation_state.conversation_id + ); + assert_ne!( + first.conversation_state.agent_continuation_id, + second.conversation_state.agent_continuation_id + ); + } }