From 4598070c2986d6022279ae94becf735789e45383 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 16 Jul 2026 18:49:30 +0800 Subject: [PATCH 1/9] fix(tools): preserve research workflow contracts --- README.md | 3 +- core/src/tools/builtin/batch.rs | 67 +++++++++++++-- core/src/tools/builtin/web_search.rs | 51 ++++++++--- core/src/tools/builtin/web_search/tests.rs | 98 ++++++++++++++++++++++ core/src/tools/mod.rs | 2 +- core/src/tools/program_tool.rs | 12 +-- core/src/tools/program_tool/tests.rs | 31 ++++++- 7 files changed, 237 insertions(+), 27 deletions(-) diff --git a/README.md b/README.md index c92e8665..0bee125a 100644 --- a/README.md +++ b/README.md @@ -452,7 +452,8 @@ repeating successful model work after a restart. `program` runs a bounded JavaScript module in QuickJS with a host-provided `ctx` surface. Time, nested tool calls, output bytes, recursion, and allowed -tools are limited. Recursive `program`, `dynamic_workflow`, and +tools are limited. Auditable workflow source is capped at 192 KiB. Recursive +`program`, `dynamic_workflow`, and `parallel_task` calls are excluded from its normal allow-list. `DynamicWorkflowRuntime` is a separate, explicit capability. A sandboxed diff --git a/core/src/tools/builtin/batch.rs b/core/src/tools/builtin/batch.rs index 0cfd5711..ca238dad 100644 --- a/core/src/tools/builtin/batch.rs +++ b/core/src/tools/builtin/batch.rs @@ -268,11 +268,44 @@ fn compact_child_metadata(metadata: Option) -> Option().into()) + .collect(), + ) + } else { + value.clone() + }; + compacted.insert(key.to_string(), value); + if serde_json::to_vec(&compacted) + .is_ok_and(|encoded| encoded.len() > MAX_CHILD_METADATA_BYTES) + { + compacted.remove(key); + } + } + } + Some(serde_json::Value::Object(compacted)) } // ============================================================================ @@ -427,6 +460,30 @@ mod tests { assert!(examples[0]["invocations"][0].get("name").is_none()); } + #[test] + fn compact_child_metadata_retains_search_routing_fields() { + let metadata = serde_json::json!({ + "status": "failed", + "engine_selection_source": "config", + "selected_engines": ["private-search", "x".repeat(5 * 1024)], + "engine_fallback": null, + "search_metrics": { + "oversized": "x".repeat(5 * 1024) + } + }); + + let compacted = compact_child_metadata(Some(metadata)).expect("compacted metadata"); + + assert_eq!(compacted["truncated"], true); + assert_eq!(compacted["status"], "failed"); + assert_eq!(compacted["engine_selection_source"], "config"); + assert_eq!(compacted["selected_engines"][0], "private-search"); + assert_eq!(compacted["selected_engines"][1].as_str().unwrap().len(), 96); + assert!(compacted.get("engine_fallback").is_some()); + assert!(compacted.get("search_metrics").is_none()); + assert!(serde_json::to_vec(&compacted).unwrap().len() <= 4 * 1024); + } + #[tokio::test] async fn test_execute_missing_invocations() { let tool = BatchTool::new(make_registry()); diff --git a/core/src/tools/builtin/web_search.rs b/core/src/tools/builtin/web_search.rs index fb309335..9c320634 100644 --- a/core/src/tools/builtin/web_search.rs +++ b/core/src/tools/builtin/web_search.rs @@ -284,6 +284,23 @@ fn add_http_engine(search: &mut Search, shortcut: &str, proxy_url: Option<&str>) } } +fn default_engine_selection<'a>( + config: Option<&'a crate::config::SearchConfig>, +) -> (Vec<&'a str>, &'static str) { + match config { + Some(config) if !config.engines.is_empty() => ( + config + .engines + .iter() + .filter(|(_, engine)| engine.enabled) + .map(|(name, _)| name.as_str()) + .collect(), + "config", + ), + _ => (vec!["ddg", "wiki"], "builtin_default"), + } +} + /// Add a headless engine using BrowserPool. fn add_headless_engine(search: &mut Search, shortcut: &str, pool: &Arc) -> bool { match shortcut.trim() { @@ -411,17 +428,14 @@ impl Tool for WebSearchTool { // Get configuration from context or use defaults let config = ctx.search_config.as_ref(); let default_timeout = config.map(|c| c.timeout).unwrap_or(10); - let default_engines: Vec<&str> = if let Some(cfg) = config { - // Build default engines list from enabled engines in config - cfg.engines - .iter() - .filter(|(_, engine_cfg)| engine_cfg.enabled) - .map(|(name, _)| name.as_str()) - .collect() + let (default_engines, default_engine_selection_source) = + default_engine_selection(config.map(Arc::as_ref)); + + let engine_selection_source = if args.get("engines").is_some() { + "request" } else { - vec!["ddg", "wiki"] + default_engine_selection_source }; - let engines: Vec<&str> = args .get("engines") .and_then(|v| { @@ -438,6 +452,10 @@ impl Tool for WebSearchTool { } }) .unwrap_or_else(|| default_engines.clone()); + let selected_engines = engines + .iter() + .map(|engine| engine.to_string()) + .collect::>(); // HTTP-only searches must not probe for or initialize a managed browser. let needs_headless = requires_headless_browser(&engines); @@ -533,7 +551,12 @@ impl Tool for WebSearchTool { if search.engine_count() == 0 { let message = format!("No valid engines found in: {:?}", engines); return Ok(ToolOutput::error(&message) - .with_error_kind(ToolErrorKind::InvalidArgument { message })); + .with_error_kind(ToolErrorKind::InvalidArgument { message }) + .with_metadata(serde_json::json!({ + "status": "failed", + "engine_selection_source": engine_selection_source, + "selected_engines": &selected_engines, + }))); } // Configure proxy if provided @@ -562,6 +585,8 @@ impl Tool for WebSearchTool { let output = ToolOutput::error(format!("Search failed: {message}")).with_metadata( serde_json::json!({ "status": "failed", + "engine_selection_source": engine_selection_source, + "selected_engines": &selected_engines, "search_metrics": search_metrics_json(&metrics), }), ); @@ -584,6 +609,8 @@ impl Tool for WebSearchTool { }) .with_metadata(serde_json::json!({ "status": "failed", + "engine_selection_source": engine_selection_source, + "selected_engines": &selected_engines, "search_metrics": search_metrics_json(&metrics), }))); } @@ -625,6 +652,8 @@ impl Tool for WebSearchTool { if results.is_empty() { let metadata = serde_json::json!({ "status": if errors.is_empty() { "complete" } else { "failed" }, + "engine_selection_source": engine_selection_source, + "selected_engines": &selected_engines, "engine_fallback": fell_back_from_headless.then_some("ddg,wiki"), "search_metrics": metrics_json, "engine_errors": engine_errors, @@ -679,6 +708,8 @@ impl Tool for WebSearchTool { Ok( ToolOutput::success(output).with_metadata(serde_json::json!({ "status": if errors.is_empty() { "complete" } else { "partial" }, + "engine_selection_source": engine_selection_source, + "selected_engines": &selected_engines, "source_anchors": source_anchors, "engine_fallback": fell_back_from_headless.then_some("ddg,wiki"), "search_metrics": metrics_json, diff --git a/core/src/tools/builtin/web_search/tests.rs b/core/src/tools/builtin/web_search/tests.rs index c41a8fe8..63953e8f 100644 --- a/core/src/tools/builtin/web_search/tests.rs +++ b/core/src/tools/builtin/web_search/tests.rs @@ -1,4 +1,5 @@ use super::*; +use crate::config::{SearchConfig, SearchEngineConfig}; use std::collections::HashMap; use std::path::PathBuf; @@ -50,6 +51,64 @@ fn latest_search_metrics_are_exposed_as_stable_metadata() { assert_eq!(metadata["latency_p99_ms"], 30); } +fn search_config(engines: HashMap) -> SearchConfig { + SearchConfig { + timeout: 10, + health: None, + engines, + headless: None, + } +} + +#[test] +fn default_engine_selection_uses_builtin_defaults_without_engine_configuration() { + let (engines, source) = default_engine_selection(None); + assert_eq!(engines, ["ddg", "wiki"]); + assert_eq!(source, "builtin_default"); + + let config = search_config(HashMap::new()); + let (engines, source) = default_engine_selection(Some(&config)); + assert_eq!(engines, ["ddg", "wiki"]); + assert_eq!(source, "builtin_default"); +} + +#[test] +fn default_engine_selection_respects_explicit_configuration() { + let config = search_config(HashMap::from([ + ( + "enabled".to_string(), + SearchEngineConfig { + enabled: true, + weight: 1.0, + timeout: None, + }, + ), + ( + "disabled".to_string(), + SearchEngineConfig { + enabled: false, + weight: 1.0, + timeout: None, + }, + ), + ])); + let (engines, source) = default_engine_selection(Some(&config)); + assert_eq!(engines, ["enabled"]); + assert_eq!(source, "config"); + + let config = search_config(HashMap::from([( + "disabled".to_string(), + SearchEngineConfig { + enabled: false, + weight: 1.0, + timeout: None, + }, + )])); + let (engines, source) = default_engine_selection(Some(&config)); + assert!(engines.is_empty()); + assert_eq!(source, "config"); +} + #[tokio::test] async fn test_web_search_missing_query() { let tool = WebSearchTool::new(); @@ -85,6 +144,45 @@ async fn test_web_search_no_valid_engines() { .unwrap(); assert!(!result.success); assert!(result.content.contains("No valid engines")); + let metadata = result.metadata.expect("search selection metadata"); + assert_eq!(metadata["status"], "failed"); + assert_eq!(metadata["engine_selection_source"], "request"); + assert_eq!( + metadata["selected_engines"], + serde_json::json!(["nonexistent"]) + ); +} + +#[tokio::test] +async fn configured_engine_selection_is_identified_in_failure_metadata() { + let tool = WebSearchTool::new(); + let ctx = ToolContext::new(PathBuf::from("/tmp")).with_search_config(SearchConfig { + timeout: 10, + health: None, + engines: HashMap::from([( + "private-search".to_string(), + SearchEngineConfig { + enabled: true, + weight: 1.0, + timeout: None, + }, + )]), + headless: None, + }); + + let result = tool + .execute(&serde_json::json!({"query": "test"}), &ctx) + .await + .unwrap(); + + assert!(!result.success); + let metadata = result.metadata.expect("search selection metadata"); + assert_eq!(metadata["status"], "failed"); + assert_eq!(metadata["engine_selection_source"], "config"); + assert_eq!( + metadata["selected_engines"], + serde_json::json!(["private-search"]) + ); } #[tokio::test] diff --git a/core/src/tools/mod.rs b/core/src/tools/mod.rs index c84c4d38..ddf7a3d8 100644 --- a/core/src/tools/mod.rs +++ b/core/src/tools/mod.rs @@ -33,7 +33,7 @@ pub use builtin::{ pub(crate) use invocation::{ registry_tool_invoker, HostDirectPolicy, InvocationOrigin, ToolInvocation, ToolInvoker, }; -pub use program_tool::ProgramTool; +pub use program_tool::{ProgramTool, MAX_PROGRAM_SCRIPT_SOURCE_BYTES}; pub use registry::ToolRegistry; pub use selector::{select_tools_for_messages, select_tools_for_prompt}; pub use task::{ diff --git a/core/src/tools/program_tool.rs b/core/src/tools/program_tool.rs index 604157b2..26a4f6ed 100644 --- a/core/src/tools/program_tool.rs +++ b/core/src/tools/program_tool.rs @@ -22,10 +22,10 @@ const DELEGATION_SCRIPT_TIMEOUT_MS: u64 = 600_000; const DEFAULT_SCRIPT_MAX_TOOL_CALLS: usize = 20; const DEFAULT_SCRIPT_MAX_OUTPUT_BYTES: usize = 64 * 1024; const PROGRAM_CANCELLATION_SETTLE_GRACE: Duration = Duration::from_millis(500); -// Engineered workflows include planner, maker, checker, recovery, and event -// projection logic in one auditable script. Keep a firm bound, but allow the -// complete workflow contract without requiring opaque minification. -const MAX_SCRIPT_SOURCE_BYTES: usize = 128 * 1024; +// Engineered workflows include planner, maker, checker, deterministic evidence +// gates, recovery, and event projection logic in one auditable script. Keep a +// firm bound while leaving enough room for explicit contracts and diagnostics. +pub const MAX_PROGRAM_SCRIPT_SOURCE_BYTES: usize = 192 * 1024; pub struct ProgramTool { fallback_invoker: Arc, @@ -158,11 +158,11 @@ async fn execute_script_program( Ok(source) => source, Err(message) => return Ok(ToolOutput::error(message)), }; - if source.len() > MAX_SCRIPT_SOURCE_BYTES { + if source.len() > MAX_PROGRAM_SCRIPT_SOURCE_BYTES { return Ok(ToolOutput::error(format!( "script source is too large: {} bytes exceeds {} bytes", source.len(), - MAX_SCRIPT_SOURCE_BYTES + MAX_PROGRAM_SCRIPT_SOURCE_BYTES ))); } if let Err(message) = validate_script_source(&source) { diff --git a/core/src/tools/program_tool/tests.rs b/core/src/tools/program_tool/tests.rs index 08b12d28..d7083605 100644 --- a/core/src/tools/program_tool/tests.rs +++ b/core/src/tools/program_tool/tests.rs @@ -133,7 +133,7 @@ async fn program_tool_rejects_unsupported_language() { #[tokio::test] async fn program_tool_rejects_source_above_engineered_workflow_limit() { let tool = ProgramTool::new(Arc::new(ToolRegistry::new(PathBuf::from("/tmp")))); - let source = "x".repeat(MAX_SCRIPT_SOURCE_BYTES + 1); + let source = "x".repeat(MAX_PROGRAM_SCRIPT_SOURCE_BYTES + 1); let output = tool .execute( &serde_json::json!({ @@ -147,9 +147,32 @@ async fn program_tool_rejects_source_above_engineered_workflow_limit() { assert!(!output.success); assert!(output.content.contains("script source is too large")); - assert!(output - .content - .contains(&format!("exceeds {} bytes", MAX_SCRIPT_SOURCE_BYTES))); + assert!(output.content.contains(&format!( + "exceeds {} bytes", + MAX_PROGRAM_SCRIPT_SOURCE_BYTES + ))); +} + +#[tokio::test] +async fn program_tool_accepts_auditable_workflow_above_legacy_limit() { + let tool = ProgramTool::new(Arc::new(ToolRegistry::new(PathBuf::from("/tmp")))); + let source = format!( + "async function run() {{ return {{ accepted: true }}; }}{}", + " ".repeat(140 * 1024) + ); + let output = tool + .execute( + &serde_json::json!({ + "type": "script", + "source": source + }), + &ToolContext::new(PathBuf::from("/tmp")), + ) + .await + .unwrap(); + + assert!(output.success, "{}", output.content); + assert!(output.content.contains("accepted"), "{}", output.content); } #[tokio::test] From 31d7466bb0de7f5f167ee2723d6e0f23c39d621c Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 14:45:53 +0800 Subject: [PATCH 2/9] fix(agent): defer model validation to session --- core/src/agent_api/agent_bootstrap.rs | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/core/src/agent_api/agent_bootstrap.rs b/core/src/agent_api/agent_bootstrap.rs index a03b82b2..30d1b4bd 100644 --- a/core/src/agent_api/agent_bootstrap.rs +++ b/core/src/agent_api/agent_bootstrap.rs @@ -51,10 +51,10 @@ pub(super) fn load_code_config(config_source: String) -> Result { } pub(super) async fn build_agent_from_config(config: CodeConfig) -> Result { - config - .default_llm_config() - .context("default_model must be set in 'provider/model' format with a valid API key")?; - + // A host may inject an `LlmClient` through `SessionOptions`, so model + // validation belongs to session construction where that override is known. + // Sessions without either a host client or a valid configured default still + // fail in `resolve_session_llm_client` with the same actionable error. let mut agent_config = base_agent_config(&config); install_global_skill_registry(&mut agent_config, &config); let (global_mcp, global_mcp_tools) = connect_global_mcp(&config).await; @@ -159,4 +159,11 @@ mod tests { let err = load_code_config("{\"default_model\":\"x\"}".to_string()).unwrap_err(); assert!(err.to_string().contains("JSON config is not supported")); } + + #[tokio::test] + async fn agent_bootstrap_allows_a_host_supplied_session_client() { + build_agent_from_config(CodeConfig::default()) + .await + .expect("bootstrap should defer model validation to the session"); + } } From 59b44c551e072e9dfeabf554b0b6506ff9bb0955 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 14:56:43 +0800 Subject: [PATCH 3/9] fix(agent): retry interrupted streams in place --- README.md | 9 ++ core/Cargo.toml | 2 + core/src/agent.rs | 4 +- core/src/agent/extra_agent_tests.rs | 212 +++++++++++++++++++++++++++ core/src/agent/llm_turn.rs | 72 ++++++++- core/src/security/event_sanitizer.rs | 124 +++++++++++++++- 6 files changed, 419 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 0bee125a..892452ad 100644 --- a/README.md +++ b/README.md @@ -404,6 +404,15 @@ reasoning, tool calls, usage, images, retries, cancellation, and streaming where the provider supports them. Structured generation uses native response formats when available and a schema-validated prompt fallback otherwise. +If an established response stream closes before its final response, Core +restarts the same LLM turn with the same message snapshot up to ten times. The +delay grows exponentially from one second and is capped at thirty seconds, with +jitter. Core re-emits `TurnStart` with the same turn number before each wait so +stream consumers can roll back provisional deltas and tool drafts in place; +only the successful assistant response is committed. Cancellation interrupts +the wait immediately, and setup/status retries remain owned by the provider +transport rather than nesting another turn-level retry loop. + MCP connections support stdio, HTTP SSE, streamable HTTP, and OAuth client credentials. Global managers can seed new sessions and refresh cached tools; session managers can connect, disconnect, and live-add or remove isolated diff --git a/core/Cargo.toml b/core/Cargo.toml index 9395be30..b7e98937 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -147,6 +147,8 @@ s3 = [ serve = ["dep:cron"] [dev-dependencies] +# Deterministically advance the exponential stream-retry backoff in async tests. +tokio = { version = "1.35", features = ["test-util"] } # HTTP mocking for the RemoteGitBackend (and any future HTTP-backed workspace # provider). Production code does not depend on wiremock. wiremock = "0.6" diff --git a/core/src/agent.rs b/core/src/agent.rs index 739197f9..ca14d49c 100644 --- a/core/src/agent.rs +++ b/core/src/agent.rs @@ -295,7 +295,9 @@ pub enum AgentEvent { description: String, }, - /// LLM turn started + /// LLM turn started. The same turn number is emitted again when an + /// interrupted response stream is retried; consumers should roll back + /// provisional output from that turn before applying replacement deltas. #[serde(rename = "turn_start")] TurnStart { turn: usize }, diff --git a/core/src/agent/extra_agent_tests.rs b/core/src/agent/extra_agent_tests.rs index b8c3b27c..55a8cb5a 100644 --- a/core/src/agent/extra_agent_tests.rs +++ b/core/src/agent/extra_agent_tests.rs @@ -5,6 +5,7 @@ use crate::queue::SessionQueueConfig; use crate::tools::ToolExecutor; use std::path::PathBuf; use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; use tokio::sync::mpsc; fn test_tool_context() -> ToolContext { @@ -3464,6 +3465,217 @@ async fn test_circuit_breaker_succeeds_if_llm_recovers() { assert_eq!(budget.records.load(Ordering::SeqCst), 1); } +struct InterruptedStreamClient { + failures_before_success: usize, + calls: AtomicUsize, + requests: std::sync::Mutex>, +} + +impl InterruptedStreamClient { + fn new(failures_before_success: usize) -> Self { + Self { + failures_before_success, + calls: AtomicUsize::new(0), + requests: std::sync::Mutex::new(Vec::new()), + } + } +} + +#[async_trait::async_trait] +impl LlmClient for InterruptedStreamClient { + async fn complete( + &self, + _messages: &[Message], + _system: Option<&str>, + _tools: &[ToolDefinition], + ) -> Result { + anyhow::bail!("an established response stream must retry in streaming mode") + } + + async fn complete_streaming( + &self, + messages: &[Message], + _system: Option<&str>, + _tools: &[ToolDefinition], + _cancel_token: tokio_util::sync::CancellationToken, + ) -> Result> { + let attempt = self.calls.fetch_add(1, Ordering::SeqCst); + self.requests + .lock() + .unwrap() + .push(serde_json::to_string(messages).unwrap()); + let failures_before_success = self.failures_before_success; + let (tx, rx) = mpsc::channel(2); + tokio::spawn(async move { + if attempt < failures_before_success { + tx.send(StreamEvent::TextDelta(format!("partial attempt {attempt}"))) + .await + .ok(); + return; + } + + let response = MockLlmClient::text_response("Recovered in place."); + tx.send(StreamEvent::TextDelta(response.text())).await.ok(); + tx.send(StreamEvent::Done(response)).await.ok(); + }); + Ok(rx) + } +} + +async fn collect_stream_retry_events( + mut rx: mpsc::Receiver, +) -> (Vec, usize) { + let mut events = Vec::new(); + let mut turn_starts = 0usize; + while let Some(event) = rx.recv().await { + let terminal = matches!(event, AgentEvent::End { .. } | AgentEvent::Error { .. }); + if matches!(event, AgentEvent::TurnStart { .. }) { + turn_starts += 1; + if turn_starts > 1 { + // Retry delays are capped at 30 seconds. Yield first so the + // producer can enter its cancellation-aware sleep. + tokio::task::yield_now().await; + tokio::time::advance(Duration::from_secs(31)).await; + tokio::task::yield_now().await; + } + } + events.push(event); + if terminal { + break; + } + } + (events, turn_starts) +} + +#[tokio::test(start_paused = true)] +async fn interrupted_stream_retries_ten_times_in_the_same_turn() { + let client = Arc::new(InterruptedStreamClient::new(10)); + let agent = AgentLoop::new( + client.clone(), + Arc::new(ToolExecutor::new("/tmp".to_string())), + test_tool_context(), + AgentConfig { + planning_mode: crate::prompts::PlanningMode::Disabled, + ..AgentConfig::default() + }, + ); + let (tx, rx) = mpsc::channel(128); + let run = tokio::spawn(async move { + agent + .execute_with_session(&[], "keep this message once", None, Some(tx), None) + .await + }); + + let (events, turn_starts) = collect_stream_retry_events(rx).await; + let result = run.await.unwrap().unwrap(); + + assert_eq!(client.calls.load(Ordering::SeqCst), 11); + assert_eq!(turn_starts, 11, "one start plus ten in-place restarts"); + assert_eq!( + events + .iter() + .filter(|event| matches!(event, AgentEvent::Error { .. })) + .count(), + 0 + ); + assert_eq!( + result + .messages + .iter() + .filter(|message| message.role == "user") + .count(), + 1, + "provider retries must not append another user message" + ); + assert_eq!( + result + .messages + .iter() + .filter(|message| message.role == "assistant") + .count(), + 1, + "only the successful response may be committed" + ); + + let requests = client.requests.lock().unwrap(); + assert_eq!(requests.len(), 11); + assert!(requests.windows(2).all(|pair| pair[0] == pair[1])); +} + +#[tokio::test(start_paused = true)] +async fn interrupted_stream_stops_after_ten_retries() { + let client = Arc::new(InterruptedStreamClient::new(usize::MAX)); + let agent = AgentLoop::new( + client.clone(), + Arc::new(ToolExecutor::new("/tmp".to_string())), + test_tool_context(), + AgentConfig { + planning_mode: crate::prompts::PlanningMode::Disabled, + ..AgentConfig::default() + }, + ); + let (tx, rx) = mpsc::channel(128); + let run = tokio::spawn(async move { + agent + .execute_with_session(&[], "fail once", None, Some(tx), None) + .await + }); + + let (events, turn_starts) = collect_stream_retry_events(rx).await; + let error = run.await.unwrap().unwrap_err().to_string(); + + assert_eq!(client.calls.load(Ordering::SeqCst), 11); + assert_eq!(turn_starts, 11); + assert!(error.contains("after 10 retries"), "{error}"); + assert_eq!( + events + .iter() + .filter(|event| matches!(event, AgentEvent::Error { .. })) + .count(), + 1, + "retry exhaustion should emit one terminal error" + ); +} + +#[tokio::test] +async fn interrupted_stream_backoff_stops_immediately_on_cancellation() { + let client = Arc::new(InterruptedStreamClient::new(usize::MAX)); + let agent = AgentLoop::new( + client.clone(), + Arc::new(ToolExecutor::new("/tmp".to_string())), + test_tool_context(), + AgentConfig { + planning_mode: crate::prompts::PlanningMode::Disabled, + ..AgentConfig::default() + }, + ); + let cancellation = tokio_util::sync::CancellationToken::new(); + let run_cancellation = cancellation.clone(); + let (tx, mut rx) = mpsc::channel(16); + let run = tokio::spawn(async move { + agent + .execute_with_session(&[], "cancel retry", None, Some(tx), Some(&run_cancellation)) + .await + }); + + let mut starts = 0usize; + while let Some(event) = rx.recv().await { + if matches!(event, AgentEvent::TurnStart { .. }) { + starts += 1; + if starts == 2 { + break; + } + } + } + cancellation.cancel(); + + let _result = tokio::time::timeout(Duration::from_millis(50), run) + .await + .expect("cancellation must interrupt exponential backoff") + .unwrap(); + assert_eq!(client.calls.load(Ordering::SeqCst), 1); +} + // ── Continuation detection tests ───────────────────────────────────────── #[test] diff --git a/core/src/agent/llm_turn.rs b/core/src/agent/llm_turn.rs index 7914b66a..bb564cbb 100644 --- a/core/src/agent/llm_turn.rs +++ b/core/src/agent/llm_turn.rs @@ -7,11 +7,21 @@ use crate::llm::{ estimate_prompt_tokens, non_retryable_llm_error_message, LlmResponse, Message, ToolCall, ToolDefinition, }; +use crate::retry::RetryConfig; use anyhow::Context; use std::time::Duration; use tokio::sync::mpsc; const DEFAULT_AUTO_COMPACT_TIMEOUT_MS: u64 = 60_000; +const MAX_STREAM_INTERRUPTION_RETRIES: u32 = 10; + +#[derive(Debug, thiserror::Error)] +#[error("LLM response stream ended before the final response")] +struct IncompleteLlmStream; + +fn stream_interruption_retry_delay(retry_index: u32) -> Duration { + RetryConfig::default().delay_for_attempt(retry_index) +} pub(super) struct LlmTurnOutput { pub(super) turn: usize, @@ -217,6 +227,7 @@ impl AgentLoop { ) -> anyhow::Result { let threshold = self.config.circuit_breaker_threshold.max(1); let mut attempt = 0u32; + let mut stream_retries = 0u32; let llm_client = self.scoped_llm_client_for_parts( request.session_id, request.event_tx, @@ -250,7 +261,41 @@ impl AgentLoop { } let non_retryable_message = non_retryable_llm_error_message(&error); - if non_retryable_message.is_none() + let stream_interrupted = error.downcast_ref::().is_some(); + if stream_interrupted + && non_retryable_message.is_none() + && stream_retries < MAX_STREAM_INTERRUPTION_RETRIES + { + let retry_index = stream_retries; + stream_retries += 1; + let delay = stream_interruption_retry_delay(retry_index); + tracing::warn!( + turn = request.turn, + attempt, + retry = stream_retries, + max_retries = MAX_STREAM_INTERRUPTION_RETRIES, + delay_ms = delay.as_millis() as u64, + error = %error, + "LLM response stream was interrupted; restarting the same turn" + ); + + // Re-emitting the same turn number is the transactional + // boundary for stream consumers: discard provisional + // deltas/tool drafts from the failed attempt, but keep + // the original user message and all prior turns. + self.emit_turn_start(request.turn, request.event_tx).await; + tokio::select! { + biased; + _ = request.cancel_token.cancelled() => { + anyhow::bail!("Operation cancelled by user") + } + _ = tokio::time::sleep(delay) => {} + } + continue; + } + + if !stream_interrupted + && non_retryable_message.is_none() && attempt < threshold && (request.event_tx.is_none() || attempt == 1) { @@ -273,6 +318,11 @@ impl AgentLoop { let msg = if let Some(message) = non_retryable_message { message.to_string() + } else if stream_interrupted { + format!( + "LLM response stream interrupted after {} retries ({} attempts): {}", + stream_retries, attempt, error + ) } else if attempt > 1 { format!( "LLM circuit breaker triggered: failed after {} attempt(s): {}", @@ -593,7 +643,7 @@ impl AgentLoop { } } } - final_response.context("Stream ended without final response") + final_response.ok_or_else(|| anyhow::Error::new(IncompleteLlmStream)) } else { self.call_non_streaming_llm(llm_client, messages, system, tools, cancel_token) .await @@ -690,3 +740,21 @@ impl AgentLoop { } } } + +#[cfg(test)] +mod stream_retry_tests { + use super::*; + + #[test] + fn stream_retry_delay_is_exponential_and_capped() { + let first = stream_interruption_retry_delay(0).as_millis(); + let second = stream_interruption_retry_delay(1).as_millis(); + let third = stream_interruption_retry_delay(2).as_millis(); + let capped = stream_interruption_retry_delay(9).as_millis(); + + assert!((750..=1_250).contains(&first)); + assert!((1_500..=2_500).contains(&second)); + assert!((3_000..=5_000).contains(&third)); + assert!((22_500..=37_500).contains(&capped)); + } +} diff --git a/core/src/security/event_sanitizer.rs b/core/src/security/event_sanitizer.rs index 0a48546e..7e86a3ea 100644 --- a/core/src/security/event_sanitizer.rs +++ b/core/src/security/event_sanitizer.rs @@ -24,25 +24,42 @@ pub(crate) struct AgentEventStreamSanitizer { tool_inputs: HashMap, tool_outputs: HashMap, current_tool_input_id: Option, + active_turn: Option, + attempt_checkpoint: Option, } +#[derive(Clone)] struct BufferedDelta { sequence: u64, value: String, } +#[derive(Clone)] struct BufferedToolInput { sequence: u64, emitted_id: Option, value: String, } +#[derive(Clone)] struct BufferedToolOutput { sequence: u64, name: String, value: String, } +#[derive(Clone)] +struct StreamAttemptCheckpoint { + next_sequence: u64, + buffered_bytes: usize, + failed_closed: bool, + text: Option, + reasoning: Option, + tool_inputs: HashMap, + tool_outputs: HashMap, + current_tool_input_id: Option, +} + #[derive(Clone, Debug, Eq, Hash, PartialEq)] enum ToolInputKey { Known(String), @@ -68,6 +85,8 @@ impl AgentEventStreamSanitizer { tool_inputs: HashMap::new(), tool_outputs: HashMap::new(), current_tool_input_id: None, + active_turn: None, + attempt_checkpoint: None, } } @@ -158,6 +177,7 @@ impl AgentEventStreamSanitizer { event => { let mut ready = Vec::new(); match &event { + AgentEvent::TurnStart { turn } => self.begin_or_restart_turn(*turn), AgentEvent::ToolStart { id, .. } => { self.bind_anonymous_tool_input(id); self.current_tool_input_id = Some(id.clone()); @@ -174,9 +194,15 @@ impl AgentEventStreamSanitizer { | AgentEvent::ConfirmationTimeout { tool_id, .. } => { self.take_tool_fields(tool_id, false, &mut ready); } - AgentEvent::TurnEnd { .. } => self.take_all_tool_inputs(&mut ready), + AgentEvent::TurnEnd { .. } => { + self.take_all_tool_inputs(&mut ready); + self.active_turn = None; + self.attempt_checkpoint = None; + } AgentEvent::End { .. } | AgentEvent::Error { .. } => { self.take_all(&mut ready); + self.active_turn = None; + self.attempt_checkpoint = None; } _ => {} } @@ -200,6 +226,46 @@ impl AgentEventStreamSanitizer { self.sanitize_pending(ready) } + /// A repeated `TurnStart` with the same turn number marks an in-place LLM + /// stream retry. Restore buffered security domains to their state before + /// that attempt so discarded partial output can never be concatenated with + /// the replacement response. + fn begin_or_restart_turn(&mut self, turn: usize) { + if self.active_turn == Some(turn) { + if let Some(checkpoint) = self.attempt_checkpoint.clone() { + self.restore_attempt_checkpoint(checkpoint); + } + return; + } + + self.active_turn = Some(turn); + self.attempt_checkpoint = Some(self.attempt_checkpoint()); + } + + fn attempt_checkpoint(&self) -> StreamAttemptCheckpoint { + StreamAttemptCheckpoint { + next_sequence: self.next_sequence, + buffered_bytes: self.buffered_bytes, + failed_closed: self.failed_closed, + text: self.text.clone(), + reasoning: self.reasoning.clone(), + tool_inputs: self.tool_inputs.clone(), + tool_outputs: self.tool_outputs.clone(), + current_tool_input_id: self.current_tool_input_id.clone(), + } + } + + fn restore_attempt_checkpoint(&mut self, checkpoint: StreamAttemptCheckpoint) { + self.next_sequence = checkpoint.next_sequence; + self.buffered_bytes = checkpoint.buffered_bytes; + self.failed_closed = checkpoint.failed_closed; + self.text = checkpoint.text; + self.reasoning = checkpoint.reasoning; + self.tool_inputs = checkpoint.tool_inputs; + self.tool_outputs = checkpoint.tool_outputs; + self.current_tool_input_id = checkpoint.current_tool_input_id; + } + fn reserve_sequence(&mut self) -> u64 { let sequence = self.next_sequence; self.next_sequence = self.next_sequence.saturating_add(1); @@ -613,6 +679,62 @@ mod tests { )); } + #[test] + fn repeated_turn_start_discards_only_the_interrupted_attempt_buffers() { + use crate::agent::AgentEvent; + + let provider: Arc = Arc::new(DefaultSecurityProvider::new()); + let mut sanitizer = AgentEventStreamSanitizer::new(Some(provider)); + let mut emitted = Vec::new(); + + emitted.extend(sanitizer.push(AgentEvent::TurnStart { turn: 1 })); + emitted.extend(sanitizer.push(AgentEvent::TextDelta { + text: "stable response. ".to_string(), + })); + emitted.extend(sanitizer.push(AgentEvent::TurnEnd { + turn: 1, + usage: crate::llm::TokenUsage::default(), + })); + + emitted.extend(sanitizer.push(AgentEvent::TurnStart { turn: 2 })); + emitted.extend(sanitizer.push(AgentEvent::TextDelta { + text: "discarded partial. ".to_string(), + })); + emitted.extend(sanitizer.push(AgentEvent::ReasoningDelta { + text: "discarded reasoning. ".to_string(), + })); + emitted.extend(sanitizer.push(AgentEvent::ToolStart { + id: "discarded-tool".to_string(), + name: "bash".to_string(), + })); + emitted.extend(sanitizer.push(AgentEvent::ToolInputDelta { + id: Some("discarded-tool".to_string()), + delta: "discarded input".to_string(), + })); + + // Core re-emits the same turn number immediately before waiting and + // restarting the provider request. + emitted.extend(sanitizer.push(AgentEvent::TurnStart { turn: 2 })); + emitted.extend(sanitizer.push(AgentEvent::TextDelta { + text: "recovered response.".to_string(), + })); + emitted.extend(sanitizer.push(AgentEvent::End { + text: "recovered response.".to_string(), + usage: crate::llm::TokenUsage::default(), + verification_summary: Box::new(crate::verification::VerificationSummary::from_reports( + &[], + )), + meta: None, + })); + + let serialized = serde_json::to_string(&emitted).unwrap(); + assert!(serialized.contains("stable response"), "{serialized}"); + assert!(serialized.contains("recovered response"), "{serialized}"); + assert!(!serialized.contains("discarded partial"), "{serialized}"); + assert!(!serialized.contains("discarded reasoning"), "{serialized}"); + assert!(!serialized.contains("discarded input"), "{serialized}"); + } + #[test] fn stream_sanitizer_drops_all_deltas_after_aggregate_limit() { let provider: Arc = Arc::new(DefaultSecurityProvider::new()); From fa20c2c5f95f019e844b7daef3f139aa9e4892a5 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 14:59:52 +0800 Subject: [PATCH 4/9] feat(search): expose Bing RSS engine --- core/src/tools/builtin/web_search.rs | 11 ++++++++--- core/src/tools/builtin/web_search/tests.rs | 8 +++++++- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/core/src/tools/builtin/web_search.rs b/core/src/tools/builtin/web_search.rs index 9c320634..e942c980 100644 --- a/core/src/tools/builtin/web_search.rs +++ b/core/src/tools/builtin/web_search.rs @@ -4,7 +4,8 @@ use crate::config::{BrowserBackend, HeadlessConfig}; use crate::tools::types::{Tool, ToolContext, ToolErrorKind, ToolOutput}; use a3s_search::a3s_use_browser::{BrowserPool, BrowserPoolConfig, BrowserProvider}; use a3s_search::engines::{ - Baidu, BingChina, BraveParser, DuckDuckGoParser, Google, So360Parser, SogouParser, Wikipedia, + Baidu, BingChina, BingParser, BraveParser, DuckDuckGoParser, Google, So360Parser, SogouParser, + Wikipedia, }; use a3s_search::proxy::{ProxyConfig, ProxyPool}; use a3s_search::WaitStrategy; @@ -264,6 +265,10 @@ fn add_http_engine(search: &mut Search, shortcut: &str, proxy_url: Option<&str>) search.add_engine(HtmlEngine::with_fetcher(BraveParser, Arc::new(fetcher()))); true } + "bing" => { + search.add_engine(HtmlEngine::with_fetcher(BingParser, Arc::new(fetcher()))); + true + } "wiki" => { search.add_engine(Wikipedia::with_http_fetcher(fetcher())); true @@ -332,7 +337,7 @@ impl Tool for WebSearchTool { fn description(&self) -> &str { "Search the web using multiple search engines. Aggregates results from multiple engines \ - (DuckDuckGo, Wikipedia, Brave, Sogou, 360, Google, Baidu, Bing China, etc.). \ + (DuckDuckGo, Wikipedia, Brave, Bing, Sogou, 360, Google, Baidu, Bing China, etc.). \ Supports proxy configuration for anti-crawler protection. Returns deduplicated and ranked results. \ Google and Baidu use a headless browser; Bing China uses its HTTP RSS endpoint." } @@ -351,7 +356,7 @@ impl Tool for WebSearchTool { "items": { "type": "string" }, - "description": "Optional. List of search engines to use. Default: [\"ddg\",\"wiki\"]. Available: ddg (DuckDuckGo), brave (Brave Search), wiki (Wikipedia), sogou (Sogou), 360 / so360 (360 Search), bing_cn (Bing China RSS), g / google (Google, headless), baidu (Baidu, headless)." + "description": "Optional. List of search engines to use. Default: [\"ddg\",\"wiki\"]. Available: ddg (DuckDuckGo), brave (Brave Search), bing (Bing RSS), wiki (Wikipedia), sogou (Sogou), 360 / so360 (360 Search), bing_cn (Bing China RSS), g / google (Google, headless), baidu (Baidu, headless)." }, "limit": { "type": "integer", diff --git a/core/src/tools/builtin/web_search/tests.rs b/core/src/tools/builtin/web_search/tests.rs index 63953e8f..4e53de84 100644 --- a/core/src/tools/builtin/web_search/tests.rs +++ b/core/src/tools/builtin/web_search/tests.rs @@ -380,6 +380,9 @@ fn test_web_search_schema_is_canonical() { // Example with engines should use array format assert!(examples[1]["engines"].is_array()); assert_eq!(examples[1]["engines"].as_array().unwrap(), &["ddg", "wiki"]); + assert!(params["properties"]["engines"]["description"] + .as_str() + .is_some_and(|description| description.contains("bing (Bing RSS)"))); } #[test] @@ -418,8 +421,11 @@ fn test_add_http_engine_valid() { assert!(add_http_engine(&mut search, "brave", None)); assert_eq!(search.engine_count(), 3); - assert!(add_http_engine(&mut search, "bing_cn", None)); + assert!(add_http_engine(&mut search, "bing", None)); assert_eq!(search.engine_count(), 4); + + assert!(add_http_engine(&mut search, "bing_cn", None)); + assert_eq!(search.engine_count(), 5); } #[test] From fc5b8769f9f19b1d46af921d447de47681834428 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 15:51:08 +0800 Subject: [PATCH 5/9] fix(search): satisfy release lint gate --- core/src/tools/builtin/web_search.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/core/src/tools/builtin/web_search.rs b/core/src/tools/builtin/web_search.rs index e942c980..28b8dd21 100644 --- a/core/src/tools/builtin/web_search.rs +++ b/core/src/tools/builtin/web_search.rs @@ -289,9 +289,9 @@ fn add_http_engine(search: &mut Search, shortcut: &str, proxy_url: Option<&str>) } } -fn default_engine_selection<'a>( - config: Option<&'a crate::config::SearchConfig>, -) -> (Vec<&'a str>, &'static str) { +fn default_engine_selection( + config: Option<&crate::config::SearchConfig>, +) -> (Vec<&str>, &'static str) { match config { Some(config) if !config.engines.is_empty() => ( config From 1546de505f6b3f4f4d5ad4e7f8160d0dabbb9249 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 15:51:09 +0800 Subject: [PATCH 6/9] test(agent): cover deferred session model validation --- core/src/agent_api/tests.rs | 21 ++++++++++++++++----- 1 file changed, 16 insertions(+), 5 deletions(-) diff --git a/core/src/agent_api/tests.rs b/core/src/agent_api/tests.rs index 941a0a8c..5e221129 100644 --- a/core/src/agent_api/tests.rs +++ b/core/src/agent_api/tests.rs @@ -1210,9 +1210,8 @@ async fn test_new_rejects_hcl_files() { assert!(msg.contains(".acl")); } -#[test] -fn test_from_config_requires_default_model() { - let rt = tokio::runtime::Runtime::new().unwrap(); +#[tokio::test] +async fn test_from_config_defers_default_model_validation_until_session() { let config = CodeConfig { providers: vec![ProviderConfig { name: "anthropic".to_string(), @@ -1224,8 +1223,20 @@ fn test_from_config_requires_default_model() { }], ..Default::default() }; - let result = rt.block_on(Agent::from_config(config)); - assert!(result.is_err()); + let agent = Agent::from_config(config) + .await + .expect("agent bootstrap should allow a host-supplied session client"); + + let result = agent + .session_async("/tmp/test-missing-default-model", None) + .await; + let error = result + .err() + .expect("session without an LLM source must fail"); + assert!( + error.to_string().contains("default_model must be set"), + "unexpected session error: {error}" + ); } #[tokio::test] From 5dc90b1003697a421d57e59b6de974557cd1aaf9 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 16:04:18 +0800 Subject: [PATCH 7/9] release: prepare a3s-code 5.3.5 --- .github/workflows/release.yml | 4 +- CHANGELOG.md | 23 ++ Cargo.lock | 2 +- core/Cargo.toml | 2 +- scripts/check_semver.sh | 5 +- sdk/node/Cargo.lock | 4 +- sdk/node/Cargo.toml | 4 +- sdk/node/examples/package-lock.json | 14 +- sdk/node/package-lock.json | 16 +- sdk/node/package.json | 14 +- sdk/python-bootstrap/README.md | 6 +- .../build/lib/a3s_code/__init__.py | 18 -- .../build/lib/a3s_code/_bootstrap.py | 213 ------------------ .../build/lib/a3s_code/py.typed | 0 sdk/python-bootstrap/pyproject.toml | 10 +- .../src/a3s_code/_bootstrap.py | 2 +- sdk/python/CHANGELOG.md | 8 + sdk/python/Cargo.lock | 4 +- sdk/python/Cargo.toml | 4 +- sdk/python/README.md | 6 +- sdk/python/pyproject.toml | 2 +- 21 files changed, 84 insertions(+), 277 deletions(-) delete mode 100644 sdk/python-bootstrap/build/lib/a3s_code/__init__.py delete mode 100644 sdk/python-bootstrap/build/lib/a3s_code/_bootstrap.py delete mode 100644 sdk/python-bootstrap/build/lib/a3s_code/py.typed diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index d5ecac80..e29880b0 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -55,10 +55,10 @@ jobs: test -n "$VERSION" bash check-version.sh "$VERSION" - - name: Check public API compatibility with v5.3.3 + - name: Check public API compatibility with v5.3.4 run: | cargo install cargo-semver-checks --version 0.48.0 --locked - bash scripts/check_semver.sh 5.3.3 + bash scripts/check_semver.sh 5.3.4 - name: Check SDK protocol and API alignment run: | diff --git a/CHANGELOG.md b/CHANGELOG.md index bd8837d2..503efb97 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,29 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [5.3.5] - 2026-07-17 + +### Added + +- Added `bing` as a first-class HTTP RSS search engine and exposed selected + engines plus their request, configuration, or built-in-default origin in + search result metadata. + +### Changed + +- Raised the bounded auditable `program` source limit to 192 KiB and retained + compact search-routing fields when oversized child metadata passes through a + batch workflow. + +### Fixed + +- Allowed hosts that supply an `LlmClient` in session options to bootstrap an + agent without a configured default model, while preserving session-time + validation when neither source is available. +- Retried an established response stream up to ten times within the same turn, + with cancellation-aware exponential backoff and transactional rollback of + provisional text, reasoning, and tool drafts between attempts. + ## [5.3.4] - 2026-07-16 ### Fixed diff --git a/Cargo.lock b/Cargo.lock index 14ae86d4..b53abff7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -15,7 +15,7 @@ source = "git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ac [[package]] name = "a3s-code-core" -version = "5.3.4" +version = "5.3.5" dependencies = [ "a3s-acl 0.2.1 (git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ace873aaf15ce22)", "a3s-common", diff --git a/core/Cargo.toml b/core/Cargo.toml index b7e98937..459a6910 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "a3s-code-core" -version = "5.3.4" +version = "5.3.5" edition = "2021" authors = ["A3S Lab Team"] license = "MIT" diff --git a/scripts/check_semver.sh b/scripts/check_semver.sh index ffec4cd6..713be62f 100644 --- a/scripts/check_semver.sh +++ b/scripts/check_semver.sh @@ -3,10 +3,13 @@ set -euo pipefail -BASELINE_VERSION="${1:-5.3.3}" +BASELINE_VERSION="${1:-5.3.4}" PACKAGE="a3s-code-core" case "$BASELINE_VERSION" in + 5.3.4) + BASELINE_SHA256="2ea4c48286d828e09fb44df83144d05b2d41db25e4695f3bdce768e7a46e0399" + ;; 5.3.3) BASELINE_SHA256="5ed95dff8354578962130615e36654f7107c7bd9a0b862fe8e9e6ae1ed1676f9" ;; diff --git a/sdk/node/Cargo.lock b/sdk/node/Cargo.lock index b947c2c2..060a9fc5 100644 --- a/sdk/node/Cargo.lock +++ b/sdk/node/Cargo.lock @@ -15,7 +15,7 @@ source = "git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ac [[package]] name = "a3s-code-core" -version = "5.3.4" +version = "5.3.5" dependencies = [ "a3s-acl 0.2.1 (git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ace873aaf15ce22)", "a3s-common", @@ -76,7 +76,7 @@ dependencies = [ [[package]] name = "a3s-code-node" -version = "5.3.4" +version = "5.3.5" dependencies = [ "a3s-code-core", "anyhow", diff --git a/sdk/node/Cargo.toml b/sdk/node/Cargo.toml index 44944972..18de7dd3 100644 --- a/sdk/node/Cargo.toml +++ b/sdk/node/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "a3s-code-node" -version = "5.3.4" +version = "5.3.5" edition = "2021" authors = ["A3S Lab Team"] license = "MIT" @@ -11,7 +11,7 @@ description = "A3S Code Node.js bindings - Native addon via napi-rs" crate-type = ["cdylib"] [dependencies] -a3s-code-core = { version = "5.3.4", path = "../../core", features = ["s3", "serve"] } +a3s-code-core = { version = "5.3.5", path = "../../core", features = ["s3", "serve"] } napi = { version = "2", features = ["async", "napi6", "serde-json"] } napi-derive = "2" tokio = { version = "1.35", features = ["full"] } diff --git a/sdk/node/examples/package-lock.json b/sdk/node/examples/package-lock.json index 09b8c6b2..c092c849 100644 --- a/sdk/node/examples/package-lock.json +++ b/sdk/node/examples/package-lock.json @@ -18,7 +18,7 @@ }, "..": { "name": "@a3s-lab/code", - "version": "5.3.4", + "version": "5.3.5", "license": "MIT", "devDependencies": { "@napi-rs/cli": "^2", @@ -27,12 +27,12 @@ "typescript": "^5.9.3" }, "optionalDependencies": { - "@a3s-lab/code-darwin-arm64": "5.3.4", - "@a3s-lab/code-linux-arm64-gnu": "5.3.4", - "@a3s-lab/code-linux-arm64-musl": "5.3.4", - "@a3s-lab/code-linux-x64-gnu": "5.3.4", - "@a3s-lab/code-linux-x64-musl": "5.3.4", - "@a3s-lab/code-win32-x64-msvc": "5.3.4" + "@a3s-lab/code-darwin-arm64": "5.3.5", + "@a3s-lab/code-linux-arm64-gnu": "5.3.5", + "@a3s-lab/code-linux-arm64-musl": "5.3.5", + "@a3s-lab/code-linux-x64-gnu": "5.3.5", + "@a3s-lab/code-linux-x64-musl": "5.3.5", + "@a3s-lab/code-win32-x64-msvc": "5.3.5" } }, "node_modules/@a3s-lab/code": { diff --git a/sdk/node/package-lock.json b/sdk/node/package-lock.json index 37f83d8a..c80a86ad 100644 --- a/sdk/node/package-lock.json +++ b/sdk/node/package-lock.json @@ -1,12 +1,12 @@ { "name": "@a3s-lab/code", - "version": "5.3.4", + "version": "5.3.5", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@a3s-lab/code", - "version": "5.3.4", + "version": "5.3.5", "license": "MIT", "devDependencies": { "@napi-rs/cli": "^2", @@ -15,12 +15,12 @@ "typescript": "^5.9.3" }, "optionalDependencies": { - "@a3s-lab/code-darwin-arm64": "5.3.4", - "@a3s-lab/code-linux-arm64-gnu": "5.3.4", - "@a3s-lab/code-linux-arm64-musl": "5.3.4", - "@a3s-lab/code-linux-x64-gnu": "5.3.4", - "@a3s-lab/code-linux-x64-musl": "5.3.4", - "@a3s-lab/code-win32-x64-msvc": "5.3.4" + "@a3s-lab/code-darwin-arm64": "5.3.5", + "@a3s-lab/code-linux-arm64-gnu": "5.3.5", + "@a3s-lab/code-linux-arm64-musl": "5.3.5", + "@a3s-lab/code-linux-x64-gnu": "5.3.5", + "@a3s-lab/code-linux-x64-musl": "5.3.5", + "@a3s-lab/code-win32-x64-msvc": "5.3.5" } }, "node_modules/@a3s-lab/code-darwin-arm64": { diff --git a/sdk/node/package.json b/sdk/node/package.json index f2c45765..9f90821c 100644 --- a/sdk/node/package.json +++ b/sdk/node/package.json @@ -1,6 +1,6 @@ { "name": "@a3s-lab/code", - "version": "5.3.4", + "version": "5.3.5", "description": "A3S Code - Native Node.js bindings for the coding-agent runtime", "main": "index.js", "types": "index.d.ts", @@ -44,11 +44,11 @@ "test:helpers": "node test-helpers.mjs" }, "optionalDependencies": { - "@a3s-lab/code-darwin-arm64": "5.3.4", - "@a3s-lab/code-linux-x64-gnu": "5.3.4", - "@a3s-lab/code-linux-x64-musl": "5.3.4", - "@a3s-lab/code-linux-arm64-gnu": "5.3.4", - "@a3s-lab/code-linux-arm64-musl": "5.3.4", - "@a3s-lab/code-win32-x64-msvc": "5.3.4" + "@a3s-lab/code-darwin-arm64": "5.3.5", + "@a3s-lab/code-linux-x64-gnu": "5.3.5", + "@a3s-lab/code-linux-x64-musl": "5.3.5", + "@a3s-lab/code-linux-arm64-gnu": "5.3.5", + "@a3s-lab/code-linux-arm64-musl": "5.3.5", + "@a3s-lab/code-win32-x64-msvc": "5.3.5" } } diff --git a/sdk/python-bootstrap/README.md b/sdk/python-bootstrap/README.md index 18878261..134a320c 100644 --- a/sdk/python-bootstrap/README.md +++ b/sdk/python-bootstrap/README.md @@ -3,7 +3,7 @@ `pip install a3s-code` ships this small pure-Python package. On first `import a3s_code` it downloads the native extension matching your interpreter and platform from the project's -[GitHub Releases](https://github.com/AI45Lab/Code/releases), verifies +[GitHub Releases](https://github.com/A3S-Lab/Code/releases), verifies the wheel's sha256 against the release manifest, extracts the compiled extension into a per-user cache, and exposes the normal `a3s_code` API. @@ -41,5 +41,7 @@ wheel directly: ```bash pip install \ - https://github.com/AI45Lab/Code/releases/download/v3.2.1/a3s_code-3.2.1-cp312-cp312-manylinux_2_28_x86_64.whl + 'https://github.com/A3S-Lab/Code/releases/download/v/a3s_code--cp312-cp312-manylinux_2_28_x86_64.whl' ``` + +Replace `` with the release to install, for example `5.3.5`. diff --git a/sdk/python-bootstrap/build/lib/a3s_code/__init__.py b/sdk/python-bootstrap/build/lib/a3s_code/__init__.py deleted file mode 100644 index 07aba873..00000000 --- a/sdk/python-bootstrap/build/lib/a3s_code/__init__.py +++ /dev/null @@ -1,18 +0,0 @@ -"""A3S Code Python SDK. - -This is the pure-Python bootstrap distributed on PyPI. The actual native -extension is fetched from GitHub Releases on first import — see -`a3s_code._bootstrap` for the loader logic and the environment variables -that customize cache location / source URL. -""" - -from . import _bootstrap as _bootstrap - -_bootstrap.ensure_native_loaded() - -# Re-import after the cache dir is on sys.path. The bootstrap extracted -# `_native..so` into a per-version cache; Python's import machinery -# picks it up from there. -from ._native import * # noqa: E402,F401,F403 - -__version__ = _bootstrap.__version__ diff --git a/sdk/python-bootstrap/build/lib/a3s_code/_bootstrap.py b/sdk/python-bootstrap/build/lib/a3s_code/_bootstrap.py deleted file mode 100644 index e1ed30e7..00000000 --- a/sdk/python-bootstrap/build/lib/a3s_code/_bootstrap.py +++ /dev/null @@ -1,213 +0,0 @@ -"""Lazy loader for the a3s-code native extension. - -This module is part of the pure-Python bootstrap published to PyPI under -the `a3s-code` name. On first import it resolves the matching native -wheel for the current interpreter/platform, downloads it from the -project's GitHub Releases, verifies the wheel's sha256 against the -release manifest, extracts the compiled `_native` extension into a -per-user cache, and prepends the cache to `sys.path` so the rest of -`a3s_code/__init__.py` can `from ._native import *` normally. - -Override the cache location via `A3S_CODE_CACHE_DIR`. Override the -release source via `A3S_CODE_RELEASES_BASE_URL` (default points at the -GitHub Releases page for `AI45Lab/Code`). Skip the integrity check via -`A3S_CODE_SKIP_HASH_CHECK=1` (not recommended outside of CI). -""" - -from __future__ import annotations - -import hashlib -import io -import json -import os -import platform -import sys -import threading -import urllib.error -import urllib.request -import zipfile -from pathlib import Path -from typing import Optional - -# Version is the bootstrap's own version, which equals the matching native -# wheel version on GH Releases. Bumped by the release workflow. -__version__ = "3.2.1" - -_DEFAULT_BASE_URL = "https://github.com/AI45Lab/Code/releases/download" -_REQUEST_TIMEOUT_S = 120 -_USER_AGENT = f"a3s-code-bootstrap/{__version__}" -_LOAD_LOCK = threading.Lock() -_LOADED = False - - -class BootstrapError(RuntimeError): - """Raised when the native extension cannot be located, downloaded, or verified.""" - - -def _base_url() -> str: - return os.environ.get("A3S_CODE_RELEASES_BASE_URL", _DEFAULT_BASE_URL).rstrip("/") - - -def _cache_root() -> Path: - override = os.environ.get("A3S_CODE_CACHE_DIR") - if override: - return Path(override).expanduser() / __version__ - base = os.environ.get("XDG_CACHE_HOME") or str(Path.home() / ".cache") - return Path(base) / "a3s-code" / __version__ - - -def _platform_tag() -> str: - """Return the wheel platform tag for the current interpreter/platform. - - Mirrors the matrix produced by `.github/workflows/publish-python.yml`. - Raises `BootstrapError` for unsupported combinations so callers can - surface a clear install hint. - """ - sys_plat = sys.platform - machine = platform.machine().lower() - if sys_plat == "darwin": - if machine in ("arm64", "aarch64"): - return "macosx_11_0_arm64" - elif sys_plat == "linux": - if machine in ("x86_64", "amd64"): - return "manylinux_2_28_x86_64" - elif sys_plat == "win32": - if machine in ("amd64", "x86_64"): - return "win_amd64" - raise BootstrapError( - f"a3s-code: no native wheel published for {sys_plat}/{machine}. " - "Supported platforms: macOS arm64, Linux x86_64 (glibc 2.28+), Windows x86_64." - ) - - -def _wheel_filename(version: str = __version__) -> str: - py_tag = f"cp{sys.version_info.major}{sys.version_info.minor}" - return f"a3s_code-{version}-{py_tag}-{py_tag}-{_platform_tag()}.whl" - - -def _release_url(filename: str, version: str = __version__) -> str: - return f"{_base_url()}/v{version}/{filename}" - - -def _http_get(url: str) -> bytes: - req = urllib.request.Request(url, headers={"User-Agent": _USER_AGENT}) - try: - with urllib.request.urlopen(req, timeout=_REQUEST_TIMEOUT_S) as resp: - return resp.read() - except urllib.error.HTTPError as exc: - raise BootstrapError(f"GET {url} failed: HTTP {exc.code}") from exc - except urllib.error.URLError as exc: - raise BootstrapError(f"GET {url} failed: {exc.reason}") from exc - - -def _expected_sha256(wheel_name: str, version: str = __version__) -> Optional[str]: - """Look up the published sha256 for `wheel_name` in the release manifest. - - Returns `None` if the manifest is unreachable — bootstrap will then - proceed without integrity verification but emit a warning. Override - with `A3S_CODE_SKIP_HASH_CHECK=1` for hermetic offline mirrors. - """ - manifest_url = f"{_base_url()}/v{version}/python-native-manifest.json" - try: - data = json.loads(_http_get(manifest_url)) - except BootstrapError as exc: - sys.stderr.write( - f"a3s-code: warning: manifest fetch failed ({exc}); skipping hash check\n" - ) - return None - for asset in data.get("assets", []): - if asset.get("filename") == wheel_name: - return asset.get("sha256") - return None - - -def _extract_native(wheel_bytes: bytes, target_dir: Path) -> Path: - """Extract the compiled `_native.*` extension from `wheel_bytes` into - `target_dir`. Returns the path to the extracted file. - """ - target_dir.mkdir(parents=True, exist_ok=True) - with zipfile.ZipFile(io.BytesIO(wheel_bytes)) as zf: - for name in zf.namelist(): - base = Path(name).name - # match _native..{so,pyd,dylib} - if base.startswith("_native.") and not base.endswith(".dist-info"): - out_path = target_dir / base - with zf.open(name) as src, out_path.open("wb") as dst: - dst.write(src.read()) - return out_path - raise BootstrapError( - "downloaded wheel did not contain a _native extension; " - "the release artifact appears to be corrupt" - ) - - -def _find_cached_native(cache_dir: Path) -> Optional[Path]: - if not cache_dir.is_dir(): - return None - for child in cache_dir.iterdir(): - if child.is_file() and child.name.startswith("_native."): - return child - return None - - -def _register_native(native_path: Path) -> None: - """Load `native_path` as the `a3s_code._native` module. - - `_native` is a compiled extension, not a regular Python file, so use - `importlib.machinery.ExtensionFileLoader` + the matching spec instead - of plain `spec_from_file_location` (which works but doesn't set the - right loader for extensions on all Python versions). - """ - import importlib.machinery - import importlib.util - - fullname = "a3s_code._native" - loader = importlib.machinery.ExtensionFileLoader(fullname, str(native_path)) - spec = importlib.util.spec_from_loader(fullname, loader, origin=str(native_path)) - if spec is None: - raise BootstrapError(f"failed to build import spec for {native_path}") - module = importlib.util.module_from_spec(spec) - sys.modules[fullname] = module - spec.loader.exec_module(module) - - -def ensure_native_loaded(version: str = __version__) -> Path: - """Idempotently ensure the `_native` extension is registered as - `a3s_code._native` in `sys.modules`. Returns the cache directory the - extension was loaded from. Safe across threads — first caller wins. - """ - global _LOADED - cache = _cache_root() - - if _LOADED: - return cache - - with _LOAD_LOCK: - if _LOADED: - return cache - - native = _find_cached_native(cache) - if native is None: - wheel_name = _wheel_filename(version) - url = _release_url(wheel_name, version) - sys.stderr.write( - f"a3s-code: fetching native wheel {wheel_name} " - f"from {url} (first import only)...\n" - ) - wheel_bytes = _http_get(url) - - if os.environ.get("A3S_CODE_SKIP_HASH_CHECK") != "1": - expected = _expected_sha256(wheel_name, version) - if expected is not None: - actual = hashlib.sha256(wheel_bytes).hexdigest() - if actual != expected: - raise BootstrapError( - f"sha256 mismatch for {wheel_name}: " - f"expected {expected}, got {actual}" - ) - - native = _extract_native(wheel_bytes, cache) - - _register_native(native) - _LOADED = True - return cache diff --git a/sdk/python-bootstrap/build/lib/a3s_code/py.typed b/sdk/python-bootstrap/build/lib/a3s_code/py.typed deleted file mode 100644 index e69de29b..00000000 diff --git a/sdk/python-bootstrap/pyproject.toml b/sdk/python-bootstrap/pyproject.toml index a48f0b41..2b946b0d 100644 --- a/sdk/python-bootstrap/pyproject.toml +++ b/sdk/python-bootstrap/pyproject.toml @@ -5,9 +5,9 @@ build-backend = "setuptools.build_meta" [project] name = "a3s-code" # Keep in sync with crates/code core release. The bootstrap loader fetches -# the matching native wheel from `https://github.com/AI45Lab/Code/releases/tag/v` +# the matching native wheel from `https://github.com/A3S-Lab/Code/releases/tag/v` # at import time. -version = "5.3.4" +version = "5.3.5" description = "A3S Code Python SDK — pure-Python bootstrap that fetches the native wheel from GitHub Releases" readme = "README.md" license = {text = "MIT"} @@ -25,9 +25,9 @@ classifiers = [ dependencies = [] [project.urls] -Homepage = "https://github.com/AI45Lab/Code" -Repository = "https://github.com/AI45Lab/Code" -Issues = "https://github.com/AI45Lab/Code/issues" +Homepage = "https://github.com/A3S-Lab/Code" +Repository = "https://github.com/A3S-Lab/Code" +Issues = "https://github.com/A3S-Lab/Code/issues" [tool.setuptools.packages.find] where = ["src"] diff --git a/sdk/python-bootstrap/src/a3s_code/_bootstrap.py b/sdk/python-bootstrap/src/a3s_code/_bootstrap.py index a4d57713..c36106d2 100644 --- a/sdk/python-bootstrap/src/a3s_code/_bootstrap.py +++ b/sdk/python-bootstrap/src/a3s_code/_bootstrap.py @@ -31,7 +31,7 @@ # Version is the bootstrap's own version, which equals the matching native # wheel version on GH Releases. Bumped by the release workflow. -__version__ = "5.3.4" +__version__ = "5.3.5" _DEFAULT_BASE_URL = "https://github.com/A3S-Lab/Code/releases/download" _REQUEST_TIMEOUT_S = 120 diff --git a/sdk/python/CHANGELOG.md b/sdk/python/CHANGELOG.md index 5189f11c..98f9c834 100644 --- a/sdk/python/CHANGELOG.md +++ b/sdk/python/CHANGELOG.md @@ -4,6 +4,14 @@ All notable changes to the A3S Code Python SDK will be documented in this file. ## [Unreleased] +## [5.3.5] - 2026-07-17 + +### Changed + +- Updated the bundled Core with host-supplied session client bootstrap, bounded + in-place response-stream retry and rollback, explicit search-routing + metadata, and the Bing RSS engine. + ## [5.3.4] - 2026-07-16 ### Fixed diff --git a/sdk/python/Cargo.lock b/sdk/python/Cargo.lock index 21e5a5e4..bb871803 100644 --- a/sdk/python/Cargo.lock +++ b/sdk/python/Cargo.lock @@ -15,7 +15,7 @@ source = "git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ac [[package]] name = "a3s-code-core" -version = "5.3.4" +version = "5.3.5" dependencies = [ "a3s-acl 0.2.1 (git+https://github.com/A3S-Lab/ACL.git?rev=6e2a6469edc0f4c61b1e588d0ace873aaf15ce22)", "a3s-common", @@ -76,7 +76,7 @@ dependencies = [ [[package]] name = "a3s-code-py" -version = "5.3.4" +version = "5.3.5" dependencies = [ "a3s-code-core", "anyhow", diff --git a/sdk/python/Cargo.toml b/sdk/python/Cargo.toml index eb3b424b..56f7a148 100644 --- a/sdk/python/Cargo.toml +++ b/sdk/python/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "a3s-code-py" -version = "5.3.4" +version = "5.3.5" edition = "2021" authors = ["A3S Lab Team"] license = "MIT" @@ -12,7 +12,7 @@ name = "a3s_code" crate-type = ["cdylib"] [dependencies] -a3s-code-core = { version = "5.3.4", path = "../../core", features = ["s3", "serve"] } +a3s-code-core = { version = "5.3.5", path = "../../core", features = ["s3", "serve"] } pyo3 = { version = "0.23", features = ["multiple-pymethods"] } tokio = { version = "1.35", features = ["full"] } serde_json = "1.0" diff --git a/sdk/python/README.md b/sdk/python/README.md index f7bbd357..39f2bb09 100644 --- a/sdk/python/README.md +++ b/sdk/python/README.md @@ -10,7 +10,7 @@ pip install a3s-code From v3.2.1 onwards the PyPI `a3s-code` package is a small pure-Python bootstrap. On first `import a3s_code` it downloads the matching native -wheel from [GitHub Releases](https://github.com/AI45Lab/Code/releases), +wheel from [GitHub Releases](https://github.com/A3S-Lab/Code/releases), verifies the wheel's sha256 against the release manifest, and caches the compiled extension under `~/.cache/a3s-code//`. Subsequent imports use the cache. @@ -25,9 +25,11 @@ interpreter directly: ```bash pip install \ - https://github.com/AI45Lab/Code/releases/download/v3.2.1/a3s_code-3.2.1-cp312-cp312-manylinux_2_28_x86_64.whl + 'https://github.com/A3S-Lab/Code/releases/download/v/a3s_code--cp312-cp312-manylinux_2_28_x86_64.whl' ``` +Replace `` with the release to install, for example `5.3.5`. + ## Quick Start ```python diff --git a/sdk/python/pyproject.toml b/sdk/python/pyproject.toml index f5c57ec2..0e45ae35 100644 --- a/sdk/python/pyproject.toml +++ b/sdk/python/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "a3s-code" -version = "5.3.4" +version = "5.3.5" description = "A3S Code - Native Python bindings for the coding-agent runtime" readme = "README.md" license = {text = "MIT"} From d085ce13d2f229d4a2d19e2680e44a9980bc47c1 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 17:03:59 +0800 Subject: [PATCH 8/9] fix(llm): surface interrupted provider streams --- CHANGELOG.md | 4 + core/src/llm/openai/streaming.rs | 15 ++- core/src/llm/openai/tests.rs | 154 +++++++++++++++++++++++++++++ sdk/node/README.md | 18 +++- sdk/python/README.md | 29 +++++- sdk/python/docs/QUICK_REFERENCE.md | 15 ++- sdk/python/docs/README.md | 12 ++- 7 files changed, 236 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 503efb97..e01a17aa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Retried an established response stream up to ten times within the same turn, with cancellation-aware exponential backoff and transactional rollback of provisional text, reasoning, and tool drafts between attempts. +- Treated OpenAI-compatible transport failures and cancellation before terminal + evidence as incomplete streams instead of synthesizing a successful partial + response; a received finish reason remains valid without a trailing + `[DONE]` marker. ## [5.3.4] - 2026-07-16 diff --git a/core/src/llm/openai/streaming.rs b/core/src/llm/openai/streaming.rs index 04a98901..8b49a35f 100644 --- a/core/src/llm/openai/streaming.rs +++ b/core/src/llm/openai/streaming.rs @@ -107,12 +107,13 @@ impl OpenAiClient { let mut first_token_ms = None; let mut saw_done = false; let mut parsed_any_event = false; + let mut stream_failed = false; loop { let chunk_result = tokio::select! { biased; - _ = stream_cancellation.cancelled() => break, - _ = tx.closed() => break, + _ = stream_cancellation.cancelled() => return, + _ = tx.closed() => return, chunk = stream.next() => match chunk { Some(chunk) => chunk, None => break, @@ -122,6 +123,7 @@ impl OpenAiClient { Ok(c) => c, Err(e) => { tracing::error!("Stream error: {}", e); + stream_failed = true; break; } }; @@ -562,6 +564,15 @@ impl OpenAiClient { } } + if stream_failed && finish_reason.is_none() { + tracing::warn!( + provider = %provider_name, + model = %request_model, + "OpenAI-compatible stream failed before terminal evidence; closing without Done" + ); + return; + } + if parsed_any_event || !text_content.is_empty() || !tool_calls.is_empty() diff --git a/core/src/llm/openai/tests.rs b/core/src/llm/openai/tests.rs index 2fee600c..bd7cb047 100644 --- a/core/src/llm/openai/tests.rs +++ b/core/src/llm/openai/tests.rs @@ -1,5 +1,6 @@ use super::*; use crate::llm::types::{Message, ToolDefinition}; +use futures::StreamExt; fn make_client() -> OpenAiClient { OpenAiClient::new("test-key".to_string(), "gpt-test".to_string()) @@ -17,6 +18,12 @@ struct MockSseHttp { struct PendingSseHttp; +struct FailingSseHttp { + chunks: Vec, +} + +struct PartialThenPendingSseHttp; + #[async_trait::async_trait] impl crate::llm::http::HttpClient for MockSseHttp { async fn post( @@ -78,6 +85,75 @@ impl crate::llm::http::HttpClient for PendingSseHttp { } } +#[async_trait::async_trait] +impl crate::llm::http::HttpClient for FailingSseHttp { + async fn post( + &self, + _url: &str, + _headers: Vec<(&str, &str)>, + _body: &serde_json::Value, + _cancel: tokio_util::sync::CancellationToken, + ) -> anyhow::Result { + anyhow::bail!("post is unused in the interrupted streaming test") + } + + async fn post_streaming( + &self, + _url: &str, + _headers: Vec<(&str, &str)>, + _body: &serde_json::Value, + _cancel: tokio_util::sync::CancellationToken, + ) -> anyhow::Result { + let mut items = self + .chunks + .iter() + .map(|chunk| Ok(bytes::Bytes::from(chunk.clone()))) + .collect::>>(); + items.push(Err(anyhow::anyhow!("connection reset"))); + Ok(crate::llm::http::StreamingHttpResponse { + status: 200, + retry_after: None, + byte_stream: Box::pin(futures::stream::iter(items)), + error_body: String::new(), + }) + } +} + +#[async_trait::async_trait] +impl crate::llm::http::HttpClient for PartialThenPendingSseHttp { + async fn post( + &self, + _url: &str, + _headers: Vec<(&str, &str)>, + _body: &serde_json::Value, + _cancel: tokio_util::sync::CancellationToken, + ) -> anyhow::Result { + anyhow::bail!("post is unused in the partial cancellation test") + } + + async fn post_streaming( + &self, + _url: &str, + _headers: Vec<(&str, &str)>, + _body: &serde_json::Value, + _cancel: tokio_util::sync::CancellationToken, + ) -> anyhow::Result { + let partial = Ok(bytes::Bytes::from_static( + br#"data: {"choices":[{"delta":{"content":"partial"}}]} + +"#, + )); + Ok(crate::llm::http::StreamingHttpResponse { + status: 200, + retry_after: None, + byte_stream: Box::pin( + futures::stream::iter(vec![partial]).chain(futures::stream::pending()), + ), + error_body: String::new(), + }) + } +} + fn glm_client(chunks: Vec) -> OpenAiClient { OpenAiClient::new("k".to_string(), "glm-test".to_string()) .with_http_client(std::sync::Arc::new(MockSseHttp { chunks })) @@ -123,6 +199,84 @@ async fn streaming_parser_closes_when_caller_cancels() { assert!(next.is_none()); } +#[tokio::test] +async fn streaming_transport_error_after_partial_delta_does_not_emit_done() { + use crate::llm::{LlmClient, StreamEvent}; + + let client = OpenAiClient::new("k".to_string(), "model".to_string()).with_http_client( + std::sync::Arc::new(FailingSseHttp { + chunks: vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n".to_string(), + ], + }), + ); + let mut rx = client + .complete_streaming( + &[Message::user("go")], + None, + &[], + tokio_util::sync::CancellationToken::new(), + ) + .await + .expect("stream opened"); + + let mut text = String::new(); + let mut saw_done = false; + while let Some(event) = rx.recv().await { + match event { + StreamEvent::TextDelta(delta) => text.push_str(&delta), + StreamEvent::Done(_) => saw_done = true, + _ => {} + } + } + + assert_eq!(text, "partial"); + assert!( + !saw_done, + "a failed transport must close without Done so the agent retries the turn" + ); +} + +#[tokio::test] +async fn streaming_transport_error_after_finish_reason_can_finalize() { + let client = OpenAiClient::new("k".to_string(), "model".to_string()).with_http_client( + std::sync::Arc::new(FailingSseHttp { + chunks: vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"complete\"}}]}\n\n".to_string(), + "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n".to_string(), + ], + }), + ); + + let response = drain_to_done(&client).await; + assert_eq!(response.text(), "complete"); + assert_eq!(response.stop_reason.as_deref(), Some("stop")); +} + +#[tokio::test] +async fn streaming_partial_response_is_not_finalized_after_cancellation() { + use crate::llm::{LlmClient, StreamEvent}; + + let client = OpenAiClient::new("k".to_string(), "model".to_string()) + .with_http_client(std::sync::Arc::new(PartialThenPendingSseHttp)); + let cancellation = tokio_util::sync::CancellationToken::new(); + let mut rx = client + .complete_streaming(&[Message::user("go")], None, &[], cancellation.clone()) + .await + .expect("stream opened"); + + assert!(matches!( + rx.recv().await, + Some(StreamEvent::TextDelta(text)) if text == "partial" + )); + cancellation.cancel(); + + let next = tokio::time::timeout(std::time::Duration::from_millis(100), rx.recv()) + .await + .expect("provider parser must stop after cancellation"); + assert!(next.is_none(), "cancellation must not synthesize Done"); +} + #[tokio::test] async fn streaming_reasoning_does_not_leak_into_content_and_keeps_tool_call() { let chunks = vec![ diff --git a/sdk/node/README.md b/sdk/node/README.md index 8a4671e1..fe8ba3f9 100644 --- a/sdk/node/README.md +++ b/sdk/node/README.md @@ -76,16 +76,32 @@ switching on `type`: future event names remain intact instead of becoming ```js const stream = await session.stream('Explain the current test failures') +let activeTurn +let attemptText = '' while (true) { const { value: event, done } = await stream.next() if (!event) break - if (event.type === 'text_delta') process.stdout.write(event.text ?? '') + if (event.type === 'turn_start') { + // The same turn number means the provider stream was retried. Discard + // provisional output from that attempt before accepting replacement deltas. + activeTurn = event.turn + attemptText = '' + } else if (event.type === 'text_delta') { + attemptText += event.text ?? '' + } else if (event.type === 'turn_end' && event.turn === activeTurn) { + process.stdout.write(attemptText) + attemptText = '' + } if (event.type === 'agent_end') console.log(event.verificationSummaryText ?? '') console.debug(event.version, event.type, event.payload, event.metadata) if (done) break } ``` +`turn_start` may repeat with the same `turn` when an established provider +stream is interrupted. Treat each turn as provisional until `turn_end`; reset +text, reasoning, and tool-call drafts when that turn restarts. + `agentEventTypesV1()` returns the ordered catalog known by this build, while the exported `AgentEventTypeV1` TypeScript type deliberately remains open for forward compatibility. diff --git a/sdk/python/README.md b/sdk/python/README.md index 39f2bb09..e849a437 100644 --- a/sdk/python/README.md +++ b/sdk/python/README.md @@ -106,14 +106,27 @@ and `exit_code` are derived from the same core projection. ```python from a3s_code import EventType +active_turn = None +attempt_text = [] for event in session.stream("Explain the current test failures"): - if event.type == EventType.TEXT_DELTA: - print(event.text or "", end="") + if event.type == EventType.TURN_START: + # A repeated turn number replaces an interrupted attempt. + active_turn = event.turn + attempt_text.clear() + elif event.type == EventType.TEXT_DELTA: + attempt_text.append(event.text or "") + elif event.type == EventType.TURN_END and event.turn == active_turn: + print("".join(attempt_text), end="", flush=True) + attempt_text.clear() elif event.type == EventType.AGENT_END: print(event.verification_summary_text or "") print(event.version, event.type, event.payload, event.metadata) ``` +`turn_start` may repeat with the same `turn` when an established provider +stream is interrupted. Treat each turn as provisional until `turn_end`; reset +text, reasoning, and tool-call drafts when that turn restarts. + `agent_event_types_v1()` returns the ordered catalog known by the native runtime. `AgentEventTypeV1` and `EventType` are generated from the core catalog; callers should still retain a default branch for future values. @@ -170,9 +183,17 @@ session = agent.session("/my-project", # Send / Stream result = session.send({"prompt": "Explain the auth module"}) +active_turn = None +attempt_text = [] for event in session.stream({"prompt": "Refactor auth"}): - if event.event_type == "text_delta": - print(event.text, end="", flush=True) + if event.event_type == "turn_start": + active_turn = event.turn + attempt_text.clear() + elif event.event_type == "text_delta": + attempt_text.append(event.text or "") + elif event.event_type == "turn_end" and event.turn == active_turn: + print("".join(attempt_text), end="", flush=True) + attempt_text.clear() # Streams with no custom history update session history and verification evidence # when the stream completes. Passing explicit history keeps the stream isolated. diff --git a/sdk/python/docs/QUICK_REFERENCE.md b/sdk/python/docs/QUICK_REFERENCE.md index 6643b36a..6929550a 100644 --- a/sdk/python/docs/QUICK_REFERENCE.md +++ b/sdk/python/docs/QUICK_REFERENCE.md @@ -13,14 +13,25 @@ session = agent.session(".") result = session.send({"prompt": "Find the code path that handles authentication."}) print(result.text) +active_turn = None +attempt_text = [] for event in session.stream({"prompt": "Refactor the tests around that code path."}): - if event.event_type == "text_delta": - print(event.text, end="", flush=True) + if event.event_type == "turn_start": + active_turn = event.turn + attempt_text.clear() + elif event.event_type == "text_delta": + attempt_text.append(event.text or "") + elif event.event_type == "turn_end" and event.turn == active_turn: + print("".join(attempt_text), end="", flush=True) + attempt_text.clear() ``` Streaming events use envelope version `1`: `event.type` is the canonical open discriminant, `event.payload` is lossless, and `event.metadata` preserves optional protocol metadata. `event.event_type` is a compatibility alias. +Repeated `turn_start` events with the same turn replace an interrupted stream +attempt, so clear all provisional text, reasoning, and tool drafts for that +turn before applying the replacement deltas. Use the short object-shaped request APIs first. They own normal model execution, built-in tools, diff --git a/sdk/python/docs/README.md b/sdk/python/docs/README.md index 8bc3b395..f7d6ffcb 100644 --- a/sdk/python/docs/README.md +++ b/sdk/python/docs/README.md @@ -119,9 +119,17 @@ session = agent.session("/my-project", # Send / Stream result = session.send({"prompt": "Explain the auth module"}) +active_turn = None +attempt_text = [] for event in session.stream({"prompt": "Refactor auth"}): - if event.event_type == "text_delta": - print(event.text, end="", flush=True) + if event.event_type == "turn_start": + active_turn = event.turn + attempt_text.clear() + elif event.event_type == "text_delta": + attempt_text.append(event.text or "") + elif event.event_type == "turn_end" and event.turn == active_turn: + print("".join(attempt_text), end="", flush=True) + attempt_text.clear() # Planning events # Prefer planning_mode="auto" | "enabled" | "disabled". The legacy planning From e93f01f8ab35e7c22f324ef8eec2a8b05b96c2c8 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 17:18:41 +0800 Subject: [PATCH 9/9] fix(llm): require terminal evidence for streams --- core/src/llm/openai/streaming.rs | 44 +++++------ core/src/llm/openai/tests.rs | 121 +++++++++++++++++++++++++++---- 2 files changed, 131 insertions(+), 34 deletions(-) diff --git a/core/src/llm/openai/streaming.rs b/core/src/llm/openai/streaming.rs index 8b49a35f..8a5990cf 100644 --- a/core/src/llm/openai/streaming.rs +++ b/core/src/llm/openai/streaming.rs @@ -105,7 +105,6 @@ impl OpenAiClient { let mut response_model = None; let mut response_object = None; let mut first_token_ms = None; - let mut saw_done = false; let mut parsed_any_event = false; let mut stream_failed = false; @@ -137,7 +136,6 @@ impl OpenAiClient { for line in event_data.lines() { if let Some(data) = crate::sse::data_field_value(line) { if data == "[DONE]" { - saw_done = true; if !text_content.is_empty() { content_blocks.push(ContentBlock::Text { text: text_content.clone(), @@ -185,7 +183,7 @@ impl OpenAiClient { }), }; let _ = tx.send(StreamEvent::Done(response)).await; - continue; + return; } if let Ok(event) = serde_json::from_str::(data) { @@ -415,10 +413,6 @@ impl OpenAiClient { } } - if saw_done { - return; - } - let trailing = buffer.trim(); if !trailing.is_empty() { if let Ok(event) = serde_json::from_str::(trailing) { @@ -564,12 +558,27 @@ impl OpenAiClient { } } - if stream_failed && finish_reason.is_none() { - tracing::warn!( - provider = %provider_name, - model = %request_model, - "OpenAI-compatible stream failed before terminal evidence; closing without Done" - ); + if finish_reason.is_none() { + if stream_failed { + tracing::warn!( + provider = %provider_name, + model = %request_model, + "OpenAI-compatible stream failed before terminal evidence; closing without Done" + ); + } else if parsed_any_event { + tracing::warn!( + provider = %provider_name, + model = %request_model, + "OpenAI-compatible stream reached EOF before terminal evidence; closing without Done" + ); + } else { + tracing::warn!( + provider = %provider_name, + model = %request_model, + trailing = %trailing.chars().take(400).collect::(), + "OpenAI-compatible stream ended without any parseable events" + ); + } return; } @@ -581,7 +590,7 @@ impl OpenAiClient { tracing::warn!( provider = %provider_name, model = %request_model, - "OpenAI-compatible stream ended without [DONE]; finalizing buffered response" + "OpenAI-compatible stream ended without [DONE] after finish reason; finalizing buffered response" ); if !text_content.is_empty() { content_blocks.push(ContentBlock::Text { @@ -627,13 +636,6 @@ impl OpenAiClient { }), }; let _ = tx.send(StreamEvent::Done(response)).await; - } else { - tracing::warn!( - provider = %provider_name, - model = %request_model, - trailing = %trailing.chars().take(400).collect::(), - "OpenAI-compatible stream ended without any parseable events" - ); } }); diff --git a/core/src/llm/openai/tests.rs b/core/src/llm/openai/tests.rs index bd7cb047..9640168f 100644 --- a/core/src/llm/openai/tests.rs +++ b/core/src/llm/openai/tests.rs @@ -22,7 +22,9 @@ struct FailingSseHttp { chunks: Vec, } -struct PartialThenPendingSseHttp; +struct ChunksThenPendingSseHttp { + chunks: Vec, +} #[async_trait::async_trait] impl crate::llm::http::HttpClient for MockSseHttp { @@ -120,7 +122,7 @@ impl crate::llm::http::HttpClient for FailingSseHttp { } #[async_trait::async_trait] -impl crate::llm::http::HttpClient for PartialThenPendingSseHttp { +impl crate::llm::http::HttpClient for ChunksThenPendingSseHttp { async fn post( &self, _url: &str, @@ -128,7 +130,7 @@ impl crate::llm::http::HttpClient for PartialThenPendingSseHttp { _body: &serde_json::Value, _cancel: tokio_util::sync::CancellationToken, ) -> anyhow::Result { - anyhow::bail!("post is unused in the partial cancellation test") + anyhow::bail!("post is unused in the pending streaming tests") } async fn post_streaming( @@ -138,17 +140,15 @@ impl crate::llm::http::HttpClient for PartialThenPendingSseHttp { _body: &serde_json::Value, _cancel: tokio_util::sync::CancellationToken, ) -> anyhow::Result { - let partial = Ok(bytes::Bytes::from_static( - br#"data: {"choices":[{"delta":{"content":"partial"}}]} - -"#, - )); + let items = self + .chunks + .iter() + .map(|chunk| Ok(bytes::Bytes::from(chunk.clone()))) + .collect::>>(); Ok(crate::llm::http::StreamingHttpResponse { status: 200, retry_after: None, - byte_stream: Box::pin( - futures::stream::iter(vec![partial]).chain(futures::stream::pending()), - ), + byte_stream: Box::pin(futures::stream::iter(items).chain(futures::stream::pending())), error_body: String::new(), }) } @@ -237,6 +237,40 @@ async fn streaming_transport_error_after_partial_delta_does_not_emit_done() { ); } +#[tokio::test] +async fn streaming_clean_eof_after_partial_delta_does_not_emit_done() { + use crate::llm::{LlmClient, StreamEvent}; + + let client = glm_client(vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n".to_string(), + ]); + let mut rx = client + .complete_streaming( + &[Message::user("go")], + None, + &[], + tokio_util::sync::CancellationToken::new(), + ) + .await + .expect("stream opened"); + + let mut text = String::new(); + let mut saw_done = false; + while let Some(event) = rx.recv().await { + match event { + StreamEvent::TextDelta(delta) => text.push_str(&delta), + StreamEvent::Done(_) => saw_done = true, + _ => {} + } + } + + assert_eq!(text, "partial"); + assert!( + !saw_done, + "EOF without protocol terminal evidence must close without Done" + ); +} + #[tokio::test] async fn streaming_transport_error_after_finish_reason_can_finalize() { let client = OpenAiClient::new("k".to_string(), "model".to_string()).with_http_client( @@ -253,12 +287,29 @@ async fn streaming_transport_error_after_finish_reason_can_finalize() { assert_eq!(response.stop_reason.as_deref(), Some("stop")); } +#[tokio::test] +async fn streaming_clean_eof_after_finish_reason_can_finalize() { + let client = glm_client(vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"complete\"}}]}\n\n".to_string(), + "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n".to_string(), + ]); + + let response = drain_to_done(&client).await; + assert_eq!(response.text(), "complete"); + assert_eq!(response.stop_reason.as_deref(), Some("stop")); +} + #[tokio::test] async fn streaming_partial_response_is_not_finalized_after_cancellation() { use crate::llm::{LlmClient, StreamEvent}; - let client = OpenAiClient::new("k".to_string(), "model".to_string()) - .with_http_client(std::sync::Arc::new(PartialThenPendingSseHttp)); + let client = OpenAiClient::new("k".to_string(), "model".to_string()).with_http_client( + std::sync::Arc::new(ChunksThenPendingSseHttp { + chunks: vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n".to_string(), + ], + }), + ); let cancellation = tokio_util::sync::CancellationToken::new(); let mut rx = client .complete_streaming(&[Message::user("go")], None, &[], cancellation.clone()) @@ -277,6 +328,50 @@ async fn streaming_partial_response_is_not_finalized_after_cancellation() { assert!(next.is_none(), "cancellation must not synthesize Done"); } +#[tokio::test] +async fn streaming_done_closes_before_pending_transport_and_emits_once() { + use crate::llm::{LlmClient, StreamEvent}; + + let client = OpenAiClient::new("k".to_string(), "model".to_string()).with_http_client( + std::sync::Arc::new(ChunksThenPendingSseHttp { + chunks: vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"complete\"}}]}\n\n".to_string(), + "data: [DONE]\n\n".to_string(), + "data: [DONE]\n\n".to_string(), + ], + }), + ); + let mut rx = client + .complete_streaming( + &[Message::user("go")], + None, + &[], + tokio_util::sync::CancellationToken::new(), + ) + .await + .expect("stream opened"); + + let events = tokio::time::timeout(std::time::Duration::from_secs(1), async move { + let mut events = Vec::new(); + while let Some(event) = rx.recv().await { + events.push(event); + } + events + }) + .await + .expect("[DONE] must close the parser without waiting for transport EOF"); + let done = events + .into_iter() + .filter_map(|event| match event { + StreamEvent::Done(response) => Some(response), + _ => None, + }) + .collect::>(); + + assert_eq!(done.len(), 1, "[DONE] must emit exactly one final response"); + assert_eq!(done[0].text(), "complete"); +} + #[tokio::test] async fn streaming_reasoning_does_not_leak_into_content_and_keeps_tool_call() { let chunks = vec![