Skip to main content

vtcode_core/core/agent/
runner.rs

1//! Agent runner for executing individual agent instances
2
3use crate::config::VTCodeConfig;
4use crate::config::constants::tools;
5use crate::config::models::{ModelId, Provider};
6use crate::config::types::{ReasoningEffortLevel, VerbosityLevel};
7use crate::core::agent::events::EventSink;
8use crate::core::agent::features::FeatureSet;
9use crate::core::agent::session_config::ResolvedSessionConfig;
10use crate::core::agent::steering::SteeringMessage;
11use crate::core::threads::{ThreadBootstrap, ThreadRuntimeHandle, build_thread_archive_metadata};
12use std::path::Path;
13
14/// Settings for the agent runner
15#[derive(Clone, Default)]
16pub struct RunnerSettings {
17    /// Reasoning effort level for the agent
18    pub reasoning_effort: Option<ReasoningEffortLevel>,
19    /// Verbosity level for output text
20    pub verbosity: Option<VerbosityLevel>,
21}
22
23use crate::core::agent::types::AgentType;
24use crate::core::loop_detector::LoopDetector;
25use crate::exec::events::ThreadEvent;
26use crate::llm::AnyClient;
27use crate::llm::client::ProviderClientAdapter;
28use crate::llm::factory::{ProviderConfig, create_provider_with_config, infer_provider_from_model};
29use crate::llm::provider as uni_provider;
30use crate::primary_agent::ActivePrimaryAgent;
31use crate::prompts::PromptContext;
32use crate::tools::ToolRegistry;
33
34use anyhow::{Context, Result, anyhow};
35use parking_lot::{Mutex, RwLock};
36use std::path::PathBuf;
37use std::str::FromStr;
38use std::sync::Arc;
39use tracing::{info, warn};
40use vtcode_config::auth::OpenAIChatGptAuthHandle;
41
42mod config_helpers;
43mod constants;
44mod continuation;
45mod contract_render;
46mod escalation;
47mod evaluator_types;
48mod execute;
49mod execute_checks;
50mod execute_helpers;
51mod helpers;
52mod orchestration;
53mod output;
54mod planner_helpers;
55mod planner_types;
56pub mod prompt_alignment;
57mod replan_helpers;
58mod retry;
59mod summarize;
60mod summary;
61mod task_setup;
62mod telemetry;
63mod tool_access;
64mod tool_args;
65mod tool_dispatch_common;
66mod tool_exec;
67mod tool_execution_guard;
68mod tool_rejection;
69mod tool_types;
70mod types;
71mod validation;
72mod workspace_detection;
73
74/// Attach an eagerly initialized MCP client to the runner's tool registry.
75///
76/// Exec and other non-interactive runner sessions do not drive the async MCP
77/// manager, so MCP tools would otherwise remain invisible for the whole run.
78async fn attach_mcp_client(
79    tool_registry: &ToolRegistry,
80    session_config: &ResolvedSessionConfig,
81    workspace: &Path,
82) -> Result<()> {
83    let config = session_config.effective();
84    let sandbox_context = if config.sandbox.enabled
85        && !matches!(config.sandbox.default_policy, vtcode_config::SandboxPolicy::DangerFullAccess)
86    {
87        let policy = crate::tools::registry::sandbox_policy_from_runtime_config(&config.sandbox, workspace)?;
88        Some(crate::mcp::McpSandboxContext::new(policy, workspace))
89    } else {
90        None
91    };
92
93    let mut client = crate::mcp::McpClient::with_sandbox_context(config.mcp.clone(), sandbox_context);
94    let startup_timeout = std::time::Duration::from_secs(config.mcp.startup_timeout_seconds.unwrap_or(30));
95    let result = tokio::time::timeout(startup_timeout, client.initialize()).await;
96    match result {
97        Ok(Ok(())) => {}
98        Ok(Err(err)) => return Err(err).context("MCP client initialization failed"),
99        Err(_) => anyhow::bail!("MCP client initialization timed out after {} seconds", startup_timeout.as_secs()),
100    }
101
102    let client = Arc::new(client);
103    tool_registry.set_mcp_client(client).await;
104    if let Err(err) = tool_registry.refresh_mcp_tools().await {
105        warn!("Failed to refresh MCP tools after attach: {err:#}");
106    }
107    info!("MCP client attached to runner tool registry");
108    Ok(())
109}
110
111#[cfg(test)]
112mod tests;
113
114type ToolArgTransform = Arc<dyn Fn(&str, serde_json::Value) -> serde_json::Value + Send + Sync>;
115
116/// Individual agent runner for executing specialized agent tasks
117pub struct AgentRunner {
118    /// Agent type and configuration
119    agent_type: AgentType,
120    /// LLM client for this agent
121    client: AnyClient,
122    /// Unified provider client (OpenAI/Anthropic/Gemini) for tool-calling
123    provider_client: Box<dyn uni_provider::LLMProvider>,
124    /// Tool registry with restricted access
125    tool_registry: ToolRegistry,
126    /// System prompt content
127    system_prompt: String,
128    /// Token-budget report for `system_prompt`, computed at session
129    /// construction time (or recomputed by `set_system_prompt`). Reused by
130    /// `compose_task_system_prompt` for non-simple tasks so the budget
131    /// warning and preflight check stay accurate without recomposing the
132    /// prompt on every turn.
133    system_prompt_report: crate::prompts::system::SystemPromptReport,
134    /// Session information
135    session_id: String,
136    /// Initial archived history used to seed the first task on this runner.
137    bootstrap_messages: Vec<crate::llm::provider::Message>,
138    /// Workspace path
139    _workspace: PathBuf,
140    /// Frozen session-scoped configuration snapshot
141    session_config: Arc<ResolvedSessionConfig>,
142    /// Tool catalogue configuration used for this runner's model surface.
143    session_tools_config: crate::tools::handlers::SessionToolsConfig,
144    /// Model identifier
145    model: String,
146    /// API key (for provider client construction in future flows)
147    _api_key: String,
148    /// Reasoning effort level for models that support it
149    reasoning_effort: Option<ReasoningEffortLevel>,
150    /// Verbosity level for output text
151    verbosity: Option<VerbosityLevel>,
152    /// Suppress stdout output when emitting structured events
153    quiet: bool,
154    /// Optional sink for streaming structured events
155    event_sink: Option<EventSink>,
156    /// Shared thread runtime state for history/event ownership
157    thread_handle: ThreadRuntimeHandle,
158    /// Maximum number of autonomous turns before halting
159    max_turns: usize,
160    /// Loop detector to prevent infinite exploration
161    loop_detector: Mutex<LoopDetector>,
162    /// Cached shell policy patterns to avoid recompilation
163
164    /// Receiver for steering messages (e.g., stop, pause)
165    steering_receiver: Mutex<Option<tokio::sync::mpsc::UnboundedReceiver<SteeringMessage>>>,
166    /// Optional restricted tool definitions used instead of the default registry projection.
167    tool_definitions_override: RwLock<Option<Vec<uni_provider::ToolDefinition>>>,
168    /// Optional argument transformer applied before tool validation/execution.
169    tool_arg_transform: Option<ToolArgTransform>,
170    local_tools_only: bool,
171    /// Active primary-agent policy to intersect with runner tools.
172    active_primary_agent: Option<ActivePrimaryAgent>,
173}
174
175impl AgentRunner {
176    /// Get the selected model for the current turn.
177    fn get_selected_model(&self) -> String {
178        self.model.clone()
179    }
180
181    fn runner_println(&self, args: std::fmt::Arguments) {
182        if !self.quiet {
183            println!("{args}");
184        }
185    }
186
187    /// Create a new agent runner.
188    pub async fn new(
189        agent_type: AgentType,
190        model: ModelId,
191        api_key: String,
192        workspace: PathBuf,
193        session_id: String,
194        settings: RunnerSettings,
195        steering_receiver: Option<tokio::sync::mpsc::UnboundedReceiver<SteeringMessage>>,
196    ) -> Result<Self> {
197        Box::pin(Self::new_internal(
198            agent_type,
199            model,
200            api_key,
201            workspace,
202            session_id,
203            settings,
204            steering_receiver,
205            ThreadBootstrap::new(None),
206            None,
207            None,
208        ))
209        .await
210    }
211
212    /// Create an agent runner with prebuilt thread bootstrap, config, and auth.
213    #[allow(
214        clippy::too_many_arguments,
215        reason = "Intentional compatibility, platform, or test-only suppression."
216    )]
217    pub async fn new_with_bootstrap(
218        agent_type: AgentType,
219        model: ModelId,
220        api_key: String,
221        workspace: PathBuf,
222        session_id: String,
223        settings: RunnerSettings,
224        steering_receiver: Option<tokio::sync::mpsc::UnboundedReceiver<SteeringMessage>>,
225        bootstrap: ThreadBootstrap,
226        vt_cfg: Option<VTCodeConfig>,
227        openai_chatgpt_auth: Option<OpenAIChatGptAuthHandle>,
228    ) -> Result<Self> {
229        Box::pin(Self::new_internal(
230            agent_type,
231            model,
232            api_key,
233            workspace,
234            session_id,
235            settings,
236            steering_receiver,
237            bootstrap,
238            vt_cfg,
239            openai_chatgpt_auth,
240        ))
241        .await
242    }
243
244    #[expect(
245        clippy::too_many_arguments,
246        reason = "Intentional compatibility, platform, test, or API-shape suppression."
247    )]
248    async fn new_internal(
249        agent_type: AgentType,
250        model: ModelId,
251        api_key: String,
252        workspace: PathBuf,
253        session_id: String,
254        settings: RunnerSettings,
255        steering_receiver: Option<tokio::sync::mpsc::UnboundedReceiver<SteeringMessage>>,
256        bootstrap: ThreadBootstrap,
257        vt_cfg: Option<VTCodeConfig>,
258        openai_chatgpt_auth: Option<OpenAIChatGptAuthHandle>,
259    ) -> Result<Self> {
260        // Load configuration once to seed system prompt and runtime policies
261        let session_config = if let Some(vt_cfg) = vt_cfg {
262            ResolvedSessionConfig::from_config(vt_cfg)
263        } else {
264            match ResolvedSessionConfig::load_from_workspace(&workspace) {
265                Ok(session_config) => session_config,
266                Err(err) => {
267                    warn!("Failed to load vtcode configuration for system prompt composition: {err:#}");
268                    ResolvedSessionConfig::from_config(VTCodeConfig::default())
269                }
270            }
271        };
272        let session_config = Arc::new(session_config);
273        let provider_name = {
274            let configured = session_config.effective().agent.provider.trim();
275            if configured.is_empty() {
276                infer_provider_from_model(&model.as_str())
277                    .map(|provider| provider.to_string())
278                    .ok_or_else(|| anyhow!("Failed to determine provider for model {model}"))?
279            } else {
280                configured.to_lowercase()
281            }
282        };
283        let provider_config = ProviderConfig {
284            api_key: Some(api_key.clone()),
285            openai_chatgpt_auth: openai_chatgpt_auth.clone(),
286            copilot_auth: Some(session_config.effective().auth.copilot.clone()),
287            base_url: None,
288            model: Some(model.to_string()),
289            prompt_cache: Some(session_config.effective().prompt_cache.clone()),
290            timeouts: None,
291            openai: Some(session_config.effective().provider.openai.clone()),
292            anthropic: Some(session_config.effective().provider.anthropic.clone()),
293            model_behavior: Some(session_config.effective().model.clone()),
294            workspace_root: Some(workspace.clone()),
295        };
296
297        let client: AnyClient = Box::new(ProviderClientAdapter::new(
298            create_provider_with_config(&provider_name, provider_config.clone())
299                .with_context(|| "Failed to create client provider")?,
300            model.to_string(),
301        ));
302        let provider_client = create_provider_with_config(&provider_name, provider_config)
303            .with_context(|| "Failed to create provider client")?;
304        if std::env::var_os("VTCODE_DEBUG_PROVIDER").is_some() {
305            eprintln!(
306                "vtcode-debug: runner provider={} client_provider={} model={}",
307                provider_name,
308                provider_client.name(),
309                model
310            );
311        }
312        let max_repeated_tool_calls = session_config.effective().tools.max_repeated_tool_calls.max(1);
313        let deferred_tool_policy = crate::tools::handlers::deferred_tool_policy_for_runtime(
314            crate::llm::factory::infer_provider(Some(&session_config.effective().agent.provider), &model.as_str()),
315            provider_client.supports_responses_compaction(&model.as_str()),
316            Some(session_config.effective()),
317        );
318        let anthropic_native_memory_enabled = crate::tools::handlers::anthropic_native_memory_enabled_for_runtime(
319            Provider::from_str(provider_client.name()).ok(),
320            &model.as_str(),
321            Some(session_config.effective()),
322        );
323        let tool_registry = ToolRegistry::new(workspace.clone()).await;
324        tool_registry.set_harness_session(session_id.clone());
325        tool_registry.set_matrix_coordinator(session_config.effective().default_primary_agent == "coordinator");
326        tool_registry.set_agent_type(agent_type.to_string());
327        tool_registry.initialize_async().await?;
328        if let Err(err) = tool_registry
329            .apply_session_runtime_config(
330                &session_config.effective().commands,
331                &session_config.effective().permissions,
332                &session_config.effective().sandbox,
333                &session_config.effective().timeouts,
334                &session_config.effective().tools,
335            )
336            .await
337        {
338            warn!("Failed to apply tool policies from config: {}", err);
339        }
340        if session_config.effective().mcp.enabled {
341            if let Err(err) = crate::mcp::validate_mcp_config(&session_config.effective().mcp) {
342                warn!("MCP configuration validation error: {err}");
343            }
344            if let Err(err) = attach_mcp_client(&tool_registry, &session_config, &workspace).await {
345                warn!("Failed to attach MCP client: {err:#}");
346            }
347        }
348        if session_config.effective().context.dynamic.enabled
349            && let Err(err) =
350                crate::context::initialize_dynamic_context(&workspace, &session_config.effective().context.dynamic)
351                    .await
352        {
353            warn!("Failed to initialize dynamic context directories: {}", err);
354        }
355        let session_tools_config = crate::tools::handlers::SessionToolsConfig {
356            surface: crate::tools::handlers::SessionSurface::AgentRunner,
357            capability_level: crate::config::types::CapabilityLevel::CodeSearch,
358            documentation_mode: session_config.effective().agent.tool_documentation_mode,
359            planning_active: tool_registry.is_planning_active(),
360            request_user_input_enabled: false,
361            model_capabilities: crate::tools::handlers::ToolModelCapabilities::for_model_name(&model.as_str()),
362            deferred_tool_policy,
363            anthropic_native_memory_enabled,
364            tool_profile: session_config.effective().tools.profile,
365        };
366        let available_tools = tool_registry
367            .model_tools(session_tools_config.clone())
368            .await
369            .into_iter()
370            .map(|tool| tool.function_name().to_string())
371            .collect::<Vec<_>>();
372        let mut prompt_context = PromptContext::from_workspace_tools(&workspace, available_tools);
373        prompt_context.set_current_directory(workspace.clone());
374        prompt_context.load_available_skills_async().await;
375        let (system_prompt, system_prompt_report) = helpers::compose_system_prompt_with_appendix(
376            workspace.as_path(),
377            session_config.effective(),
378            &prompt_context,
379        )
380        .await?;
381        let loop_detector = LoopDetector::with_max_repeated_calls(max_repeated_tool_calls);
382        let bootstrap_messages = bootstrap.messages.clone();
383        let mut bootstrap = bootstrap;
384        if bootstrap.metadata.is_none() {
385            bootstrap.metadata = Some(build_thread_archive_metadata(
386                workspace.as_path(),
387                &model.as_str(),
388                &session_config.effective().agent.provider,
389                &session_config.effective().agent.theme,
390                settings
391                    .reasoning_effort
392                    .unwrap_or(session_config.effective().agent.reasoning_effort)
393                    .as_str(),
394            ));
395        }
396        let thread_handle =
397            crate::core::threads::ThreadManager::new().start_thread_with_identifier(session_id.clone(), bootstrap);
398        let max_turns = session_config.effective().automation.full_auto.max_turns.max(1);
399        if session_config.effective().default_primary_agent == "coordinator" {
400            let controller =
401                Box::pin(crate::subagents::SubagentController::new(crate::subagents::SubagentControllerConfig {
402                    workspace_root: workspace.clone(),
403                    parent_session_id: session_id.clone(),
404                    parent_model: model.to_string(),
405                    parent_provider: provider_name,
406                    parent_reasoning_effort: settings
407                        .reasoning_effort
408                        .unwrap_or(session_config.effective().agent.reasoning_effort),
409                    api_key: api_key.clone(),
410                    vt_cfg: session_config.effective().clone(),
411                    openai_chatgpt_auth,
412                    depth: 0,
413                    workspace_gated: session_config
414                        .effective()
415                        .workspace_lifecycle_hooks
416                        .as_ref()
417                        .is_some_and(|hooks| !hooks.is_empty()),
418                    exec_sessions: tool_registry.exec_session_manager(),
419                    pty_manager: tool_registry.pty_manager().clone(),
420                    managed_background_runtime: false,
421                }))
422                .await?;
423            tool_registry.set_subagent_controller(Arc::new(controller));
424        }
425
426        Ok(Self {
427            agent_type,
428            client,
429            provider_client,
430            tool_registry,
431            system_prompt,
432            system_prompt_report,
433            session_id,
434            bootstrap_messages,
435            _workspace: workspace,
436            session_config,
437            session_tools_config,
438            model: model.to_string(),
439            _api_key: api_key,
440            reasoning_effort: settings.reasoning_effort,
441            verbosity: settings.verbosity,
442            quiet: false,
443            event_sink: None,
444            thread_handle,
445            max_turns,
446            loop_detector: Mutex::new(loop_detector),
447            steering_receiver: Mutex::new(steering_receiver),
448            tool_definitions_override: RwLock::new(None),
449            tool_arg_transform: None,
450            local_tools_only: false,
451            active_primary_agent: None,
452        })
453    }
454
455    /// Enable or disable console output for this runner.
456    pub fn set_quiet(&mut self, quiet: bool) {
457        self.quiet = quiet;
458    }
459
460    /// Configure the loop detector for subagent mode with tighter read-only
461    /// budgets and earlier navigation streak intervention.
462    pub fn set_subagent_mode(&self, is_subagent: bool) {
463        self.loop_detector.lock().set_subagent_mode(is_subagent);
464    }
465
466    /// Attach a subagent controller to this runner's tool registry so the
467    /// subagent-lifecycle tools (`agent` and its aliases) become available to
468    /// the model. Used to enable nested delegation: the attached controller
469    /// carries the child-scoped depth so the depth check in `spawn_with_spec`
470    /// still governs grandchild spawns.
471    pub fn set_subagent_controller(&self, controller: Arc<crate::subagents::SubagentController>) {
472        self.tool_registry.set_subagent_controller(controller);
473    }
474
475    /// Snapshot the runner-owned conversation messages for archive persistence.
476    pub fn session_messages(&self) -> Vec<crate::llm::provider::Message> {
477        self.thread_handle.messages()
478    }
479
480    /// Clone the underlying thread handle so callers can capture snapshots.
481    pub fn thread_handle(&self) -> ThreadRuntimeHandle {
482        self.thread_handle.clone()
483    }
484
485    /// Workspace root this runner operates within.
486    pub fn workspace(&self) -> &Path {
487        &self._workspace
488    }
489
490    /// Enable read-only planning workflow for the underlying tool registry.
491    pub fn enable_planning(&self) {
492        self.tool_registry.enable_planning();
493    }
494
495    /// Disable read-only planning workflow for the underlying tool registry.
496    pub fn disable_planning(&self) {
497        self.tool_registry.disable_planning();
498    }
499
500    /// Attach a callback that will be invoked for each structured event as it is recorded.
501    pub fn set_event_handler<F>(&mut self, handler: F)
502    where
503        F: FnMut(&ThreadEvent) + Send + 'static,
504    {
505        self.event_sink = Some(Arc::new(Mutex::new(Box::new(handler))));
506    }
507
508    /// Remove any previously registered structured event callback.
509    pub fn clear_event_handler(&mut self) {
510        self.event_sink = None;
511    }
512
513    /// Override the composed system prompt for downstream embedders.
514    ///
515    /// Recomputes `system_prompt_report` against the overridden text so the
516    /// budget warning and preflight check stay accurate even when the prompt
517    /// bypasses the normal section-based composition pipeline.
518    pub fn set_system_prompt(&mut self, system_prompt: impl Into<String>) {
519        self.system_prompt = system_prompt.into();
520        self.system_prompt_report = crate::prompts::system::SystemPromptReport::measure(
521            &self.system_prompt,
522            self.config().agent.max_system_prompt_tokens,
523        );
524    }
525
526    /// Clone the underlying tool registry so embedders can register custom tools.
527    pub fn tool_registry(&self) -> ToolRegistry {
528        self.tool_registry.clone()
529    }
530
531    /// Override the tool definitions used for LLM requests instead of the default registry projection.
532    pub fn set_tool_definitions_override(&mut self, definitions: Vec<uni_provider::ToolDefinition>) {
533        *self.tool_definitions_override.write() = Some(definitions);
534    }
535
536    /// Clear any previously set tool definition override, restoring the default registry projection.
537    pub fn clear_tool_definitions_override(&mut self) {
538        *self.tool_definitions_override.write() = None;
539    }
540
541    /// Set an argument transformer applied to tool calls before validation and execution.
542    pub fn set_tool_arg_transform(&mut self, transform: ToolArgTransform) {
543        self.tool_arg_transform = Some(transform);
544    }
545
546    /// Restrict skill child execution using live trusted registration metadata.
547    pub fn restrict_to_local_tools(&mut self) {
548        self.local_tools_only = true;
549    }
550
551    /// Clear any previously set tool argument transformer.
552    pub fn clear_tool_arg_transform(&mut self) {
553        self.tool_arg_transform = None;
554    }
555
556    /// Apply active primary-agent tool and permission policy to this runner.
557    pub fn set_active_primary_agent(&mut self, active_primary_agent: ActivePrimaryAgent) {
558        self.tool_registry
559            .set_matrix_coordinator(active_primary_agent.identity.name == "coordinator");
560        self.active_primary_agent = Some(active_primary_agent);
561    }
562
563    /// Enable full-auto execution with the provided allow-list.
564    pub async fn enable_full_auto(&mut self, allowed_tools: &[String]) {
565        let mut permission_catalog_config = self.session_tools_config.clone();
566        permission_catalog_config.planning_active = allowed_tools
567            .iter()
568            .any(|tool_name| crate::tools::names::canonical_tool_name(tool_name) == tools::CODE_SEARCH);
569        self.tool_registry
570            .enable_full_auto_permission_for_session(allowed_tools, permission_catalog_config)
571            .await;
572    }
573
574    /// Restrict an allow-list to tools suitable for strict review-only runs.
575    pub async fn review_tool_allowlist(&self, allowed_tools: &[String]) -> Vec<String> {
576        let review_candidates = if allowed_tools.iter().any(|tool| tool.trim() == tools::WILDCARD_ALL) {
577            let mut candidates = self.tool_registry.available_tools().await;
578            if !candidates.iter().any(|tool| tool == tools::CODE_SEARCH) {
579                candidates.push(tools::CODE_SEARCH.to_string());
580            }
581            candidates
582        } else {
583            allowed_tools.to_vec()
584        };
585
586        review_candidates
587            .iter()
588            .filter(|tool_name| {
589                let canonical = crate::tools::names::canonical_tool_name(tool_name);
590
591                !matches!(canonical, tools::REQUEST_USER_INPUT | tools::TASK_TRACKER | tools::START_PLANNING)
592                    && !self.tool_registry.is_mutating_tool(tool_name)
593            })
594            .cloned()
595            .collect()
596    }
597
598    fn features(&self) -> FeatureSet {
599        self.session_config.features().clone()
600    }
601}