1use super::AgentRunner;
2use super::continuation::{CompletionAssessment, VerificationResult};
3use super::escalation::{EscalationDecision, EscalationGate};
4use super::execute_helpers::{
5 discard_refused_assistant_message, emit_blocked_handoff_events, prepare_responses_request_messages,
6 record_terminal_turn_event, stop_reason_from_finish_reason, summarize_verification_output,
7};
8use super::helpers::detect_textual_exec_tool_call;
9use super::orchestration::EvaluatorGateOutcome;
10use super::prompt_alignment;
11use crate::config::constants::tools;
12use crate::config::models::{ModelId, Provider as ModelProvider};
13use crate::config::tool_loop_limit_reached;
14use crate::config::types::{ReasoningEffortLevel, SystemPromptMode, VerbosityLevel};
15use crate::config::{build_session_affinity_prompt_cache_key, session_affinity_key_enabled};
16use crate::core::agent::blocked_handoff::{BlockedHandoffResume, write_blocked_handoff_with_resume};
17use crate::core::agent::completion::{check_completion_candidate, check_for_response_loop};
18use crate::core::agent::events::ExecEventRecorder;
19use crate::core::agent::harness_artifacts::existing_harness_artifact_paths;
20use crate::core::agent::harness_kernel::{
21 HarnessRequestPlanInput, SessionToolCatalogSnapshot, build_harness_request_plan,
22};
23use crate::core::agent::hash_utils::stable_system_prefix_hash;
24use crate::core::agent::refusal;
25use crate::core::agent::runtime::{AgentRuntime, RuntimeControl};
26use crate::core::agent::session::AgentSessionState;
27use crate::core::agent::state::{
28 clear_old_tool_results, normalize_history_for_request_shared, should_apply_local_tool_result_clearing,
29};
30use crate::core::agent::task::{ContextItem, Task, TaskOutcome, TaskResults};
31use crate::exec::events::HarnessEventKind;
32use crate::llm::provider::{Message, ToolCall, ToolChoice, ToolDefinition, supports_responses_chaining};
33use crate::llm::providers::gemini::wire::Part;
34use crate::prompts::{
35 PromptContext, RuntimePromptContract, append_runtime_mode_sections, append_runtime_tool_prompt_sections_for_model,
36 upsert_harness_limits_section,
37};
38use crate::utils::colors::style;
39use anyhow::{Context, Result};
40use serde_json::json;
41use std::sync::Arc;
42use tracing::{debug, warn};
43
44pub(super) struct RuntimePromptBundle {
45 prompt_policy_hash: u64,
46 pub(super) request_envelope: crate::core::agent::request_envelope::SessionRequestEnvelope,
47 tool_snapshot: SessionToolCatalogSnapshot,
48 tool_def_tokens: u64,
51 pub(super) system_prompt_report: crate::prompts::system::SystemPromptReport,
57}
58
59impl RuntimePromptBundle {
60 fn has_same_provider_visible_prefix(&self, other: &Self) -> bool {
67 self.prompt_policy_hash == other.prompt_policy_hash
68 && self.tool_snapshot.tool_catalog_hash == other.tool_snapshot.tool_catalog_hash
69 && self.tool_snapshot.planning_active == other.tool_snapshot.planning_active
70 && self.tool_snapshot.request_user_input_enabled == other.tool_snapshot.request_user_input_enabled
71 && self.tool_snapshot.active_tool_names == other.tool_snapshot.active_tool_names
72 && self.request_envelope.instruction_digest() == other.request_envelope.instruction_digest()
73 && self.request_envelope.catalog_hash() == other.request_envelope.catalog_hash()
74 }
75}
76
77enum AssessmentResolution {
82 Break,
84 ForceContinue,
86 VerifyNotHandled,
89}
90
91struct MatrixExecutionCleanup(Option<Arc<crate::subagents::SubagentController>>);
93impl Drop for MatrixExecutionCleanup {
94 fn drop(&mut self) {
95 if let Some(controller) = self.0.take().filter(|controller| controller.matrix_is_driving()) {
96 controller.request_matrix_stop();
97 if let Ok(handle) = tokio::runtime::Handle::try_current() {
98 handle.spawn(async move {
99 if let Err(error) = controller.cancel_matrix().await {
100 tracing::warn!(%error, "failed to persist aborted matrix cancellation");
101 }
102 });
103 }
104 }
105 }
106}
107
108impl AgentRunner {
109 fn runtime_prompt_policy_hash(&self) -> u64 {
110 let model = self.get_selected_model();
111 let cfg = self.config();
112 let identity = crate::core::agent::hash_utils::PromptCapabilityIdentity::resolve(
113 self.provider_client.as_ref(),
114 &model,
115 self.provider_client
116 .sampling_overrides(&model)
117 .reasoning_effort
118 .or(self.reasoning_effort),
119 0,
120 );
121 crate::core::agent::hash_utils::hash_value(&(
122 identity,
123 cfg.context.max_context_tokens,
124 cfg.agent.max_system_prompt_tokens,
125 cfg.agent.harness.max_budget_usd.map(f64::to_bits),
126 format!("{:?}", cfg.agent.shell_prompt_profile),
127 cfg.agent.harness.max_tool_calls_per_turn,
128 cfg.agent.harness.max_tool_wall_clock_secs,
129 cfg.agent.harness.max_tool_retries,
130 ))
131 }
132
133 async fn compose_task_system_prompt(
134 &self,
135 prompt_tools: Arc<Vec<ToolDefinition>>,
136 is_simple_task: bool,
137 ) -> Result<(String, crate::prompts::system::SystemPromptReport)> {
138 if !is_simple_task {
139 return Ok((self.system_prompt.clone(), self.system_prompt_report.clone()));
140 }
141
142 let mut config = self.config().clone();
143 config.agent.system_prompt_mode = SystemPromptMode::Minimal;
144 let mut prompt_context = PromptContext::from_workspace_tools(
145 self._workspace.as_path(),
146 prompt_tools.iter().map(|tool| tool.function_name().to_string()),
147 );
148 prompt_context.set_current_directory(self._workspace.clone());
149 prompt_context.load_available_skills_async().await;
150
151 let (prompt, report) =
152 super::helpers::compose_system_prompt_with_appendix(self._workspace.as_path(), &config, &prompt_context)
153 .await?;
154
155 Ok((prompt, report))
156 }
157
158 async fn build_runtime_prompt_bundle(&self, is_simple_task: bool) -> Result<RuntimePromptBundle> {
159 let tool_snapshot = self.build_universal_tool_snapshot().await?;
160 let request_tools = tool_snapshot.snapshot.clone();
161 let prompt_tools = request_tools.clone().unwrap_or_else(|| Arc::new(Vec::new()));
162 let (mut system_prompt, mut system_prompt_report) =
163 self.compose_task_system_prompt(prompt_tools, is_simple_task).await?;
164 self.append_active_primary_agent_context(&mut system_prompt);
165
166 let planning_active = self.tool_registry.is_planning_active();
167 let request_user_input_enabled = self.features().request_user_input_enabled(planning_active, false);
168 let full_auto_active = self.tool_registry.current_full_auto_allowlist().await.is_some();
169
170 append_runtime_mode_sections(
171 &mut system_prompt,
172 RuntimePromptContract {
173 full_auto: full_auto_active,
174 planning_active,
175 request_user_input_enabled,
176 },
177 );
178 upsert_harness_limits_section(
179 &mut system_prompt,
180 self.config().agent.harness.max_tool_calls_per_turn,
181 self.config().agent.harness.max_tool_wall_clock_secs,
182 self.config().agent.harness.max_tool_retries,
183 );
184 let shell_profile = self.config().agent.shell_prompt_profile.resolve_for_current_platform();
185 append_runtime_tool_prompt_sections_for_model(
186 &mut system_prompt,
187 &tool_snapshot,
188 true,
189 shell_profile,
190 self.provider_client.as_ref(),
191 &self.get_selected_model(),
192 Some(self.config()),
193 );
194
195 let tool_def_tokens = request_tools
196 .as_deref()
197 .map(|tools| crate::llm::usage_cost::estimate_tool_definition_tokens(tools))
198 .unwrap_or(0);
199 let tool_count = request_tools.as_deref().map_or(0, Vec::len);
200 self.tool_registry
201 .metrics_collector()
202 .record_sdk_tool_definition_tokens(tool_def_tokens);
203 debug!(tool_def_tokens, tool_count, "tool definition overhead");
204
205 let instruction_digest = stable_system_prefix_hash(&system_prompt);
206 let capability_digest = crate::core::agent::hash_utils::PromptCapabilityIdentity::resolve(
207 self.provider_client.as_ref(),
208 &self.get_selected_model(),
209 self.reasoning_effort,
210 tool_snapshot.epoch,
211 )
212 .digest();
213 let system_instruction_prefix_hash =
214 crate::core::agent::hash_utils::hash_value(&(instruction_digest, capability_digest));
215 system_prompt_report.token_estimate = crate::prompts::system::estimate_token_count(&system_prompt);
216 system_prompt_report.over_budget =
217 system_prompt_report.token_estimate > self.config().agent.max_system_prompt_tokens;
218 let request_envelope = crate::core::agent::request_envelope::SessionRequestEnvelope::with_prefix_hash(
219 format!("{}-catalog-{}-{}", self.session_id, tool_snapshot.version, tool_snapshot.epoch),
220 system_prompt,
221 request_tools.as_deref().map_or_else(Vec::new, |tools| tools.clone()),
222 instruction_digest,
223 system_instruction_prefix_hash,
224 );
225 Ok(RuntimePromptBundle {
226 prompt_policy_hash: self.runtime_prompt_policy_hash(),
227 request_envelope,
228 tool_snapshot,
229 tool_def_tokens,
230 system_prompt_report,
231 })
232 }
233
234 fn append_active_primary_agent_context(&self, system_prompt: &mut String) {
235 let Some(active_primary_agent) = self.active_primary_agent.as_ref() else {
236 return;
237 };
238 crate::prompts::apply_coordinator_role_guidance(
239 system_prompt,
240 active_primary_agent.identity.name == "coordinator",
241 );
242
243 system_prompt.push_str("\n\n## Active Primary Agent Runtime State\n");
244 system_prompt.push_str("- Active agent: ");
245 system_prompt.push_str(&active_primary_agent.display_name);
246 system_prompt.push_str("\n- Spec name: ");
247 system_prompt.push_str(&active_primary_agent.identity.name);
248 if let Some(model) = active_primary_agent
249 .model
250 .as_deref()
251 .map(str::trim)
252 .filter(|model| !model.is_empty() && !model.eq_ignore_ascii_case("inherit"))
253 {
254 system_prompt.push_str("\n- Agent model: ");
255 system_prompt.push_str(model);
256 }
257 if let Some(reasoning_effort) = active_primary_agent.reasoning_effort.as_ref().map(|e| e.as_str()) {
258 system_prompt.push_str("\n- Agent reasoning effort: ");
259 system_prompt.push_str(reasoning_effort);
260 }
261 system_prompt.push_str("\n\n## Active Primary Agent Instructions\n");
262 system_prompt.push_str(active_primary_agent.instructions.trim());
263 }
264
265 pub(super) async fn build_validated_runtime_prompt_bundle(
266 &self,
267 is_simple_task: bool,
268 ) -> Result<RuntimePromptBundle> {
269 let mut runner = self;
270 let bundle = runner.build_runtime_prompt_bundle(is_simple_task).await?;
271 prompt_alignment::rebuild_once_on_alignment_mismatch(
272 &mut runner,
273 bundle,
274 |runner| Box::pin((*runner).build_runtime_prompt_bundle(is_simple_task)),
275 |runner, bundle| {
276 let planning_active = runner.tool_registry.is_planning_active();
277 let request_user_input_enabled = runner.features().request_user_input_enabled(planning_active, false);
278 prompt_alignment::validate_prompt_catalog_alignment(
279 bundle.request_envelope.system_prompt().as_ref(),
280 &bundle.tool_snapshot,
281 planning_active,
282 request_user_input_enabled,
283 )
284 },
285 "prompt/catalog alignment mismatch; rebuilding runtime prompt bundle",
286 "prompt/catalog alignment mismatch persisted after rebuild",
287 )
288 .await
289 }
290
291 async fn refresh_runtime_prompt_bundle_if_catalog_changed(
292 &self,
293 bundle: &mut RuntimePromptBundle,
294 is_simple_task: bool,
295 ) -> Result<bool> {
296 let current_version = self.tool_registry.tool_catalog_state().current_version();
297 let current_policy = self.runtime_prompt_policy_hash();
298 if current_version == bundle.tool_snapshot.version && current_policy == bundle.prompt_policy_hash {
299 return Ok(false);
300 }
301
302 debug!(
303 old_version = bundle.tool_snapshot.version,
304 new_version = current_version,
305 "Tool catalog changed mid-task; refreshing runtime prompt bundle"
306 );
307 let new_bundle = self.build_validated_runtime_prompt_bundle(is_simple_task).await?;
308 if new_bundle.has_same_provider_visible_prefix(bundle) {
314 debug!(
315 old_version = bundle.tool_snapshot.version,
316 new_version = new_bundle.tool_snapshot.version,
317 "Tool catalog version bumped without provider-visible change; preserving request envelope"
318 );
319 bundle.tool_snapshot = new_bundle.tool_snapshot;
320 return Ok(false);
321 }
322 *bundle = new_bundle;
323 Ok(true)
324 }
325
326 async fn resolve_completion_acceptance(
327 &mut self,
328 effective_task: &Task,
329 session_state: &mut AgentSessionState,
330 event_recorder: &mut ExecEventRecorder,
331 orchestration_enabled: bool,
332 verification_results: &[VerificationResult],
333 revision_rounds_used: &mut usize,
334 max_revision_rounds: usize,
335 should_write_blocked_handoff: &mut bool,
336 ) -> Result<bool> {
337 if !orchestration_enabled {
338 session_state.is_completed = true;
339 session_state.outcome = TaskOutcome::Success;
340 return Ok(true);
341 }
342
343 match self
344 .apply_evaluator_gate(
345 effective_task,
346 session_state,
347 event_recorder,
348 verification_results,
349 revision_rounds_used,
350 max_revision_rounds,
351 )
352 .await?
353 {
354 EvaluatorGateOutcome::Accept => {
355 session_state.is_completed = true;
356 session_state.outcome = TaskOutcome::Success;
357 Ok(true)
358 }
359 EvaluatorGateOutcome::Continue { prompt } => {
360 session_state.add_user_message(prompt);
361 Ok(false)
362 }
363 EvaluatorGateOutcome::Exhausted { reason } => {
364 session_state.outcome = TaskOutcome::failed(reason, vec![], None, None);
365 *should_write_blocked_handoff = true;
366 Ok(true)
367 }
368 }
369 }
370
371 #[allow(
388 clippy::too_many_arguments,
389 reason = "Intentional compatibility, platform, or test-only suppression."
390 )]
391 async fn resolve_completion_assessment(
392 &mut self,
393 assessment: CompletionAssessment,
394 verification_results: &[VerificationResult],
395 effective_task: &Task,
396 runtime: &mut AgentRuntime,
397 event_recorder: &mut ExecEventRecorder,
398 orchestration_enabled: bool,
399 revision_rounds_used: &mut usize,
400 max_revision_rounds: usize,
401 should_write_blocked_handoff: &mut bool,
402 ) -> Result<AssessmentResolution> {
403 match assessment {
404 CompletionAssessment::Accept => {
405 if self
406 .resolve_completion_acceptance(
407 effective_task,
408 &mut runtime.state,
409 event_recorder,
410 orchestration_enabled,
411 verification_results,
412 revision_rounds_used,
413 max_revision_rounds,
414 should_write_blocked_handoff,
415 )
416 .await?
417 {
418 return Ok(AssessmentResolution::Break);
419 }
420 Ok(AssessmentResolution::ForceContinue)
421 }
422 CompletionAssessment::SkipAccept { reason } => {
423 event_recorder.harness_event(
424 HarnessEventKind::ContinuationSkipped,
425 Some(reason),
426 None,
427 None,
428 None,
429 None,
430 None,
431 );
432 runtime.state.is_completed = true;
433 runtime.state.outcome = TaskOutcome::Success;
434 Ok(AssessmentResolution::Break)
435 }
436 CompletionAssessment::Continue { reason, prompt } => {
437 self.emit_continuation_started(event_recorder, reason);
438 runtime.state.add_user_message(prompt);
439 Ok(AssessmentResolution::ForceContinue)
440 }
441 CompletionAssessment::Verify { .. } => Ok(AssessmentResolution::VerifyNotHandled),
444 }
445 }
446
447 async fn run_verification_commands(
448 &self,
449 commands: &[String],
450 event_recorder: &mut ExecEventRecorder,
451 ) -> Result<Vec<VerificationResult>> {
452 let mut results = Vec::with_capacity(commands.len());
453 for command in commands {
454 let command_event = event_recorder.command_started(command);
455 let payload = json!({
456 "action": "run",
457 "command": command,
458 "workdir": self._workspace.display().to_string(),
459 "yield_time_ms": 1000,
460 });
461 let result = self.tool_registry.execute_harness_command_session(payload).await?;
462 let exit_code = result
463 .get("exit_code")
464 .and_then(serde_json::Value::as_i64)
465 .map(|value| value as i32);
466 let success = exit_code.unwrap_or(0) == 0;
467 let output = summarize_verification_output(&result);
468 event_recorder.command_finished(
469 &command_event,
470 if success {
471 crate::exec::events::CommandExecutionStatus::Completed
472 } else {
473 crate::exec::events::CommandExecutionStatus::Failed
474 },
475 exit_code,
476 &output,
477 );
478 results.push(VerificationResult {
479 command: command.clone(),
480 success,
481 exit_code,
482 output,
483 });
484 if !success {
485 break;
486 }
487 }
488 Ok(results)
489 }
490
491 fn emit_continuation_started(&self, event_recorder: &mut ExecEventRecorder, reason: String) {
495 event_recorder.harness_event(HarnessEventKind::ContinuationStarted, Some(reason), None, None, None, None, None);
496 }
497
498 fn emit_verification_outcome(
503 &self,
504 event_recorder: &mut ExecEventRecorder,
505 commands: &[String],
506 results: &[VerificationResult],
507 ) {
508 if let Some(failure) = results.iter().find(|result| !result.success) {
509 event_recorder.harness_event(
510 HarnessEventKind::VerificationFailed,
511 Some(super::continuation::build_verification_failure_payload(failure)),
512 Some(failure.command.clone()),
513 None,
514 failure.exit_code,
515 None,
516 None,
517 );
518 } else {
519 event_recorder.harness_event(
520 HarnessEventKind::VerificationPassed,
521 Some(format!("Verification passed: {}", commands.join(", "))),
522 commands.last().cloned(),
523 None,
524 Some(0),
525 None,
526 None,
527 );
528 }
529 }
530
531 pub async fn execute_task(&mut self, task: &Task, contexts: &[ContextItem]) -> Result<TaskResults> {
533 let _matrix_cleanup = MatrixExecutionCleanup(self.tool_registry.subagent_controller());
534 self.tool_registry.begin_tracker_request(false);
535 self.tool_registry.begin_patch_recovery_turn();
536 let setup = self.prepare_task_execution(task, contexts).await?;
539
540 let agent_prefix = setup.agent_prefix;
541 let session_store_handle = setup.session_store_handle;
542 let mut event_recorder = setup.event_recorder;
543 let run_started_at = setup.run_started_at;
544 let is_simple_task = setup.is_simple_task;
545 let mut prompt_bundle = setup.prompt_bundle;
546 let preserve_recent_turns = setup.preserve_recent_turns;
547 let max_tool_loops = setup.max_tool_loops;
548 let max_context_tokens = setup.max_context_tokens;
549 let mut runtime = setup.runtime;
550 runtime.state.stats.provider_name = self.provider_client.name().to_string();
553 let mut continuation_controller = setup.continuation_controller;
554 let effective_task = setup.effective_task;
555 let orchestration_enabled = setup.orchestration_enabled;
556 let mut budget_warning_emitted = false;
557 let mut session_costs = crate::llm::usage_cost::SessionCostAccumulator::default();
558 let max_budget_usd = setup.max_budget_usd;
559 let max_revision_rounds = setup.max_revision_rounds;
560 let mut revision_rounds_used = 0usize;
561 let mut should_write_blocked_handoff = false;
562
563 let result: Result<_> = {
564 for turn in 0..self.max_turns {
565 if matches!(runtime.poll_turn_control().await, RuntimeControl::StopRequested) {
566 self.runner_println(format_args!(
567 "{} {}",
568 agent_prefix,
569 style("Stopped by steering signal.").red().bold()
570 ));
571 runtime.state.outcome = TaskOutcome::Cancelled;
572 break;
573 }
574
575 if let Some(input) = runtime.run_until_idle() {
576 self.runner_println(format_args!(
577 "{} {}: {}",
578 agent_prefix,
579 style("Follow-up Input").cyan().bold(),
580 input
581 ));
582 }
583
584 if runtime.state.is_completed {
585 break;
586 }
587
588 self.runner_println(format_args!(
589 "{} {} is processing turn {}...",
590 agent_prefix,
591 style("(PROC)").cyan().bold(),
592 turn + 1
593 ));
594
595 let turn_model = self.get_selected_model();
596 self.refresh_runtime_prompt_bundle_if_catalog_changed(&mut prompt_bundle, is_simple_task)
597 .await?;
598 let provider_name = self.provider_client.name().to_string();
599 if let Err(error) =
600 crate::llm::usage_cost::require_budget_pricing(&provider_name, &turn_model, max_budget_usd)
601 {
602 let message = error.to_string();
603 event_recorder.turn_blocked(vtcode_exec_events::TurnBlockedEvent {
604 completed_at: None,
605 message: message.clone(),
606 last_tool: None,
607 blocked_streak: 0,
608 blocked_total: 0,
609 consecutive_cap: 0,
610 total_cap: 0,
611 recovery_active: false,
612 usage: Some(runtime.state.stats.total_usage.clone()),
613 });
614 should_write_blocked_handoff = true;
615 runtime.state.outcome = TaskOutcome::failed(
616 message,
617 Vec::new(),
618 Some(
619 "Choose a model route with complete pricing or remove the USD budget explicitly"
620 .to_string(),
621 ),
622 None,
623 );
624 break;
625 }
626
627 if std::env::var_os("VTCODE_DEBUG_PROVIDER").is_some() {
628 tracing::debug!(
629 provider_client = self.provider_client.name(),
630 turn_model,
631 "Provider debug turn selection"
632 );
633 }
634 let sampling_overrides = self.provider_client.sampling_overrides(&turn_model);
635 let turn_reasoning = if is_simple_task
636 && self
637 .provider_client
638 .supported_reasoning_efforts(&turn_model)
639 .contains(&ReasoningEffortLevel::Minimal.as_str())
640 {
641 Some(ReasoningEffortLevel::Minimal)
642 } else {
643 sampling_overrides.reasoning_effort.or(self.reasoning_effort)
644 };
645 let turn_verbosity = if is_simple_task {
646 Some(VerbosityLevel::Low)
647 } else {
648 self.verbosity
649 };
650 let max_tokens = sampling_overrides
651 .max_tokens
652 .or(if is_simple_task { Some(800) } else { Some(2000) });
653
654 let context_budget = crate::compaction::effective_context_budget(
655 Some(self.config()),
656 self.provider_client.as_ref(),
657 &turn_model,
658 );
659 runtime.state.constraints.max_context_tokens = context_budget;
660 let reserved_output_tokens = max_tokens
661 .map_or(crate::compaction::memory_envelope::DEFAULT_OUTPUT_RESERVE_TOKENS, |value| value as usize);
662 let prompt_overhead_tokens = (prompt_bundle.system_prompt_report.token_estimate as usize)
663 .saturating_add(prompt_bundle.tool_def_tokens as usize);
664 tracing::info!(
665 turn, model = %turn_model, context_budget,
666 reserved_output_tokens, prompt_overhead_tokens,
667 "Resolved per-turn context budget denominator"
668 );
669 self.maybe_auto_compact(
670 &mut runtime.state,
671 &mut event_recorder,
672 &turn_model,
673 preserve_recent_turns,
674 prompt_overhead_tokens,
675 reserved_output_tokens,
676 &mut prompt_bundle.request_envelope,
677 )
678 .await;
679 let (fits, estimated, budget) = runtime.state.preflight_token_check(
680 prompt_bundle.system_prompt_report.token_estimate as usize,
681 prompt_bundle.tool_def_tokens as usize,
682 reserved_output_tokens,
683 );
684 if !fits {
685 tracing::warn!(
690 estimated,
691 budget,
692 "Pre-flight token check failed: prompt exceeds context budget after compaction"
693 );
694 #[allow(
695 clippy::cast_sign_loss,
696 reason = "Intentional compatibility, platform, or test-only suppression."
697 )]
698 let pct = (estimated as f64 / budget.max(1) as f64 * 100.0) as u32;
699 runtime
700 .state
701 .warnings
702 .push(format!("Pre-flight check: {pct}% of context budget used before LLM call"));
703 }
704
705 let parallel_tool_config = if self.provider_client.supports_parallel_tool_config(&turn_model) {
706 Some(Box::new(crate::llm::provider::ParallelToolConfig::anthropic_optimized()))
707 } else {
708 None
709 };
710
711 let provider_kind = turn_model
712 .parse::<ModelId>()
713 .map(|model| model.provider())
714 .unwrap_or(ModelProvider::Gemini);
715
716 if matches!(provider_kind, ModelProvider::Gemini)
717 && runtime.state.conversation.len() > runtime.state.last_processed_message_idx
718 {
719 let new_messages: Vec<Message> = runtime.state.conversation
723 [runtime.state.last_processed_message_idx..]
724 .iter()
725 .map(|content| {
726 let mut text = String::new();
727 for part in &content.parts {
728 if let Part::Text { text: part_text, .. } = part {
729 if !text.is_empty() {
730 text.push('\n');
731 }
732 text.push_str(part_text);
733 }
734 }
735 match content.role.as_str() {
736 "model" => Message::assistant(text),
737 _ => Message::user(text),
738 }
739 })
740 .collect();
741 let batch_tokens: usize = new_messages.iter().map(|m| m.estimate_tokens()).sum();
742 runtime.state.adjust_token_count(batch_tokens as isize);
743 runtime.state.messages_mut().extend(new_messages);
744 runtime.state.last_processed_message_idx = runtime.state.conversation.len();
745 }
746
747 let reasoning_mapping = turn_reasoning
748 .map(|requested| {
749 crate::llm::reasoning_effort::ReasoningEffortMapper::resolve(
750 self.provider_client.as_ref(),
751 &turn_model,
752 requested,
753 self.config().agent.allow_reasoning_effort_downgrade,
754 )
755 })
756 .transpose();
757 let reasoning_effort = match reasoning_mapping {
758 Ok(mapping) => mapping,
759 Err(error) => {
760 let message = error.to_string();
761 event_recorder.turn_blocked(vtcode_exec_events::TurnBlockedEvent {
762 completed_at: None,
763 message: message.clone(),
764 last_tool: None,
765 blocked_streak: 0,
766 blocked_total: 0,
767 consecutive_cap: 0,
768 total_cap: 0,
769 recovery_active: false,
770 usage: Some(runtime.state.stats.total_usage.clone()),
771 });
772 should_write_blocked_handoff = true;
773 runtime.state.outcome = TaskOutcome::failed(
774 message,
775 Vec::new(),
776 Some("Choose a supported reasoning effort or explicitly enable downgrade".to_string()),
777 None,
778 );
779 break;
780 }
781 }
782 .map(|mapping| {
783 if mapping.degraded() {
784 tracing::warn!(requested = %mapping.requested, effective = %mapping.effective,
785 model = %turn_model, "Harness reasoning effort explicitly downgraded");
786 }
787 mapping.effective
788 })
789 .filter(|effort| *effort != ReasoningEffortLevel::None);
790 let reasoning_active = reasoning_effort.is_some_and(|effort| {
791 !matches!(effort, ReasoningEffortLevel::None | ReasoningEffortLevel::Unknown)
792 });
793 let capability_prefix_hash = crate::core::agent::hash_utils::PromptCapabilityIdentity::resolve(
794 self.provider_client.as_ref(),
795 &turn_model,
796 reasoning_effort,
797 prompt_bundle.tool_snapshot.epoch,
798 )
799 .prefix_hash(prompt_bundle.request_envelope.prefix_hash());
800
801 if runtime.state.note_reasoning_effort_change(reasoning_effort) {
806 let message = "Reasoning effort changed mid-task; provider prompt cache \
807 will be invalidated and the next request re-pays full \
808 input cost."
809 .to_string();
810 tracing::warn!("{message}");
811 runtime.state.push_warning(message);
812 }
813
814 if runtime.state.note_model_change(&turn_model) {
820 let message = "Model changed mid-task; provider prompt cache \
821 will be invalidated and the next request re-pays full \
822 input cost."
823 .to_string();
824 tracing::warn!("{message}");
825 runtime.state.push_warning(message);
826 }
827
828 let anthropic_shaped = matches!(provider_kind, ModelProvider::Anthropic | ModelProvider::Minimax);
829 let temperature = if sampling_overrides.suppresses_sampling(anthropic_shaped, reasoning_active) {
830 None
831 } else {
832 Some(sampling_overrides.temperature.unwrap_or(self.config().agent.temperature))
833 };
834 let mut top_p_override = sampling_overrides.top_p;
835 let mut top_k_override = sampling_overrides.top_k;
836 if sampling_overrides.suppresses_sampling(anthropic_shaped, reasoning_active) {
837 top_k_override = None;
842 top_p_override = top_p_override.filter(|value| *value >= 0.95);
843 }
844
845 let normalized_messages = normalize_history_for_request_shared(Arc::clone(&runtime.state.messages));
846 let normalized_messages = {
850 let clearing = &self.config().agent.harness.tool_result_clearing;
851 if should_apply_local_tool_result_clearing(&provider_name, false, clearing.enabled) {
858 Arc::new(clear_old_tool_results(
859 normalized_messages.as_slice(),
860 clearing.trigger_tokens,
861 clearing.keep_tool_uses,
862 clearing.clear_at_least_tokens,
863 clearing.clear_tool_inputs,
864 ))
865 } else {
866 normalized_messages
867 }
868 };
869 let (request_messages, previous_response_id) = prepare_responses_request_messages(
870 &mut runtime.state.previous_response_chains,
871 &provider_name,
872 self.provider_client.supports_responses_compaction(&turn_model),
873 &turn_model,
874 &normalized_messages,
875 );
876 let request_messages = match request_messages {
877 std::borrow::Cow::Borrowed(_) => Arc::clone(&normalized_messages),
878 std::borrow::Cow::Owned(messages) => Arc::new(messages),
879 };
880 let request = build_harness_request_plan(HarnessRequestPlanInput {
881 messages: request_messages,
882 system_prompt: prompt_bundle.request_envelope.system_prompt(),
883 tools: (!prompt_bundle.request_envelope.ordered_tools().is_empty())
884 .then(|| prompt_bundle.request_envelope.ordered_tools()),
885 model: turn_model.clone(),
886 max_tokens,
887 temperature,
888 top_p: top_p_override,
889 top_k: top_k_override,
890 presence_penalty: if sampling_overrides.suppresses_sampling(anthropic_shaped, reasoning_active) {
891 None
892 } else {
893 sampling_overrides.presence_penalty
894 },
895 frequency_penalty: if sampling_overrides.suppresses_sampling(anthropic_shaped, reasoning_active) {
896 None
897 } else {
898 sampling_overrides.frequency_penalty
899 },
900 stream: self.provider_client.supports_streaming(),
901 tool_choice: (provider_name.eq_ignore_ascii_case("openai")
902 && !prompt_bundle.tool_snapshot.active_tool_names.is_empty())
903 .then(|| {
904 ToolChoice::allowed_tools_auto(prompt_bundle.tool_snapshot.active_tool_names.as_ref().clone())
905 }),
906 parallel_tool_config,
907 reasoning_effort,
908 verbosity: turn_verbosity,
909 metadata: None,
910 context_management: None,
911 previous_response_id,
912 prompt_cache_key: build_session_affinity_prompt_cache_key(
926 &provider_name,
927 session_affinity_key_enabled(
928 &provider_name,
929 self.config().prompt_cache.enabled,
930 self.config().prompt_cache.providers.openai.enabled,
931 ),
932 &self.config().prompt_cache.providers.openai.prompt_cache_key_mode,
933 Some(&self.session_id),
934 ),
935 prompt_cache_profile: None,
936 tool_catalog_hash: prompt_bundle.request_envelope.catalog_hash(),
937 system_prompt_prefix_hash: Some(capability_prefix_hash),
938 })
939 .request;
940 self.validate_llm_request(&request)?;
944 let previous_response_chain_present = request.previous_response_id.is_some();
945 let sent_messages = Arc::clone(&request.messages);
950 let streaming_timeout = self
953 .config()
954 .timeouts
955 .ceiling_duration(self.config().timeouts.streaming_ceiling_seconds);
956
957 if let Some(threshold) = self
962 .config()
963 .prompt_cache
964 .gap_threshold_secs(self.config().agent.provider.as_str())
965 {
966 let threshold = std::time::Duration::from_secs(threshold);
967 if runtime.state.stats.total_usage.cached_input_tokens > 0
968 && let Some(elapsed) = runtime.state.cache_gap_exceeds(threshold)
969 {
970 let gap = crate::llm::request_gap::format_gap(elapsed);
971 let message = format!(
972 "~{gap} since the last request; the provider prompt cache has likely expired, so this request may re-pay full input cost."
973 );
974 tracing::warn!("{message}");
975 runtime.state.push_warning(message);
976 }
977 }
978 runtime.state.note_request_sent();
979
980 let turn_output = runtime
981 .run_turn_once(&mut self.provider_client, request, streaming_timeout)
982 .await?;
983 super::tool_dispatch_common::drain_and_record_runtime_events(&mut runtime, &mut event_recorder);
984 let response = turn_output.response;
985 runtime.state.stop_reason = Some(stop_reason_from_finish_reason(&response.finish_reason));
986 let refused = refusal::is_refusal(&response);
987
988 if !refused && !runtime.state.is_completed && runtime.state.record_progress_hash_and_check_stagnation()
992 {
993 let nudge = "It looks like you're repeating the same response. \
994 If you're stuck, try a different approach: break the \
995 problem into smaller steps, use different tools, or \
996 declare the task complete if you have enough information.";
997 self.runner_println(format_args!(
998 "{} {}",
999 agent_prefix,
1000 style("[STAGNATION WARNING]").yellow().bold(),
1001 ));
1002 runtime.state.add_user_message(nudge.into());
1003 }
1004 if !refused
1005 && supports_responses_chaining(
1006 &provider_name,
1007 self.provider_client.supports_responses_compaction(&turn_model),
1008 )
1009 {
1010 runtime.state.set_previous_response_chain_shared(
1011 &provider_name,
1012 &turn_model,
1013 response.request_id.as_deref(),
1014 Arc::clone(&sent_messages),
1015 );
1016 }
1017 let turn_cost = response.usage.as_ref().and_then(|usage| {
1018 let normalized = crate::llm::usage_cost::normalized_turn_usage(&provider_name, usage);
1019 crate::llm::usage_cost::estimate_session_costs(&provider_name, &turn_model, &normalized)
1020 });
1021 match session_costs.record(turn_cost) {
1022 Some(estimate) => {
1023 runtime.state.total_cost_usd = Some(estimate.effective_usd);
1024 let threshold = self.config().agent.harness.budget_warning_threshold;
1025 match crate::llm::usage_cost::BudgetStatus::classify(
1026 estimate.raw_usd,
1027 max_budget_usd,
1028 threshold,
1029 ) {
1030 crate::llm::usage_cost::BudgetStatus::Exceeded { max, .. } => {
1031 runtime.state.outcome = TaskOutcome::budget_limit_reached(max, estimate.raw_usd);
1032 break;
1033 }
1034 crate::llm::usage_cost::BudgetStatus::Warning { max, .. } if !budget_warning_emitted => {
1035 budget_warning_emitted = true;
1036 warn!(
1037 provider = %self.config().agent.provider,
1038 model = %turn_model,
1039 cost_usd = estimate.raw_usd,
1040 max_budget_usd = max,
1041 "Session cost approaching budget limit"
1042 );
1043 runtime.state.push_warning(format!(
1044 "Session cost ${:.4} has reached {:.0}% of the ${max:.2} budget. {}",
1045 estimate.raw_usd,
1046 threshold * 100.0,
1047 runtime.state.stats.total_usage.cache_summary()
1048 ));
1049 }
1050 _ => {}
1051 }
1052 }
1053 None => {
1054 runtime.state.total_cost_usd = None;
1055 if max_budget_usd.is_some() {
1056 runtime.state.outcome = TaskOutcome::failed(
1057 format!(
1058 "Pricing or usage became unavailable for `{provider_name}/{turn_model}`; budget enforcement stopped execution"
1059 ),
1060 Vec::new(),
1061 Some("Select a model route with complete pricing".to_string()),
1062 None,
1063 );
1064 should_write_blocked_handoff = true;
1065 break;
1066 }
1067 }
1068 }
1069 if refused {
1074 let reason = refusal::refusal_reason(&response);
1075 discard_refused_assistant_message(&mut runtime.state);
1076 runtime.state.clear_previous_response_chain_for(&provider_name, &turn_model);
1079 self.runner_println(format_args!(
1080 "{} {} {}",
1081 agent_prefix,
1082 style("(REFUSED)").red().bold(),
1083 reason
1084 ));
1085 runtime.state.outcome = TaskOutcome::refused(reason);
1086 break;
1087 }
1088
1089 self.runner_println(format_args!(
1090 "{} {} {} received response, processing...",
1091 agent_prefix,
1092 self.agent_type,
1093 style("(RECV)").green().bold()
1094 ));
1095
1096 self.warn_on_empty_response(
1097 &agent_prefix,
1098 response.content.as_deref().unwrap_or(""),
1099 response.tool_calls.as_ref().is_some_and(|tool_calls| !tool_calls.is_empty()),
1100 );
1101
1102 let response_text = response.content_string();
1103 if !response_text.trim().is_empty() {
1104 self.emit_final_assistant_message(&self.agent_type, &response_text);
1105 }
1106
1107 let mut effective_tool_calls = response.tool_calls.clone();
1108 let mut forced_continuation = false;
1109
1110 if effective_tool_calls.is_none()
1111 && response.content_text().len() < 150
1112 && let Some(args_value) = detect_textual_exec_tool_call(response.content_text())
1113 {
1114 effective_tool_calls = Some(vec![ToolCall::function(
1115 format!("call_text_{turn}"),
1116 tools::EXEC_COMMAND.to_string(),
1117 args_value.to_string(),
1118 )]);
1119 }
1120
1121 let is_gemini = matches!(provider_kind, ModelProvider::Gemini);
1122
1123 #[allow(
1130 unused_assignments,
1131 reason = "Intentional compatibility, platform, or test-only suppression."
1132 )] if self.config().agent.harness.confidence_escalation.enabled
1134 && !runtime.state.is_completed
1135 && let Some(tool_calls) = effective_tool_calls.as_ref()
1136 && !tool_calls.is_empty()
1137 {
1138 let escalation_config = &self.config().agent.harness.confidence_escalation;
1139 let orchestration_mode = if orchestration_enabled {
1140 "plan_build_evaluate"
1141 } else {
1142 "single"
1143 };
1144
1145 let cost_estimate = runtime.state.total_cost_usd;
1147
1148 let result = EscalationGate::decide(
1149 tool_calls,
1150 escalation_config,
1151 &runtime.state.error_recovery.lock(),
1152 cost_estimate,
1153 orchestration_mode,
1154 );
1155
1156 if result.any_escalated {
1157 let reasons: Vec<String> = result
1158 .decisions
1159 .iter()
1160 .filter_map(|d| match d {
1161 EscalationDecision::Escalate { reason, .. } => Some(reason.clone()),
1162 EscalationDecision::Proceed => None,
1163 })
1164 .collect();
1165 let summary = reasons.join("; ");
1166
1167 self.runner_println(format_args!(
1168 "{} {}: {}",
1169 agent_prefix,
1170 style("[ESCALATION REQUIRED]").red().bold(),
1171 summary
1172 ));
1173
1174 event_recorder.harness_event(
1175 HarnessEventKind::EscalationTriggered,
1176 Some(summary.clone()),
1177 None,
1178 None,
1179 None,
1180 None,
1181 None,
1182 );
1183
1184 let esc_count = runtime.state.consecutive_escalations;
1185
1186 if esc_count < escalation_config.max_replan_attempts {
1188 runtime.state.consecutive_escalations += 1;
1189 let replan_msg = format!(
1190 "The following tool calls were blocked by the safety \
1191 escalation gate:\n\n{summary}\n\n\
1192 Please try a different approach. You can use alternative \
1193 tools, break the task into smaller steps, or provide a \
1194 direct answer based on what you have already learned."
1195 );
1196 runtime.state.add_user_message(replan_msg);
1197 effective_tool_calls = None;
1199 forced_continuation = true;
1200 }
1201 else if esc_count < escalation_config.max_total_escalations
1203 && escalation_config.prompt_user_on_exhaust
1204 {
1205 runtime.state.consecutive_escalations += 1;
1206 let user_msg = format!(
1207 "I've tried multiple approaches but the safety escalation \
1208 gate keeps blocking my tool calls:\n\n{summary}\n\n\
1209 Could you provide guidance on how to proceed? \
1210 What approach should I use instead?"
1211 );
1212 runtime.state.add_user_message(user_msg);
1213 should_write_blocked_handoff = true;
1214 runtime.state.outcome = TaskOutcome::escalated(summary, "multi_tool".to_string());
1215 break;
1216 }
1217 else {
1219 runtime.state.consecutive_escalations += 1;
1220 should_write_blocked_handoff = true;
1221 runtime.state.outcome = TaskOutcome::failed(
1222 format!("Escalation chain exhausted: {summary}"),
1223 Vec::new(),
1224 Some(
1225 "The safety escalation gate repeatedly blocked tool calls. \
1226 Consider disabling the gate or adjusting the confidence \
1227 threshold if this action should be permitted."
1228 .into(),
1229 ),
1230 None,
1231 );
1232 break;
1233 }
1234 } else {
1235 runtime.state.consecutive_escalations = 0;
1237
1238 event_recorder.harness_event(
1239 HarnessEventKind::EscalationBypassed,
1240 Some("All tool calls passed escalation gate".into()),
1241 None,
1242 None,
1243 None,
1244 None,
1245 None,
1246 );
1247 }
1248 }
1249
1250 if self
1251 .active_primary_agent
1252 .as_ref()
1253 .is_some_and(|agent| agent.identity.name == "coordinator")
1254 && effective_tool_calls.as_ref().is_none_or(|calls| calls.is_empty())
1255 && let Some(controller) = self.tool_registry.subagent_controller()
1256 && let Some(snapshot) = controller.matrix_snapshot().await
1257 {
1258 use crate::exec::events::matrix::MatrixLifecycle;
1259 match snapshot.lifecycle {
1260 MatrixLifecycle::Running | MatrixLifecycle::Verifying => {
1261 let settled = controller
1262 .wait_matrix_idle()
1263 .await
1264 .context("matrix disappeared before completion")?;
1265 anyhow::ensure!(
1266 !matches!(settled.lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying),
1267 "matrix stopped without a durable result; reconcile owned cleanup before continuation"
1268 );
1269 runtime.state.is_completed = false;
1270 runtime.state.add_user_message(format!(
1271 "Durable matrix results: {}. Report outcomes from this snapshot; matrix success requires Succeeded. Other states require a coordinator decision.",
1272 serde_json::to_string(&settled)?
1273 ));
1274 forced_continuation = true;
1275 }
1276 MatrixLifecycle::Succeeded => {
1277 runtime.state.is_completed = true;
1278 runtime.state.outcome = TaskOutcome::Success;
1279 break;
1280 }
1281 MatrixLifecycle::Cancelled => {
1282 runtime.state.outcome = TaskOutcome::Cancelled;
1283 break;
1284 }
1285 MatrixLifecycle::Blocked => {
1286 runtime.state.outcome = TaskOutcome::Failed {
1287 reason: "matrix requires a coordinator decision or confirmed owned cleanup".into(),
1288 accomplished: vec![],
1289 recovery_suggestion: Some(format!(
1290 "Resume session and inspect matrix {} status",
1291 snapshot.spec.id
1292 )),
1293 checkpoint_path: None,
1294 };
1295 should_write_blocked_handoff = true;
1296 break;
1297 }
1298 MatrixLifecycle::Created | MatrixLifecycle::Paused => {}
1299 }
1300 }
1301
1302 if !forced_continuation
1303 && !runtime.state.is_completed
1304 && effective_tool_calls.as_ref().is_none_or(|tool_calls| tool_calls.is_empty())
1305 && !response.content_text().is_empty()
1306 {
1307 if check_for_response_loop(response.content_text(), &mut runtime.state) {
1308 self.runner_println(format_args!(
1309 "[{}] {}",
1310 self.agent_type,
1311 style("Repetitive assistant response detected. Breaking potential loop.")
1312 .red()
1313 .bold()
1314 ));
1315 runtime.state.outcome = TaskOutcome::LoopDetected;
1316 break;
1317 }
1318
1319 if check_completion_candidate(response.content_text()) {
1320 self.runner_println(format_args!(
1321 "[{}] {}",
1322 self.agent_type,
1323 style("Completion candidate detected; checking tracker and verification state.")
1324 .green()
1325 .bold()
1326 ));
1327 let assessment = continuation_controller
1328 .assess_completion(
1329 &effective_task,
1330 &runtime.state,
1331 self.tool_registry.tracker_adopted_for_request(),
1332 )
1333 .await?;
1334
1335 if let CompletionAssessment::Verify { commands } = &assessment {
1339 event_recorder.harness_event(
1340 HarnessEventKind::VerificationStarted,
1341 Some(format!("Running verification: {}", commands.join(", "))),
1342 commands.first().cloned(),
1343 None,
1344 None,
1345 None,
1346 None,
1347 );
1348 let verification_results =
1349 self.run_verification_commands(commands, &mut event_recorder).await?;
1350 self.emit_verification_outcome(&mut event_recorder, commands, &verification_results);
1351
1352 let post_verification =
1353 continuation_controller.after_verification(&verification_results).await?;
1354 match self
1355 .resolve_completion_assessment(
1356 post_verification,
1357 &verification_results,
1358 &effective_task,
1359 &mut runtime,
1360 &mut event_recorder,
1361 orchestration_enabled,
1362 &mut revision_rounds_used,
1363 max_revision_rounds,
1364 &mut should_write_blocked_handoff,
1365 )
1366 .await?
1367 {
1368 AssessmentResolution::Break => break,
1369 AssessmentResolution::ForceContinue => {
1370 forced_continuation = true;
1371 }
1372 AssessmentResolution::VerifyNotHandled => {}
1374 }
1375 } else {
1376 match self
1377 .resolve_completion_assessment(
1378 assessment,
1379 &[],
1380 &effective_task,
1381 &mut runtime,
1382 &mut event_recorder,
1383 orchestration_enabled,
1384 &mut revision_rounds_used,
1385 max_revision_rounds,
1386 &mut should_write_blocked_handoff,
1387 )
1388 .await?
1389 {
1390 AssessmentResolution::Break => break,
1391 AssessmentResolution::ForceContinue => {
1392 forced_continuation = true;
1393 }
1394 AssessmentResolution::VerifyNotHandled => {
1395 return Err(anyhow::anyhow!("unexpected VerifyNotHandled from assess_completion"));
1398 }
1399 }
1400 }
1401 } else {
1402 let tracker_auto_continue = self.config().agent.harness.continuation.auto_continue_tracker;
1411 let idle_limit = self.config().agent.idle_turn_limit;
1412 let idle_limit_hit = runtime.state.consecutive_idle_turns >= idle_limit;
1413 let asks_user = response.content_text().contains('?');
1414 if tracker_auto_continue && !idle_limit_hit && !asks_user {
1415 let status_text = response.content_text();
1416 let safety_handoff =
1417 crate::core::agent::completion::tracker_final_text_is_safety_handoff(status_text);
1418 if !safety_handoff {
1419 let incomplete = continuation_controller
1420 .incomplete_tracker_labels(self.tool_registry.tracker_adopted_for_request())
1421 .await?;
1422 if !incomplete.is_empty() {
1423 let joined = incomplete.join(", ");
1424 let reason = format!("Task tracker is incomplete: {joined}.");
1425 if super::continuation::tracker_status_force_continue_eligible(
1426 tracker_auto_continue,
1427 idle_limit_hit,
1428 asks_user,
1429 &reason,
1430 ) {
1431 let prompt = super::continuation::tracker_incomplete_continue_prompt(&joined);
1432 self.runner_println(format_args!(
1433 "[{}] {}: {}",
1434 self.agent_type,
1435 style("[TRACKER CONTINUE]").yellow().bold(),
1436 reason
1437 ));
1438 runtime.state.add_user_message(prompt);
1439 forced_continuation = true;
1440 }
1441 }
1442 }
1443 }
1444 }
1445 }
1446
1447 if let Some(tool_calls) = effective_tool_calls
1448 .as_ref()
1449 .filter(|tool_calls| !tool_calls.is_empty())
1450 .cloned()
1451 {
1452 self.execute_tool_call_batches(
1453 tool_calls,
1454 &mut runtime,
1455 &mut event_recorder,
1456 &agent_prefix,
1457 is_gemini,
1458 previous_response_chain_present,
1459 )
1460 .await?;
1461 super::tool_dispatch_common::drain_and_record_runtime_events(&mut runtime, &mut event_recorder);
1462
1463 if let Some(outcome) = self.tool_registry.matrix_worker_outcome() {
1464 runtime.state.outcome = if outcome == crate::exec::events::matrix::MatrixOutcome::Success {
1465 TaskOutcome::Success
1466 } else {
1467 TaskOutcome::failed(format!("Matrix worker reported {outcome:?}"), vec![], None, None)
1468 };
1469 runtime.state.is_completed = true;
1470 break;
1471 }
1472
1473 if self
1474 .active_primary_agent
1475 .as_ref()
1476 .is_some_and(|agent| agent.identity.name == "coordinator")
1477 && let Some(controller) = self.tool_registry.subagent_controller()
1478 && let Some(snapshot) = controller.matrix_snapshot().await
1479 && (matches!(
1480 snapshot.lifecycle,
1481 crate::exec::events::matrix::MatrixLifecycle::Running
1482 | crate::exec::events::matrix::MatrixLifecycle::Verifying
1483 ) || controller.matrix_is_driving())
1484 {
1485 let settled = loop {
1488 tokio::select! {
1489 settled = controller.wait_matrix_idle() => break settled.context("matrix disappeared before completion")?,
1490 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {
1491 if matches!(runtime.poll_turn_control().await, RuntimeControl::StopRequested) {
1492 controller.cancel_matrix().await?;
1493 runtime.state.outcome = TaskOutcome::Cancelled;
1494 }
1495 }
1496 }
1497 };
1498 anyhow::ensure!(
1499 !matches!(
1500 settled.lifecycle,
1501 crate::exec::events::matrix::MatrixLifecycle::Running
1502 | crate::exec::events::matrix::MatrixLifecycle::Verifying
1503 ),
1504 "matrix stopped without a durable result; reconcile owned cleanup before continuation"
1505 );
1506 if matches!(runtime.state.outcome, TaskOutcome::Cancelled) {
1507 break;
1508 }
1509 runtime.state.add_user_message(format!(
1510 "Durable matrix results: {}. Report success only for Succeeded; other states require a coordinator decision.",
1511 serde_json::to_string(&settled)?
1512 ));
1513 }
1514 }
1515
1516 if self
1519 .refresh_runtime_prompt_bundle_if_catalog_changed(&mut prompt_bundle, is_simple_task)
1520 .await?
1521 {
1522 tracing::debug!("runtime prompt bundle refreshed after tool catalog mutation");
1523 }
1524
1525 if !runtime.state.turn_tool_observations.is_empty() {
1527 let observations = std::mem::take(&mut runtime.state.turn_tool_observations);
1528 for observation in &observations {
1529 event_recorder.record_tool_outcome(
1530 &observation.tool_name,
1531 observation.attempts,
1532 observation.duration_ms,
1533 observation.error_category.as_ref().map(vtcode_commons::ErrorCategory::as_str),
1534 );
1535 }
1536 }
1537
1538 let had_effective_shell_tool_call = effective_tool_calls.as_ref().is_some_and(|calls| {
1539 calls.iter().any(|call| {
1540 call.function.as_ref().map(|function| function.name.as_str()) == Some(tools::UNIFIED_EXEC)
1541 })
1542 });
1543 let had_tool_call = response.tool_calls.as_ref().is_some_and(|tool_calls| !tool_calls.is_empty())
1544 || had_effective_shell_tool_call;
1545
1546 if had_tool_call {
1547 let loops = runtime.state.register_tool_loop();
1548 if tool_loop_limit_reached(loops, runtime.state.constraints.max_tool_loops) {
1549 let warning_message = format!(
1550 "You have reached the tool-call iteration limit of {}. \
1551 This typically means you are in a loop — repeatedly calling tools \
1552 without making progress toward the task goal.\n\n\
1553 To proceed:\n\
1554 1. Review what you have already learned from previous tool outputs.\n\
1555 2. Synthesize your findings into a concrete answer or implementation.\n\
1556 3. If you need more information, use a different approach or tool.\n\
1557 4. If you are truly stuck, explain what you have accomplished so far \
1558 and what is blocking you.",
1559 runtime.state.constraints.max_tool_loops
1560 );
1561 self.record_warning(
1562 &agent_prefix,
1563 &mut runtime.state,
1564 &mut event_recorder,
1565 warning_message.clone(),
1566 );
1567 runtime.state.add_user_message(warning_message);
1568 runtime.state.mark_tool_loop_limit_hit();
1569 break;
1570 }
1571 runtime.state.consecutive_idle_turns = 0;
1572 } else {
1573 runtime.state.reset_tool_loop_guard();
1574 if forced_continuation {
1575 runtime.state.consecutive_idle_turns = 0;
1576 } else if !runtime.state.is_completed {
1577 runtime.state.consecutive_idle_turns = runtime.state.consecutive_idle_turns.saturating_add(1);
1578 let idle_turn_limit = self.config().agent.idle_turn_limit;
1579 if runtime.state.consecutive_idle_turns >= idle_turn_limit {
1580 let warning_message = format!(
1581 "No tool calls or completion for {} consecutive turns. \
1582 The agent appears to be idle — it is responding without \
1583 taking actions or declaring the task complete.\n\n\
1584 To proceed:\n\
1585 1. Take concrete action using the available tools.\n\
1586 2. If you have enough information, present your solution.\n\
1587 3. If you are waiting for something, explain the situation.",
1588 runtime.state.consecutive_idle_turns
1589 );
1590 self.record_warning(
1591 &agent_prefix,
1592 &mut runtime.state,
1593 &mut event_recorder,
1594 warning_message.clone(),
1595 );
1596 runtime.state.add_user_message(warning_message);
1597 runtime.state.outcome = TaskOutcome::StoppedNoAction;
1598 break;
1599 }
1600 }
1601 }
1602
1603 let should_continue = forced_continuation
1604 || had_tool_call
1605 || runtime.has_pending_follow_up_inputs()
1606 || (!runtime.state.is_completed && (turn + 1) < self.max_turns);
1607
1608 if !should_continue {
1609 if runtime.state.is_completed {
1610 runtime.state.outcome = TaskOutcome::Success;
1611 } else if (turn + 1) >= self.max_turns {
1612 runtime.state.outcome = TaskOutcome::turn_limit_reached(self.max_turns, turn + 1);
1613 } else {
1614 runtime.state.outcome = TaskOutcome::StoppedNoAction;
1615 }
1616 break;
1617 }
1618 }
1619
1620 runtime.state.finalize_outcome(self.max_turns);
1621
1622 let total_duration_ms = run_started_at.elapsed().as_millis();
1623
1624 self.runner_println(format_args!("{agent_prefix} Done"));
1626
1627 let average_turn_duration_ms = if runtime.state.turn_count > 0 {
1629 Some(runtime.state.turn_total_ms as f64 / runtime.state.turn_count as f64)
1630 } else {
1631 None
1632 };
1633
1634 let max_turn_duration_ms = if runtime.state.turn_count > 0 {
1635 Some(runtime.state.turn_max_ms)
1636 } else {
1637 None
1638 };
1639
1640 let outcome = runtime.state.outcome.clone();
1641 self.thread_handle.replace_messages((*runtime.state.messages).clone());
1642 let summary = self.generate_task_summary(
1643 &effective_task,
1644 &runtime.state.modified_files,
1645 &runtime.state.executed_commands,
1646 &runtime.state.warnings,
1647 &runtime.state.messages,
1648 runtime.state.stats.turns_executed,
1649 runtime.state.max_tool_loop_streak,
1650 max_tool_loops,
1651 outcome,
1652 total_duration_ms,
1653 average_turn_duration_ms,
1654 max_turn_duration_ms,
1655 &runtime.state.stats.total_usage,
1656 );
1657
1658 if !summary.trim().is_empty() {
1659 event_recorder.agent_message(&summary);
1661 self.runner_println(format_args!(
1663 "\n{} Agent Task Summary\n{}",
1664 style("[Task]").cyan().bold(),
1665 summary
1666 ));
1667 }
1668
1669 let runtime_agent_config = self.core_agent_config();
1670 if let Err(err) = crate::persistent_memory::finalize_persistent_memory(
1671 &runtime_agent_config,
1672 Some(self.config()),
1673 &runtime.state.messages,
1674 &self.session_id,
1675 )
1676 .await
1677 {
1678 warn!(
1679 error = %err,
1680 session_id = %self.session_id,
1681 "Failed to update persistent memory"
1682 );
1683 }
1684
1685 if runtime.state.outcome.is_hard_block() || should_write_blocked_handoff {
1686 let relevant_paths = existing_harness_artifact_paths(&self._workspace);
1687 match write_blocked_handoff_with_resume(
1688 &self._workspace,
1689 &self.session_id,
1690 runtime.state.outcome.code(),
1691 &runtime.state.outcome.description(),
1692 &relevant_paths,
1693 BlockedHandoffResume::Unavailable(
1694 "Resume unavailable because the legacy agent runner does not create a session archive.",
1695 ),
1696 self.tool_registry.is_planning_active(),
1697 ) {
1698 Ok(artifacts) => emit_blocked_handoff_events(
1699 &mut event_recorder,
1700 &runtime.state.outcome.description(),
1701 &artifacts.current_path,
1702 &artifacts.archive_path,
1703 ),
1704 Err(err) => warn!(
1705 error = %err,
1706 session_id = %self.session_id,
1707 "Failed to persist blocked handoff"
1708 ),
1709 }
1710 }
1711
1712 let total_usage = runtime.state.stats.total_usage.clone();
1713 record_terminal_turn_event(&mut event_recorder, &runtime.state.outcome, total_usage.clone());
1714 event_recorder.thread_completed(
1715 &self.session_id,
1716 runtime.state.outcome.thread_completion_subtype(),
1717 runtime.state.outcome.code(),
1718 runtime.state.outcome.is_success().then_some(summary.as_str()),
1719 runtime.state.stop_reason.as_deref(),
1720 total_usage,
1721 session_costs
1722 .total()
1723 .and_then(|cost| serde_json::Number::from_f64(cost.effective_usd)),
1724 runtime.state.stats.turns_executed,
1725 );
1726 let thread_events = event_recorder.take_events();
1727 let steering_receiver = runtime.take_steering_receiver();
1728 let state = std::mem::replace(
1729 &mut runtime.state,
1730 AgentSessionState::new(self.session_id.clone(), self.max_turns, max_tool_loops, max_context_tokens),
1731 );
1732
1733 Ok((state.into_results(summary, thread_events, total_duration_ms), steering_receiver))
1734 };
1735
1736 let result = match result {
1737 Ok((task_results, steering_receiver)) => {
1738 *self.steering_receiver.lock() = steering_receiver;
1739 Ok(task_results)
1740 }
1741 Err(err) => {
1742 event_recorder.thread_failed(
1743 &self.session_id,
1744 &err.to_string(),
1745 runtime.state.stats.turns_executed.max(1),
1746 );
1747 *self.steering_receiver.lock() = runtime.take_steering_receiver();
1748 Err(err)
1749 }
1750 };
1751
1752 let persistence_result = if let Some(session_store_handle) = session_store_handle {
1753 session_store_handle
1754 .close()
1755 .await
1756 .context("authoritative session event persistence failed")
1757 } else {
1758 Ok(())
1759 };
1760
1761 self.tool_registry.set_harness_task(None);
1762 persistence_result?;
1763 result
1764 }
1765}
1766
1767#[cfg(test)]
1768mod tests {
1769 use super::{prepare_responses_request_messages, record_terminal_turn_event, tool_loop_limit_reached};
1770 use crate::core::agent::events::ExecEventRecorder;
1771 use crate::core::agent::session::AgentSessionState;
1772 use crate::core::agent::state::should_apply_local_tool_result_clearing;
1773 use crate::core::agent::task::TaskOutcome;
1774 use crate::exec::events::ThreadEvent;
1775 use crate::llm::provider::{Message, records_responses_continuation_state};
1776
1777 #[test]
1778 fn headless_local_clearing_gate_ignores_provider_context_edit_capability() {
1779 assert!(should_apply_local_tool_result_clearing("anthropic", false, true));
1784 assert!(!should_apply_local_tool_result_clearing("anthropic", false, false));
1785 assert!(!should_apply_local_tool_result_clearing("anthropic", true, true));
1788 assert!(should_apply_local_tool_result_clearing("openai", true, true));
1789 }
1790
1791 #[test]
1792 fn failed_outcome_emits_only_turn_failed() {
1793 let mut recorder = ExecEventRecorder::new("thread", None, None);
1794 recorder.turn_started();
1795
1796 record_terminal_turn_event(
1797 &mut recorder,
1798 &TaskOutcome::failed("boom".to_string(), vec![], None, None),
1799 Default::default(),
1800 );
1801
1802 let events = recorder.into_events();
1803 assert_eq!(
1804 events
1805 .iter()
1806 .filter(|event| matches!(event, ThreadEvent::TurnFailed(_)))
1807 .count(),
1808 1
1809 );
1810 assert_eq!(
1811 events
1812 .iter()
1813 .filter(|event| matches!(event, ThreadEvent::TurnCompleted(_)))
1814 .count(),
1815 0
1816 );
1817 }
1818
1819 #[test]
1820 fn successful_outcome_emits_only_turn_completed() {
1821 let mut recorder = ExecEventRecorder::new("thread", None, None);
1822 recorder.turn_started();
1823
1824 record_terminal_turn_event(&mut recorder, &TaskOutcome::Success, Default::default());
1825
1826 let events = recorder.into_events();
1827 assert_eq!(
1828 events
1829 .iter()
1830 .filter(|event| matches!(event, ThreadEvent::TurnCompleted(_)))
1831 .count(),
1832 1
1833 );
1834 assert_eq!(
1835 events
1836 .iter()
1837 .filter(|event| matches!(event, ThreadEvent::TurnFailed(_)))
1838 .count(),
1839 0
1840 );
1841 }
1842
1843 #[test]
1844 fn disabled_tool_loop_limit_never_trips() {
1845 assert!(!tool_loop_limit_reached(1, 0));
1846 assert!(!tool_loop_limit_reached(32, 0));
1847 }
1848
1849 #[test]
1850 fn openai_prepare_responses_request_messages_keeps_full_history_without_previous_response_id() {
1851 let mut state = AgentSessionState::new("session".to_string(), 4, 4, 16_000);
1852 let prior_messages = vec![Message::user("hello".to_string())];
1853 let current_messages = vec![
1854 Message::user("hello".to_string()),
1855 Message::user("continue".to_string()),
1856 ];
1857 state.set_previous_response_chain("openai", "gpt-5.6-sol", Some("resp_123"), prior_messages);
1858
1859 let (request_messages, previous_response_id) = prepare_responses_request_messages(
1860 &mut state.previous_response_chains,
1861 "openai",
1862 false,
1863 "gpt-5.6-sol",
1864 ¤t_messages,
1865 );
1866
1867 assert_eq!(previous_response_id, None);
1868 assert_eq!(request_messages.as_ref(), current_messages.as_slice());
1869 }
1870
1871 #[test]
1872 fn openai_runner_success_path_does_not_record_previous_response_chain() {
1873 let mut state = AgentSessionState::new("session".to_string(), 4, 4, 16_000);
1874 let messages = vec![Message::user("hello".to_string())];
1875
1876 if records_responses_continuation_state("openai", true) {
1877 state.set_previous_response_chain("openai", "gpt-5.6-sol", Some("resp_123"), messages);
1878 }
1879
1880 assert_eq!(state.previous_response_chain_for("openai", "gpt-5.6-sol"), None);
1881 }
1882
1883 #[test]
1884 fn compatible_runner_success_path_does_not_record_previous_response_chain() {
1885 let mut state = AgentSessionState::new("session".to_string(), 4, 4, 16_000);
1886 let messages = vec![Message::user("hello".to_string())];
1887
1888 if records_responses_continuation_state("mycorp", true) {
1889 state.set_previous_response_chain("mycorp", "gpt-5.6-sol", Some("resp_123"), messages);
1890 }
1891
1892 assert_eq!(state.previous_response_chain_for("mycorp", "gpt-5.6-sol"), None);
1893 }
1894
1895 #[test]
1896 fn gemini_prepare_responses_request_messages_keeps_full_history() {
1897 let mut state = AgentSessionState::new("session".to_string(), 4, 4, 16_000);
1898 let prior_messages = vec![Message::user("hello".to_string())];
1899 let current_messages = vec![
1900 Message::user("hello".to_string()),
1901 Message::user("continue".to_string()),
1902 ];
1903 state.set_previous_response_chain("gemini", "gemini-2.5-pro", Some("resp_123"), prior_messages);
1904
1905 let (request_messages, previous_response_id) = prepare_responses_request_messages(
1906 &mut state.previous_response_chains,
1907 "gemini",
1908 false,
1909 "gemini-2.5-pro",
1910 ¤t_messages,
1911 );
1912
1913 assert_eq!(previous_response_id.as_deref(), Some("resp_123"));
1914 assert_eq!(request_messages.as_ref(), current_messages.as_slice());
1915 }
1916
1917 #[test]
1918 fn compatible_prepare_responses_request_messages_keeps_custom_provider_stateless() {
1919 let mut state = AgentSessionState::new("session".to_string(), 4, 4, 16_000);
1920 let prior_messages = vec![Message::user("hello".to_string())];
1921 let current_messages = vec![
1922 Message::user("hello".to_string()),
1923 Message::user("continue".to_string()),
1924 ];
1925 state.set_previous_response_chain("mycorp", "gpt-5.6-sol", Some("resp_123"), prior_messages);
1926
1927 let (request_messages, previous_response_id) = prepare_responses_request_messages(
1928 &mut state.previous_response_chains,
1929 "mycorp",
1930 true,
1931 "gpt-5.6-sol",
1932 ¤t_messages,
1933 );
1934
1935 assert_eq!(previous_response_id, None);
1936 assert_eq!(request_messages.as_ref(), current_messages.as_slice());
1937 }
1938}