use super::*;
use super::{
activity::{SubagentActivitySink, subagent_task_activity_metadata},
cwd::resolve_child_cwd,
dto::{MAX_SUBAGENT_TASK_CONTEXT_BYTES, MAX_SUBAGENT_TASK_INTENT_BYTES},
scheduler::{
PreparedSubagentTask, SCHEDULER_EVENT_CHANNEL_BOUND, SchedulerTaskReporter,
SchedulerWaitConfig, SharedTaskQueue, SharedTaskResults, SubagentScheduler, WorkerEvent,
output_from_results, send_terminal_worker_event,
},
worker::{
SUBAGENT_PROVIDER_STREAM_NO_SEMANTIC_PROGRESS_TIMEOUT, SUBAGENT_RESULT_ERROR_CHAR_LIMIT,
SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT, SUBAGENT_TRUNCATION_MARKER,
changed_files_from_tool_results, child_agent_for_task, failed_result,
failed_result_with_session_and_output, failed_subagent_session_snapshot, subagent_prompt,
truncate_string_field, truncate_subagents_output,
},
};
use crate::{
agent::{AgentOutputSink, AgentSession, cancellation::AgentCancellation},
config::{HookDefinition, HookFailurePolicy, HookSettings},
hooks::HookRuntime,
instructions::{InstructionFile, InstructionSourceKind},
output::{
ActivityEvent, ActivityId, ActivityKind, ActivitySender, ActivityStatus, OutputEvent,
},
providers::{Provider, ProviderEvent, ProviderRequest, ToolCall, Usage},
sessions::{SessionEvent, SessionManager},
skills::{Skill, SkillDiscovery, filter_enabled_skills},
tools::{ToolResult, ToolRuntime},
};
use serde_json::json;
use std::{
collections::{BTreeMap, BTreeSet},
path::{Path, PathBuf},
sync::atomic::{AtomicBool, AtomicUsize, Ordering},
sync::{Arc, Mutex},
thread,
time::Duration,
};
#[test]
fn subagent_activity_sink_emits_context_usage_for_task() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-g1"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-g1/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation: AgentCancellation::default(),
};
sink.output_event(OutputEvent::ContextUsage {
current_tokens: 12_345,
max_tokens: 128_000,
reasoning_tokens: Some(99),
source: crate::output::ContextUsageSource::ProviderExact,
request_sequence: 7,
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
&events[0],
ActivityEvent::UsageUpdate {
id,
current_tokens: 12_345,
max_tokens: 128_000,
reasoning_tokens: Some(99),
source: crate::output::ContextUsageSource::ProviderExact,
request_sequence: 7,
} if id.as_str() == "task-g1"
));
}
#[test]
fn subagent_activity_sink_delivers_only_terminal_canceled_after_cancel() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| captured.lock().unwrap().push(event));
let (cancellation, handle) = AgentCancellation::default().child_token();
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-cancel"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-cancel/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation,
};
sink.assistant_delta("partial").unwrap();
handle.cancel();
sink.assistant_delta("late").unwrap();
sink.finish_assistant(ActivityStatus::Failed);
sink.finish_assistant(ActivityStatus::Failed);
sink.activity_event(ActivityEvent::Delta {
id: ActivityId::new("tool-cancel"),
preview: "late tool output".to_string(),
})
.unwrap();
sink.activity_event(ActivityEvent::Finished {
id: ActivityId::new("tool-cancel"),
status: ActivityStatus::Canceled,
metadata: None,
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 4);
assert!(matches!(
&events[0],
ActivityEvent::Started {
status: ActivityStatus::Running,
..
}
));
assert!(matches!(
&events[1],
ActivityEvent::Delta { preview, .. } if preview == "partial"
));
assert!(matches!(
&events[2],
ActivityEvent::Finished {
status: ActivityStatus::Canceled,
metadata: None,
..
}
));
assert!(matches!(
&events[3],
ActivityEvent::Finished {
id,
status: ActivityStatus::Canceled,
metadata: None,
} if id.as_str() == "tool-cancel"
));
}
#[test]
fn display_safety_preserves_tool_start_pair_during_cancellation() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| captured.lock().unwrap().push(event));
let (cancellation, handle) = AgentCancellation::default().child_token();
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-cancel-tool"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-cancel-tool/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation,
};
let call = ToolCall {
id: "child-tool".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"printf retained","timeout":5}),
};
handle.cancel();
sink.activity_event(ActivityEvent::Started {
id: ActivityId::new("child-tool"),
parent_id: Some(ActivityId::new("task-cancel-tool")),
kind: ActivityKind::Tool,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new("bash printf retained"),
})
.unwrap();
sink.output_event(OutputEvent::ToolStarted {
call: Box::new(call),
label: "bash printf retained".to_string(),
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 2);
assert!(matches!(
&events[0],
ActivityEvent::Started {
id,
parent_id: Some(parent_id),
kind: ActivityKind::Tool,
..
} if id.as_str() == "child-tool" && parent_id.as_str() == "task-cancel-tool"
));
assert!(matches!(
&events[1],
ActivityEvent::ToolStartedDetail { id, detail }
if id.as_str() == "child-tool"
&& detail.tool_name.as_ref() == "bash"
&& detail.params == json!({"command":"printf retained","timeout":5})
));
}
#[test]
fn round_four_reconciles_nested_tool_result_after_cancellation() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| captured.lock().unwrap().push(event));
let (cancellation, handle) = AgentCancellation::default().child_token();
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-cancel-result"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-cancel-result/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation,
};
let call = ToolCall {
id: "child-result".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"printf done"}),
};
sink.activity_event(ActivityEvent::Started {
id: ActivityId::new("child-result"),
parent_id: Some(ActivityId::new("task-cancel-result")),
kind: ActivityKind::Tool,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new("bash printf done"),
})
.unwrap();
handle.cancel();
sink.output_event(OutputEvent::ToolResult {
call: Box::new(call),
result: Box::new(ToolResult {
tool_name: "bash".to_string(),
success: true,
content: "done".to_string(),
metadata: json!({"exit_code":0}),
display: crate::tools::ToolResultDisplay::default(),
}),
summary: Box::new(crate::output::ToolDisplaySummary {
tool_name: "bash".to_string(),
label: "bash printf done".to_string(),
status: crate::output::ToolStatus::Success,
unicode_mark: "✓",
ascii_mark: "OK",
metadata: Vec::new(),
}),
})
.unwrap();
sink.activity_event(ActivityEvent::Finished {
id: ActivityId::new("child-result"),
status: ActivityStatus::Success,
metadata: None,
})
.unwrap();
let events = events.lock().unwrap();
assert!(matches!(
&events[1],
ActivityEvent::ToolResultDetail { id, detail }
if id.as_str() == "child-result"
&& detail.status == ActivityStatus::Success
&& detail.output.as_ref() == "done"
));
assert!(matches!(
&events[2],
ActivityEvent::Finished { id, status: ActivityStatus::Success, .. }
if id.as_str() == "child-result"
));
}
#[test]
fn subagent_activity_sink_groups_consecutive_reasoning_summary_lines() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-g1"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-g1/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation: AgentCancellation::default(),
};
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "checked inputs\nverified output".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::ContextUsage {
current_tokens: 12,
max_tokens: 100,
reasoning_tokens: None,
source: crate::output::ContextUsageSource::ProviderExact,
request_sequence: 1,
})
.unwrap();
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "planned fix".to_string(),
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 3);
assert!(matches!(events[1], ActivityEvent::UsageUpdate { .. }));
assert!(matches!(
&events[0],
ActivityEvent::Started {
id,
parent_id: Some(parent_id),
kind: ActivityKind::Assistant,
status: ActivityStatus::Success,
metadata,
} if id.as_str() == "task-g1/reasoning/1"
&& parent_id.as_str() == "task-g1"
&& metadata.label == "reasoning summaries ×2"
&& metadata.detail.as_deref() == Some("1. checked inputs\n2. verified output")
));
assert!(matches!(
&events[2],
ActivityEvent::Started {
id,
parent_id: Some(parent_id),
kind: ActivityKind::Assistant,
status: ActivityStatus::Success,
metadata,
} if id.as_str() == "task-g1/reasoning/1"
&& parent_id.as_str() == "task-g1"
&& metadata.label == "reasoning summaries ×3"
&& metadata.detail.as_deref() == Some("1. checked inputs\n2. verified output\n3. planned fix")
));
}
#[test]
fn subagent_activity_sink_splits_reasoning_groups_on_activity_boundaries() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-g1"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-g1/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation: AgentCancellation::default(),
};
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "first".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::ToolStarted {
call: Box::new(ToolCall {
id: "child_tool".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"true","timeout":5}),
}),
label: "bash true".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "second\nthird".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantDelta {
text: "answer".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "fourth".to_string(),
})
.unwrap();
let events = events.lock().unwrap();
let reasoning = events
.iter()
.filter_map(|event| match event {
ActivityEvent::Started { id, metadata, .. } if id.as_str().contains("/reasoning/") => {
Some((id.as_str().to_string(), metadata.label.clone()))
}
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
reasoning,
vec![
(
"task-g1/reasoning/1".to_string(),
"reasoning summaries ×1".to_string()
),
(
"task-g1/reasoning/2".to_string(),
"reasoning summaries ×2".to_string()
),
(
"task-g1/reasoning/3".to_string(),
"reasoning summaries ×1".to_string()
),
]
);
}
#[test]
fn subagent_activity_sink_tool_result_final_preview() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-g1"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-g1/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation: AgentCancellation::default(),
};
let call = ToolCall {
id: "child_tool".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"false","timeout":5}),
};
let result = ToolResult {
tool_name: "bash".to_string(),
success: false,
content: "stderr: boom".to_string(),
metadata: json!({"exit_code":1}),
display: crate::tools::ToolResultDisplay::default(),
};
let summary = crate::output::tool_display_summary(&call, &result);
sink.output_event(OutputEvent::ToolResult {
call: Box::new(call),
result: Box::new(result),
summary: Box::new(summary),
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
&events[0],
ActivityEvent::ToolResultDetail { id, detail }
if id.as_str() == "child_tool"
&& detail.status == ActivityStatus::Failed
&& detail.output.as_ref() == "stderr: boom"
));
assert!(!matches!(events[0], ActivityEvent::Delta { .. }));
}
#[test]
fn subagent_activity_sink_edit_tool_result_final_preview_includes_diff() {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let sender: ActivitySender = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
let mut sink = SubagentActivitySink {
parent_id: ActivityId::new("task-g1"),
activity_sender: Some(sender),
assistant_id: ActivityId::new("task-g1/assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: None,
cancellation: AgentCancellation::default(),
};
let call = ToolCall {
id: "child_edit".to_string(),
name: "hash_edit".to_string(),
arguments: json!({"input":"[src/lib.rs#ABCD]\nSWAP 1.=1:\n+new line"}),
};
let result = ToolResult {
tool_name: "hash_edit".to_string(),
success: true,
content: "applied 1 edits".to_string(),
metadata: json!({"path":"src/lib.rs","edits":1}),
display: crate::tools::ToolResultDisplay {
edit_diff: Some(
"diff --git a/src/lib.rs b/src/lib.rs\n--- a/src/lib.rs\n+++ b/src/lib.rs\n@@ -1 +1 @@\n-old line\n+new line\n"
.to_string(),
),
},
};
let summary = crate::output::tool_display_summary(&call, &result);
sink.output_event(OutputEvent::ToolResult {
call: Box::new(call),
result: Box::new(result),
summary: Box::new(summary),
})
.unwrap();
let events = events.lock().unwrap();
assert_eq!(events.len(), 1);
assert!(matches!(
&events[0],
ActivityEvent::ToolResultDetail { id, detail }
if id.as_str() == "child_edit"
&& detail.status == ActivityStatus::Success
&& detail.output.as_ref() == "applied 1 edits"
&& detail.applied_diff.as_deref().is_some_and(|diff| diff.contains("diff --git"))
));
}
struct CountingProvider {
active: AtomicUsize,
max_active: AtomicUsize,
fail_on: String,
requests: Mutex<Vec<ProviderRequest>>,
}
impl CountingProvider {
fn new(fail_on: &str) -> Self {
Self {
active: AtomicUsize::new(0),
max_active: AtomicUsize::new(0),
fail_on: fail_on.to_string(),
requests: Mutex::new(Vec::new()),
}
}
}
impl Provider for CountingProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.requests.lock().unwrap().push(request.clone());
let now = self.active.fetch_add(1, Ordering::SeqCst) + 1;
self.max_active.fetch_max(now, Ordering::SeqCst);
thread::sleep(Duration::from_millis(20));
self.active.fetch_sub(1, Ordering::SeqCst);
let messages = request.messages();
let user = messages
.last()
.map(|m| m.content.as_str())
.unwrap_or_default();
if user.contains(&self.fail_on) {
anyhow::bail!("planned failure");
}
on_event(ProviderEvent::TextDelta(
user.lines().next().unwrap_or_default().to_string(),
))?;
on_event(ProviderEvent::Done)?;
Ok(())
}
}
struct TokenUsageProvider;
struct LargeOutputProvider;
impl Provider for LargeOutputProvider {
fn stream_cancellable(
&self,
_request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
on_event(ProviderEvent::TextDelta(
"x".repeat(SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT + 1024),
))?;
on_event(ProviderEvent::Done)?;
Ok(())
}
}
struct PartialThenFailProvider;
struct SchemaRetryProvider {
requests: Mutex<Vec<ProviderRequest>>,
always_invalid: bool,
invalid_output: String,
}
impl SchemaRetryProvider {
fn new(always_invalid: bool) -> Self {
Self::new_with_invalid_output(always_invalid, "not json")
}
fn new_with_invalid_output(always_invalid: bool, invalid_output: &str) -> Self {
Self {
requests: Mutex::new(Vec::new()),
always_invalid,
invalid_output: invalid_output.to_string(),
}
}
}
impl Provider for SchemaRetryProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let attempt = self.requests.lock().unwrap().len();
self.requests.lock().unwrap().push(request);
if self.always_invalid || attempt == 0 {
on_event(ProviderEvent::TextDelta(self.invalid_output.clone()))?;
} else {
on_event(ProviderEvent::TextDelta(
json!({
"phase": "IMPLEMENT",
"status": "COMPLETE",
"summary": "implemented",
"artifacts": [],
"verification": ["cargo test"],
"risks": [],
"changed_files": ["src/lib.rs"],
"tests_run": ["cargo test"],
"implementation_notes": ["fixed"]
})
.to_string(),
))?;
}
on_event(ProviderEvent::Done)?;
Ok(())
}
}
impl Provider for PartialThenFailProvider {
fn stream_cancellable(
&self,
_request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
on_event(ProviderEvent::TextDelta(
"checkpoint before failure".to_string(),
))?;
anyhow::bail!("planned partial failure")
}
}
impl Provider for TokenUsageProvider {
fn stream_cancellable(
&self,
_request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
on_event(ProviderEvent::Usage(Usage {
input: 40,
output: 2,
cache_read: 0,
cache_write: 0,
total: 123,
reasoning_tokens: None,
}))?;
on_event(ProviderEvent::TextDelta("done".to_string()))?;
on_event(ProviderEvent::Done)?;
Ok(())
}
}
fn config(provider: Arc<dyn Provider>, cwd: &Path) -> SubagentRunConfig {
SubagentRunConfig {
parent_agent: AgentSession::new("model", &[], &crate::skills::SkillDiscovery::default()),
provider,
provider_override: None,
parent_tools: ToolRuntime::new(cwd).unwrap(),
parent_cwd: cwd.to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: BTreeMap::new(),
subagent_profiles_prompt: None,
sessions_root: None,
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: None,
semantic_progress_timeout: SUBAGENT_PROVIDER_STREAM_NO_SEMANTIC_PROGRESS_TIMEOUT,
schema_validation_max_retries: 2,
}
}
fn completed_result(id: &str, total_tokens: Option<u64>) -> SubagentTaskResult {
SubagentTaskResult {
id: id.to_string(),
status: SubagentStatus::Completed,
intent: format!("intent {id}"),
agent: None,
identity: None,
cwd: PathBuf::from("."),
session_id: None,
session_path: None,
total_tokens,
changed_files: Vec::new(),
output: "done".to_string(),
structured_output: None,
output_truncated: false,
error: None,
}
}
struct WritingProvider {
requests: Mutex<Vec<ProviderRequest>>,
}
impl WritingProvider {
fn new() -> Self {
Self {
requests: Mutex::new(Vec::new()),
}
}
}
impl Provider for WritingProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let has_tool_result = !request.tool_results().is_empty();
self.requests.lock().unwrap().push(request);
if has_tool_result {
on_event(ProviderEvent::TextDelta("done".to_string()))?;
} else {
on_event(ProviderEvent::ToolCall(ToolCall {
id: "write_1".to_string(),
name: "write".to_string(),
arguments: json!({"path":"child.txt","content":"made by child"}),
}))?;
}
on_event(ProviderEvent::Done)?;
Ok(())
}
}
struct PolicyToolProvider {
requests: Mutex<Vec<ProviderRequest>>,
}
impl PolicyToolProvider {
fn new() -> Self {
Self {
requests: Mutex::new(Vec::new()),
}
}
}
impl Provider for PolicyToolProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let has_tool_result = !request.tool_results().is_empty();
self.requests.lock().unwrap().push(request);
if has_tool_result {
on_event(ProviderEvent::TextDelta("child done".to_string()))?;
} else {
on_event(ProviderEvent::ToolCall(ToolCall {
id: "write_policy_1".to_string(),
name: "write".to_string(),
arguments: json!({"path":"policy.txt","content":"target ran"}),
}))?;
}
on_event(ProviderEvent::Done)?;
Ok(())
}
}
fn profile(id: &str, prompt: &str) -> profiles::SubagentProfile {
profiles::SubagentProfile {
id: id.to_string(),
name: id.to_string(),
description: format!("{id} description"),
model: None,
reasoning: None,
path: PathBuf::from(format!("{id}.md")),
prompt: prompt.to_string(),
output_schema: None,
}
}
#[test]
fn subagent_schema_valid_output_sets_structured_output() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(SchemaRetryProvider::new(false));
let mut cfg = config(provider.clone(), temp.path());
cfg.profiles.insert(
"tars-code-writing-execution".to_string(),
profile("tars-code-writing-execution", "IMPLEMENT PROFILE"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "implement".into(),
agent: None,
identity: Some("tars-code-writing-execution".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Completed);
assert_eq!(
result.structured_output.as_ref().unwrap()["phase"],
"IMPLEMENT"
);
assert_eq!(provider.requests.lock().unwrap().len(), 2);
}
#[test]
fn subagent_schema_invalid_output_retries_with_feedback() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(SchemaRetryProvider::new(false));
let mut cfg = config(provider.clone(), temp.path());
cfg.profiles.insert(
"tars-code-writing-execution".to_string(),
profile("tars-code-writing-execution", "IMPLEMENT PROFILE"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "implement".into(),
agent: None,
identity: Some("tars-code-writing-execution".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.completed, 1);
let requests = provider.requests.lock().unwrap();
assert_eq!(requests.len(), 2);
let first_prompt = requests[0].messages().last().unwrap().content.clone();
assert!(
first_prompt.contains("Output schema contract"),
"{first_prompt}"
);
let retry_prompt = requests[1].messages().last().unwrap().content.clone();
assert!(
retry_prompt.contains("schema_validation_error"),
"{retry_prompt}"
);
assert!(
retry_prompt.contains("Return only a JSON object"),
"{retry_prompt}"
);
}
#[test]
fn subagent_schema_max_retry_exhaustion_fails() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(SchemaRetryProvider::new(true));
let mut cfg = config(provider.clone(), temp.path());
cfg.schema_validation_max_retries = 1;
cfg.profiles.insert(
"tars-code-writing-execution".to_string(),
profile("tars-code-writing-execution", "IMPLEMENT PROFILE"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "implement".into(),
agent: None,
identity: Some("tars-code-writing-execution".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Failed);
assert!(
result
.error
.as_deref()
.unwrap()
.contains("schema_validation_error")
);
assert_eq!(provider.requests.lock().unwrap().len(), 2);
}
#[test]
fn subagent_schema_exhaustion_does_not_return_raw_session_snapshot() {
let temp = tempfile::TempDir::new().unwrap();
let secret_output = "token=leak123";
let provider = Arc::new(SchemaRetryProvider::new_with_invalid_output(
true,
secret_output,
));
let mut cfg = config(provider.clone(), temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.schema_validation_max_retries = 1;
cfg.profiles.insert(
"tars-code-writing-execution".to_string(),
profile("tars-code-writing-execution", "IMPLEMENT PROFILE"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "implement".into(),
agent: None,
identity: Some("tars-code-writing-execution".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Failed);
let error = result.error.as_deref().unwrap();
assert!(error.contains("schema_validation_error"), "{error}");
assert!(result.output.is_empty(), "{}", result.output);
assert!(!result.output.contains(secret_output), "{}", result.output);
assert!(!result.output.contains("leak123"), "{}", result.output);
assert!(
result
.session_path
.as_ref()
.is_some_and(|path| path.exists())
);
assert_eq!(provider.requests.lock().unwrap().len(), 2);
}
#[test]
fn subagents_schema_retry_reuses_child_session() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(SchemaRetryProvider::new(false));
let mut cfg = config(provider.clone(), temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.profiles.insert(
"tars-code-writing-execution".to_string(),
profile("tars-code-writing-execution", "IMPLEMENT PROFILE"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "implement".into(),
agent: None,
identity: Some("tars-code-writing-execution".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
let session_path = result.session_path.as_ref().expect("session path");
let jsonl = std::fs::read_to_string(session_path).unwrap();
assert_eq!(result.status, SubagentStatus::Completed);
assert!(jsonl.contains("schema_validation_error"), "{jsonl}");
assert_eq!(provider.requests.lock().unwrap().len(), 2);
}
#[test]
fn run_subagents_caps_large_child_output_in_result() {
let temp = tempfile::TempDir::new().unwrap();
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "large child output".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(Arc::new(LargeOutputProvider), temp.path()),
)
.unwrap();
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Completed);
assert!(result.output_truncated);
assert!(result.output.chars().count() <= SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT);
assert!(
result
.output
.contains(SUBAGENT_TRUNCATION_MARKER.trim_start())
);
}
#[test]
fn failed_subagent_result_includes_partial_session_snapshot() {
let temp = tempfile::TempDir::new().unwrap();
let mut cfg = config(Arc::new(PartialThenFailProvider), temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "partial fail".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
assert_eq!(result.status, SubagentStatus::Failed);
assert!(
result.output.contains("checkpoint before failure"),
"{}",
result.output
);
assert!(
result
.session_path
.as_ref()
.is_some_and(|path| path.exists())
);
}
#[test]
fn failed_subagent_result_uses_authoritative_output_once() {
let temp = tempfile::TempDir::new().unwrap();
let session = SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let events = [
SessionEvent::new(
"assistant_chunk",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"generated assistant text"}),
),
SessionEvent::new(
"assistant_output",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"generated assistant text"}),
),
];
for event in &events {
session.append(event).unwrap();
}
let output = failed_subagent_session_snapshot(Some(&session)).unwrap();
let result = failed_result_with_session_and_output(
"child".to_string(),
SubagentTask {
intent: "failed child".to_string(),
agent: None,
identity: None,
context: None,
cwd: None,
},
temp.path().to_path_buf(),
Some(session.id().to_string()),
Some(session.path().to_path_buf()),
"planned failure".to_string(),
Some(output),
);
assert_eq!(result.output.matches("generated assistant text").count(), 1);
}
#[test]
fn failed_subagent_result_redacts_sensitive_error_text() {
let result = failed_result(
"g1".to_string(),
SubagentTask {
intent: "redact failure".to_string(),
agent: None,
identity: None,
context: None,
cwd: None,
},
PathBuf::from("."),
format!("child failed with api_key = sk-{}", "x".repeat(24)),
);
let error = result.error.as_deref().unwrap();
assert!(!error.contains(&format!("sk-{}", "x".repeat(24))));
assert!(error.contains("<redacted>"), "{error}");
}
#[test]
fn truncate_string_field_caps_tiny_limits() {
for limit in 0..=3 {
let mut value = "abcdef".to_string();
assert!(truncate_string_field(&mut value, limit));
assert!(
value.chars().count() <= limit,
"limit={limit} value={value:?}"
);
}
let mut value = "abcdef".to_string();
assert!(truncate_string_field(&mut value, 1));
assert_eq!(value, "\n");
let mut empty_limit = "abcdef".to_string();
assert!(truncate_string_field(&mut empty_limit, 0));
assert!(empty_limit.is_empty());
}
#[test]
fn child_agent_inherits_configured_markdown_before_identity_prompt() {
let temp = tempfile::TempDir::new().unwrap();
let instructions = vec![
InstructionFile {
kind: InstructionSourceKind::Repository,
path: PathBuf::from("/repo/AGENTS.md"),
content: "repo rules".to_string(),
},
InstructionFile {
kind: InstructionSourceKind::Configured,
path: PathBuf::from("/shared/team.md"),
content: "configured shared rules".to_string(),
},
];
let mut config = config(Arc::new(CountingProvider::new("never")), temp.path());
config.parent_agent = AgentSession::new(
"model",
&instructions,
&crate::skills::SkillDiscovery::default(),
);
config.profiles.insert(
"reviewer".to_string(),
profile("reviewer", "IDENTITY REVIEWER BODY"),
);
let task = SubagentTask {
intent: "inspect".to_string(),
agent: None,
identity: Some("reviewer".to_string()),
context: None,
cwd: None,
};
let child = child_agent_for_task(&task, temp.path(), &config).unwrap();
let prompt = child.agent.system_prompt();
let repo_index = prompt.find("repo rules").unwrap();
let configured_index = prompt.find("configured shared rules").unwrap();
let identity_index = prompt.find("IDENTITY REVIEWER BODY").unwrap();
assert!(repo_index < configured_index, "{prompt}");
assert!(configured_index < identity_index, "{prompt}");
assert!(prompt.contains("/shared/team.md"), "{prompt}");
}
struct NestedSubagentsProvider {
requests: Mutex<Vec<ProviderRequest>>,
respect_schema_hiding: bool,
}
struct ChildCancelingProvider {
cancel: Arc<AtomicBool>,
requests: Mutex<Vec<ProviderRequest>>,
}
struct DrainOnCancellationProvider {
entered: Mutex<Option<std::sync::mpsc::Sender<()>>>,
release: Mutex<std::sync::mpsc::Receiver<()>>,
observed_child_cancellation: AtomicBool,
requests: Mutex<Vec<ProviderRequest>>,
}
impl DrainOnCancellationProvider {
fn new(entered: std::sync::mpsc::Sender<()>, release: std::sync::mpsc::Receiver<()>) -> Self {
Self {
entered: Mutex::new(Some(entered)),
release: Mutex::new(release),
observed_child_cancellation: AtomicBool::new(false),
requests: Mutex::new(Vec::new()),
}
}
}
struct StallingAfterToolResultProvider {
requests: Mutex<Vec<ProviderRequest>>,
saw_continuation: AtomicBool,
saw_cancellation: AtomicBool,
attempted_post_cancel_event: AtomicBool,
unblock: AtomicBool,
}
impl StallingAfterToolResultProvider {
fn new() -> Self {
Self {
requests: Mutex::new(Vec::new()),
saw_continuation: AtomicBool::new(false),
saw_cancellation: AtomicBool::new(false),
attempted_post_cancel_event: AtomicBool::new(false),
unblock: AtomicBool::new(false),
}
}
fn unblock(&self) {
self.unblock.store(true, Ordering::SeqCst);
}
}
impl Provider for StallingAfterToolResultProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let has_tool_result = !request.tool_results().is_empty();
self.requests.lock().unwrap().push(request);
if has_tool_result {
self.saw_continuation.store(true, Ordering::SeqCst);
while !self.unblock.load(Ordering::SeqCst) {
if cancellation.is_canceled() {
self.saw_cancellation.store(true, Ordering::SeqCst);
self.attempted_post_cancel_event
.store(true, Ordering::SeqCst);
on_event(ProviderEvent::ToolCall(ToolCall {
id: "bash_after_cancel".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"printf unsafe","timeout":5}),
}))?;
return Ok(());
}
thread::sleep(Duration::from_millis(10));
}
on_event(ProviderEvent::Done)?;
} else {
on_event(ProviderEvent::ToolCall(ToolCall {
id: "bash_stall_1".to_string(),
name: "bash".to_string(),
arguments: json!({"command":"printf 'child bash ok'","timeout":5}),
}))?;
on_event(ProviderEvent::Done)?;
}
Ok(())
}
}
impl ChildCancelingProvider {
fn new(cancel: Arc<AtomicBool>) -> Self {
Self {
cancel,
requests: Mutex::new(Vec::new()),
}
}
}
impl Provider for ChildCancelingProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.requests.lock().unwrap().push(request);
on_event(ProviderEvent::TextDelta("before cancel".to_string()))?;
self.cancel.store(true, Ordering::SeqCst);
on_event(ProviderEvent::TextDelta("after cancel".to_string()))?;
on_event(ProviderEvent::Done)?;
Ok(())
}
}
impl Provider for DrainOnCancellationProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.requests.lock().unwrap().push(request);
if let Some(entered) = self.entered.lock().unwrap().take() {
let _ = entered.send(());
}
while !cancellation.is_canceled() {
thread::sleep(Duration::from_millis(5));
}
self.observed_child_cancellation
.store(true, Ordering::SeqCst);
self.release
.lock()
.unwrap()
.recv_timeout(Duration::from_secs(1))
.unwrap();
on_event(ProviderEvent::Done)?;
Ok(())
}
}
impl NestedSubagentsProvider {
fn new() -> Self {
Self {
requests: Mutex::new(Vec::new()),
respect_schema_hiding: true,
}
}
fn ignoring_schema_hiding() -> Self {
Self {
requests: Mutex::new(Vec::new()),
respect_schema_hiding: false,
}
}
}
impl Provider for NestedSubagentsProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let has_tool_result = !request.tool_results().is_empty();
self.requests.lock().unwrap().push(request.clone());
if self.respect_schema_hiding && !request.subagents_tool_enabled() {
on_event(ProviderEvent::TextDelta("schema hidden".to_string()))?;
} else if has_tool_result {
on_event(ProviderEvent::TextDelta("nested complete".to_string()))?;
} else {
on_event(ProviderEvent::ToolCall(ToolCall {
id: "nested_1".to_string(),
name: "subagents".to_string(),
arguments: json!({"tasks":[{"intent":"nested"}],"concurrency":1}),
}))?;
}
on_event(ProviderEvent::Done)?;
Ok(())
}
}
#[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();
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() == "bash_stall_1" && detail.status == ActivityStatus::Success)));
assert!(!events.iter().any(|event| matches!(event, ActivityEvent::FinalPreview { id, .. } if id.as_str() == "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();
}
#[cfg(unix)]
#[test]
fn inherited_subagent_hook_payload_uses_child_metadata() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(WritingProvider::new());
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.inherited_hooks = Some(
HookRuntime::new(
temp.path(),
HookSettings {
enabled: true,
before_tool: vec![HookDefinition {
label: Some("capture-child".into()),
command: "cat > child-hook.json".into(),
include_tools: vec!["write".into()],
..HookDefinition::default()
}],
..HookSettings::default()
},
true,
)
.unwrap(),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "write from child".into(),
agent: Some("reviewer".into()),
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.completed, 1);
let result = &output.results[0];
let session_path = result.session_path.as_ref().expect("child session path");
let payload: serde_json::Value = serde_json::from_str(
&std::fs::read_to_string(temp.path().join("child-hook.json")).unwrap(),
)
.unwrap();
assert_eq!(payload["context"]["invocation_mode"], "subagent");
assert_eq!(payload["context"]["subagent"], true);
assert_eq!(payload["context"]["agent_id"], "reviewer");
assert_eq!(
payload["context"]["session_id"],
result.session_id.as_deref().unwrap()
);
assert_eq!(
payload["context"]["session_path"],
session_path.display().to_string()
);
assert_eq!(payload["context"]["turn_id"], "turn-0");
assert!(
payload["context"]["message_id"]
.as_str()
.unwrap()
.contains("write_1")
);
assert_eq!(payload["affected_paths"][0]["path"], "child.txt");
}
#[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));
}
#[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 error = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "cancel child".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap_err();
assert!(
error.to_string().to_lowercase().contains("cancel"),
"{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 error = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "wait for cancellation".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap_err();
returned_in_thread.store(true, Ordering::SeqCst);
error
});
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 error = handle.join().unwrap();
assert!(
error.to_string().to_lowercase().contains("cancel"),
"{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 error = result.unwrap_err();
assert!(provider.observed_child_cancellation.load(Ordering::SeqCst));
let error_text = format!("{error:#}");
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 child_prompt_reuses_parent_agents_filtered_skills_and_tools() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let instructions = vec![InstructionFile {
kind: InstructionSourceKind::Repository,
path: PathBuf::from("AGENTS.md"),
content: "AGENTS.md rules".to_string(),
}];
let mut discovered_skills = SkillDiscovery::default();
for name in ["test-skill", "disabled-skill"] {
discovered_skills.skills.insert(
name.to_string(),
Skill {
name: name.to_string(),
path: PathBuf::from(format!(".agents/skills/{name}/SKILL.md")),
frontmatter: BTreeMap::new(),
body: String::new(),
},
);
}
let skills = filter_enabled_skills(
&discovered_skills,
&BTreeSet::from(["disabled-skill".to_string()]),
);
let agent = AgentSession::new("model", &instructions, &skills).with_context_budget(
crate::context::ContextBudget {
enabled: false,
..Default::default()
},
);
let cfg = SubagentRunConfig {
parent_agent: agent,
provider: provider.clone(),
provider_override: None,
parent_tools: ToolRuntime::new(temp.path()).unwrap(),
parent_cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: BTreeMap::new(),
subagent_profiles_prompt: None,
sessions_root: None,
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: None,
semantic_progress_timeout: SUBAGENT_PROVIDER_STREAM_NO_SEMANTIC_PROGRESS_TIMEOUT,
schema_validation_max_retries: 2,
};
run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![SubagentTask {
intent: "inspect".into(),
agent: Some("label".into()),
identity: None,
context: Some("ctx".into()),
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
assert!(
requests[0].messages()[0]
.content
.contains("AGENTS.md rules")
);
assert!(requests[0].messages()[0].content.contains("test-skill"));
assert!(!requests[0].messages()[0].content.contains("disabled-skill"));
assert!(requests[0].messages()[0].content.contains("subagents"));
assert!(
requests[0].messages()[1]
.content
.contains("Agent label/persona: label")
);
}
#[test]
fn subagent_with_subagents_tool_receives_available_identity_list() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "PROFILE BODY MUST NOT LEAK IN LIST"),
);
cfg.subagent_profiles_prompt = Some(
"<Subagent-Identities>\n- `frontend-dev` — UI work\n</Subagent-Identities>".to_string(),
);
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "inspect available identities".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert!(system.contains("<Subagent-Identities>"), "{system}");
assert!(system.contains("`frontend-dev` — UI work"), "{system}");
assert!(
!system.contains("PROFILE BODY MUST NOT LEAK IN LIST"),
"{system}"
);
}
#[test]
fn subagent_at_subagents_tool_depth_limit_does_not_receive_identity_list() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
cfg.parent_tools = ToolRuntime::new_with_settings(
temp.path(),
crate::tools::ToolSettings {
subagents: crate::tools::SubagentsToolSettings {
max_depth: 2,
..crate::tools::SubagentsToolSettings::default()
},
..crate::tools::ToolSettings::default()
},
)
.unwrap();
cfg.depth = 1;
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "PROFILE BODY"),
);
cfg.subagent_profiles_prompt = Some(
"<Subagent-Identities>\n- `frontend-dev` — UI work\n</Subagent-Identities>".to_string(),
);
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "at max depth".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert!(!system.contains("<Subagent-Identities>"), "{system}");
assert!(!requests[0].subagents_tool_enabled());
}
#[test]
fn selected_identity_receives_profile_prompt_and_available_identity_list() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "PROFILE BODY ONLY FOR FRONTEND"),
);
cfg.subagent_profiles_prompt = Some(
"<Subagent-Identities>\n- `frontend-dev` — UI work\n</Subagent-Identities>".to_string(),
);
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "inspect ui".into(),
agent: None,
identity: Some("frontend-dev".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert!(
system.contains("PROFILE BODY ONLY FOR FRONTEND"),
"{system}"
);
assert!(system.contains("<Subagent-Identities>"), "{system}");
assert!(system.contains("`frontend-dev` — UI work"), "{system}");
}
#[test]
fn selected_identity_appends_profile_prompt_to_child_system_only() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "PROFILE BODY ONLY FOR FRONTEND"),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "inspect ui".into(),
agent: Some("legacy label".into()),
identity: Some("frontend-dev".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.results[0].identity.as_deref(), Some("frontend-dev"));
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert!(system.contains("PROFILE BODY ONLY FOR FRONTEND"));
assert!(
!requests[0].messages()[1]
.content
.contains("PROFILE BODY ONLY FOR FRONTEND")
);
assert!(
requests[0].messages()[1]
.content
.contains("Agent label/persona: legacy label")
);
}
#[test]
fn selected_identity_model_and_reasoning_override_route_child_provider_request() {
let temp = tempfile::TempDir::new().unwrap();
let paths = crate::config::McPaths::from_root(temp.path().join("mc"));
let mut settings = crate::config::Settings::default();
settings.custom_providers.insert(
"local-ai".to_string(),
crate::config::CustomProviderConfig {
label: "Local AI".to_string(),
base_url: "http://localhost:8080/v1".to_string(),
api_key_env_var: None,
models_dev_provider: None,
use_responses_endpoint: false,
supports_text_verbosity: false,
reasoning_protocol: crate::config::CustomReasoningProtocol::default(),
extra_models: Vec::new(),
},
);
crate::config::write_settings(&paths, &settings).unwrap();
let mut entry = crate::model_catalog::ModelCatalogEntry::new("local-ai", "gpt-test");
entry.reasoning_efforts = Some(crate::thinking::ThinkingLevel::EFFORT_GENERIC.to_vec());
crate::model_catalog::write_catalog_cache_for_configured_provider(&paths, "local-ai", &[entry])
.unwrap();
let parent_provider = Arc::new(CountingProvider::new("never"));
let override_provider = Arc::new(CountingProvider::new("never"));
let resolver_provider = override_provider.clone();
let expected_cwd = temp.path().canonicalize().unwrap();
let mut cfg = config(parent_provider.clone(), temp.path());
cfg.provider_override = Some(SubagentProviderOverride::new(
paths,
Arc::new(move |selection, cwd| {
assert_eq!(cwd.canonicalize().unwrap(), expected_cwd);
assert_eq!(selection.provider, "local-ai");
assert_eq!(selection.model, "gpt-test");
let provider: Arc<dyn Provider> = resolver_provider.clone();
Ok(crate::subagents::ResolvedProviderOverride {
provider,
scope: crate::thinking::ThinkingCapabilityScope::Custom(
crate::config::CustomReasoningProtocol::GptLike,
),
})
}),
));
let mut frontend = profile("frontend-dev", "PROFILE BODY ONLY FOR FRONTEND");
frontend.model = Some(profiles::SubagentModelOverride {
provider: "local-ai".to_string(),
model: "gpt-test".to_string(),
});
frontend.reasoning = Some(crate::thinking::ThinkingLevel::High);
cfg.profiles.insert("frontend-dev".to_string(), frontend);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "inspect ui".into(),
agent: None,
identity: Some("frontend-dev".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.results[0].status, SubagentStatus::Completed);
assert!(parent_provider.requests.lock().unwrap().is_empty());
let requests = override_provider.requests.lock().unwrap();
assert_eq!(requests.len(), 1);
let request = &requests[0];
assert_eq!(request.model, "gpt-test");
assert_eq!(request.thinking_level, crate::thinking::ThinkingLevel::High);
assert!(request.send_default_reasoning_summary());
assert!(
request.messages()[0]
.content
.contains("PROFILE BODY ONLY FOR FRONTEND")
);
}
#[test]
fn selected_primary_agent_is_excluded_from_subagent_identity_child_prompt() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
let primary = crate::primary_agents::PrimaryAgentProfile {
id: "tars".to_string(),
name: "TARS".to_string(),
description: "Tactical".to_string(),
path: PathBuf::new(),
prompt: "PRIMARY AGENT SECRET BODY MUST NOT LEAK".to_string(),
};
let main_agent = crate::agent::runner::append_primary_agent_to_main_prompt(
cfg.parent_agent.clone(),
Some(&primary),
);
assert!(
main_agent
.system_prompt()
.contains("PRIMARY AGENT SECRET BODY")
);
assert!(
!cfg.parent_agent
.system_prompt()
.contains("PRIMARY AGENT SECRET BODY")
);
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "SUBAGENT IDENTITY BODY"),
);
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "inspect ui".into(),
agent: None,
identity: Some("frontend-dev".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert!(system.contains("SUBAGENT IDENTITY BODY"), "{system}");
assert!(!system.contains("PRIMARY AGENT SECRET BODY"), "{system}");
assert!(
!requests[0].messages()[1]
.content
.contains("PRIMARY AGENT SECRET BODY")
);
}
#[test]
fn no_identity_preserves_base_child_prompt_without_profiles_or_discovery_metadata() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
let base_system_prompt = cfg.parent_agent.system_prompt().to_string();
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "PROFILE BODY MUST NOT LEAK"),
);
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "generic".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let requests = provider.requests.lock().unwrap();
let system = &requests[0].messages()[0].content;
assert_eq!(system, &base_system_prompt);
assert!(!system.contains("PROFILE BODY MUST NOT LEAK"));
assert!(!system.contains("Discoverable subagent identities"));
}
#[test]
fn unknown_identity_fails_without_provider_call() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "generic".into(),
agent: None,
identity: Some("missing".into()),
context: None,
cwd: None,
}],
},
config(provider.clone(), temp.path()),
)
.unwrap();
assert_eq!(output.summary.failed, 1);
assert_eq!(output.results[0].identity.as_deref(), Some("missing"));
assert!(
output.results[0]
.error
.as_deref()
.unwrap()
.contains("identity 'missing' is unavailable")
);
assert!(provider.requests.lock().unwrap().is_empty());
}
#[test]
fn failed_subagent_emits_one_failed_task_finish() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
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.activity_sender = Some(Arc::new(move |event| captured.lock().unwrap().push(event)));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "generic".into(),
agent: None,
identity: Some("missing".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.failed, 1);
assert!(provider.requests.lock().unwrap().is_empty());
let events = events.lock().unwrap();
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::Started { id, status: ActivityStatus::Running, .. } if id.as_str() == "subagents/g1")));
}
#[cfg(unix)]
#[test]
fn discover_subagent_profiles_rejects_symlinked_root() {
let temp = tempfile::TempDir::new().unwrap();
let real = temp.path().join("real");
std::fs::create_dir(&real).unwrap();
std::fs::write(
real.join("frontend.md"),
"---\nname: Frontend\ndescription: UI work\n---\nPrompt body\n",
)
.unwrap();
let linked = temp.path().join("linked");
std::os::unix::fs::symlink(&real, &linked).unwrap();
let discovery = profiles::discover_subagent_profiles(&linked);
assert!(discovery.profiles.is_empty(), "{:?}", discovery.profiles);
assert_eq!(discovery.diagnostics.len(), 1);
assert!(
discovery.diagnostics[0].message.contains("symlink"),
"{}",
discovery.diagnostics[0].message
);
}
#[test]
fn discover_subagent_profiles_missing_root_is_empty() {
let temp = tempfile::TempDir::new().unwrap();
let discovery = profiles::discover_subagent_profiles(&temp.path().join("missing"));
assert!(discovery.profiles.is_empty());
assert!(
discovery.diagnostics.is_empty(),
"{:?}",
discovery.diagnostics
);
}
#[test]
fn parent_identity_metadata_can_be_appended_without_exposing_bodies_or_mutating_base() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(
temp.path().join("frontend.md"),
"---\nname: Frontend\ndescription: UI work\n---\nSECRET PROFILE BODY\n",
)
.unwrap();
let discovery = profiles::discover_subagent_profiles(temp.path());
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let base_prompt = agent.system_prompt().to_string();
let parent = agent.with_appended_system_prompt(
&profiles::render_subagent_profiles_prompt(None, &discovery)
.unwrap()
.unwrap(),
);
assert_eq!(agent.system_prompt(), base_prompt);
assert!(parent.system_prompt().contains("frontend"));
assert!(parent.system_prompt().contains("UI work"));
assert!(!parent.system_prompt().contains("SECRET PROFILE BODY"));
assert!(
!agent
.system_prompt()
.contains("Discoverable subagent identities")
);
}
#[test]
fn dispatch_returns_structured_json() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let result = dispatch_subagents(
json!({"tasks":[{"intent":"one","identity":"frontend-dev"}],"concurrency":1}),
{
let mut cfg = config(provider, temp.path());
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "Frontend prompt"),
);
cfg
},
);
assert!(result.success, "{}", result.content);
assert!(result.content.contains("\"id\": \"g1\""));
assert!(result.content.contains("\"identity\": \"frontend-dev\""));
}
#[test]
fn dispatch_truncates_large_result_fields_without_changing_json_schema() {
let mut output = SubagentsOutput {
summary: SubagentsSummary {
total: 2,
completed: 1,
failed: 1,
total_tokens: None,
},
results: vec![
SubagentTaskResult {
id: "g1".to_string(),
status: SubagentStatus::Completed,
intent: "large output".to_string(),
agent: None,
identity: None,
cwd: PathBuf::from("/tmp"),
session_id: Some("child-session-1".to_string()),
session_path: Some(PathBuf::from("/tmp/subagents/child-session-1.jsonl")),
total_tokens: None,
changed_files: vec![PathBuf::from("src/lib.rs")],
output: "x".repeat(SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT + 100),
structured_output: None,
output_truncated: false,
error: None,
},
SubagentTaskResult {
id: "g2".to_string(),
status: SubagentStatus::Failed,
intent: "large error".to_string(),
agent: None,
identity: None,
cwd: PathBuf::from("/tmp"),
session_id: Some("child-session-2".to_string()),
session_path: Some(PathBuf::from("/tmp/subagents/child-session-2.jsonl")),
total_tokens: None,
changed_files: vec![PathBuf::from("src/main.rs")],
output: String::new(),
structured_output: None,
output_truncated: false,
error: Some("e".repeat(SUBAGENT_RESULT_ERROR_CHAR_LIMIT + 100)),
},
],
};
truncate_subagents_output(&mut output);
let json = serde_json::to_string_pretty(&output).unwrap();
let value: Value = serde_json::from_str(&json).unwrap();
assert_eq!(value["summary"]["total"], 2);
assert_eq!(value["results"].as_array().unwrap().len(), 2);
assert_eq!(value["results"][0]["id"], "g1");
assert_eq!(value["results"][0]["status"], "completed");
assert_eq!(
value["results"][0]["session_path"],
"/tmp/subagents/child-session-1.jsonl"
);
assert_eq!(value["results"][0]["changed_files"][0], "src/lib.rs");
assert_eq!(value["results"][0]["output_truncated"], true);
assert!(value["results"][0]["error"].is_null());
assert!(
value["results"][0]["output"]
.as_str()
.unwrap()
.contains(SUBAGENT_TRUNCATION_MARKER.trim_start())
);
assert_eq!(value["results"][1]["id"], "g2");
assert_eq!(value["results"][1]["status"], "failed");
assert_eq!(
value["results"][1]["session_path"],
"/tmp/subagents/child-session-2.jsonl"
);
assert_eq!(value["results"][1]["changed_files"][0], "src/main.rs");
assert_eq!(value["results"][1]["output_truncated"], true);
assert!(
value["results"][1]["error"]
.as_str()
.unwrap()
.contains(SUBAGENT_TRUNCATION_MARKER.trim_start())
);
}
#[test]
fn subagents_structured_output_serializes() {
let output = SubagentsOutput {
summary: SubagentsSummary {
total: 1,
completed: 1,
failed: 0,
total_tokens: None,
},
results: vec![SubagentTaskResult {
id: "g1".to_string(),
status: SubagentStatus::Completed,
intent: "structured".to_string(),
agent: None,
identity: Some("tars-code-writing-execution".to_string()),
cwd: PathBuf::from("."),
session_id: None,
session_path: None,
total_tokens: None,
changed_files: Vec::new(),
output: "{}".to_string(),
structured_output: Some(json!({"phase": "IMPLEMENT"})),
output_truncated: false,
error: None,
}],
};
let json = serde_json::to_string_pretty(&SerializableSubagentsOutput::from(&output)).unwrap();
let value: Value = serde_json::from_str(&json).unwrap();
assert_eq!(
value["results"][0]["structured_output"]["phase"],
"IMPLEMENT"
);
}
#[cfg(unix)]
#[test]
fn subagents_output_serializes_non_utf8_paths_lossily() {
use std::ffi::OsString;
use std::os::unix::ffi::OsStringExt;
let non_utf8 = PathBuf::from(OsString::from_vec(b"bad-\xFF-path".to_vec()));
let output = SubagentsOutput {
summary: SubagentsSummary {
total: 1,
completed: 1,
failed: 0,
total_tokens: None,
},
results: vec![SubagentTaskResult {
id: "g1".to_string(),
status: SubagentStatus::Completed,
intent: "non utf8 paths".to_string(),
agent: None,
identity: None,
cwd: non_utf8.clone(),
session_id: Some("session".to_string()),
session_path: Some(non_utf8.clone()),
total_tokens: None,
changed_files: vec![non_utf8],
output: "done".to_string(),
structured_output: None,
output_truncated: false,
error: None,
}],
};
let json = serde_json::to_string_pretty(&SerializableSubagentsOutput::from(&output)).unwrap();
let value: Value = serde_json::from_str(&json).unwrap();
assert_ne!(json, "{}");
assert_eq!(value["summary"]["total"], 1);
assert!(value["results"][0]["cwd"].as_str().unwrap().contains('�'));
assert!(
value["results"][0]["session_path"]
.as_str()
.unwrap()
.contains('�')
);
assert!(
value["results"][0]["changed_files"][0]
.as_str()
.unwrap()
.contains('�')
);
}
#[test]
fn resolve_child_cwd_allows_absolute_by_default_and_false_rejects_outside() {
let parent = tempfile::TempDir::new().unwrap();
let outside = tempfile::TempDir::new().unwrap();
let allowed = resolve_child_cwd(parent.path(), Some(outside.path()), true).unwrap();
assert_eq!(allowed, outside.path().canonicalize().unwrap());
let absolute_err = resolve_child_cwd(parent.path(), Some(outside.path()), false)
.unwrap_err()
.to_string();
assert!(absolute_err.contains("tools.subagents.absolute_paths"));
let relative_err = resolve_child_cwd(parent.path(), Some(Path::new("..")), true)
.unwrap_err()
.to_string();
assert!(relative_err.contains("escapes parent cwd"));
}
#[test]
fn subagent_result_references_child_session_artifact() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
let output = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![SubagentTask {
intent: "artifact".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let result = &output.results[0];
assert!(result.session_id.is_some());
let session_path = result.session_path.as_ref().unwrap();
assert!(session_path.starts_with(temp.path().join("sessions/subagents")));
assert!(session_path.exists());
assert!(
std::fs::read_to_string(session_path)
.unwrap()
.contains("user_input")
);
}
#[test]
fn subagent_result_includes_changed_files_from_child_write() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(WritingProvider::new());
let output = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![SubagentTask {
intent: "write child file".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(provider.clone(), temp.path()),
)
.unwrap();
let changed_file = temp.path().join("child.txt").canonicalize().unwrap();
assert_eq!(
std::fs::read_to_string(&changed_file).unwrap(),
"made by child"
);
assert_eq!(output.results[0].changed_files, vec![changed_file]);
let requests = provider.requests.lock().unwrap();
assert_eq!(requests.len(), 2);
assert!(
requests
.iter()
.all(|request| request.semantic_progress_timeout()
== Some(SUBAGENT_PROVIDER_STREAM_NO_SEMANTIC_PROGRESS_TIMEOUT))
);
}
#[cfg(unix)]
#[test]
fn inherited_hooks_run_child_write_with_child_cwd_and_session_only() {
let temp = tempfile::TempDir::new().unwrap();
let child_dir = temp.path().join("child");
std::fs::create_dir(&child_dir).unwrap();
let provider = Arc::new(WritingProvider::new());
let runtime = HookRuntime::new(
temp.path(),
HookSettings {
enabled: true,
before_tool: vec![HookDefinition {
label: Some("child-cwd-capture".into()),
command: "cat > hook-payload.json".into(),
include_tools: vec!["write".into()],
..HookDefinition::default()
}],
..HookSettings::default()
},
true,
)
.unwrap();
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.inherited_hooks = Some(runtime);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "write child file".into(),
agent: None,
identity: None,
context: None,
cwd: Some(PathBuf::from("child")),
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.failed, 0);
assert!(!temp.path().join("hook-payload.json").exists());
let payload_text = std::fs::read_to_string(child_dir.join("hook-payload.json")).unwrap();
let payload: Value = serde_json::from_str(&payload_text).unwrap();
let expected_child_cwd = child_dir.canonicalize().unwrap().display().to_string();
assert_eq!(payload["cwd"].as_str(), Some(expected_child_cwd.as_str()));
assert_eq!(
payload["context"]["cwd"].as_str(),
Some(expected_child_cwd.as_str())
);
assert_eq!(payload["context"]["invocation_mode"], "subagent");
assert_eq!(payload["context"]["subagent"], true);
assert!(payload["context"]["session_id"].as_str().is_some());
let session_path = payload["context"]["session_path"].as_str().unwrap();
assert!(
session_path.contains("/sessions/subagents/"),
"{session_path}"
);
assert!(session_path.ends_with(".jsonl"), "{session_path}");
let session_jsonl =
std::fs::read_to_string(output.results[0].session_path.as_ref().unwrap()).unwrap();
assert!(session_jsonl.contains("hook_lifecycle"), "{session_jsonl}");
assert!(
session_jsonl.contains("child-cwd-capture"),
"{session_jsonl}"
);
let parent_result_content = serde_json::to_string(&output).unwrap();
assert!(!parent_result_content.contains("hook_lifecycle"));
assert!(!parent_result_content.contains("child-cwd-capture"));
}
#[cfg(unix)]
#[test]
fn absent_inherited_hooks_leave_child_session_without_hook_records() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(WritingProvider::new());
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.inherited_hooks = None;
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "write child file".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let session_jsonl =
std::fs::read_to_string(output.results[0].session_path.as_ref().unwrap()).unwrap();
assert!(!session_jsonl.contains("hook_lifecycle"), "{session_jsonl}");
assert!(
!session_jsonl.contains("hook_diagnostic"),
"{session_jsonl}"
);
}
#[cfg(unix)]
#[test]
fn inherited_hook_failure_policies_preserve_local_records_and_parent_safety() {
for (policy, target_runs, failed) in [
(HookFailurePolicy::Ignore, true, false),
(HookFailurePolicy::Warn, true, false),
(HookFailurePolicy::Block, false, false),
(HookFailurePolicy::Fail, false, true),
] {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(PolicyToolProvider::new());
let runtime = HookRuntime::new(
temp.path(),
HookSettings {
enabled: true,
before_tool: vec![HookDefinition {
label: Some("CHILD_HOOK_LABEL_LEAK_MARKER".into()),
command: "printf CHILD_STDOUT_POISON; printf CHILD_STDERR_POISON >&2; exit 7"
.into(),
failure_policy: Some(policy),
include_tools: vec!["write".into()],
..HookDefinition::default()
}],
..HookSettings::default()
},
true,
)
.unwrap();
let mut cfg = config(provider.clone(), temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.inherited_hooks = Some(runtime);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: format!("policy {}", policy.as_str()),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(
output.results[0].status == SubagentStatus::Failed,
failed,
"{policy:?}"
);
assert_eq!(
temp.path().join("policy.txt").exists(),
target_runs,
"{policy:?}"
);
let session_jsonl =
std::fs::read_to_string(output.results[0].session_path.as_ref().unwrap()).unwrap();
assert!(
session_jsonl.contains("hook_lifecycle"),
"{policy:?}: {session_jsonl}"
);
assert!(
session_jsonl.contains("CHILD_HOOK_LABEL_LEAK_MARKER"),
"{policy:?}: {session_jsonl}"
);
if matches!(
policy,
HookFailurePolicy::Warn | HookFailurePolicy::Block | HookFailurePolicy::Fail
) {
assert!(
session_jsonl.contains("hook_diagnostic"),
"{policy:?}: {session_jsonl}"
);
}
let parent_result_content = serde_json::to_string(&output).unwrap();
for marker in [
"CHILD_HOOK_LABEL_LEAK_MARKER",
"CHILD_STDOUT_POISON",
"CHILD_STDERR_POISON",
"hook_lifecycle",
"hook_diagnostic",
] {
assert!(
!parent_result_content.contains(marker),
"{policy:?}: marker {marker} leaked in {parent_result_content}"
);
}
if policy == HookFailurePolicy::Fail {
assert_eq!(
output.results[0].error.as_deref(),
Some(
"child subagent stopped by local hook policy; inspect child session JSONL for local hook details"
)
);
}
}
}
#[cfg(unix)]
#[test]
fn inherited_hook_activity_is_local_only_and_respects_visibility_setting() {
for show_in_tui in [true, false] {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(WritingProvider::new());
let events: Arc<Mutex<Vec<ActivityEvent>>> = Arc::new(Mutex::new(Vec::new()));
let captured = Arc::clone(&events);
let runtime = HookRuntime::new(
temp.path(),
HookSettings {
enabled: true,
show_in_tui,
before_tool: vec![HookDefinition {
label: Some("CHILD_TUI_HOOK_MARKER".into()),
command: "true".into(),
include_tools: vec!["write".into()],
..HookDefinition::default()
}],
..HookSettings::default()
},
true,
)
.unwrap();
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
cfg.inherited_hooks = Some(runtime);
cfg.activity_sender = Some(Arc::new(move |event| captured.lock().unwrap().push(event)));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "visible hook activity".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let session_jsonl =
std::fs::read_to_string(output.results[0].session_path.as_ref().unwrap()).unwrap();
assert!(session_jsonl.contains("hook_lifecycle"));
let has_hook_activity = events.lock().unwrap().iter().any(|event| {
matches!(
event,
ActivityEvent::Started {
kind: ActivityKind::Hook,
..
}
)
});
assert_eq!(has_hook_activity, show_in_tui);
let parent_result_content = serde_json::to_string(&output).unwrap();
assert!(!parent_result_content.contains("CHILD_TUI_HOOK_MARKER"));
assert!(!parent_result_content.contains("hook_lifecycle"));
}
}
#[test]
fn subagent_task_activity_metadata_includes_prompt_detail() {
let task = SubagentTask {
intent: "inspect parent detail payload".into(),
agent: Some("issue-43-plan".into()),
identity: Some("frontend-dev".into()),
context: Some("first context line\nsecond context line".into()),
cwd: None,
};
let metadata = subagent_task_activity_metadata("g1", &task, 2);
let expected_prompt = subagent_prompt("g1", &task, None);
assert_eq!(
metadata.label,
"g1 · depth 2 · inspect parent detail payload"
);
assert_eq!(
metadata.fields,
vec![
("depth".to_string(), "2".to_string()),
("identity".to_string(), "frontend-dev".to_string()),
("agent".to_string(), "issue-43-plan".to_string()),
]
);
assert_eq!(metadata.detail.as_deref(), Some(expected_prompt.as_str()));
assert!(expected_prompt.contains("Subagent g1 task intent:\ninspect parent detail payload"));
assert!(expected_prompt.contains("Agent label/persona: issue-43-plan"));
assert!(expected_prompt.contains("Task context:\nfirst context line\nsecond context line"));
}
#[test]
fn subagent_task_activity_metadata_includes_identity_and_agent_when_present() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let events: Arc<Mutex<Vec<ActivityEvent>>> = Arc::new(Mutex::new(Vec::new()));
let captured = Arc::clone(&events);
let mut cfg = config(provider, temp.path());
cfg.profiles.insert(
"frontend-dev".to_string(),
profile("frontend-dev", "frontend prompt"),
);
cfg.activity_sender = Some(Arc::new(move |event| captured.lock().unwrap().push(event)));
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "one".into(),
agent: Some("reviewer".into()),
identity: Some("frontend-dev".into()),
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let events = events.lock().unwrap();
let task_metadata = events
.iter()
.find_map(|event| match event {
ActivityEvent::Started {
kind: ActivityKind::SubagentTask,
metadata,
..
} => Some(metadata),
_ => None,
})
.expect("subagent task start event");
assert_eq!(task_metadata.label, "g1 · depth 1 · one");
assert_eq!(
task_metadata.fields,
vec![
("depth".to_string(), "1".to_string()),
("identity".to_string(), "frontend-dev".to_string()),
("agent".to_string(), "reviewer".to_string()),
]
);
}
#[test]
fn child_tasks_attach_directly_to_parent_activity_without_synthetic_batch() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let events: Arc<Mutex<Vec<ActivityEvent>>> = Arc::new(Mutex::new(Vec::new()));
let captured = Arc::clone(&events);
let mut cfg = config(provider, temp.path());
cfg.parent_activity_id = Some(ActivityId::new("tool-parent"));
cfg.activity_sender = Some(Arc::new(move |event| captured.lock().unwrap().push(event)));
run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "one".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
let events = events.lock().unwrap();
assert!(!events.iter().any(|event| matches!(
event,
ActivityEvent::Started {
kind: ActivityKind::SubagentBatch,
..
}
)));
assert!(events.iter().any(|event| matches!(event, ActivityEvent::Started { id, parent_id: Some(parent), kind: ActivityKind::SubagentTask, .. } if id.as_str() == "tool-parent/g1" && parent.as_str() == "tool-parent")));
assert!(events.iter().any(|event| matches!(event, ActivityEvent::Finished { id, status: ActivityStatus::Success, .. } if id.as_str() == "tool-parent/g1")));
assert!(events.iter().any(|event| matches!(event, ActivityEvent::Finished { id, status: ActivityStatus::Success, .. } if id.as_str() == "tool-parent/g1/assistant")));
}
#[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}");
}
}
#[test]
fn child_session_creation_failure_is_reported() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider.clone(), temp.path());
let file_root = temp.path().join("not-a-dir");
std::fs::write(&file_root, "file").unwrap();
cfg.sessions_root = Some(file_root);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "artifact".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.failed, 1);
assert_eq!(provider.requests.lock().unwrap().len(), 0);
let error = output.results[0].error.as_deref().unwrap();
assert!(error.contains("session creation"), "{error}");
}
#[test]
fn changed_files_reports_successful_write_and_hash_edit_only_sorted_unique() {
let results = vec![
ToolResult {
tool_name: "bash".into(),
success: true,
content: String::new(),
metadata: json!({"path":"z.txt"}),
display: crate::tools::ToolResultDisplay::default(),
},
ToolResult {
tool_name: "write".into(),
success: false,
content: String::new(),
metadata: json!({"path":"b.txt"}),
display: crate::tools::ToolResultDisplay::default(),
},
ToolResult {
tool_name: "write".into(),
success: true,
content: String::new(),
metadata: json!({"path":"b.txt"}),
display: crate::tools::ToolResultDisplay::default(),
},
ToolResult {
tool_name: "hash_edit".into(),
success: true,
content: String::new(),
metadata: json!({"path":"a.txt"}),
display: crate::tools::ToolResultDisplay::default(),
},
ToolResult {
tool_name: "write".into(),
success: true,
content: String::new(),
metadata: json!({"path":"b.txt"}),
display: crate::tools::ToolResultDisplay::default(),
},
];
assert_eq!(
changed_files_from_tool_results(&results),
vec![PathBuf::from("a.txt"), PathBuf::from("b.txt")]
);
}
#[test]
fn nested_subagents_tool_call_is_rejected() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(NestedSubagentsProvider::ignoring_schema_hiding());
let mut cfg = config(provider.clone(), temp.path());
cfg.parent_tools = ToolRuntime::new_with_settings(
temp.path(),
crate::tools::ToolSettings {
subagents: crate::tools::SubagentsToolSettings {
max_depth: 1,
..crate::tools::SubagentsToolSettings::default()
},
..crate::tools::ToolSettings::default()
},
)
.unwrap();
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "launch nested".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.completed, 1);
let requests = provider.requests.lock().unwrap();
let all_tool_results = requests
.iter()
.flat_map(|request| request.tool_results())
.collect::<Vec<_>>();
let nested_result = all_tool_results
.iter()
.find(|result| result.call_id == "nested_1")
.expect("nested tool result");
assert!(!nested_result.success);
assert!(nested_result.output.contains("nesting limit"));
}
#[test]
fn default_max_depth_allows_one_nested_batch_and_hides_schema_at_limit() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(NestedSubagentsProvider::new());
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "launch nested".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
config(provider.clone(), temp.path()),
)
.unwrap();
assert_eq!(output.summary.completed, 1);
let requests = provider.requests.lock().unwrap();
assert!(
requests.iter().any(ProviderRequest::subagents_tool_enabled),
"child request before max depth should expose subagents"
);
assert!(
requests
.iter()
.any(|request| !request.subagents_tool_enabled()),
"grandchild request at max depth should hide subagents"
);
let all_tool_results = requests
.iter()
.flat_map(|request| request.tool_results())
.collect::<Vec<_>>();
assert!(all_tool_results.iter().any(|result| {
result.call_id == "nested_1" && result.success && result.output.contains("schema hidden")
}));
}
#[test]
fn nested_child_at_subagents_tool_depth_limit_does_not_inherit_identity_list() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(NestedSubagentsProvider::new());
let mut cfg = config(provider.clone(), temp.path());
cfg.subagent_profiles_prompt = Some(
"<Subagent-Identities>\n- `frontend-dev` — UI work\n</Subagent-Identities>".to_string(),
);
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "launch nested".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap();
assert_eq!(output.summary.completed, 1);
let requests = provider.requests.lock().unwrap();
let enabled_systems = requests
.iter()
.filter(|request| request.subagents_tool_enabled())
.map(|request| request.messages()[0].content.clone())
.collect::<Vec<_>>();
assert!(!enabled_systems.is_empty());
assert!(
enabled_systems
.iter()
.all(|system| system.contains("<Subagent-Identities>")),
"{enabled_systems:?}"
);
let disabled_systems = requests
.iter()
.filter(|request| !request.subagents_tool_enabled())
.map(|request| request.messages()[0].content.clone())
.collect::<Vec<_>>();
assert!(!disabled_systems.is_empty());
assert!(
disabled_systems
.iter()
.all(|system| !system.contains("<Subagent-Identities>")),
"{disabled_systems:?}"
);
}
#[test]
fn child_cwd_write_stays_inside_requested_child_directory() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::create_dir(temp.path().join("child-dir")).unwrap();
let provider = Arc::new(WritingProvider::new());
let output = run_subagents(
SubagentsArgs {
concurrency: Some(1),
tasks: vec![SubagentTask {
intent: "write child file".into(),
agent: None,
identity: None,
context: None,
cwd: Some(PathBuf::from("child-dir")),
}],
},
config(provider, temp.path()),
)
.unwrap();
let changed_file = temp
.path()
.join("child-dir/child.txt")
.canonicalize()
.unwrap();
assert_eq!(
std::fs::read_to_string(&changed_file).unwrap(),
"made by child"
);
assert_eq!(
output.results[0].cwd,
temp.path().join("child-dir").canonicalize().unwrap()
);
assert_eq!(output.results[0].changed_files, vec![changed_file]);
}
#[test]
fn resolve_child_cwd_rejects_missing_non_directory_and_symlink_escape() {
let parent = tempfile::TempDir::new().unwrap();
let missing = resolve_child_cwd(parent.path(), Some(Path::new("missing")), true)
.unwrap_err()
.to_string();
assert!(
missing.contains("No such file") || missing.contains("not found"),
"{missing}"
);
std::fs::write(parent.path().join("file-cwd"), "file").unwrap();
let file_err = resolve_child_cwd(parent.path(), Some(Path::new("file-cwd")), true)
.unwrap_err()
.to_string();
assert!(file_err.contains("not a directory"), "{file_err}");
#[cfg(unix)]
{
use std::os::unix::fs::symlink;
let outside = tempfile::TempDir::new().unwrap();
symlink(outside.path(), parent.path().join("escape-link")).unwrap();
let link_err = resolve_child_cwd(parent.path(), Some(Path::new("escape-link")), true)
.unwrap_err()
.to_string();
assert!(link_err.contains("escapes parent cwd"), "{link_err}");
}
}
#[test]
fn multiple_child_sessions_are_unique() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider, temp.path());
cfg.sessions_root = Some(temp.path().join("sessions"));
let output = run_subagents(
SubagentsArgs {
concurrency: Some(2),
tasks: 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,
},
],
},
cfg,
)
.unwrap();
let first = &output.results[0];
let second = &output.results[1];
assert!(first.session_id.is_some() && second.session_id.is_some());
assert_ne!(first.session_id, second.session_id);
assert_ne!(first.session_path, second.session_path);
for path in [
first.session_path.as_ref().unwrap(),
second.session_path.as_ref().unwrap(),
] {
assert!(path.starts_with(temp.path().join("sessions/subagents")));
assert!(path.exists());
}
}
#[test]
fn nesting_cap_rejects_recursive_batches() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(CountingProvider::new("never"));
let mut cfg = config(provider, temp.path());
cfg.depth = cfg.parent_tools.subagents_max_depth();
let err = run_subagents(
SubagentsArgs {
concurrency: None,
tasks: vec![SubagentTask {
intent: "one".into(),
agent: None,
identity: None,
context: None,
cwd: None,
}],
},
cfg,
)
.unwrap_err()
.to_string();
assert!(err.contains("nesting limit"));
}