pub(super) mod activity;
pub(super) mod config;
pub(super) mod cwd;
pub(super) mod dto;
pub(super) mod output_schema;
pub(crate) mod profiles;
pub(super) mod scheduler;
pub(super) mod worker;
#[cfg(test)]
mod tests;
pub use config::SubagentRunConfig;
pub(crate) use config::{ResolvedProviderOverride, SubagentProviderOverride};
#[allow(unused_imports)]
pub(crate) use dto::{
DEFAULT_SUBAGENT_CONCURRENCY, DEFAULT_SUBAGENT_MAX_DEPTH, MAX_SUBAGENT_CONCURRENCY,
MAX_SUBAGENT_MAX_DEPTH, MAX_SUBAGENT_TASKS, SubagentStatus, SubagentTask, SubagentTaskResult,
SubagentsArgs, SubagentsOutput, SubagentsSummary,
};
use crate::{
output::{ActivityEvent, ActivityId, ActivityKind, ActivityMetadata, ActivityStatus},
tools::{
ToolResult, ToolResultDisplay,
contract::{metadata_key as meta, tool_name},
},
};
use activity::emit_activity;
use scheduler::{SchedulerWaitConfig, SubagentScheduler};
use serde::Serialize;
use serde_json::{Value, json};
pub fn dispatch_subagents(arguments: Value, config: SubagentRunConfig) -> ToolResult {
match serde_json::from_value::<SubagentsArgs>(arguments)
.map_err(anyhow::Error::from)
.and_then(|args| run_subagents(args, config))
{
Ok(mut output) => {
worker::truncate_subagents_output(&mut output);
let serializable_output = SerializableSubagentsOutput::from(&output);
match serde_json::to_string_pretty(&serializable_output) {
Ok(content) => ToolResult {
tool_name: tool_name::SUBAGENTS.to_string(),
success: output.summary.failed == 0,
content,
metadata: json!({(meta::SUMMARY): output.summary}),
display: ToolResultDisplay::default(),
},
Err(error) => ToolResult {
tool_name: tool_name::SUBAGENTS.to_string(),
success: false,
content: format!("subagents output serialization failed: {error}"),
metadata: json!({}),
display: ToolResultDisplay::default(),
},
}
}
Err(error) => ToolResult {
tool_name: tool_name::SUBAGENTS.to_string(),
success: false,
content: error.to_string(),
metadata: json!({}),
display: ToolResultDisplay::default(),
},
}
}
#[derive(Serialize)]
struct SerializableSubagentsOutput<'a> {
summary: &'a SubagentsSummary,
results: Vec<SerializableSubagentTaskResult<'a>>,
}
#[derive(Serialize)]
struct SerializableSubagentTaskResult<'a> {
id: &'a str,
status: &'a SubagentStatus,
intent: &'a str,
agent: &'a Option<String>,
identity: &'a Option<String>,
cwd: String,
session_id: &'a Option<String>,
session_path: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
total_tokens: Option<u64>,
changed_files: Vec<String>,
output: &'a str,
output_truncated: bool,
error: &'a Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
structured_output: &'a Option<Value>,
}
impl<'a> From<&'a SubagentsOutput> for SerializableSubagentsOutput<'a> {
fn from(output: &'a SubagentsOutput) -> Self {
Self {
summary: &output.summary,
results: output
.results
.iter()
.map(SerializableSubagentTaskResult::from)
.collect(),
}
}
}
impl<'a> From<&'a SubagentTaskResult> for SerializableSubagentTaskResult<'a> {
fn from(result: &'a SubagentTaskResult) -> Self {
Self {
id: &result.id,
status: &result.status,
intent: &result.intent,
agent: &result.agent,
identity: &result.identity,
cwd: result.cwd.to_string_lossy().into_owned(),
session_id: &result.session_id,
session_path: result
.session_path
.as_ref()
.map(|path| path.to_string_lossy().into_owned()),
total_tokens: result.total_tokens,
changed_files: result
.changed_files
.iter()
.map(|path| path.to_string_lossy().into_owned())
.collect(),
output: &result.output,
output_truncated: result.output_truncated,
error: &result.error,
structured_output: &result.structured_output,
}
}
}
pub fn run_subagents(
args: SubagentsArgs,
config: SubagentRunConfig,
) -> anyhow::Result<SubagentsOutput> {
run_subagents_with_wait(args, config, SchedulerWaitConfig::default())
}
fn run_subagents_with_wait(
args: SubagentsArgs,
config: SubagentRunConfig,
wait: SchedulerWaitConfig,
) -> anyhow::Result<SubagentsOutput> {
if config.depth >= config.parent_tools.subagents_max_depth() {
anyhow::bail!(
"subagents nesting limit reached (depth {}, max_depth {})",
config.depth,
config.parent_tools.subagents_max_depth()
);
}
config.cancellation.check()?;
let concurrency = args.validated_concurrency()?.min(args.tasks.len());
let synthetic_batch = config.parent_activity_id.is_none();
let batch_id = config
.parent_activity_id
.clone()
.unwrap_or_else(|| ActivityId::new("subagents"));
if synthetic_batch {
emit_activity(
&config,
ActivityEvent::Started {
id: batch_id.clone(),
parent_id: None,
kind: ActivityKind::SubagentBatch,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new(format!("subagents ยท {} tasks", args.tasks.len())),
},
);
}
let mut scheduler =
SubagentScheduler::new(args.tasks, concurrency, config.clone(), batch_id.clone())?;
scheduler.wait = wait;
let result = scheduler.run();
if synthetic_batch {
emit_activity(
&config,
ActivityEvent::Finished {
id: batch_id,
status: match &result {
Ok(output) if output.summary.failed == 0 => ActivityStatus::Success,
Ok(_) => ActivityStatus::Failed,
Err(_) => ActivityStatus::Failed,
},
metadata: None,
},
);
}
result
}