1#![allow(
2 unused_imports,
3 reason = "Intentional compatibility, platform, or test-only suppression."
4)]
5mod background;
8mod config;
9mod constants;
10mod discovery;
11pub mod matrix;
12mod model;
13mod prompt;
14mod types;
15
16pub use background::{
19 background_record_id, build_background_subagent_command, extract_tail_lines, load_archive_preview,
20 subagent_display_label,
21};
22pub use config::{
23 ResolvedAgentRuntimeView, build_child_config, compose_subagent_instructions, filter_child_tools,
24 normalize_background_child_max_turns, normalize_child_max_turns, prepare_child_runtime_config,
25};
26pub use discovery::discover_controller_subagents;
27pub use model::{
28 agent_type_for_spec, load_memory_appendix, load_memory_appendix_async, load_primary_memory_appendix,
29 load_primary_memory_appendix_async,
30};
31pub use prompt::{
32 contains_explicit_delegation_request, contains_explicit_model_request, delegated_task_requires_clarification,
33 extract_explicit_agent_mentions, normalize_requested_model_override, request_prompt, sanitize_subagent_input_items,
34};
35pub use types::{
36 BackgroundCompletionEvent, BackgroundRecord, BackgroundSubprocessEntry, BackgroundSubprocessSnapshot,
37 BackgroundSubprocessStatus, ChildRecord, ChildRunResult, ControllerState, PersistedBackgroundRecord,
38 PersistedBackgroundState, SendInputRequest, SpawnAgentRequest, SpawnBackgroundSubprocessRequest,
39 StatusEntryBuilder, SubagentInputItem, SubagentStatus, SubagentStatusEntry, SubagentThreadSnapshot,
40 TurnDelegationHints,
41};
42
43pub fn is_subagent_tool(name: &str) -> bool {
50 SUBAGENT_TOOL_NAMES.contains(&name)
51}
52
53#[derive(Clone, Default)]
54pub(super) struct BackgroundLaunchOverrides {
55 prompt: Option<String>,
56 max_turns: Option<usize>,
57 model_override: Option<String>,
58 reasoning_override: Option<String>,
59}
60
61#[derive(Clone, Default)]
62pub(super) struct PreparedDelegationContext {
63 requested_agent: Option<String>,
64 explicit_mentions: Vec<String>,
65 explicit_request: bool,
66}
67
68#[derive(Debug, Clone)]
73pub struct VerificationResult {
74 pub approved: bool,
76 pub issues: Vec<String>,
78 pub reasoning: String,
80}
81
82use anyhow::{Context, Result, anyhow, bail};
85use chrono::Utc;
86use futures::future::select_all;
87use parking_lot::Mutex as ParkingMutex;
88use std::collections::VecDeque;
89use std::path::PathBuf;
90use std::sync::Arc;
91use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
92use tokio::sync::{Mutex, Notify, RwLock, broadcast};
93use tokio::task::JoinHandle;
94use tokio_util::sync::CancellationToken;
95
96use crate::config::VTCodeConfig;
97use crate::config::types::ReasoningEffortLevel;
98use crate::core::agent::runner::{AgentRunner, RunnerSettings};
99use crate::core::agent::task::Task;
100use crate::core::threads::{ThreadBootstrap, ThreadId, ThreadRuntimeHandle, ThreadSnapshot};
101use crate::hooks::{LifecycleHookEngine, SessionStartTrigger};
102use crate::llm::provider::Message;
103use crate::tools::exec_session::ExecSessionManager;
104use crate::tools::pty::{PtyManager, PtySize};
105use crate::utils::session_archive::{SessionArchive, find_session_by_identifier};
106use vtcode_config::SubagentSpec;
107use vtcode_config::auth::OpenAIChatGptAuthHandle;
108
109use self::background::*;
110use self::config::*;
111use self::constants::*;
112use self::model::*;
113use vtcode_config::subagents::SUBAGENT_HARD_CONCURRENCY_LIMIT;
114
115const BACKGROUND_COMPLETION_CHANNEL_CAPACITY: usize = 64;
116
117struct BackgroundCompletionChannel {
118 sender: broadcast::Sender<BackgroundCompletionEvent>,
119 parent_sender: broadcast::Sender<BackgroundCompletionEvent>,
120 pending_parent_events: VecDeque<BackgroundCompletionEvent>,
121 parent_subscribed: bool,
122}
123
124impl BackgroundCompletionChannel {
125 fn new() -> Self {
126 let (sender, _) = broadcast::channel(BACKGROUND_COMPLETION_CHANNEL_CAPACITY);
127 let (parent_sender, _) = broadcast::channel(BACKGROUND_COMPLETION_CHANNEL_CAPACITY);
128 Self {
129 sender,
130 parent_sender,
131 pending_parent_events: VecDeque::new(),
132 parent_subscribed: false,
133 }
134 }
135
136 fn subscribe(&self) -> broadcast::Receiver<BackgroundCompletionEvent> {
137 self.sender.subscribe()
138 }
139
140 fn subscribe_parent(&mut self) -> (broadcast::Receiver<BackgroundCompletionEvent>, bool) {
141 let receiver = self.parent_sender.subscribe();
142 self.parent_subscribed = true;
143 let replayed = !self.pending_parent_events.is_empty();
144 while let Some(event) = self.pending_parent_events.pop_front() {
145 let _ = self.parent_sender.send(event);
146 }
147 (receiver, replayed)
148 }
149
150 fn publish(&mut self, event: BackgroundCompletionEvent) {
151 if !self.parent_subscribed {
152 if self.pending_parent_events.len() >= BACKGROUND_COMPLETION_CHANNEL_CAPACITY {
153 self.pending_parent_events.pop_front();
154 tracing::warn!(
155 capacity = BACKGROUND_COMPLETION_CHANNEL_CAPACITY,
156 "Dropping oldest undelivered parent background completion"
157 );
158 }
159 self.pending_parent_events.push_back(event.clone());
160 } else {
161 let _ = self.parent_sender.send(event.clone());
162 }
163 let _ = self.sender.send(event);
164 }
165}
166
167#[derive(Clone)]
171pub struct SubagentControllerConfig {
172 pub workspace_root: PathBuf,
174 pub parent_session_id: String,
176 pub parent_model: String,
178 pub parent_provider: String,
180 pub parent_reasoning_effort: ReasoningEffortLevel,
182 pub api_key: String,
184 pub vt_cfg: VTCodeConfig,
186 pub openai_chatgpt_auth: Option<OpenAIChatGptAuthHandle>,
188 pub depth: usize,
190 pub workspace_gated: bool,
197 pub exec_sessions: ExecSessionManager,
199 pub pty_manager: PtyManager,
201 pub managed_background_runtime: bool,
203}
204
205pub struct SubagentController {
209 admission: Arc<tokio::sync::Semaphore>,
210 matrix: Arc<matrix::MatrixRuntime>,
211 config: Arc<SubagentControllerConfig>,
212 parent_session_id: Arc<RwLock<String>>,
213 lifecycle_hooks: Option<LifecycleHookEngine>,
214 state: Arc<RwLock<ControllerState>>,
215 shutdown_requested: Arc<AtomicBool>,
216 closing: Arc<AtomicBool>,
222 background_completion_channel: Arc<ParkingMutex<BackgroundCompletionChannel>>,
223 background_completion_notify: Arc<Notify>,
224 background_completion_shutdown: CancellationToken,
225 background_completion_monitor: Arc<Mutex<Option<JoinHandle<()>>>>,
226 background_completion_owners: Arc<AtomicUsize>,
228 background_completion_monitor_owner: bool,
231}
232
233impl Clone for SubagentController {
234 fn clone(&self) -> Self {
235 self.background_completion_owners.fetch_add(1, Ordering::Relaxed);
236 Self {
237 admission: Arc::clone(&self.admission),
238 matrix: Arc::clone(&self.matrix),
239 config: Arc::clone(&self.config),
240 parent_session_id: Arc::clone(&self.parent_session_id),
241 lifecycle_hooks: self.lifecycle_hooks.clone(),
242 state: Arc::clone(&self.state),
243 shutdown_requested: Arc::clone(&self.shutdown_requested),
244 closing: Arc::clone(&self.closing),
245 background_completion_channel: Arc::clone(&self.background_completion_channel),
246 background_completion_notify: Arc::clone(&self.background_completion_notify),
247 background_completion_shutdown: self.background_completion_shutdown.clone(),
248 background_completion_monitor: Arc::clone(&self.background_completion_monitor),
249 background_completion_owners: Arc::clone(&self.background_completion_owners),
250 background_completion_monitor_owner: true,
251 }
252 }
253}
254
255impl Drop for SubagentController {
256 fn drop(&mut self) {
257 if !self.background_completion_monitor_owner
258 || self.background_completion_owners.fetch_sub(1, Ordering::AcqRel) != 1
259 {
260 return;
261 }
262
263 self.background_completion_shutdown.cancel();
264 if let Ok(mut monitor_slot) = self.background_completion_monitor.try_lock()
265 && let Some(monitor) = monitor_slot.take()
266 {
267 monitor.abort();
268 }
269 }
270}
271
272impl SubagentController {
273 pub async fn new(config: SubagentControllerConfig) -> Result<Self> {
275 let discovered = discover_controller_subagents(&config.workspace_root).await?;
276 Box::pin(Self::new_with_discovered(config, discovered)).await
280 }
281
282 pub async fn new_with_discovered(
288 config: SubagentControllerConfig,
289 discovered: vtcode_config::DiscoveredSubagents,
290 ) -> Result<Self> {
291 let workspace_gated = config.workspace_gated;
292 let lifecycle_hooks = LifecycleHookEngine::new_with_session_gated(
293 config.workspace_root.clone(),
294 &config.vt_cfg.hooks,
295 SessionStartTrigger::Startup,
296 config.parent_session_id.clone(),
297 workspace_gated,
298 )?;
299 if let Some(engine) = lifecycle_hooks.as_ref() {
300 crate::hooks::lifecycle::restore_workspace_hook_approval(engine, &config.workspace_root).await;
301 }
302 let background_children = load_background_state(&config.workspace_root)
303 .await?
304 .records
305 .into_iter()
306 .map(|record| (record.id.clone(), BackgroundRecord::from_persisted(record)))
307 .collect();
308 let controller = Self {
309 admission: Arc::new(tokio::sync::Semaphore::new(
310 config.vt_cfg.subagents.max_concurrent.min(SUBAGENT_HARD_CONCURRENCY_LIMIT),
311 )),
312 matrix: Arc::new(matrix::MatrixRuntime::default()),
313 parent_session_id: Arc::new(RwLock::new(config.parent_session_id.clone())),
314 lifecycle_hooks,
315 config: Arc::new(config),
316 state: Arc::new(RwLock::new(ControllerState {
317 discovered,
318 parent_messages: Vec::new(),
319 turn_hints: TurnDelegationHints::default(),
320 children: std::collections::BTreeMap::new(),
321 background_children,
322 background_completion_identities: VecDeque::new(),
323 })),
324 shutdown_requested: Arc::new(AtomicBool::new(false)),
325 closing: Arc::new(AtomicBool::new(false)),
326 background_completion_channel: Arc::new(ParkingMutex::new(BackgroundCompletionChannel::new())),
327 background_completion_notify: Arc::new(Notify::new()),
328 background_completion_shutdown: CancellationToken::new(),
329 background_completion_monitor: Arc::new(Mutex::new(None)),
330 background_completion_owners: Arc::new(AtomicUsize::new(1)),
331 background_completion_monitor_owner: true,
332 };
333 controller.start_background_completion_monitor().await;
334 Ok(controller)
335 }
336
337 pub fn subscribe_background_completions(&self) -> broadcast::Receiver<BackgroundCompletionEvent> {
339 self.background_completion_channel.lock().subscribe()
340 }
341
342 pub fn subscribe_parent_background_completions(&self) -> broadcast::Receiver<BackgroundCompletionEvent> {
345 let (receiver, replayed) = self.background_completion_channel.lock().subscribe_parent();
346 if replayed {
347 self.background_completion_notify.notify_one();
348 }
349 receiver
350 }
351
352 pub fn background_completion_notify(&self) -> Arc<Notify> {
354 Arc::clone(&self.background_completion_notify)
355 }
356
357 pub async fn reload(&self) -> Result<()> {
359 let discovered = discover_controller_subagents(&self.config.workspace_root).await?;
360 self.state.write().await.discovered = discovered;
361 Ok(())
362 }
363
364 pub async fn set_parent_messages(&self, messages: &[Message]) {
366 let cloned = messages.to_vec();
367 self.state.write().await.parent_messages = cloned;
368 }
369
370 pub async fn set_turn_delegation_hints_from_input(&self, input: &str) -> Vec<String> {
372 let mut state = self.state.write().await;
373 let explicit_mentions = extract_explicit_agent_mentions(input, state.discovered.effective.as_slice());
374 let explicit_request = contains_explicit_delegation_request(input, explicit_mentions.as_slice());
375 state.turn_hints = TurnDelegationHints {
376 explicit_mentions: explicit_mentions.clone(),
377 explicit_request,
378 current_input: input.to_string(),
379 };
380 explicit_mentions
381 }
382
383 pub async fn clear_turn_delegation_hints(&self) {
385 self.state.write().await.turn_hints = TurnDelegationHints::default();
386 }
387
388 pub async fn set_parent_session_id(&self, session_id: impl Into<String>) {
390 *self.parent_session_id.write().await = session_id.into();
391 }
392
393 pub async fn effective_specs(&self) -> Vec<SubagentSpec> {
395 self.state.read().await.discovered.effective.clone()
396 }
397
398 pub async fn shadowed_specs(&self) -> Vec<SubagentSpec> {
400 self.state.read().await.discovered.shadowed.clone()
401 }
402
403 pub async fn status_entries(&self) -> Vec<SubagentStatusEntry> {
405 let state = self.state.read().await;
406 let mut entries = state.children.values().map(ChildRecord::build_status_entry).collect::<Vec<_>>();
407 drop(state);
408 entries.extend(self.matrix_projection_entries().await);
409 entries
410 }
411}
412
413mod controller_background_ops;
416mod controller_child_loop;
417mod controller_helpers;
418mod controller_spawn_run;
419mod controller_verify;
420
421#[allow(
422 unused_imports,
423 reason = "Intentional compatibility, platform, or test-only suppression."
424)]
425pub(super) use controller_helpers::*;
426
427#[cfg(test)]
430mod tests;