use super::*;
#[tokio::test(flavor = "multi_thread")]
async fn test_spawn_send_task_wait_close_lifecycle() {
let ma = make_ma_runtime();
let path = ma
.spawn_child(
"worker",
"child system prompt".to_string(),
0,
false,
vec![],
)
.await
.expect("spawn child");
assert_eq!(path, "root/worker");
let agents = ma.list_agents();
assert_eq!(agents.len(), 1);
assert_eq!(agents[0].agent_path, "root/worker");
assert_eq!(agents[0].task, None);
assert!(
ma.send_message("root/worker", "heads up".to_string())
.unwrap()
);
assert!(
ma.send_task("root/worker", "do the thing".to_string(), false)
.unwrap()
);
let agents = ma.list_agents();
assert_eq!(agents[0].task.as_deref(), Some("do the thing"));
let result = ma.wait_for_result(Some("root/worker"), 2000).await;
assert_eq!(result.status, "ok");
assert_eq!(result.result.as_deref(), Some("child ok"));
let close = ma.close_agent("root/worker").unwrap();
assert!(close.closed);
assert_eq!(close.message, "agent closed");
let close2 = ma.close_agent("root/worker").unwrap();
assert!(!close2.closed);
assert_eq!(close2.message, "agent not found");
}
#[tokio::test(flavor = "multi_thread")]
async fn test_spawn_child_with_history_defaults_to_none() {
let ma = make_ma_runtime();
let path = ma
.spawn_child_with_history(
"w2",
"prompt".to_string(),
false,
None,
None,
&agent_base::SessionId::new(0),
)
.await
.expect("spawn with history");
assert_eq!(path, "root/w2");
}
#[tokio::test(flavor = "multi_thread")]
async fn test_error_paths() {
let ma = make_ma_runtime();
assert_eq!(
ma.send_task("root/ghost", "x".to_string(), false)
.unwrap_err(),
"agent not found"
);
assert!(ma.send_message("worker", "x".to_string()).is_err());
assert!(ma.send_message("", "x".to_string()).is_err());
let r = ma.wait_for_result(Some("worker"), 10).await;
assert_eq!(r.status, "error");
let r2 = ma.wait_for_result(None, 50).await;
assert_eq!(r2.status, "timeout");
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_pending_when_no_result() {
let ma = make_ma_runtime();
let _path = ma
.spawn_child(
"worker",
"child system prompt".to_string(),
0,
false,
vec![],
)
.await
.expect("spawn child");
let result = ma.try_wait(Some("root/worker"));
assert_eq!(result.status, "pending");
assert!(result.result.is_none());
assert_eq!(result.agent_path.as_deref(), Some("root/worker"));
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_returns_result_after_post() {
use crate::multi_agent::mailbox::{MailboxResult, MailboxStatus};
let ma = make_ma_runtime();
let _path = ma
.spawn_child(
"worker",
"child system prompt".to_string(),
0,
false,
vec![],
)
.await
.expect("spawn child");
ma.mailbox().post_result(MailboxResult {
agent_path: crate::multi_agent::path::AgentPath::parse("root/worker").unwrap(),
status: MailboxStatus::Ok,
result: Some("done!".to_string()),
denied_tools: vec![],
});
let result = ma.try_wait(Some("root/worker"));
assert_eq!(result.status, "ok");
assert_eq!(result.result.as_deref(), Some("done!"));
assert_eq!(result.agent_path.as_deref(), Some("root/worker"));
assert!(!result.has_more);
let result2 = ma.try_wait(Some("root/worker"));
assert_eq!(result2.status, "pending");
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_any_returns_pending_when_empty() {
let ma = make_ma_runtime();
let result = ma.try_wait(None);
assert_eq!(result.status, "pending");
assert!(result.result.is_none());
assert!(result.agent_path.is_none());
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_invalid_path_returns_error() {
let ma = make_ma_runtime();
let result = ma.try_wait(Some("worker"));
assert_eq!(result.status, "error");
assert!(
result
.result
.as_deref()
.unwrap()
.contains("invalid agent path")
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_closed_after_agent_gone() {
let ma = make_ma_runtime();
let path = ma
.spawn_child("worker", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task(&path, "do it".to_string(), false).unwrap();
let result = ma.wait_for_result(Some("root/worker"), 2000).await;
assert_eq!(result.status, "ok");
ma.close_agent("root/worker").unwrap();
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let result = ma.try_wait(Some("root/worker"));
assert_eq!(result.status, "closed");
}
#[tokio::test(flavor = "multi_thread")]
async fn test_try_wait_any_returns_first_available() {
use crate::multi_agent::mailbox::{MailboxResult, MailboxStatus};
let ma = make_ma_runtime();
ma.spawn_child("a", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn a");
ma.spawn_child("b", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn b");
ma.mailbox().post_result(MailboxResult {
agent_path: crate::multi_agent::path::AgentPath::parse("root/b").unwrap(),
status: MailboxStatus::Ok,
result: Some("b done".to_string()),
denied_tools: vec![],
});
let result = ma.try_wait(None);
assert_eq!(result.status, "ok");
assert_eq!(result.agent_path.as_deref(), Some("root/b"));
assert_eq!(result.result.as_deref(), Some("b done"));
let result2 = ma.try_wait(None);
assert_eq!(result2.status, "pending");
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_wait_for_result_basic() {
let ma = make_ma_runtime();
ma.spawn_child("w", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/w", "go".to_string(), false).unwrap();
let result = ma.wait_for_result(Some("root/w"), 5000).await;
assert_eq!(result.status, "ok");
assert_eq!(result.result.as_deref(), Some("child ok"));
assert_eq!(result.agent_path.as_deref(), Some("root/w"));
assert!(!result.has_more);
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_wait_for_result_any() {
let ma = make_ma_runtime();
ma.spawn_child("x", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/x", "go".to_string(), false).unwrap();
let result = ma.wait_for_result(None, 5000).await;
assert_eq!(result.status, "ok");
assert_eq!(result.agent_path.as_deref(), Some("root/x"));
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_wait_for_result_timeout() {
let ma = make_ma_runtime();
ma.spawn_child("slow", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
let result = ma.wait_for_result(Some("root/slow"), 100).await;
assert_eq!(result.status, "timeout");
assert!(result.result.is_none());
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_wait_for_result_has_more() {
use crate::multi_agent::mailbox::{MailboxResult, MailboxStatus};
let ma = make_ma_runtime();
ma.spawn_child("a", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn a");
ma.spawn_child("b", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn b");
ma.mailbox().post_result(MailboxResult {
agent_path: crate::multi_agent::path::AgentPath::parse("root/a").unwrap(),
status: MailboxStatus::Ok,
result: Some("a done".to_string()),
denied_tools: vec![],
});
ma.mailbox().post_result(MailboxResult {
agent_path: crate::multi_agent::path::AgentPath::parse("root/b").unwrap(),
status: MailboxStatus::Ok,
result: Some("b done".to_string()),
denied_tools: vec![],
});
let result = ma.wait_for_result(Some("root/a"), 2000).await;
assert_eq!(result.status, "ok");
assert!(result.has_more, "should report more results pending");
let result2 = ma.wait_for_result(Some("root/b"), 2000).await;
assert_eq!(result2.status, "ok");
assert!(!result2.has_more);
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_list_agents_shows_done_after_result_delivered() {
let ma = make_ma_runtime();
ma.spawn_child("w", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/w", "task".to_string(), false).unwrap();
let result = ma.wait_for_result(Some("root/w"), 2000).await;
assert_eq!(result.status, "ok");
let agents = ma.list_agents();
assert_eq!(
agents[0].status, "done",
"a child that delivered its result must not report running"
);
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_close_idle_child_delivers_closed_event() {
let ma = make_ma_runtime();
ma.spawn_child("w", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/w", "task".to_string(), false).unwrap();
let result = ma.wait_for_result(Some("root/w"), 2000).await;
assert_eq!(result.status, "ok");
let (_handle, mut rx) = ma.start_watcher();
let close = ma.close_agent("root/w").unwrap();
assert!(close.closed);
let cr = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("timed out waiting for closed event")
.expect("watcher channel closed");
match cr {
ChildResultEvent::Progress {
agent_path, status, ..
} => {
assert_eq!(agent_path, "root/w");
assert_eq!(status, "closed");
}
other => panic!("expected Progress(closed), got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn test_queued_task_does_not_fire_premature_batch() {
let ma = make_ma_runtime();
let (_handle, mut rx) = ma.start_watcher();
ma.spawn_child("w", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/w", "task one".to_string(), false)
.unwrap();
ma.send_task("root/w", "task two".to_string(), false)
.unwrap();
let mut batches = Vec::new();
loop {
let cr = tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv())
.await
.expect("timed out waiting for watcher events")
.expect("watcher channel closed");
if let ChildResultEvent::Batch { reports } = cr {
batches.push(reports);
break;
}
}
assert_eq!(batches.len(), 1);
assert_eq!(
batches[0].len(),
2,
"first batch must carry both queued-task reports"
);
while let Ok(Some(cr)) =
tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()).await
{
assert!(
!matches!(cr, ChildResultEvent::Batch { .. }),
"straggler batch after the flush: {:?}",
cr
);
}
ma.cancel_all();
}
#[tokio::test(flavor = "multi_thread")]
async fn test_close_running_child_result_still_delivered() {
let ma = make_ma_runtime_with(Arc::new(DelayedStub));
let (_handle, mut rx) = ma.start_watcher();
ma.spawn_child("w", "prompt".to_string(), 0, false, vec![])
.await
.expect("spawn");
ma.send_task("root/w", "task".to_string(), false).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let close = ma.close_agent("root/w").unwrap();
assert!(close.closed, "child was mid-task, close must land");
let mut saw_progress_ok = false;
let mut saw_batch = false;
for _ in 0..2 {
let cr = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv())
.await
.expect("timed out waiting for late result")
.expect("watcher channel closed");
match cr {
ChildResultEvent::Progress {
agent_path, status, ..
} => {
assert_eq!(agent_path, "root/w");
assert_eq!(status, "ok");
saw_progress_ok = true;
}
ChildResultEvent::Batch { reports } => {
assert_eq!(reports.len(), 1, "redundant Closed must not ride along");
assert_eq!(reports[0].status, "ok");
assert_eq!(reports[0].result.as_deref(), Some("late ok"));
saw_batch = true;
}
}
}
assert!(saw_progress_ok, "late result must surface as Progress");
assert!(saw_batch, "late result must surface in the batch");
if let Ok(Some(extra)) =
tokio::time::timeout(std::time::Duration::from_millis(300), rx.recv()).await
{
match extra {
ChildResultEvent::Progress { status, .. } => {
assert_eq!(status, "closed", "only a lone Closed may follow the flush");
}
other => panic!("unexpected extra event {other:?}"),
}
}
}