Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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: |
Expand Down
27 changes: 27 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,33 @@ 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.
- 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

### Fixed
Expand Down
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 11 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -452,7 +461,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
Expand Down
4 changes: 3 additions & 1 deletion core/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -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"
Expand Down
4 changes: 3 additions & 1 deletion core/src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 },

Expand Down
212 changes: 212 additions & 0 deletions core/src/agent/extra_agent_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<Vec<String>>,
}

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<LlmResponse> {
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<mpsc::Receiver<StreamEvent>> {
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<AgentEvent>,
) -> (Vec<AgentEvent>, 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]
Expand Down
Loading
Loading