1use anyhow::{Context, Result, anyhow};
4use serde_json::{Value, json};
5use std::borrow::Cow;
6use std::future::Future;
7use std::time::{Duration, Instant};
8use tracing::{trace, warn};
9use vtcode_commons::ErrorCategory;
10
11use crate::config::constants::tools;
12use crate::core::agent::harness_kernel::PreparedToolCall;
13use crate::core::memory_pool::SizeRecommendation;
14use crate::tool_policy::ToolExecutionDecision;
15use crate::tools::error_messages::agent_execution;
16use crate::tools::request_response::{ToolCallRequest, ToolCallResponse};
17use crate::tools::tool_intent;
18use crate::tools::unified_error::UnifiedErrorKind;
19use crate::tools::unified_error::UnifiedToolError;
20use crate::ui::search::fuzzy_match;
21
22use super::assembly::public_tool_name_candidates;
23use super::execution_kernel;
24use super::reentrancy::ToolReentrancyGuard;
25use super::{
26 ExecSettlementMode, ExecutionPolicySnapshot, ToolErrorType, ToolExecutionError, ToolExecutionOutcome,
27 ToolExecutionRecord, ToolExecutionRequest, ToolRegistry,
28};
29use vtcode_config::constants::execution::{LOOP_THROTTLE_MAX_MS, LOOP_THROTTLE_REGISTRY_BASE_MS};
30
31const LOOP_HARD_BLOCK_REPEAT_COUNT: usize = 5;
35
36impl ToolRegistry {
37 fn annotate_timeout_error_payload(
38 payload: &mut Value,
39 timeout_category: &str,
40 timeout_ms: u64,
41 circuit_breaker: bool,
42 ) {
43 if let Some(obj) = payload.get_mut("error").and_then(|value| value.as_object_mut()) {
44 obj.insert("timeout_category".into(), Value::String(timeout_category.to_string()));
45 obj.insert("timeout_ms".into(), Value::from(timeout_ms));
46 obj.insert("circuit_breaker".into(), Value::Bool(circuit_breaker));
47 }
48 }
49
50 pub fn safety_gateway(&self) -> std::sync::Arc<crate::tools::safety_gateway::SafetyGateway> {
51 std::sync::Arc::clone(&self.safety_gateway)
52 }
53
54 pub fn execute_public_tool_request(
57 &self,
58 request: ToolExecutionRequest,
59 ) -> impl Future<Output = ToolExecutionOutcome> + '_ {
60 self.execute_tool_request_internal(request)
61 }
62
63 pub async fn execute_prepared_public_tool_request(
64 &self,
65 prepared: &PreparedToolCall,
66 policy: ExecutionPolicySnapshot,
67 ) -> ToolExecutionOutcome {
68 let request = ToolExecutionRequest::new(prepared.canonical_name.clone(), prepared.effective_args.clone())
69 .with_policy(
70 policy
71 .with_prevalidated(prepared.already_preflighted)
72 .with_safety_prevalidated(false),
73 );
74 self.execute_tool_request_internal(request).await
75 }
76
77 async fn should_skip_loop_detection_for_exec_continuation(&self, tool_name: &str, args: &Value) -> bool {
78 if tool_name == tools::WRITE_STDIN {
79 return matches!(
80 crate::tools::command_args::write_stdin_dispatch(args),
81 Ok(crate::tools::command_args::WriteStdinDispatch::Poll
82 | crate::tools::command_args::WriteStdinDispatch::Wait,)
83 );
84 }
85
86 if tool_name != tools::UNIFIED_EXEC {
87 return false;
88 }
89
90 if !tool_intent::command_session_action_in(args, &["poll", "continue"]) {
91 return false;
92 }
93 if tool_intent::command_session_action_is(args, "continue")
94 && crate::tools::command_args::interactive_input_text(args).is_some()
95 {
96 return false;
97 }
98
99 let Some(session_id) = crate::tools::command_args::session_id_text(args) else {
100 return false;
101 };
102
103 matches!(self.exec_session_completed(session_id).await, Ok(None))
104 }
105
106 async fn public_tool_catalog_for_error(&self, requested_name: &str) -> (Vec<String>, Vec<String>) {
107 let mut tool_names = self.available_tools().await;
108 tool_names.sort_unstable();
109 tool_names.dedup();
110
111 let requested_candidates = public_tool_name_candidates(requested_name);
112 let mut similar_tools = Vec::new();
113
114 if let Ok(resolved) = self.resolve_public_tool_name_sync(requested_name)
115 && tool_names.iter().any(|tool| tool == &resolved)
116 {
117 similar_tools.push(resolved);
118 }
119
120 for tool in &tool_names {
121 if similar_tools.len() >= 3 {
122 break;
123 }
124
125 if similar_tools.iter().any(|candidate| candidate == tool) {
126 continue;
127 }
128
129 if requested_candidates.iter().any(|candidate| fuzzy_match(candidate, tool)) {
130 similar_tools.push(tool.clone());
131 }
132 }
133
134 (tool_names, similar_tools)
135 }
136
137 pub fn preflight_validate_call(&self, name: &str, args: &Value) -> Result<super::ToolPreflightOutcome> {
138 execution_kernel::preflight_validate_call(self, name, args)
139 }
140
141 pub fn preflight_validate_harness_call(&self, name: &str, args: &Value) -> Result<super::ToolPreflightOutcome> {
145 execution_kernel::preflight_validate_call_with_mode(self, name, args, execution_kernel::DispatchMode::Harness)
146 }
147
148 pub fn admit_public_tool_call(&self, name: &str, args: &Value) -> Result<PreparedToolCall> {
149 let preflight = self.preflight_validate_harness_call(name, args)?;
150 Ok(PreparedToolCall::new(
151 preflight.normalized_tool_name,
152 preflight.readonly_classification,
153 preflight.parallel_safe_after_preflight,
154 preflight.effective_args,
155 ))
156 }
157
158 pub async fn execute_tool(&self, name: &str, args: Value) -> Result<Value> {
159 self.execute_tool_ref(name, &args).await
160 }
161
162 pub async fn execute_public_tool_ref(&self, name: &str, args: &Value) -> Result<Value> {
164 self.execute_public_tool_ref_internal(name, args, false).await
165 }
166
167 pub async fn execute_tool_ref(&self, name: &str, args: &Value) -> Result<Value> {
170 self.execute_tool_ref_internal(name, args, false, ExecSettlementMode::Manual)
171 .await
172 }
173
174 pub async fn execute_tool_ref_prevalidated(&self, name: &str, args: &Value) -> Result<Value> {
179 self.execute_tool_ref_internal(name, args, true, ExecSettlementMode::Manual)
180 .await
181 }
182
183 pub async fn execute_public_tool_ref_prevalidated(&self, name: &str, args: &Value) -> Result<Value> {
185 self.execute_public_tool_ref_prevalidated_with_mode(name, args, ExecSettlementMode::Manual)
186 .await
187 }
188
189 #[doc(hidden)]
190 pub async fn execute_public_tool_ref_prevalidated_with_mode(
191 &self,
192 name: &str,
193 args: &Value,
194 exec_settlement_mode: ExecSettlementMode,
195 ) -> Result<Value> {
196 self.execute_public_tool_ref_internal_with_mode(name, args, true, exec_settlement_mode)
197 .await
198 }
199
200 pub async fn execute_prepared_public_tool_ref_with_mode(
201 &self,
202 prepared: &PreparedToolCall,
203 exec_settlement_mode: ExecSettlementMode,
204 ) -> Result<Value> {
205 self.execute_public_tool_ref_dispatch(
209 prepared.canonical_name.as_str(),
210 &prepared.effective_args,
211 prepared.already_preflighted,
212 execution_kernel::DispatchMode::Harness,
213 exec_settlement_mode,
214 )
215 .await
216 }
217
218 async fn execute_public_tool_ref_internal(&self, name: &str, args: &Value, prevalidated: bool) -> Result<Value> {
219 self.execute_public_tool_ref_dispatch(
220 name,
221 args,
222 prevalidated,
223 execution_kernel::DispatchMode::ModelPublic,
224 ExecSettlementMode::Manual,
225 )
226 .await
227 }
228
229 async fn execute_public_tool_ref_internal_with_mode(
230 &self,
231 name: &str,
232 args: &Value,
233 prevalidated: bool,
234 exec_settlement_mode: ExecSettlementMode,
235 ) -> Result<Value> {
236 self.execute_public_tool_ref_dispatch(
237 name,
238 args,
239 prevalidated,
240 execution_kernel::DispatchMode::ModelPublic,
241 exec_settlement_mode,
242 )
243 .await
244 }
245
246 pub(super) async fn execute_public_tool_ref_dispatch(
256 &self,
257 name: &str,
258 args: &Value,
259 prevalidated: bool,
260 dispatch_mode: execution_kernel::DispatchMode,
261 exec_settlement_mode: ExecSettlementMode,
262 ) -> Result<Value> {
263 let routed_name = execution_kernel::resolve_dispatch_target(self, name, dispatch_mode)
264 .map_err(|err| anyhow!(err.to_string()))?;
265 let effective_args = execution_kernel::remap_public_file_operation_alias_args(name, routed_name.as_str(), args)
266 .or_else(|| execution_kernel::remap_consolidated_action_alias_args(name, routed_name.as_str(), args));
267 self.execute_tool_ref_internal(
268 routed_name.as_str(),
269 effective_args.as_ref().unwrap_or(args),
270 prevalidated,
271 exec_settlement_mode,
272 )
273 .await
274 }
275
276 async fn execute_tool_ref_internal(
277 &self,
278 name: &str,
279 args: &Value,
280 prevalidated: bool,
281 exec_settlement_mode: ExecSettlementMode,
282 ) -> Result<Value> {
283 let resolved = self.resolve_tool_name_with_display(name);
284 self.enforce_matrix_role(&resolved.canonical)?;
285 if let Err(error) = crate::core::agent::snapshots::declare_prompt_edit(
286 self.harness_context_snapshot().session_id,
287 name.to_owned(),
288 args.clone(),
289 )
290 .await
291 {
292 tracing::warn!(
293 tool = %name,
294 error = %error,
295 "Checkpoint pre-image capture failed; rewind may not restore this edit"
296 );
297 }
298 let _pool_guard = if self.optimization_config.memory_pool.enabled {
300 Some(self.memory_pool.get_string())
301 } else {
302 None
303 };
304
305 if self.optimization_config.memory_pool.enabled {
307 let recommendation = self.memory_pool.auto_tune(&self.optimization_config.memory_pool);
308
309 if !matches!(
311 (
312 recommendation.string_size_recommendation,
313 recommendation.value_size_recommendation,
314 recommendation.vec_size_recommendation
315 ),
316 (SizeRecommendation::Maintain, SizeRecommendation::Maintain, SizeRecommendation::Maintain)
317 ) {
318 tracing::debug!(
319 "Memory pool tuning recommendation: string={:?}, value={:?}, vec={:?}, allocations_avoided={}",
320 recommendation.string_size_recommendation,
321 recommendation.value_size_recommendation,
322 recommendation.vec_size_recommendation,
323 recommendation.total_allocations_avoided
324 );
325 }
326 }
327
328 let resolved_name = self.resolve_tool_name_with_display(name);
329 let tool_name = resolved_name.canonical;
330 let tool_name_owned = tool_name.clone();
331 let display_name = resolved_name.display;
332
333 let cached_tool = if self.optimization_config.tool_registry.use_optimized_registry {
339 let cache = self.hot_tool_cache.read();
340 cache.peek(&tool_name).cloned()
341 } else {
342 None
343 };
344
345 if let Some(tool_arc) = cached_tool.as_ref()
347 && self.optimization_config.tool_registry.use_optimized_registry
348 && tool_name != name
349 {
350 self.hot_tool_cache.write().put(tool_name.clone(), tool_arc.clone());
352 }
353
354 let execution_args = self.prepare_execution_args(&tool_name, args)?;
355 let is_verification_command = execution_args.is_verification_command;
356 let max_output_tokens = execution_args.max_output_tokens;
357 let args = execution_args.handler_args.as_ref();
358 let requested_name = name.to_string();
359
360 let args_for_recording = args.clone();
362 let context_snapshot = self.harness_context_snapshot();
364 let record_failure = |tool_name: String,
365 is_mcp_tool: bool,
366 mcp_provider: Option<String>,
367 args: Value,
368 error_msg: String,
369 timeout_category: Option<String>,
370 base_timeout_ms: Option<u64>,
371 adaptive_timeout_ms: Option<u64>,
372 effective_timeout_ms: Option<u64>,
373 circuit_breaker: bool| {
374 self.execution_history.add_record(ToolExecutionRecord::failure(
375 tool_name,
376 requested_name.clone(),
377 is_mcp_tool,
378 mcp_provider,
379 args,
380 error_msg,
381 context_snapshot.clone(),
382 timeout_category,
383 base_timeout_ms,
384 adaptive_timeout_ms,
385 effective_timeout_ms,
386 circuit_breaker,
387 ));
388 };
389
390 let allow_parallel_sibling = prevalidated && tool_intent::is_parallel_safe_call(&tool_name, args);
391 let _reentrancy_guard = match ToolReentrancyGuard::enter(&tool_name, allow_parallel_sibling) {
392 Ok(guard) => guard,
393 Err(violation) => {
394 let reentry_count = violation.tool_reentry_count + 1;
395 let error_message = format!(
396 "Reentrancy guard: tool '{}' is already running in this call stack, so this recursive call was blocked. \
397 Repeating the same call is blocked the same way; change the control flow or use a different tool.\n\
398 Current stack depth: {}. Re-entry count for this tool in the current task: {}.\n\
399 Stack trace: {}",
400 display_name, violation.stack_depth, reentry_count, violation.stack_trace
401 );
402 let error = ToolExecutionError::new(
403 tool_name_owned.clone(),
404 ToolErrorType::PolicyViolation,
405 error_message.clone(),
406 );
407 let mut payload = error.to_json_value();
408 if let Some(obj) = payload.as_object_mut() {
409 obj.insert("reentrant_call_blocked".into(), json!(true));
410 obj.insert("stack_depth".into(), json!(violation.stack_depth));
411 obj.insert("reentry_count".into(), json!(reentry_count));
412 obj.insert("tool".into(), json!(display_name));
413 obj.insert("stack_trace".into(), json!(violation.stack_trace));
414 }
415 record_failure(
416 tool_name_owned.clone(),
417 false,
418 None,
419 args_for_recording.clone(),
420 error_message.clone(),
421 None,
422 None,
423 None,
424 None,
425 false,
426 );
427 return Err(anyhow!(error_message).context("tool reentrancy blocked"));
428 }
429 };
430
431 let (intent, readonly_classification) = if prevalidated {
437 #[cfg(debug_assertions)]
438 {
439 if let Err(err) = execution_kernel::preflight_validate_resolved_call(self, &tool_name, args)
440 && !agent_execution::is_planning_active_denial(&err.to_string())
441 {
442 debug_assert!(false, "prevalidated execution received invalid call for '{tool_name}': {err}");
443 }
444 }
445 let intent = tool_intent::classify_tool_intent(&tool_name, args);
446 (intent, !intent.mutating)
447 } else {
448 match execution_kernel::preflight_validate_resolved_call(self, &tool_name, args) {
449 Ok(outcome) => (outcome.intent, outcome.readonly_classification),
450 Err(err) => {
451 let err_msg = err.to_string();
452 record_failure(
453 tool_name_owned.clone(),
454 false,
455 None,
456 args_for_recording.clone(),
457 err_msg,
458 None,
459 None,
460 None,
461 None,
462 false,
463 );
464 return Err(err);
465 }
466 }
467 };
468
469 if readonly_classification {
470 trace!(tool = %tool_name, "Validation classified tool as read-only");
471 }
472
473 if self.is_planning_active() && !self.is_planning_active_allowed_with_intent(&tool_name, args, &intent) {
477 let error_msg = agent_execution::planning_workflow_denial_message(&display_name);
478 record_failure(
479 tool_name_owned.clone(),
480 false,
481 None,
482 args_for_recording.clone(),
483 error_msg.clone(),
484 None,
485 None,
486 None,
487 None,
488 false,
489 );
490 return Err(anyhow!(error_msg).context(agent_execution::PLANNING_DENIED_CONTEXT));
491 }
492
493 let shared_circuit_breaker = self.shared_circuit_breaker();
494 if let Some(breaker) = shared_circuit_breaker.as_ref()
495 && !breaker.allow_request_for_tool(&tool_name)
496 {
497 let diagnostics = breaker.get_diagnostics(&tool_name);
498 let retry_after = diagnostics
499 .remaining_backoff
500 .map(|backoff| format!(" retry_after={}s.", backoff.as_secs()))
501 .unwrap_or_default();
502 let error_msg = format!(
503 "Tool '{display_name}' is temporarily disabled due to high failure rate (Circuit Breaker OPEN).{retry_after}"
504 );
505 self.execution_history.add_record(
506 ToolExecutionRecord::failure(
507 tool_name_owned.clone(),
508 requested_name.clone(),
509 false,
510 None,
511 args_for_recording.clone(),
512 error_msg.clone(),
513 context_snapshot.clone(),
514 None,
515 None,
516 None,
517 None,
518 true,
519 )
520 .with_circuit_breaker_state(format!("{:?}", diagnostics.status))
521 .with_retry_after(diagnostics.remaining_backoff),
522 );
523 return Err(anyhow!(error_msg).context("tool denied by circuit breaker"));
524 }
525
526 let timeout_category = self.timeout_category_for_args(&tool_name, args).await;
527
528 if let Some(backoff) = self.should_circuit_break(timeout_category) {
529 warn!(
530 tool = %tool_name,
531 category = %timeout_category.label(),
532 delay_ms = %backoff.as_millis(),
533 "Circuit breaker active for tool category; backing off before execution"
534 );
535 tokio::time::sleep(backoff).await;
536 }
537
538 let execution_span = tracing::debug_span!(
539 "tool_execution",
540 tool = %tool_name,
541 requested = %name,
542 session_id = %context_snapshot.session_id,
543 task_id = %context_snapshot.task_id.as_deref().unwrap_or("")
544 );
545 let _span_guard = execution_span.enter();
546
547 trace!(
548 tool = %tool_name,
549 session_id = %context_snapshot.session_id,
550 task_id = %context_snapshot.task_id.as_deref().unwrap_or(""),
551 "Executing tool with harness context"
552 );
553
554 if tool_name != name {
555 trace!(
556 requested = %name,
557 canonical = %tool_name,
558 "Resolved tool alias to canonical name"
559 );
560 }
561
562 let base_timeout_ms = self
563 .timeout_policy
564 .read()
565 .ceiling_for(timeout_category)
566 .map(|d| d.as_millis() as u64);
567 let adaptive_timeout_ms = self
568 .resiliency
569 .lock()
570 .adaptive_timeout_ceiling
571 .get(&timeout_category)
572 .filter(|d| d.as_millis() > 0)
573 .map(|d| d.as_millis() as u64);
574 let timeout_category_label = Some(timeout_category.label().to_string());
575
576 if let Some(rate_limit) = self.execution_history.rate_limit_per_minute() {
577 let calls_last_minute = self.execution_history.calls_in_window(Duration::from_secs(60));
578 if calls_last_minute >= rate_limit {
579 warn!(
580 tool = %tool_name_owned,
581 requested = %requested_name,
582 calls_last_minute,
583 rate_limit,
584 "Execution history rate-limit threshold exceeded (observability-only)"
585 );
586 }
587 }
588
589 let fresh_patch_read = self.consume_patch_recovery_read(&tool_name, args);
590 let matrix_verification = self
593 .matrix_worker
594 .read()
595 .as_ref()
596 .is_some_and(|worker| worker.assignment.phase == crate::exec::events::matrix::MatrixPhase::Verify);
597 let reusable_result = readonly_classification
598 && !matrix_verification
599 && !matches!(tool_name.as_str(), tools::RECORD_DECISION | tools::TASK_TRACKER | tools::MATRIX);
600 let skip_loop_detection = self.should_skip_loop_detection_for_exec_continuation(&tool_name, args).await;
601 if skip_loop_detection {
602 trace!(
603 tool = %tool_name,
604 "Skipping identical-call loop detection for stateful exec continuation"
605 );
606 }
607
608 if reusable_result && !is_verification_command && !skip_loop_detection && !fresh_patch_read {
623 let fast_reuse_max_age = Duration::from_secs(60);
624 let fast_reused = self
625 .execution_history
626 .find_recent_spooled_result(&tool_name, args, fast_reuse_max_age)
627 .or_else(|| {
628 self.execution_history
629 .find_recent_successful_result(&tool_name, args, fast_reuse_max_age)
630 });
631 if let Some(mut reused_value) = fast_reused {
632 if let Some(obj) = reused_value.as_object_mut() {
633 obj.insert("reused_recent_result".into(), json!(true));
634 obj.insert("tool".into(), json!(display_name));
635 let reused_spooled = obj.get("spool_path").and_then(|v| v.as_str()).is_some();
636 let note = if reused_spooled {
637 "Reusing a recent spooled output for this identical read-only call. Continue from the spool file instead of re-running the tool."
638 } else {
639 "Reusing a recent successful output for this identical read-only call."
640 };
641 obj.insert("reused_result_note".into(), json!(note));
642 }
643 self.execution_history.add_record(ToolExecutionRecord::success(
659 tool_name.clone(),
660 requested_name.clone(),
661 false,
662 None,
663 args_for_recording.clone(),
664 reused_value.clone(),
665 context_snapshot.clone(),
666 timeout_category_label.clone(),
667 base_timeout_ms,
668 adaptive_timeout_ms,
669 None,
670 false,
671 ));
672 trace!(
673 tool = %tool_name,
674 "Fast-reusing recent successful read-only result"
675 );
676 return Ok(reused_value);
677 }
678 }
679
680 let loop_limit = if skip_loop_detection {
682 0
683 } else {
684 self.execution_history.loop_limit_for(&tool_name, args)
685 };
686 let loop_result = if skip_loop_detection {
687 crate::tools::registry::execution_history::LoopDetectionResult {
688 detected: false,
689 repeat_count: 0,
690 tool_name: tool_name.clone(),
691 }
692 } else {
693 self.execution_history.detect_loop(&tool_name, args)
694 };
695 if loop_result.detected && loop_result.repeat_count > 1 {
696 let delay_ms = (LOOP_THROTTLE_REGISTRY_BASE_MS * loop_result.repeat_count as u64).min(LOOP_THROTTLE_MAX_MS);
697 if delay_ms > 0 {
698 tokio::time::sleep(Duration::from_millis(delay_ms)).await;
699 }
700 }
701 if loop_limit > 0 && loop_result.detected {
702 warn!(
703 tool = %tool_name,
704 repeats = loop_result.repeat_count,
705 "Loop detected: agent calling same tool with identical parameters {} times",
706 loop_result.repeat_count
707 );
708 if loop_result.repeat_count >= loop_limit {
709 let hard_block = loop_result.repeat_count >= LOOP_HARD_BLOCK_REPEAT_COUNT;
714
715 if reusable_result && !hard_block && !fresh_patch_read {
716 let reuse_max_age = Duration::from_secs(120);
717 let reused = self
718 .execution_history
719 .find_recent_spooled_result(&tool_name, args, reuse_max_age)
720 .or_else(|| {
721 self.execution_history
722 .find_recent_successful_result(&tool_name, args, reuse_max_age)
723 });
724 if let Some(mut reused_value) = reused {
725 if let Some(obj) = reused_value.as_object_mut() {
726 obj.insert("reused_recent_result".into(), json!(true));
727 obj.insert("loop_detected".into(), json!(true));
728 obj.insert("repeat_count".into(), json!(loop_result.repeat_count));
729 obj.insert("limit".into(), json!(loop_limit));
730 obj.insert("tool".into(), json!(display_name));
731 let reused_spooled = obj.get("spool_path").and_then(|v| v.as_str()).is_some();
732 let note = if reused_spooled {
733 "Loop detected: this identical read-only call has been repeated, so the earlier result was reused. The full output is in the spool file and the conversation history; further repeats return an error instead."
734 } else {
735 "Loop detected: this identical read-only call has been repeated with no new information, so the earlier result was reused. It is already in the conversation history; further repeats return an error instead."
736 };
737 obj.insert("loop_detected_note".into(), json!(note));
738 }
739 return Ok(reused_value);
740 }
741 }
742
743 let delay_ms =
744 (LOOP_THROTTLE_REGISTRY_BASE_MS * loop_result.repeat_count as u64).min(LOOP_THROTTLE_MAX_MS);
745 if delay_ms > 0 {
746 tokio::time::sleep(Duration::from_millis(delay_ms)).await;
747 }
748
749 let error = ToolExecutionError::new(
750 tool_name_owned.clone(),
751 ToolErrorType::PolicyViolation,
752 agent_execution::loop_detection_block_message(&display_name, loop_result.repeat_count as u64, None),
753 );
754 let mut payload = error.to_json_value();
755 if let Some(obj) = payload.as_object_mut() {
756 obj.insert("loop_detected".into(), json!(true));
757 obj.insert("repeat_count".into(), json!(loop_result.repeat_count));
758 obj.insert("limit".into(), json!(loop_limit));
759 obj.insert("tool".into(), json!(display_name));
760 obj.insert(
761 "next_action".into(),
762 json!("Identical calls to this tool are blocked. Use the data already in the conversation history; for different information, change the arguments or use another tool."),
763 );
764 }
765
766 record_failure(
767 tool_name_owned,
768 false,
769 None,
770 args_for_recording,
771 "Tool call blocked due to repeated identical invocations".to_string(),
772 timeout_category_label.clone(),
773 base_timeout_ms,
774 adaptive_timeout_ms,
775 None,
776 false,
777 );
778
779 return Ok(payload);
780 }
781 }
782
783 let full_auto_denied = {
784 let gateway = self.policy_gateway.clone();
785 let tool_name_ref = &tool_name;
786 async move { gateway.is_denied_in_full_auto(tool_name_ref).await }
787 };
788 let full_auto_denied = full_auto_denied.await;
789 if full_auto_denied {
790 let _error = ToolExecutionError::new(
791 tool_name_owned.clone(),
792 ToolErrorType::PolicyViolation,
793 format!("Tool '{display_name}' is not permitted while full-auto permission review is active"),
794 );
795
796 record_failure(
797 tool_name_owned.clone(),
798 false,
799 None,
800 args_for_recording.clone(),
801 "Tool execution denied by policy".to_string(),
802 timeout_category_label.clone(),
803 base_timeout_ms,
804 adaptive_timeout_ms,
805 None,
806 false,
807 );
808
809 return Err(anyhow!("Tool '{display_name}' is not permitted while full-auto permission review is active")
810 .context("tool denied by full-auto allowlist"));
811 }
812
813 let skip_policy_prompt = self.policy_gateway.take_preapproved(&tool_name).await;
814
815 let decision = if skip_policy_prompt {
816 ToolExecutionDecision::Allowed
817 } else {
818 self.policy_gateway.should_execute_tool(&tool_name).await?
819 };
820
821 if !decision.is_allowed() {
822 let error_msg = match decision {
823 ToolExecutionDecision::DeniedWithFeedback(feedback) => {
824 format!("Tool '{display_name}' denied by user: {feedback}")
825 }
826 _ => format!("Tool '{display_name}' execution denied by policy"),
827 };
828
829 let _error =
830 ToolExecutionError::new(tool_name_owned.clone(), ToolErrorType::PolicyViolation, error_msg.clone());
831
832 record_failure(
833 tool_name_owned.clone(),
834 false,
835 None,
836 args_for_recording.clone(),
837 error_msg.clone(),
838 timeout_category_label.clone(),
839 base_timeout_ms,
840 adaptive_timeout_ms,
841 None,
842 false,
843 );
844
845 return Err(anyhow!("{error_msg}").context("tool denied by policy"));
846 }
847
848 let gateway = self.policy_gateway.clone();
849 let constrained_result = gateway.apply_policy_constraints(&tool_name, args).await;
850 let mut args = match constrained_result {
851 Ok(processed_args) => processed_args,
852 Err(err) => {
853 let error = ToolExecutionError::with_original_error(
854 tool_name_owned.clone(),
855 ToolErrorType::InvalidParameters,
856 "Failed to apply policy constraints".to_string(),
857 err.to_string(),
858 );
859
860 record_failure(
861 tool_name_owned,
862 false,
863 None,
864 args_for_recording,
865 format!("Failed to apply policy constraints: {err}"),
866 timeout_category_label.clone(),
867 base_timeout_ms,
868 adaptive_timeout_ms,
869 None,
870 false,
871 );
872
873 return Err(anyhow!(error.to_json_value()).context("tool denied by policy constraints"));
874 }
875 };
876
877 let super::execution_stages::ExecutionRoute {
878 route:
879 super::execution_stages::ToolRoute {
880 needs_pty,
881 tool_exists,
882 is_mcp: is_mcp_tool,
883 mcp_provider,
884 mcp_tool_name,
885 },
886 mcp_lookup_error,
887 } = self.resolve_execution_route(name, &tool_name).await;
888
889 if !tool_exists {
891 if let Some(err) = mcp_lookup_error {
892 let error = ToolExecutionError::with_original_error(
893 tool_name_owned.clone(),
894 ToolErrorType::ExecutionError,
895 format!("Failed to resolve MCP tool '{display_name}': {err}"),
896 err.to_string(),
897 );
898
899 record_failure(
900 tool_name_owned,
901 is_mcp_tool,
902 mcp_provider.clone(),
903 args_for_recording,
904 format!("Failed to resolve MCP tool '{display_name}': {err}"),
905 timeout_category_label.clone(),
906 base_timeout_ms,
907 adaptive_timeout_ms,
908 None,
909 false,
910 );
911
912 return Ok(error.to_json_value());
913 }
914
915 let (all_tool_names, similar_tools) = self.public_tool_catalog_for_error(name).await;
916 let suggestion = if !similar_tools.is_empty() {
917 format!(" Did you mean: {}?", similar_tools.join(", "))
918 } else {
919 String::new()
920 };
921 let available_tool_list = all_tool_names.join(", ");
922 let message = format!("Unknown tool: {display_name}. Available tools: {available_tool_list}.{suggestion}");
923 let error = ToolExecutionError::new(tool_name_owned.clone(), ToolErrorType::ToolNotFound, message.clone());
924
925 record_failure(
926 tool_name_owned,
927 is_mcp_tool,
928 mcp_provider.clone(),
929 args_for_recording,
930 message,
931 timeout_category_label.clone(),
932 base_timeout_ms,
933 adaptive_timeout_ms,
934 None,
935 false,
936 );
937
938 return Ok(error.to_json_value());
939 }
940
941 if is_mcp_tool && !self.mcp_circuit_breaker.allow_request() {
943 let diag = self.mcp_circuit_breaker.diagnostics();
944 let error = ToolExecutionError::new(
945 tool_name_owned.clone(),
946 ToolErrorType::ExecutionError,
947 format!("MCP circuit breaker {:?}; skipping execution", diag.status),
948 );
949 let payload = json!({
950 "error": error.to_json_value(),
951 "circuit_breaker_state": format!("{:?}", diag.status),
952 "consecutive_failures": diag.consecutive_failures,
953 "note": "MCP provider circuit breaker open; execution skipped",
954 "last_failed_at_ago_ms": diag.last_failure_time
955 .map(|ts| ts.elapsed().as_millis() as u64),
956 "current_timeout_seconds": diag.current_timeout.as_secs(),
957 "mcp_provider": mcp_provider,
958 });
959 warn!(
960 tool = %tool_name_owned,
961 payload = %payload,
962 "Skipping MCP tool execution due to circuit breaker"
963 );
964 self.execution_history.add_record(
965 ToolExecutionRecord::failure(
966 tool_name_owned,
967 requested_name.clone(),
968 is_mcp_tool,
969 mcp_provider.clone(),
970 args_for_recording,
971 format!("MCP circuit breaker {:?}; execution skipped", diag.status),
972 context_snapshot.clone(),
973 timeout_category_label.clone(),
974 base_timeout_ms,
975 adaptive_timeout_ms,
976 None,
977 false,
978 )
979 .with_circuit_breaker_state(format!("{:?}", diag.status))
980 .with_retry_after(diag.retry_after),
981 );
982 return Ok(payload);
983 }
984
985 trace!(
986 tool = %tool_name,
987 requested = %name,
988 is_mcp = is_mcp_tool,
989 uses_pty = needs_pty,
990 alias = %if tool_name == name { "" } else { name },
991 mcp_provider = %mcp_provider.as_deref().unwrap_or(""),
992 "Resolved tool route"
993 );
994
995 let _pty_guard = if needs_pty {
997 match self.start_pty_session() {
998 Ok(guard) => Some(guard),
999 Err(err) => {
1000 let error = ToolExecutionError::with_original_error(
1001 tool_name_owned.clone(),
1002 ToolErrorType::ExecutionError,
1003 "Failed to start PTY session".to_string(),
1004 err.to_string(),
1005 );
1006
1007 record_failure(
1008 tool_name_owned,
1009 is_mcp_tool,
1010 mcp_provider.clone(),
1011 args_for_recording,
1012 "Failed to start PTY session".to_string(),
1013 timeout_category_label.clone(),
1014 base_timeout_ms,
1015 adaptive_timeout_ms,
1016 None,
1017 false,
1018 );
1019
1020 return Ok(error.to_json_value());
1021 }
1022 }
1023 } else {
1024 None
1025 };
1026
1027 let execution_started_at = Instant::now();
1030 let effective_timeout = self.effective_timeout_for_call(timeout_category, &args);
1034 let effective_timeout_ms = effective_timeout.map(|d| d.as_millis() as u64);
1035
1036 let fail_open = self.optimization_config.tool_registry.middleware_fail_open;
1037 let middleware_req = ToolCallRequest {
1038 id: requested_name.clone(),
1039 tool_name: tool_name.as_str().into(),
1040 args: args.clone(),
1041 metadata: None,
1042 };
1043 if let Err(err) = self.middleware.before_execute_opt(&middleware_req, fail_open).await {
1044 if !fail_open {
1045 let error_msg = format!("Middleware denied execution: {err}");
1046 record_failure(
1047 tool_name_owned.clone(),
1048 is_mcp_tool,
1049 mcp_provider.clone(),
1050 args_for_recording.clone(),
1051 error_msg.clone(),
1052 timeout_category_label.clone(),
1053 base_timeout_ms,
1054 adaptive_timeout_ms,
1055 None,
1056 false,
1057 );
1058 return Err(anyhow!(error_msg).context("tool denied by middleware"));
1059 }
1060 }
1061
1062 if fresh_patch_read
1065 && !tool_intent::is_command_run_tool_call(&tool_name, &args)
1066 && let Some(object) = args.as_object_mut()
1067 {
1068 object.insert(
1069 crate::tools::file_ops::PATCH_READ_CACHE_NONCE.to_string(),
1070 json!(uuid::Uuid::new_v4().to_string()),
1071 );
1072 }
1073 let exec_future = async {
1074 if is_mcp_tool {
1075 let mcp_name = mcp_tool_name
1076 .as_deref()
1077 .context("MCP tool routing inconsistency: resolved MCP tool name missing")?;
1078 self.execute_mcp_tool(mcp_name, args).await
1079 } else if exec_settlement_mode.settle_noninteractive()
1080 && matches!(tool_name.as_str(), tools::UNIFIED_EXEC | tools::EXEC_COMMAND | tools::EXEC_PTY_CMD)
1081 {
1082 let exec_args = match tool_name.as_str() {
1083 tools::EXEC_COMMAND => super::executors::normalize_command_session_run_alias_args(&args, false)?,
1084 tools::EXEC_PTY_CMD => super::executors::normalize_command_session_run_alias_args(&args, true)?,
1085 _ => args.clone(),
1086 };
1087 if self.optimization_config.memory_pool.enabled {
1088 let _execution_guard = self.memory_pool.get_value();
1089 let _string_guard = self.memory_pool.get_string();
1090 let _vec_guard = self.memory_pool.get_vec();
1091 self.execute_command_session_internal(exec_args, exec_settlement_mode).await
1092 } else {
1093 self.execute_command_session_internal(exec_args, exec_settlement_mode).await
1094 }
1095 } else if exec_settlement_mode.settle_noninteractive() && tool_name == tools::WRITE_STDIN {
1096 self.execute_write_stdin(args, exec_settlement_mode).await
1097 } else if let Some(registration) = self.inventory.registration_for(&tool_name) {
1098 self.execute_registered_handler(&tool_name, ®istration, args, cached_tool.as_ref())
1099 .await
1100 } else {
1101 let (tool_names, similar_tools) = self.public_tool_catalog_for_error(&requested_name).await;
1104 let available_tool_list = tool_names.join(", ");
1105
1106 let suggestion = if !similar_tools.is_empty() {
1107 format!(" Did you mean: {}?", similar_tools.join(", "))
1108 } else {
1109 String::new()
1110 };
1111
1112 let error_msg = format!(
1113 "Tool '{display_name}' not found in registry. Available tools: {available_tool_list}.{suggestion}"
1114 );
1115
1116 let error =
1117 ToolExecutionError::new(tool_name_owned.clone(), ToolErrorType::ToolNotFound, error_msg.clone());
1118
1119 record_failure(
1120 tool_name_owned.clone(),
1121 is_mcp_tool,
1122 mcp_provider.clone(),
1123 args_for_recording.clone(),
1124 error_msg,
1125 timeout_category_label.clone(),
1126 base_timeout_ms,
1127 adaptive_timeout_ms,
1128 effective_timeout_ms,
1129 false,
1130 );
1131
1132 Ok(error.to_json_value())
1133 }
1134 };
1135
1136 let result = if let Some(limit) = effective_timeout {
1137 trace!(
1138 tool = %tool_name_owned,
1139 category = %timeout_category.label(),
1140 timeout_ms = %limit.as_millis(),
1141 "Executing tool with effective timeout"
1142 );
1143 match tokio::time::timeout(limit, exec_future).await {
1144 Ok(res) => res,
1145 Err(_) => {
1146 let timeout_ms = limit.as_millis() as u64;
1147 let tripped = self.record_tool_failure(timeout_category);
1148 if tripped {
1149 warn!(
1150 tool = %tool_name_owned,
1151 category = %timeout_category.label(),
1152 "Tool circuit breaker tripped after consecutive timeout failures"
1153 );
1154 }
1155 let retry_after = self.should_circuit_break(timeout_category);
1156
1157 let mut timeout_error = ToolExecutionError::new(
1158 tool_name_owned.clone(),
1159 ToolErrorType::Timeout,
1160 format!(
1161 "Operation '{}' exceeded the {} timeout ceiling ({}s)",
1162 tool_name_owned,
1163 timeout_category.label(),
1164 limit.as_secs()
1165 ),
1166 )
1167 .with_tool_call_context(&tool_name_owned, &args_for_recording)
1168 .with_surface("tool_registry")
1169 .with_debug_metadata("timeout_category", timeout_category.label())
1170 .with_debug_metadata("timeout_ms", timeout_ms.to_string());
1171
1172 if tool_name_owned == tools::UNIFIED_EXEC {
1173 timeout_error.recovery_suggestions = vec![
1174 Cow::Borrowed("Use write_stdin with empty chars to poll command progress"),
1175 Cow::Borrowed("Use exec_command with a fresh command if the original session is stale"),
1176 Cow::Borrowed("Ask for manual cleanup if a stale session is still active"),
1177 ];
1178 }
1179
1180 if let Some(delay) = retry_after {
1181 timeout_error.retry_after_ms = Some(delay.as_millis().min(u128::from(u64::MAX)) as u64);
1182 }
1183
1184 let mut timeout_payload = timeout_error.to_json_value();
1185 Self::annotate_timeout_error_payload(
1186 &mut timeout_payload,
1187 timeout_category.label(),
1188 timeout_ms,
1189 tripped,
1190 );
1191
1192 if let Some(breaker) = shared_circuit_breaker.as_ref() {
1193 breaker.record_failure_category_for_tool(&tool_name_owned, ErrorCategory::Timeout);
1194 }
1195 if is_mcp_tool {
1196 self.mcp_circuit_breaker.record_failure_category(ErrorCategory::Timeout);
1197 }
1198 record_failure(
1199 tool_name_owned,
1200 is_mcp_tool,
1201 mcp_provider,
1202 args_for_recording,
1203 timeout_error.user_message(),
1204 timeout_category_label.clone(),
1205 base_timeout_ms,
1206 adaptive_timeout_ms,
1207 Some(timeout_ms),
1208 tripped,
1209 );
1210 return Ok(timeout_payload);
1211 }
1212 }
1213 } else {
1214 exec_future.await
1215 };
1216
1217 match result {
1222 Ok(value) => {
1223 if let Some(breaker) = shared_circuit_breaker.as_ref() {
1224 breaker.record_success_for_tool(&tool_name_owned);
1225 }
1226 if is_mcp_tool {
1227 self.mcp_circuit_breaker.record_success();
1228 }
1229 self.reset_tool_failure(timeout_category);
1230 let should_decay = {
1231 let mut state = self.resiliency.lock();
1232 let success_streak = state.adaptive_tuning.success_streak;
1233 if let Some(counter) = state.success_trackers.get_mut(&timeout_category) {
1234 *counter = counter.saturating_add(1);
1235 let counter_val = *counter;
1236 if counter_val >= success_streak {
1237 *counter = 0;
1238 true
1239 } else {
1240 false
1241 }
1242 } else {
1243 false
1244 }
1245 };
1246 if should_decay {
1247 self.decay_adaptive_timeout(timeout_category);
1248 }
1249 self.record_tool_latency(timeout_category, execution_started_at.elapsed());
1250 let super::execution_results::ExecutionOutput { normalized_value, structured_error } = self
1251 .prepare_execution_output(
1252 &tool_name_owned,
1253 &args_for_recording,
1254 value,
1255 is_mcp_tool,
1256 max_output_tokens,
1257 )
1258 .await;
1259
1260 if !readonly_classification {
1261 self.invalidate_mutated_reads(&tool_name_owned, &args_for_recording);
1262 }
1263
1264 if let Some(error_msg) = structured_error {
1265 self.execution_history.add_record(ToolExecutionRecord::failure(
1266 tool_name_owned,
1267 requested_name,
1268 is_mcp_tool,
1269 mcp_provider,
1270 args_for_recording,
1271 error_msg,
1272 context_snapshot.clone(),
1273 timeout_category_label.clone(),
1274 base_timeout_ms,
1275 adaptive_timeout_ms,
1276 effective_timeout_ms,
1277 false,
1278 ));
1279 } else {
1280 self.execution_history.add_record(ToolExecutionRecord::success(
1281 tool_name_owned,
1282 requested_name,
1283 is_mcp_tool,
1284 mcp_provider,
1285 args_for_recording,
1286 normalized_value.clone(),
1287 context_snapshot.clone(),
1288 timeout_category_label.clone(),
1289 base_timeout_ms,
1290 adaptive_timeout_ms,
1291 effective_timeout_ms,
1292 false,
1293 ));
1294 }
1295
1296 let _ = self
1297 .middleware
1298 .after_execute(
1299 &middleware_req,
1300 &ToolCallResponse {
1301 id: middleware_req.id.clone(),
1302 success: true,
1303 result: Some(normalized_value.clone()),
1304 error: None,
1305 duration_ms: Some(execution_started_at.elapsed().as_millis() as u64),
1306 cache_hit: None,
1307 },
1308 )
1309 .await;
1310
1311 Ok(normalized_value)
1312 }
1313 Err(err) => {
1314 if err.to_string().contains("tool reentrancy blocked") {
1318 return Err(err);
1319 }
1320
1321 let error = ToolExecutionError::from_anyhow(
1322 tool_name_owned.clone(),
1323 &err,
1324 0,
1325 false,
1326 false,
1327 Some("tool_registry"),
1328 )
1329 .with_tool_call_context(&tool_name_owned, &args_for_recording);
1330 self.grant_patch_recovery_read(&error).await;
1331 let error_category = error.category;
1332 if error.circuit_breaker_impact
1333 && let Some(breaker) = shared_circuit_breaker.as_ref()
1334 {
1335 breaker.record_failure_category_for_tool(&tool_name_owned, error_category);
1336 }
1337 if error.circuit_breaker_impact && is_mcp_tool {
1338 self.mcp_circuit_breaker.record_failure_category(error_category);
1339 }
1340
1341 let tripped = if error.circuit_breaker_impact {
1342 let tripped = self.record_tool_failure(timeout_category);
1343 if tripped {
1344 warn!(
1345 tool = %tool_name_owned,
1346 category = %timeout_category.label(),
1347 "Tool circuit breaker tripped after consecutive failures"
1348 );
1349 }
1350 tripped
1351 } else {
1352 false
1353 };
1354
1355 let mut payload = error.to_json_value();
1356 Self::annotate_timeout_error_payload(
1357 &mut payload,
1358 timeout_category.label(),
1359 effective_timeout_ms.unwrap_or(0),
1360 tripped,
1361 );
1362
1363 record_failure(
1364 tool_name_owned,
1365 is_mcp_tool,
1366 mcp_provider,
1367 args_for_recording,
1368 format!("Tool execution failed: {err}"),
1369 timeout_category_label.clone(),
1370 base_timeout_ms,
1371 adaptive_timeout_ms,
1372 effective_timeout_ms,
1373 tripped,
1374 );
1375
1376 let _ = self
1377 .middleware
1378 .on_error(
1379 &middleware_req,
1380 &UnifiedToolError::new(
1381 UnifiedErrorKind::from(vtcode_commons::classify_anyhow_error(&err)),
1382 err.to_string(),
1383 ),
1384 )
1385 .await;
1386
1387 Ok(payload)
1388 }
1389 }
1390 }
1391}