use super::post_tool_recovery::complete_turn_after_failed_tool_free_recovery;
use super::post_tool_recovery::prepare_post_tool_tool_free_recovery;
use super::post_tool_recovery::{ensure_post_tool_resume_directive, has_tool_response_since};
use super::{
HarnessUsage, PENDING_VERIFICATION_BLOCK_REASON, PLANNING_RECOVERY_SYNTHESIS_FALLBACK,
POST_TOOL_CONTEXT_COMPACTION_FAILED_REASON, POST_TOOL_RECOVERY_REASON, POST_TOOL_RECOVERY_REASON_PLAN_MODE,
POST_TOOL_RESUME_DIRECTIVE, POST_TOOL_TOOL_ENABLED_RETRY_DIRECTIVE, PostToolFailureRecovery,
RECOVERY_CONTRACT_VIOLATION_REASON, RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER, accumulate_turn_usage,
blocked_turn_final_response, count_assistant_text_responses_for_guard, count_assistant_text_responses_in_turn,
current_turn_preserve_index, ensure_blocked_turn_response, finalize_turn, has_turn_usage,
maybe_recover_after_post_tool_llm_failure, normalize_tool_free_recovery_break_outcome, run_turn_loop,
};
use std::fs;
use std::path::Path;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use crate::agent::runloop::unified::planning_workflow::recovery::{
PLANNING_SYNTHESIS_TRUNCATED_CONDENSE_DIRECTIVE, plan_synthesis_was_truncated,
};
use crate::agent::runloop::unified::planning_workflow_state::PlanningWorkflowSessionState;
use crate::agent::runloop::unified::turn::context::TurnLoopResult;
use crate::agent::runloop::unified::turn::turn_processing::test_support::TestTurnProcessingBacking;
use anyhow::anyhow;
use serde_json::json;
use vtcode_core::config::constants::tools as tool_names;
use vtcode_core::exec::events::{ThreadEvent, ThreadItemDetails, VersionedThreadEvent};
use vtcode_core::llm::provider as uni;
use vtcode_core::utils::ansi::AnsiRenderer;
use vtcode_ui::tui::app::InlineHandle;
const PENDING_VERIFICATION_RESPONSE_MARKER: &str = "Inspection-only checks do not clear the verification gate";
const CONTEXT_CAPACITY_RESPONSE_MARKER: &str = "context capacity or compaction failed";
fn final_answer_text(history: &[uni::Message]) -> String {
history
.iter()
.filter(|message| message.phase == Some(uni::AssistantPhase::FinalAnswer))
.map(|message| message.content.as_text())
.next_back()
.expect("blocked turn should retain a final assistant response")
.to_string()
}
fn assert_blocked_response_surfaces(
backing: &mut TestTurnProcessingBacking,
history: &[uni::Message],
harness_path: &Path,
response_marker: &str,
) {
let response = final_answer_text(history);
assert!(response.contains(response_marker), "unexpected blocked response: {response}");
assert_eq!(
history
.iter()
.filter(|message| message.phase == Some(uni::AssistantPhase::FinalAnswer))
.count(),
1,
"blocked recovery must append exactly one final assistant message"
);
let rendered = backing.rendered_inline_output();
assert_eq!(
rendered.matches(response_marker).count(),
1,
"blocked recovery must render the final response exactly once: {rendered}"
);
let harness = fs::read_to_string(harness_path).expect("read harness events");
let events = harness
.lines()
.map(|line| {
serde_json::from_str::<VersionedThreadEvent>(line)
.expect("blocked recovery harness output should use the versioned event contract")
.into_event()
})
.collect::<Vec<_>>();
let agent_messages = events
.iter()
.filter(|event| {
let ThreadEvent::ItemCompleted(item) = event else {
return false;
};
let ThreadItemDetails::AgentMessage(message) = &item.item.details else {
return false;
};
message.text.contains(response_marker)
})
.count();
assert_eq!(agent_messages, 1, "blocked recovery must emit one agent_message item: {harness}");
assert_eq!(
events
.iter()
.filter(|event| matches!(event, ThreadEvent::TurnFailed(_)))
.count(),
1,
"blocked recovery must emit one turn.failed event: {harness}"
);
assert!(
!events.iter().any(|event| matches!(event, ThreadEvent::TurnCompleted(_))),
"blocked recovery must not emit turn.completed: {harness}"
);
}
#[test]
fn current_turn_preserve_index_keeps_user_request_before_transient_notes() {
let history = vec![
uni::Message::assistant("older response".to_string()),
uni::Message::user("apply the requested fix".to_string()),
uni::Message::system("transient recovery note".to_string()),
];
assert_eq!(current_turn_preserve_index(&history, history.len()), 1);
}
#[test]
fn recovery_synthesis_fallback_says_no_tool_call_was_applied() {
assert!(RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER.contains("no tool call applied"));
}
#[test]
fn has_tool_response_since_detects_new_tool_message() {
let messages = vec![
uni::Message::user("run script".to_string()),
uni::Message::assistant("".to_string()),
uni::Message::tool_response("call_1".to_string(), "ok".to_string()),
];
assert!(has_tool_response_since(&messages, 1));
}
#[test]
fn has_tool_response_since_ignores_non_tool_messages() {
let messages = vec![
uni::Message::user("hello".to_string()),
uni::Message::assistant("done".to_string()),
];
assert!(!has_tool_response_since(&messages, 0));
}
#[test]
fn has_tool_response_since_handles_baseline_past_end() {
let messages = vec![uni::Message::tool_response("call_1".to_string(), "ok".to_string())];
assert!(!has_tool_response_since(&messages, 10));
}
#[test]
fn ensure_post_tool_resume_directive_is_idempotent_near_history_tail() {
let mut history = vec![
uni::Message::user("run cargo nextest".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"success\":false}".to_string()),
];
ensure_post_tool_resume_directive(&mut history);
ensure_post_tool_resume_directive(&mut history);
let directive_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RESUME_DIRECTIVE
})
.count();
assert_eq!(directive_count, 1);
}
#[test]
fn prepare_post_tool_tool_free_recovery_is_idempotent_near_history_tail() {
let mut history = vec![
uni::Message::user("summarize the existing tool outputs".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"ok\":true}".to_string()),
];
prepare_post_tool_tool_free_recovery(&mut history, POST_TOOL_RECOVERY_REASON);
prepare_post_tool_tool_free_recovery(&mut history, POST_TOOL_RECOVERY_REASON);
let resume_directive_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RESUME_DIRECTIVE
})
.count();
assert_eq!(resume_directive_count, 0);
let recovery_reason_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RECOVERY_REASON
})
.count();
assert_eq!(recovery_reason_count, 1);
}
#[test]
fn retryable_post_tool_follow_up_failure_schedules_one_tool_enabled_recovery() {
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let handle = InlineHandle::new_for_tests(tx);
let mut renderer = AnsiRenderer::with_inline_ui(handle, Default::default());
let mut history = vec![
uni::Message::user("run cargo nextest".to_string()),
uni::Message::assistant("".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"critical_note\":\"reuse output\"}".to_string()),
];
let action = maybe_recover_after_post_tool_llm_failure(
&mut renderer,
&mut history,
&anyhow!("Network error"),
2,
1,
"streaming",
true,
true,
false,
)
.expect("recovery should succeed");
assert_eq!(action, PostToolFailureRecovery::RetryToolEnabled);
let action_again = maybe_recover_after_post_tool_llm_failure(
&mut renderer,
&mut history,
&anyhow!("Network error"),
3,
1,
"streaming",
false,
false,
false,
)
.expect("repeat recovery should succeed");
assert_eq!(action_again, PostToolFailureRecovery::StopAfterDirective);
let enabled_retry_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System
&& message.content.as_text() == POST_TOOL_TOOL_ENABLED_RETRY_DIRECTIVE
})
.count();
assert_eq!(enabled_retry_count, 1);
let directive_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RESUME_DIRECTIVE
})
.count();
assert_eq!(directive_count, 1);
let recovery_reason_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RECOVERY_REASON
})
.count();
assert_eq!(recovery_reason_count, 0);
}
#[test]
fn plan_mode_recovery_uses_plan_aware_directive() {
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let handle = InlineHandle::new_for_tests(tx);
let mut renderer = AnsiRenderer::with_inline_ui(handle, Default::default());
let mut history = vec![
uni::Message::user("plan launch-time optimization".to_string()),
uni::Message::assistant("".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"ok\":true}".to_string()),
];
let action = maybe_recover_after_post_tool_llm_failure(
&mut renderer,
&mut history,
&anyhow!("Network error"),
2,
1,
"streaming",
true,
false,
true,
)
.expect("recovery should succeed");
assert_eq!(action, PostToolFailureRecovery::RetryToolFree);
let plan_directive_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RECOVERY_REASON_PLAN_MODE
})
.count();
assert_eq!(plan_directive_count, 1);
let generic_directive_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RECOVERY_REASON
})
.count();
assert_eq!(generic_directive_count, 0);
}
#[test]
fn retryable_post_tool_follow_up_failure_stops_after_recovery_pass_is_spent() {
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let handle = InlineHandle::new_for_tests(tx);
let mut renderer = AnsiRenderer::with_inline_ui(handle, Default::default());
let mut history = vec![
uni::Message::user("summarize the tool output".to_string()),
uni::Message::assistant("".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"ok\":true}".to_string()),
];
let action = maybe_recover_after_post_tool_llm_failure(
&mut renderer,
&mut history,
&anyhow!("Network error"),
2,
1,
"streaming",
false,
false,
false,
)
.expect("recovery classification should succeed");
assert_eq!(action, PostToolFailureRecovery::StopAfterDirective);
assert!(!history.iter().any(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RECOVERY_REASON
}));
assert!(history.iter().any(|message| {
message.role == uni::MessageRole::System && message.content.as_text() == POST_TOOL_RESUME_DIRECTIVE
}));
}
#[test]
fn post_tool_follow_up_failure_chain_consumes_tool_free_recovery_pass() {
use crate::agent::runloop::unified::run_loop_context::{HarnessTurnState, TurnId, TurnRunId};
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let handle = InlineHandle::new_for_tests(tx);
let mut renderer = AnsiRenderer::with_inline_ui(handle, Default::default());
let mut history = vec![
uni::Message::user("run cargo nextest".to_string()),
uni::Message::assistant("".to_string()),
uni::Message::tool_response("call_1".to_string(), "{\"critical_note\":\"reuse output\"}".to_string()),
];
let mut state = HarnessTurnState::new(TurnRunId("run-1".to_string()), TurnId("turn-1".to_string()), 4, 10, 1);
assert!(!state.is_recovery_active());
let action = maybe_recover_after_post_tool_llm_failure(
&mut renderer,
&mut history,
&anyhow!("Network error"),
2,
1,
"execute_llm_request",
true,
false,
false,
)
.expect("recovery classification should succeed");
assert_eq!(action, PostToolFailureRecovery::RetryToolFree);
assert!(state.switch_to_tool_free_recovery());
assert!(state.consume_recovery_pass(), "consume_recovery_pass must succeed after switch from Inactive");
assert!(state.recovery_is_tool_free());
}
#[tokio::test]
async fn empty_model_response_after_recovery_is_visible_and_blocked() {
let mut backing = TestTurnProcessingBacking::new(4).await;
backing.activate_tool_free_recovery_for_test("post-tool follow-up failure");
let mut history = vec![uni::Message::user("summarize the tool outputs".to_string())];
let outcome = run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("recovery should produce a visible fallback");
assert!(matches!(outcome.result, TurnLoopResult::Blocked { .. }));
assert!(outcome.final_response_was_fallback);
let final_text = history
.iter()
.rev()
.find(|message| message.role == uni::MessageRole::Assistant)
.map(|message| message.content.as_text().trim().to_string())
.unwrap_or_default();
assert!(!final_text.is_empty(), "recovery must not leave an empty final response");
}
#[test]
fn blocked_turn_final_response_explains_pending_verification() {
let response = blocked_turn_final_response(PENDING_VERIFICATION_BLOCK_REASON);
assert!(response.contains("Inspection-only checks do not clear the verification gate"));
assert!(response.contains("cargo check --locked"));
assert!(response.contains("cargo nextest run"));
}
#[test]
fn blocked_turn_final_response_explains_context_capacity_failure() {
let response =
blocked_turn_final_response(&format!("recovery failed: {POST_TOOL_CONTEXT_COMPACTION_FAILED_REASON}"));
assert!(response.contains("context capacity"));
assert!(response.contains("retained"));
assert!(response.contains("resume"));
assert!(response.contains("switch models"));
}
#[test]
fn blocked_turn_final_response_is_not_suppressed_by_prior_event_state() {
assert!(
blocked_turn_final_response(&format!("prior event; {PENDING_VERIFICATION_BLOCK_REASON}"))
.contains("Inspection-only checks")
);
}
#[test]
fn blocked_turn_final_response_has_generic_fallback() {
let response = blocked_turn_final_response("some other blocked reason");
assert!(response.contains("blocked"));
assert!(response.contains("resume"));
}
#[tokio::test]
async fn complete_turn_after_failed_tool_free_recovery_appends_fallback_once() {
let mut history = vec![uni::Message::user("summarize".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
Some(&anyhow!("Network error")),
None,
None,
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
let fallback_count = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER
})
.count();
assert_eq!(fallback_count, 1);
let outcome_again =
complete_turn_after_failed_tool_free_recovery(&mut history, "test.stage", None, None, None, None).await;
assert!(matches!(outcome_again, TurnLoopResult::Completed { .. }));
let fallback_count_again = history
.iter()
.filter(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER
})
.count();
assert_eq!(fallback_count_again, 1);
}
#[tokio::test]
async fn complete_turn_after_failed_tool_free_recovery_prefers_salvaged_prose() {
let mut history = vec![uni::Message::user("summarize".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
None,
Some("Here is the launch-time plan: reduce config IO.".to_string()),
None,
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
let last = history.last().unwrap();
assert_eq!(last.role, uni::MessageRole::Assistant);
assert_eq!(last.phase, Some(uni::AssistantPhase::FinalAnswer));
let text = last.content.as_text();
assert!(text.contains("reduce config IO"));
assert!(text != RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER);
let mut history = vec![uni::Message::user("summarize".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
None,
Some(" \n".to_string()),
None,
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert_eq!(history.last().unwrap().content.as_text(), RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER);
}
#[tokio::test]
async fn normalize_tool_free_recovery_break_outcome_converts_contract_violation_to_completed() {
let mut history = vec![uni::Message::user("summarize".to_string())];
let outcome = normalize_tool_free_recovery_break_outcome(
&mut history,
TurnLoopResult::Blocked {
reason: Some(RECOVERY_CONTRACT_VIOLATION_REASON.to_string()),
},
true,
None,
None,
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert!(history.iter().any(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER
}));
}
#[tokio::test]
async fn normalize_tool_free_recovery_break_outcome_keeps_non_recovery_blocked_result() {
let mut history = vec![uni::Message::user("summarize".to_string())];
let outcome = normalize_tool_free_recovery_break_outcome(
&mut history,
TurnLoopResult::Blocked {
reason: Some("Stopped after reaching budget limit.".to_string()),
},
true,
None,
None,
None,
)
.await;
assert!(matches!(
outcome,
TurnLoopResult::Blocked {
reason: Some(ref reason)
} if reason == "Stopped after reaching budget limit."
));
assert!(!history.iter().any(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER
}));
}
#[tokio::test]
async fn plan_mode_recovery_fallback_marks_interview_pending_and_preserves_research() {
use vtcode_core::core::interfaces::session::PlanningEntrySource;
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.enter(PlanningEntrySource::UserRequest);
assert!(!plan_session.interview_pending());
let mut history = vec![uni::Message::user("plan launch-time optimization".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
Some(&anyhow!("Network error")),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert!(plan_session.interview_pending());
assert!(history.iter().any(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text().contains(PLANNING_RECOVERY_SYNTHESIS_FALLBACK)
}));
assert!(!history.iter().any(|message| {
message.role == uni::MessageRole::Assistant
&& message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message.content.as_text() == RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER
}));
}
#[tokio::test]
async fn plan_mode_recovery_exhausted_finalizes_instead_of_reforcing_interview() {
use vtcode_core::core::interfaces::session::PlanningEntrySource;
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.enter(PlanningEntrySource::UserRequest);
plan_session.mark_recovery_exhausted();
assert!(!plan_session.interview_pending());
let mut history = vec![uni::Message::user("plan launch-time optimization".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
Some(&anyhow!("context length exceeded")),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert!(!plan_session.interview_pending());
let last = history.last().unwrap();
assert_eq!(last.role, uni::MessageRole::Assistant);
assert_eq!(last.phase, Some(uni::AssistantPhase::FinalAnswer));
let text = last.content.as_text();
assert!(text.contains("Plan synthesis failed after repeated recovery attempts"));
assert!(
!text.contains("Do NOT attempt more tool calls"),
"model directive must not leak into the user-visible final answer"
);
assert!(!text.contains("`implement`"), "no approval hint is allowed without a valid draft");
assert!(text.to_ascii_lowercase().contains("keep planning"));
}
#[tokio::test]
async fn plan_mode_recovery_rejects_non_plan_salvage() {
use vtcode_core::core::interfaces::session::PlanningEntrySource;
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.enter(PlanningEntrySource::UserRequest);
let mut history = vec![uni::Message::user("plan launch-time optimization".to_string())];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
None,
Some("Partial plan: batch config reads.".to_string()),
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert!(plan_session.interview_pending());
let last = history.last().unwrap();
assert!(last.content.as_text().contains("final synthesis failed"));
assert!(!last.content.as_text().contains("batch config reads"));
}
#[tokio::test]
async fn plan_mode_recovery_fallback_lists_files_read_when_present() {
use vtcode_core::core::interfaces::session::PlanningEntrySource;
let mut plan_session = PlanningWorkflowSessionState::default();
plan_session.enter(PlanningEntrySource::UserRequest);
let mut history = vec![
uni::Message::user("plan launch-time optimization".to_string()),
uni::Message::tool_response(
"call_1".to_string(),
"{\"path\": \"src/main.rs\", \"content\": \"...\"}".to_string(),
),
uni::Message::tool_response(
"call_2".to_string(),
"{\"path\": \"src/startup/mod.rs\", \"content\": \"...\"}".to_string(),
),
];
let outcome = complete_turn_after_failed_tool_free_recovery(
&mut history,
"test.stage",
Some(&anyhow!("Network error")),
None,
Some(&mut plan_session),
None,
)
.await;
assert!(matches!(outcome, TurnLoopResult::Completed { .. }));
assert!(plan_session.interview_pending());
let last = history.last().unwrap();
let text = last.content.as_text();
assert!(text.contains("Files already read this turn"));
assert!(text.contains("src/main.rs"));
assert!(text.contains("src/startup/mod.rs"));
assert!(text.contains(PLANNING_RECOVERY_SYNTHESIS_FALLBACK));
}
#[test]
fn accumulate_turn_usage_merges_prompt_completion_and_cached_tokens() {
let mut total = HarnessUsage::default();
accumulate_turn_usage(
"openai",
&mut total,
&Some(uni::Usage {
prompt_tokens: 100,
completion_tokens: 20,
total_tokens: 120,
cached_prompt_tokens: Some(15),
cache_creation_tokens: None,
cache_read_tokens: Some(15),
iterations: None,
}),
);
accumulate_turn_usage(
"openai",
&mut total,
&Some(uni::Usage {
prompt_tokens: 40,
completion_tokens: 10,
total_tokens: 50,
cached_prompt_tokens: None,
cache_creation_tokens: None,
cache_read_tokens: None,
iterations: None,
}),
);
assert_eq!(total.input_tokens, 140);
assert_eq!(total.cached_input_tokens, 15);
assert_eq!(total.output_tokens, 30);
assert!(has_turn_usage(&total));
}
#[test]
fn accumulate_turn_usage_normalizes_anthropic_exclusive_input() {
let mut total = HarnessUsage::default();
accumulate_turn_usage(
"anthropic",
&mut total,
&Some(uni::Usage {
prompt_tokens: 100,
completion_tokens: 20,
total_tokens: 120,
cached_prompt_tokens: None,
cache_creation_tokens: Some(50),
cache_read_tokens: Some(400),
iterations: None,
}),
);
assert_eq!(total.input_tokens, 550);
assert_eq!(total.cached_input_tokens, 400);
assert_eq!(total.cache_creation_tokens, 50);
assert_eq!(total.output_tokens, 20);
}
#[tokio::test]
async fn turn_loop_preserves_legacy_loop_detector_state() {
let mut backing = TestTurnProcessingBacking::new(4).await;
backing.set_loop_limit(tool_names::READ_FILE, 2);
let seeded_args = json!({"path":"sample.txt"});
assert!(backing.record_tool_call(tool_names::READ_FILE, &seeded_args).is_none());
let _ = backing.record_tool_call(tool_names::READ_FILE, &seeded_args);
let warning = backing.record_tool_call(tool_names::READ_FILE, &seeded_args);
assert!(warning.is_some());
assert!(backing.is_hard_limit_exceeded(tool_names::READ_FILE));
let mut history = vec![uni::Message::user("continue".to_string())];
run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("turn loop should complete");
assert!(backing.is_hard_limit_exceeded(tool_names::READ_FILE));
}
#[tokio::test]
async fn anti_blind_guard_stops_outer_loop_after_two_pending_stale_plan_pause_responses() {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone)]
struct RepeatedTextAfterMutationsProvider {
requests: Arc<AtomicUsize>,
text_responses: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for RepeatedTextAfterMutationsProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
let request_number = self.requests.fetch_add(1, Ordering::SeqCst);
let response = if request_number < 4 {
let path = format!("anti-blind-regression-{request_number}.txt");
let patch =
format!("*** Begin Patch\n*** Add File: {path}\n+mutation {request_number}\n*** End Patch\n");
uni::LLMResponse {
content: None,
model: request.model.clone(),
tool_calls: Some(vec![uni::ToolCall::function(
format!("mutation-{request_number}"),
tool_names::APPLY_PATCH.to_string(),
json!({"patch": patch}).to_string(),
)]),
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
}
} else {
self.text_responses.fetch_add(1, Ordering::SeqCst);
uni::LLMResponse {
content: Some(
"Implementation is paused because tool use is disabled. Wait for the next turn.".to_string(),
),
model: request.model.clone(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
}
};
Ok(response)
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let requests = Arc::new(AtomicUsize::new(0));
let text_responses = Arc::new(AtomicUsize::new(0));
let mut backing = TestTurnProcessingBacking::new(8).await;
backing.set_provider(Box::new(RepeatedTextAfterMutationsProvider {
requests: requests.clone(),
text_responses: text_responses.clone(),
}));
let mut history = vec![uni::Message::user("apply and verify the requested change".to_string())];
let turn_context = backing.turn_loop_context();
turn_context.harness_state.set_approved_plan_execution(true);
let outcome = run_turn_loop(&mut history, turn_context)
.await
.expect("turn loop should stop at the pending-verification response cap");
assert!(matches!(outcome.result, TurnLoopResult::Blocked { .. }));
assert!(matches!(
outcome.result,
TurnLoopResult::Blocked {
reason: Some(ref reason)
} if reason == "Turn blocked after repeated unverified assistant responses; verification is still pending."
));
assert_eq!(text_responses.load(Ordering::SeqCst), 2);
assert_eq!(requests.load(Ordering::SeqCst), 6, "the provider must not receive a third pending text request");
assert!(outcome.final_response_was_fallback);
assert!(history.iter().any(|message| {
message.phase == Some(uni::AssistantPhase::FinalAnswer)
&& message
.content
.as_text()
.contains("Inspection-only checks do not clear the verification gate")
}));
assert_eq!(
history
.iter()
.filter(|message| message.phase == Some(uni::AssistantPhase::FinalAnswer))
.count(),
1
);
assert!(!history.iter().any(|message| {
message
.content
.as_text()
.contains("Implementation is paused because tool use is disabled.")
}));
}
#[tokio::test]
async fn blocked_anti_blind_recovery_publishes_one_actionable_handoff() {
#[derive(Clone)]
struct VerificationRecoveryProvider {
requests: Arc<AtomicUsize>,
steps: Arc<Mutex<Vec<String>>>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for VerificationRecoveryProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
let request_number = self.requests.fetch_add(1, Ordering::SeqCst);
let tool_call = |label: &str, tool_name: &str, args: serde_json::Value| {
self.steps.lock().expect("step trace lock").push(label.to_string());
uni::LLMResponse {
content: None,
model: request.model.clone(),
tool_calls: Some(vec![uni::ToolCall::function(
format!("verification-{request_number}"),
tool_name.to_string(),
args.to_string(),
)]),
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
}
};
let patch = |path: &str, contents: &str| {
format!("*** Begin Patch\n*** Add File: {path}\n+{contents}\n*** End Patch\n")
};
let response = match request_number {
0..=2 => tool_call(
"successful_edit",
tool_names::APPLY_PATCH,
json!({"patch": patch(&format!("anti-blind-sequence-{request_number}.txt"), "effective edit")}),
),
3 => tool_call(
"failed_patch",
tool_names::APPLY_PATCH,
json!({
"patch": "*** Begin Patch\n*** Update File: missing-target.txt\n@@\n-old\n+new\n*** End Patch\n"
}),
),
4 => tool_call(
"inspection",
tool_names::EXEC_COMMAND,
json!({"cmd": "rg -n 'effective edit' . || true"}),
),
5 => tool_call(
"link_check",
tool_names::EXEC_COMMAND,
json!({"cmd": "rg -n '\\[[^]]+\\]\\([^)]*\\)' . || true"}),
),
6 => tool_call("diff_check", tool_names::EXEC_COMMAND, json!({"cmd": "git diff --check"})),
7 => tool_call(
"successful_edit",
tool_names::APPLY_PATCH,
json!({"patch": patch("anti-blind-sequence-final.txt", "last effective edit")}),
),
_ => {
self.steps.lock().expect("step trace lock").push("unverified_text".to_string());
uni::LLMResponse {
content: Some("The edits are complete, but verification was not run.".to_string()),
model: request.model.clone(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
}
}
};
Ok(response)
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let requests = Arc::new(AtomicUsize::new(0));
let steps = Arc::new(Mutex::new(Vec::new()));
let mut backing = TestTurnProcessingBacking::new(16).await;
let harness_path = backing.enable_harness_emitter();
backing.set_provider(Box::new(VerificationRecoveryProvider { requests: requests.clone(), steps: steps.clone() }));
let mut history = vec![uni::Message::user("apply the change and verify it".to_string())];
let turn_context = backing.turn_loop_context();
turn_context.harness_state.set_approved_plan_execution(true);
let outcome = run_turn_loop(&mut history, turn_context)
.await
.expect("anti-blind recovery should return a blocked outcome");
assert!(matches!(
outcome.result,
TurnLoopResult::Blocked {
reason: Some(ref reason)
} if reason == PENDING_VERIFICATION_BLOCK_REASON
));
assert_eq!(requests.load(Ordering::SeqCst), 10, "two pending text responses are the terminal cap");
assert_eq!(
steps.lock().expect("step trace lock").as_slice(),
[
"successful_edit",
"successful_edit",
"successful_edit",
"failed_patch",
"inspection",
"link_check",
"diff_check",
"successful_edit",
"unverified_text",
"unverified_text",
]
);
assert_blocked_response_surfaces(&mut backing, &history, &harness_path, PENDING_VERIFICATION_RESPONSE_MARKER);
}
#[tokio::test]
async fn context_capacity_blocked_recovery_publishes_one_actionable_handoff() {
#[derive(Clone)]
struct ContextCapacityProvider {
requests: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for ContextCapacityProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
if self.requests.fetch_add(1, Ordering::SeqCst) == 0 {
let patch =
"*** Begin Patch\n*** Add File: context-capacity-sequence.txt\n+effective edit\n*** End Patch\n";
return Ok(uni::LLMResponse {
content: None,
model: request.model,
tool_calls: Some(vec![uni::ToolCall::function(
"context-capacity-edit".to_string(),
tool_names::APPLY_PATCH.to_string(),
json!({"patch": patch}).to_string(),
)]),
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
});
}
Err(uni::LLMError::InvalidRequest {
message: "maximum context length is 114688 tokens".to_string(),
metadata: None,
})
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let requests = Arc::new(AtomicUsize::new(0));
let mut backing = TestTurnProcessingBacking::new(8).await;
let harness_path = backing.enable_harness_emitter();
backing.set_provider(Box::new(ContextCapacityProvider { requests: requests.clone() }));
let mut history = vec![uni::Message::user("apply the change".to_string())];
let outcome = run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("context-capacity recovery should return a blocked outcome");
assert!(matches!(
outcome.result,
TurnLoopResult::Blocked {
reason: Some(ref reason)
} if reason == POST_TOOL_CONTEXT_COMPACTION_FAILED_REASON
));
assert_eq!(requests.load(Ordering::SeqCst), 2, "context failure should stop after the bounded retry path");
assert!(history.iter().any(|message| {
message.role == uni::MessageRole::System && message.content.as_text().contains(POST_TOOL_RESUME_DIRECTIVE)
}));
assert_blocked_response_surfaces(&mut backing, &history, &harness_path, CONTEXT_CAPACITY_RESPONSE_MARKER);
}
#[tokio::test]
async fn blocked_recovery_does_not_duplicate_prior_harness_agent_message() {
let prior_text = "A streamed progress message was already published.";
let mut backing = TestTurnProcessingBacking::new(4).await;
let harness_path = backing.enable_harness_emitter();
backing.emit_harness_assistant_message_for_test(prior_text);
let mut history = vec![uni::Message::user("resume the request".to_string())];
let blocked = TurnLoopResult::Blocked {
reason: Some(PENDING_VERIFICATION_BLOCK_REASON.to_string()),
};
{
let mut context = backing.turn_loop_context();
context.harness_state.mark_final_response_event_emitted();
ensure_blocked_turn_response(&mut context, &mut history, 1, PENDING_VERIFICATION_BLOCK_REASON)
.expect("blocked recovery handoff");
finalize_turn(&mut context, &history, &blocked, &HarnessUsage::default()).await;
}
assert_eq!(
history
.iter()
.filter(|message| message.phase == Some(uni::AssistantPhase::FinalAnswer))
.count(),
1
);
let final_text = final_answer_text(&history);
assert!(final_text.contains(PENDING_VERIFICATION_RESPONSE_MARKER));
let rendered = backing.rendered_inline_output();
assert_eq!(rendered.matches(PENDING_VERIFICATION_RESPONSE_MARKER).count(), 1);
let harness = fs::read_to_string(harness_path).expect("read harness events");
let events = harness
.lines()
.map(|line| {
serde_json::from_str::<VersionedThreadEvent>(line)
.expect("versioned harness event")
.into_event()
})
.collect::<Vec<_>>();
let agent_messages = events
.iter()
.filter_map(|event| {
let ThreadEvent::ItemCompleted(item) = event else {
return None;
};
let ThreadItemDetails::AgentMessage(message) = &item.item.details else {
return None;
};
Some(message.text.as_str())
})
.collect::<Vec<_>>();
assert_eq!(agent_messages, vec![prior_text]);
assert_eq!(
events
.iter()
.filter(|event| matches!(event, ThreadEvent::TurnFailed(_)))
.count(),
1
);
assert!(!events.iter().any(|event| matches!(event, ThreadEvent::TurnCompleted(_))));
}
#[tokio::test]
async fn stale_plan_pause_without_mutations_consumes_text_response_budget() {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone)]
struct RepeatedStalePauseProvider {
requests: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for RepeatedStalePauseProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
self.requests.fetch_add(1, Ordering::SeqCst);
Ok(uni::LLMResponse {
content: Some(
"Implementation is paused because tool use is disabled. Wait for the next turn.".to_string(),
),
model: request.model,
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
})
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let requests = Arc::new(AtomicUsize::new(0));
let mut backing = TestTurnProcessingBacking::new(8).await;
backing.set_provider(Box::new(RepeatedStalePauseProvider { requests: requests.clone() }));
let mut history = vec![uni::Message::user("continue the approved implementation".to_string())];
let turn_context = backing.turn_loop_context();
turn_context.harness_state.set_approved_plan_execution(true);
let outcome = run_turn_loop(&mut history, turn_context)
.await
.expect("turn loop should stop at the discarded text response cap");
assert!(matches!(outcome.result, TurnLoopResult::Blocked { .. }));
assert!(outcome.final_response_was_fallback);
assert_eq!(requests.load(Ordering::SeqCst), 2);
assert!(!history.iter().any(|message| {
message
.content
.as_text()
.contains("Implementation is paused because tool use is disabled.")
}));
}
#[test]
fn count_assistant_text_responses_in_turn_zero_for_empty_history() {
let history: Vec<uni::Message> = Vec::new();
assert_eq!(count_assistant_text_responses_in_turn(&history, 0), 0);
}
#[test]
#[allow(
clippy::vec_init_then_push,
reason = "Intentional compatibility, platform, or test-only suppression."
)]
fn count_assistant_text_responses_in_turn_skips_tool_call_messages() {
let mut history: Vec<uni::Message> = Vec::new();
history.push(uni::Message::assistant_with_tools(
String::new(),
vec![uni::ToolCall::function(
"tool_call_0".to_string(),
"code_search".to_string(),
"{}".to_string(),
)],
));
history.push(uni::Message::assistant("Functions and structs.".to_string()));
history.push(uni::Message::system("Tools disabled.".to_string()));
history.push(uni::Message::assistant("Functions and structs again.".to_string()));
history.push(uni::Message::assistant(String::new()));
history.push(uni::Message::assistant(" \n ".to_string()));
assert_eq!(count_assistant_text_responses_in_turn(&history, 0), 2);
}
#[test]
fn count_assistant_text_responses_in_turn_ignores_history_before_baseline() {
let mut history = vec![
uni::Message::user("previous request".to_string()),
uni::Message::assistant("previous answer one".to_string()),
uni::Message::assistant("previous answer two".to_string()),
];
let turn_history_start_len = history.len();
assert_eq!(
count_assistant_text_responses_in_turn(&history, turn_history_start_len),
0,
"historical assistant text before the current turn must not count"
);
assert_eq!(
count_assistant_text_responses_for_guard(&history, turn_history_start_len, 0),
0,
"the guard must not count historical assistant text when the per-turn floor is empty"
);
history.push(uni::Message::assistant("current answer one".to_string()));
assert_eq!(count_assistant_text_responses_in_turn(&history, turn_history_start_len), 1);
history.push(uni::Message::assistant_with_tools(
String::new(),
vec![uni::ToolCall::function(
"tool_call_0".to_string(),
"code_search".to_string(),
"{}".to_string(),
)],
));
history.push(uni::Message::assistant("current answer two".to_string()));
assert_eq!(
count_assistant_text_responses_in_turn(&history, turn_history_start_len),
super::MAX_ASSISTANT_TEXT_RESPONSES_PER_TURN,
"current-turn assistant text after the baseline still trips the cap"
);
}
#[test]
fn count_assistant_text_responses_in_turn_counts_after_compaction_rebase() {
let stale_turn_history_start_len = 5;
let mut compacted_history = vec![
uni::Message::system("Compacted history summary.".to_string()),
uni::Message::user("current request".to_string()),
];
let rebased_turn_history_start_len = compacted_history.len();
compacted_history.push(uni::Message::assistant("current answer after compaction".to_string()));
assert_eq!(
count_assistant_text_responses_in_turn(&compacted_history, stale_turn_history_start_len),
0,
"a stale pre-compaction baseline misses newly appended assistant text"
);
assert_eq!(
count_assistant_text_responses_in_turn(&compacted_history, rebased_turn_history_start_len),
1,
"rebasing to the compacted length counts current-turn assistant text promptly"
);
}
#[test]
fn count_assistant_text_responses_for_guard_preserves_pre_compaction_turn_floor() {
let compacted_history = vec![
uni::Message::system("Compacted history summary.".to_string()),
uni::Message::user("current request".to_string()),
];
let rebased_turn_history_start_len = compacted_history.len();
let recorded_text_responses_in_turn = 1;
assert_eq!(
count_assistant_text_responses_in_turn(&compacted_history, rebased_turn_history_start_len),
0,
"history slice cannot see same-turn assistant text removed by compaction"
);
assert_eq!(
count_assistant_text_responses_for_guard(
&compacted_history,
rebased_turn_history_start_len,
recorded_text_responses_in_turn,
),
recorded_text_responses_in_turn,
"guard uses the per-turn counter as a compaction-safe floor"
);
}
#[test]
fn count_assistant_text_responses_for_guard_counts_post_compaction_growth_above_floor() {
let mut compacted_history = vec![
uni::Message::system("Compacted history summary.".to_string()),
uni::Message::user("current request".to_string()),
];
let rebased_turn_history_start_len = compacted_history.len();
let recorded_text_responses_in_turn = 1;
compacted_history.push(uni::Message::assistant("current answer after compaction".to_string()));
compacted_history.push(uni::Message::assistant("second current answer after compaction".to_string()));
assert_eq!(
count_assistant_text_responses_for_guard(
&compacted_history,
rebased_turn_history_start_len,
recorded_text_responses_in_turn,
),
2,
"new assistant text appended after compaction still counts promptly"
);
}
#[test]
fn count_assistant_text_responses_in_turn_matches_observed_pattern() {
let mut history: Vec<uni::Message> = Vec::new();
for _ in 0..4 {
history.push(uni::Message::assistant(
"# Functions and Structs in crates/codegen/vtcode-core/src/tools/registry\n\
\n\
The directory has 70 files, 23 structs, 69 functions, 11 enums.\n\
\n\
## Structs (23)\n\
\n\
| File | Struct |\n\
|---|---|\n\
| mod.rs | ToolRegistry |\n\
| distributed.rs | ToolConfigSnapshot |\n\
... (and many more rows)\n"
.to_string(),
));
}
assert_eq!(count_assistant_text_responses_in_turn(&history, 0), 4);
assert!(
count_assistant_text_responses_in_turn(&history, 0) >= super::MAX_ASSISTANT_TEXT_RESPONSES_PER_TURN,
"anti-runaway guard would trip on this history"
);
}
#[tokio::test]
async fn tool_free_recovery_retries_on_contract_violation_then_salvages() {
use std::sync::{Arc, Mutex};
#[derive(Clone)]
struct ContractViolationProvider {
requests: Arc<Mutex<usize>>,
content: String,
}
#[async_trait::async_trait]
impl uni::LLMProvider for ContractViolationProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
*self.requests.lock().expect("requests lock") += 1;
Ok(uni::LLMResponse {
content: Some(self.content.clone()),
model: request.model.clone(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
})
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let mut backing = TestTurnProcessingBacking::new(4).await;
backing.activate_tool_free_recovery_for_test("post-tool follow-up failure");
let markup = "Here is my plan: the change was not applied because tools were disabled. \
</tool_call> Please re-run with tools enabled.";
let requests = Arc::new(Mutex::new(0usize));
backing.set_provider(Box::new(ContractViolationProvider {
requests: requests.clone(),
content: markup.to_string(),
}));
let mut history = vec![uni::Message::user("summarize the tool outputs".to_string())];
run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("turn loop should complete after recovery retries");
assert_eq!(
*requests.lock().expect("requests lock"),
super::MAX_RECOVERY_RETRIES as usize + 1,
"recovery must retry exactly MAX_RECOVERY_RETRIES times before falling back"
);
let final_text = history
.iter()
.rev()
.find(|m| m.role == uni::MessageRole::Assistant)
.map(|m| m.content.as_text().to_string())
.unwrap_or_default();
assert!(final_text.contains("Here is my plan:"), "expected salvaged prose, got: {final_text}");
assert!(
!final_text.contains(RECOVERY_SYNTHESIS_FALLBACK_FINAL_ANSWER),
"must not emit canned fallback when salvage is available"
);
assert!(backing.recovery_is_tool_free());
}
#[test]
fn plan_synthesis_truncated_detects_unclosed_proposed_plan() {
let truncated = uni::LLMResponse {
content: Some("<proposed_plan>\n# Improve launch time\n## Steps\n1. Fix warmup -> src/main.rs".to_string()),
model: "noop".to_string(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Length,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
};
assert!(
plan_synthesis_was_truncated(&truncated),
"unclosed <proposed_plan> with Length finish must be detected as truncated"
);
let complete = uni::LLMResponse {
content: Some("<proposed_plan>\n# Title\n## Steps\n1. x\n</proposed_plan>".to_string()),
model: "noop".to_string(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Length,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
};
assert!(!plan_synthesis_was_truncated(&complete), "closed <proposed_plan> must not be flagged as truncated");
let normal = uni::LLMResponse {
content: Some("<proposed_plan>\n# Title\n</proposed_plan>".to_string()),
model: "noop".to_string(),
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
};
assert!(!plan_synthesis_was_truncated(&normal), "Stop-finished plan must not be flagged as truncated");
}
#[tokio::test]
async fn planning_synthesis_truncated_retries_with_compact_spec() {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone)]
struct TruncateThenCompactProvider {
calls: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for TruncateThenCompactProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
let n = self.calls.fetch_add(1, Ordering::SeqCst);
let (content, finish_reason) = if n == 0 {
(
"<proposed_plan>\n# Improve launch time\n## Summary\nMake cold start faster.\n## Steps\n1. Fix warmup -> src/main.rs -> verify: build".to_string(),
uni::FinishReason::Length,
)
} else {
(
"Plan condensed: warmup path in src/main.rs fixed; rebuild to verify.".to_string(),
uni::FinishReason::Stop,
)
};
Ok(uni::LLMResponse {
content: Some(content),
model: request.model.clone(),
tool_calls: None,
usage: None,
finish_reason,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
})
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let calls = Arc::new(AtomicUsize::new(0));
let mut backing = TestTurnProcessingBacking::new(4).await;
backing.activate_planning_for_test();
backing.set_provider(Box::new(TruncateThenCompactProvider { calls: calls.clone() }));
let mut history = vec![uni::Message::user("make a plan to improve launch time".to_string())];
run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("turn loop must complete after condensing the truncated plan");
assert_eq!(calls.load(Ordering::SeqCst), 2, "must re-run synthesis exactly once after truncation, not loop");
assert!(
history.iter().any(|message| message
.content
.as_text()
.contains(PLANNING_SYNTHESIS_TRUNCATED_CONDENSE_DIRECTIVE)),
"condense directive must be injected after a truncated plan"
);
let final_text = history
.iter()
.rev()
.find(|message| message.role == uni::MessageRole::Assistant)
.map(|message| message.content.as_text().to_string())
.unwrap_or_default();
assert!(final_text.contains("Plan condensed:"), "final answer must be the compact retry, got: {final_text}");
assert!(
!final_text.contains("Fix warmup -> src/main.rs -> verify: build"),
"final answer must not be the truncated draft"
);
}
#[tokio::test]
async fn approval_input_without_plan_synthesizes_before_approval() {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone)]
struct PlanProvider {
calls: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl uni::LLMProvider for PlanProvider {
fn name(&self) -> &str {
"openai"
}
fn supports_streaming(&self) -> bool {
false
}
async fn generate(&self, request: uni::LLMRequest) -> Result<uni::LLMResponse, uni::LLMError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(uni::LLMResponse {
content: Some(
"<proposed_plan>\nSummary: optimize one startup subsystem.\n1. Measure -> src/main.rs -> verify: cargo check --locked\nValidation: run the focused startup test.\nAssumptions: preserve public APIs.\n</proposed_plan>"
.to_string(),
),
model: request.model,
tool_calls: None,
usage: None,
finish_reason: uni::FinishReason::Stop,
reasoning: None,
reasoning_details: None,
organization_id: None,
request_id: None,
tool_references: Vec::new(),
compaction: None,
})
}
fn supported_models(&self) -> Vec<String> {
vec!["noop-model".to_string()]
}
fn validate_request(&self, _request: &uni::LLMRequest) -> Result<(), uni::LLMError> {
Ok(())
}
}
let calls = Arc::new(AtomicUsize::new(0));
let mut backing = TestTurnProcessingBacking::new(4).await;
backing.activate_planning_for_test();
backing.set_provider(Box::new(PlanProvider { calls: calls.clone() }));
let mut history = vec![uni::Message::user("yes".to_string())];
run_turn_loop(&mut history, backing.turn_loop_context())
.await
.expect("missing-plan approval should continue to synthesis");
assert_eq!(calls.load(Ordering::SeqCst), 1, "the missing-plan path must synthesize exactly once");
assert!(
history
.iter()
.any(|message| { message.content.as_text().contains("no completed plan draft exists yet") })
);
}