use super::*;
use crate::multi_agent::capability::{CapabilityResolution, ChildToolCapability};
use crate::multi_agent::registry::AgentStatus;
impl MultiAgentRuntime {
pub async fn spawn_child(
&self,
name: &str,
system_prompt: String,
depth: i32,
full_permission: bool,
parent_messages: Vec<agent_base::ChatMessage>,
) -> Result<String, String> {
let config = ChildConfig {
system_prompt: Some(system_prompt),
full_permission: Some(full_permission),
..Default::default()
};
let prepared = self
.spawn_inner(name, depth, &config, parent_messages, None)
.await
.map_err(|e| match e {
AgentError::ConfigError(s) => s,
other => other.to_string(),
})?;
let path = prepared.path.clone();
self.spawn_ready(prepared);
Ok(path.to_string())
}
#[allow(clippy::too_many_arguments)]
async fn spawn_inner(
&self,
name: &str,
depth: i32,
config: &ChildConfig,
parent_messages: Vec<agent_base::ChatMessage>,
capability: Option<&ChildToolCapability>,
) -> Result<PreparedChild, AgentError> {
let path = AgentPath::root().join(name);
let ticket = self
.control
.budget()
.try_reserve_spawn()
.map_err(|e| AgentError::ConfigError(e.to_string()))?;
let recycled = self.recycle_predecessor(&path);
{
let mut registry = self.registry.lock().unwrap();
registry
.can_spawn(depth)
.map_err(|e| AgentError::ConfigError(e.to_string()))?;
registry
.register(&path, depth)
.map_err(|e| AgentError::ConfigError(e.to_string()))?;
}
let slot = match self.limiter.try_acquire() {
Ok(slot) => slot,
Err(e) => {
self.registry.lock().unwrap().close(&path);
return Err(AgentError::ConfigError(e.to_string()));
}
};
let child_mailbox = match self.mailbox.register(&path) {
Some(mb) => mb,
None => {
self.registry.lock().unwrap().close(&path);
return Err(AgentError::ConfigError(
"mailbox already exists".to_string(),
));
}
};
let (child_runtime, spawned_tools, resolution) = self
.build_child_runtime_with_config(
config,
self.spawn_permission(config.full_permission),
capability,
&path.to_string(),
)
.await
.map_err(|e| {
self.registry.lock().unwrap().close(&path);
self.mailbox.unregister(&path);
match e {
AgentError::ToolNotFound { .. } => e,
other => {
AgentError::ConfigError(format!("failed to build child runtime: {other}"))
}
}
})?;
let session_id = child_runtime.create_session().await;
self.prefill_child_session(&child_runtime, &session_id, &parent_messages)
.await
.map_err(|e| {
self.registry.lock().unwrap().close(&path);
self.mailbox.unregister(&path);
AgentError::ConfigError(format!("failed to prefill child session: {e}"))
})?;
let child_cancel = self.root_cancel.child_token();
{
let mut cancels = self.child_cancels.lock().unwrap();
cancels.insert(path.clone(), child_cancel.clone());
}
let generation = {
let mut gens = self.child_generations.lock().unwrap();
let g = gens.entry(path.clone()).or_insert(0);
*g += 1;
*g
};
Ok(PreparedChild {
path,
child_mailbox,
child_runtime,
session_id,
child_cancel,
slot,
ticket,
spawned_tools,
resolution,
generation,
recycled,
})
}
fn recycle_predecessor(&self, path: &AgentPath) -> bool {
let finished = {
let registry = self.registry.lock().unwrap();
matches!(
registry.get(path).map(|e| e.status()),
Some(AgentStatus::Done) | Some(AgentStatus::Closed)
)
};
if !finished {
return false;
}
{
let mut gens = self.child_generations.lock().unwrap();
*gens.entry(path.clone()).or_insert(0) += 1;
}
let old_token = self.child_cancels.lock().unwrap().remove(path);
if let Some(token) = old_token {
token.cancel();
}
self.registry.lock().unwrap().close(path);
self.mailbox.unregister(path);
self.write_gate.release_all(path.to_string().as_str());
self.spawned_tools.lock().unwrap().remove(&path.to_string());
true
}
fn spawn_ready(&self, prepared: PreparedChild) {
let PreparedChild {
path,
child_mailbox,
child_runtime,
session_id,
child_cancel,
slot,
ticket,
spawned_tools: _spawned_tools,
resolution: _resolution,
generation,
recycled: _recycled,
} = prepared;
let task_timeout = self.task_timeout;
let agent_path = path.clone();
let mailbox_for_task = self.mailbox.clone();
let registry_for_task = self.registry.clone();
let event_tx = self.event_tx.lock().unwrap().clone();
let write_gate_for_loop = Arc::clone(&self.write_gate);
let cleanup = ChildCleanup {
_slot: slot,
mailbox: self.mailbox.clone(),
registry: self.registry.clone(),
child_cancels: self.child_cancels.clone(),
path: agent_path.clone(),
write_gate: Arc::clone(&self.write_gate),
spawned_tools: Arc::clone(&self.spawned_tools),
child_generations: Arc::clone(&self.child_generations),
generation,
};
self.join_set.lock().unwrap().spawn(async move {
let _cleanup = cleanup;
run_child_loop(
child_mailbox,
child_runtime,
session_id,
agent_path,
mailbox_for_task,
registry_for_task,
event_tx,
child_cancel,
task_timeout,
write_gate_for_loop,
)
.await;
});
ticket.commit();
}
#[allow(clippy::too_many_arguments)] pub async fn spawn_child_with_history(
&self,
name: &str,
system_prompt: String,
full_permission: bool,
fork_history: Option<String>,
model: Option<String>,
capability: Option<ChildToolCapability>,
parent_session_id: &SessionId,
) -> Result<SpawnEcho, String> {
let parent_messages = self
.resolve_fork_history(fork_history, parent_session_id)
.await;
let config = ChildConfig {
system_prompt: Some(system_prompt),
full_permission: Some(full_permission),
model,
..Default::default()
};
self.spawn_with_config_forked(name.to_string(), config, parent_messages, capability)
.await
.map_err(|e| match e {
AgentError::ConfigError(s) => s,
other => other.to_string(),
})
.map(|spawned| SpawnEcho {
agent_path: spawned.path.to_string(),
registered_tools: spawned.spawned_tools.clone(),
degraded_reason: spawned.resolution().degraded_reason.clone(),
recycled: spawned.recycled,
})
}
pub fn child(self: &Arc<Self>) -> ChildBuilder {
ChildBuilder::new(Arc::clone(self))
}
#[cfg(test)]
pub(crate) async fn spawn_with_config(
&self,
name: String,
config: ChildConfig,
) -> Result<SpawnedChild, AgentError> {
self.spawn_with_config_forked(name, config, Vec::new(), None)
.await
}
pub(crate) async fn spawn_with_config_forked(
&self,
name: String,
config: ChildConfig,
parent_messages: Vec<agent_base::ChatMessage>,
capability: Option<ChildToolCapability>,
) -> Result<SpawnedChild, AgentError> {
if config.system_prompt.as_deref().unwrap_or("").is_empty() {
return Err(AgentError::ConfigError(
"ChildConfig.system_prompt is required (set it directly or use a preset)".into(),
));
}
let prepared = self
.spawn_inner(&name, 1, &config, parent_messages, capability.as_ref())
.await?;
let path = prepared.path.clone();
let spawned_tools = prepared.spawned_tools.clone();
let resolution = prepared.resolution.clone();
let prepared_recycled = prepared.recycled;
self.spawn_ready(prepared);
if spawned_tools.iter().any(|t| self.write_tools.contains(t)) {
self.spawned_tools
.lock()
.unwrap()
.insert(path.to_string(), spawned_tools.iter().cloned().collect());
}
Ok(SpawnedChild {
path,
spawned_tools,
resolution,
recycled: prepared_recycled,
})
}
}
struct PreparedChild {
path: AgentPath,
child_mailbox: ChildMailbox,
child_runtime: AgentRuntime,
session_id: SessionId,
child_cancel: CancellationToken,
slot: ExecutionSlot,
ticket: SpawnTicket,
spawned_tools: BTreeSet<String>,
resolution: CapabilityResolution,
generation: u64,
recycled: bool,
}
pub(crate) struct SpawnedChild {
path: AgentPath,
spawned_tools: BTreeSet<String>,
resolution: CapabilityResolution,
recycled: bool,
}
impl SpawnedChild {
pub(crate) fn agent_path(&self) -> &AgentPath {
&self.path
}
pub(crate) fn spawned_tools(&self) -> &BTreeSet<String> {
&self.spawned_tools
}
pub(crate) fn resolution(&self) -> &CapabilityResolution {
&self.resolution
}
}
#[derive(Clone, Debug)]
pub struct SpawnEcho {
pub agent_path: String,
pub registered_tools: BTreeSet<String>,
pub degraded_reason: Option<String>,
pub recycled: bool,
}
struct ChildCleanup {
_slot: ExecutionSlot,
mailbox: Arc<MailboxHub>,
registry: Arc<Mutex<AgentRegistry>>,
child_cancels: Arc<Mutex<HashMap<AgentPath, CancellationToken>>>,
path: AgentPath,
write_gate: Arc<crate::multi_agent::write_gate::WorkspaceWriteGate>,
spawned_tools: Arc<Mutex<HashMap<String, Vec<String>>>>,
child_generations: Arc<Mutex<HashMap<AgentPath, u64>>>,
generation: u64,
}
impl Drop for ChildCleanup {
fn drop(&mut self) {
{
let mut gens = self.child_generations.lock().unwrap();
if gens.get(&self.path).copied() != Some(self.generation) {
return; }
gens.remove(&self.path);
}
self.mailbox.post_result(MailboxResult {
agent_path: self.path.clone(),
status: MailboxStatus::Closed,
result: None,
denied_tools: vec![],
});
self.registry.lock().unwrap().close(&self.path);
self.mailbox.unregister(&self.path);
self.child_cancels.lock().unwrap().remove(&self.path);
self.write_gate.release_all(self.path.to_string().as_str());
self.spawned_tools
.lock()
.unwrap()
.remove(&self.path.to_string());
}
}
async fn execute_child_task(
child_runtime: &AgentRuntime,
session_id: &SessionId,
task: &crate::multi_agent::mailbox::MailboxTask,
task_timeout: Option<Duration>,
) -> (MailboxStatus, Option<String>, Vec<String>) {
let input = outcome::build_child_input(task);
let run = child_runtime.run_turn_collect(session_id.clone(), &input);
let result = if let Some(dur) = task_timeout {
match tokio::time::timeout(dur, run).await {
Ok(r) => r,
Err(_elapsed) => {
child_runtime.cancel();
return (
MailboxStatus::Error,
Some(format!("task timed out after {dur:?}")),
vec![],
);
}
}
} else {
run.await
};
match result {
Ok((events, outcome)) => (
MailboxStatus::Ok,
Some(outcome::build_child_result(&outcome, &events)),
outcome::collect_denied_tools(&events),
),
Err(e) => (MailboxStatus::Error, Some(e.to_string()), vec![]),
}
}
#[allow(clippy::too_many_arguments)] async fn run_child_loop(
child_mailbox: ChildMailbox,
child_runtime: AgentRuntime,
session_id: SessionId,
agent_path: AgentPath,
mailbox: Arc<MailboxHub>,
registry: Arc<Mutex<AgentRegistry>>,
event_tx: Option<tokio::sync::mpsc::UnboundedSender<RuntimeEvent>>,
child_cancel: CancellationToken,
task_timeout: Option<Duration>,
write_gate: Arc<crate::multi_agent::write_gate::WorkspaceWriteGate>,
) {
let mut task_rx = child_mailbox.task_rx;
if let Some(tx) = event_tx {
let mut child_events = child_runtime.subscribe_runtime_events();
let bridge_path = agent_path.to_string();
let bridge_cancel = child_cancel.clone();
let bridge_registry = registry.clone();
let bridge_agent_path = agent_path.clone();
tokio::spawn(async move {
loop {
tokio::select! {
_ = bridge_cancel.cancelled() => break,
event = child_events.recv() => {
match event {
Ok(event) => {
if matches!(event, RuntimeEvent::ToolCallStarted { .. }) {
bridge_registry
.lock()
.unwrap()
.record_tool_call(&bridge_agent_path);
}
if matches!(event,
RuntimeEvent::RunFinished { .. }
| RuntimeEvent::RunCancelled { .. }
| RuntimeEvent::AwaitingApproval { .. }) {
continue;
}
let _ = tx.send(event.with_agent_id(bridge_path.as_str()));
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
tracing::warn!(
subagent = %bridge_path,
lagged = n,
"child event bridge lagged"
);
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
}
}
});
}
loop {
tokio::select! {
_ = child_cancel.cancelled() => {
break;
}
task = task_rx.recv() => {
match task {
Some(task) => {
registry.lock().unwrap().note_dequeued(&agent_path);
let (status, result_text, denied_tools) =
execute_child_task(
&child_runtime,
&session_id,
&task,
task_timeout,
)
.await;
write_gate.release_all(agent_path.to_string().as_str());
registry.lock().unwrap().note_posted(&agent_path);
mailbox.post_result(MailboxResult {
agent_path: agent_path.clone(),
status,
result: result_text,
denied_tools,
});
}
None => break, }
}
}
}
}