Skip to main content

vtcode_core/subagents/
controller_spawn_run.rs

1#![allow(
2    unused_imports,
3    reason = "Intentional compatibility, platform, or test-only suppression."
4)]
5use anyhow::{Context, Result, anyhow, bail};
6use chrono::Utc;
7use futures::future::{BoxFuture, select_all};
8use std::collections::VecDeque;
9use std::path::PathBuf;
10use std::sync::Arc;
11use std::sync::atomic::{AtomicBool, Ordering};
12use tokio::sync::{Notify, RwLock};
13
14use crate::config::VTCodeConfig;
15use crate::config::types::ReasoningEffortLevel;
16use crate::core::agent::runner::{AgentRunner, RunnerSettings};
17use crate::core::agent::task::Task;
18use crate::core::threads::{ThreadBootstrap, ThreadId, ThreadRuntimeHandle, ThreadSnapshot};
19use crate::hooks::{LifecycleHookEngine, SessionStartTrigger};
20use crate::llm::provider::Message;
21use crate::tools::exec_session::ExecSessionManager;
22use crate::tools::pty::{PtyManager, PtySize};
23use crate::utils::session_archive::{SessionArchive, find_session_by_identifier};
24use vtcode_config::SubagentSpec;
25use vtcode_config::auth::OpenAIChatGptAuthHandle;
26
27use self::background::*;
28use self::config::*;
29use self::constants::*;
30use self::discovery::discover_controller_subagents;
31use self::model::*;
32use vtcode_config::subagents::SUBAGENT_HARD_CONCURRENCY_LIMIT;
33
34#[allow(
35    unused_imports,
36    reason = "Intentional compatibility, platform, or test-only suppression."
37)]
38use super::*;
39
40/// A dropped blocking launch must not leave its newly created worktree behind.
41struct LaunchWorktree {
42    root: PathBuf,
43    name: String,
44    path: PathBuf,
45    permit: Option<tokio::sync::OwnedSemaphorePermit>,
46    armed: bool,
47}
48
49impl Drop for LaunchWorktree {
50    fn drop(&mut self) {
51        if self.armed
52            && let Err(error) = crate::git::WorktreeManager::new(&self.root).remove(&self.name)
53        {
54            tracing::warn!(%error, worktree = %self.path.display(), "failed to roll back cancelled subagent launch");
55        }
56    }
57}
58
59async fn wait_for_child_stop(handle: JoinHandle<()>, child_id: &str) -> Result<()> {
60    match handle.await {
61        Ok(()) => Ok(()),
62        Err(error) if error.is_cancelled() => Ok(()),
63        Err(error) => Err(error).with_context(|| format!("subagent {child_id} launch failed to stop cleanly")),
64    }
65}
66
67impl SubagentController {
68    /// Spawns a new subagent child process from a [`SpawnAgentRequest`].
69    pub async fn spawn(&self, request: SpawnAgentRequest) -> Result<SubagentStatusEntry> {
70        self.ensure_ordinary_delegation_allowed()?;
71        let mut request = request;
72        let delegation = self
73            .prepare_delegation_context(
74                request.agent_type.clone(),
75                &mut request.items,
76                &mut request.model,
77                "spawn_agent",
78            )
79            .await?;
80        let spec = self.resolve_requested_spec(delegation.requested_agent.as_deref()).await?;
81        let prompt = self.prepare_delegation_prompt(
82            &spec,
83            &delegation,
84            &request.message,
85            &request.items,
86            "spawn_agent",
87            "spawning the subagent",
88        )?;
89        self.spawn_with_spec(
90            spec,
91            prompt,
92            request.fork_context,
93            request.background,
94            request.max_turns,
95            request.model,
96            request.reasoning_effort,
97        )
98        .await
99    }
100
101    /// Spawns a background subprocess for a subagent marked `background: true`.
102    pub async fn spawn_background_subprocess(
103        &self,
104        request: SpawnBackgroundSubprocessRequest,
105    ) -> Result<BackgroundSubprocessEntry> {
106        if self.config.managed_background_runtime {
107            bail!("managed background subprocesses cannot launch nested background subprocesses");
108        }
109        if !self.config.vt_cfg.subagents.background.enabled {
110            bail!("Background subagents are disabled by configuration");
111        }
112
113        let mut request = request;
114        let delegation = self
115            .prepare_delegation_context(
116                request.agent_type.clone(),
117                &mut request.items,
118                &mut request.model,
119                "spawn_background_subprocess",
120            )
121            .await?;
122        let spec = self.resolve_requested_spec(delegation.requested_agent.as_deref()).await?;
123        if !spec.background {
124            bail!(
125                "spawn_background_subprocess requires an agent with `background: true`; '{}' is a normal delegated child agent. Use spawn_agent instead.",
126                spec.name
127            );
128        }
129        let prompt = self.prepare_delegation_prompt(
130            &spec,
131            &delegation,
132            &request.message,
133            &request.items,
134            "spawn_background_subprocess",
135            "launching the background subprocess",
136        )?;
137        let desired_max_turns = normalize_background_child_max_turns(request.max_turns.or(spec.max_turns), true);
138        let desired_model_override = request.model.clone().or_else(|| spec.model.clone());
139        let desired_reasoning_override = request
140            .reasoning_effort
141            .clone()
142            .or_else(|| spec.reasoning_effort.as_ref().map(|e| e.as_str().to_string()));
143
144        let record_id = background_record_id(spec.name.as_str());
145        let _ = self.refresh_background_processes().await?;
146        {
147            let state = self.state.read().await;
148            if let Some(record) = state.background_children.get(&record_id)
149                && record.desired_enabled
150                && record.status.is_active()
151            {
152                let conflicts = Self::active_background_launch_conflicts(
153                    record,
154                    prompt.as_str(),
155                    desired_max_turns,
156                    desired_model_override.as_deref(),
157                    desired_reasoning_override.as_deref(),
158                );
159                if !conflicts.is_empty() {
160                    bail!(
161                        "spawn_background_subprocess found active background subprocess '{}' with different {}. Stop or restart the existing subprocess before changing its launch settings.",
162                        spec.name,
163                        conflicts.join(", "),
164                    );
165                }
166                return Ok(record.build_status_entry());
167            }
168        }
169
170        self.ensure_background_record_running(
171            spec.name.as_str(),
172            Some(record_id.as_str()),
173            0,
174            Some(BackgroundLaunchOverrides {
175                prompt: Some(prompt),
176                max_turns: request.max_turns,
177                model_override: request.model,
178                reasoning_override: request.reasoning_effort,
179            }),
180        )
181        .await
182    }
183
184    /// Spawns a subagent with a custom [`SubagentSpec`] that must be read-only.
185    pub async fn spawn_custom(&self, spec: SubagentSpec, request: SpawnAgentRequest) -> Result<SubagentStatusEntry> {
186        if !spec.is_subagent() {
187            bail!("custom subagent spawn only supports subagent-capable specs; '{}' is primary-only", spec.name);
188        }
189
190        if !spec.is_read_only() {
191            bail!(
192                "custom subagent spawn only supports read-only specs; '{}' exposes write-capable behavior",
193                spec.name
194            );
195        }
196
197        let mut request = request;
198        sanitize_subagent_input_items(&mut request.items);
199
200        let prompt = request_prompt(&request.message, &request.items)
201            .or_else(|| spec.initial_prompt.clone())
202            .filter(|value| !value.trim().is_empty())
203            .ok_or_else(|| anyhow!("custom subagent spawn requires a task message or items"))?;
204        if delegated_task_requires_clarification(&prompt) {
205            bail!(
206                "custom subagent task for '{}' is too vague ('{}'). Provide a specific delegated task before spawning the subagent.",
207                spec.name,
208                prompt.trim()
209            );
210        }
211
212        self.spawn_with_spec(
213            spec,
214            prompt,
215            request.fork_context,
216            request.background,
217            request.max_turns,
218            request.model,
219            request.reasoning_effort,
220        )
221        .await
222    }
223
224    /// Sends additional input to a running or queued subagent.
225    pub async fn send_input(&self, request: SendInputRequest) -> Result<SubagentStatusEntry> {
226        self.ensure_ordinary_delegation_allowed()?;
227        let prompt = request_prompt(&request.message, &request.items)
228            .ok_or_else(|| anyhow!("send_input requires a message or items"))?;
229
230        let maybe_restart = {
231            let mut state = self.state.write().await;
232            let record = state
233                .children
234                .get_mut(&request.target)
235                .ok_or_else(|| anyhow!("Unknown subagent id {}", request.target))?;
236
237            if record.status == SubagentStatus::Closed {
238                bail!("Subagent {} is closed", request.target);
239            }
240
241            record.updated_at = Utc::now();
242            record.last_prompt = Some(prompt.clone());
243
244            if request.interrupt {
245                if let Some(handle) = record.handle.as_ref() {
246                    handle.abort();
247                }
248                record.status = SubagentStatus::Queued;
249                record.queued_prompts.clear();
250                record.queued_prompts.push_back(prompt.clone());
251                true
252            } else if !record.status.is_terminal() {
253                record.status = SubagentStatus::Waiting;
254                record.queued_prompts.push_back(prompt.clone());
255                false
256            } else {
257                record.status = SubagentStatus::Queued;
258                record.queued_prompts.push_back(prompt.clone());
259                true
260            }
261        };
262
263        if maybe_restart {
264            self.restart_child(&request.target).await?;
265        }
266
267        self.status_for(&request.target).await
268    }
269
270    /// Resumes a closed or errored subagent and its descendants by re-queuing
271    /// their prompts. Cascades recursively through child-scoped controllers so
272    /// grandchildren closed by `close_tree` are resumed too, not merely
273    /// un-gated.
274    pub async fn resume(&self, target: &str) -> Result<SubagentStatusEntry> {
275        self.ensure_ordinary_delegation_allowed()?;
276        self.resume_tree(target).await
277    }
278
279    /// Recursively reopens `target` and every descendant across the whole
280    /// delegation tree (boxed for async recursion, mirroring `close_tree`).
281    fn resume_tree(&self, target: &str) -> BoxFuture<'static, Result<SubagentStatusEntry>> {
282        let self_owned = self.clone();
283        let target_owned = target.to_string();
284        Box::pin(async move {
285            if self_owned.shutdown_requested.load(Ordering::Relaxed) {
286                bail!("Subagent controller is shutting down; cannot resume subagents");
287            }
288            let subtree_ids = self_owned.collect_spawn_subtree_ids(&target_owned).await?;
289            let mut restart_ids = Vec::new();
290            for node_id in subtree_ids.iter() {
291                if self_owned.reopen_single(node_id.as_str()).await? {
292                    restart_ids.push(node_id.clone());
293                }
294            }
295            // Recursively resume each child-scoped controller's descendants,
296            // but only for controllers whose owning record is inside the
297            // resumed subtree (mirroring close_tree's sibling isolation).
298            // Grandchildren live in the child controller's state, so they are
299            // reopened through that controller to keep node ownership correct.
300            let nested = self_owned.nested_controllers_in_subtree(&subtree_ids).await;
301            for (controller, session_id) in nested {
302                let ids = controller.spawn_child_ids_for_parent(&session_id).await;
303                for id in ids {
304                    if let Err(err) = controller.resume_tree(&id).await {
305                        tracing::warn!(node_id = id.as_str(), error = %err, "Failed to resume nested subagent subtree");
306                    }
307                }
308            }
309            // Re-check shutdown before launching: a concurrent signal_shutdown
310            // between reopen and restart would otherwise start a child after
311            // shutdown aborted the subtree.
312            if self_owned.shutdown_requested.load(Ordering::Relaxed) {
313                bail!("Subagent controller is shutting down; cannot resume subagents");
314            }
315            for restart_id in restart_ids {
316                self_owned.restart_child(&restart_id).await?;
317            }
318            self_owned.status_for(&target_owned).await
319        })
320    }
321
322    /// Closes a subagent and all its descendants, aborting any in-flight work.
323    pub async fn close(&self, target: &str) -> Result<SubagentStatusEntry> {
324        // `close_tree` walks the full nesting tree (including grandchildren
325        // spawned through child-scoped controllers) and closes bottom-up.
326        self.close_tree(target).await
327    }
328
329    /// Closes `target` and every descendant across the whole delegation tree.
330    ///
331    /// Unlike [`Self::close`] this is recursive over child-scoped controllers,
332    /// so it works for arbitrary `max_depth`. The recursion is boxed to satisfy
333    /// Rust's async-fn recursion requirement.
334    ///
335    /// Ordering prevents a race where a still-running child spawns a new
336    /// descendant after the descendant snapshot: each child-scoped controller
337    /// is marked as closing (which rejects new spawns via `spawn_with_spec`)
338    /// before the owning child's handle is aborted and its subtree is closed.
339    fn close_tree(&self, target: &str) -> BoxFuture<'static, Result<SubagentStatusEntry>> {
340        let self_owned = self.clone();
341        let target_owned = target.to_string();
342        Box::pin(async move {
343            let subtree_ids = self_owned.collect_spawn_subtree_ids(&target_owned).await?;
344            let nested = self_owned.nested_controllers_in_subtree(&subtree_ids).await;
345            // Mark every child-scoped controller in the subtree as closing so
346            // it rejects new grandchild spawns before we start aborting.
347            for (controller, _) in &nested {
348                controller.begin_close().await;
349            }
350            // Close the target's own subtree (deepest first) on this
351            // controller. Aborting the owning child's handle first stops it
352            // from spawning further descendants.
353            let subtree_ids_for_rescan = subtree_ids.clone();
354            for node_id in subtree_ids.into_iter().rev() {
355                self_owned.close_single(node_id.as_str()).await?;
356            }
357            // ...then recursively close each child-scoped controller's subtree.
358            // Grandchildren live in the child controller's state, not ours, so
359            // they are closed through that controller to keep node ownership
360            // correct. Only controllers of nodes inside the closed subtree are
361            // cascaded, so closing one subagent never kills a sibling's
362            // descendants.
363            for (controller, session_id) in nested {
364                let ids = controller.spawn_child_ids_for_parent(&session_id).await;
365                for id in ids {
366                    if let Err(err) = controller.close_tree(&id).await {
367                        tracing::warn!(node_id = id.as_str(), error = %err, "Failed to close nested subagent subtree");
368                    }
369                }
370            }
371            // A child may have attached a freshly created controller between
372            // the snapshot above and its handle abort. Re-scan and cascade
373            // again so no grandchild outlives the close; the second pass is a
374            // no-op in the common case because close_tree is idempotent.
375            let late = self_owned.nested_controllers_in_subtree(&subtree_ids_for_rescan).await;
376            for (controller, session_id) in late {
377                controller.begin_close().await;
378                let ids = controller.spawn_child_ids_for_parent(&session_id).await;
379                for id in ids {
380                    if let Err(err) = controller.close_tree(&id).await {
381                        tracing::warn!(node_id = id.as_str(), error = %err, "Failed to close late nested subagent subtree");
382                    }
383                }
384            }
385            self_owned.status_for(&target_owned).await
386        })
387    }
388
389    /// Blocks until one of the target subagents reaches a terminal state or the timeout expires.
390    pub async fn wait(&self, targets: &[String], timeout_ms: Option<u64>) -> Result<Option<SubagentStatusEntry>> {
391        for target in targets {
392            if let Ok(entry) = self.status_for(target).await
393                && entry.status.is_terminal()
394            {
395                return Ok(Some(entry));
396            }
397        }
398
399        let timeout = std::time::Duration::from_millis(
400            timeout_ms.unwrap_or_else(|| self.config.vt_cfg.subagents.default_timeout_seconds.saturating_mul(1000)),
401        );
402        let deadline = tokio::time::Instant::now() + timeout;
403
404        loop {
405            // Collect notify handles from child records.
406            let notifies = {
407                let state = self.state.read().await;
408                targets
409                    .iter()
410                    .filter_map(|target| state.children.get(target).map(|record| record.notify.clone()))
411                    .collect::<Vec<_>>()
412            };
413            if notifies.is_empty() {
414                return Ok(None);
415            }
416
417            // Register notified() futures BEFORE checking terminal status.
418            // This prevents a Tokio Notify race condition: if apply_result()
419            // calls notify_waiters() between a status check and future
420            // creation, the notification is permanently lost. By registering
421            // futures first, any concurrent notification either:
422            //   (a) arrives before we poll the future → stored as a permit,
423            //       select! returns immediately, loop re-checks status, or
424            //   (b) arrives after we start waiting → wakes the future normally.
425            let wait_any = select_all(
426                notifies
427                    .into_iter()
428                    .map(|notify| Box::pin(async move { notify.notified().await }))
429                    .collect::<Vec<_>>(),
430            );
431            tokio::pin!(wait_any);
432
433            // Now check if any target is already terminal.
434            for target in targets {
435                if let Ok(entry) = self.status_for(target).await
436                    && entry.status.is_terminal()
437                {
438                    return Ok(Some(entry));
439                }
440            }
441
442            let sleep = tokio::time::sleep_until(deadline);
443            tokio::pin!(sleep);
444
445            tokio::select! {
446                _ = &mut sleep => return Ok(None),
447                _ = &mut wait_any => {}
448            }
449        }
450    }
451
452    /// Returns the current status of a tracked subagent by its target id.
453    pub async fn status_for(&self, target: &str) -> Result<SubagentStatusEntry> {
454        let state = self.state.read().await;
455        let record = state
456            .children
457            .get(target)
458            .ok_or_else(|| anyhow!("Unknown subagent id {target}"))?;
459        Ok(record.build_status_entry())
460    }
461
462    pub(super) async fn spawn_child_ids_for_parent(&self, parent_thread_id: &str) -> Vec<String> {
463        let state = self.state.read().await;
464        let mut child_ids = state
465            .children
466            .values()
467            .filter(|record| record.parent_thread_id == parent_thread_id)
468            .map(|record| record.id.clone())
469            .collect::<Vec<_>>();
470        child_ids.sort();
471        child_ids
472    }
473
474    pub(super) async fn collect_spawn_subtree_ids(&self, root_thread_id: &str) -> Result<Vec<String>> {
475        let mut subtree_ids = Vec::new();
476        let mut stack = vec![root_thread_id.to_string()];
477
478        while let Some(thread_id) = stack.pop() {
479            subtree_ids.push(thread_id.clone());
480            let child_ids = self.spawn_child_ids_for_parent(&thread_id).await;
481            for child_id in child_ids.into_iter().rev() {
482                stack.push(child_id);
483            }
484        }
485
486        Ok(subtree_ids)
487    }
488
489    /// Collects the child-scoped controllers owned by records inside `subtree_ids`,
490    /// paired with each owning record's session id. Shared by `resume_tree` and
491    /// `close_tree` so subtree isolation cannot diverge between them.
492    async fn nested_controllers_in_subtree(&self, subtree_ids: &[String]) -> Vec<(Arc<SubagentController>, String)> {
493        let subtree_set = subtree_ids.iter().collect::<std::collections::HashSet<_>>();
494        let state = self.state.read().await;
495        state
496            .children
497            .iter()
498            .filter(|(id, _)| subtree_set.contains(id))
499            .filter_map(|(_, record)| {
500                record
501                    .child_controller
502                    .clone()
503                    .map(|controller| (controller, record.session_id.clone()))
504            })
505            .collect()
506    }
507
508    pub(super) async fn reopen_single(&self, target: &str) -> Result<bool> {
509        let child_controller = {
510            let mut state = self.state.write().await;
511            let record = state
512                .children
513                .get_mut(target)
514                .ok_or_else(|| anyhow!("Unknown subagent id {target}"))?;
515            if !record.status.is_terminal() {
516                return Ok(false);
517            }
518            let prompt = record
519                .last_prompt
520                .clone()
521                .unwrap_or_else(|| "Continue the delegated task from the existing context.".to_string());
522            record.status = SubagentStatus::Queued;
523            record.updated_at = Utc::now();
524            record.completed_at = None;
525            record.error = None;
526            record.summary = None;
527            if record.queued_prompts.is_empty() {
528                record.queued_prompts.push_back(prompt);
529            }
530            record.child_controller.clone()
531        };
532        // Reopening a subtree reverses the transient `begin_close` on any
533        // child-scoped controller so a resumed child can delegate again (and
534        // its controller resumes saving background state).
535        if let Some(controller) = child_controller {
536            controller.end_close().await;
537        }
538        Ok(true)
539    }
540
541    async fn close_single(&self, target: &str) -> Result<SubagentStatusEntry> {
542        let mut state = self.state.write().await;
543        let record = state
544            .children
545            .get_mut(target)
546            .ok_or_else(|| anyhow!("Unknown subagent id {target}"))?;
547        if record.status == SubagentStatus::Closed {
548            return Ok(record.build_status_entry());
549        }
550        let handle = record.handle.take();
551        if let Some(handle) = handle.as_ref() {
552            handle.abort();
553        }
554        record.status = SubagentStatus::Closed;
555        record.updated_at = Utc::now();
556        record.completed_at = Some(Utc::now());
557        record.notify.notify_waiters();
558        let entry = record.build_status_entry();
559        drop(state);
560        if let Some(handle) = handle {
561            wait_for_child_stop(handle, target).await?;
562        }
563        Ok(entry)
564    }
565
566    pub(super) async fn background_status_for(&self, target: &str) -> Result<BackgroundSubprocessEntry> {
567        let state = self.state.read().await;
568        let record = state
569            .background_children
570            .get(target)
571            .ok_or_else(|| anyhow!("Unknown background subprocess {target}"))?;
572        Ok(record.build_status_entry())
573    }
574
575    pub(super) async fn ensure_background_record_running(
576        &self,
577        agent_name: &str,
578        stable_id: Option<&str>,
579        restart_attempts: u8,
580        overrides: Option<BackgroundLaunchOverrides>,
581    ) -> Result<BackgroundSubprocessEntry> {
582        let spec = self
583            .resolve_requested_spec(Some(agent_name))
584            .await
585            .with_context(|| format!("Failed to resolve background subagent '{agent_name}'"))?;
586        let record_id = stable_id
587            .map(ToOwned::to_owned)
588            .unwrap_or_else(|| background_record_id(agent_name));
589        let previous_record = {
590            let state = self.state.read().await;
591            state.background_children.get(&record_id).map(|record| {
592                (
593                    record.created_at,
594                    record.prompt.clone(),
595                    record.max_turns,
596                    record.model_override.clone(),
597                    record.reasoning_override.clone(),
598                )
599            })
600        };
601        let parent_session_id = self.parent_session_id.read().await.clone();
602        let session_id = format!(
603            "{}-{}-{}",
604            sanitize_component(parent_session_id.as_str()),
605            sanitize_component(record_id.as_str()),
606            Utc::now().format("%Y%m%dT%H%M%S%3fZ")
607        );
608        let exec_session_id = format!("exec-{session_id}");
609        let (created_at, previous_prompt, previous_max_turns, previous_model_override, previous_reasoning_override) =
610            previous_record.unwrap_or((Utc::now(), String::new(), None, None, None));
611        let prompt = overrides
612            .as_ref()
613            .and_then(|overrides| overrides.prompt.clone())
614            .filter(|value| !value.trim().is_empty())
615            .or_else(|| (!previous_prompt.trim().is_empty()).then_some(previous_prompt))
616            .or_else(|| spec.initial_prompt.clone())
617            .filter(|value| !value.trim().is_empty())
618            .unwrap_or_else(|| {
619                format!(
620                    "You are the VT Code background subagent `{}`, started without a specific task. Inspect the workspace at a high level, reply with a short readiness summary (what the project is and what you are set up to do), then end your turn; the process keeps running until it is stopped.",
621                    spec.name
622                )
623            });
624        let max_turns = normalize_background_child_max_turns(
625            overrides
626                .as_ref()
627                .and_then(|overrides| overrides.max_turns)
628                .or(previous_max_turns)
629                .or(spec.max_turns),
630            true,
631        );
632        let model_override = overrides
633            .as_ref()
634            .and_then(|overrides| overrides.model_override.clone())
635            .or(previous_model_override)
636            .or_else(|| spec.model.clone());
637        let reasoning_override = overrides
638            .as_ref()
639            .and_then(|overrides| overrides.reasoning_override.clone())
640            .or(previous_reasoning_override)
641            .or_else(|| spec.reasoning_effort.as_ref().map(|e| e.as_str().to_string()));
642
643        {
644            let mut state = self.state.write().await;
645            state.background_children.insert(
646                record_id.clone(),
647                BackgroundRecord {
648                    exit_code: None,
649                    termination_requested: false,
650                    id: record_id.clone(),
651                    agent_name: spec.name.clone(),
652                    display_label: subagent_display_label(&spec),
653                    description: spec.description.clone(),
654                    source: spec.source.label(),
655                    color: spec.color.clone(),
656                    session_id: session_id.clone(),
657                    exec_session_id: exec_session_id.clone(),
658                    desired_enabled: true,
659                    status: BackgroundSubprocessStatus::Starting,
660                    created_at,
661                    updated_at: Utc::now(),
662                    started_at: None,
663                    ended_at: None,
664                    pid: None,
665                    prompt: prompt.clone(),
666                    summary: Some("Starting background subagent".to_string()),
667                    error: None,
668                    archive_path: None,
669                    transcript_path: None,
670                    max_turns,
671                    model_override: model_override.clone(),
672                    reasoning_override: reasoning_override.clone(),
673                    restart_attempts,
674                },
675            );
676        }
677
678        let launch = build_background_launch_spec(
679            &self.config.workspace_root,
680            spec.name.as_str(),
681            parent_session_id.as_str(),
682            session_id.as_str(),
683            prompt.as_str(),
684            max_turns,
685            model_override.as_deref(),
686            reasoning_override.as_deref(),
687        )?;
688        let metadata = if launch.use_pty {
689            self.config
690                .exec_sessions
691                .create_pty_session_for_managed_background(
692                    exec_session_id.clone().into(),
693                    launch.command,
694                    self.config.workspace_root.clone(),
695                    PtySize {
696                        rows: 24,
697                        cols: 80,
698                        pixel_width: 0,
699                        pixel_height: 0,
700                    },
701                    hashbrown::HashMap::new(),
702                    None,
703                    hashbrown::HashMap::new(),
704                    false,
705                )
706                .await
707        } else {
708            self.config
709                .exec_sessions
710                .create_pipe_session_for_managed_background(
711                    exec_session_id.clone().into(),
712                    launch.command,
713                    self.config.workspace_root.clone(),
714                    hashbrown::HashMap::new(),
715                )
716                .await
717        }
718        .with_context(|| format!("Failed to spawn background subprocess for subagent '{}'", spec.name))?;
719
720        tracing::info!(
721            agent_name = spec.name.as_str(),
722            record_id = record_id.as_str(),
723            exec_session_id = exec_session_id.as_str(),
724            pid = metadata.child_pid,
725            "Spawned background subagent subprocess"
726        );
727
728        {
729            let mut state = self.state.write().await;
730            let record = state
731                .background_children
732                .get_mut(&record_id)
733                .ok_or_else(|| anyhow!("Unknown background subprocess {record_id}"))?;
734            finalize_background_launch(
735                record,
736                exec_session_id.as_str(),
737                metadata.child_pid,
738                metadata.started_at,
739                Utc::now(),
740            );
741        }
742
743        self.save_background_state().await?;
744        self.background_status_for(&record_id).await
745    }
746
747    pub(super) async fn refresh_background_archive_metadata(&self, target: &str) -> Result<()> {
748        let session_id = {
749            let state = self.state.read().await;
750            state
751                .background_children
752                .get(target)
753                .map(|record| record.session_id.clone())
754                .ok_or_else(|| anyhow!("Unknown background subprocess {target}"))?
755        };
756
757        if let Some(listing) = find_session_by_identifier(&session_id).await? {
758            let mut state = self.state.write().await;
759            if let Some(record) = state.background_children.get_mut(target) {
760                record.archive_path = Some(listing.path.clone());
761                record.transcript_path = Some(listing.path);
762            }
763        }
764
765        Ok(())
766    }
767
768    /// Marks this controller as closing so [`Self::spawn`] rejects new spawns.
769    ///
770    /// Used before aborting a subtree so a still-running child cannot spawn a
771    /// descendant between the descendant snapshot and the handle abort. Unlike
772    /// [`Self::signal_shutdown`] this uses the transient `closing` flag, which
773    /// is cleared again when the subtree is reopened, so a resumed child can
774    /// delegate and its controller keeps saving background state.
775    pub(super) async fn begin_close(&self) {
776        self.closing.store(true, Ordering::Relaxed);
777    }
778
779    /// Clears the transient close-in-progress flag. Called when a closed
780    /// subtree is reopened so its child-scoped controller can delegate again.
781    pub(super) async fn end_close(&self) {
782        self.closing.store(false, Ordering::Relaxed);
783    }
784
785    /// Signal that the program is shutting down. Subsequent calls to
786    /// `save_background_state` will be skipped. All running child handles
787    /// are aborted so subagent tasks do not outlive the parent session.
788    pub async fn signal_shutdown(&self) {
789        if let Err(error) = self.cancel_matrix().await {
790            self.matrix.cancellation.read().cancel();
791            tracing::warn!(%error, "failed to persist matrix shutdown cancellation");
792        }
793        self.shutdown_requested.store(true, Ordering::Relaxed);
794        self.stop_background_completion_monitor().await;
795        let nested = {
796            let mut state = self.state.write().await;
797            let mut nested = Vec::new();
798            for record in state.children.values_mut() {
799                if let Some(handle) = record.handle.take() {
800                    handle.abort();
801                }
802                record.status = SubagentStatus::Closed;
803                record.completed_at = Some(Utc::now());
804                record.notify.notify_waiters();
805                if let Some(controller) = record.child_controller.clone() {
806                    nested.push((controller, record.session_id.clone()));
807                }
808            }
809            nested
810        };
811        // Mark every child-scoped controller as permanently shut down (so its
812        // background state is not saved) and as closing (so a race cannot let
813        // a fresh grandchild outlive shutdown), then cascade.
814        for (controller, _) in &nested {
815            controller.shutdown_requested.store(true, Ordering::Relaxed);
816            controller.begin_close().await;
817            controller.stop_background_completion_monitor().await;
818        }
819        // Cascade shutdown to child-scoped controllers so grandchildren tasks
820        // are aborted too; otherwise their tokio tasks keep running detached.
821        for (controller, session_id) in nested {
822            let ids = controller.spawn_child_ids_for_parent(&session_id).await;
823            for id in ids {
824                if let Err(err) = controller.close_tree(&id).await {
825                    tracing::warn!(node_id = id.as_str(), error = %err, "Failed to close nested subagent subtree during shutdown");
826                }
827            }
828        }
829    }
830
831    pub(super) async fn save_background_state(&self) -> Result<()> {
832        if self.shutdown_requested.load(Ordering::Relaxed) {
833            return Ok(());
834        }
835        let records = {
836            let state = self.state.read().await;
837            state
838                .background_children
839                .values()
840                .cloned()
841                .map(BackgroundRecord::into_persisted)
842                .collect()
843        };
844        persist_background_state(&self.config.workspace_root, records).await
845    }
846
847    pub(super) async fn find_spec(&self, candidate: &str) -> Option<SubagentSpec> {
848        self.state
849            .read()
850            .await
851            .discovered
852            .effective
853            .iter()
854            .find(|spec| spec.is_subagent() && spec.matches_name(candidate))
855            .cloned()
856    }
857
858    pub(super) async fn resolve_requested_spec(&self, requested: Option<&str>) -> Result<SubagentSpec> {
859        let requested = requested.unwrap_or("default");
860        self.find_spec(requested)
861            .await
862            .ok_or_else(|| anyhow!("Unknown subagent type {requested}"))
863    }
864
865    async fn prepare_delegation_context(
866        &self,
867        requested_agent: Option<String>,
868        items: &mut Vec<SubagentInputItem>,
869        model: &mut Option<String>,
870        tool_name: &'static str,
871    ) -> Result<PreparedDelegationContext> {
872        let state = self.state.read().await;
873        sanitize_subagent_input_items(items);
874        *model = normalize_requested_model_override(model.take(), &state.turn_hints.current_input);
875        let requested_agent = if let Some(agent_type) = requested_agent {
876            Some(agent_type)
877        } else {
878            match state.turn_hints.explicit_mentions.as_slice() {
879                [] => None,
880                [single] => Some(single.clone()),
881                mentions => {
882                    bail!(
883                        "{} omitted agent_type, but the user explicitly selected multiple agents: {}. Specify agent_type explicitly.",
884                        tool_name,
885                        mentions.join(", ")
886                    );
887                }
888            }
889        };
890        Ok(PreparedDelegationContext {
891            requested_agent,
892            explicit_mentions: state.turn_hints.explicit_mentions.clone(),
893            explicit_request: state.turn_hints.explicit_request,
894        })
895    }
896
897    fn prepare_delegation_prompt(
898        &self,
899        spec: &SubagentSpec,
900        delegation: &PreparedDelegationContext,
901        message: &Option<String>,
902        items: &[SubagentInputItem],
903        tool_name: &'static str,
904        launch_phrase: &'static str,
905    ) -> Result<String> {
906        if let Some(explicit) = delegation.explicit_mentions.first()
907            && delegation.explicit_mentions.len() == 1
908            && !spec.matches_name(explicit)
909        {
910            bail!(
911                "{} requested agent_type '{}', but the user explicitly selected '{}'. Use the selected agent or ask the user to clarify.",
912                tool_name,
913                spec.name,
914                explicit
915            );
916        }
917        if !spec.is_read_only() && !delegation.explicit_request && delegation.requested_agent.is_none() {
918            bail!(
919                "{} cannot launch write-capable agent '{}' without an explicit delegation signal from the current user turn. Ask the user to mention the agent, say 'delegate'/'spawn', or request parallel work.",
920                tool_name,
921                spec.name
922            );
923        }
924        if spec.is_read_only() && !self.config.vt_cfg.subagents.auto_delegate_read_only && !delegation.explicit_request
925        {
926            bail!(
927                "{} cannot proactively launch read-only agent '{}' because `subagents.auto_delegate_read_only` is disabled and the current user turn did not explicitly request delegation.",
928                tool_name,
929                spec.name
930            );
931        }
932        let prompt = request_prompt(message, items)
933            .or_else(|| spec.initial_prompt.clone())
934            .filter(|value| !value.trim().is_empty())
935            .ok_or_else(|| anyhow!("{tool_name} requires a task message or items"))?;
936        if delegated_task_requires_clarification(&prompt) {
937            bail!(
938                "{} task for '{}' is too vague ('{}'). Ask the user for a specific delegated task before {}.",
939                tool_name,
940                spec.name,
941                prompt.trim(),
942                launch_phrase
943            );
944        }
945        Ok(prompt)
946    }
947
948    fn active_background_launch_conflicts(
949        record: &BackgroundRecord,
950        prompt: &str,
951        max_turns: Option<usize>,
952        model_override: Option<&str>,
953        reasoning_override: Option<&str>,
954    ) -> Vec<&'static str> {
955        let mut conflicts = Vec::new();
956        if record.prompt != prompt {
957            conflicts.push("prompt");
958        }
959        if record.max_turns != max_turns {
960            conflicts.push("max_turns");
961        }
962        if record.model_override.as_deref() != model_override {
963            conflicts.push("model");
964        }
965        if record.reasoning_override.as_deref() != reasoning_override {
966            conflicts.push("reasoning_effort");
967        }
968        conflicts
969    }
970
971    async fn spawn_with_spec(
972        &self,
973        spec: SubagentSpec,
974        prompt: String,
975        fork_context: bool,
976        background: bool,
977        max_turns: Option<usize>,
978        model_override: Option<String>,
979        reasoning_override: Option<String>,
980    ) -> Result<SubagentStatusEntry> {
981        if !self.config.vt_cfg.subagents.enabled {
982            bail!("Subagents are disabled by configuration");
983        }
984        if self.shutdown_requested.load(Ordering::Relaxed) || self.closing.load(Ordering::Relaxed) {
985            bail!("Subagent controller is shutting down; cannot spawn new subagents");
986        }
987        if self.config.depth.saturating_add(1) > self.config.vt_cfg.subagents.max_depth {
988            bail!("Subagent depth limit reached (max_depth={})", self.config.vt_cfg.subagents.max_depth);
989        }
990        if self.config.depth > 0 && spec.isolation == Some(vtcode_config::IsolationMode::Worktree) {
991            bail!(
992                "Subagent '{}' requests isolation=worktree, but nested worktree isolation is not supported \
993                 (child-scoped controllers operate inside the parent's worktree). Use isolation=worktree only \
994                 at the root delegation level.",
995                spec.name
996            );
997        }
998        let is_background_child = background;
999        let child_max_turns = normalize_background_child_max_turns(max_turns.or(spec.max_turns), is_background_child);
1000        let (_, _, effective_config) = prepare_child_runtime_config(
1001            &self.config.vt_cfg,
1002            &spec,
1003            self.config.parent_model.as_str(),
1004            self.config.parent_provider.as_str(),
1005            self.config.parent_reasoning_effort,
1006            child_max_turns,
1007            model_override.as_deref(),
1008            reasoning_override.as_deref(),
1009            !spec.is_read_only() && self.config.depth.saturating_add(2) <= self.config.vt_cfg.subagents.max_depth,
1010            resolve_effective_subagent_model,
1011        )?;
1012
1013        // Reserve admission before any side effect. The permit is owned by the
1014        // launched task and is released on every launch-error/cancellation path.
1015        let permit = Arc::clone(&self.admission)
1016            .try_acquire_owned()
1017            .context("Subagent concurrency limit reached")?;
1018        let launch_id = uuid::Uuid::new_v4();
1019        {
1020            let state = self.state.read().await;
1021            let active = state.children.values().filter(|record| !record.status.is_terminal()).count();
1022            let cap = self.config.vt_cfg.subagents.max_concurrent.min(SUBAGENT_HARD_CONCURRENCY_LIMIT);
1023            if active >= cap {
1024                bail!("Subagent concurrency limit reached (max_concurrent={cap})");
1025            }
1026        }
1027        self.ensure_ordinary_delegation_allowed()?;
1028        // Create a worktree for isolation if requested.
1029        let (mut worktree_guard, permit) = if spec.isolation == Some(vtcode_config::IsolationMode::Worktree) {
1030            let workspace_root = self.config.workspace_root.clone();
1031            let worktree_name = format!("{}-{launch_id}", sanitize_component(spec.name.as_str()));
1032            let worktree_name_for_error = worktree_name.clone();
1033            let worktree_result = tokio::task::spawn_blocking(move || {
1034                let path = crate::git::WorktreeManager::new(&workspace_root).create(&worktree_name)?;
1035                Ok::<_, anyhow::Error>(LaunchWorktree {
1036                    root: workspace_root,
1037                    name: worktree_name,
1038                    path,
1039                    permit: Some(permit),
1040                    armed: true,
1041                })
1042            })
1043            .await
1044            .context("Worktree creation task panicked")?;
1045            let mut guard = worktree_result
1046                .with_context(|| format!("Failed to create worktree for subagent '{}'", worktree_name_for_error))?;
1047            let permit = guard.permit.take().context("worktree launch reservation missing")?;
1048            (Some(guard), permit)
1049        } else {
1050            (None, permit)
1051        };
1052        let worktree_path = worktree_guard.as_ref().map(|guard| guard.path.clone());
1053
1054        let id = format!("agent-{}-{launch_id}", sanitize_component(spec.name.as_str()));
1055        let parent_session_id = self.parent_session_id.read().await.clone();
1056        let session_id =
1057            format!("{}-{}", sanitize_component(parent_session_id.as_str()), sanitize_component(id.as_str()));
1058        let display_label = subagent_display_label(&spec);
1059        let notify = Arc::new(Notify::new());
1060        let mut state = self.state.write().await;
1061        // Re-check the close gates while the write lock is held. The earlier
1062        // check can race with a concurrent `close_tree` that snapshots and
1063        // closes the subtree between this spawn's admission and its record
1064        // insertion; checking after acquisition closes that window.
1065        if self.shutdown_requested.load(Ordering::Relaxed) || self.closing.load(Ordering::Relaxed) {
1066            drop(state);
1067            self.remove_failed_launch_worktree(worktree_path.as_ref()).await?;
1068            if let Some(guard) = worktree_guard.as_mut() {
1069                guard.armed = false;
1070            }
1071            bail!("Subagent controller is shutting down; cannot spawn new subagents");
1072        }
1073        let initial_messages = if fork_context {
1074            state.parent_messages.clone()
1075        } else {
1076            Vec::new()
1077        };
1078        let entry = ChildRecord {
1079            id: id.clone(),
1080            session_id,
1081            parent_thread_id: parent_session_id,
1082            spec: spec.clone(),
1083            display_label,
1084            status: SubagentStatus::Queued,
1085            background: is_background_child,
1086            depth: self.config.depth.saturating_add(1),
1087            created_at: Utc::now(),
1088            updated_at: Utc::now(),
1089            completed_at: None,
1090            summary: None,
1091            error: None,
1092            archive_metadata: None,
1093            archive_path: None,
1094            transcript_path: None,
1095            effective_config: Some(effective_config),
1096            stored_messages: initial_messages,
1097            last_prompt: Some(prompt.clone()),
1098            queued_prompts: VecDeque::from([prompt]),
1099            max_turns: child_max_turns,
1100            model_override,
1101            reasoning_override,
1102            thread_handle: None,
1103            handle: None,
1104            notify,
1105            worktree_path,
1106            child_controller: None,
1107        };
1108        state.children.insert(id.clone(), entry);
1109        drop(state);
1110
1111        if let Err(error) = self.launch_child_reserved(id.as_str(), permit).await {
1112            let worktree = self
1113                .state
1114                .write()
1115                .await
1116                .children
1117                .remove(&id)
1118                .and_then(|record| record.worktree_path);
1119            self.remove_failed_launch_worktree(worktree.as_ref())
1120                .await
1121                .with_context(|| format!("launch failed: {error:#}; worktree rollback failed"))?;
1122            if let Some(guard) = worktree_guard.as_mut() {
1123                guard.armed = false;
1124            }
1125            return Err(error);
1126        }
1127        if let Some(guard) = worktree_guard.as_mut() {
1128            guard.armed = false;
1129        }
1130        self.status_for(&id).await
1131    }
1132
1133    async fn remove_failed_launch_worktree(&self, path: Option<&std::path::PathBuf>) -> Result<()> {
1134        if let Some(path) = path {
1135            let name = path
1136                .file_name()
1137                .and_then(|name| name.to_str())
1138                .context("invalid owned worktree name")?
1139                .to_owned();
1140            let root = self.config.workspace_root.clone();
1141            tokio::task::spawn_blocking(move || crate::git::WorktreeManager::new(root).remove(&name))
1142                .await
1143                .context("failed launch worktree cleanup panicked")??;
1144        }
1145        Ok(())
1146    }
1147
1148    async fn restart_child(&self, target: &str) -> Result<()> {
1149        let finishing_handle = {
1150            let mut state = self.state.write().await;
1151            let record = state
1152                .children
1153                .get_mut(target)
1154                .ok_or_else(|| anyhow!("Unknown subagent id {target}"))?;
1155            if record.queued_prompts.is_empty()
1156                && let Some(prompt) = record.last_prompt.clone()
1157            {
1158                record.queued_prompts.push_back(prompt);
1159            }
1160            if record.queued_prompts.is_empty() {
1161                bail!("Subagent {target} has no queued input");
1162            }
1163            record.handle.take()
1164        };
1165        let launch_result = async {
1166            // Abort schedules cancellation; await the owned launch before
1167            // reserving its replacement's concurrency slot.
1168            if let Some(handle) = finishing_handle {
1169                wait_for_child_stop(handle, target).await?;
1170            }
1171            self.launch_child(target).await
1172        }
1173        .await;
1174        if let Err(error) = &launch_result {
1175            let mut state = self.state.write().await;
1176            if let Some(record) = state.children.get_mut(target)
1177                && record.status == SubagentStatus::Queued
1178                && record.handle.as_ref().is_none_or(|handle| handle.is_finished())
1179            {
1180                record.status = SubagentStatus::Failed;
1181                record.error = Some(format!("{error:#}"));
1182                record.summary = None;
1183                record.updated_at = Utc::now();
1184                record.completed_at = Some(record.updated_at);
1185                record.notify.notify_waiters();
1186            }
1187        }
1188        launch_result
1189    }
1190}
1191
1192pub(super) fn finalize_background_launch(
1193    record: &mut BackgroundRecord,
1194    expected_exec_session_id: &str,
1195    child_pid: Option<u32>,
1196    started_at: Option<chrono::DateTime<Utc>>,
1197    updated_at: chrono::DateTime<Utc>,
1198) -> bool {
1199    if record.exec_session_id != expected_exec_session_id
1200        || !matches!(record.status, BackgroundSubprocessStatus::Starting)
1201    {
1202        return false;
1203    }
1204    record.pid = child_pid;
1205    record.started_at = started_at;
1206    record.status = BackgroundSubprocessStatus::Running;
1207    record.updated_at = updated_at;
1208    record.ended_at = None;
1209    record.error = None;
1210    record.summary = Some("Background subagent is running".to_string());
1211    true
1212}