use super::support::*;
#[test]
fn subagents_stalled_child_returns_failed_result_and_preserves_child_tool_result() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(StallingAfterToolResultProvider::new());
let events: Arc<Mutex<Vec<ActivityEvent>>> = Arc::new(Mutex::new(Vec::new()));
let captured = Arc::clone(&events);
let mut cfg = config(provider.clone(), temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.activity_sender = Some(Arc::new(move |event| captured.lock().unwrap().push(event)));
let output = run_subagents_with_wait(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![
SubagentTask {
intent: "run bash then stall".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "must not start after sibling stall".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
],
},
cfg,
SchedulerWaitConfig {
poll_interval: Duration::from_millis(10),
stall_after: Duration::from_millis(500),
},
)
.unwrap();
assert_eq!(output.summary.total, 2);
assert_eq!(output.summary.completed, 0);
assert_eq!(output.summary.failed, 2);
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Failed);
assert!(
result
.error
.as_deref()
.unwrap()
.contains("stalled with no activity"),
"{:?}",
result.error
);
let session_path = result.session_path.as_ref().expect("child session path");
assert!(session_path.exists());
assert_eq!(output.results[1].status, SubagentStatus::Failed);
assert!(output.results[1].session_path.is_none());
assert!(
output.results[1]
.error
.as_deref()
.unwrap()
.contains("did not start before scheduler workers stalled")
);
assert!(provider.saw_continuation.load(Ordering::SeqCst));
let requests = provider.requests.lock().unwrap();
assert_eq!(requests.len(), 2);
let continuation_results = requests[1].tool_results();
assert_eq!(continuation_results.len(), 1);
assert_eq!(continuation_results[0].call_id, "bash_stall_1");
assert!(continuation_results[0].success);
assert!(continuation_results[0].output.contains("child bash ok"));
drop(requests);
let session_jsonl = std::fs::read_to_string(session_path).unwrap();
assert!(session_jsonl.contains("tool_call"), "{session_jsonl}");
assert!(session_jsonl.contains("tool_result"), "{session_jsonl}");
assert!(session_jsonl.contains("bash_stall_1"), "{session_jsonl}");
assert!(session_jsonl.contains("child bash ok"), "{session_jsonl}");
assert!(
!session_jsonl.contains("hook_lifecycle"),
"child subagent tool calls with inherited_hooks=None should not run hooks: {session_jsonl}"
);
let events = events.lock().unwrap();
for task_id in ["subagents/g1", "subagents/g2"] {
let lifecycle = events
.iter()
.filter(|event| match event {
ActivityEvent::Started { id, .. } | ActivityEvent::Finished { id, .. } => {
id.as_str() == task_id
}
_ => false,
})
.collect::<Vec<_>>();
assert!(lifecycle.len() >= 2, "{task_id}: {lifecycle:?}");
assert!(
matches!(lifecycle[0], ActivityEvent::Started { parent_id: Some(parent), kind: ActivityKind::SubagentTask, .. } if parent.as_str() == "subagents")
);
assert!(matches!(
lifecycle.last(),
Some(ActivityEvent::Finished {
status: ActivityStatus::Failed,
..
})
));
}
assert_eq!(
events
.iter()
.filter(|event| matches!(event, ActivityEvent::Finished { id, status: ActivityStatus::Failed, .. } if id.as_str() == "subagents/g1"))
.count(),
1
);
assert!(events.iter().any(|event| matches!(event, ActivityEvent::ToolResultDetail { id, detail } if id.as_str() == "subagents/g1/bash_stall_1" && detail.status == ActivityStatus::Success)));
assert!(!events.iter().any(|event| matches!(event, ActivityEvent::FinalPreview { id, .. } if id.as_str() == "subagents/g1/bash_after_cancel")));
drop(events);
let deadline = std::time::Instant::now() + Duration::from_secs(1);
while !provider.saw_cancellation.load(Ordering::SeqCst) && std::time::Instant::now() < deadline
{
thread::sleep(Duration::from_millis(10));
}
assert!(provider.saw_cancellation.load(Ordering::SeqCst));
assert!(provider.attempted_post_cancel_event.load(Ordering::SeqCst));
thread::sleep(Duration::from_millis(100));
assert_eq!(provider.requests.lock().unwrap().len(), 2);
let session_jsonl_after_cancel = std::fs::read_to_string(session_path).unwrap();
assert!(
!session_jsonl_after_cancel.contains("bash_after_cancel"),
"{session_jsonl_after_cancel}"
);
provider.unblock();
}
#[test]
fn scheduler_preserves_order_and_isolates_failures() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("fail"));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(2),
tasks: vec![
SubagentTask {
intent: "first".into(),
agent: Some("a".into()),
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "fail this".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "third".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
],
},
config(provider.clone(), temp.path()),
)
.unwrap();
assert_eq!(output.summary.total, 3);
assert_eq!(output.summary.completed, 2);
assert_eq!(output.summary.failed, 1);
assert_eq!(output.results[0].id, "g1");
assert_eq!(output.results[1].id, "g2");
assert_eq!(output.results[2].id, "g3");
assert_eq!(output.results[1].status, SubagentStatus::Failed);
assert!(provider.max_active.load(Ordering::SeqCst) <= 2);
}
#[test]
fn subagents_summary_omits_total_tokens_when_any_result_usage_unknown() {
let output = output_from_results(vec![
completed_result("g1", Some(100)),
completed_result("g2", None),
]);
assert_eq!(output.summary.completed, 2);
assert_eq!(output.summary.failed, 0);
assert_eq!(output.summary.total_tokens, None);
assert_eq!(output.results[0].total_tokens, Some(100));
assert_eq!(output.results[1].total_tokens, None);
}
#[test]
fn subagents_summary_sums_total_tokens_when_all_results_known() {
let output = output_from_results(vec![
completed_result("g1", Some(100)),
completed_result("g2", Some(23)),
]);
assert_eq!(output.summary.completed, 2);
assert_eq!(output.summary.failed, 0);
assert_eq!(output.summary.total_tokens, Some(123));
}
#[test]
fn subagents_output_includes_total_tokens_when_provider_reports_usage() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(TokenUsageProvider);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(2),
tasks: vec![
SubagentTask {
intent: "first".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "second".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
],
},
config(provider, temp.path()),
)
.unwrap();
assert_eq!(output.summary.completed, 2);
assert_eq!(output.summary.total_tokens, Some(246));
assert_eq!(output.results[0].total_tokens, Some(123));
assert_eq!(output.results[1].total_tokens, Some(123));
for result in &output.results {
let usage = result.usage.expect("usage from provider");
assert_eq!(usage.whole_run.effective_input, 40);
assert_eq!(usage.whole_run.output, 2);
assert_eq!(usage.latest, Some(usage.whole_run));
}
}
#[test]
fn scheduler_cancellation_returns_collected_results_in_order() {
let temp = tempfile::TempDir::new().unwrap();
let cancel = Arc::new(AtomicBool::new(true));
let mut scheduler = SubagentScheduler::new(
vec![
SubagentTask {
intent: "first".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "second".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "third".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
],
1,
{
let mut cfg = config(Arc::new(CountingProvider::new("never")), temp.path());
cfg.cancellation = AgentCancellation::new(cancel);
cfg
},
ActivityId::new("subagents"),
)
.unwrap();
let mut results = vec![None; 3];
let mut remaining = 3;
scheduler
.apply_worker_event_for_test(
WorkerEvent::TaskFinished {
index: 0,
result: Box::new(completed_result("g1", Some(7))),
},
&mut results,
&mut remaining,
)
.unwrap();
scheduler
.apply_worker_event_for_test(
WorkerEvent::TaskFinished {
index: 1,
result: Box::new(failed_result(
"g2".to_string(),
SubagentTask {
intent: "second".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
temp.path().to_path_buf(),
"child failure".into(),
)),
},
&mut results,
&mut remaining,
)
.unwrap();
let output = scheduler
.output_after_cancellation_for_test(&mut results, None)
.unwrap();
assert_eq!(
output
.results
.iter()
.map(|result| result.id.as_str())
.collect::<Vec<_>>(),
vec!["g1", "g2", "g3"]
);
assert_eq!(output.summary.total, 3);
assert_eq!(output.summary.completed, 1);
assert_eq!(output.summary.failed, 2);
assert_eq!(output.results[0].total_tokens, Some(7));
assert_eq!(output.results[1].error.as_deref(), Some("child failure"));
assert_eq!(output.results[2].error.as_deref(), Some("prompt canceled"));
}
#[test]
fn parent_cancellation_reaches_child_agent_run() {
let temp = tempfile::TempDir::new().unwrap();
let cancel = Arc::new(AtomicBool::new(false));
let provider = Arc::new(ChildCancelingProvider::new(Arc::clone(&cancel)));
let mut cfg = config(provider.clone(), temp.path());
cfg.cancellation = AgentCancellation::new(cancel);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "cancel child".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.total, 1);
assert_eq!(output.summary.completed, 0);
assert_eq!(output.summary.failed, 1);
assert!(
output.results[0]
.error
.as_deref()
.unwrap()
.to_lowercase()
.contains("cancel"),
"{:?}",
output.results[0].error
);
assert_eq!(provider.requests.lock().unwrap().len(), 1);
}
#[test]
fn parent_cancellation_drains_worker_before_returning() {
let temp = tempfile::TempDir::new().unwrap();
let parent_cancel = Arc::new(AtomicBool::new(false));
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let provider = Arc::new(DrainOnCancellationProvider::new(entered_tx, release_rx));
let mut cfg = config(provider.clone(), temp.path());
cfg.cancellation = AgentCancellation::new(Arc::clone(&parent_cancel));
let returned = Arc::new(AtomicBool::new(false));
let returned_in_thread = Arc::clone(&returned);
let handle = thread::spawn(move || {
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "wait for cancellation".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
returned_in_thread.store(true, Ordering::SeqCst);
output
});
entered_rx.recv_timeout(Duration::from_secs(1)).unwrap();
parent_cancel.store(true, Ordering::SeqCst);
let deadline = std::time::Instant::now() + Duration::from_secs(1);
while !provider.observed_child_cancellation.load(Ordering::SeqCst)
&& std::time::Instant::now() < deadline
{
thread::sleep(Duration::from_millis(5));
}
assert!(provider.observed_child_cancellation.load(Ordering::SeqCst));
assert!(!returned.load(Ordering::SeqCst));
release_tx.send(()).unwrap();
let output = handle.join().unwrap();
assert_eq!(output.summary.failed, 1);
assert!(
output.results[0]
.error
.as_deref()
.unwrap()
.to_lowercase()
.contains("cancel"),
"{:?}",
output.results[0].error
);
assert!(returned.load(Ordering::SeqCst));
assert_eq!(provider.requests.lock().unwrap().len(), 1);
}
#[test]
fn parent_cancellation_blocked_worker_released_before_return() {
let temp = tempfile::TempDir::new().unwrap();
let parent_cancel = Arc::new(AtomicBool::new(false));
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let provider = Arc::new(DrainOnCancellationProvider::new(entered_tx, release_rx));
let mut cfg = config(provider.clone(), temp.path());
cfg.cancellation = AgentCancellation::new(Arc::clone(&parent_cancel));
let (returned_tx, returned_rx) = std::sync::mpsc::channel();
let handle = thread::spawn(move || {
let result = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "stay blocked after cancellation".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
);
returned_tx.send(result).unwrap();
});
entered_rx.recv_timeout(Duration::from_secs(1)).unwrap();
let cancellation_started = std::time::Instant::now();
parent_cancel.store(true, Ordering::SeqCst);
assert!(
returned_rx
.recv_timeout(Duration::from_millis(450))
.is_err(),
"scheduler returned before 500ms shutdown grace period"
);
let result = returned_rx
.recv_timeout(Duration::from_secs(2))
.expect("scheduler must return at its shutdown deadline");
let elapsed = cancellation_started.elapsed();
let output = result.unwrap();
assert_eq!(output.summary.failed, 1);
assert!(provider.observed_child_cancellation.load(Ordering::SeqCst));
let error_text = output.results[0].error.as_deref().unwrap();
assert!(
error_text.contains("subagent worker cleanup incomplete after 500ms grace period"),
"{error_text}"
);
assert!(
error_text.contains("detached worker still occupies capacity until it exits"),
"{error_text}"
);
assert!(error_text.contains("prompt canceled"), "{error_text}");
release_tx.send(()).unwrap();
handle.join().unwrap();
assert!(
elapsed >= Duration::from_millis(450),
"scheduler returned before 500ms drain: {elapsed:?}"
);
assert!(elapsed < Duration::from_secs(2), "elapsed={elapsed:?}");
}
#[test]
fn stalled_child_scheduler_drains_worker_before_returning() {
let temp = tempfile::TempDir::new().unwrap();
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let provider = Arc::new(DrainOnCancellationProvider::new(entered_tx, release_rx));
let cfg = config(provider.clone(), temp.path());
let returned = Arc::new(AtomicBool::new(false));
let returned_in_thread = Arc::clone(&returned);
let handle = thread::spawn(move || {
let output = run_subagents_with_wait(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "stall until canceled".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
SchedulerWaitConfig {
poll_interval: Duration::from_millis(10),
stall_after: Duration::from_millis(40),
},
)
.unwrap();
returned_in_thread.store(true, Ordering::SeqCst);
output
});
entered_rx.recv_timeout(Duration::from_secs(1)).unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(1);
while !provider.observed_child_cancellation.load(Ordering::SeqCst)
&& std::time::Instant::now() < deadline
{
thread::sleep(Duration::from_millis(5));
}
assert!(provider.observed_child_cancellation.load(Ordering::SeqCst));
release_tx.send(()).unwrap();
let output = handle.join().unwrap();
assert!(returned.load(Ordering::SeqCst));
assert_eq!(output.summary.failed, 1);
assert!(
output.results[0]
.error
.as_deref()
.unwrap()
.contains("stalled with no activity")
);
}
#[test]
fn scheduler_cancellation_drain_retains_task_finished_before_deadline() {
let temp = tempfile::TempDir::new().unwrap();
let mut scheduler = SubagentScheduler::new(
vec![SubagentTask {
intent: "race completion with cancellation".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
1,
config(Arc::new(CountingProvider::new("never")), temp.path()),
ActivityId::new("subagents"),
)
.unwrap();
let (sender, receiver) = std::sync::mpsc::sync_channel(1);
let (release_tx, release_rx) = std::sync::mpsc::channel();
let (event_sent_tx, event_sent_rx) = std::sync::mpsc::channel();
let (done_tx, done_rx) = std::sync::mpsc::channel();
let worker = thread::spawn(move || {
send_terminal_worker_event(
&sender,
WorkerEvent::TaskFinished {
index: 0,
result: Box::new(completed_result("g1", Some(7))),
},
);
event_sent_tx.send(()).unwrap();
release_rx.recv().unwrap();
done_tx.send(()).unwrap();
});
event_sent_rx
.recv_timeout(Duration::from_millis(100))
.expect("TaskFinished must be queued before cancellation drain");
let mut workers = vec![worker];
let mut results = vec![None];
let warning = scheduler
.cancel_and_drain_workers(&mut workers, &receiver, &mut results)
.unwrap()
.expect("blocked worker cleanup warning");
assert!(warning.contains("cleanup incomplete"), "{warning}");
let result = results[0].as_ref().expect("completed result retained");
assert_eq!(result.status, SubagentStatus::Completed);
assert_eq!(result.total_tokens, Some(7));
release_tx.send(()).unwrap();
done_rx.recv_timeout(Duration::from_secs(1)).unwrap();
}
#[test]
fn scheduler_session_events_are_required_and_survive_backpressure() {
let (sender, receiver) = std::sync::mpsc::sync_channel(SCHEDULER_EVENT_CHANNEL_BOUND);
let reporter = SchedulerTaskReporter { index: 0, sender };
for _ in 0..SCHEDULER_EVENT_CHANNEL_BOUND {
reporter.progress();
}
let handle = thread::spawn(move || {
reporter.session("session-id".to_string(), PathBuf::from("session.jsonl"));
});
let _ = receiver.recv().unwrap();
handle.join().unwrap();
let session = receiver
.try_iter()
.find_map(|event| match event {
WorkerEvent::TaskSession {
session_id,
session_path,
..
} => Some((session_id, session_path)),
_ => None,
})
.unwrap();
assert_eq!(session.0, "session-id");
assert_eq!(session.1, PathBuf::from("session.jsonl"));
}
#[test]
fn scheduler_progress_events_are_bounded_and_droppable() {
let (sender, receiver) = std::sync::mpsc::sync_channel(SCHEDULER_EVENT_CHANNEL_BOUND);
let reporter = SchedulerTaskReporter { index: 0, sender };
for _ in 0..(SCHEDULER_EVENT_CHANNEL_BOUND * 2) {
reporter.progress();
}
let mut progress_events = 0;
while let Ok(event) = receiver.try_recv() {
assert!(matches!(event, WorkerEvent::TaskProgress { index: 0, .. }));
progress_events += 1;
}
assert_eq!(progress_events, SCHEDULER_EVENT_CHANNEL_BOUND);
}
#[test]
fn scheduler_terminal_events_survive_backpressure() {
let (sender, receiver) = std::sync::mpsc::sync_channel(SCHEDULER_EVENT_CHANNEL_BOUND);
let reporter = SchedulerTaskReporter {
index: 0,
sender: sender.clone(),
};
for _ in 0..SCHEDULER_EVENT_CHANNEL_BOUND {
reporter.progress();
}
let handle = thread::spawn(move || {
send_terminal_worker_event(
&sender,
WorkerEvent::TaskFinished {
index: 0,
result: Box::new(completed_result("g1", None)),
},
);
});
let mut saw_finished = false;
for _ in 0..=SCHEDULER_EVENT_CHANNEL_BOUND {
match receiver.recv_timeout(Duration::from_secs(5)).unwrap() {
WorkerEvent::TaskFinished { index, .. } => {
assert_eq!(index, 0);
saw_finished = true;
break;
}
WorkerEvent::TaskProgress { .. } => {}
_ => panic!("unexpected worker event"),
}
}
handle.join().unwrap();
assert!(saw_finished);
let (sender, receiver) = std::sync::mpsc::sync_channel(SCHEDULER_EVENT_CHANNEL_BOUND);
let reporter = SchedulerTaskReporter {
index: 0,
sender: sender.clone(),
};
for _ in 0..SCHEDULER_EVENT_CHANNEL_BOUND {
reporter.progress();
}
let handle = thread::spawn(move || {
send_terminal_worker_event(
&sender,
WorkerEvent::WorkerFailed("worker failed".to_string()),
);
});
let mut saw_failed = false;
for _ in 0..=SCHEDULER_EVENT_CHANNEL_BOUND {
match receiver.recv_timeout(Duration::from_secs(5)).unwrap() {
WorkerEvent::WorkerFailed(error) => {
assert_eq!(error, "worker failed");
saw_failed = true;
break;
}
WorkerEvent::TaskProgress { .. } => {}
_ => panic!("unexpected worker event"),
}
}
handle.join().unwrap();
assert!(saw_failed);
}
#[test]
fn worker_failure_cancels_unfinished_siblings_without_changing_child_failures() {
let temp = tempfile::TempDir::new().unwrap();
let mut scheduler = SubagentScheduler::new(
vec![
SubagentTask {
intent: "one".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
SubagentTask {
intent: "two".into(),
agent: None,
identity: None,
context: None,
cwd: None,
},
],
2,
config(Arc::new(CountingProvider::new("never")), temp.path()),
ActivityId::new("subagents"),
)
.unwrap();
let (error, cancellations) =
scheduler.apply_worker_failure_for_test("worker setup failed".to_string());
assert!(error.contains("worker setup failed"), "{error}");
assert_eq!(cancellations, vec![true, true]);
}
#[test]
fn scheduler_mutex_poisoning_returns_errors() {
let (cancellation, cancel_handle) = AgentCancellation::default().child_token();
let task = PreparedSubagentTask {
index: 0,
id: "g1".to_string(),
task: SubagentTask {
intent: "one".to_string(),
agent: None,
identity: None,
context: None,
cwd: None,
},
cwd: PathBuf::from("."),
cancellation,
cancel_handle,
};
let queue: SharedTaskQueue = Arc::new(Mutex::new(vec![task].into_iter()));
let poisoned_queue = Arc::clone(&queue);
let _ = thread::spawn(move || {
let _guard = poisoned_queue.lock().unwrap();
panic!("poison queue");
})
.join();
let queue_error = SubagentScheduler::next_task(&queue)
.unwrap_err()
.to_string();
assert!(queue_error.contains("subagent queue mutex poisoned"));
let results: SharedTaskResults = Arc::new(Mutex::new(vec![None]));
let poisoned_results = Arc::clone(&results);
let _ = thread::spawn(move || {
let _guard = poisoned_results.lock().unwrap();
panic!("poison results");
})
.join();
let result_error = SubagentScheduler::record_result(
&results,
0,
failed_result(
"g1".to_string(),
SubagentTask {
intent: "one".to_string(),
agent: None,
identity: None,
context: None,
cwd: None,
},
PathBuf::from("."),
"failed".to_string(),
),
)
.unwrap_err()
.to_string();
assert!(result_error.contains("subagent results mutex poisoned"));
}
#[test]
fn scheduler_missing_result_guard_returns_error() {
let results: SharedTaskResults = Arc::new(Mutex::new(vec![None]));
let error = SubagentScheduler::finish_results(results)
.unwrap_err()
.to_string();
assert!(error.contains("subagent result missing at index 0"));
}
#[test]
fn validates_min_one_and_max_ten_tasks() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let empty_err = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![],
},
config(provider.clone(), temp.path()),
)
.unwrap_err()
.to_string();
assert!(empty_err.contains("at least one task"));
let err = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: (0..11)
.map(|i| SubagentTask {
intent: format!("t{i}"),
agent: None,
identity: None,
context: None,
cwd: None,
})
.collect(),
},
config(provider, temp.path()),
)
.unwrap_err()
.to_string();
assert!(err.contains("at most 10"));
}
#[test]
fn validates_subagent_task_text_size_caps() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let intent_err = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "x".repeat(MAX_SUBAGENT_TASK_INTENT_BYTES + 1),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(provider.clone(), temp.path()),
)
.unwrap_err()
.to_string();
assert!(
intent_err.contains(&format!(
"intent exceeds {MAX_SUBAGENT_TASK_INTENT_BYTES} bytes"
)),
"{intent_err}"
);
let context_err = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "within cap".to_string(),
agent: None,
identity: None,
context: Some("x".repeat(MAX_SUBAGENT_TASK_CONTEXT_BYTES + 1)),
cwd: None,
}],
},
config(provider, temp.path()),
)
.unwrap_err()
.to_string();
assert!(
context_err.contains(&format!(
"context exceeds {MAX_SUBAGENT_TASK_CONTEXT_BYTES} bytes"
)),
"{context_err}"
);
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "x".repeat(MAX_SUBAGENT_TASK_INTENT_BYTES),
agent: None,
identity: None,
context: Some("x".repeat(MAX_SUBAGENT_TASK_CONTEXT_BYTES)),
cwd: None,
}],
}
.validated_concurrency()
.unwrap();
}
#[test]
fn rejects_empty_task_intent() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let err = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![SubagentTask {
intent: " ".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(provider, temp.path()),
)
.unwrap_err()
.to_string();
assert!(err.contains("intent must not be empty"));
}
#[test]
fn rejects_out_of_range_concurrency() {
for concurrency in [0, 5] {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let err = run_subagents(
SubagentsArgs {
concurrency: Some(concurrency),
tasks: vec![SubagentTask {
intent: "one".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(provider, temp.path()),
)
.unwrap_err()
.to_string();
assert!(err.contains("concurrency must be between 1 and 4"), "{err}");
}
}