Skip to main content

vtcode_core/core/agent/runner/
execute.rs

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    /// Estimated token overhead of `request_tools`, computed once per
49    /// snapshot (see [`crate::llm::usage_cost::estimate_tool_definition_tokens`]).
50    tool_def_tokens: u64,
51    /// Token-budget report for the composed system instruction sections
52    /// (measured before the runtime-mode/tool-guideline/harness-limits
53    /// sections appended in `build_runtime_prompt_bundle`, which are runtime
54    /// scaffolding rather than trimmable prompt layers). Drives the
55    /// preflight token check and the over-budget user warning.
56    pub(super) system_prompt_report: crate::prompts::system::SystemPromptReport,
57}
58
59impl RuntimePromptBundle {
60    /// Provider-visible prefix unchanged across a catalog version bump.
61    ///
62    /// Compares only wire-visible identity (stable prompt digest, tool bytes,
63    /// active set, planning flags). Version/epoch and token estimates are
64    /// bookkeeping and intentionally excluded so no-op MCP refreshes preserve
65    /// the frozen envelope and keep the provider prefix cache hot.
66    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
77/// Outcome of [`AgentRunner::resolve_completion_assessment`].
78///
79/// Collapses the duplicated `CompletionAssessment` handling (pre- and
80/// post-verification) into a single signal the turn loop reacts to.
81enum AssessmentResolution {
82    /// Assessment resolved to acceptance or skip — the turn loop should break.
83    Break,
84    /// Assessment requested more work — the turn loop should force a continuation.
85    ForceContinue,
86    /// Assessment is `Verify`; the caller runs verification and re-dispatches the
87    /// `after_verification` result through the same helper.
88    VerifyNotHandled,
89}
90
91/// An aborted headless parent cannot leave scheduler-owned commands running.
92struct 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        // No-op refresh: version/epoch bumped without provider-visible change
309        // (e.g. MCP refresh returning identical tools). Preserve the frozen
310        // envelope, token estimates, and reports so the provider prefix cache
311        // stays hot; only advance snapshot bookkeeping so the next version
312        // check passes.
313        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    /// Resolve a [`CompletionAssessment`] produced by either
372    /// [`super::continuation::ContinuationController::assess_completion`]
373    /// (pre-verification) or
374    /// [`super::continuation::ContinuationController::after_verification`]
375    /// (post-verification) into a
376    /// single [`AssessmentResolution`] the turn loop reacts to.
377    ///
378    /// This consolidates the previously duplicated match arms for `Accept`,
379    /// `SkipAccept`, and `Continue`. `Verify` is returned as
380    /// [`AssessmentResolution::VerifyNotHandled`] so the caller can run
381    /// verification commands and re-dispatch the `after_verification` result
382    /// through this same helper.
383    ///
384    /// Note: `SkipAccept` intentionally bypasses the evaluator gate (mirroring
385    /// the original pre-verification path) regardless of which controller call
386    /// produced it.
387    #[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            // Handled by the caller, which has access to the continuation
442            // controller and runs verification before re-dispatching.
443            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    /// Emit the `ContinuationStarted` harness event for a forced-continuation
492    /// assessment. Centralizes the (reason-only) event shape used by both the
493    /// pre- and post-verification continuation paths.
494    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    /// Emit `VerificationFailed` (for the first failing command) or
499    /// `VerificationPassed` based on the verification run. Reuses
500    /// [`super::continuation::build_verification_failure_payload`] so the failure
501    /// headline stays in sync with the continuation prompt builder.
502    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    /// Execute a task with this agent
532    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        // Phase 1: Setup — harness alignment, conversation building, session init,
537        // orchestration planning. Extracted to `prepare_task_execution` for testability.
538        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        // Normalize usage from the provider actually serving this session;
551        // configuration may contain an alias (or be inferred from the model).
552        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                    // Advisory only: the estimate can overshoot (un-droppable
686                    // messages, suppressed compaction), and the provider owns
687                    // the authoritative context limit. Aborting here turned a
688                    // recoverable situation into a dead `exec`/auto run.
689                    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                    // Collect new messages first to avoid borrow conflict:
720                    // conversation is immutably borrowed in the loop, but
721                    // adjust_token_count/push need &mut state.
722                    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                // Reasoning-effort-change advisory (Phase E4): a mid-task
802                // change to the reasoning effort alters the request prefix,
803                // which invalidates the provider prompt cache for the next
804                // request.
805                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                // Model-change advisory: prompt caches are unique per model,
815                // so a mid-task switch rebuilds the cache at full input cost
816                // even when the rest of the prefix is unchanged. Prefer
817                // resolving the model up front; when a switch is unavoidable,
818                // expect one full-price request before hits resume.
819                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                    // Keep the payload inside what our Anthropic reasoning
838                    // validator accepts: `validate_reasoning_constraints`
839                    // rejects any top_k and requires top_p in [0.95, 1.0]
840                    // while extended thinking is active.
841                    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                // Local stand-in for Anthropic `clear_tool_uses` on routes
847                // without context edits (same gate as the interactive
848                // runloop). Request-only: durable history is untouched.
849                let normalized_messages = {
850                    let clearing = &self.config().agent.harness.tool_result_clearing;
851                    // Headless requests never attach native context
852                    // management (`context_management: None` in the
853                    // HarnessRequestPlanInput below), so the gate must not
854                    // consume the provider capability: for Anthropic that
855                    // combination would apply neither native nor local
856                    // clearing and old tool bodies would ride every request.
857                    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                    // Keep the wire key stable per session. OpenAI routes by
913                    // (prefix hash + key); the key must stay consistent across
914                    // requests sharing a prefix ("Use prompt_cache_key
915                    // consistently"). Per-turn suffixes (capability/catalog
916                    // hashes) fragment routing buckets without benefit: distinct
917                    // prefixes already route distinctly via their prefix hash,
918                    // and a suffix that changes on tool-catalog churn busts an
919                    // otherwise reusable routing affinity every few turns.
920                    // Prefix identity stays tracked per request via
921                    // `tool_catalog_hash` / `system_prompt_prefix_hash` below.
922                    // OpenRouter/xAI receive namespaced session lineage for
923                    // documented sticky routing even when the OpenAI cache
924                    // block is off.
925                    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                // Cheap pre-flight: catch malformed requests (empty system
941                // prompt, no messages, duplicate tool names, missing required
942                // properties) before paying for an API round-trip.
943                self.validate_llm_request(&request)?;
944                let previous_response_chain_present = request.previous_response_id.is_some();
945                // O(1) Arc bump: keeps the sent history available for
946                // set_previous_response_chain while the request still holds it
947                // for providers that validate messages (e.g. MiMo) during
948                // stream().
949                let sent_messages = Arc::clone(&request.messages);
950                // Compute timeout before the call to avoid simultaneous mutable/immutable
951                // borrows of `self` (provider_client vs config).
952                let streaming_timeout = self
953                    .config()
954                    .timeouts
955                    .ceiling_duration(self.config().timeouts.streaming_ceiling_seconds);
956
957                // Cache-gap advisory (Phase E1): warn once per gap when the
958                // provider prompt cache has likely expired since the last
959                // request, so this request may unexpectedly re-pay full
960                // input cost.
961                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                // --- Progress stagnation detection ---
989                // If the assistant produces near-identical responses across consecutive
990                // turns (no tool calls, no progress), inject a nudge to break the loop.
991                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                // A refusal is terminal for this prompt: idle recovery would
1070                // resend it and be refused again, and any partial output or
1071                // tool calls were cut off by the provider, so neither is
1072                // committed as an answer nor executed.
1073                if refused {
1074                    let reason = refusal::refusal_reason(&response);
1075                    discard_refused_assistant_message(&mut runtime.state);
1076                    // The refused response is not part of the kept history, so
1077                    // a stored continuation id would point past it.
1078                    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                // --- Confidence-based escalation gate + chain ---
1124                // Evaluate tool calls before dispatching them.  Implements the
1125                // escalation chain: re-plan → prompt user → abort with partial
1126                // results.  (Steps 1–2 already exist at the tool-exec level:
1127                // auto-retry and alternative-tool fallback are handled in
1128                // tool_exec.rs and execution_facade.rs.)
1129                #[allow(
1130                    unused_assignments,
1131                    reason = "Intentional compatibility, platform, or test-only suppression."
1132                )] // compiler keeps flags but this is clear
1133                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                    // Reuse the per-turn cost estimate from the budget check above.
1146                    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                        // --- Step 3: Re-plan via conversation injection ---
1187                        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                            // Skip dispatching blocked tool calls — let the LLM re-plan
1198                            effective_tool_calls = None;
1199                            forced_continuation = true;
1200                        }
1201                        // --- Step 4: Prompt user for guidance ---
1202                        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                        // --- Step 5: Abort with partial results ---
1218                        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                        // Tool calls passed escalation gate — reset chain counter
1236                        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                        // Verify requires running verification commands before
1336                        // re-dispatching the after_verification result through the
1337                        // same helper. All other variants are resolved directly.
1338                        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                                // `after_verification` never yields `Verify`.
1373                                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                                    // Verify is handled in the if-branch above;
1396                                    // the helper only returns this for Verify.
1397                                    return Err(anyhow::anyhow!("unexpected VerifyNotHandled from assess_completion"));
1398                                }
1399                            }
1400                        }
1401                    } else {
1402                        // Tracker-aware status continuation: text-only responses
1403                        // that are not completion candidates still must not end
1404                        // the run while the task tracker has incomplete steps.
1405                        //
1406                        // Use a read-only tracker probe — never `assess_completion`
1407                        // — so status text cannot create/complete the internal
1408                        // scaffold or invent a pending verify step under
1409                        // ContinuationPolicy::All.
1410                        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                        // Headless sessions await owned work before another model turn;
1486                        // status polling must not exhaust the coordinator's loop budget.
1487                        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                // Refresh tool definitions if the catalog was mutated during tool
1517                // execution (e.g. tools.load / tools.unload / skill activation).
1518                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                // --- Emit tool latency events ---
1526                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            // Agent execution completed
1625            self.runner_println(format_args!("{agent_prefix} Done"));
1626
1627            // Generate meaningful summary based on agent actions
1628            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                // Record summary as agent message for event stream
1660                event_recorder.agent_message(&summary);
1661                // Also display summary prominently for immediate visibility in TUI transcript
1662                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        // Headless requests never attach native context management, so the
1780        // gate must run with context_edits == false even for Anthropic;
1781        // otherwise headless Anthropic gets neither native nor local
1782        // clearing and old tool bodies ride every request.
1783        assert!(should_apply_local_tool_result_clearing("anthropic", false, true));
1784        assert!(!should_apply_local_tool_result_clearing("anthropic", false, false));
1785        // Interactive semantics stay pinned: native edits exclude the local
1786        // pass, other providers use it regardless of capability.
1787        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            &current_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            &current_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            &current_messages,
1933        );
1934
1935        assert_eq!(previous_response_id, None);
1936        assert_eq!(request_messages.as_ref(), current_messages.as_slice());
1937    }
1938}