use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use super::outcome::format_child_result;
use crate::focus::ProgressSummarizer;
use crate::multi_agent::mailbox::{MailboxHub, MailboxResult, MailboxStatus};
use crate::multi_agent::path::AgentPath;
use crate::multi_agent::registry::AgentRegistry;
#[derive(Clone, Debug)]
pub struct ChildReport {
pub agent_path: String,
pub status: String,
pub result: Option<String>,
pub message: String,
}
#[derive(Clone, Debug)]
pub enum ChildResultEvent {
Progress {
agent_path: String,
status: String,
summary: Option<String>,
},
Batch {
reports: Vec<ChildReport>,
},
}
impl ChildResultEvent {
pub fn agent_path(&self) -> &str {
match self {
Self::Progress { agent_path, .. } => agent_path,
Self::Batch { reports } => reports
.first()
.map(|r| r.agent_path.as_str())
.unwrap_or_default(),
}
}
}
const MAX_ERROR_REASON_CHARS: usize = 160;
fn strip_ansi(text: &str) -> String {
let mut out = String::with_capacity(text.len());
let mut chars = text.chars().peekable();
while let Some(c) = chars.next() {
if c == '\x1b' && chars.peek() == Some(&'[') {
chars.next();
while let Some(&n) = chars.peek() {
chars.next();
if ('@'..='~').contains(&n) {
break;
}
}
} else {
out.push(c);
}
}
out
}
fn error_reason(result_text: Option<&str>) -> Option<String> {
let line = result_text
.map(strip_ansi)
.as_deref()?
.lines()
.map(|l| {
l.chars()
.map(|c| if c.is_control() { ' ' } else { c })
.collect::<String>()
})
.map(|l| l.trim().to_string())
.find(|l| !l.is_empty())?;
let mut reason: String = line.chars().take(MAX_ERROR_REASON_CHARS).collect();
if line.chars().nth(MAX_ERROR_REASON_CHARS).is_some() {
reason.push('…');
}
Some(reason)
}
pub(crate) const HELD_BATCH_REAP_AFTER: Duration = Duration::from_secs(90);
pub fn spawn_watcher(
mailbox: Arc<MailboxHub>,
registry: Arc<Mutex<AgentRegistry>>,
summarizer: Option<Arc<ProgressSummarizer>>,
child_result_tx: Option<mpsc::UnboundedSender<ChildResultEvent>>,
cancel: CancellationToken,
held_reap_after: Duration,
) -> tokio::task::JoinHandle<()> {
let mut seq_rx = mailbox.subscribe_seq();
tokio::spawn(async move {
let mut batch: Vec<MailboxResult> = Vec::new();
let mut batch_first_at: Option<Instant> = None;
let mut reap_tick = tokio::time::interval(held_reap_after);
reap_tick.tick().await;
loop {
tokio::select! {
biased;
_ = cancel.cancelled() => {
force_handover(&mut batch, ®istry, &child_result_tx);
break;
}
_ = reap_tick.tick() => {
let held = batch_first_at.is_some_and(|t| t.elapsed() >= held_reap_after);
if !held {
continue;
}
let any_working = registry
.lock()
.unwrap()
.list()
.iter()
.any(|e| !e.closing && (e.in_flight || e.queue_len > 0));
if any_working {
continue;
}
tracing::info!(
batch_len = batch.len(),
held_for_secs = batch_first_at.map(|t| t.elapsed().as_secs()),
"watcher: held-batch reaper firing — quiescence unachievable, forcing handover"
);
force_handover(&mut batch, ®istry, &child_result_tx);
batch_first_at = None;
}
result = seq_rx.changed() => {
if result.is_err() {
force_handover(&mut batch, ®istry, &child_result_tx);
break;
}
while let Some(r) = mailbox.try_recv_any() {
if matches!(r.status, MailboxStatus::Closed) {
let already_accounted = batch.iter().any(|b| {
b.agent_path == r.agent_path
&& !matches!(b.status, MailboxStatus::Closed)
});
if already_accounted {
tracing::debug!(
agent = %r.agent_path,
"dropping redundant Closed notification"
);
continue;
}
}
if let Some(tx) = &child_result_tx {
let status = status_str(&r.status);
match (&summarizer, &r.status) {
(_, MailboxStatus::Error) => {
let _ = tx.send(ChildResultEvent::Progress {
agent_path: r.agent_path.to_string(),
status: status.to_string(),
summary: error_reason(r.result.as_deref()),
});
}
(Some(s), MailboxStatus::Ok) => {
let _ = tx.send(ChildResultEvent::Progress {
agent_path: r.agent_path.to_string(),
status: status.to_string(),
summary: None,
});
let task = registry
.lock()
.unwrap()
.get(&r.agent_path)
.and_then(|e| e.task.clone());
let agent_name = r.agent_path.name().to_string();
let agent_path = r.agent_path.to_string();
let status = status.to_string();
let result_text = r.result.clone();
let summarizer = Arc::clone(s);
let tx = tx.clone();
tokio::spawn(async move {
let summary = match summarizer
.summarize(
&agent_name,
&status,
task.as_deref(),
result_text.as_deref(),
)
.await
{
Some(text) => Some(text),
None => {
tracing::debug!(
agent = %agent_path,
"progress summary unavailable — plain notice already sent"
);
return;
}
};
let _ = tx.send(ChildResultEvent::Progress {
agent_path,
status,
summary,
});
});
}
_ => {
let _ = tx.send(ChildResultEvent::Progress {
agent_path: r.agent_path.to_string(),
status: status.to_string(),
summary: None,
});
}
}
}
if batch_first_at.is_none() {
batch_first_at = Some(Instant::now());
}
batch.push(r);
}
{
let reg = registry.lock().unwrap();
let q = reg.quiescent();
if !q {
let snapshot = reg.snapshot();
tracing::info!(
batch_len = batch.len(),
agent_count = snapshot.agents.len(),
agents = ?snapshot.agents.iter().map(|a| {
format!("{}: status={}, tool_calls={}, pending={}", a.path, a.status, a.tool_calls, a.pending_results)
}).collect::<Vec<_>>(),
"watcher: quiescent=false, holding batch"
);
}
}
if batch.is_empty() || !registry.lock().unwrap().quiescent() {
continue;
}
let any_real = batch
.iter()
.any(|b| !matches!(b.status, MailboxStatus::Closed));
if !any_real {
batch.clear();
batch_first_at = None;
continue;
}
let batch_paths: Vec<AgentPath> = {
let mut seen = std::collections::BTreeSet::new();
batch
.iter()
.filter(|r| !matches!(r.status, MailboxStatus::Closed))
.filter(|r| seen.insert(r.agent_path.to_string()))
.map(|r| r.agent_path.clone())
.collect()
};
{
let mut reg = registry.lock().unwrap();
for path in &batch_paths {
reg.note_batch_handed_over(path);
}
}
let reports = batch
.drain(..)
.map(|r| format_child_result(&r))
.collect();
batch_first_at = None;
if let Some(tx) = &child_result_tx {
let _ = tx.send(ChildResultEvent::Batch { reports });
}
}
}
}
})
}
fn status_str(status: &MailboxStatus) -> &'static str {
match status {
MailboxStatus::Ok => "ok",
MailboxStatus::Error => "error",
MailboxStatus::Closed => "closed",
}
}
fn force_handover(
batch: &mut Vec<MailboxResult>,
registry: &Arc<Mutex<AgentRegistry>>,
child_result_tx: &Option<mpsc::UnboundedSender<ChildResultEvent>>,
) {
if batch.is_empty() {
return;
}
let any_real = batch
.iter()
.any(|b| !matches!(b.status, MailboxStatus::Closed));
if !any_real {
batch.clear();
return;
}
let batch_paths: Vec<AgentPath> = {
let mut seen = std::collections::BTreeSet::new();
batch
.iter()
.filter(|r| !matches!(r.status, MailboxStatus::Closed))
.filter(|r| seen.insert(r.agent_path.to_string()))
.map(|r| r.agent_path.clone())
.collect()
};
{
let mut reg = registry.lock().unwrap();
for path in &batch_paths {
reg.note_batch_handed_over(path);
}
}
let reports = batch.drain(..).map(|r| format_child_result(&r)).collect();
if let Some(tx) = child_result_tx {
let _ = tx.send(ChildResultEvent::Batch { reports });
}
}
pub fn spawn_watcher_with_watchdog(
mailbox: Arc<MailboxHub>,
registry: Arc<Mutex<AgentRegistry>>,
summarizer: Option<Arc<ProgressSummarizer>>,
child_result_tx: Option<mpsc::UnboundedSender<ChildResultEvent>>,
cancel: CancellationToken,
held_reap_after: Duration,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut restart_count: u32 = 0;
loop {
if cancel.is_cancelled() {
break;
}
let handle = spawn_watcher(
mailbox.clone(),
registry.clone(),
summarizer.clone(),
child_result_tx.clone(),
cancel.clone(),
held_reap_after,
);
match handle.await {
Ok(()) => {
break;
}
Err(join_err) if join_err.is_panic() => {
restart_count += 1;
let panic_info = join_err
.try_into_panic()
.ok()
.and_then(|p| {
p.downcast_ref::<&str>()
.map(|s| s.to_string())
.or_else(|| p.downcast_ref::<String>().cloned())
})
.unwrap_or_else(|| "unknown panic".to_string());
tracing::warn!(
restart_count,
panic_info = %panic_info,
"watcher task panicked, restarting"
);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
Err(join_err) => {
tracing::debug!(error = %join_err, "watcher task join error, exiting watchdog");
break;
}
}
}
if restart_count > 0 {
tracing::info!(restart_count, "watcher watchdog exiting after restarts");
}
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::multi_agent::config::MultiAgentConfig;
use crate::multi_agent::mailbox::MailboxResult;
use crate::multi_agent::path::AgentPath;
struct Fixture {
mailbox: Arc<MailboxHub>,
registry: Arc<Mutex<AgentRegistry>>,
rx: mpsc::UnboundedReceiver<ChildResultEvent>,
cancel: CancellationToken,
#[allow(dead_code)]
handle: tokio::task::JoinHandle<()>,
}
fn fixture() -> Fixture {
fixture_with_summarizer(None)
}
fn fixture_with_summarizer(summarizer: Option<Arc<ProgressSummarizer>>) -> Fixture {
fixture_with(Config {
summarizer,
held_reap_after: std::time::Duration::from_secs(1),
})
}
struct Config {
summarizer: Option<Arc<ProgressSummarizer>>,
held_reap_after: std::time::Duration,
}
fn fixture_with(cfg: Config) -> Fixture {
let mailbox = Arc::new(MailboxHub::new());
let registry = Arc::new(Mutex::new(AgentRegistry::new(MultiAgentConfig::enabled())));
let (tx, rx) = mpsc::unbounded_channel();
let cancel = CancellationToken::new();
let handle = spawn_watcher_with_watchdog(
mailbox.clone(),
registry.clone(),
cfg.summarizer,
Some(tx),
cancel.clone(),
cfg.held_reap_after,
);
Fixture {
mailbox,
registry,
rx,
cancel,
handle,
}
}
fn spawn_running(fx: &Fixture, name: &str) {
let path = AgentPath::root().join(name);
fx.mailbox.register(&path);
fx.registry.lock().unwrap().register(&path, 1).unwrap();
fx.registry.lock().unwrap().note_enqueued(&path);
fx.registry.lock().unwrap().note_dequeued(&path);
}
fn finish_and_post(fx: &Fixture, name: &str, status: MailboxStatus, text: Option<&str>) {
let path = AgentPath::root().join(name);
fx.registry.lock().unwrap().note_posted(&path);
fx.mailbox.register(&path);
fx.mailbox.post_result(MailboxResult {
agent_path: path,
status,
result: text.map(|s| s.to_string()),
denied_tools: vec![],
});
}
async fn next_event(fx: &mut Fixture) -> ChildResultEvent {
tokio::time::timeout(std::time::Duration::from_secs(2), fx.rx.recv())
.await
.expect("event timeout")
.expect("channel closed")
}
async fn assert_no_event(fx: &mut Fixture) {
let got = tokio::time::timeout(std::time::Duration::from_millis(200), fx.rx.recv()).await;
assert!(
got.is_err(),
"expected no event, got {:?}",
got.ok().flatten()
);
}
#[tokio::test]
async fn lone_result_progress_then_batch() {
let mut fx = fixture();
finish_and_post(&fx, "worker", MailboxStatus::Ok, Some("done!"));
match next_event(&mut fx).await {
ChildResultEvent::Progress {
agent_path, status, ..
} => {
assert_eq!(agent_path, "root/worker");
assert_eq!(status, "ok");
}
other => panic!("expected Progress, got {other:?}"),
}
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
assert_eq!(reports[0].agent_path, "root/worker");
assert_eq!(reports[0].status, "ok");
assert_eq!(reports[0].result.as_deref(), Some("done!"));
assert!(reports[0].message.contains("done!"));
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn batch_fire_clears_pending_results() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
let pending = |fx: &Fixture, name: &str| {
fx.registry
.lock()
.unwrap()
.snapshot()
.agents
.into_iter()
.find(|a| a.path == format!("root/{name}"))
.expect("agent registered")
.pending_results
};
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
finish_and_post(&fx, "b", MailboxStatus::Ok, Some("b done"));
next_event(&mut fx).await;
next_event(&mut fx).await;
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => assert_eq!(reports.len(), 2),
other => panic!("expected Batch, got {other:?}"),
}
assert_eq!(pending(&fx, "a"), 0, "handed over with the batch");
assert_eq!(pending(&fx, "b"), 0, "handed over with the batch");
fx.cancel.cancel();
}
#[tokio::test]
async fn batch_held_until_last_child_finishes() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
match next_event(&mut fx).await {
ChildResultEvent::Progress { agent_path, .. } => {
assert_eq!(agent_path, "root/a");
}
other => panic!("expected Progress for a, got {other:?}"),
}
assert_no_event(&mut fx).await;
finish_and_post(&fx, "b", MailboxStatus::Ok, Some("b done"));
match next_event(&mut fx).await {
ChildResultEvent::Progress { agent_path, .. } => {
assert_eq!(agent_path, "root/b");
}
other => panic!("expected Progress for b, got {other:?}"),
}
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
let paths: Vec<&str> = reports.iter().map(|r| r.agent_path.as_str()).collect();
assert_eq!(paths, vec!["root/a", "root/b"]);
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn closed_sibling_rides_along_in_batch() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; assert_no_event(&mut fx).await;
finish_and_post(&fx, "b", MailboxStatus::Closed, None);
let _ = next_event(&mut fx).await; match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 2);
assert_eq!(reports[0].status, "ok");
assert_eq!(reports[1].status, "closed");
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn closed_only_batch_never_wakes_parent() {
let mut fx = fixture();
spawn_running(&fx, "w");
finish_and_post(&fx, "w", MailboxStatus::Closed, None);
match next_event(&mut fx).await {
ChildResultEvent::Progress { status, .. } => assert_eq!(status, "closed"),
other => panic!("expected Progress, got {other:?}"),
}
assert_no_event(&mut fx).await; fx.cancel.cancel();
}
#[tokio::test]
async fn redundant_closed_is_dropped_while_held() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await;
fx.mailbox.post_result(MailboxResult {
agent_path: AgentPath::root().join("a"),
status: MailboxStatus::Closed,
result: None,
denied_tools: vec![],
});
assert_no_event(&mut fx).await;
finish_and_post(&fx, "b", MailboxStatus::Ok, Some("b done"));
let _ = next_event(&mut fx).await; match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 2, "a + b, no duplicate of either");
assert!(
reports.iter().all(|r| r.status == "ok"),
"redundant Closed must not ride along"
);
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn stale_closed_after_flush_does_not_wake() {
let mut fx = fixture();
fx.mailbox.register(&AgentPath::root().join("a"));
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; let _ = next_event(&mut fx).await;
fx.mailbox.post_result(MailboxResult {
agent_path: AgentPath::root().join("a"),
status: MailboxStatus::Closed,
result: None,
denied_tools: vec![],
});
match next_event(&mut fx).await {
ChildResultEvent::Progress { .. } => {}
other => panic!("expected Progress, got {other:?}"),
}
assert_no_event(&mut fx).await; fx.cancel.cancel();
}
#[tokio::test]
async fn batch_flushed_results_are_not_repeated() {
let mut fx = fixture();
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("gen1"));
let _ = next_event(&mut fx).await;
let _ = next_event(&mut fx).await;
finish_and_post(&fx, "b", MailboxStatus::Ok, Some("gen2"));
let _ = next_event(&mut fx).await;
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
assert!(reports[0].message.contains("gen2"));
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn force_closed_child_does_not_block_batch() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; assert_no_event(&mut fx).await;
fx.registry
.lock()
.unwrap()
.note_closing(&AgentPath::root().join("b"));
finish_and_post(&fx, "b", MailboxStatus::Closed, None);
let _ = next_event(&mut fx).await;
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
let paths: Vec<&str> = reports.iter().map(|r| r.agent_path.as_str()).collect();
assert!(
paths.contains(&"root/a"),
"a's held report must be delivered, got {paths:?}"
);
}
other => panic!("expected Batch after close, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn watcher_shutdown_flush_delivers_held_batch() {
let mut fx = fixture();
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; assert_no_event(&mut fx).await;
fx.cancel.cancel();
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
let paths: Vec<&str> = reports.iter().map(|r| r.agent_path.as_str()).collect();
assert_eq!(paths, vec!["root/a"], "held report flushed on shutdown");
}
other => panic!("expected Batch on shutdown flush, got {other:?}"),
}
}
#[tokio::test]
async fn held_batch_reaper_fires_when_nothing_can_post() {
let mut fx = fixture_with(Config {
summarizer: None,
held_reap_after: std::time::Duration::from_millis(100),
});
for name in ["a", "b"] {
spawn_running(&fx, name);
}
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; assert_no_event(&mut fx).await;
fx.registry
.lock()
.unwrap()
.note_closing(&AgentPath::root().join("b"));
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
let paths: Vec<&str> = reports.iter().map(|r| r.agent_path.as_str()).collect();
assert_eq!(paths, vec!["root/a"], "reaper must deliver the held report");
}
other => panic!("expected reaper Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn watcher_exits_on_cancel() {
let mailbox = Arc::new(MailboxHub::new());
let registry = Arc::new(Mutex::new(AgentRegistry::new(MultiAgentConfig::enabled())));
let (tx, _rx) = mpsc::unbounded_channel();
let cancel = CancellationToken::new();
let handle = spawn_watcher(
mailbox.clone(),
registry,
None,
Some(tx),
cancel.clone(),
HELD_BATCH_REAP_AFTER,
);
cancel.cancel();
let result = tokio::time::timeout(std::time::Duration::from_secs(2), handle).await;
assert!(result.is_ok(), "watcher should exit on cancel");
}
#[tokio::test]
async fn watchdog_exits_on_cancel() {
let mailbox = Arc::new(MailboxHub::new());
let registry = Arc::new(Mutex::new(AgentRegistry::new(MultiAgentConfig::enabled())));
let (tx, _rx) = mpsc::unbounded_channel();
let cancel = CancellationToken::new();
let handle = spawn_watcher_with_watchdog(
mailbox.clone(),
registry,
None,
Some(tx),
cancel.clone(),
HELD_BATCH_REAP_AFTER,
);
cancel.cancel();
let result = tokio::time::timeout(std::time::Duration::from_secs(3), handle).await;
assert!(result.is_ok(), "watchdog should exit on cancel");
}
#[tokio::test]
async fn concurrent_results_all_land_in_one_batch() {
let mut fx = fixture();
fx.mailbox.register(&AgentPath::root().join("a"));
fx.mailbox.register(&AgentPath::root().join("b"));
fx.mailbox.post_result(MailboxResult {
agent_path: AgentPath::root().join("a"),
status: MailboxStatus::Ok,
result: Some("a done".to_string()),
denied_tools: vec![],
});
fx.mailbox.post_result(MailboxResult {
agent_path: AgentPath::root().join("b"),
status: MailboxStatus::Error,
result: Some("b failed".to_string()),
denied_tools: vec![],
});
let mut statuses = Vec::new();
for _ in 0..2 {
match next_event(&mut fx).await {
ChildResultEvent::Progress {
agent_path, status, ..
} => {
assert_ne!(agent_path, "root/none");
statuses.push(status);
}
other => panic!("expected Progress, got {other:?}"),
}
}
assert!(statuses.contains(&"ok".to_string()));
assert!(statuses.contains(&"error".to_string()));
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 2);
assert!(reports.iter().any(|r| r.status == "ok"));
assert!(reports.iter().any(|r| r.status == "error"));
}
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
struct SummaryStub {
delay: std::time::Duration,
}
#[async_trait::async_trait]
impl agent_base::llm_trait::LlmProvider for SummaryStub {
async fn stream(
&self,
_request: agent_base::llm_trait::ChatRequest,
) -> Result<agent_base::llm_trait::ChatStream, agent_base::llm_trait::LlmError> {
Ok(agent_base::llm_trait::ChatStream::new(Box::pin(
futures_util::stream::empty(),
)))
}
async fn chat(
&self,
_request: agent_base::llm_trait::ChatRequest,
) -> Result<agent_base::llm_trait::ChatResponse, agent_base::llm_trait::LlmError> {
if !self.delay.is_zero() {
tokio::time::sleep(self.delay).await;
}
Ok(agent_base::llm_trait::ChatResponse {
content: r#"{"summary": "mock 摘要"}"#.to_string(),
tool_calls: vec![],
usage: Default::default(),
finish_reason: agent_base::llm_trait::response::FinishReason::Stop,
raw: None,
reasoning_content: None,
thinking_signature: None,
})
}
fn capabilities(&self) -> agent_base::llm_trait::Capabilities {
Default::default()
}
fn info(&self) -> agent_base::llm_trait::ProviderInfo {
agent_base::llm_trait::ProviderInfo {
name: "summary-stub".to_string(),
model: "stub".to_string(),
version: None,
}
}
}
#[tokio::test]
async fn progress_summary_flows_from_summarizer() {
let summarizer = Arc::new(ProgressSummarizer::new(
Arc::new(SummaryStub {
delay: std::time::Duration::ZERO,
}),
std::time::Duration::from_secs(5),
));
let mut fx = fixture_with_summarizer(Some(summarizer));
spawn_running(&fx, "a");
fx.registry
.lock()
.unwrap()
.set_task(&AgentPath::root().join("a"), "分析任务".to_string());
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("report text"));
let mut saw_plain = false;
let mut saw_summary = false;
let mut batch_ok = false;
for _ in 0..3 {
match next_event(&mut fx).await {
ChildResultEvent::Progress {
status, summary, ..
} => {
assert_eq!(status, "ok");
match summary.as_deref() {
None => saw_plain = true,
Some("mock 摘要") => saw_summary = true,
other => panic!("unexpected summary {other:?}"),
}
}
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
assert_eq!(reports[0].result.as_deref(), Some("report text"));
batch_ok = true;
}
}
}
assert!(
saw_plain,
"synchronous plain Progress must arrive first-class"
);
assert!(saw_summary, "Progress with summary must arrive");
assert!(batch_ok, "Batch with full report must arrive");
fx.cancel.cancel();
}
#[tokio::test]
async fn slow_summary_does_not_delay_batch() {
let summarizer = Arc::new(ProgressSummarizer::new(
Arc::new(SummaryStub {
delay: std::time::Duration::from_millis(500),
}),
std::time::Duration::from_secs(5),
));
let mut fx = fixture_with_summarizer(Some(summarizer));
spawn_running(&fx, "a");
let started = std::time::Instant::now();
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("report text"));
match next_event(&mut fx).await {
ChildResultEvent::Progress { summary, .. } => {
assert!(summary.is_none(), "first notice must be the plain one");
}
other => panic!("expected plain Progress first, got {other:?}"),
}
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
}
ChildResultEvent::Progress {
summary: Some(_), ..
} => panic!("summary raced ahead of Batch"),
ChildResultEvent::Progress { summary: None, .. } => {
panic!("only one plain Progress expected")
}
}
assert!(
started.elapsed() < std::time::Duration::from_millis(400),
"Batch must not wait for the summary (took {:?})",
started.elapsed()
);
match next_event(&mut fx).await {
ChildResultEvent::Progress { summary, .. } => {
assert_eq!(summary.as_deref(), Some("mock 摘要"));
}
other => panic!("expected late Progress, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn closed_progress_skips_summarizer() {
let summarizer = Arc::new(ProgressSummarizer::new(
Arc::new(SummaryStub {
delay: std::time::Duration::ZERO,
}),
std::time::Duration::from_secs(5),
));
let mut fx = fixture_with_summarizer(Some(summarizer));
spawn_running(&fx, "w");
finish_and_post(&fx, "w", MailboxStatus::Closed, None);
match next_event(&mut fx).await {
ChildResultEvent::Progress {
status, summary, ..
} => {
assert_eq!(status, "closed");
assert!(summary.is_none(), "closed needs no LLM call");
}
other => panic!("expected Progress, got {other:?}"),
}
assert_no_event(&mut fx).await; fx.cancel.cancel();
}
#[tokio::test]
async fn error_progress_carries_raw_reason_not_a_paraphrase() {
let summarizer = Arc::new(ProgressSummarizer::new(
Arc::new(SummaryStub {
delay: std::time::Duration::ZERO,
}),
std::time::Duration::from_secs(5),
));
let mut fx = fixture_with_summarizer(Some(summarizer));
spawn_running(&fx, "e");
finish_and_post(
&fx,
"e",
MailboxStatus::Error,
Some(
"LLM call failed: HTTP request failed: error sending request for url (https://example.invalid/v1/messages)",
),
);
match next_event(&mut fx).await {
ChildResultEvent::Progress {
status, summary, ..
} => {
assert_eq!(status, "error");
let reason = summary.expect("error Progress must carry the raw reason");
assert!(reason.starts_with("LLM call failed"), "{reason}");
assert!(reason.contains("error sending request"), "{reason}");
assert!(!reason.contains("mock 摘要"), "paraphrase leaked: {reason}");
}
other => panic!("expected Progress, got {other:?}"),
}
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
assert_eq!(
reports[0].result.as_deref(),
Some(
"LLM call failed: HTTP request failed: error sending request for url (https://example.invalid/v1/messages)"
)
);
}
other => panic!("expected Batch, got {other:?}"),
}
assert_no_event(&mut fx).await;
fx.cancel.cancel();
}
#[tokio::test]
async fn error_reason_truncates_and_skips_blank_lines() {
let long = format!("{}…tail", "x".repeat(300));
assert_eq!(
error_reason(Some("\n \nboom: first line is the meat")),
Some("boom: first line is the meat".to_string())
);
assert_eq!(
error_reason(Some("first line\nsecond line\nthird")),
Some("first line".to_string())
);
let at = "x".repeat(MAX_ERROR_REASON_CHARS);
let r = error_reason(Some(&at)).unwrap();
assert_eq!(r.chars().count(), MAX_ERROR_REASON_CHARS);
assert!(!r.ends_with('…'));
let over = "x".repeat(MAX_ERROR_REASON_CHARS + 1);
let r = error_reason(Some(&over)).unwrap();
assert_eq!(r.chars().count(), MAX_ERROR_REASON_CHARS + 1);
assert!(r.ends_with('…'));
let truncated = error_reason(Some(&long)).unwrap();
assert_eq!(truncated.chars().count(), 161, "160 chars + ellipsis");
assert!(truncated.ends_with('…'));
assert_eq!(error_reason(Some(" \n\n ")), None);
assert_eq!(error_reason(None), None);
}
#[tokio::test]
async fn error_reason_neutralizes_control_chars_and_ansi() {
assert_eq!(
error_reason(Some("\x1b[31mError: boom\x1b[0m")),
Some("Error: boom".to_string())
);
assert_eq!(error_reason(Some("A\rB")), Some("A B".to_string()));
assert_eq!(
error_reason(Some("error:\tboom")),
Some("error: boom".to_string())
);
assert_eq!(
error_reason(Some("\x07\nreal reason")),
Some("real reason".to_string())
);
let dirty = format!("\x1b[1m{}", "y".repeat(300));
let r = error_reason(Some(&dirty)).unwrap();
assert_eq!(r.chars().count(), 161);
assert!(r.ends_with('…'));
}
#[tokio::test]
async fn error_progress_with_none_result_falls_back_to_plain_notice() {
let summarizer = Arc::new(ProgressSummarizer::new(
Arc::new(SummaryStub {
delay: std::time::Duration::ZERO,
}),
std::time::Duration::from_secs(5),
));
let mut fx = fixture_with_summarizer(Some(summarizer));
spawn_running(&fx, "e");
finish_and_post(&fx, "e", MailboxStatus::Error, None);
match next_event(&mut fx).await {
ChildResultEvent::Progress {
status, summary, ..
} => {
assert_eq!(status, "error");
assert!(
summary.is_none(),
"None error text must fall back to a plain notice"
);
}
other => panic!("expected Progress, got {other:?}"),
}
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1);
assert_eq!(reports[0].result, None);
}
other => panic!("expected Batch, got {other:?}"),
}
assert_no_event(&mut fx).await;
fx.cancel.cancel();
}
#[tokio::test]
async fn fresh_spawn_without_delivery_blocks_the_batch() {
let mut fx = fixture();
spawn_running(&fx, "a");
let b = AgentPath::root().join("b");
fx.mailbox.register(&b);
fx.registry.lock().unwrap().register(&b, 1).unwrap();
finish_and_post(&fx, "a", MailboxStatus::Ok, Some("a done"));
let _ = next_event(&mut fx).await; assert_no_event(&mut fx).await;
{
let mut reg = fx.registry.lock().unwrap();
reg.note_enqueued(&b);
reg.note_dequeued(&b);
}
finish_and_post(&fx, "b", MailboxStatus::Ok, Some("b done"));
let _ = next_event(&mut fx).await;
match next_event(&mut fx).await {
ChildResultEvent::Batch { reports } => assert_eq!(
reports
.iter()
.map(|r| r.agent_path.as_str())
.collect::<Vec<_>>(),
vec!["root/a", "root/b"]
),
other => panic!("expected Batch, got {other:?}"),
}
fx.cancel.cancel();
}
#[tokio::test]
async fn close_during_queued_task_drops_phantom_queue_from_quiescence() {
let mut fx = fixture();
let b = AgentPath::root().join("b");
fx.mailbox.register(&b);
{
let mut reg = fx.registry.lock().unwrap();
reg.register(&b, 1).unwrap();
reg.note_enqueued(&b);
}
assert!(!fx.registry.lock().unwrap().quiescent());
fx.registry.lock().unwrap().close(&b); assert!(
fx.registry.lock().unwrap().quiescent(),
"phantom queue must not outlive the entry"
);
finish_and_post(&fx, "b", MailboxStatus::Closed, None);
match next_event(&mut fx).await {
ChildResultEvent::Progress { status, .. } => assert_eq!(status, "closed"),
other => panic!("expected Progress(closed), got {other:?}"),
}
assert_no_event(&mut fx).await;
fx.cancel.cancel();
}
}