Skip to main content

agent_works/multi_agent/
runtime.rs

1//! Multi-agent runtime — coordinates sub-agent lifecycle, event bridging, and
2//! cancellation.
3//!
4//! The [`MultiAgentRuntime`] is the central coordinator. It is created once during
5//! builder setup and shared via `Arc` to all 6 multi-agent tools.
6
7use std::collections::HashMap;
8use std::sync::{Arc, Mutex};
9
10use agent_base::{
11    AgentBuilder, AgentResult, AgentRuntime, AllowAllApprovalHandler, ApprovalHandler,
12    DenyAllApprovalHandler, DenyAllToolPolicy, Language, RunOutcome, RuntimeEvent, SessionId,
13    StreamClient, Tool, ToolPolicy,
14};
15use tokio::task::JoinSet;
16use tokio_util::sync::CancellationToken;
17
18use super::config::{ChildPermissionMode, MultiAgentConfig};
19use super::mailbox::{ChildMailbox, MailboxHub, MailboxResult, MailboxStatus, MailboxTask};
20use super::path::AgentPath;
21use super::registry::{AgentRegistry, AgentStatus};
22
23// ---------------------------------------------------------------------------
24// MultiAgentRuntime
25// ---------------------------------------------------------------------------
26
27/// Coordinates sub-agent lifecycle, event bridging, and cancellation.
28///
29/// Created once during builder setup and shared via `Arc` to all 6 multi-agent
30/// tools. Each tool calls methods on the runtime to spawn, communicate with, or
31/// close sub-agents.
32pub struct MultiAgentRuntime {
33    /// Agent lifecycle registry (spawn/close/query).
34    registry: Mutex<AgentRegistry>,
35
36    /// Inter-agent message hub.
37    mailbox: Arc<MailboxHub>,
38
39    /// Shared LLM client (from parent agent).
40    client: Arc<dyn StreamClient>,
41
42    /// Business tools to register on child agents (NOT the 6 multi-agent tools).
43    business_tools: Vec<Arc<dyn Tool>>,
44
45    /// Channel to the bridge task that emits events on parent's event bus.
46    event_tx: Mutex<Option<tokio::sync::mpsc::UnboundedSender<RuntimeEvent>>>,
47
48    /// Root cancellation token (propagates to all children).
49    root_cancel: CancellationToken,
50
51    /// JoinSet tracking all child agent tasks.
52    join_set: Mutex<JoinSet<()>>,
53
54    /// Per-child cancellation tokens.
55    child_cancels: Mutex<HashMap<AgentPath, CancellationToken>>,
56
57    /// Error recovery strategy (inherited from parent).
58    error_recovery: Option<Arc<dyn agent_base::ToolErrorRecovery>>,
59
60    /// Language preference.
61    language: Language,
62
63    /// Permission mode for spawned child agents (resolved per spawn).
64    child_permission_mode: ChildPermissionMode,
65
66    /// Parent's tool policy, inherited by "no permission" children.
67    tool_policy: Option<Arc<dyn ToolPolicy>>,
68
69    /// Parent session manager — for fork_history (child context inheritance).
70    session_manager: Mutex<Option<Arc<agent_base::engine::SessionManager>>>,
71}
72
73impl MultiAgentRuntime {
74    /// Create a new multi-agent runtime.
75    ///
76    /// This is called internally by the builder. Tools receive an `Arc<Self>`.
77    pub fn new(
78        config: MultiAgentConfig,
79        client: Arc<dyn StreamClient>,
80        business_tools: Vec<Arc<dyn Tool>>,
81        root_cancel: CancellationToken,
82        error_recovery: Option<Arc<dyn agent_base::ToolErrorRecovery>>,
83        language: Language,
84        tool_policy: Option<Arc<dyn ToolPolicy>>,
85    ) -> Self {
86        let child_permission_mode = config.child_permission_mode;
87        Self {
88            registry: Mutex::new(AgentRegistry::new(config)),
89            mailbox: Arc::new(MailboxHub::new()),
90            client,
91            business_tools,
92            event_tx: Mutex::new(None),
93            root_cancel,
94            join_set: Mutex::new(JoinSet::new()),
95            child_cancels: Mutex::new(HashMap::new()),
96            error_recovery,
97            language,
98            child_permission_mode,
99            tool_policy,
100            session_manager: Mutex::new(None),
101        }
102    }
103
104    /// Set the event sender for bridging child events to parent.
105    ///
106    /// Called by the builder after creating the bridge channel.
107    pub fn set_event_sender(&self, tx: tokio::sync::mpsc::UnboundedSender<RuntimeEvent>) {
108        *self.event_tx.lock().unwrap() = Some(tx);
109    }
110
111    /// Set the parent session manager for fork_history support.
112    ///
113    /// Called by the builder after creating the runtime.
114    pub fn set_session_manager(&self, session_manager: Arc<agent_base::engine::SessionManager>) {
115        *self.session_manager.lock().unwrap() = Some(session_manager);
116    }
117
118    /// Spawn a child agent at the given path with a specific system prompt.
119    ///
120    /// This is called by the `spawn_agent` tool. It:
121    /// 1. Checks spawn limits
122    /// 2. Registers the agent in the registry
123    /// 3. Creates a mailbox
124    /// 4. Builds a child AgentRuntime
125    /// 5. Spawns a tokio task for the child's event loop
126    /// 6. Returns the AgentPath
127    ///
128    /// `parent_messages` optionally provides context from the parent session
129    /// for fork_history support.
130    ///
131    /// # Errors
132    ///
133    /// Returns a string error message if spawning fails (limits exceeded, etc.).
134    pub async fn spawn_child(
135        &self,
136        name: &str,
137        system_prompt: String,
138        depth: i32,
139        tool_count: usize,
140        full_permission: bool,
141        parent_messages: Vec<agent_base::ChatMessage>,
142    ) -> Result<String, String> {
143        let path = AgentPath::root().join(name);
144
145        // 1. Check limits and register
146        {
147            let mut registry = self.registry.lock().unwrap();
148            registry.can_spawn(depth).map_err(|e| e.to_string())?;
149            registry
150                .register(&path, depth, tool_count)
151                .map_err(|e| e.to_string())?;
152        }
153
154        // 2. Create mailbox
155        let child_mailbox = self
156            .mailbox
157            .register(&path)
158            .ok_or_else(|| "mailbox already exists".to_string())?;
159
160        // 3. Build child AgentRuntime (roll back registry+mailbox on failure)
161        let child_runtime = self
162            .build_child_runtime(system_prompt, self.effective_permission(full_permission))
163            .await
164            .map_err(|e| {
165                self.registry.lock().unwrap().close(&path);
166                self.mailbox.unregister(&path);
167                format!("failed to build child runtime: {}", e)
168            })?;
169
170        // 4. Create session for child and pre-fill with parent context
171        let session_id = child_runtime.create_session().await;
172        self.prefill_child_session(&child_runtime, &session_id, &parent_messages)
173            .await
174            .map_err(|e| {
175                self.registry.lock().unwrap().close(&path);
176                self.mailbox.unregister(&path);
177                format!("failed to prefill child session: {}", e)
178            })?;
179
180        // 5. Create child cancellation token
181        let child_cancel = self.root_cancel.child_token();
182        {
183            let mut cancels = self.child_cancels.lock().unwrap();
184            cancels.insert(path.clone(), child_cancel.clone());
185        }
186
187        // 6. Spawn child agent event loop
188        let agent_path = path.clone();
189        let mailbox_for_task = self.mailbox.clone();
190        let mailbox_for_close = self.mailbox.clone();
191        let event_tx = self.event_tx.lock().unwrap().clone();
192        let registry_agent_path = path.clone();
193
194        self.join_set.lock().unwrap().spawn(async move {
195            run_child_loop(
196                child_mailbox,
197                child_runtime,
198                session_id,
199                agent_path.clone(),
200                mailbox_for_task,
201                event_tx,
202                child_cancel,
203            )
204            .await;
205
206            // Post close notification when loop exits
207            mailbox_for_close.post_result(MailboxResult {
208                agent_path,
209                status: MailboxStatus::Closed,
210                result: None,
211                denied_tools: vec![],
212            });
213        });
214
215        self.registry
216            .lock()
217            .unwrap()
218            .set_status(&registry_agent_path, AgentStatus::Idle);
219
220        Ok(path.to_string())
221    }
222
223    /// Spawn a child agent with fork_history support.
224    ///
225    /// `fork_history`: "none" (default), "all", or a number N for last N turns.
226    /// `parent_session_id`: the parent agent's session ID.
227    #[allow(clippy::too_many_arguments)] // spawn config is naturally positional
228    pub async fn spawn_child_with_history(
229        &self,
230        name: &str,
231        system_prompt: String,
232        depth: i32,
233        tool_count: usize,
234        full_permission: bool,
235        fork_history: Option<String>,
236        parent_session_id: &SessionId,
237    ) -> Result<String, String> {
238        let parent_messages = self
239            .resolve_fork_history(fork_history, parent_session_id)
240            .await;
241        self.spawn_child(
242            name,
243            system_prompt,
244            depth,
245            tool_count,
246            full_permission,
247            parent_messages,
248        )
249        .await
250    }
251
252    /// Resolve fork_history parameter into a list of parent ChatMessages.
253    pub(crate) async fn resolve_fork_history(
254        &self,
255        fork_history: Option<String>,
256        parent_session_id: &SessionId,
257    ) -> Vec<agent_base::ChatMessage> {
258        use agent_base::ChatMessage;
259        let mode = match fork_history.as_deref() {
260            None | Some("none") => return vec![],
261            Some(s) => s,
262        };
263
264        let sm = match self.session_manager.lock().unwrap().as_ref() {
265            Some(sm) => sm.clone(),
266            None => {
267                tracing::warn!("fork_history requested but no session_manager set");
268                return vec![];
269            }
270        };
271
272        // Get all messages from parent session
273        let all_messages = match sm.session_or_err(parent_session_id).await {
274            Ok(session) => session.chat_messages().to_vec(),
275            Err(e) => {
276                tracing::warn!(session_id = parent_session_id.id, error = %e, "failed to load parent session for fork_history");
277                return vec![];
278            }
279        };
280
281        if all_messages.is_empty() {
282            return vec![];
283        }
284
285        // Filter out system messages (child has its own system prompt)
286        let non_system: Vec<ChatMessage> = all_messages
287            .into_iter()
288            .filter(|m| !matches!(m, ChatMessage::System { .. }))
289            .collect();
290
291        match mode {
292            "all" => non_system,
293            n_str => {
294                // Parse N: number of recent user/assistant message pairs (turns)
295                let n: usize = match n_str.parse() {
296                    Ok(n) if n > 0 => n,
297                    _ => {
298                        tracing::warn!(
299                            fork_history = n_str,
300                            "invalid fork_history value, treating as 'none'"
301                        );
302                        return vec![];
303                    }
304                };
305
306                // Count turns from the end (each turn = user message followed by response)
307                let mut turns = 0usize;
308                let mut cutoff = non_system.len();
309                for (i, msg) in non_system.iter().enumerate().rev() {
310                    if matches!(msg, ChatMessage::User { .. }) {
311                        turns += 1;
312                        if turns >= n {
313                            cutoff = i;
314                            break;
315                        }
316                    }
317                }
318                non_system[cutoff..].to_vec()
319            }
320        }
321    }
322
323    /// Send a message to a child agent (no execution trigger).
324    ///
325    /// Called by `send_message` tool.
326    pub fn send_message(&self, agent_path: &str, message: String) -> Result<bool, String> {
327        let path = self.parse_path(agent_path)?;
328        Ok(self.mailbox.send_message(&path, message))
329    }
330
331    /// Send a task to a child agent (triggers execution).
332    ///
333    /// Called by `followup_task` tool. Updates status to Running.
334    pub fn send_task(
335        &self,
336        agent_path: &str,
337        task: String,
338        interrupt: bool,
339    ) -> Result<bool, String> {
340        let path = self.parse_path(agent_path)?;
341        if !self.mailbox.contains(&path) {
342            return Err("agent not found".to_string());
343        }
344        let sent = self.mailbox.send_task(&path, task, interrupt);
345        if sent {
346            self.registry
347                .lock()
348                .unwrap()
349                .set_status(&path, AgentStatus::Running);
350        }
351        Ok(sent)
352    }
353
354    /// Wait for a result from any or a specific child agent.
355    ///
356    /// Called by `wait_agent` tool. Blocks until a result arrives or timeout.
357    pub async fn wait_for_result(&self, agent_path: Option<&str>, timeout_ms: u64) -> WaitResult {
358        let filter_path = match agent_path {
359            Some(s) => match AgentPath::parse(s) {
360                Some(p) => Some(p),
361                None => {
362                    return WaitResult {
363                        status: "error".to_string(),
364                        result: Some(format!("invalid agent path: {}", s)),
365                        agent_path: None,
366                        has_more: false,
367                        denied_tools: vec![],
368                    };
369                }
370            },
371            None => None,
372        };
373
374        let mut seq = self.mailbox.subscribe_seq();
375        let deadline = tokio::time::Instant::now() + tokio::time::Duration::from_millis(timeout_ms);
376
377        loop {
378            // Check for existing results
379            let result = match &filter_path {
380                Some(path) => self.mailbox.try_recv_result(path),
381                None => self.mailbox.try_recv_any(),
382            };
383
384            if let Some(r) = result {
385                let has_more = self.mailbox.total_pending_results() > 0;
386                let (status_str, result_text) = match r.status {
387                    MailboxStatus::Ok => ("ok".to_string(), r.result),
388                    MailboxStatus::Error => ("error".to_string(), r.result),
389                    MailboxStatus::Closed => ("closed".to_string(), r.result),
390                };
391                return WaitResult {
392                    status: status_str,
393                    result: result_text,
394                    agent_path: Some(r.agent_path.to_string()),
395                    has_more,
396                    denied_tools: r.denied_tools,
397                };
398            }
399
400            // Wait for sequence number change or timeout
401            let now = tokio::time::Instant::now();
402            if now >= deadline {
403                return WaitResult {
404                    status: "timeout".to_string(),
405                    result: None,
406                    agent_path: None,
407                    has_more: false,
408                    denied_tools: vec![],
409                };
410            }
411
412            let remaining = deadline - now;
413            tokio::select! {
414                _ = seq.changed() => {
415                    // Sequence changed — loop back to check results
416                    continue;
417                }
418                _ = tokio::time::sleep(remaining) => {
419                    return WaitResult {
420                        status: "timeout".to_string(),
421                        result: None,
422                        agent_path: None,
423                        has_more: false,
424                        denied_tools: vec![],
425                    };
426                }
427            }
428        }
429    }
430
431    /// Close a child agent.
432    ///
433    /// Called by `close_agent` tool. Cancels the child's task, removes from
434    /// registry, and posts a Closed result.
435    pub fn close_agent(&self, agent_path: &str) -> Result<CloseResult, String> {
436        let path = self.parse_path(agent_path)?;
437
438        // Get previous status
439        let previous_status = {
440            let registry = self.registry.lock().unwrap();
441            registry
442                .get(&path)
443                .map(|e| format!("{:?}", e.status).to_lowercase())
444                .unwrap_or_else(|| "unknown".to_string())
445        };
446
447        // Cancel child token
448        {
449            let mut cancels = self.child_cancels.lock().unwrap();
450            if let Some(token) = cancels.remove(&path) {
451                token.cancel();
452            }
453        }
454
455        // Close in registry
456        let existed = { self.registry.lock().unwrap().close(&path).is_some() };
457
458        // Unregister mailbox
459        self.mailbox.unregister(&path);
460
461        Ok(CloseResult {
462            closed: existed,
463            previous_status,
464            message: if existed {
465                "agent closed".to_string()
466            } else {
467                "agent not found".to_string()
468            },
469        })
470    }
471
472    /// List all active sub-agents.
473    ///
474    /// Called by `list_agents` tool.
475    pub fn list_agents(&self) -> Vec<AgentInfo> {
476        let registry = self.registry.lock().unwrap();
477        registry
478            .list()
479            .into_iter()
480            .map(|e| AgentInfo {
481                agent_path: e.path.to_string(),
482                status: format!("{:?}", e.status).to_lowercase(),
483                tool_count: e.tool_count,
484            })
485            .collect()
486    }
487
488    /// Get the mailbox hub (for tools that need it directly).
489    pub fn mailbox(&self) -> &Arc<MailboxHub> {
490        &self.mailbox
491    }
492
493    /// Get reference to the registry.
494    pub fn registry(&self) -> &Mutex<AgentRegistry> {
495        &self.registry
496    }
497
498    /// Cancel all child agents.
499    pub fn cancel_all(&self) {
500        let mut cancels = self.child_cancels.lock().unwrap();
501        for (_, token) in cancels.drain() {
502            token.cancel();
503        }
504    }
505}
506
507impl Drop for MultiAgentRuntime {
508    fn drop(&mut self) {
509        self.cancel_all();
510        // Drain any already-completed join handles to detect panics
511        let mut js = self.join_set.lock().unwrap();
512        while let Some(result) = js.try_join_next() {
513            if let Err(e) = result
514                && e.is_panic()
515            {
516                tracing::error!(
517                    error = %e,
518                    "child agent task panicked"
519                );
520            }
521        }
522    }
523}
524
525impl MultiAgentRuntime {
526    fn parse_path(&self, s: &str) -> Result<AgentPath, String> {
527        AgentPath::parse(s).ok_or_else(|| format!("invalid agent path: '{}'", s))
528    }
529
530    async fn build_child_runtime(
531        &self,
532        system_prompt: String,
533        full_permission: bool,
534    ) -> AgentResult<AgentRuntime> {
535        let (prompt, policy, approval): (
536            String,
537            Option<Arc<dyn ToolPolicy>>,
538            Arc<dyn ApprovalHandler>,
539        ) = if full_permission {
540            // Full: no tool policy → every tool auto-approves (= current behaviour).
541            (system_prompt, None, Arc::new(AllowAllApprovalHandler))
542        } else {
543            // None: policy = parent's (if any) or DenyAllToolPolicy fallback.
544            let note = "If a tool call is rejected for lack of permission, explain in your final answer that you lacked permission for that action.";
545            let policy: Arc<dyn ToolPolicy> = match &self.tool_policy {
546                Some(p) => p.clone(),
547                None => Arc::new(DenyAllToolPolicy),
548            };
549            (
550                format!("{}\n\n{}", system_prompt, note),
551                Some(policy),
552                Arc::new(DenyAllApprovalHandler),
553            )
554        };
555
556        let mut builder = AgentBuilder::new(self.client.clone())
557            .system_prompt(prompt)
558            .approval_handler(approval)
559            .language(self.language.clone());
560
561        if let Some(p) = policy {
562            builder = builder.tool_policy(p);
563        }
564
565        // Register business tools (NOT multi-agent tools)
566        for tool in &self.business_tools {
567            builder = builder.register_tool_arc(tool.clone());
568        }
569
570        if let Some(ref recovery) = self.error_recovery {
571            builder = builder.error_recovery(recovery.clone());
572        }
573
574        builder.build()
575    }
576
577    /// Resolve the effective full-permission flag for a spawn, applying the
578    /// configured [`ChildPermissionMode`]. `Full`/`None` override the LLM-supplied
579    /// flag; `PerSpawn` lets the LLM decide.
580    fn effective_permission(&self, full_permission: bool) -> bool {
581        match self.child_permission_mode {
582            ChildPermissionMode::Full => true,
583            ChildPermissionMode::None => false,
584            ChildPermissionMode::PerSpawn => full_permission,
585        }
586    }
587
588    /// Pre-fill a child session with parent conversation context (fork_history).
589    ///
590    /// Skips system messages and tool-call-only assistant messages. Assistant text
591    /// responses and tool results are stored as system messages with labels so the
592    /// child sees the context without confusing role semantics.
593    pub(crate) async fn prefill_child_session(
594        &self,
595        child_runtime: &AgentRuntime,
596        session_id: &SessionId,
597        parent_messages: &[agent_base::ChatMessage],
598    ) -> AgentResult<()> {
599        use agent_base::ChatMessage;
600
601        for msg in parent_messages {
602            match msg {
603                ChatMessage::User { content, .. } => {
604                    child_runtime.add_user_message(session_id, content).await?;
605                }
606                ChatMessage::Assistant {
607                    content: Some(text),
608                    ..
609                } => {
610                    child_runtime
611                        .add_system_message(
612                            session_id,
613                            format!("[Parent assistant response]: {}", text),
614                        )
615                        .await?;
616                }
617                ChatMessage::Assistant { tool_calls, .. } if tool_calls.is_some() => {
618                    // Skip tool-call-only messages — parent's tool decisions
619                    // don't make sense in the child's context.
620                }
621                ChatMessage::Tool {
622                    tool_call_id,
623                    content,
624                } => {
625                    child_runtime
626                        .add_system_message(
627                            session_id,
628                            format!("[Parent tool result ({}): {}]", tool_call_id, content),
629                        )
630                        .await?;
631                }
632                _ => {} // Skip system messages and empty assistant
633            }
634        }
635
636        Ok(())
637    }
638}
639
640// ---------------------------------------------------------------------------
641// Result types
642// ---------------------------------------------------------------------------
643
644/// Result from `wait_for_result()`.
645#[derive(Clone, Debug)]
646pub struct WaitResult {
647    pub status: String,
648    pub result: Option<String>,
649    pub agent_path: Option<String>,
650    pub has_more: bool,
651    /// Tools the child attempted but was denied permission to call.
652    pub denied_tools: Vec<String>,
653}
654
655/// Result from `close_agent()`.
656#[derive(Clone, Debug)]
657pub struct CloseResult {
658    pub closed: bool,
659    pub previous_status: String,
660    pub message: String,
661}
662
663/// Agent info for `list_agents()`.
664#[derive(Clone, Debug, serde::Serialize)]
665pub struct AgentInfo {
666    pub agent_path: String,
667    pub status: String,
668    pub tool_count: usize,
669}
670
671// ---------------------------------------------------------------------------
672// Child agent event loop
673// ---------------------------------------------------------------------------
674
675/// Run the child agent's main event loop.
676///
677/// This function runs inside a tokio task spawned by [`MultiAgentRuntime::spawn_child`].
678/// It:
679/// 1. Subscribes to child agent events and bridges them to parent
680/// 2. Listens for tasks from the mailbox
681/// 3. Executes each task via `run_turn`
682/// 4. Posts results back via the mailbox
683async fn run_child_loop(
684    child_mailbox: ChildMailbox,
685    child_runtime: AgentRuntime,
686    session_id: SessionId,
687    agent_path: AgentPath,
688    mailbox: Arc<MailboxHub>,
689    event_tx: Option<tokio::sync::mpsc::UnboundedSender<RuntimeEvent>>,
690    child_cancel: CancellationToken,
691) {
692    let mut task_rx = child_mailbox.task_rx;
693
694    // Spawn event bridging: forward child events to parent, tagging agent_id
695    if let Some(tx) = event_tx {
696        let mut child_events = child_runtime.subscribe_runtime_events();
697        let bridge_path = agent_path.to_string();
698        let bridge_cancel = child_cancel.clone();
699
700        tokio::spawn(async move {
701            loop {
702                tokio::select! {
703                    _ = bridge_cancel.cancelled() => break,
704                    event = child_events.recv() => {
705                        match event {
706                            Ok(event) => {
707                                if matches!(event,
708                                    RuntimeEvent::RunFinished { .. }
709                                    | RuntimeEvent::RunCancelled { .. }
710                                    | RuntimeEvent::AwaitingApproval { .. }) {
711                                    continue;
712                                }
713                                let _ = tx.send(event.with_agent_id(bridge_path.as_str()));
714                            }
715                            Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
716                                tracing::warn!(
717                                    subagent = %bridge_path,
718                                    lagged = n,
719                                    "child event bridge lagged"
720                                );
721                            }
722                            Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
723                        }
724                    }
725                }
726            }
727        });
728    }
729
730    // Main task loop
731    loop {
732        tokio::select! {
733            _ = child_cancel.cancelled() => {
734                break;
735            }
736            task = task_rx.recv() => {
737                match task {
738                    Some(task) => {
739                        let input = build_child_input(&task);
740                        let result = child_runtime.run_turn_collect(
741                            session_id.clone(),
742                            &input,
743                        ).await;
744
745                        match result {
746                            Ok((events, outcome)) => {
747                                let result_text = build_child_result(&outcome, &events);
748                                let denied_tools = collect_denied_tools(&events);
749                                mailbox.post_result(MailboxResult {
750                                    agent_path: agent_path.clone(),
751                                    status: MailboxStatus::Ok,
752                                    result: Some(result_text),
753                                    denied_tools,
754                                });
755                            }
756                            Err(e) => {
757                                mailbox.post_result(MailboxResult {
758                                    agent_path: agent_path.clone(),
759                                    status: MailboxStatus::Error,
760                                    result: Some(e.to_string()),
761                                    denied_tools: vec![],
762                                });
763                            }
764                        }
765                    }
766                    None => break, // task channel closed
767                }
768            }
769        }
770    }
771}
772
773/// Build the input text for a child agent from a mailbox task.
774fn build_child_input(task: &MailboxTask) -> String {
775    if task.pending_messages.is_empty() {
776        task.task.clone()
777    } else {
778        let mut parts: Vec<String> = Vec::new();
779        for msg in &task.pending_messages {
780            parts.push(format!("[Message]: {}", msg));
781        }
782        parts.push(format!("[Task]: {}", task.task));
783        parts.join("\n\n")
784    }
785}
786
787/// Extract a human-readable summary from a run outcome.
788fn summarize_outcome(outcome: &RunOutcome) -> String {
789    match outcome {
790        RunOutcome::Completed => "task completed".to_string(),
791        RunOutcome::Failed { error } => format!("task failed: {}", error),
792        RunOutcome::MaxTurnsExceeded { turns } => {
793            format!("max turns exceeded ({} turns)", turns)
794        }
795        RunOutcome::Cancelled => "cancelled".to_string(),
796    }
797}
798
799/// Extract the child agent's own assistant text from its collected events.
800///
801/// Excludes sub-sub-agent text (events tagged with an `agent_id`), so the parent
802/// only sees this child's direct answer.
803fn extract_assistant_text(events: &[RuntimeEvent]) -> String {
804    let mut text = String::new();
805    for event in events {
806        if let RuntimeEvent::TextDelta {
807            text: delta,
808            agent_id,
809            ..
810        } = event
811            && agent_id.is_none()
812        {
813            text.push_str(delta);
814        }
815    }
816    text
817}
818
819/// Collect the names of tools the child attempted but was denied permission to
820/// call, from a run's collected events.
821///
822/// Filters to the child's own events (`agent_id == None`) so grandchild denials
823/// — if children ever gain multi-agent tools — are not mis-attributed to the
824/// child. Only the child's direct denials matter to the parent.
825fn collect_denied_tools(events: &[RuntimeEvent]) -> Vec<String> {
826    events
827        .iter()
828        .filter_map(|e| match e {
829            RuntimeEvent::ToolCallFinished {
830                tool_name,
831                denied: true,
832                agent_id: None,
833                ..
834            } => Some(tool_name.clone()),
835            _ => None,
836        })
837        .collect()
838}
839
840/// Build the result string posted back to the parent via the mailbox.
841///
842/// For a completed task this returns the child's actual final answer (so the
843/// parent learns what the child concluded — e.g. "no permission" reports) rather
844/// than a coarse "task completed". Other outcomes keep the coarse summary.
845fn build_child_result(outcome: &RunOutcome, events: &[RuntimeEvent]) -> String {
846    match outcome {
847        RunOutcome::Completed => {
848            let text = extract_assistant_text(events);
849            if text.trim().is_empty() {
850                summarize_outcome(outcome)
851            } else {
852                text
853            }
854        }
855        _ => summarize_outcome(outcome),
856    }
857}
858
859// ---------------------------------------------------------------------------
860// Tests
861// ---------------------------------------------------------------------------
862
863#[cfg(test)]
864mod tests {
865    use super::*;
866    use agent_base::RunOutcome;
867
868    // ── summarize_outcome ──
869
870    #[test]
871    fn test_summarize_completed() {
872        let s = summarize_outcome(&RunOutcome::Completed);
873        assert_eq!(s, "task completed");
874    }
875
876    #[test]
877    fn test_summarize_failed() {
878        let outcome = RunOutcome::Failed {
879            error: "connection refused".to_string(),
880        };
881        let s = summarize_outcome(&outcome);
882        assert_eq!(s, "task failed: connection refused");
883    }
884
885    #[test]
886    fn test_summarize_max_turns() {
887        let outcome = RunOutcome::MaxTurnsExceeded { turns: 42 };
888        let s = summarize_outcome(&outcome);
889        assert!(s.contains("max turns exceeded"));
890        assert!(s.contains("42"));
891    }
892
893    #[test]
894    fn test_summarize_cancelled() {
895        let s = summarize_outcome(&RunOutcome::Cancelled);
896        assert_eq!(s, "cancelled");
897    }
898
899    // ── build_child_result / extract_assistant_text ──
900
901    fn text_delta(text: &str, agent_id: Option<&str>) -> agent_base::RuntimeEvent {
902        agent_base::RuntimeEvent::TextDelta {
903            session_id: agent_base::SessionId::new(1),
904            text: text.to_string(),
905            agent_id: agent_id.map(|s| s.to_string()),
906            trace_id: None,
907        }
908    }
909
910    #[test]
911    fn test_build_child_result_completed_returns_final_text() {
912        let events = vec![text_delta("I couldn't ", None), text_delta("delete.", None)];
913        assert_eq!(
914            build_child_result(&RunOutcome::Completed, &events),
915            "I couldn't delete."
916        );
917    }
918
919    #[test]
920    fn test_build_child_result_completed_falls_back_when_no_text() {
921        assert_eq!(
922            build_child_result(&RunOutcome::Completed, &[]),
923            "task completed"
924        );
925    }
926
927    #[test]
928    fn test_extract_assistant_text_ignores_subagent_text() {
929        let events = vec![
930            text_delta("root answer", None),
931            text_delta("grandchild", Some("root/child/grandchild")),
932        ];
933        assert_eq!(extract_assistant_text(&events), "root answer");
934    }
935
936    #[test]
937    fn test_build_child_result_failed_keeps_error() {
938        let outcome = RunOutcome::Failed {
939            error: "boom".to_string(),
940        };
941        assert_eq!(build_child_result(&outcome, &[]), "task failed: boom");
942    }
943
944    // ── collect_denied_tools ──
945
946    fn tool_finished(tool_name: &str, denied: bool) -> agent_base::RuntimeEvent {
947        agent_base::RuntimeEvent::ToolCallFinished {
948            session_id: agent_base::SessionId::new(1),
949            tool_name: tool_name.to_string(),
950            summary: "summary".to_string(),
951            agent_id: None,
952            trace_id: None,
953            denied,
954        }
955    }
956
957    #[test]
958    fn test_collect_denied_tools_filters_denied_only() {
959        let events = vec![
960            tool_finished("read_file", false),
961            tool_finished("delete_file", true),
962            tool_finished("shell", true),
963        ];
964        assert_eq!(
965            collect_denied_tools(&events),
966            vec!["delete_file".to_string(), "shell".to_string()]
967        );
968    }
969
970    #[test]
971    fn test_collect_denied_tools_empty_when_no_denials() {
972        let events = vec![
973            tool_finished("read_file", false),
974            text_delta("all good", None),
975        ];
976        assert!(collect_denied_tools(&events).is_empty());
977    }
978
979    #[test]
980    fn test_collect_denied_tools_excludes_grandchild_denials() {
981        // A grandchild's denial carries an agent_id and must not be attributed to
982        // the child.
983        let events = vec![
984            agent_base::RuntimeEvent::ToolCallFinished {
985                session_id: agent_base::SessionId::new(1),
986                tool_name: "grandchild_tool".to_string(),
987                summary: "summary".to_string(),
988                agent_id: Some("root/child/grandchild".to_string()),
989                trace_id: None,
990                denied: true,
991            },
992            tool_finished("child_tool", true),
993        ];
994        assert_eq!(
995            collect_denied_tools(&events),
996            vec!["child_tool".to_string()]
997        );
998    }
999
1000    // ── build_child_input ──
1001
1002    #[test]
1003    fn test_build_child_input_task_only() {
1004        let task = MailboxTask {
1005            task: "do work".into(),
1006            interrupt: true,
1007            pending_messages: vec![],
1008        };
1009        let out = build_child_input(&task);
1010        assert_eq!(out, "do work");
1011    }
1012
1013    #[test]
1014    fn test_build_child_input_with_pending_messages() {
1015        let task = MailboxTask {
1016            task: "do work".into(),
1017            interrupt: false,
1018            pending_messages: vec!["context 1".into(), "context 2".into()],
1019        };
1020        let out = build_child_input(&task);
1021        assert!(out.contains("[Message]: context 1"));
1022        assert!(out.contains("[Message]: context 2"));
1023        assert!(out.contains("[Task]: do work"));
1024        // Messages come before task
1025        let msg_pos = out.find("[Message]:").unwrap();
1026        let task_pos = out.find("[Task]:").unwrap();
1027        assert!(msg_pos < task_pos, "messages should precede task");
1028    }
1029
1030    #[test]
1031    fn test_build_child_input_single_message() {
1032        let task = MailboxTask {
1033            task: "final task".into(),
1034            interrupt: true,
1035            pending_messages: vec!["hint".into()],
1036        };
1037        let out = build_child_input(&task);
1038        assert_eq!(out, "[Message]: hint\n\n[Task]: final task");
1039    }
1040
1041    // ── fork_history: resolve_fork_history ──
1042
1043    /// Mock LLM client for fork_history tests (minimal — never called).
1044    #[derive(Clone)]
1045    struct NoopLlmClient;
1046
1047    #[async_trait::async_trait]
1048    impl agent_base::LlmClient for NoopLlmClient {
1049        async fn chat(
1050            &self,
1051            _messages: &[agent_base::ChatMessage],
1052            _tools: &[serde_json::Value],
1053            _reasoning: Option<&agent_base::ReasoningConfig>,
1054            _response_format: Option<&agent_base::ResponseFormat>,
1055        ) -> agent_base::AgentResult<serde_json::Value> {
1056            unimplemented!()
1057        }
1058
1059        async fn chat_stream(
1060            &self,
1061            _messages: &[agent_base::ChatMessage],
1062            _tools: &[serde_json::Value],
1063            _reasoning: Option<&agent_base::ReasoningConfig>,
1064            _response_format: Option<&agent_base::ResponseFormat>,
1065        ) -> agent_base::AgentResult<
1066            std::pin::Pin<
1067                Box<
1068                    dyn futures_core::Stream<
1069                            Item = agent_base::AgentResult<agent_base::StreamChunk>,
1070                        > + Send,
1071                >,
1072            >,
1073        > {
1074            unimplemented!()
1075        }
1076
1077        fn capabilities(&self) -> agent_base::LlmCapabilities {
1078            agent_base::LlmCapabilities {
1079                supports_streaming: true,
1080                supports_tools: false,
1081                supports_vision: false,
1082                supports_thinking: false,
1083                max_context_tokens: None,
1084                max_output_tokens: None,
1085            }
1086        }
1087    }
1088
1089    /// Build a MultiAgentRuntime with a parent runtime that has a populated session.
1090    async fn setup_fork_history_test(
1091        parent_messages: Vec<agent_base::ChatMessage>,
1092    ) -> (Arc<MultiAgentRuntime>, agent_base::SessionId) {
1093        use tokio_util::sync::CancellationToken;
1094
1095        let llm = agent_base::llm::adapt(Arc::new(NoopLlmClient));
1096        let parent_runtime = agent_base::AgentBuilder::new(llm)
1097            .build()
1098            .expect("build parent runtime");
1099        let parent_sid = parent_runtime.create_session().await;
1100
1101        // Push messages directly into the session's chat_messages vector so
1102        // we can use proper Assistant/Tool variants (not just System).
1103        parent_runtime
1104            .with_session_mut(&parent_sid, |session| {
1105                session.chat_messages_mut().extend(parent_messages.clone());
1106            })
1107            .await
1108            .unwrap();
1109
1110        let session_manager = Arc::new(parent_runtime.session_manager().clone());
1111
1112        let ma_runtime = Arc::new(MultiAgentRuntime::new(
1113            MultiAgentConfig::enabled(),
1114            agent_base::llm::adapt(Arc::new(NoopLlmClient)),
1115            vec![],
1116            CancellationToken::new(),
1117            None,
1118            agent_base::Language::En,
1119            None,
1120        ));
1121        ma_runtime.set_session_manager(session_manager);
1122
1123        (ma_runtime, parent_sid)
1124    }
1125
1126    #[tokio::test]
1127    async fn resolve_fork_history_none_returns_empty() {
1128        let messages = vec![agent_base::ChatMessage::User {
1129            content: "hello".into(),
1130            images: vec![],
1131            ephemeral: false,
1132        }];
1133        let (ma, parent_sid) = setup_fork_history_test(messages).await;
1134
1135        // None
1136        let result = ma.resolve_fork_history(None, &parent_sid).await;
1137        assert!(result.is_empty());
1138
1139        // Some("none")
1140        let result = ma
1141            .resolve_fork_history(Some("none".to_string()), &parent_sid)
1142            .await;
1143        assert!(result.is_empty());
1144    }
1145
1146    #[tokio::test]
1147    async fn resolve_fork_history_all_returns_all_non_system() {
1148        let messages = vec![
1149            agent_base::ChatMessage::User {
1150                content: "question 1".into(),
1151                images: vec![],
1152                ephemeral: false,
1153            },
1154            agent_base::ChatMessage::Assistant {
1155                content: Some("answer 1".into()),
1156                reasoning_content: None,
1157                tool_calls: None,
1158            },
1159            agent_base::ChatMessage::User {
1160                content: "question 2".into(),
1161                images: vec![],
1162                ephemeral: false,
1163            },
1164            agent_base::ChatMessage::Assistant {
1165                content: Some("answer 2".into()),
1166                reasoning_content: None,
1167                tool_calls: None,
1168            },
1169        ];
1170        let (ma, parent_sid) = setup_fork_history_test(messages).await;
1171
1172        let result = ma
1173            .resolve_fork_history(Some("all".to_string()), &parent_sid)
1174            .await;
1175
1176        // Should have 4 messages (2 user + 2 assistant) — system messages are filtered out
1177        assert_eq!(result.len(), 4);
1178        assert!(matches!(result[0], agent_base::ChatMessage::User { .. }));
1179        assert!(matches!(
1180            result[1],
1181            agent_base::ChatMessage::Assistant { .. }
1182        ));
1183        assert!(matches!(result[2], agent_base::ChatMessage::User { .. }));
1184        assert!(matches!(
1185            result[3],
1186            agent_base::ChatMessage::Assistant { .. }
1187        ));
1188    }
1189
1190    #[tokio::test]
1191    async fn resolve_fork_history_n_turns() {
1192        // 3 turns: 3 user messages, 3 assistant responses
1193        let messages = vec![
1194            agent_base::ChatMessage::User {
1195                content: "q1".into(),
1196                images: vec![],
1197                ephemeral: false,
1198            },
1199            agent_base::ChatMessage::Assistant {
1200                content: Some("a1".into()),
1201                reasoning_content: None,
1202                tool_calls: None,
1203            },
1204            agent_base::ChatMessage::User {
1205                content: "q2".into(),
1206                images: vec![],
1207                ephemeral: false,
1208            },
1209            agent_base::ChatMessage::Assistant {
1210                content: Some("a2".into()),
1211                reasoning_content: None,
1212                tool_calls: None,
1213            },
1214            agent_base::ChatMessage::User {
1215                content: "q3".into(),
1216                images: vec![],
1217                ephemeral: false,
1218            },
1219            agent_base::ChatMessage::Assistant {
1220                content: Some("a3".into()),
1221                reasoning_content: None,
1222                tool_calls: None,
1223            },
1224        ];
1225        let (ma, parent_sid) = setup_fork_history_test(messages).await;
1226
1227        // Last 1 turn
1228        let result = ma
1229            .resolve_fork_history(Some("1".to_string()), &parent_sid)
1230            .await;
1231        assert_eq!(result.len(), 2, "1 turn = user q3 + assistant a3");
1232        assert!(matches!(result[0], agent_base::ChatMessage::User { .. }));
1233        assert_eq!(extract_user_content(&result[0]), "q3");
1234
1235        // Last 2 turns
1236        let result = ma
1237            .resolve_fork_history(Some("2".to_string()), &parent_sid)
1238            .await;
1239        assert_eq!(result.len(), 4, "2 turns = q2,a2,q3,a3");
1240    }
1241
1242    #[tokio::test]
1243    async fn resolve_fork_history_invalid_number_treats_as_none() {
1244        let messages = vec![agent_base::ChatMessage::User {
1245            content: "hello".into(),
1246            images: vec![],
1247            ephemeral: false,
1248        }];
1249        let (ma, parent_sid) = setup_fork_history_test(messages).await;
1250
1251        // Invalid number → empty
1252        let result = ma
1253            .resolve_fork_history(Some("not-a-number".to_string()), &parent_sid)
1254            .await;
1255        assert!(result.is_empty());
1256
1257        // Zero → empty
1258        let result = ma
1259            .resolve_fork_history(Some("0".to_string()), &parent_sid)
1260            .await;
1261        assert!(result.is_empty());
1262    }
1263
1264    #[tokio::test]
1265    async fn resolve_fork_history_no_session_manager_returns_empty() {
1266        use tokio_util::sync::CancellationToken;
1267
1268        let ma_runtime = MultiAgentRuntime::new(
1269            MultiAgentConfig::enabled(),
1270            agent_base::llm::adapt(Arc::new(NoopLlmClient)),
1271            vec![],
1272            CancellationToken::new(),
1273            None,
1274            agent_base::Language::En,
1275            None,
1276        );
1277        // session_manager is NOT set
1278
1279        let sid = agent_base::SessionId::new(9999);
1280        let result = ma_runtime
1281            .resolve_fork_history(Some("all".to_string()), &sid)
1282            .await;
1283        assert!(result.is_empty());
1284    }
1285
1286    #[tokio::test]
1287    async fn resolve_fork_history_empty_session_returns_empty() {
1288        let (ma, parent_sid) = setup_fork_history_test(vec![]).await;
1289
1290        let result = ma
1291            .resolve_fork_history(Some("all".to_string()), &parent_sid)
1292            .await;
1293        assert!(result.is_empty());
1294    }
1295
1296    // ── fork_history: prefill_child_session ──
1297
1298    #[tokio::test]
1299    async fn prefill_child_session_user_and_assistant() {
1300        let llm = agent_base::llm::adapt(Arc::new(NoopLlmClient));
1301        let child_runtime = agent_base::AgentBuilder::new(llm)
1302            .build()
1303            .expect("build child runtime");
1304        let child_sid = child_runtime.create_session().await;
1305
1306        let parent_messages = vec![
1307            agent_base::ChatMessage::User {
1308                content: "user question".into(),
1309                images: vec![],
1310                ephemeral: false,
1311            },
1312            agent_base::ChatMessage::Assistant {
1313                content: Some("assistant reply".into()),
1314                reasoning_content: None,
1315                tool_calls: None,
1316            },
1317            agent_base::ChatMessage::Tool {
1318                tool_call_id: "call_123".into(),
1319                content: "tool output".into(),
1320            },
1321        ];
1322
1323        // Create a minimal MultiAgentRuntime just to call prefill_child_session
1324        use tokio_util::sync::CancellationToken;
1325        let ma_runtime = MultiAgentRuntime::new(
1326            MultiAgentConfig::enabled(),
1327            agent_base::llm::adapt(Arc::new(NoopLlmClient)),
1328            vec![],
1329            CancellationToken::new(),
1330            None,
1331            agent_base::Language::En,
1332            None,
1333        );
1334
1335        ma_runtime
1336            .prefill_child_session(&child_runtime, &child_sid, &parent_messages)
1337            .await
1338            .expect("prefill should succeed");
1339
1340        // Verify the child session contains the pre-filled messages
1341        let session = child_runtime
1342            .session(&child_sid)
1343            .await
1344            .expect("session exists");
1345        let msgs = session.chat_messages().to_vec();
1346
1347        // Should have: user msg + system msg (assistant) + system msg (tool)
1348        assert_eq!(msgs.len(), 3);
1349        assert!(matches!(msgs[0], agent_base::ChatMessage::User { .. }));
1350        assert!(matches!(msgs[1], agent_base::ChatMessage::System { .. }));
1351        assert!(matches!(msgs[2], agent_base::ChatMessage::System { .. }));
1352    }
1353
1354    #[tokio::test]
1355    async fn prefill_child_session_tool_call_only_skipped() {
1356        let llm = agent_base::llm::adapt(Arc::new(NoopLlmClient));
1357        let child_runtime = agent_base::AgentBuilder::new(llm)
1358            .build()
1359            .expect("build child runtime");
1360        let child_sid = child_runtime.create_session().await;
1361
1362        // Assistant message with only tool_calls (no text content) should be skipped
1363        let parent_messages = vec![
1364            agent_base::ChatMessage::User {
1365                content: "do something".into(),
1366                images: vec![],
1367                ephemeral: false,
1368            },
1369            agent_base::ChatMessage::Assistant {
1370                content: None, // no text — tool call only
1371                reasoning_content: None,
1372                tool_calls: Some(vec![]),
1373            },
1374        ];
1375
1376        use tokio_util::sync::CancellationToken;
1377        let ma_runtime = MultiAgentRuntime::new(
1378            MultiAgentConfig::enabled(),
1379            agent_base::llm::adapt(Arc::new(NoopLlmClient)),
1380            vec![],
1381            CancellationToken::new(),
1382            None,
1383            agent_base::Language::En,
1384            None,
1385        );
1386
1387        ma_runtime
1388            .prefill_child_session(&child_runtime, &child_sid, &parent_messages)
1389            .await
1390            .expect("prefill should succeed");
1391
1392        let session = child_runtime
1393            .session(&child_sid)
1394            .await
1395            .expect("session exists");
1396        let msgs = session.chat_messages().to_vec();
1397
1398        // Only the user message — tool-call-only assistant should be skipped
1399        assert_eq!(msgs.len(), 1);
1400        assert!(matches!(msgs[0], agent_base::ChatMessage::User { .. }));
1401    }
1402
1403    #[tokio::test]
1404    async fn prefill_child_session_empty_vec_noop() {
1405        let llm = agent_base::llm::adapt(Arc::new(NoopLlmClient));
1406        let child_runtime = agent_base::AgentBuilder::new(llm)
1407            .build()
1408            .expect("build child runtime");
1409        let child_sid = child_runtime.create_session().await;
1410
1411        use tokio_util::sync::CancellationToken;
1412        let ma_runtime = MultiAgentRuntime::new(
1413            MultiAgentConfig::enabled(),
1414            agent_base::llm::adapt(Arc::new(NoopLlmClient)),
1415            vec![],
1416            CancellationToken::new(),
1417            None,
1418            agent_base::Language::En,
1419            None,
1420        );
1421
1422        ma_runtime
1423            .prefill_child_session(&child_runtime, &child_sid, &[])
1424            .await
1425            .expect("prefill should succeed");
1426
1427        let session = child_runtime
1428            .session(&child_sid)
1429            .await
1430            .expect("session exists");
1431        let msgs = session.chat_messages().to_vec();
1432
1433        // System prompt is added but we don't assert exact count — just that no user/injected msgs
1434        assert!(msgs.is_empty() || matches!(msgs[0], agent_base::ChatMessage::System { .. }));
1435    }
1436
1437    fn extract_user_content(msg: &agent_base::ChatMessage) -> &str {
1438        match msg {
1439            agent_base::ChatMessage::User { content, .. } => content.as_str(),
1440            _ => "",
1441        }
1442    }
1443
1444    // ── lifecycle: spawn / send / wait / close ──
1445
1446    struct StreamingStub;
1447
1448    #[async_trait::async_trait]
1449    impl agent_base::StreamClient for StreamingStub {
1450        async fn stream(
1451            &self,
1452            _messages: &[agent_base::ChatMessage],
1453            _tools: &[serde_json::Value],
1454            _reasoning: Option<&agent_base::ReasoningConfig>,
1455            _response_format: Option<&agent_base::ResponseFormat>,
1456        ) -> agent_base::AgentResult<
1457            std::pin::Pin<
1458                Box<
1459                    dyn futures_core::Stream<
1460                            Item = agent_base::AgentResult<agent_base::StreamChunk>,
1461                        > + Send,
1462                >,
1463            >,
1464        > {
1465            Ok(Box::pin(futures_util::stream::iter(vec![
1466                Ok(agent_base::StreamChunk::Text("child ok".to_string())),
1467                Ok(agent_base::StreamChunk::Stop {
1468                    finish_reason: Some("stop".to_string()),
1469                }),
1470            ])))
1471        }
1472
1473        fn capabilities(&self) -> agent_base::LlmCapabilities {
1474            agent_base::LlmCapabilities::default()
1475        }
1476    }
1477
1478    fn make_ma_runtime() -> Arc<MultiAgentRuntime> {
1479        let client: Arc<dyn agent_base::StreamClient> = Arc::new(StreamingStub);
1480        Arc::new(MultiAgentRuntime::new(
1481            MultiAgentConfig::enabled(),
1482            client,
1483            vec![],
1484            tokio_util::sync::CancellationToken::new(),
1485            None,
1486            agent_base::Language::En,
1487            None,
1488        ))
1489    }
1490
1491    #[tokio::test(flavor = "multi_thread")]
1492    async fn test_spawn_send_task_wait_close_lifecycle() {
1493        let ma = make_ma_runtime();
1494
1495        let path = ma
1496            .spawn_child(
1497                "worker",
1498                "child system prompt".to_string(),
1499                0,
1500                0,
1501                false,
1502                vec![],
1503            )
1504            .await
1505            .expect("spawn child");
1506        assert_eq!(path, "root/worker");
1507
1508        let agents = ma.list_agents();
1509        assert_eq!(agents.len(), 1);
1510        assert_eq!(agents[0].agent_path, "root/worker");
1511
1512        // send_message posts a pending message (no execution trigger).
1513        assert!(
1514            ma.send_message("root/worker", "heads up".to_string())
1515                .unwrap()
1516        );
1517
1518        // send_task triggers execution; the child completes via the stub stream.
1519        assert!(
1520            ma.send_task("root/worker", "do the thing".to_string(), false)
1521                .unwrap()
1522        );
1523
1524        let result = ma.wait_for_result(Some("root/worker"), 2000).await;
1525        assert_eq!(result.status, "ok");
1526        assert_eq!(result.result.as_deref(), Some("child ok"));
1527
1528        let close = ma.close_agent("root/worker").unwrap();
1529        assert!(close.closed);
1530        assert_eq!(close.message, "agent closed");
1531
1532        // Closing a second time reports not found.
1533        let close2 = ma.close_agent("root/worker").unwrap();
1534        assert!(!close2.closed);
1535        assert_eq!(close2.message, "agent not found");
1536    }
1537
1538    #[tokio::test(flavor = "multi_thread")]
1539    async fn test_spawn_child_with_history_defaults_to_none() {
1540        let ma = make_ma_runtime();
1541        // No session manager set → fork_history resolves to empty; spawn still succeeds.
1542        let path = ma
1543            .spawn_child_with_history(
1544                "w2",
1545                "prompt".to_string(),
1546                0,
1547                0,
1548                false,
1549                None,
1550                &agent_base::SessionId::new(0),
1551            )
1552            .await
1553            .expect("spawn with history");
1554        assert_eq!(path, "root/w2");
1555    }
1556
1557    #[tokio::test(flavor = "multi_thread")]
1558    async fn test_error_paths() {
1559        let ma = make_ma_runtime();
1560
1561        // Valid path, unknown agent → "agent not found".
1562        assert_eq!(
1563            ma.send_task("root/ghost", "x".to_string(), false)
1564                .unwrap_err(),
1565            "agent not found"
1566        );
1567
1568        // Invalid paths (must start with "root", no empty segments).
1569        assert!(ma.send_message("worker", "x".to_string()).is_err());
1570        assert!(ma.send_message("", "x".to_string()).is_err());
1571
1572        // wait_for_result with invalid path.
1573        let r = ma.wait_for_result(Some("worker"), 10).await;
1574        assert_eq!(r.status, "error");
1575
1576        // wait_for_result with no results times out.
1577        let r2 = ma.wait_for_result(None, 50).await;
1578        assert_eq!(r2.status, "timeout");
1579
1580        // cancel_all is a no-op when no children are running.
1581        ma.cancel_all();
1582    }
1583
1584    // ── child permission mode ──
1585
1586    fn make_runtime_full(
1587        mode: ChildPermissionMode,
1588        policy: Option<Arc<dyn ToolPolicy>>,
1589    ) -> Arc<MultiAgentRuntime> {
1590        let config = MultiAgentConfig {
1591            child_permission_mode: mode,
1592            ..MultiAgentConfig::enabled()
1593        };
1594        Arc::new(MultiAgentRuntime::new(
1595            config,
1596            Arc::new(StreamingStub),
1597            vec![],
1598            tokio_util::sync::CancellationToken::new(),
1599            None,
1600            agent_base::Language::En,
1601            policy,
1602        ))
1603    }
1604
1605    #[test]
1606    fn effective_permission_respects_mode() {
1607        let full = make_runtime_full(ChildPermissionMode::Full, None);
1608        assert!(full.effective_permission(false));
1609        assert!(full.effective_permission(true));
1610
1611        let none = make_runtime_full(ChildPermissionMode::None, None);
1612        assert!(!none.effective_permission(false));
1613        assert!(!none.effective_permission(true));
1614
1615        let per_spawn = make_runtime_full(ChildPermissionMode::PerSpawn, None);
1616        assert!(per_spawn.effective_permission(true));
1617        assert!(!per_spawn.effective_permission(false));
1618    }
1619
1620    #[tokio::test]
1621    async fn build_child_runtime_full_carries_no_policy() {
1622        let ma = make_ma_runtime();
1623        let child = ma
1624            .build_child_runtime("prompt".to_string(), true)
1625            .await
1626            .expect("build child");
1627        assert!(child.tool_policy().is_none());
1628    }
1629
1630    #[tokio::test]
1631    async fn build_child_runtime_none_falls_back_to_deny_all() {
1632        // Parent has no tool policy → child falls back to DenyAllToolPolicy.
1633        let ma = make_ma_runtime();
1634        let child = ma
1635            .build_child_runtime("prompt".to_string(), false)
1636            .await
1637            .expect("build child");
1638        assert!(child.tool_policy().is_some());
1639    }
1640
1641    #[tokio::test]
1642    async fn build_child_runtime_none_inherits_parent_policy() {
1643        // Parent has a tool policy → child inherits the same allocation.
1644        let parent_policy: Arc<dyn ToolPolicy> = Arc::new(DenyAllToolPolicy);
1645        let ma = make_runtime_full(ChildPermissionMode::None, Some(parent_policy.clone()));
1646        let child = ma
1647            .build_child_runtime("prompt".to_string(), false)
1648            .await
1649            .expect("build child");
1650        let child_policy = child.tool_policy().expect("child should carry a policy");
1651        assert!(Arc::ptr_eq(&parent_policy, child_policy));
1652    }
1653
1654    // ── end-to-end: denied_tools flows child → parent ──
1655
1656    struct NoopReadFileTool;
1657
1658    #[async_trait::async_trait]
1659    impl Tool for NoopReadFileTool {
1660        fn name(&self) -> &'static str {
1661            "read_file"
1662        }
1663
1664        fn description(&self) -> &'static str {
1665            "Read a file's contents"
1666        }
1667
1668        fn schema(&self) -> serde_json::Value {
1669            serde_json::json!({
1670                "type": "object",
1671                "properties": { "path": { "type": "string" } }
1672            })
1673        }
1674
1675        async fn call(
1676            &self,
1677            _args: &serde_json::Value,
1678            _ctx: &agent_base::ToolContext,
1679        ) -> agent_base::AgentResult<Vec<agent_base::Content>> {
1680            Ok(vec![agent_base::Content::text("contents")])
1681        }
1682    }
1683
1684    /// Scripted client: the first turn requests the `read_file` tool (which the
1685    /// child is denied); any later turn emits a plain text answer.
1686    struct DenialScriptedClient {
1687        turn: std::sync::atomic::AtomicUsize,
1688    }
1689
1690    #[async_trait::async_trait]
1691    impl agent_base::StreamClient for DenialScriptedClient {
1692        async fn stream(
1693            &self,
1694            _messages: &[agent_base::ChatMessage],
1695            _tools: &[serde_json::Value],
1696            _reasoning: Option<&agent_base::ReasoningConfig>,
1697            _response_format: Option<&agent_base::ResponseFormat>,
1698        ) -> agent_base::AgentResult<
1699            std::pin::Pin<
1700                Box<
1701                    dyn futures_core::Stream<
1702                            Item = agent_base::AgentResult<agent_base::StreamChunk>,
1703                        > + Send,
1704                >,
1705            >,
1706        > {
1707            let n = self.turn.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1708            let chunks: Vec<agent_base::AgentResult<agent_base::StreamChunk>> = if n == 0 {
1709                vec![
1710                    Ok(agent_base::StreamChunk::ToolCall(serde_json::json!({
1711                        "delta": {
1712                            "tool_calls": [{
1713                                "id": "call_1",
1714                                "function": {
1715                                    "name": "read_file",
1716                                    "arguments": "{\"path\":\"/etc/passwd\"}"
1717                                }
1718                            }]
1719                        }
1720                    }))),
1721                    Ok(agent_base::StreamChunk::Stop {
1722                        finish_reason: Some("tool_calls".to_string()),
1723                    }),
1724                ]
1725            } else {
1726                vec![
1727                    Ok(agent_base::StreamChunk::Text(
1728                        "I lack permission.".to_string(),
1729                    )),
1730                    Ok(agent_base::StreamChunk::Stop {
1731                        finish_reason: Some("stop".to_string()),
1732                    }),
1733                ]
1734            };
1735            Ok(Box::pin(futures_util::stream::iter(chunks)))
1736        }
1737
1738        fn capabilities(&self) -> agent_base::LlmCapabilities {
1739            agent_base::LlmCapabilities::default()
1740        }
1741    }
1742
1743    #[tokio::test(flavor = "multi_thread")]
1744    async fn test_child_denied_tool_reaches_parent_via_wait() {
1745        // `None` mode + no parent policy → child carries DenyAllToolPolicy, so
1746        // every tool it attempts is denied. Assert the denied tool name flows
1747        // end-to-end: collect_denied_tools → mailbox → wait_for_result.
1748        let config = MultiAgentConfig {
1749            child_permission_mode: ChildPermissionMode::None,
1750            ..MultiAgentConfig::enabled()
1751        };
1752        let ma = Arc::new(MultiAgentRuntime::new(
1753            config,
1754            Arc::new(DenialScriptedClient {
1755                turn: std::sync::atomic::AtomicUsize::new(0),
1756            }),
1757            vec![Arc::new(NoopReadFileTool) as Arc<dyn Tool>],
1758            tokio_util::sync::CancellationToken::new(),
1759            None,
1760            agent_base::Language::En,
1761            None,
1762        ));
1763
1764        let path = ma
1765            .spawn_child(
1766                "worker",
1767                "child system prompt".to_string(),
1768                0,
1769                1,
1770                false,
1771                vec![],
1772            )
1773            .await
1774            .expect("spawn child");
1775        assert_eq!(path, "root/worker");
1776
1777        ma.send_task("root/worker", "read the file".to_string(), false)
1778            .unwrap();
1779
1780        let result = ma.wait_for_result(Some("root/worker"), 3000).await;
1781        assert_eq!(result.status, "ok");
1782        assert_eq!(result.denied_tools, vec!["read_file".to_string()]);
1783    }
1784}