use anyhow::Result;
use vtcode_commons::ErrorCategory;
use vtcode_core::llm::provider as uni;
use vtcode_core::tools::handlers::planning_workflow::{PlanningWorkflowState, persist_plan_draft};
use vtcode_core::utils::ansi::{AnsiRenderer, MessageStyle};
use super::{
MAX_POST_TOOL_RECOVERY_CYCLES, PLANNING_RECOVERY_SYNTHESIS_FALLBACK,
PLANNING_RECOVERY_SYNTHESIS_FALLBACK_NO_INTERVIEW, POST_TOOL_RECOVERY_REASON, POST_TOOL_RECOVERY_REASON_PLAN_MODE,
POST_TOOL_RESUME_DIRECTIVE, RECOVERY_CONTRACT_VIOLATION_REASON, RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER,
};
use crate::agent::runloop::unified::plan_blocks::extract_any_plan;
use crate::agent::runloop::unified::planning_workflow_state::{
PlanningWorkflowSessionState, short_confirmation_hint_with_fallback,
};
use crate::agent::runloop::unified::run_loop_context::HarnessTurnState;
use crate::agent::runloop::unified::turn::context::TurnLoopResult;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum PostToolFailureRecovery {
NotApplicable,
RetryToolFree,
StopAfterDirective,
}
pub(super) fn has_tool_response_since(messages: &[uni::Message], baseline_len: usize) -> bool {
messages
.get(baseline_len..)
.is_some_and(|recent| recent.iter().any(|msg| msg.role == uni::MessageRole::Tool))
}
fn ensure_recent_system_message(working_history: &mut Vec<uni::Message>, content: &str) {
let already_present = working_history
.iter()
.rev()
.take(3)
.any(|message| message.role == uni::MessageRole::System && message.content.as_text() == content);
if already_present {
return;
}
working_history.push(uni::Message::system(content.to_string()));
}
pub(super) fn ensure_post_tool_resume_directive(working_history: &mut Vec<uni::Message>) {
ensure_recent_system_message(working_history, POST_TOOL_RESUME_DIRECTIVE);
}
pub(crate) fn prepare_post_tool_tool_free_recovery(working_history: &mut Vec<uni::Message>, reason: &str) {
ensure_recent_system_message(working_history, reason);
}
pub(super) fn maybe_recover_after_post_tool_llm_failure(
renderer: &mut AnsiRenderer,
working_history: &mut Vec<uni::Message>,
err: &anyhow::Error,
step_count: usize,
turn_history_start_len: usize,
failure_stage: &'static str,
allow_tool_free_retry: bool,
planning_active: bool,
) -> Result<PostToolFailureRecovery> {
let has_partial_tool_progress = has_tool_response_since(working_history, turn_history_start_len);
if !has_partial_tool_progress {
return Ok(PostToolFailureRecovery::NotApplicable);
}
let err_cat = vtcode_commons::classify_anyhow_error(err);
let transient_hint = if err_cat.is_retryable() {
" (transient — may resolve on retry)"
} else {
""
};
let summary =
format!("Tool execution completed, but the model follow-up failed{transient_hint}. Output above is valid.",);
renderer.line(MessageStyle::Info, &summary)?;
renderer.line(MessageStyle::Info, &format!("Follow-up error category: {}", err_cat.user_label()))?;
if !err_cat.is_retryable() {
renderer.line(
MessageStyle::Info,
"Tip: rerun with a narrower prompt or switch provider/model for the follow-up.",
)?;
}
let should_retry =
allow_tool_free_retry && (err_cat.is_retryable() || matches!(err_cat, ErrorCategory::ExecutionError));
let action = if should_retry {
let reason = if planning_active {
POST_TOOL_RECOVERY_REASON_PLAN_MODE
} else {
POST_TOOL_RECOVERY_REASON
};
prepare_post_tool_tool_free_recovery(working_history, reason);
renderer.line(
MessageStyle::Info,
"[!] Follow-up failed after tool execution; scheduling a final tool-free recovery pass.",
)?;
PostToolFailureRecovery::RetryToolFree
} else {
ensure_post_tool_resume_directive(working_history);
PostToolFailureRecovery::StopAfterDirective
};
tracing::warn!(
error = %err,
step = step_count,
stage = failure_stage,
category = ?err_cat,
retryable = err_cat.is_retryable(),
recovery_action = ?action,
"Recovered turn after post-tool LLM phase failure"
);
Ok(action)
}
fn gather_files_read_this_turn(working_history: &[uni::Message]) -> Vec<String> {
let mut files = Vec::new();
let mut seen = std::collections::HashSet::new();
for msg in working_history.iter() {
if msg.role != uni::MessageRole::Tool {
continue;
}
let text = msg.content.as_text();
if let Ok(val) = serde_json::from_str::<serde_json::Value>(&text) {
if let Some(path) = val.get("path").and_then(serde_json::Value::as_str) {
if seen.insert(path.to_string()) {
files.push(path.to_string());
}
}
}
}
files
}
fn plan_mode_recovery_fallback(
salvaged_text: Option<String>,
structured_message: &str,
working_history: &[uni::Message],
) -> String {
if let Some(text) = salvaged_text
&& text.trim().contains("<proposed_plan")
{
text.trim().to_string()
} else {
build_recovery_fallback(working_history, structured_message)
}
}
fn build_recovery_fallback(working_history: &[uni::Message], lead_in: &str) -> String {
let files_read = gather_files_read_this_turn(working_history);
if files_read.is_empty() {
lead_in.to_string()
} else {
format!(
"{lead_in}\n\nFiles already read this turn (do NOT re-read):\n{}",
files_read.iter().map(|f| format!(" - {f}")).collect::<Vec<_>>().join("\n")
)
}
}
pub(super) async fn complete_turn_after_failed_tool_free_recovery(
working_history: &mut Vec<uni::Message>,
failure_stage: &str,
err: Option<&anyhow::Error>,
salvaged_text: Option<String>,
plan_session: Option<&mut PlanningWorkflowSessionState>,
plan_state: Option<&PlanningWorkflowState>,
) -> TurnLoopResult {
if let (Some(state), Some(salvaged)) = (plan_state, salvaged_text.as_ref()) {
if let Some(plan_text) = extract_any_plan(salvaged).plan_text {
if state.get_plan_file().await.is_some() {
if let Err(e) = persist_plan_draft(state, &plan_text).await {
tracing::warn!(
error = %e,
"plan-mode recovery: failed to persist salvaged plan to session plan file"
);
}
}
}
}
let is_transient_error = err
.map(|e| vtcode_commons::classify_anyhow_error(e).is_retryable())
.unwrap_or(false);
if let Some(plan_session) = plan_session {
if plan_session.is_budget_exhausted()
|| plan_session.is_recovery_exhausted()
|| plan_session.is_interview_denied()
{
let finalize_message = if plan_session.is_budget_exhausted() {
super::PLANNING_BUDGET_EXHAUSTED_USER_NOTICE
} else if plan_session.is_recovery_exhausted() {
super::PLANNING_RECOVERY_EXHAUSTED_USER_NOTICE
} else {
PLANNING_RECOVERY_SYNTHESIS_FALLBACK_NO_INTERVIEW
};
let mut planning_fallback = plan_mode_recovery_fallback(salvaged_text, finalize_message, working_history);
planning_fallback.push_str("\n\n");
planning_fallback.push_str(&short_confirmation_hint_with_fallback());
push_final_answer_if_absent(working_history, &planning_fallback);
tracing::warn!(
stage = failure_stage,
budget_exhausted = plan_session.is_budget_exhausted(),
recovery_exhausted = plan_session.is_recovery_exhausted(),
interview_denied = plan_session.is_interview_denied(),
"Plan-mode tool-free recovery failed; finalizing plan from gathered evidence."
);
return TurnLoopResult::Completed;
}
plan_session.mark_interview_pending();
let planning_fallback =
plan_mode_recovery_fallback(salvaged_text, PLANNING_RECOVERY_SYNTHESIS_FALLBACK, working_history);
push_final_answer_if_absent(working_history, &planning_fallback);
tracing::warn!(
stage = failure_stage,
transient_error = is_transient_error,
"Plan-mode tool-free recovery failed; marking interview pending for next turn."
);
return TurnLoopResult::Completed;
}
if let Some(salvaged) = salvaged_text.filter(|text| !text.trim().is_empty()) {
let answer = format!(
"[!] Recovery synthesis was interrupted; best-effort answer below \
(tool-call markup removed):\n\n{salvaged}"
);
push_final_answer_if_absent(working_history, &answer);
tracing::warn!(
stage = failure_stage,
"Tool-free recovery failed; concluding turn with salvaged synthesis prose."
);
return TurnLoopResult::Completed;
}
let fallback = build_recovery_fallback(working_history, RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER);
push_final_answer_if_absent(working_history, &fallback);
tracing::warn!(
stage = failure_stage,
error = ?err,
"Final tool-free recovery pass failed; concluding turn with deterministic fallback answer."
);
TurnLoopResult::Completed
}
fn push_final_answer_if_absent(working_history: &mut Vec<uni::Message>, text: &str) {
let already_present = working_history.iter().rev().take(3).any(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == text
});
if !already_present {
working_history
.push(uni::Message::assistant(text.to_string()).with_phase(Some(uni::AssistantPhase::FinalAnswer)));
}
}
pub(super) async fn normalize_tool_free_recovery_break_outcome(
working_history: &mut Vec<uni::Message>,
outcome_result: TurnLoopResult,
tool_free_recovery: bool,
salvaged_text: Option<String>,
plan_session: Option<&mut PlanningWorkflowSessionState>,
plan_state: Option<&PlanningWorkflowState>,
) -> TurnLoopResult {
let should_fallback = tool_free_recovery
&& matches!(
outcome_result,
TurnLoopResult::Blocked {
reason: Some(ref reason)
} if reason == RECOVERY_CONTRACT_VIOLATION_REASON
);
if should_fallback {
return complete_turn_after_failed_tool_free_recovery(
working_history,
"handle_turn_processing_result.tool_free_recovery_contract_violation",
None,
salvaged_text,
plan_session,
plan_state,
)
.await;
}
outcome_result
}
#[derive(Debug)]
pub(super) enum PostToolFailureAction {
Continue,
Break(TurnLoopResult),
Fallthrough,
}
pub(super) struct PostToolRecoveryContext<'a> {
pub renderer: &'a mut AnsiRenderer,
pub working_history: &'a mut Vec<uni::Message>,
pub harness_state: &'a mut HarnessTurnState,
pub plan_session: Option<&'a mut PlanningWorkflowSessionState>,
pub plan_state: Option<&'a PlanningWorkflowState>,
pub err: &'a anyhow::Error,
pub step_count: usize,
pub turn_history_start_len: usize,
pub stage: &'static str,
pub tool_free_recovery: bool,
}
pub(super) async fn dispatch_post_tool_failure(ctx: PostToolRecoveryContext<'_>) -> Result<PostToolFailureAction> {
let PostToolRecoveryContext {
renderer,
working_history,
harness_state,
mut plan_session,
plan_state,
err,
step_count,
turn_history_start_len,
stage,
tool_free_recovery,
} = ctx;
let planning_active = plan_session.is_some();
if (harness_state.wall_clock_exhausted_emitted
|| harness_state.wall_clock_exhausted()
|| harness_state.tool_budget_exhausted_emitted)
&& let Some(session) = plan_session.as_deref_mut()
{
session.mark_recovery_exhausted();
}
let recovery = maybe_recover_after_post_tool_llm_failure(
renderer,
working_history,
err,
step_count,
turn_history_start_len,
stage,
!tool_free_recovery,
planning_active,
)?;
match recovery {
PostToolFailureRecovery::NotApplicable => {
if tool_free_recovery {
let salvaged = harness_state.take_recovery_rejected_synthesis();
let direct_stage = concat_compact(stage, ".direct_tool_free_failure");
let result = complete_turn_after_failed_tool_free_recovery(
working_history,
&direct_stage,
Some(err),
salvaged,
plan_session,
plan_state,
)
.await;
Ok(PostToolFailureAction::Break(result))
} else {
Ok(PostToolFailureAction::Fallthrough)
}
}
PostToolFailureRecovery::RetryToolFree => {
let salvaged = harness_state.take_recovery_rejected_synthesis();
let cycle_stage = concat_compact(stage, ".recovery_cycle_cap");
if let Some(r) = check_recovery_cycle_cap(
harness_state.post_tool_recovery_cycles(),
working_history,
&cycle_stage,
err,
salvaged,
plan_session,
plan_state,
)
.await
{
return Ok(PostToolFailureAction::Break(r));
}
harness_state.increment_post_tool_recovery_cycle();
harness_state.switch_to_tool_free_recovery();
Ok(PostToolFailureAction::Continue)
}
PostToolFailureRecovery::StopAfterDirective => {
let result = if tool_free_recovery {
let salvaged = harness_state.take_recovery_rejected_synthesis();
let directive_stage = concat_compact(stage, ".stop_after_directive");
complete_turn_after_failed_tool_free_recovery(
working_history,
&directive_stage,
Some(err),
salvaged,
plan_session,
plan_state,
)
.await
} else {
TurnLoopResult::Completed
};
Ok(PostToolFailureAction::Break(result))
}
}
}
fn concat_compact(a: &str, b: &str) -> String {
let mut buf = String::with_capacity(a.len() + b.len());
buf.push_str(a);
buf.push_str(b);
buf
}
async fn check_recovery_cycle_cap(
cycles: u8,
working_history: &mut Vec<uni::Message>,
stage: &str,
err: &anyhow::Error,
salvaged_text: Option<String>,
mut plan_session: Option<&mut PlanningWorkflowSessionState>,
plan_state: Option<&PlanningWorkflowState>,
) -> Option<TurnLoopResult> {
if cycles >= MAX_POST_TOOL_RECOVERY_CYCLES {
tracing::warn!(
cycles,
"Post-tool recovery cycle cap reached; concluding turn \
with deterministic fallback answer"
);
if let Some(plan_session) = plan_session.as_deref_mut() {
plan_session.mark_recovery_exhausted();
}
return Some(
complete_turn_after_failed_tool_free_recovery(
working_history,
stage,
Some(err),
salvaged_text,
plan_session,
plan_state,
)
.await,
);
}
None
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::runloop::unified::planning_workflow_state::PlanningWorkflowSessionState;
use vtcode_commons::llm::LLMError;
fn transient_err() -> anyhow::Error {
anyhow::Error::new(LLMError::Network {
message: "simulated network blip".to_string(),
metadata: None,
})
}
#[tokio::test]
async fn tool_free_recovery_keeps_planning_alive_on_transient_error() {
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
Some(&transient_err()),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
assert!(
plan_session.interview_pending(),
"transient error must keep planning alive by re-forcing the interview"
);
}
#[tokio::test]
async fn tool_free_recovery_keeps_planning_alive_on_non_transient_error() {
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
let err = anyhow::Error::new(LLMError::InvalidRequest { message: "bad request".to_string(), metadata: None });
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
Some(&err),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
assert!(
plan_session.interview_pending(),
"any tool-free recovery failure must keep planning alive (not dead-end)"
);
}
#[tokio::test]
async fn dispatch_marks_recovery_exhausted_when_wall_clock_exhausted_in_plan_mode() {
use crate::agent::runloop::unified::run_loop_context::{HarnessTurnState, TurnId, TurnRunId};
use vtcode_core::utils::ansi::AnsiRenderer;
let mut renderer = AnsiRenderer::stdout();
let mut working_history: Vec<uni::Message> = Vec::new();
let mut harness_state =
HarnessTurnState::new(TurnRunId("test-run".to_string()), TurnId("test-turn".to_string()), 4, 600, 0);
harness_state.wall_clock_exhausted_emitted = true;
let mut plan_session = PlanningWorkflowSessionState::default();
let err = transient_err();
let action = dispatch_post_tool_failure(PostToolRecoveryContext {
renderer: &mut renderer,
working_history: &mut working_history,
harness_state: &mut harness_state,
plan_session: Some(&mut plan_session),
plan_state: None,
err: &err,
step_count: 1,
turn_history_start_len: 0,
stage: "stage",
tool_free_recovery: true,
})
.await
.expect("dispatch must not error");
assert!(
plan_session.is_recovery_exhausted(),
"wall-clock exhaustion during planning must mark the session \
recovery-exhausted so the plan finalizes instead of looping"
);
assert!(!plan_session.interview_pending(), "must not re-force the interview after wall-clock exhaustion");
assert!(matches!(action, PostToolFailureAction::Break(_)));
}
#[tokio::test]
async fn tool_free_recovery_finalizes_when_budget_exhausted() {
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.mark_budget_exhausted();
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
Some(&transient_err()),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
assert!(
!plan_session.interview_pending(),
"budget-exhausted must not re-force the interview (would loop forever)"
);
assert!(
working_history.iter().any(|m| m.role == uni::MessageRole::Assistant),
"budget-exhausted must finalize the plan with a fallback answer"
);
}
#[tokio::test]
async fn plan_mode_recovery_rejects_garbled_tool_call_salvage() {
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
let garbled = "I have enough to plan. <invoke name=\"unified_search\"> \
read more files</invoke> Here is my half-baked plan.";
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
None,
Some(garbled.to_string()),
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
assert!(plan_session.interview_pending(), "non-exhausted plan failure must re-force the interview");
let text = working_history
.iter()
.rev()
.find(|m| m.role == uni::MessageRole::Assistant)
.expect("a final answer must be pushed")
.content
.as_text();
assert!(
text.contains("final synthesis failed"),
"plan-mode fallback must be the structured message, not garbled salvage: {text}"
);
assert!(!text.contains("unified_search"), "garbled tool-call salvage must not leak into the plan: {text}");
}
#[tokio::test]
async fn plan_mode_recovery_keeps_partial_proposed_plan_salvage() {
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
let partial_plan =
"<proposed_plan>\n- Action: add caching -> src/cache.rs\n verify: cargo test\n</proposed_plan>";
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
None,
Some(partial_plan.to_string()),
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
let text = working_history
.iter()
.rev()
.find(|m| m.role == uni::MessageRole::Assistant)
.expect("a final answer must be pushed")
.content
.as_text();
assert!(text.contains("<proposed_plan"), "a real partial plan must be kept as the plan: {text}");
}
#[tokio::test]
async fn plan_mode_recovery_persists_salvaged_proposed_plan_to_session_file() {
use tempfile::TempDir;
use vtcode_core::tools::handlers::planning_workflow::PlanningWorkflowState;
let temp_dir = TempDir::new().unwrap();
let state = PlanningWorkflowState::new(temp_dir.path().to_path_buf());
let plan_file = state.plans_dir().join("recovered-plan.md");
state.set_plan_file(Some(plan_file.clone())).await;
let mut working_history: Vec<uni::Message> = Vec::new();
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.mark_budget_exhausted();
let salvaged = "<proposed_plan>\n- Action: add caching -> src/cache.rs\n verify: cargo test\n</proposed_plan>";
let result = complete_turn_after_failed_tool_free_recovery(
&mut working_history,
"stage",
Some(&transient_err()),
Some(salvaged.to_string()),
Some(&mut plan_session),
Some(&state),
)
.await;
assert!(matches!(result, TurnLoopResult::Completed));
let content =
std::fs::read_to_string(&plan_file).expect("salvaged plan must be persisted to the session plan file");
assert!(
content.contains("add caching"),
"salvaged plan must be written to the session plan file, got: {content}"
);
}
}