#![allow(dead_code)]
use crate::agent::AgentRunner;
use crate::error::{OrchestratorError, Result};
use crate::history::{AcceptanceAttempt, OutputCollector};
use crate::openspec::Change;
use tracing::{info, warn};
use super::output::OutputHandler;
const ACCEPTANCE_OUTPUT_FALLBACK: &str = "No acceptance output captured";
pub const MAX_ACCEPTANCE_RETRY_CYCLES: u32 = 10;
pub const MISSING_VERDICT_DIAGNOSTIC: &str = "Missing acceptance verdict: acceptance command \
exited without emitting a canonical verdict (protocol failure; status-only or waiting \
output is not a verdict)";
pub const MAX_MISSING_VERDICT_RETRIES: u32 = 2;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MissingVerdictRetry {
pub attempt: u32,
pub max: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MissingVerdictRetryDecision {
Retry(MissingVerdictRetry),
Exhausted { attempts: u32, max: u32 },
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct MissingVerdictRetryState {
consecutive: u32,
}
impl MissingVerdictRetryState {
pub fn consecutive(&self) -> u32 {
self.consecutive
}
pub fn record_canonical_verdict(&mut self) {
self.consecutive = 0;
}
pub fn record_missing_verdict(&mut self) -> MissingVerdictRetryDecision {
self.consecutive = self.consecutive.saturating_add(1);
if self.consecutive <= MAX_MISSING_VERDICT_RETRIES {
MissingVerdictRetryDecision::Retry(MissingVerdictRetry {
attempt: self.consecutive,
max: MAX_MISSING_VERDICT_RETRIES,
})
} else {
MissingVerdictRetryDecision::Exhausted {
attempts: self.consecutive,
max: MAX_MISSING_VERDICT_RETRIES,
}
}
}
}
fn bounded_missing_verdict_evidence(findings: &[String]) -> String {
let evidence = findings
.iter()
.take(5)
.cloned()
.collect::<Vec<_>>()
.join(" | ");
if evidence.is_empty() {
"no acceptance output captured".to_string()
} else {
evidence
}
}
pub fn missing_verdict_retry_progress(retry: MissingVerdictRetry, findings: &[String]) -> String {
format!(
"Acceptance completed without a canonical verdict; retrying acceptance \
(protocol retry {}/{}). Evidence: {}",
retry.attempt,
retry.max,
bounded_missing_verdict_evidence(findings)
)
}
pub fn missing_verdict_exhausted_error(attempts: u32, max: u32, findings: &[String]) -> String {
format!(
"Acceptance completed without a canonical verdict (missing-verdict protocol failure); \
status-only or waiting output is not a verdict. Exhausted {attempts} consecutive \
attempts after {max} protocol retries. Evidence: {}",
bounded_missing_verdict_evidence(findings)
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MissingVerdictRetryStep {
Retry {
retry: MissingVerdictRetry,
progress: String,
},
Exhausted { error: String },
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct MissingVerdictRetryDriver {
state: MissingVerdictRetryState,
pending: Option<MissingVerdictRetry>,
}
impl MissingVerdictRetryDriver {
pub fn take_protocol_retry(&mut self) -> Option<MissingVerdictRetry> {
self.pending.take()
}
pub fn consecutive_missing_verdicts(&self) -> u32 {
self.state.consecutive()
}
pub fn observe_canonical_verdict(&mut self) {
self.state.record_canonical_verdict();
self.pending = None;
}
pub fn observe_missing_verdict(&mut self, findings: &[String]) -> MissingVerdictRetryStep {
match self.state.record_missing_verdict() {
MissingVerdictRetryDecision::Retry(retry) => {
self.pending = Some(retry);
MissingVerdictRetryStep::Retry {
retry,
progress: missing_verdict_retry_progress(retry, findings),
}
}
MissingVerdictRetryDecision::Exhausted { attempts, max } => {
self.pending = None;
MissingVerdictRetryStep::Exhausted {
error: missing_verdict_exhausted_error(attempts, max, findings),
}
}
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NormalizedFinding {
pub identity: String,
pub text: String,
pub external: bool,
}
pub fn normalize_findings(findings: &[String]) -> Vec<NormalizedFinding> {
fn rule_kind(text: &str) -> &'static str {
if ["test", "coverage", "verification", "evidence"]
.iter()
.any(|word| text.contains(word))
{
"verification"
} else if ["spec", "proposal", "requirement"]
.iter()
.any(|word| text.contains(word))
{
"specification"
} else if ["task", "checklist", "truthful"]
.iter()
.any(|word| text.contains(word))
{
"task-truthfulness"
} else if text.contains("dirty working tree") {
"workspace-cleanliness"
} else {
"implementation"
}
}
let mut normalized = findings
.iter()
.filter_map(|finding| {
let normalized = finding.split_whitespace().collect::<Vec<_>>().join(" ");
(!normalized.is_empty()).then(|| {
let lower = normalized.to_ascii_lowercase();
let finding_code = lower
.split_whitespace()
.next()
.filter(|word| word.starts_with('[') && word.ends_with(']'));
let path_token = lower
.split_whitespace()
.find(|word| {
word.contains('/') || word.ends_with(".rs") || word.ends_with(".md")
})
.unwrap_or("");
let path = path_token
.trim_matches(|character: char| {
matches!(character, '`' | '(' | ')' | '[' | ']' | ',' | '.' | ';')
})
.split(':')
.next()
.unwrap_or("");
let external = path.is_empty()
&& !lower.contains("fix ")
&& !lower.contains("repair ")
&& [
"external non-mockable",
"non-mockable external",
"external prerequisite",
"external service outage",
"missing non-mockable external credential",
]
.iter()
.any(|needle| lower.contains(needle));
let scope = if external { "external" } else { "repository" };
NormalizedFinding {
identity: finding_code.map_or_else(
|| {
let location = if path.is_empty() {
lower.as_str()
} else {
path
};
format!("{scope}|{location}|{}", rule_kind(&lower))
},
|code| format!("{scope}|code|{code}"),
),
text: normalized,
external,
}
})
})
.collect::<Vec<_>>();
normalized.sort_by(|left, right| left.identity.cmp(&right.identity));
normalized.dedup_by(|left, right| left.identity == right.identity);
normalized
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceRetryDecision {
Retry {
reason: &'static str,
},
Stall {
reason: &'static str,
external_blockers: Vec<String>,
},
}
pub fn repository_findings(findings: &[String]) -> Vec<String> {
findings
.iter()
.filter(|finding| {
normalize_findings(std::slice::from_ref(*finding))
.first()
.is_some_and(|normalized| !normalized.external)
})
.cloned()
.collect()
}
pub fn semantic_progress_fingerprint(workspace: &std::path::Path) -> std::io::Result<String> {
fn include(path: &str) -> bool {
!path.starts_with(".git/")
&& !path.starts_with(".cflx/")
&& !path.contains("/APPLY_BLOCKED/")
&& !path.starts_with("logs/")
&& !path.starts_with("history/")
&& (path.starts_with("src/")
|| path.starts_with("tests/")
|| path.starts_with("config/")
|| path.starts_with("openspec/specs/")
|| path.contains("/specs/")
|| path == ".cflx.jsonc"
|| path.ends_with("/.cflx.jsonc")
|| path.ends_with("Cargo.toml")
|| path.ends_with("tasks.md"))
}
fn visit(
root: &std::path::Path,
directory: &std::path::Path,
output: &mut Vec<(String, Vec<u8>)>,
) -> std::io::Result<()> {
for entry in std::fs::read_dir(directory)? {
let entry = entry?;
let path = entry.path();
if path.is_dir() {
visit(root, &path, output)?;
continue;
}
let relative = path
.strip_prefix(root)
.unwrap()
.to_string_lossy()
.replace('\\', "/");
if include(&relative) {
let mut contents = std::fs::read(path)?;
if relative.ends_with("tasks.md") {
let text = String::from_utf8_lossy(&contents);
contents = text
.split("\n## Current Acceptance Follow-up")
.next()
.unwrap_or(&text)
.split("\n## Acceptance #")
.next()
.unwrap_or(&text)
.as_bytes()
.to_vec();
}
output.push((relative, contents));
}
}
Ok(())
}
let mut files = Vec::new();
visit(workspace, workspace, &mut files)?;
files.sort_by(|left, right| left.0.cmp(&right.0));
let hash = files
.into_iter()
.flat_map(|(path, bytes)| path.into_bytes().into_iter().chain(bytes))
.fold(0xcbf29ce484222325u64, |hash, byte| {
(hash ^ byte as u64).wrapping_mul(0x100000001b3)
});
Ok(format!("{hash:016x}"))
}
pub fn decide_acceptance_retry(
previous_identities: &[String],
previous_fingerprint: Option<&str>,
findings: &[NormalizedFinding],
semantic_fingerprint: &str,
cycle_count: u32,
) -> AcceptanceRetryDecision {
let identities = findings
.iter()
.map(|finding| finding.identity.clone())
.collect::<Vec<_>>();
let external_blockers = findings
.iter()
.filter(|finding| finding.external)
.map(|finding| finding.identity.clone())
.collect();
if cycle_count >= MAX_ACCEPTANCE_RETRY_CYCLES {
return AcceptanceRetryDecision::Stall {
reason: "acceptance_cycle_limit_exhausted",
external_blockers,
};
}
if !findings.is_empty() && findings.iter().all(|finding| finding.external) {
return AcceptanceRetryDecision::Stall {
reason: "external_acceptance_blocker",
external_blockers,
};
}
if previous_identities.is_empty() {
return AcceptanceRetryDecision::Retry {
reason: "first_acceptance_failure",
};
}
if previous_identities != identities || previous_fingerprint != Some(semantic_fingerprint) {
return AcceptanceRetryDecision::Retry {
reason: "finding_or_semantic_progress_changed",
};
}
AcceptanceRetryDecision::Stall {
reason: "repeated_acceptance_findings",
external_blockers,
}
}
pub fn build_acceptance_tail_findings(
stdout_tail: Option<String>,
stderr_tail: Option<String>,
) -> Vec<String> {
let stdout = stdout_tail.filter(|text| !text.trim().is_empty());
let stderr = stderr_tail.filter(|text| !text.trim().is_empty());
let selected = stdout
.or(stderr)
.unwrap_or_else(|| ACCEPTANCE_OUTPUT_FALLBACK.to_string());
let lines = selected
.lines()
.filter(|line| {
let trimmed = line.trim();
!trimmed.is_empty()
&& !trimmed.starts_with("ACCEPTANCE:")
&& !trimmed.starts_with("FINDINGS:")
})
.map(|line| line.to_string())
.collect::<Vec<_>>();
if lines.is_empty() {
vec![ACCEPTANCE_OUTPUT_FALLBACK.to_string()]
} else {
lines
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceResult {
Pass,
Fail { findings: Vec<String> },
Continue,
Gated,
CommandFailed {
error: String,
findings: Vec<String>,
},
PermissionStalled {
blocker: crate::events::StalledBlocker,
},
MissingVerdict { findings: Vec<String> },
Cancelled,
}
impl AcceptanceResult {
pub fn is_pass(&self) -> bool {
matches!(self, AcceptanceResult::Pass)
}
pub fn is_canonical_verdict(&self) -> bool {
matches!(
self,
AcceptanceResult::Pass
| AcceptanceResult::Fail { .. }
| AcceptanceResult::Continue
| AcceptanceResult::Gated
| AcceptanceResult::PermissionStalled { .. }
)
}
}
#[allow(clippy::too_many_arguments)]
pub async fn acceptance_test_streaming<O, F>(
change: &Change,
agent: &mut AgentRunner,
ai_runner: &crate::ai_command_runner::AiCommandRunner,
_config: &crate::config::OrchestratorConfig,
output: &O,
cancel_check: F,
protocol_retry: Option<MissingVerdictRetry>,
) -> Result<(AcceptanceResult, u32, String)>
where
O: OutputHandler,
F: Fn() -> bool,
{
use crate::agent::OutputLine;
info!("Running acceptance test for: {}", change.id);
output.on_info(&format!("Acceptance test: {}", change.id));
let commit_hash = crate::vcs::git::commands::get_current_commit(".")
.await
.ok();
let base_branch = crate::vcs::git::commands::get_current_branch(".")
.await
.ok()
.flatten();
let (mut child, mut output_rx, start_time, command) = agent
.run_acceptance_streaming_with_runner(
&change.id,
ai_runner,
None,
base_branch.as_deref(),
protocol_retry,
)
.await?;
output.on_info(&format!("Acceptance started: {}", change.id));
output.on_info(&format!(
" {}",
crate::events::command_log_summary(&command)
));
let mut output_collector = OutputCollector::new();
let mut full_stdout = String::new();
const MARKER_GRACE_PERIOD: std::time::Duration = std::time::Duration::from_secs(30);
let mut marker_detected = false;
let mut verdict_stream_detector = crate::acceptance::VerdictStreamDetector::default();
let mut marker_deadline: Option<tokio::time::Instant> = None;
let mut early_terminated = false;
loop {
let recv_future = output_rx.recv();
let line = if let Some(deadline) = marker_deadline {
match tokio::time::timeout_at(deadline, recv_future).await {
Ok(Some(line)) => line,
Ok(None) => break, Err(_) => {
warn!(
"Acceptance marker grace period expired for {}, terminating process",
change.id
);
let _ = child.terminate();
early_terminated = true;
break;
}
}
} else {
match tokio::time::timeout(std::time::Duration::from_millis(50), recv_future).await {
Ok(Some(line)) => line,
Ok(None) => break, Err(_) => {
if cancel_check() {
warn!("Acceptance test cancelled while waiting for output");
output.on_warn("Acceptance test cancelled");
let _ = child.terminate();
return Ok((AcceptanceResult::Cancelled, 0, command));
}
continue;
}
}
};
if cancel_check() {
warn!("Acceptance test cancelled for: {}", change.id);
output.on_warn("Acceptance test cancelled");
let _ = child.terminate();
return Ok((AcceptanceResult::Cancelled, 0, command));
}
match line {
OutputLine::Stdout(s) => {
output_collector.add_stdout(&s);
full_stdout.push_str(&s);
full_stdout.push('\n');
output.on_stdout(&s);
if !marker_detected && verdict_stream_detector.detect(&s).is_some() {
marker_detected = true;
marker_deadline = Some(tokio::time::Instant::now() + MARKER_GRACE_PERIOD);
info!(
"Acceptance canonical verdict detected for {}, starting {}s grace period",
change.id,
MARKER_GRACE_PERIOD.as_secs()
);
}
}
OutputLine::Stderr(s) => {
output_collector.add_stderr(&s);
output.on_agent_stderr(&s);
}
}
}
let status = loop {
if cancel_check() {
warn!(
"Acceptance test cancelled while waiting for child status for: {}",
change.id
);
output.on_warn("Acceptance test cancelled");
let _ = child.terminate();
return Ok((AcceptanceResult::Cancelled, 0, command));
}
match tokio::time::timeout(std::time::Duration::from_millis(50), child.wait()).await {
Ok(status) => {
break status.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Failed to wait for acceptance command for change '{}': {}",
change.id, e
))
})?;
}
Err(_) => continue,
}
};
let stdout_tail = output_collector.stdout_tail();
let stderr_tail = output_collector.stderr_tail();
let tail_findings = build_acceptance_tail_findings(stdout_tail.clone(), stderr_tail.clone());
let verdict_finalized_run = early_terminated && marker_detected;
if !status.success() && !verdict_finalized_run {
let error_msg = format!(
"Acceptance command failed with exit code: {:?}",
status.code()
);
let attempt_number = agent.next_acceptance_attempt_number(&change.id);
let attempt = AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(tail_findings.clone()),
exit_code: status.code(),
stdout_tail,
stderr_tail,
commit_hash: commit_hash.clone(),
};
agent.record_acceptance_attempt(&change.id, attempt);
output.on_error(&error_msg);
return Ok((
AcceptanceResult::CommandFailed {
error: error_msg,
findings: tail_findings,
},
attempt_number,
command,
));
}
let parsed_result = crate::acceptance::parse_acceptance_output(&full_stdout);
let (result, passed) = match parsed_result {
crate::acceptance::AcceptanceResult::Pass => {
info!("Acceptance test passed for: {}", change.id);
output.on_info("Acceptance test: PASS");
(AcceptanceResult::Pass, true)
}
crate::acceptance::AcceptanceResult::Fail {
findings: parsed_findings,
} => {
info!("Acceptance test failed for: {}", change.id);
output.on_warn("Acceptance test: FAIL");
let findings = if parsed_findings.is_empty() {
vec!["Investigate acceptance failure and apply the required fix".to_string()]
} else {
parsed_findings
};
(AcceptanceResult::Fail { findings }, false)
}
crate::acceptance::AcceptanceResult::Continue => {
info!("Acceptance requires continuation for: {}", change.id);
output.on_info("Acceptance test: CONTINUE");
(AcceptanceResult::Continue, false)
}
crate::acceptance::AcceptanceResult::Gated => {
info!("Acceptance gated for: {}", change.id);
output.on_warn("Acceptance test: GATED");
(AcceptanceResult::Gated, false)
}
crate::acceptance::AcceptanceResult::MissingVerdict => {
warn!(
"Acceptance completed without a canonical verdict for: {} (missing-verdict protocol failure)",
change.id
);
output.on_error("Acceptance test: MISSING VERDICT (protocol failure — the acceptance command exited without a canonical verdict; status-only or waiting output is not a verdict)");
(
AcceptanceResult::MissingVerdict {
findings: tail_findings.clone(),
},
false,
)
}
};
let history_findings = match &result {
AcceptanceResult::Fail { findings } => Some(findings.clone()),
AcceptanceResult::Continue => {
Some(vec!["Investigation incomplete - continue later".to_string()])
}
AcceptanceResult::Gated => Some(vec!["Implementation blocker detected".to_string()]),
AcceptanceResult::Pass => None,
AcceptanceResult::MissingVerdict { findings } => {
let mut evidence = vec![MISSING_VERDICT_DIAGNOSTIC.to_string()];
evidence.extend(findings.iter().cloned());
Some(evidence)
}
AcceptanceResult::CommandFailed { .. }
| AcceptanceResult::PermissionStalled { .. }
| AcceptanceResult::Cancelled => Some(tail_findings.clone()),
};
let attempt_number = agent.next_acceptance_attempt_number(&change.id);
let attempt = AcceptanceAttempt {
attempt: attempt_number,
passed,
duration: start_time.elapsed(),
findings: history_findings,
exit_code: status.code(),
stdout_tail,
stderr_tail,
commit_hash: commit_hash.clone(),
};
agent.record_acceptance_attempt(&change.id, attempt);
match &result {
AcceptanceResult::Fail { findings } => {
if !findings.is_empty() {
agent.record_acceptance_follow_up(&change.id, attempt_number, findings.clone());
}
}
AcceptanceResult::Pass => agent.clear_acceptance_follow_up(&change.id),
_ => {}
}
Ok((result, attempt_number, command))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn missing_verdict_budget_allows_two_retries_then_exhausts() {
let mut state = MissingVerdictRetryState::default();
let expected = [
MissingVerdictRetryDecision::Retry(MissingVerdictRetry { attempt: 1, max: 2 }),
MissingVerdictRetryDecision::Retry(MissingVerdictRetry { attempt: 2, max: 2 }),
MissingVerdictRetryDecision::Exhausted {
attempts: 3,
max: 2,
},
];
for (index, want) in expected.iter().enumerate() {
assert_eq!(
state.record_missing_verdict(),
*want,
"consecutive missing verdict #{} must route as {:?}",
index + 1,
want
);
}
assert_eq!(MAX_MISSING_VERDICT_RETRIES, 2);
assert_eq!(
state.consecutive(),
3,
"a fourth protocol retry must never be offered after exhaustion"
);
assert!(matches!(
state.record_missing_verdict(),
MissingVerdictRetryDecision::Exhausted { .. }
));
}
#[test]
fn missing_verdict_canonical_verdict_resets_consecutive_sequence() {
let mut state = MissingVerdictRetryState::default();
assert!(matches!(
state.record_missing_verdict(),
MissingVerdictRetryDecision::Retry(MissingVerdictRetry { attempt: 1, .. })
));
assert!(matches!(
state.record_missing_verdict(),
MissingVerdictRetryDecision::Retry(MissingVerdictRetry { attempt: 2, .. })
));
state.record_canonical_verdict();
assert_eq!(state.consecutive(), 0);
assert_eq!(
state.record_missing_verdict(),
MissingVerdictRetryDecision::Retry(MissingVerdictRetry { attempt: 1, max: 2 }),
"a later missing verdict must start a fresh protocol-retry sequence"
);
}
#[test]
fn missing_verdict_budget_is_independent_of_configured_continue_count() {
for configured_continues in [0u32, 1, 5, 50] {
let mut state = MissingVerdictRetryState::default();
let mut retries = 0u32;
loop {
match state.record_missing_verdict() {
MissingVerdictRetryDecision::Retry(retry) => {
assert_eq!(retry.max, MAX_MISSING_VERDICT_RETRIES);
retries += 1;
}
MissingVerdictRetryDecision::Exhausted { attempts, max } => {
assert_eq!(attempts, MAX_MISSING_VERDICT_RETRIES + 1);
assert_eq!(max, MAX_MISSING_VERDICT_RETRIES);
break;
}
}
}
assert_eq!(
retries, MAX_MISSING_VERDICT_RETRIES,
"configured acceptance_max_continues={configured_continues} must not change the \
dedicated missing-verdict budget"
);
}
}
#[test]
fn missing_verdict_diagnostics_are_bounded_and_distinguish_progress_from_terminal() {
let findings = (0..10)
.map(|index| format!("evidence line {index}"))
.collect::<Vec<_>>();
let progress =
missing_verdict_retry_progress(MissingVerdictRetry { attempt: 1, max: 2 }, &findings);
assert!(progress.contains("protocol retry 1/2"));
assert!(progress.contains("evidence line 4"));
assert!(
!progress.contains("evidence line 5"),
"evidence must stay bounded to the first five findings, got {progress}"
);
assert!(
!progress.to_ascii_lowercase().contains("protocol failure"),
"an available retry must not read as a terminal failure: {progress}"
);
let terminal = missing_verdict_exhausted_error(3, 2, &findings);
assert!(terminal.contains("missing-verdict protocol failure"));
assert!(terminal.contains("Exhausted 3 consecutive attempts after 2 protocol retries"));
assert!(terminal.contains("evidence line 4"));
assert!(!terminal.contains("evidence line 5"));
assert!(
missing_verdict_exhausted_error(3, 2, &[]).contains("no acceptance output captured"),
"empty evidence must still produce an actionable diagnostic"
);
}
fn replay_missing_verdict_sequence(
sequence: &[AcceptanceResult],
) -> (
Vec<Option<MissingVerdictRetry>>,
Option<AcceptanceResult>,
Option<String>,
) {
let mut protocol = MissingVerdictRetryDriver::default();
let mut received = Vec::new();
for result in sequence {
received.push(protocol.take_protocol_retry());
match result {
AcceptanceResult::MissingVerdict { findings } => {
match protocol.observe_missing_verdict(findings) {
MissingVerdictRetryStep::Retry { .. } => continue,
MissingVerdictRetryStep::Exhausted { error } => {
return (received, None, Some(error))
}
}
}
canonical => {
protocol.observe_canonical_verdict();
return (received, Some(canonical.clone()), None);
}
}
}
(received, None, None)
}
fn missing_verdict(evidence: &str) -> AcceptanceResult {
AcceptanceResult::MissingVerdict {
findings: vec![evidence.to_string()],
}
}
#[test]
fn missing_verdict_driver_retries_twice_then_accepts_canonical_pass() {
let (received, canonical, error) = replay_missing_verdict_sequence(&[
missing_verdict("waiting for verification"),
missing_verdict("still waiting"),
AcceptanceResult::Pass,
]);
assert_eq!(received.len(), 3, "acceptance must be invoked three times");
assert_eq!(
received,
vec![
None,
Some(MissingVerdictRetry { attempt: 1, max: 2 }),
Some(MissingVerdictRetry { attempt: 2, max: 2 }),
],
"only the retries may carry a continuation marker"
);
assert_eq!(canonical, Some(AcceptanceResult::Pass));
assert!(
error.is_none(),
"a canonical verdict within budget must not produce a terminal error"
);
}
#[test]
fn missing_verdict_driver_exhausts_after_three_consecutive_missing_verdicts() {
let (received, canonical, error) = replay_missing_verdict_sequence(&[
missing_verdict("waiting one"),
missing_verdict("waiting two"),
missing_verdict("waiting three"),
missing_verdict("never reached"),
]);
assert_eq!(
received.len(),
3,
"no fourth protocol retry may start after exhaustion"
);
assert!(canonical.is_none());
let error = error.expect("third consecutive missing verdict must be terminal");
assert!(error.contains("missing-verdict protocol failure"));
assert!(error.contains("Exhausted 3 consecutive attempts after 2 protocol retries"));
assert!(
error.contains("waiting three"),
"terminal diagnostic must carry bounded evidence, got {error}"
);
}
#[test]
fn missing_verdict_driver_treats_every_canonical_outcome_as_a_reset() {
for canonical in [
AcceptanceResult::Pass,
AcceptanceResult::Fail {
findings: vec!["src/lib.rs:1 fix".to_string()],
},
AcceptanceResult::Continue,
AcceptanceResult::Gated,
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_infrastructure("denied"),
},
] {
let mut protocol = MissingVerdictRetryDriver::default();
assert!(matches!(
protocol.observe_missing_verdict(&["waiting".to_string()]),
MissingVerdictRetryStep::Retry { .. }
));
assert!(protocol.take_protocol_retry().is_some());
protocol.observe_canonical_verdict();
assert_eq!(
protocol.consecutive_missing_verdicts(),
0,
"{canonical:?} must reset the consecutive protocol counter"
);
assert!(
protocol.take_protocol_retry().is_none(),
"{canonical:?} must clear any pending continuation marker"
);
assert!(matches!(
protocol.observe_missing_verdict(&["waiting again".to_string()]),
MissingVerdictRetryStep::Retry {
retry: MissingVerdictRetry { attempt: 1, .. },
..
}
));
}
}
#[test]
fn acceptance_routing_matrix_keeps_missing_verdict_distinct() {
let cases: [(AcceptanceResult, bool, bool); 8] = [
(AcceptanceResult::Pass, true, false),
(
AcceptanceResult::Fail {
findings: vec!["src/lib.rs:1 fix".to_string()],
},
true,
false,
),
(AcceptanceResult::Continue, true, false),
(AcceptanceResult::Gated, true, false),
(
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_infrastructure("denied"),
},
true,
false,
),
(missing_verdict("waiting"), false, true),
(
AcceptanceResult::CommandFailed {
error: "exit code 1".to_string(),
findings: vec!["boom".to_string()],
},
false,
false,
),
(AcceptanceResult::Cancelled, false, false),
];
for (result, canonical, missing) in cases {
assert_eq!(
result.is_canonical_verdict(),
canonical,
"{result:?} canonical-verdict classification"
);
assert_eq!(
matches!(result, AcceptanceResult::MissingVerdict { .. }),
missing,
"{result:?} must not be confused with a missing verdict"
);
assert!(
!(canonical && missing),
"{result:?} cannot be both canonical and a protocol failure"
);
}
let mut protocol = MissingVerdictRetryDriver::default();
for result in [
AcceptanceResult::CommandFailed {
error: "exit code 1".to_string(),
findings: Vec::new(),
},
AcceptanceResult::Cancelled,
] {
assert!(!result.is_canonical_verdict());
assert!(protocol.take_protocol_retry().is_none());
assert_eq!(protocol.consecutive_missing_verdicts(), 0);
}
assert!(matches!(
protocol.observe_missing_verdict(&["waiting".to_string()]),
MissingVerdictRetryStep::Retry {
retry: MissingVerdictRetry { attempt: 1, max: 2 },
..
}
));
}
#[test]
fn missing_verdict_driver_has_serial_and_parallel_routing_parity() {
let sequence = [
missing_verdict("waiting"),
missing_verdict("waiting"),
AcceptanceResult::Fail {
findings: vec!["src/lib.rs:1 missing coverage".to_string()],
},
];
let serial = replay_missing_verdict_sequence(&sequence);
let parallel = replay_missing_verdict_sequence(&sequence);
assert_eq!(serial, parallel);
assert_eq!(serial.0.len(), 3);
assert!(matches!(serial.1, Some(AcceptanceResult::Fail { .. })));
assert!(serial.2.is_none());
}
#[test]
fn semantic_fingerprint_excludes_runtime_bookkeeping() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::create_dir_all(temp.path().join("src")).unwrap();
std::fs::write(temp.path().join("src/lib.rs"), "one").unwrap();
let before = semantic_progress_fingerprint(temp.path()).unwrap();
std::fs::create_dir_all(temp.path().join(".cflx")).unwrap();
std::fs::write(temp.path().join(".cflx/runtime.json"), "runtime").unwrap();
assert_eq!(before, semantic_progress_fingerprint(temp.path()).unwrap());
std::fs::write(temp.path().join("src/lib.rs"), "two").unwrap();
assert_ne!(before, semantic_progress_fingerprint(temp.path()).unwrap());
}
#[test]
fn finding_identity_prefers_code_and_uses_structural_fallback() {
let coded = normalize_findings(&[
"[MISSING_RETRY_TEST] old evidence at src/run.rs:10".into(),
"[MISSING_RETRY_TEST] changed summary at tests/run.rs:99".into(),
]);
assert_eq!(coded.len(), 1);
assert_eq!(coded[0].identity, "repository|code|[missing_retry_test]");
let changed_detail = normalize_findings(&[
"Missing retry test at src/run.rs:10 because the branch is uncovered".into(),
"Regression coverage absent in src/run.rs:77; add a focused test".into(),
]);
assert_eq!(changed_detail.len(), 1);
assert_eq!(
changed_detail[0].identity,
"repository|src/run.rs|verification"
);
let distinct = normalize_findings(&[
"Missing test at src/run.rs:10".into(),
"Incorrect implementation at src/run.rs:11".into(),
"Missing test at src/other.rs:10".into(),
]);
assert_eq!(
distinct.len(),
3,
"rule and location must prevent collisions"
);
}
#[test]
fn retry_decision_normalizes_order_whitespace_duplicates_and_stalls_repeats() {
let findings = normalize_findings(&[
" src/lib.rs:10 missing test ".to_string(),
"src/lib.rs:11 missing test".to_string(),
]);
assert_eq!(findings.len(), 1);
let decision = decide_acceptance_retry(
&[findings[0].identity.clone()],
Some("unchanged"),
&findings,
"unchanged",
2,
);
assert!(matches!(
decision,
AcceptanceRetryDecision::Stall {
reason: "repeated_acceptance_findings",
..
}
));
}
#[test]
fn retry_decision_stalls_external_only_and_allows_progress_changed() {
let findings = normalize_findings(&["external service outage".to_string()]);
assert!(findings[0].external);
assert!(matches!(
decide_acceptance_retry(&[], None, &findings, "one", 1),
AcceptanceRetryDecision::Stall {
reason: "external_acceptance_blocker",
..
}
));
assert!(
matches!(decide_acceptance_retry(&[], None, &findings, "one", MAX_ACCEPTANCE_RETRY_CYCLES), AcceptanceRetryDecision::Stall { reason: "acceptance_cycle_limit_exhausted", external_blockers } if external_blockers.len() == 1)
);
}
#[test]
fn semantic_fingerprint_tracks_change_specs_and_jsonc_but_excludes_runtime_follow_up() {
let temp = tempfile::TempDir::new().unwrap();
let tasks = temp.path().join("openspec/changes/example/tasks.md");
let spec = temp
.path()
.join("openspec/changes/example/specs/runtime/spec.md");
std::fs::create_dir_all(spec.parent().unwrap()).unwrap();
std::fs::write(&tasks, "## Implementation Tasks\n- [x] work\n").unwrap();
std::fs::write(&spec, "requirement one").unwrap();
std::fs::write(temp.path().join(".cflx.jsonc"), "{ \"mode\": 1 }").unwrap();
let before = semantic_progress_fingerprint(temp.path()).unwrap();
std::fs::write(&tasks, "## Implementation Tasks\n- [x] work\n\n## Acceptance #2 Failure Follow-up\n- [ ] runtime finding\n").unwrap();
assert_eq!(before, semantic_progress_fingerprint(temp.path()).unwrap());
std::fs::write(&spec, "requirement two").unwrap();
assert_ne!(before, semantic_progress_fingerprint(temp.path()).unwrap());
std::fs::write(temp.path().join(".cflx.jsonc"), "{ \"mode\": 2 }").unwrap();
assert_ne!(before, semantic_progress_fingerprint(temp.path()).unwrap());
}
#[test]
fn generic_credential_and_unavailable_errors_remain_repository_fixable() {
let findings = normalize_findings(&[
"missing API key in test fixture".to_string(),
"src/client.rs: rate limit retry missing".to_string(),
"network unreachable: fix retry handling".to_string(),
"dns resolution failed while repairing src/client.rs".to_string(),
"missing non-mockable external credential".to_string(),
]);
assert_eq!(
findings.iter().filter(|finding| finding.external).count(),
1
);
assert!(findings[0].identity.starts_with("external|"));
assert_eq!(
repository_findings(&[
"missing API key in test fixture".to_string(),
"src/client.rs: rate limit retry missing".to_string(),
"network unreachable: fix retry handling".to_string(),
"dns resolution failed while repairing src/client.rs".to_string(),
"missing non-mockable external credential".to_string(),
])
.len(),
4
);
}
#[test]
fn alternating_continue_and_fail_keeps_fail_retry_history_deterministic() {
let findings = normalize_findings(&["src/lib.rs:10 missing regression coverage".into()]);
let identities = findings
.iter()
.map(|finding| finding.identity.clone())
.collect::<Vec<_>>();
assert!(matches!(
decide_acceptance_retry(&[], None, &findings, "unchanged", 1),
AcceptanceRetryDecision::Retry {
reason: "first_acceptance_failure"
}
));
assert!(matches!(
decide_acceptance_retry(&identities, Some("unchanged"), &findings, "unchanged", 2),
AcceptanceRetryDecision::Stall {
reason: "repeated_acceptance_findings",
..
}
));
}
#[test]
fn serial_and_parallel_same_inputs_have_retry_outcome_parity() {
let findings = normalize_findings(&[
"src/lib.rs:10 missing regression coverage".into(),
"external non-mockable prerequisite unavailable".into(),
]);
let previous = findings
.iter()
.map(|finding| finding.identity.clone())
.collect::<Vec<_>>();
let serial = decide_acceptance_retry(&previous, Some("same"), &findings, "same", 2);
let parallel = decide_acceptance_retry(&previous, Some("same"), &findings, "same", 2);
assert_eq!(serial, parallel);
assert!(matches!(
serial,
AcceptanceRetryDecision::Stall {
reason: "repeated_acceptance_findings",
ref external_blockers
} if external_blockers.len() == 1
));
}
#[test]
fn retry_decision_handles_mixed_and_findingless_failures() {
let mixed = normalize_findings(&[
"src/lib.rs:1 fix test".into(),
"external service outage".into(),
]);
assert!(matches!(
decide_acceptance_retry(&[], None, &mixed, "one", 1),
AcceptanceRetryDecision::Retry { .. }
));
assert!(matches!(
decide_acceptance_retry(&[], None, &[], "one", 1),
AcceptanceRetryDecision::Retry { .. }
));
assert_eq!(
repository_findings(&[
"src/lib.rs:1 fix test".into(),
"external service outage".into(),
]),
vec!["src/lib.rs:1 fix test"]
);
}
#[test]
fn test_build_acceptance_tail_findings_prefers_stdout() {
let findings = build_acceptance_tail_findings(
Some("stdout line 1\nstdout line 2".to_string()),
Some("stderr line".to_string()),
);
assert_eq!(findings, vec!["stdout line 1", "stdout line 2"]);
}
#[test]
fn test_build_acceptance_tail_findings_falls_back_to_stderr() {
let findings =
build_acceptance_tail_findings(Some(" ".to_string()), Some("stderr".to_string()));
assert_eq!(findings, vec!["stderr"]);
}
#[test]
fn test_build_acceptance_tail_findings_fallback_message() {
let findings = build_acceptance_tail_findings(None, Some("\n\n".to_string()));
assert_eq!(findings, vec!["No acceptance output captured"]);
}
#[test]
fn test_acceptance_result_is_pass() {
assert!(AcceptanceResult::Pass.is_pass());
assert!(!AcceptanceResult::Fail {
findings: vec!["error".to_string()]
}
.is_pass());
assert!(!AcceptanceResult::CommandFailed {
error: "test".to_string(),
findings: vec!["failure".to_string()],
}
.is_pass());
assert!(!AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_infrastructure("permission denied"),
}
.is_pass());
assert!(!AcceptanceResult::MissingVerdict {
findings: vec!["status-only output".to_string()],
}
.is_pass());
assert!(!AcceptanceResult::Cancelled.is_pass());
assert!(!AcceptanceResult::Gated.is_pass());
}
#[test]
fn test_build_acceptance_tail_findings_filters_acceptance_marker() {
let findings = build_acceptance_tail_findings(
Some("line 1\nACCEPTANCE: FAIL\nline 2".to_string()),
None,
);
assert_eq!(findings, vec!["line 1", "line 2"]);
}
#[test]
fn test_build_acceptance_tail_findings_filters_findings_line() {
let findings = build_acceptance_tail_findings(
Some("error 1\nFINDINGS:\n- item 1\n- item 2".to_string()),
None,
);
assert_eq!(findings, vec!["error 1", "- item 1", "- item 2"]);
}
#[test]
fn test_build_acceptance_tail_findings_filters_both_markers() {
let findings = build_acceptance_tail_findings(
Some("ACCEPTANCE: FAIL\nFINDINGS:\nactual error\nanother line".to_string()),
None,
);
assert_eq!(findings, vec!["actual error", "another line"]);
}
#[test]
fn test_tail_findings_includes_preamble_parse_does_not() {
let stdout = "preamble\nACCEPTANCE: FAIL\nFINDINGS:\n- Finding 1\n- Finding 2\npostamble"
.to_string();
let tail = build_acceptance_tail_findings(Some(stdout.clone()), None);
assert!(tail.iter().any(|l| l.contains("preamble")));
assert!(tail.iter().any(|l| l.contains("postamble")));
assert!(tail.iter().any(|l| l.contains("Finding 1")));
match crate::acceptance::parse_acceptance_output(&stdout) {
crate::acceptance::AcceptanceResult::Fail { findings } => {
assert_eq!(findings, vec!["Finding 1", "Finding 2"]);
assert!(!findings.iter().any(|f| f.contains("preamble")));
assert!(!findings.iter().any(|f| f.contains("postamble")));
}
_ => panic!("Expected Fail"),
}
}
#[test]
fn test_parse_findings_is_preferred_source_for_fail_result() {
let stdout =
"ACCEPTANCE: FAIL\nFINDINGS:\n- src/foo.rs:10 issue A\n- src/bar.rs:5 issue B\n"
.to_string();
match crate::acceptance::parse_acceptance_output(&stdout) {
crate::acceptance::AcceptanceResult::Fail { findings } => {
assert_eq!(findings.len(), 2);
assert_eq!(findings[0], "src/foo.rs:10 issue A");
assert_eq!(findings[1], "src/bar.rs:5 issue B");
}
_ => panic!("Expected Fail"),
}
}
}