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
40struct 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 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 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 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 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 pub async fn resume(&self, target: &str) -> Result<SubagentStatusEntry> {
275 self.ensure_ordinary_delegation_allowed()?;
276 self.resume_tree(target).await
277 }
278
279 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 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 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 pub async fn close(&self, target: &str) -> Result<SubagentStatusEntry> {
324 self.close_tree(target).await
327 }
328
329 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 for (controller, _) in &nested {
348 controller.begin_close().await;
349 }
350 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 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 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 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 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 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 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 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 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 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 pub(super) async fn begin_close(&self) {
776 self.closing.store(true, Ordering::Relaxed);
777 }
778
779 pub(super) async fn end_close(&self) {
782 self.closing.store(false, Ordering::Relaxed);
783 }
784
785 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 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 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 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 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 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 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}