#![allow(dead_code)]
pub mod verification_evidence;
pub mod execution_manifest;
pub mod gate_execution;
pub mod manifest_builder;
pub mod evidence_location;
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;
pub const MAX_ACCEPTANCE_PROTOCOL_RETRIES: u32 = MAX_MISSING_VERDICT_RETRIES;
pub const BARE_BLOCKER_DIAGNOSTIC: &str = "Bare acceptance blocker: acceptance emitted a gated \
compatibility token without a validated structured blocker payload (protocol failure; a \
stalled hold requires an explicit supported category, concrete evidence, next action, and \
resumability)";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AcceptanceProtocolError {
MissingVerdict,
BareBlocker,
MalformedFinding,
}
impl AcceptanceProtocolError {
pub fn label(self) -> &'static str {
match self {
AcceptanceProtocolError::MissingVerdict => "missing-verdict",
AcceptanceProtocolError::BareBlocker => "bare-blocker",
AcceptanceProtocolError::MalformedFinding => "malformed-finding",
}
}
}
pub const MALFORMED_FINDING_DIAGNOSTIC: &str =
"Malformed acceptance finding: acceptance emitted a \
FAIL verdict with a structured finding that is missing a stable id, severity, summary, \
evidence, required_changes, or verification (protocol failure; runtime will not reduce it to \
a path-only repair instruction)";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AcceptanceProtocolRetry {
pub kind: AcceptanceProtocolError,
pub attempt: u32,
pub max: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MissingVerdictRetryDecision {
Retry(AcceptanceProtocolRetry),
Exhausted {
kind: AcceptanceProtocolError,
attempts: u32,
max: u32,
},
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ProtocolRetryCounter {
consecutive: u32,
}
impl ProtocolRetryCounter {
pub fn consecutive(&self) -> u32 {
self.consecutive
}
pub fn reset(&mut self) {
self.consecutive = 0;
}
pub fn record(&mut self, kind: AcceptanceProtocolError) -> MissingVerdictRetryDecision {
self.consecutive = self.consecutive.saturating_add(1);
if self.consecutive <= MAX_ACCEPTANCE_PROTOCOL_RETRIES {
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry {
kind,
attempt: self.consecutive,
max: MAX_ACCEPTANCE_PROTOCOL_RETRIES,
})
} else {
MissingVerdictRetryDecision::Exhausted {
kind,
attempts: self.consecutive,
max: MAX_ACCEPTANCE_PROTOCOL_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: AcceptanceProtocolRetry,
findings: &[String],
) -> String {
let cause = match retry.kind {
AcceptanceProtocolError::MissingVerdict => "without a canonical verdict",
AcceptanceProtocolError::BareBlocker => "with a gated token but no validated blocker",
AcceptanceProtocolError::MalformedFinding => {
"with a FAIL verdict whose structured finding did not validate"
}
};
format!(
"Acceptance completed {cause}; 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 {
protocol_exhausted_error(
AcceptanceProtocolError::MissingVerdict,
attempts,
max,
findings,
)
}
pub fn protocol_exhausted_error(
kind: AcceptanceProtocolError,
attempts: u32,
max: u32,
findings: &[String],
) -> String {
let cause = match kind {
AcceptanceProtocolError::MissingVerdict => {
"Acceptance completed without a canonical verdict (missing-verdict protocol failure); \
status-only or waiting output is not a verdict."
}
AcceptanceProtocolError::BareBlocker => {
"Acceptance emitted a gated compatibility token without a validated structured blocker \
(bare-blocker protocol failure); a stalled hold requires an explicit supported \
category, concrete evidence, next action, and resumability."
}
AcceptanceProtocolError::MalformedFinding => {
"Acceptance emitted a FAIL verdict whose structured finding did not validate \
(malformed-finding protocol failure); a structured finding requires a stable id, a \
major/minor severity, a summary, concrete evidence, and repository-relative \
required_changes and verification entries."
}
};
format!(
"{cause} Exhausted {attempts} consecutive attempts after {max} protocol retries. \
Evidence: {}",
bounded_missing_verdict_evidence(findings)
)
}
pub const MAX_ACCEPTANCE_COMMAND_RETRIES: u32 = 2;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AcceptanceCommandDiagnostic {
pub error: String,
pub exit_code: Option<i32>,
pub stdout_tail: Option<String>,
pub stderr_tail: Option<String>,
}
impl AcceptanceCommandDiagnostic {
pub fn summary(&self) -> String {
const MAX_TAIL_CHARS: usize = 400;
fn condense(tail: &str) -> String {
let single_line = tail.split_whitespace().collect::<Vec<_>>().join(" ");
match single_line.char_indices().nth(MAX_TAIL_CHARS) {
Some((idx, _)) => format!("{}...", &single_line[..idx]),
None => single_line,
}
}
let mut parts = vec![self.error.clone()];
if let Some(code) = self.exit_code {
parts.push(format!("exit_code: {}", code));
}
if let Some(stdout) = self.stdout_tail.as_deref().filter(|t| !t.trim().is_empty()) {
parts.push(format!("stdout: {}", condense(stdout)));
}
if let Some(stderr) = self.stderr_tail.as_deref().filter(|t| !t.trim().is_empty()) {
parts.push(format!("stderr: {}", condense(stderr)));
}
parts.join(" | ")
}
}
pub fn classify_acceptance_command_denial(
diagnostic: &AcceptanceCommandDiagnostic,
) -> Option<crate::permission::PermissionDenial> {
crate::permission::classify_permission_denial(&[
diagnostic.stdout_tail.as_deref(),
diagnostic.stderr_tail.as_deref(),
])
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceCommandRetryDecision {
Retry {
attempt: u32,
max_retries: u32,
diagnostic: AcceptanceCommandDiagnostic,
},
Exhausted {
attempts: u32,
max_retries: u32,
diagnostic: AcceptanceCommandDiagnostic,
},
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct AcceptanceCommandRetryCounter {
consecutive_failures: u32,
}
impl AcceptanceCommandRetryCounter {
pub fn consecutive_failures(&self) -> u32 {
self.consecutive_failures
}
pub fn reset(&mut self) {
self.consecutive_failures = 0;
}
pub fn record_command_failure(
&mut self,
diagnostic: AcceptanceCommandDiagnostic,
) -> AcceptanceCommandRetryDecision {
self.consecutive_failures = self.consecutive_failures.saturating_add(1);
if self.consecutive_failures <= MAX_ACCEPTANCE_COMMAND_RETRIES {
AcceptanceCommandRetryDecision::Retry {
attempt: self.consecutive_failures,
max_retries: MAX_ACCEPTANCE_COMMAND_RETRIES,
diagnostic,
}
} else {
AcceptanceCommandRetryDecision::Exhausted {
attempts: self.consecutive_failures,
max_retries: MAX_ACCEPTANCE_COMMAND_RETRIES,
diagnostic,
}
}
}
}
pub fn acceptance_command_retry_progress(
attempt: u32,
max_retries: u32,
diagnostic: &AcceptanceCommandDiagnostic,
) -> String {
format!(
"Acceptance command did not complete (command-failure recovery {attempt}/{max_retries}); \
re-running only the configured acceptance command against the same applied workspace. \
Evidence: {}",
diagnostic.summary()
)
}
pub fn acceptance_command_exhausted_error(
attempts: u32,
max_retries: u32,
diagnostic: &AcceptanceCommandDiagnostic,
) -> String {
format!(
"Acceptance command failed to complete on {attempts} consecutive attempts after \
{max_retries} command-failure retries. Evidence: {}",
diagnostic.summary()
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceCommandRecovery {
Retry { attempt: u32, progress: String },
Exhausted { attempts: u32, error: String },
}
pub fn decide_acceptance_command_failure(
counter: &mut AcceptanceCommandRetryCounter,
agent: &mut AgentRunner,
change_id: &str,
diagnostic: AcceptanceCommandDiagnostic,
) -> AcceptanceCommandRecovery {
match counter.record_command_failure(diagnostic) {
AcceptanceCommandRetryDecision::Retry {
attempt,
max_retries,
diagnostic,
} => {
let progress = acceptance_command_retry_progress(attempt, max_retries, &diagnostic);
agent.set_acceptance_command_recovery(change_id, diagnostic);
AcceptanceCommandRecovery::Retry { attempt, progress }
}
AcceptanceCommandRetryDecision::Exhausted {
attempts,
max_retries,
diagnostic,
} => {
agent.clear_acceptance_command_recovery(change_id);
AcceptanceCommandRecovery::Exhausted {
attempts,
error: acceptance_command_exhausted_error(attempts, max_retries, &diagnostic),
}
}
}
}
pub fn observe_completed_acceptance_invocation(
counter: &mut AcceptanceCommandRetryCounter,
agent: &mut AgentRunner,
change_id: &str,
) {
counter.reset();
agent.clear_acceptance_command_recovery(change_id);
}
pub fn observe_acceptance_invocation_result(
counter: &mut AcceptanceCommandRetryCounter,
agent: &mut AgentRunner,
change_id: &str,
result: &AcceptanceResult,
) {
if result.permits_acceptance_retry() {
observe_completed_acceptance_invocation(counter, agent, change_id);
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MissingVerdictRetryStep {
Retry {
retry: AcceptanceProtocolRetry,
progress: String,
},
Exhausted { error: String },
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AcceptanceProtocolDriver {
missing_verdict: ProtocolRetryCounter,
bare_blocker: ProtocolRetryCounter,
malformed_finding: ProtocolRetryCounter,
pending: Option<AcceptanceProtocolRetry>,
}
impl AcceptanceProtocolDriver {
pub fn take_protocol_retry(&mut self) -> Option<AcceptanceProtocolRetry> {
self.pending.take()
}
pub fn consecutive_missing_verdicts(&self) -> u32 {
self.missing_verdict.consecutive()
}
pub fn consecutive_bare_blockers(&self) -> u32 {
self.bare_blocker.consecutive()
}
pub fn observe_canonical_verdict(&mut self) {
self.missing_verdict.reset();
self.bare_blocker.reset();
self.malformed_finding.reset();
self.pending = None;
}
pub fn consecutive_malformed_findings(&self) -> u32 {
self.malformed_finding.consecutive()
}
pub fn observe_malformed_finding(
&mut self,
rejection: &crate::acceptance::FindingRejection,
) -> MissingVerdictRetryStep {
let evidence = vec![MALFORMED_FINDING_DIAGNOSTIC.to_string(), rejection.reason()];
let decision = self
.malformed_finding
.record(AcceptanceProtocolError::MalformedFinding);
self.step_from(decision, &evidence)
}
pub fn observe_missing_verdict(&mut self, findings: &[String]) -> MissingVerdictRetryStep {
let decision = self
.missing_verdict
.record(AcceptanceProtocolError::MissingVerdict);
self.step_from(decision, findings)
}
pub fn observe_bare_blocker(
&mut self,
rejection: &crate::acceptance::BlockerRejection,
) -> MissingVerdictRetryStep {
let evidence = vec![BARE_BLOCKER_DIAGNOSTIC.to_string(), rejection.reason()];
let decision = self
.bare_blocker
.record(AcceptanceProtocolError::BareBlocker);
self.step_from(decision, &evidence)
}
fn step_from(
&mut self,
decision: MissingVerdictRetryDecision,
findings: &[String],
) -> MissingVerdictRetryStep {
match decision {
MissingVerdictRetryDecision::Retry(retry) => {
self.pending = Some(retry);
MissingVerdictRetryStep::Retry {
retry,
progress: missing_verdict_retry_progress(retry, findings),
}
}
MissingVerdictRetryDecision::Exhausted {
kind,
attempts,
max,
} => {
self.pending = None;
MissingVerdictRetryStep::Exhausted {
error: protocol_exhausted_error(kind, attempts, max, findings),
}
}
}
}
}
pub const GENERIC_ACCEPTANCE_FAIL_FINDING: &str =
"Investigate acceptance failure and apply the required fix";
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum AcceptanceCommandMode {
#[default]
Normal,
Escalation,
}
impl AcceptanceCommandMode {
pub fn is_escalation(self) -> bool {
matches!(self, AcceptanceCommandMode::Escalation)
}
pub fn label(self) -> &'static str {
match self {
AcceptanceCommandMode::Normal => "acceptance_command",
AcceptanceCommandMode::Escalation => "acceptance_escalation_command",
}
}
}
pub fn acceptance_command_template(
config: &crate::config::OrchestratorConfig,
command_mode: AcceptanceCommandMode,
) -> Result<&str> {
match command_mode {
AcceptanceCommandMode::Normal => config.get_acceptance_command(),
AcceptanceCommandMode::Escalation => {
config.get_acceptance_escalation_command().ok_or_else(|| {
OrchestratorError::ConfigLoad(
"Missing optional config: acceptance_escalation_command. Acceptance \
escalation was selected without a configured alternate reviewer command."
.to_string(),
)
})
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InvalidAcceptanceResult {
MissingVerdict,
BareBlocker,
MalformedFinding,
EmptyFail,
}
impl InvalidAcceptanceResult {
pub fn label(self) -> &'static str {
match self {
InvalidAcceptanceResult::MissingVerdict => "missing-verdict",
InvalidAcceptanceResult::BareBlocker => "bare-blocker",
InvalidAcceptanceResult::MalformedFinding => "malformed-finding",
InvalidAcceptanceResult::EmptyFail => "empty-fail",
}
}
pub fn has_protocol_retry_budget(self) -> bool {
!matches!(self, InvalidAcceptanceResult::EmptyFail)
}
}
pub fn is_generic_empty_fail(findings: &[crate::acceptance::AcceptanceFinding]) -> bool {
match findings {
[] => true,
[only] => {
only.structured_payload().is_none() && only.text() == GENERIC_ACCEPTANCE_FAIL_FINDING
}
_ => false,
}
}
pub fn classify_invalid_acceptance_result(
result: &AcceptanceResult,
) -> Option<InvalidAcceptanceResult> {
match result {
AcceptanceResult::MissingVerdict { .. } => Some(InvalidAcceptanceResult::MissingVerdict),
AcceptanceResult::BareBlocker { .. } => Some(InvalidAcceptanceResult::BareBlocker),
AcceptanceResult::MalformedFinding { .. } => {
Some(InvalidAcceptanceResult::MalformedFinding)
}
AcceptanceResult::Fail { findings } => {
is_generic_empty_fail(findings).then_some(InvalidAcceptanceResult::EmptyFail)
}
AcceptanceResult::Pass
| AcceptanceResult::Continue
| AcceptanceResult::Stalled { .. }
| AcceptanceResult::PermissionStalled { .. }
| AcceptanceResult::CommandFailed { .. }
| AcceptanceResult::RuntimeLimit { .. }
| AcceptanceResult::ExecutionHold { .. }
| AcceptanceResult::Cancelled => None,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EscalationDeclineReason {
NoCommandConfigured,
BelowThreshold { after_invalid_results: u32 },
UseCapExhausted { max_uses_per_sequence: u32 },
}
impl EscalationDeclineReason {
fn describe(self) -> String {
match self {
EscalationDeclineReason::NoCommandConfigured => {
"no acceptance_escalation_command is configured".to_string()
}
EscalationDeclineReason::BelowThreshold {
after_invalid_results,
} => format!(
"the configured threshold of {after_invalid_results} consecutive invalid results \
has not been reached"
),
EscalationDeclineReason::UseCapExhausted {
max_uses_per_sequence,
} => format!(
"the escalation budget of {max_uses_per_sequence} use(s) per invalid-result \
sequence is spent"
),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceEscalationOutcome {
NotObserved,
SequenceReset,
Escalated {
kind: InvalidAcceptanceResult,
consecutive_invalid: u32,
use_index: u32,
max_uses_per_sequence: u32,
},
Retained {
kind: InvalidAcceptanceResult,
consecutive_invalid: u32,
reason: EscalationDeclineReason,
},
}
impl AcceptanceEscalationOutcome {
pub fn escalation_selected(&self) -> bool {
matches!(self, AcceptanceEscalationOutcome::Escalated { .. })
}
pub fn diagnostic(&self) -> Option<String> {
match self {
AcceptanceEscalationOutcome::NotObserved
| AcceptanceEscalationOutcome::SequenceReset => None,
AcceptanceEscalationOutcome::Escalated {
kind,
consecutive_invalid,
use_index,
max_uses_per_sequence,
} => Some(format!(
"Acceptance produced an invalid result ({}); the next acceptance-only retry uses \
acceptance_escalation_command (invalid results: {}, escalation use \
{}/{}).",
kind.label(),
consecutive_invalid,
use_index,
max_uses_per_sequence
)),
AcceptanceEscalationOutcome::Retained {
kind,
consecutive_invalid,
reason,
} => Some(format!(
"Acceptance produced an invalid result ({}, consecutive invalid results: {}); \
keeping the normal acceptance_command because {}.",
kind.label(),
consecutive_invalid,
reason.describe()
)),
}
}
}
pub fn escalates_empty_fail(outcome: &AcceptanceEscalationOutcome) -> bool {
matches!(
outcome,
AcceptanceEscalationOutcome::Escalated {
kind: InvalidAcceptanceResult::EmptyFail,
..
}
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AcceptanceEscalationDriver {
command_configured: bool,
after_invalid_results: u32,
max_uses_per_sequence: u32,
consecutive_invalid: u32,
uses_in_sequence: u32,
pending_escalation: bool,
last_revision: Option<String>,
}
impl AcceptanceEscalationDriver {
pub fn from_config(config: &crate::config::OrchestratorConfig) -> Self {
let policy = config.get_acceptance_escalation();
Self {
command_configured: config.get_acceptance_escalation_command().is_some(),
after_invalid_results: policy.after_invalid_results(),
max_uses_per_sequence: policy.max_uses_per_sequence(),
consecutive_invalid: 0,
uses_in_sequence: 0,
pending_escalation: false,
last_revision: None,
}
}
pub fn consecutive_invalid(&self) -> u32 {
self.consecutive_invalid
}
pub fn uses_in_sequence(&self) -> u32 {
self.uses_in_sequence
}
pub fn take_command_mode(&mut self) -> AcceptanceCommandMode {
if std::mem::take(&mut self.pending_escalation) {
AcceptanceCommandMode::Escalation
} else {
AcceptanceCommandMode::Normal
}
}
pub fn observe(
&mut self,
result: &AcceptanceResult,
revision: Option<&str>,
) -> AcceptanceEscalationOutcome {
if let Some(revision) = revision {
if self
.last_revision
.as_deref()
.is_some_and(|previous| previous != revision)
{
self.reset_sequence();
}
self.last_revision = Some(revision.to_string());
}
if !result.permits_acceptance_retry()
|| matches!(result, AcceptanceResult::CommandFailed { .. })
{
return AcceptanceEscalationOutcome::NotObserved;
}
let Some(kind) = classify_invalid_acceptance_result(result) else {
self.reset_sequence();
return AcceptanceEscalationOutcome::SequenceReset;
};
self.consecutive_invalid = self.consecutive_invalid.saturating_add(1);
let decline = if !self.command_configured {
Some(EscalationDeclineReason::NoCommandConfigured)
} else if self.consecutive_invalid < self.after_invalid_results {
Some(EscalationDeclineReason::BelowThreshold {
after_invalid_results: self.after_invalid_results,
})
} else if self.uses_in_sequence >= self.max_uses_per_sequence {
Some(EscalationDeclineReason::UseCapExhausted {
max_uses_per_sequence: self.max_uses_per_sequence,
})
} else {
None
};
match decline {
Some(reason) => AcceptanceEscalationOutcome::Retained {
kind,
consecutive_invalid: self.consecutive_invalid,
reason,
},
None => {
self.uses_in_sequence = self.uses_in_sequence.saturating_add(1);
self.pending_escalation = true;
AcceptanceEscalationOutcome::Escalated {
kind,
consecutive_invalid: self.consecutive_invalid,
use_index: self.uses_in_sequence,
max_uses_per_sequence: self.max_uses_per_sequence,
}
}
}
}
fn reset_sequence(&mut self) {
self.consecutive_invalid = 0;
self.uses_in_sequence = 0;
self.pending_escalation = false;
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceBlockerDecision {
ProtocolRetry {
retry: AcceptanceProtocolRetry,
progress: String,
},
ProtocolExhausted { error: String },
ExternalBlocker {
blocker: crate::acceptance::AcceptanceBlocker,
},
}
pub fn decide_acceptance_blocker(
driver: &mut AcceptanceProtocolDriver,
result: &AcceptanceResult,
) -> Option<AcceptanceBlockerDecision> {
match result {
AcceptanceResult::BareBlocker { rejection } => {
Some(match driver.observe_bare_blocker(rejection) {
MissingVerdictRetryStep::Retry { retry, progress } => {
AcceptanceBlockerDecision::ProtocolRetry { retry, progress }
}
MissingVerdictRetryStep::Exhausted { error } => {
AcceptanceBlockerDecision::ProtocolExhausted { error }
}
})
}
AcceptanceResult::MalformedFinding { rejection } => {
Some(match driver.observe_malformed_finding(rejection) {
MissingVerdictRetryStep::Retry { retry, progress } => {
AcceptanceBlockerDecision::ProtocolRetry { retry, progress }
}
MissingVerdictRetryStep::Exhausted { error } => {
AcceptanceBlockerDecision::ProtocolExhausted { error }
}
})
}
AcceptanceResult::Stalled { blocker } => {
driver.observe_canonical_verdict();
Some(AcceptanceBlockerDecision::ExternalBlocker {
blocker: blocker.clone(),
})
}
_ => None,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NormalizedFinding {
pub identity: String,
pub text: String,
pub external: bool,
pub finding: crate::acceptance::AcceptanceFinding,
}
pub fn normalize_findings(
findings: &[crate::acceptance::AcceptanceFinding],
) -> Vec<NormalizedFinding> {
let mut normalized = findings
.iter()
.filter_map(|finding| {
let structured = finding.structured_payload()?;
Some(NormalizedFinding {
identity: format!("repository|id|{}", structured.id),
text: finding.text().to_string(),
external: false,
finding: finding.clone(),
})
})
.collect::<Vec<_>>();
normalized.extend(normalize_legacy_findings(
findings
.iter()
.filter(|finding| finding.structured_payload().is_none()),
));
normalized.sort_by(|left, right| left.identity.cmp(&right.identity));
normalized.dedup_by(|left, right| left.identity == right.identity);
normalized
}
fn normalize_legacy_findings<'a, I>(findings: I) -> Vec<NormalizedFinding>
where
I: IntoIterator<Item = &'a crate::acceptance::AcceptanceFinding>,
{
let texts = findings
.into_iter()
.map(|finding| finding.text().to_string())
.collect::<Vec<_>>();
normalize_legacy_texts(&texts)
}
pub fn normalize_legacy_texts(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}"),
),
finding: crate::acceptance::AcceptanceFinding::legacy(normalized.clone()),
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: &[crate::acceptance::AcceptanceFinding],
) -> Vec<crate::acceptance::AcceptanceFinding> {
findings
.iter()
.filter(|finding| {
finding.structured_payload().is_some()
|| normalize_legacy_texts(&[finding.text().to_string()])
.first()
.is_some_and(|normalized| !normalized.external)
})
.cloned()
.collect()
}
pub fn semantic_progress_fingerprint(workspace: &std::path::Path) -> std::io::Result<String> {
fn task_file_format(path: &str) -> Option<crate::task_file::TaskFileFormat> {
path.rsplit('/')
.next()
.and_then(crate::task_file::TaskFileFormat::from_file_name)
}
fn include(path: &str) -> bool {
!path.starts_with(".git/")
&& !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")
|| task_file_format(path).is_some())
}
fn strip_runtime_follow_up(
format: crate::task_file::TaskFileFormat,
contents: Vec<u8>,
) -> Vec<u8> {
match format {
crate::task_file::TaskFileFormat::Markdown => {
let text = String::from_utf8_lossy(&contents);
text.split("\n## Current Acceptance Follow-up")
.next()
.unwrap_or(&text)
.split("\n## Acceptance #")
.next()
.unwrap_or(&text)
.as_bytes()
.to_vec()
}
crate::task_file::TaskFileFormat::Json => {
match serde_json::from_slice::<serde_json::Value>(&contents) {
Ok(serde_json::Value::Object(mut root)) => {
root.remove(crate::task_file::FOLLOW_UP_KEY);
serde_json::to_vec(&serde_json::Value::Object(root)).unwrap_or(contents)
}
_ => contents,
}
}
}
}
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 let Some(format) = task_file_format(&relative) {
contents = strip_runtime_follow_up(format, contents);
}
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 const REMEDIATION_MISMATCH_REASON: &str = "acceptance_remediation_mismatch";
pub const REPEATED_FINDING_REASON: &str = "repeated_acceptance_finding";
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RemediationCoverage {
pub changed_files: Vec<String>,
pub missing_required: Vec<(String, String)>,
pub missing_verification: Vec<(String, String)>,
pub unrelated_files: Vec<String>,
pub legacy_finding_texts: Vec<String>,
}
impl RemediationCoverage {
pub fn is_complete(&self) -> bool {
self.missing_required.is_empty() && self.missing_verification.is_empty()
}
pub fn uncovered(&self) -> Vec<String> {
self.missing_required
.iter()
.map(|(id, file)| format!("{id}: required_changes {file}"))
.chain(
self.missing_verification
.iter()
.map(|(id, file)| format!("{id}: verification {file}")),
)
.collect()
}
}
fn normalize_changed_path(raw: &str) -> Option<String> {
let trimmed = raw.trim();
if trimmed.is_empty() {
return None;
}
crate::acceptance::normalize_repository_path(trimmed)
}
pub fn normalize_changed_files<I, S>(paths: I) -> Vec<String>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
let mut files = paths
.into_iter()
.flat_map(|path| {
path.as_ref()
.split(" -> ")
.filter_map(normalize_changed_path)
.collect::<Vec<_>>()
})
.collect::<Vec<_>>();
files.sort();
files.dedup();
files
}
pub fn evaluate_remediation_coverage(
findings: &[crate::acceptance::AcceptanceFinding],
changed_files: &[String],
) -> RemediationCoverage {
let changed = normalize_changed_files(changed_files);
let mut coverage = RemediationCoverage {
changed_files: changed.clone(),
..RemediationCoverage::default()
};
let mut declared = std::collections::BTreeSet::new();
for finding in findings {
if !finding.declares_paths() {
coverage
.legacy_finding_texts
.push(finding.text().to_string());
continue;
}
let Some(structured) = finding.structured_payload() else {
continue;
};
for file in structured.required_files() {
declared.insert(file.clone());
if !changed.contains(&file) {
coverage
.missing_required
.push((structured.id.clone(), file));
}
}
for file in structured.verification_files() {
declared.insert(file.clone());
if !changed.contains(&file) {
coverage
.missing_verification
.push((structured.id.clone(), file));
}
}
}
coverage.unrelated_files = changed
.into_iter()
.filter(|file| !declared.contains(file))
.collect();
coverage
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct FindingRepairLedger {
repaired: std::collections::BTreeSet<String>,
occurrences: std::collections::BTreeMap<String, u32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum FindingRepairDecision {
Repair {
identities: Vec<String>,
},
Stop {
reason: &'static str,
repeated_identities: Vec<String>,
},
}
impl FindingRepairLedger {
pub fn occurrences(&self, identity: &str) -> u32 {
self.occurrences.get(identity).copied().unwrap_or(0)
}
pub fn occurrence_report(&self) -> Vec<(String, u32)> {
self.occurrences
.iter()
.map(|(identity, count)| (identity.clone(), *count))
.collect()
}
pub fn has_consumed_repair(&self, identity: &str) -> bool {
self.repaired.contains(identity)
}
pub fn observe_fail(&mut self, findings: &[NormalizedFinding]) -> FindingRepairDecision {
let identities = findings
.iter()
.filter(|finding| !finding.external)
.map(|finding| finding.identity.clone())
.collect::<Vec<_>>();
let open = identities
.iter()
.cloned()
.collect::<std::collections::BTreeSet<_>>();
self.occurrences
.retain(|identity, _| open.contains(identity));
self.repaired.retain(|identity| open.contains(identity));
for identity in &identities {
*self.occurrences.entry(identity.clone()).or_insert(0) += 1;
}
let repeated = identities
.iter()
.filter(|identity| self.repaired.contains(*identity))
.cloned()
.collect::<Vec<_>>();
if !repeated.is_empty() {
return FindingRepairDecision::Stop {
reason: REPEATED_FINDING_REASON,
repeated_identities: repeated,
};
}
FindingRepairDecision::Repair { identities }
}
pub fn record_repair_dispatched(&mut self, identities: &[String]) {
for identity in identities {
self.repaired.insert(identity.clone());
}
}
pub fn reset_for_explicit_retry(&mut self) {
self.repaired.clear();
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AcceptanceRepairStop {
pub change_id: String,
pub reason: &'static str,
pub findings: Vec<crate::acceptance::AcceptanceFinding>,
pub occurrences: Vec<(String, u32)>,
pub repeated_identities: Vec<String>,
pub fail_revision: Option<String>,
pub apply_revision: Option<String>,
pub required_files: Vec<String>,
pub verification_files: Vec<String>,
pub coverage: RemediationCoverage,
pub remediation_evidence: Vec<String>,
pub resumable: bool,
pub next_action: String,
}
impl AcceptanceRepairStop {
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"change_id": self.change_id,
"stop_reason": self.reason,
"findings": self.findings.iter().map(|finding| finding.to_json()).collect::<Vec<_>>(),
"finding_occurrences": self
.occurrences
.iter()
.map(|(identity, count)| serde_json::json!({"identity": identity, "occurrences": count}))
.collect::<Vec<_>>(),
"repeated_identities": self.repeated_identities,
"fail_revision": self.fail_revision,
"apply_revision": self.apply_revision,
"required_files": self.required_files,
"verification_files": self.verification_files,
"changed_files": self.coverage.changed_files,
"uncovered_files": self.coverage.uncovered(),
"coverage_complete": self.coverage.is_complete(),
"unrelated_files": self.coverage.unrelated_files,
"legacy_findings_without_declared_paths": self.coverage.legacy_finding_texts,
"remediation_evidence": self.remediation_evidence,
"resumable": self.resumable,
"next_action": self.next_action,
"proves_completion": false,
"proves_acceptance_pass": false,
"proves_archive_readiness": false,
})
}
pub fn occurrence_identities(&self) -> Vec<String> {
self.occurrences
.iter()
.map(|(identity, _)| identity.clone())
.collect()
}
pub fn summary(&self) -> String {
format!(
"Acceptance stopped automatic repair for {} ({}). {} Diagnostics: {}",
self.change_id,
self.reason,
self.next_action,
self.to_json()
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RepairGateDecision {
Proceed,
Stop(Box<AcceptanceRepairStop>),
}
pub fn decide_repair_gate(
change_id: &str,
findings: &[crate::acceptance::AcceptanceFinding],
ledger: &FindingRepairLedger,
fail_revision: Option<&str>,
apply_revision: Option<&str>,
changed_files: &[String],
remediation_evidence: &[String],
) -> RepairGateDecision {
let coverage = evaluate_remediation_coverage(findings, changed_files);
if coverage.is_complete() {
return RepairGateDecision::Proceed;
}
let uncovered = coverage.uncovered().join(", ");
RepairGateDecision::Stop(Box::new(AcceptanceRepairStop {
change_id: change_id.to_string(),
reason: REMEDIATION_MISMATCH_REASON,
findings: findings.to_vec(),
occurrences: ledger.occurrence_report(),
repeated_identities: Vec::new(),
fail_revision: fail_revision.map(str::to_string),
apply_revision: apply_revision.map(str::to_string),
required_files: declared_files(findings, true),
verification_files: declared_files(findings, false),
coverage,
remediation_evidence: remediation_evidence.to_vec(),
resumable: true,
next_action: format!(
"Change the declared files still missing from the repair delta ({uncovered}), then \
retry explicitly. Unrelated or comment-only changes cannot satisfy the finding \
contract, and coverage alone never proves the finding is resolved."
),
}))
}
#[allow(clippy::too_many_arguments)]
pub fn repeated_finding_stop(
change_id: &str,
findings: &[crate::acceptance::AcceptanceFinding],
ledger: &FindingRepairLedger,
repeated_identities: Vec<String>,
fail_revision: Option<&str>,
apply_revision: Option<&str>,
changed_files: &[String],
remediation_evidence: &[String],
) -> AcceptanceRepairStop {
AcceptanceRepairStop {
change_id: change_id.to_string(),
reason: REPEATED_FINDING_REASON,
findings: findings.to_vec(),
occurrences: ledger.occurrence_report(),
next_action: format!(
"Finding {} is still open after its one automatic repair. Review the finding and the \
recorded remediation evidence, then retry explicitly; unrelated changed files do not \
grant another automatic repair.",
repeated_identities.join(", ")
),
repeated_identities,
fail_revision: fail_revision.map(str::to_string),
apply_revision: apply_revision.map(str::to_string),
required_files: declared_files(findings, true),
verification_files: declared_files(findings, false),
coverage: evaluate_remediation_coverage(findings, changed_files),
remediation_evidence: remediation_evidence.to_vec(),
resumable: true,
}
}
fn declared_files(
findings: &[crate::acceptance::AcceptanceFinding],
required: bool,
) -> Vec<String> {
let mut files = findings
.iter()
.filter_map(|finding| finding.structured_payload())
.flat_map(|finding| {
if required {
finding.required_files()
} else {
finding.verification_files()
}
})
.collect::<Vec<_>>();
files.sort();
files.dedup();
files
}
pub async fn collect_repair_gate_inputs(
workspace: &std::path::Path,
change_id: &str,
fail_revision: Option<&str>,
) -> (Vec<String>, Option<String>, Vec<String>) {
let changed_files = match fail_revision {
Some(revision) => crate::vcs::git::commands::get_changed_files_since(workspace, revision)
.await
.unwrap_or_default(),
None => Vec::new(),
};
let apply_revision = crate::vcs::git::commands::get_current_commit(workspace)
.await
.ok();
let remediation_evidence =
crate::task_parser::resolve_acceptance_follow_up_tasks_path(change_id, workspace)
.ok()
.and_then(|path| crate::task_parser::read_acceptance_follow_up_evidence(&path).ok())
.unwrap_or_default();
(changed_files, apply_revision, remediation_evidence)
}
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 struct AcceptanceRuntimeLimit {
pub limit_secs: u64,
pub cleanup_confirmed: bool,
pub cleanup_diagnostics: String,
}
impl AcceptanceRuntimeLimit {
pub fn permits_retry(&self) -> bool {
false
}
pub fn summary(&self, change_id: &str) -> String {
let cleanup = if self.cleanup_confirmed {
"owned process group confirmed quiescent".to_string()
} else {
format!(
"owned process group NOT confirmed quiescent — {}",
self.cleanup_diagnostics
)
};
format!(
"Acceptance for '{}' was terminated by its absolute runtime limit of {}s \
(`acceptance_max_runtime_secs`); {}. This invocation is not retried automatically: \
it is not a PASS, not an external block, not an inactivity timeout, and not a \
missing-verdict continuation. Reduce the proposal's Acceptance scope or raise \
`acceptance_max_runtime_secs`, then retry explicitly.",
change_id, self.limit_secs, cleanup
)
}
}
pub fn classify_acceptance_runtime_limit(
termination: crate::process_manager::CommandTermination,
limit_secs: u64,
cleanup: &crate::process_manager::ProcessGroupCleanupReport,
) -> Option<AcceptanceRuntimeLimit> {
if !termination.is_runtime_limit() {
return None;
}
Some(AcceptanceRuntimeLimit {
limit_secs,
cleanup_confirmed: cleanup.is_confirmed(),
cleanup_diagnostics: cleanup.diagnostics(),
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AcceptanceResult {
Pass,
Fail {
findings: Vec<crate::acceptance::AcceptanceFinding>,
},
Continue,
MalformedFinding {
rejection: crate::acceptance::FindingRejection,
},
BareBlocker {
rejection: crate::acceptance::BlockerRejection,
},
Stalled {
blocker: crate::acceptance::AcceptanceBlocker,
},
CommandFailed {
error: String,
findings: Vec<String>,
diagnostic: AcceptanceCommandDiagnostic,
},
PermissionStalled {
blocker: crate::events::StalledBlocker,
},
MissingVerdict { findings: Vec<String> },
RuntimeLimit { limit: AcceptanceRuntimeLimit },
ExecutionHold {
hold: execution_manifest::AcceptanceExecutionHold,
},
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::Stalled { .. }
| AcceptanceResult::PermissionStalled { .. }
)
}
pub fn is_runtime_limit(&self) -> bool {
matches!(self, AcceptanceResult::RuntimeLimit { .. })
}
pub fn permits_acceptance_retry(&self) -> bool {
!matches!(
self,
AcceptanceResult::RuntimeLimit { .. }
| AcceptanceResult::Cancelled
| AcceptanceResult::ExecutionHold { .. }
)
}
pub fn is_execution_hold(&self) -> bool {
matches!(self, AcceptanceResult::ExecutionHold { .. })
}
}
#[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<AcceptanceProtocolRetry>,
previous_denial: Option<&crate::permission::PermissionDenial>,
command_mode: AcceptanceCommandMode,
) -> 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,
command_mode,
)
.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()
);
output.on_error(&error_msg);
let diagnostic = AcceptanceCommandDiagnostic {
error: error_msg.clone(),
exit_code: status.code(),
stdout_tail,
stderr_tail,
};
if let Some(denial) = classify_acceptance_command_denial(&diagnostic) {
let end_commit_hash = crate::vcs::git::commands::get_current_commit(".")
.await
.ok();
let revision_unchanged = commit_hash == end_commit_hash;
let repeated_unresolved = previous_denial
.is_some_and(|previous| previous.signature() == denial.signature())
&& revision_unchanged;
if repeated_unresolved {
warn!(
change_id = %change.id,
category = denial.category.as_str(),
denied_target = %denial.denied_target,
"Repeated unresolved permission/tool policy denial during acceptance; entering non-terminal hold without a corrective command retry"
);
return Ok((
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::permission_denial(
"acceptance",
&denial,
),
},
0,
command,
));
}
}
return Ok((
AcceptanceResult::CommandFailed {
error: error_msg,
findings: tail_findings,
diagnostic,
},
0,
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() {
crate::acceptance::legacy_findings([GENERIC_ACCEPTANCE_FAIL_FINDING])
} 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::Stalled { blocker } => {
info!(
"Acceptance reported a validated external blocker for: {} (category {})",
change.id, blocker.category
);
output.on_warn(&format!(
"Acceptance test: STALLED ({}) - {}",
blocker.category, blocker.next_action
));
(AcceptanceResult::Stalled { blocker }, false)
}
crate::acceptance::AcceptanceResult::MalformedFinding { rejection } => {
warn!(
"Acceptance emitted a malformed structured finding for: {} ({})",
change.id,
rejection.reason()
);
output.on_error(&format!(
"Acceptance test: MALFORMED FINDING (protocol failure — {})",
rejection.reason()
));
(AcceptanceResult::MalformedFinding { rejection }, false)
}
crate::acceptance::AcceptanceResult::BareBlocker { rejection } => {
warn!(
"Acceptance emitted a bare blocker compatibility token for: {} ({})",
change.id,
rejection.reason()
);
output.on_error(&format!(
"Acceptance test: BARE BLOCKER (protocol failure — {})",
rejection.reason()
));
(AcceptanceResult::BareBlocker { rejection }, 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(crate::acceptance::legacy_findings([
"Investigation incomplete - continue later",
])),
AcceptanceResult::Stalled { blocker } => {
let mut evidence = vec![format!(
"Validated external blocker (category {}): {}",
blocker.category, blocker.next_action
)];
evidence.extend(blocker.evidence.iter().cloned());
Some(crate::acceptance::legacy_findings(evidence))
}
AcceptanceResult::BareBlocker { rejection } => Some(crate::acceptance::legacy_findings([
BARE_BLOCKER_DIAGNOSTIC.to_string(),
rejection.reason(),
])),
AcceptanceResult::MalformedFinding { rejection } => {
Some(crate::acceptance::legacy_findings([
MALFORMED_FINDING_DIAGNOSTIC.to_string(),
rejection.reason(),
]))
}
AcceptanceResult::Pass => None,
AcceptanceResult::MissingVerdict { findings } => {
let mut evidence = vec![MISSING_VERDICT_DIAGNOSTIC.to_string()];
evidence.extend(findings.iter().cloned());
Some(crate::acceptance::legacy_findings(evidence))
}
AcceptanceResult::RuntimeLimit { .. } => None,
AcceptanceResult::ExecutionHold { .. } => None,
AcceptanceResult::CommandFailed { .. }
| AcceptanceResult::PermissionStalled { .. }
| AcceptanceResult::Cancelled => {
Some(crate::acceptance::legacy_findings(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::*;
mod runtime_limit {
use super::super::*;
use crate::process_manager::{
CommandTermination, ProcessGroupCleanupReport, ProcessGroupQuiescence,
};
const ACCEPTANCE_DEFAULT: u64 = 1_800;
const COMMON_DEFAULT: u64 = 10_800;
fn confirmed_cleanup() -> ProcessGroupCleanupReport {
ProcessGroupCleanupReport::for_test(
ProcessGroupQuiescence::Confirmed,
Some(4242),
"group empty after SIGTERM",
)
}
fn unprovable_cleanup() -> ProcessGroupCleanupReport {
ProcessGroupCleanupReport::for_test(
ProcessGroupQuiescence::Unverifiable,
Some(4242),
"a SIGTERM-immune descendant outlived the verification budget",
)
}
#[test]
fn routing_bounds_acceptance_without_moving_other_command_classes() {
use crate::ai_command_runner::AiCommandRunner;
use crate::command_queue::ACCEPTANCE_OPERATION_TYPE;
use crate::config::OrchestratorConfig;
use std::sync::Arc;
use tokio::sync::Mutex;
let config = OrchestratorConfig {
command_max_runtime_secs: Some(COMMON_DEFAULT),
acceptance_max_runtime_secs: Some(ACCEPTANCE_DEFAULT),
command_inactivity_timeout_secs: Some(900),
..OrchestratorConfig::default()
};
let runner =
AiCommandRunner::from_orchestrator_config(&config, Arc::new(Mutex::new(None)));
let queue = runner.queue_config();
assert_eq!(
queue.effective_max_runtime_secs(Some(ACCEPTANCE_OPERATION_TYPE)),
ACCEPTANCE_DEFAULT,
"the Acceptance invocation runs under the dedicated limit"
);
for class in [
Some("apply"),
Some("archive"),
Some("analyze"),
Some("resolve"),
Some("cleanup-review"),
None,
] {
assert_eq!(
queue.effective_max_runtime_secs(class),
COMMON_DEFAULT,
"{class:?} keeps `command_max_runtime_secs`"
);
}
assert_eq!(
queue.max_runtime_secs, COMMON_DEFAULT,
"the stored common limit is never rewritten for one command class"
);
}
#[test]
fn only_runtime_limit_termination_classifies_as_one() {
for termination in [
CommandTermination::Exited,
CommandTermination::InactivityTimeout,
CommandTermination::Cancelled,
CommandTermination::NotStarted,
] {
assert!(
classify_acceptance_runtime_limit(
termination,
ACCEPTANCE_DEFAULT,
&confirmed_cleanup(),
)
.is_none(),
"{} must keep its own routing",
termination.as_str()
);
}
let limit = classify_acceptance_runtime_limit(
CommandTermination::RuntimeLimit,
ACCEPTANCE_DEFAULT,
&confirmed_cleanup(),
)
.expect("runtime-limit termination classifies as a runtime limit");
assert_eq!(limit.limit_secs, ACCEPTANCE_DEFAULT);
assert!(limit.cleanup_confirmed);
}
#[test]
fn expiry_is_terminal_for_the_invocation_and_says_so() {
let limit = classify_acceptance_runtime_limit(
CommandTermination::RuntimeLimit,
ACCEPTANCE_DEFAULT,
&confirmed_cleanup(),
)
.expect("classified");
assert!(
!limit.permits_retry(),
"a timed-out Acceptance invocation is never retried automatically"
);
let summary = limit.summary("change-a");
assert!(summary.contains("1800s"), "the configured limit: {summary}");
assert!(
summary.contains("acceptance_max_runtime_secs"),
"the knob an operator has to change: {summary}"
);
assert!(
summary.contains("not retried automatically"),
"the retry decision must be visible: {summary}"
);
assert!(
summary.contains("change-a"),
"the affected proposal: {summary}"
);
}
#[test]
fn cleanup_failure_is_not_hidden() {
let limit = classify_acceptance_runtime_limit(
CommandTermination::RuntimeLimit,
ACCEPTANCE_DEFAULT,
&unprovable_cleanup(),
)
.expect("classified");
assert!(
!limit.cleanup_confirmed,
"an unswept owned group must never count as quiescent"
);
assert!(
limit.cleanup_diagnostics.contains("SIGTERM-immune"),
"the runner's own evidence must survive: {}",
limit.cleanup_diagnostics
);
let summary = limit.summary("change-a");
assert!(
summary.contains("NOT confirmed quiescent"),
"unproven cleanup must be stated, not implied: {summary}"
);
assert!(
summary.contains("SIGTERM-immune"),
"the actionable cleanup diagnostics must reach the operator: {summary}"
);
assert!(
!limit.permits_retry(),
"unprovable cleanup never converts the limit into a retry"
);
}
#[test]
fn runtime_limit_is_its_own_terminal_result() {
let limit = classify_acceptance_runtime_limit(
CommandTermination::RuntimeLimit,
ACCEPTANCE_DEFAULT,
&confirmed_cleanup(),
)
.expect("classified");
let result = AcceptanceResult::RuntimeLimit { limit };
assert!(result.is_runtime_limit());
assert!(
!result.is_pass(),
"a terminated invocation can never be a PASS"
);
assert!(
!result.is_canonical_verdict(),
"no verdict was emitted, so no verdict protocol may consume it"
);
assert!(
!result.permits_acceptance_retry(),
"the invocation that the limit stopped is not re-dispatched"
);
assert!(
!matches!(
result,
AcceptanceResult::MissingVerdict { .. }
| AcceptanceResult::CommandFailed { .. }
| AcceptanceResult::Stalled { .. }
| AcceptanceResult::BareBlocker { .. }
| AcceptanceResult::PermissionStalled { .. }
| AcceptanceResult::Cancelled
),
"runtime-limit expiry is not a missing verdict, a command failure, \
an external block, or a cancellation"
);
}
#[test]
fn cancellation_and_runtime_limit_stay_distinct() {
assert!(!AcceptanceResult::Cancelled.is_runtime_limit());
assert!(!AcceptanceResult::Cancelled.permits_acceptance_retry());
assert!(
AcceptanceResult::Pass.permits_acceptance_retry(),
"only deliberate terminations close acceptance re-dispatch"
);
assert!(AcceptanceResult::MissingVerdict {
findings: Vec::new()
}
.permits_acceptance_retry());
}
}
mod command_recovery {
use super::super::*;
fn diagnostic(marker: &str) -> AcceptanceCommandDiagnostic {
AcceptanceCommandDiagnostic {
error: format!("Acceptance command failed with exit code: Some(1) [{marker}]"),
exit_code: Some(1),
stdout_tail: Some(format!("stdout {marker}")),
stderr_tail: Some(format!("stderr {marker}")),
}
}
fn agent() -> AgentRunner {
AgentRunner::new(crate::config::OrchestratorConfig::default())
}
#[test]
fn two_retries_follow_the_initial_failure_and_the_third_exhausts() {
let mut counter = AcceptanceCommandRetryCounter::default();
assert!(matches!(
counter.record_command_failure(diagnostic("a")),
AcceptanceCommandRetryDecision::Retry {
attempt: 1,
max_retries: 2,
..
}
));
assert!(matches!(
counter.record_command_failure(diagnostic("b")),
AcceptanceCommandRetryDecision::Retry {
attempt: 2,
max_retries: 2,
..
}
));
let third = counter.record_command_failure(diagnostic("c"));
let AcceptanceCommandRetryDecision::Exhausted {
attempts,
max_retries,
diagnostic,
} = third
else {
panic!("the third consecutive command failure must exhaust the budget");
};
assert_eq!((attempts, max_retries), (3, MAX_ACCEPTANCE_COMMAND_RETRIES));
assert!(
diagnostic.error.contains("[c]"),
"exhaustion reports the latest evidence, not the first: {diagnostic:?}"
);
}
#[test]
fn the_terminal_error_carries_attempt_count_and_latest_bounded_evidence() {
let mut counter = AcceptanceCommandRetryCounter::default();
counter.record_command_failure(diagnostic("a"));
counter.record_command_failure(diagnostic("b"));
let AcceptanceCommandRetryDecision::Exhausted {
attempts,
max_retries,
diagnostic,
} = counter.record_command_failure(diagnostic("c"))
else {
panic!("expected exhaustion");
};
let error = acceptance_command_exhausted_error(attempts, max_retries, &diagnostic);
assert!(error.contains("3 consecutive attempts"), "{error}");
assert!(error.contains("2 command-failure retries"), "{error}");
assert!(error.contains("exit_code: 1"), "{error}");
assert!(error.contains("stderr c"), "{error}");
assert!(
!error.contains("stderr a"),
"only the latest evidence is reported: {error}"
);
}
#[test]
fn the_terminal_error_reports_exit_code_and_both_distinct_bounded_tails() {
let diagnostic = AcceptanceCommandDiagnostic {
error: "Acceptance command failed with exit code: Some(42)".to_string(),
exit_code: Some(42),
stdout_tail: Some(format!("the-actual-diagnosis {}", "s".repeat(10_000))),
stderr_tail: Some(format!("unrelated-warning {}", "e".repeat(10_000))),
};
let error =
acceptance_command_exhausted_error(3, MAX_ACCEPTANCE_COMMAND_RETRIES, &diagnostic);
assert!(error.contains("exit_code: 42"), "{error}");
assert!(
error.contains("stdout: ") && error.contains("the-actual-diagnosis"),
"the stdout tail must survive a present stderr: {error}"
);
assert!(
error.contains("stderr: ") && error.contains("unrelated-warning"),
"the stderr tail must be reported alongside stdout: {error}"
);
assert!(
error.len() < 1_200,
"both tails must stay bounded, got {} chars",
error.len()
);
}
#[test]
fn retry_progress_also_reports_both_streams() {
let diagnostic = AcceptanceCommandDiagnostic {
error: "Acceptance command failed with exit code: Some(7)".to_string(),
exit_code: Some(7),
stdout_tail: Some("stdout-evidence".to_string()),
stderr_tail: Some("stderr-evidence".to_string()),
};
let progress =
acceptance_command_retry_progress(1, MAX_ACCEPTANCE_COMMAND_RETRIES, &diagnostic);
assert!(progress.contains("stdout-evidence"), "{progress}");
assert!(progress.contains("stderr-evidence"), "{progress}");
}
#[test]
fn any_completed_non_command_failure_result_resets_the_sequence() {
let mut counter = AcceptanceCommandRetryCounter::default();
let mut agent = agent();
for _ in 0..3 {
counter.record_command_failure(diagnostic("a"));
counter.record_command_failure(diagnostic("b"));
assert_eq!(counter.consecutive_failures(), 2);
observe_completed_acceptance_invocation(&mut counter, &mut agent, "change-a");
assert_eq!(counter.consecutive_failures(), 0);
assert!(agent.acceptance_command_recovery("change-a").is_none());
}
}
#[test]
fn protocol_completion_breaks_command_failure_consecutiveness() {
let mut counter = AcceptanceCommandRetryCounter::default();
let mut protocol = AcceptanceProtocolDriver::default();
let mut agent = agent();
assert!(matches!(
counter.record_command_failure(diagnostic("first")),
AcceptanceCommandRetryDecision::Retry { attempt: 1, .. }
));
observe_completed_acceptance_invocation(&mut counter, &mut agent, "change-a");
let missing = protocol.observe_missing_verdict(&["waiting".to_string()]);
assert!(matches!(missing, MissingVerdictRetryStep::Retry { .. }));
assert!(
matches!(
counter.record_command_failure(diagnostic("second")),
AcceptanceCommandRetryDecision::Retry { attempt: 1, .. }
),
"the final command failure starts a new sequence"
);
assert_eq!(
protocol.consecutive_missing_verdicts(),
1,
"the missing-verdict budget keeps its own independent accounting"
);
}
#[test]
fn budgets_stay_independent_when_one_is_exhausted() {
let mut counter = AcceptanceCommandRetryCounter::default();
let mut protocol = AcceptanceProtocolDriver::default();
for _ in 0..3 {
counter.record_command_failure(diagnostic("boom"));
}
assert_eq!(counter.consecutive_failures(), 3);
assert_eq!(
protocol.consecutive_missing_verdicts(),
0,
"command failures never consume the protocol budget"
);
assert_eq!(protocol.consecutive_bare_blockers(), 0);
assert_eq!(protocol.consecutive_malformed_findings(), 0);
assert!(
protocol.take_protocol_retry().is_none(),
"command failures never schedule a protocol continuation"
);
}
#[test]
fn the_shared_policy_stores_latest_only_prompt_context_and_clears_on_exhaustion() {
let mut counter = AcceptanceCommandRetryCounter::default();
let mut agent = agent();
let first = decide_acceptance_command_failure(
&mut counter,
&mut agent,
"change-a",
diagnostic("a"),
);
assert!(matches!(
first,
AcceptanceCommandRecovery::Retry { attempt: 1, .. }
));
assert!(agent
.acceptance_command_recovery("change-a")
.expect("a retry stores the latest diagnosis")
.error
.contains("[a]"));
decide_acceptance_command_failure(
&mut counter,
&mut agent,
"change-a",
diagnostic("b"),
);
let stored = agent
.acceptance_command_recovery("change-a")
.expect("still recovering");
assert!(stored.error.contains("[b]"), "{stored:?}");
assert!(
!stored.error.contains("[a]"),
"prior command failures are never replayed: {stored:?}"
);
let third = decide_acceptance_command_failure(
&mut counter,
&mut agent,
"change-a",
diagnostic("c"),
);
let AcceptanceCommandRecovery::Exhausted { attempts, error } = third else {
panic!("the third consecutive failure must exhaust");
};
assert_eq!(attempts, 3);
assert!(error.contains("[c]"), "{error}");
assert!(
agent.acceptance_command_recovery("change-a").is_none(),
"no further command retry follows exhaustion, so the context is cleared"
);
}
#[test]
fn command_recovery_context_is_tracked_per_change() {
let mut counter_a = AcceptanceCommandRetryCounter::default();
let mut counter_b = AcceptanceCommandRetryCounter::default();
let mut agent = agent();
decide_acceptance_command_failure(
&mut counter_a,
&mut agent,
"change-a",
diagnostic("a"),
);
decide_acceptance_command_failure(
&mut counter_b,
&mut agent,
"change-b",
diagnostic("b"),
);
observe_completed_acceptance_invocation(&mut counter_a, &mut agent, "change-a");
assert!(agent.acceptance_command_recovery("change-a").is_none());
assert!(
agent.acceptance_command_recovery("change-b").is_some(),
"one change's reset must not clear another change's context"
);
}
#[test]
fn the_progress_message_is_worded_as_recovery_not_failure() {
let mut counter = AcceptanceCommandRetryCounter::default();
let AcceptanceCommandRetryDecision::Retry {
attempt,
max_retries,
diagnostic,
} = counter.record_command_failure(diagnostic("a"))
else {
panic!("expected a retry");
};
let progress = acceptance_command_retry_progress(attempt, max_retries, &diagnostic);
assert!(
progress.contains("command-failure recovery 1/2"),
"{progress}"
);
assert!(
progress.contains("re-running only the configured acceptance command"),
"{progress}"
);
}
}
fn sample_blocker() -> crate::acceptance::AcceptanceBlocker {
crate::acceptance::AcceptanceBlocker {
category: "pending_verification".to_string(),
evidence: vec!["managed verification job 42 is still running".to_string()],
unblock_condition: "managed verification job 42 reports a terminal result".to_string(),
next_action: "wait for job 42 then retry acceptance".to_string(),
resumable: true,
prerequisite_owner: None,
evidence_ids: Vec::new(),
}
}
fn bare_blocker() -> AcceptanceResult {
AcceptanceResult::BareBlocker {
rejection: crate::acceptance::BlockerRejection::Missing,
}
}
#[test]
fn missing_verdict_budget_allows_two_retries_then_exhausts() {
let mut state = ProtocolRetryCounter::default();
let expected = [
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
attempt: 1,
max: 2,
}),
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
attempt: 2,
max: 2,
}),
MissingVerdictRetryDecision::Exhausted {
kind: AcceptanceProtocolError::MissingVerdict,
attempts: 3,
max: 2,
},
];
for (index, want) in expected.iter().enumerate() {
assert_eq!(
state.record(AcceptanceProtocolError::MissingVerdict),
*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(AcceptanceProtocolError::MissingVerdict),
MissingVerdictRetryDecision::Exhausted { .. }
));
}
#[test]
fn missing_verdict_canonical_verdict_resets_consecutive_sequence() {
let mut state = ProtocolRetryCounter::default();
assert!(matches!(
state.record(AcceptanceProtocolError::MissingVerdict),
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry { attempt: 1, .. })
));
assert!(matches!(
state.record(AcceptanceProtocolError::MissingVerdict),
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry { attempt: 2, .. })
));
state.reset();
assert_eq!(state.consecutive(), 0);
assert_eq!(
state.record(AcceptanceProtocolError::MissingVerdict),
MissingVerdictRetryDecision::Retry(AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
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 = ProtocolRetryCounter::default();
let mut retries = 0u32;
loop {
match state.record(AcceptanceProtocolError::MissingVerdict) {
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(
AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
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<AcceptanceProtocolRetry>>,
Option<AcceptanceResult>,
Option<String>,
) {
let mut protocol = AcceptanceProtocolDriver::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(AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
attempt: 1,
max: 2
}),
Some(AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
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().into()],
},
AcceptanceResult::Continue,
AcceptanceResult::Stalled {
blocker: sample_blocker(),
},
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_external(
"pending_verification",
"denied",
),
},
] {
let mut protocol = AcceptanceProtocolDriver::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: AcceptanceProtocolRetry { attempt: 1, .. },
..
}
));
}
}
fn drive_blockers(sequence: &[AcceptanceResult]) -> Vec<AcceptanceBlockerDecision> {
let mut driver = AcceptanceProtocolDriver::default();
let mut decisions = Vec::new();
for result in sequence {
if let Some(decision) = decide_acceptance_blocker(&mut driver, result) {
decisions.push(decision);
} else {
driver.observe_canonical_verdict();
}
}
decisions
}
#[test]
fn bare_blocker_budget_allows_two_retries_then_exhausts() {
let decisions = drive_blockers(&[bare_blocker(), bare_blocker(), bare_blocker()]);
assert_eq!(decisions.len(), 3);
for (index, expected_attempt) in [1u32, 2].iter().enumerate() {
match &decisions[index] {
AcceptanceBlockerDecision::ProtocolRetry { retry, progress } => {
assert_eq!(retry.kind, AcceptanceProtocolError::BareBlocker);
assert_eq!(retry.attempt, *expected_attempt);
assert_eq!(retry.max, MAX_ACCEPTANCE_PROTOCOL_RETRIES);
assert!(progress.contains("no validated blocker"), "{progress}");
assert!(
progress.contains("retrying acceptance")
&& progress.contains(&format!("protocol retry {expected_attempt}/2")),
"an available retry must read as progress: {progress}"
);
assert!(
!progress.contains("Exhausted"),
"an available retry must be distinguishable from the terminal \
diagnostic: {progress}"
);
}
other => panic!("attempt {expected_attempt} must retry, got {other:?}"),
}
}
match &decisions[2] {
AcceptanceBlockerDecision::ProtocolExhausted { error } => {
assert!(error.contains("bare-blocker protocol failure"), "{error}");
assert!(
error.contains("Exhausted 3 consecutive attempts after 2 protocol retries"),
"{error}"
);
}
other => panic!("third consecutive bare blocker must be terminal, got {other:?}"),
}
}
#[test]
fn canonical_verdict_resets_the_bare_blocker_sequence() {
for canonical in [
AcceptanceResult::Pass,
AcceptanceResult::Fail {
findings: vec!["src/lib.rs:1 fix".to_string().into()],
},
AcceptanceResult::Continue,
AcceptanceResult::Stalled {
blocker: sample_blocker(),
},
] {
let mut driver = AcceptanceProtocolDriver::default();
for _ in 0..2 {
assert!(matches!(
decide_acceptance_blocker(&mut driver, &bare_blocker()),
Some(AcceptanceBlockerDecision::ProtocolRetry { .. })
));
}
assert_eq!(driver.consecutive_bare_blockers(), 2);
if decide_acceptance_blocker(&mut driver, &canonical).is_none() {
driver.observe_canonical_verdict();
}
assert_eq!(
driver.consecutive_bare_blockers(),
0,
"{canonical:?} must reset the consecutive bare-blocker counter"
);
assert!(
driver.take_protocol_retry().is_none(),
"{canonical:?} must clear any pending continuation marker"
);
assert!(
matches!(
decide_acceptance_blocker(&mut driver, &bare_blocker()),
Some(AcceptanceBlockerDecision::ProtocolRetry {
retry: AcceptanceProtocolRetry { attempt: 1, .. },
..
})
),
"a later bare blocker must start a fresh full budget"
);
}
}
#[test]
fn bare_blocker_and_missing_verdict_budgets_are_independent_but_both_bounded() {
let mut driver = AcceptanceProtocolDriver::default();
assert!(matches!(
decide_acceptance_blocker(&mut driver, &bare_blocker()),
Some(AcceptanceBlockerDecision::ProtocolRetry { .. })
));
assert!(matches!(
driver.observe_missing_verdict(&["waiting".to_string()]),
MissingVerdictRetryStep::Retry { .. }
));
assert_eq!(driver.consecutive_bare_blockers(), 1);
assert_eq!(driver.consecutive_missing_verdicts(), 1);
assert!(matches!(
decide_acceptance_blocker(&mut driver, &bare_blocker()),
Some(AcceptanceBlockerDecision::ProtocolRetry { .. })
));
assert!(matches!(
decide_acceptance_blocker(&mut driver, &bare_blocker()),
Some(AcceptanceBlockerDecision::ProtocolExhausted { .. })
));
assert!(matches!(
driver.observe_missing_verdict(&["waiting".to_string()]),
MissingVerdictRetryStep::Retry { .. }
));
assert!(matches!(
driver.observe_missing_verdict(&["waiting".to_string()]),
MissingVerdictRetryStep::Exhausted { .. }
));
}
#[test]
fn every_invalid_blocker_payload_routes_to_bounded_protocol_retry() {
use crate::acceptance::BlockerRejection;
let rejections = [
BlockerRejection::Missing,
BlockerRejection::NotAnObject,
BlockerRejection::MissingCategory,
BlockerRejection::UnsupportedCategory("flaky_test".to_string()),
BlockerRejection::EmptyEvidence,
BlockerRejection::MissingNextAction,
BlockerRejection::MissingResumable,
];
for rejection in rejections {
let mut driver = AcceptanceProtocolDriver::default();
match decide_acceptance_blocker(
&mut driver,
&AcceptanceResult::BareBlocker {
rejection: rejection.clone(),
},
) {
Some(AcceptanceBlockerDecision::ProtocolRetry { progress, .. }) => {
assert!(
progress.contains(&rejection.reason()),
"diagnostic must name the defect ({rejection:?}): {progress}"
);
}
other => panic!("{rejection:?} must retry, got {other:?}"),
}
}
}
#[test]
fn validated_blocker_stalls_with_its_explicit_category_preserved() {
let blocker = crate::acceptance::AcceptanceBlocker {
category: "human_decision".to_string(),
evidence: vec!["the credential token auth story needs an owner decision".to_string()],
unblock_condition: "the architecture owner records a decision".to_string(),
next_action: "owner decides the approach".to_string(),
resumable: false,
prerequisite_owner: Some("architecture".to_string()),
evidence_ids: vec!["adr-7".to_string()],
};
let mut driver = AcceptanceProtocolDriver::default();
match decide_acceptance_blocker(
&mut driver,
&AcceptanceResult::Stalled {
blocker: blocker.clone(),
},
) {
Some(AcceptanceBlockerDecision::ExternalBlocker { blocker: stalled }) => {
assert_eq!(stalled, blocker);
assert_eq!(
stalled.category, "human_decision",
"credential/token/auth prose must not override the explicit category"
);
}
other => panic!("a validated blocker must stall, got {other:?}"),
}
}
#[test]
fn blocker_decision_api_ignores_non_blocker_results() {
let mut driver = AcceptanceProtocolDriver::default();
for result in [
AcceptanceResult::Pass,
AcceptanceResult::Fail {
findings: Vec::new(),
},
AcceptanceResult::Continue,
AcceptanceResult::MissingVerdict {
findings: Vec::new(),
},
AcceptanceResult::CommandFailed {
error: "exit 1".to_string(),
findings: Vec::new(),
diagnostic: AcceptanceCommandDiagnostic::default(),
},
AcceptanceResult::Cancelled,
] {
assert!(
decide_acceptance_blocker(&mut driver, &result).is_none(),
"{result:?} must keep its own routing"
);
}
}
#[test]
fn bare_blocker_decisions_have_caller_parity() {
let sequences: [Vec<AcceptanceResult>; 3] = [
vec![bare_blocker(), bare_blocker(), AcceptanceResult::Pass],
vec![bare_blocker(), bare_blocker(), bare_blocker()],
vec![
bare_blocker(),
AcceptanceResult::Stalled {
blocker: sample_blocker(),
},
bare_blocker(),
],
];
for sequence in sequences {
let first = drive_blockers(&sequence);
let parallel = drive_blockers(&sequence);
assert_eq!(
first, parallel,
"equivalent observations must produce equivalent decisions for {sequence:?}"
);
}
}
#[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().into()],
},
true,
false,
),
(AcceptanceResult::Continue, true, false),
(
AcceptanceResult::Stalled {
blocker: sample_blocker(),
},
true,
false,
),
(
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_external(
"pending_verification",
"denied",
),
},
true,
false,
),
(missing_verdict("waiting"), false, true),
(
AcceptanceResult::CommandFailed {
error: "exit code 1".to_string(),
findings: vec!["boom".to_string()],
diagnostic: AcceptanceCommandDiagnostic::default(),
},
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 = AcceptanceProtocolDriver::default();
for result in [
AcceptanceResult::CommandFailed {
error: "exit code 1".to_string(),
findings: Vec::new(),
diagnostic: AcceptanceCommandDiagnostic::default(),
},
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: AcceptanceProtocolRetry {
kind: AcceptanceProtocolError::MissingVerdict,
attempt: 1,
max: 2
},
..
}
));
}
#[test]
fn missing_verdict_driver_has_caller_routing_parity() {
let sequence = [
missing_verdict("waiting"),
missing_verdict("waiting"),
AcceptanceResult::Fail {
findings: vec!["src/lib.rs:1 missing coverage".to_string().into()],
},
];
let first = replay_missing_verdict_sequence(&sequence);
let parallel = replay_missing_verdict_sequence(&sequence);
assert_eq!(first, parallel);
assert_eq!(first.0.len(), 3);
assert!(matches!(first.1, Some(AcceptanceResult::Fail { .. })));
assert!(first.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().into(),
"src/lib.rs:11 missing test".to_string().into(),
]);
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().into()]);
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 semantic_fingerprint_tracks_json_tasks_but_excludes_runtime_follow_up() {
fn document(status: &str, follow_up: &str) -> String {
format!(
r#"{{"schema_version":1,"tasks":[{{"id":"a","title":"A","status":"{status}","section":"implementation"}}]{follow_up}}}"#
)
}
let temp = tempfile::TempDir::new().unwrap();
let tasks = temp.path().join("openspec/changes/example/tasks.json");
std::fs::create_dir_all(tasks.parent().unwrap()).unwrap();
std::fs::write(&tasks, document("pending", "")).unwrap();
let before = semantic_progress_fingerprint(temp.path()).unwrap();
std::fs::write(
&tasks,
document(
"pending",
r#","acceptance_follow_up":{"attempt":2,"findings":[{"identity":"repository|src/a.rs|verification","text":"Add the missing test","remediation_claimed":false}]}"#,
),
)
.unwrap();
assert_eq!(before, semantic_progress_fingerprint(temp.path()).unwrap());
std::fs::write(&tasks, document("completed", "")).unwrap();
assert_ne!(before, semantic_progress_fingerprint(temp.path()).unwrap());
}
#[test]
fn semantic_fingerprint_hashes_malformed_json_tasks_verbatim() {
let temp = tempfile::TempDir::new().unwrap();
let tasks = temp.path().join("openspec/changes/example/tasks.json");
std::fs::create_dir_all(tasks.parent().unwrap()).unwrap();
std::fs::write(&tasks, "{not json").unwrap();
let before = semantic_progress_fingerprint(temp.path()).unwrap();
std::fs::write(
&tasks,
r#"{"schema_version":1,"tasks":[{"id":"a","title":"A","status":"pending","section":"implementation"}]}"#,
)
.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().into(),
"src/client.rs: rate limit retry missing".to_string().into(),
"network unreachable: fix retry handling".to_string().into(),
"dns resolution failed while repairing src/client.rs"
.to_string()
.into(),
"missing non-mockable external credential"
.to_string()
.into(),
]);
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().into(),
"src/client.rs: rate limit retry missing".to_string().into(),
"network unreachable: fix retry handling".to_string().into(),
"dns resolution failed while repairing src/client.rs"
.to_string()
.into(),
"missing non-mockable external credential"
.to_string()
.into(),
])
.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 identical_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 first = decide_acceptance_retry(&previous, Some("same"), &findings, "same", 2);
let parallel = decide_acceptance_retry(&previous, Some("same"), &findings, "same", 2);
assert_eq!(first, parallel);
assert!(matches!(
first,
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().into()]
}
.is_pass());
assert!(!AcceptanceResult::CommandFailed {
error: "test".to_string(),
findings: vec!["failure".to_string()],
diagnostic: AcceptanceCommandDiagnostic::default(),
}
.is_pass());
assert!(!AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::acceptance_external(
"pending_verification",
"permission denied"
),
}
.is_pass());
assert!(!AcceptanceResult::MissingVerdict {
findings: vec!["status-only output".to_string()],
}
.is_pass());
assert!(!AcceptanceResult::Cancelled.is_pass());
assert!(!AcceptanceResult::Stalled {
blocker: sample_blocker()
}
.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"),
}
}
use crate::acceptance::AcceptanceFinding;
fn finding(id: &str, implementation: &str, verification: &str) -> AcceptanceFinding {
AcceptanceFinding::structured(crate::acceptance::RepositoryFinding {
id: id.to_string(),
severity: crate::acceptance::FindingSeverity::Minor,
summary: "Challenge and proof leakage is not tested by value".to_string(),
evidence: vec!["relay exposes counts but not issued values".to_string()],
required_changes: vec![crate::acceptance::FindingFileExpectation {
file: implementation.to_string(),
description: "Expose issued challenge and presented proof values".to_string(),
}],
verification: vec![crate::acceptance::FindingFileExpectation {
file: verification.to_string(),
description: "Assert recorded values are absent from audit output".to_string(),
}],
})
}
fn secret_value_finding() -> AcceptanceFinding {
finding(
"acceptance-secret-value-scan",
"tests/support/relay.ts",
"runtime/recovery.integration.test.ts",
)
}
fn changed(paths: &[&str]) -> Vec<String> {
paths.iter().map(|path| path.to_string()).collect()
}
#[test]
fn complete_coverage_permits_semantic_review_without_claiming_resolution() {
let findings = [secret_value_finding()];
let coverage = evaluate_remediation_coverage(
&findings,
&changed(&[
"tests/support/relay.ts",
"runtime/recovery.integration.test.ts",
]),
);
assert!(coverage.is_complete());
assert!(coverage.uncovered().is_empty());
assert!(coverage.unrelated_files.is_empty());
assert_eq!(
decide_repair_gate(
"change-a",
&findings,
&FindingRepairLedger::default(),
Some("fail-rev"),
Some("apply-rev"),
&changed(&[
"tests/support/relay.ts",
"runtime/recovery.integration.test.ts",
]),
&[],
),
RepairGateDecision::Proceed
);
}
#[test]
fn missing_implementation_or_verification_path_fails_coverage() {
let findings = [secret_value_finding()];
let missing_implementation = evaluate_remediation_coverage(
&findings,
&changed(&["runtime/recovery.integration.test.ts"]),
);
assert!(!missing_implementation.is_complete());
assert_eq!(
missing_implementation.missing_required,
[(
"acceptance-secret-value-scan".to_string(),
"tests/support/relay.ts".to_string()
)]
);
assert!(missing_implementation.missing_verification.is_empty());
let missing_verification =
evaluate_remediation_coverage(&findings, &changed(&["tests/support/relay.ts"]));
assert!(!missing_verification.is_complete());
assert_eq!(
missing_verification.missing_verification,
[(
"acceptance-secret-value-scan".to_string(),
"runtime/recovery.integration.test.ts".to_string()
)]
);
}
#[test]
fn calibration_only_change_stops_before_acceptance_with_evidence() {
let findings = [secret_value_finding()];
let delta = changed(&["tests/calibration.test.ts", "src/unrelated.rs"]);
let RepairGateDecision::Stop(stop) = decide_repair_gate(
"change-a",
&findings,
&FindingRepairLedger::default(),
Some("fail-rev"),
Some("apply-rev"),
&delta,
&["adjusted calibration threshold".to_string()],
) else {
panic!("calibration-only repair must not authorize another acceptance run");
};
assert_eq!(stop.reason, REMEDIATION_MISMATCH_REASON);
assert!(stop.resumable);
assert_eq!(stop.fail_revision.as_deref(), Some("fail-rev"));
assert_eq!(stop.apply_revision.as_deref(), Some("apply-rev"));
assert_eq!(stop.required_files, ["tests/support/relay.ts"]);
assert_eq!(
stop.verification_files,
["runtime/recovery.integration.test.ts"]
);
assert_eq!(
stop.coverage.unrelated_files,
["src/unrelated.rs", "tests/calibration.test.ts"],
"every changed file is retained as a diagnostic, sorted"
);
assert_eq!(stop.coverage.changed_files.len(), delta.len());
assert_eq!(
stop.remediation_evidence,
["adjusted calibration threshold"]
);
assert_eq!(stop.findings, findings);
let json = stop.to_json();
assert_eq!(json["coverage_complete"], false);
assert_eq!(json["proves_acceptance_pass"], false);
assert_eq!(json["proves_archive_readiness"], false);
assert_eq!(json["proves_completion"], false);
}
#[test]
fn coverage_normalizes_renames_and_untracked_paths() {
let findings = [secret_value_finding()];
let coverage = evaluate_remediation_coverage(
&findings,
&changed(&[
"tests/support/old-relay.ts -> tests/support/relay.ts",
"./runtime/recovery.integration.test.ts",
]),
);
assert!(coverage.is_complete(), "{coverage:?}");
assert_eq!(
coverage.unrelated_files,
["tests/support/old-relay.ts".to_string()]
);
}
#[test]
fn legacy_findings_without_declared_paths_keep_compatibility_behavior() {
let findings = crate::acceptance::legacy_findings(["src/a.rs:10 missing coverage"]);
let coverage = evaluate_remediation_coverage(&findings, &changed(&["docs/readme.md"]));
assert!(
coverage.is_complete(),
"a legacy finding declares no path set, so strict coverage cannot apply"
);
assert_eq!(
coverage.legacy_finding_texts,
["src/a.rs:10 missing coverage"]
);
assert_eq!(
decide_repair_gate(
"change-a",
&findings,
&FindingRepairLedger::default(),
None,
None,
&changed(&["docs/readme.md"]),
&[],
),
RepairGateDecision::Proceed
);
}
#[test]
fn empty_repair_delta_cannot_satisfy_a_declared_finding() {
let findings = [secret_value_finding()];
assert!(matches!(
decide_repair_gate(
"change-a",
&findings,
&FindingRepairLedger::default(),
Some("fail-rev"),
Some("fail-rev"),
&[],
&[],
),
RepairGateDecision::Stop(_)
));
}
fn observe(
ledger: &mut FindingRepairLedger,
findings: &[AcceptanceFinding],
) -> FindingRepairDecision {
ledger.observe_fail(&normalize_findings(findings))
}
#[test]
fn first_observation_grants_one_repair_opportunity() {
let mut ledger = FindingRepairLedger::default();
let findings = [secret_value_finding()];
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &findings) else {
panic!("first observation must allow one repair");
};
assert_eq!(identities, ["repository|id|acceptance-secret-value-scan"]);
assert_eq!(
ledger.occurrences("repository|id|acceptance-secret-value-scan"),
1
);
assert!(!ledger.has_consumed_repair("repository|id|acceptance-secret-value-scan"));
ledger.record_repair_dispatched(&identities);
assert!(ledger.has_consumed_repair("repository|id|acceptance-secret-value-scan"));
}
#[test]
fn repeated_id_stops_before_a_second_repair_regardless_of_progress() {
let mut ledger = FindingRepairLedger::default();
let findings = [secret_value_finding()];
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &findings) else {
panic!("first observation must allow one repair");
};
ledger.record_repair_dispatched(&identities);
let restated = AcceptanceFinding::structured(crate::acceptance::RepositoryFinding {
id: "acceptance-secret-value-scan".to_string(),
severity: crate::acceptance::FindingSeverity::Major,
summary: "Rewritten summary at a new line".to_string(),
evidence: vec!["different evidence at src/other.rs:99".to_string()],
required_changes: vec![crate::acceptance::FindingFileExpectation {
file: "src/other.rs".to_string(),
description: "different described change".to_string(),
}],
verification: vec![crate::acceptance::FindingFileExpectation {
file: "tests/other.rs".to_string(),
description: "different described proof".to_string(),
}],
});
let FindingRepairDecision::Stop {
reason,
repeated_identities,
} = observe(&mut ledger, std::slice::from_ref(&restated))
else {
panic!("a repeated ID must stop automatic repair");
};
assert_eq!(reason, REPEATED_FINDING_REASON);
assert_eq!(
repeated_identities,
["repository|id|acceptance-secret-value-scan"]
);
assert_eq!(
ledger.occurrences("repository|id|acceptance-secret-value-scan"),
2
);
let stop = repeated_finding_stop(
"change-a",
&[restated],
&ledger,
repeated_identities,
Some("fail-rev"),
Some("apply-rev"),
&changed(&["src/other.rs", "tests/other.rs", "docs/notes.md"]),
&["claimed a repair".to_string()],
);
assert_eq!(stop.reason, REPEATED_FINDING_REASON);
assert!(
stop.coverage.is_complete(),
"coverage passing is irrelevant"
);
assert_eq!(
stop.occurrences,
[("repository|id|acceptance-secret-value-scan".to_string(), 2)]
);
assert!(stop.resumable);
assert_eq!(stop.remediation_evidence, ["claimed a repair"]);
}
#[test]
fn a_new_finding_id_receives_its_own_repair_opportunity() {
let mut ledger = FindingRepairLedger::default();
let first = [secret_value_finding()];
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &first) else {
panic!("first observation must allow one repair");
};
ledger.record_repair_dispatched(&identities);
let second = [finding(
"acceptance-new-defect",
"src/new.rs",
"tests/new.rs",
)];
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &second) else {
panic!("a new ID gets its own opportunity");
};
assert_eq!(identities, ["repository|id|acceptance-new-defect"]);
assert_eq!(
ledger.occurrences("repository|id|acceptance-secret-value-scan"),
0,
"an ID acceptance stopped reporting is closed, not carried forward"
);
}
#[test]
fn mixed_repeated_and_new_ids_stop_atomically_with_every_finding_retained() {
let mut ledger = FindingRepairLedger::default();
let repeated = secret_value_finding();
let FindingRepairDecision::Repair { identities } =
observe(&mut ledger, std::slice::from_ref(&repeated))
else {
panic!("first observation must allow one repair");
};
ledger.record_repair_dispatched(&identities);
let fresh = finding("acceptance-new-defect", "src/new.rs", "tests/new.rs");
let mixed = [repeated, fresh];
let FindingRepairDecision::Stop {
reason,
repeated_identities,
} = observe(&mut ledger, &mixed)
else {
panic!("a mixed payload containing a repeated ID must stop atomically");
};
assert_eq!(reason, REPEATED_FINDING_REASON);
assert_eq!(
repeated_identities,
["repository|id|acceptance-secret-value-scan"],
"only the repeated ID is the stop reason"
);
let stop = repeated_finding_stop(
"change-a",
&mixed,
&ledger,
repeated_identities,
None,
None,
&[],
&[],
);
assert_eq!(stop.findings.len(), 2, "diagnostics retain every finding");
assert_eq!(
stop.occurrence_identities(),
[
"repository|id|acceptance-new-defect".to_string(),
"repository|id|acceptance-secret-value-scan".to_string(),
]
);
}
#[test]
fn explicit_retry_reset_releases_the_automatic_budget_only() {
let mut ledger = FindingRepairLedger::default();
let findings = [secret_value_finding()];
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &findings) else {
panic!("first observation must allow one repair");
};
ledger.record_repair_dispatched(&identities);
ledger.reset_for_explicit_retry();
assert!(!ledger.has_consumed_repair("repository|id|acceptance-secret-value-scan"));
assert_eq!(
ledger.occurrences("repository|id|acceptance-secret-value-scan"),
1,
"prior occurrence evidence is preserved for review"
);
assert!(matches!(
observe(&mut ledger, &findings),
FindingRepairDecision::Repair { .. }
));
}
#[test]
fn legacy_findings_use_the_fallback_identity_for_the_repair_budget() {
let mut ledger = FindingRepairLedger::default();
let findings = crate::acceptance::legacy_findings(["Missing retry test at src/run.rs:10"]);
let FindingRepairDecision::Repair { identities } = observe(&mut ledger, &findings) else {
panic!("first observation must allow one repair");
};
assert_eq!(identities, ["repository|src/run.rs|verification"]);
ledger.record_repair_dispatched(&identities);
let restated =
crate::acceptance::legacy_findings(["Regression coverage absent in src/run.rs:99"]);
assert!(matches!(
observe(&mut ledger, &restated),
FindingRepairDecision::Stop { .. }
));
}
#[test]
fn structured_finding_identity_is_the_reviewer_id_not_its_prose() {
let first = normalize_findings(&[secret_value_finding()]);
let restated = normalize_findings(&[AcceptanceFinding::structured(
crate::acceptance::RepositoryFinding {
id: "acceptance-secret-value-scan".to_string(),
severity: crate::acceptance::FindingSeverity::Major,
summary: "totally different words".to_string(),
evidence: vec!["totally different evidence".to_string()],
required_changes: vec![crate::acceptance::FindingFileExpectation {
file: "src/elsewhere.rs".to_string(),
description: "d".to_string(),
}],
verification: vec![crate::acceptance::FindingFileExpectation {
file: "tests/elsewhere.rs".to_string(),
description: "d".to_string(),
}],
},
)]);
assert_eq!(first[0].identity, restated[0].identity);
assert_ne!(
first[0].text, restated[0].text,
"identity is stable while the payload stays free to change"
);
}
mod acceptance_escalation {
use super::*;
use crate::acceptance::{
AcceptanceFinding, BlockerRejection, FindingFileExpectation, FindingRejection,
FindingSeverity, RepositoryFinding,
};
use crate::config::{AcceptanceEscalationConfig, OrchestratorConfig};
fn config_with(command: Option<&str>, policy: Option<(u32, u32)>) -> OrchestratorConfig {
OrchestratorConfig {
acceptance_command: Some("accept {change_id} {prompt}".to_string()),
acceptance_escalation_command: command.map(str::to_string),
acceptance_escalation: policy.map(|(after, max)| AcceptanceEscalationConfig {
after_invalid_results: Some(after),
max_uses_per_sequence: Some(max),
}),
..Default::default()
}
}
fn driver(command: Option<&str>, policy: Option<(u32, u32)>) -> AcceptanceEscalationDriver {
AcceptanceEscalationDriver::from_config(&config_with(command, policy))
}
fn default_driver() -> AcceptanceEscalationDriver {
driver(Some("deep-accept {change_id} {prompt}"), None)
}
fn empty_fail() -> AcceptanceResult {
AcceptanceResult::Fail {
findings: crate::acceptance::legacy_findings([GENERIC_ACCEPTANCE_FAIL_FINDING]),
}
}
fn actionable_fail() -> AcceptanceResult {
AcceptanceResult::Fail {
findings: vec![AcceptanceFinding::structured(RepositoryFinding {
id: "finding-1".to_string(),
severity: FindingSeverity::Major,
summary: "escalation policy is unbounded".to_string(),
evidence: vec!["src/orchestration/acceptance.rs:1".to_string()],
required_changes: vec![FindingFileExpectation {
file: "src/orchestration/acceptance.rs".to_string(),
description: "bound the sequence".to_string(),
}],
verification: vec![FindingFileExpectation {
file: "src/orchestration/acceptance.rs".to_string(),
description: "prove the bound".to_string(),
}],
})],
}
}
fn missing_verdict() -> AcceptanceResult {
AcceptanceResult::MissingVerdict {
findings: vec!["waiting for verification".to_string()],
}
}
fn bare_blocker() -> AcceptanceResult {
AcceptanceResult::BareBlocker {
rejection: BlockerRejection::Missing,
}
}
fn malformed_finding() -> AcceptanceResult {
AcceptanceResult::MalformedFinding {
rejection: FindingRejection::MissingId,
}
}
#[test]
fn every_eligible_invalid_class_selects_the_alternate_reviewer() {
for (result, expected) in [
(missing_verdict(), InvalidAcceptanceResult::MissingVerdict),
(bare_blocker(), InvalidAcceptanceResult::BareBlocker),
(
malformed_finding(),
InvalidAcceptanceResult::MalformedFinding,
),
(empty_fail(), InvalidAcceptanceResult::EmptyFail),
] {
assert_eq!(
classify_invalid_acceptance_result(&result),
Some(expected),
"{expected:?} must be an eligible invalid result"
);
let mut escalation = default_driver();
let outcome = escalation.observe(&result, None);
assert!(
outcome.escalation_selected(),
"{expected:?} must select escalation under the default policy: {outcome:?}"
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
}
}
#[test]
fn actionable_fail_never_escalates_and_resets_the_sequence() {
assert_eq!(classify_invalid_acceptance_result(&actionable_fail()), None);
let mut escalation = default_driver();
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
let outcome = escalation.observe(&actionable_fail(), None);
assert_eq!(outcome, AcceptanceEscalationOutcome::SequenceReset);
assert_eq!(escalation.consecutive_invalid(), 0);
assert_eq!(escalation.uses_in_sequence(), 0);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal,
"a real defect report must not leave an alternate reviewer pending"
);
}
#[test]
fn only_the_exact_generic_fallback_counts_as_an_empty_fail() {
assert!(is_generic_empty_fail(&[]));
assert!(is_generic_empty_fail(&crate::acceptance::legacy_findings(
[GENERIC_ACCEPTANCE_FAIL_FINDING]
)));
assert!(!is_generic_empty_fail(&crate::acceptance::legacy_findings(
["Investigate acceptance failure and apply the required fix now"]
)));
assert!(!is_generic_empty_fail(&crate::acceptance::legacy_findings(
[
GENERIC_ACCEPTANCE_FAIL_FINDING,
GENERIC_ACCEPTANCE_FAIL_FINDING,
]
)));
}
#[test]
fn excluded_results_never_escalate() {
let excluded = [
AcceptanceResult::Pass,
AcceptanceResult::Continue,
actionable_fail(),
AcceptanceResult::Stalled {
blocker: crate::acceptance::AcceptanceBlocker {
category: "credential".to_string(),
evidence: vec!["missing token".to_string()],
unblock_condition: "token provisioned".to_string(),
next_action: "ask the owner".to_string(),
resumable: true,
prerequisite_owner: None,
evidence_ids: Vec::new(),
},
},
AcceptanceResult::PermissionStalled {
blocker: crate::events::StalledBlocker::permission_denial(
"acceptance",
&crate::permission::PermissionDenial {
category: crate::permission::PermissionDenialCategory::ToolAccess,
denied_target: "git".to_string(),
evidence: "denied".to_string(),
},
),
},
AcceptanceResult::CommandFailed {
error: "exit 1".to_string(),
findings: vec!["boom".to_string()],
diagnostic: AcceptanceCommandDiagnostic::default(),
},
AcceptanceResult::RuntimeLimit {
limit: AcceptanceRuntimeLimit {
limit_secs: 1800,
cleanup_confirmed: true,
cleanup_diagnostics: String::new(),
},
},
AcceptanceResult::Cancelled,
];
for result in excluded {
assert_eq!(
classify_invalid_acceptance_result(&result),
None,
"{result:?} must not be an eligible invalid result"
);
let mut escalation = default_driver();
let outcome = escalation.observe(&result, None);
assert!(
!outcome.escalation_selected(),
"{result:?} must not select escalation: {outcome:?}"
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
}
}
#[test]
fn non_completing_invocations_preserve_the_pending_selection() {
for result in [
AcceptanceResult::CommandFailed {
error: "exit 1".to_string(),
findings: Vec::new(),
diagnostic: AcceptanceCommandDiagnostic::default(),
},
AcceptanceResult::Cancelled,
AcceptanceResult::RuntimeLimit {
limit: AcceptanceRuntimeLimit {
limit_secs: 1800,
cleanup_confirmed: true,
cleanup_diagnostics: String::new(),
},
},
] {
let mut escalation = default_driver();
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
assert_eq!(
escalation.observe(&result, None),
AcceptanceEscalationOutcome::NotObserved,
"{result:?} must not be observed against the invalid-result sequence"
);
assert_eq!(escalation.consecutive_invalid(), 1);
assert_eq!(escalation.uses_in_sequence(), 1);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation,
"the pending alternate reviewer must survive an invocation that never ran"
);
}
}
#[test]
fn the_command_mode_is_consumed_by_exactly_one_invocation() {
let mut escalation = default_driver();
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal,
"one selection must not survive into a second invocation"
);
}
#[test]
fn a_missing_escalation_command_preserves_normal_retry_behavior() {
let mut escalation = driver(None, None);
let outcome = escalation.observe(&missing_verdict(), None);
assert_eq!(
outcome,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::MissingVerdict,
consecutive_invalid: 1,
reason: EscalationDeclineReason::NoCommandConfigured,
}
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
}
#[test]
fn escalation_waits_for_the_configured_threshold() {
let mut escalation = driver(Some("deep-accept {prompt}"), Some((3, 1)));
for expected in 1..3 {
let outcome = escalation.observe(&missing_verdict(), None);
assert_eq!(
outcome,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::MissingVerdict,
consecutive_invalid: expected,
reason: EscalationDeclineReason::BelowThreshold {
after_invalid_results: 3,
},
}
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
}
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
}
#[test]
fn the_use_cap_bounds_one_invalid_result_sequence() {
let mut escalation = default_driver();
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
let outcome = escalation.observe(&missing_verdict(), None);
assert_eq!(
outcome,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::MissingVerdict,
consecutive_invalid: 2,
reason: EscalationDeclineReason::UseCapExhausted {
max_uses_per_sequence: 1,
},
}
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
}
#[test]
fn a_larger_cap_permits_exactly_that_many_uses() {
let mut escalation = driver(Some("deep-accept {prompt}"), Some((1, 2)));
for expected_use in 1..=2 {
let outcome = escalation.observe(&bare_blocker(), None);
assert!(outcome.escalation_selected(), "{outcome:?}");
assert_eq!(escalation.uses_in_sequence(), expected_use);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
}
assert!(!escalation
.observe(&bare_blocker(), None)
.escalation_selected());
}
#[test]
fn a_non_invalid_result_resets_the_sequence() {
for reset_result in [
AcceptanceResult::Pass,
AcceptanceResult::Continue,
actionable_fail(),
] {
let mut escalation = default_driver();
assert!(escalation
.observe(&missing_verdict(), None)
.escalation_selected());
let _ = escalation.take_command_mode();
assert!(!escalation
.observe(&malformed_finding(), None)
.escalation_selected());
assert_eq!(
escalation.observe(&reset_result, None),
AcceptanceEscalationOutcome::SequenceReset
);
assert_eq!(escalation.consecutive_invalid(), 0);
assert_eq!(escalation.uses_in_sequence(), 0);
assert!(
escalation
.observe(&bare_blocker(), None)
.escalation_selected(),
"a new sequence must get a fresh escalation budget after {reset_result:?}"
);
}
}
#[test]
fn empty_fail_below_threshold_keeps_its_sequence_across_the_apply_round() {
let mut escalation = driver(Some("deep-accept {prompt}"), Some((2, 1)));
let first = escalation.observe(&empty_fail(), Some("rev-a"));
assert_eq!(
first,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::EmptyFail,
consecutive_invalid: 1,
reason: EscalationDeclineReason::BelowThreshold {
after_invalid_results: 2,
},
},
"below threshold the empty FAIL follows the existing generic fallback"
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
let second = escalation.observe(&empty_fail(), Some("rev-a"));
assert!(
second.escalation_selected(),
"the retained count must let a later invalid result reach the threshold: {second:?}"
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Escalation
);
}
#[test]
fn a_revision_change_resets_the_sequence() {
let mut escalation = driver(Some("deep-accept {prompt}"), Some((2, 1)));
assert!(!escalation
.observe(&empty_fail(), Some("rev-a"))
.escalation_selected());
let after_apply = escalation.observe(&empty_fail(), Some("rev-b"));
assert_eq!(
after_apply,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::EmptyFail,
consecutive_invalid: 1,
reason: EscalationDeclineReason::BelowThreshold {
after_invalid_results: 2,
},
},
"a moved revision must start a fresh sequence"
);
assert_eq!(
escalation.take_command_mode(),
AcceptanceCommandMode::Normal
);
}
#[test]
fn only_a_selected_empty_fail_escalation_replaces_the_generic_repair() {
let escalated = |kind| AcceptanceEscalationOutcome::Escalated {
kind,
consecutive_invalid: 1,
use_index: 1,
max_uses_per_sequence: 1,
};
assert!(
escalates_empty_fail(&escalated(InvalidAcceptanceResult::EmptyFail)),
"an escalated empty FAIL is routed as an acceptance-only alternate review"
);
for kind in [
InvalidAcceptanceResult::MissingVerdict,
InvalidAcceptanceResult::BareBlocker,
InvalidAcceptanceResult::MalformedFinding,
] {
assert!(
!escalates_empty_fail(&escalated(kind)),
"{kind:?} already owns an acceptance-only retry; its routing must not change"
);
}
for outcome in [
AcceptanceEscalationOutcome::NotObserved,
AcceptanceEscalationOutcome::SequenceReset,
AcceptanceEscalationOutcome::Retained {
kind: InvalidAcceptanceResult::EmptyFail,
consecutive_invalid: 1,
reason: EscalationDeclineReason::BelowThreshold {
after_invalid_results: 2,
},
},
] {
assert!(
!escalates_empty_fail(&outcome),
"{outcome:?} selected no alternate reviewer, so routing is unchanged"
);
}
}
#[test]
fn the_routing_decision_reads_the_drivers_own_observation() {
let mut escalation = driver(Some("deep-accept {prompt}"), Some((2, 1)));
let below_threshold = escalation.observe(&empty_fail(), Some("rev-a"));
assert!(
!escalates_empty_fail(&below_threshold),
"below threshold the empty FAIL keeps the generic FAIL-to-Apply repair"
);
let at_threshold = escalation.observe(&empty_fail(), Some("rev-a"));
assert!(
escalates_empty_fail(&at_threshold),
"at the threshold the empty FAIL becomes an acceptance-only alternate review"
);
let mut protocol_class = driver(Some("deep-accept {prompt}"), Some((1, 1)));
let escalated_missing_verdict =
protocol_class.observe(&missing_verdict(), Some("rev-a"));
assert!(escalated_missing_verdict.escalation_selected());
assert!(
!escalates_empty_fail(&escalated_missing_verdict),
"a missing verdict escalates the reviewer, not the routing"
);
}
#[test]
fn a_fresh_driver_carries_no_escalation_state() {
let config = config_with(Some("deep-accept {prompt}"), None);
let mut first = AcceptanceEscalationDriver::from_config(&config);
assert!(first
.observe(&missing_verdict(), Some("rev-a"))
.escalation_selected());
assert_eq!(first.consecutive_invalid(), 1);
assert_eq!(first.uses_in_sequence(), 1);
let mut restarted = AcceptanceEscalationDriver::from_config(&config);
assert_eq!(restarted.consecutive_invalid(), 0);
assert_eq!(restarted.uses_in_sequence(), 0);
assert_eq!(
restarted.take_command_mode(),
AcceptanceCommandMode::Normal,
"a restarted run must begin with the ordinary acceptance command"
);
}
#[test]
fn only_the_two_decisions_produce_a_diagnostic() {
assert!(AcceptanceEscalationOutcome::NotObserved
.diagnostic()
.is_none());
assert!(AcceptanceEscalationOutcome::SequenceReset
.diagnostic()
.is_none());
let mut escalation = default_driver();
let escalated = escalation.observe(&malformed_finding(), None).diagnostic();
let escalated = escalated.expect("an escalation decision must be reportable");
assert!(escalated.contains("malformed-finding"), "{escalated}");
assert!(
escalated.contains("acceptance_escalation_command"),
"{escalated}"
);
let retained = escalation.observe(&malformed_finding(), None).diagnostic();
let retained = retained.expect("a declined escalation must be reportable");
assert!(retained.contains("acceptance_command"), "{retained}");
assert!(retained.contains("budget"), "{retained}");
}
#[test]
fn the_command_mode_names_its_configuration_key() {
assert_eq!(AcceptanceCommandMode::Normal.label(), "acceptance_command");
assert_eq!(
AcceptanceCommandMode::Escalation.label(),
"acceptance_escalation_command"
);
assert!(!AcceptanceCommandMode::Normal.is_escalation());
assert!(AcceptanceCommandMode::Escalation.is_escalation());
}
#[test]
fn the_command_mode_resolves_exactly_one_reviewer_template() {
let config = config_with(Some("deep-accept {change_id} {prompt}"), None);
assert_eq!(
acceptance_command_template(&config, AcceptanceCommandMode::Normal)
.expect("the normal reviewer is a required command"),
"accept {change_id} {prompt}"
);
assert_eq!(
acceptance_command_template(&config, AcceptanceCommandMode::Escalation)
.expect("the escalation reviewer is configured"),
"deep-accept {change_id} {prompt}"
);
}
#[test]
fn escalation_without_a_configured_command_is_a_configuration_error() {
let config = config_with(None, None);
assert_eq!(
acceptance_command_template(&config, AcceptanceCommandMode::Normal)
.expect("the normal reviewer is unaffected"),
"accept {change_id} {prompt}"
);
let message = acceptance_command_template(&config, AcceptanceCommandMode::Escalation)
.expect_err("an unconfigured alternate reviewer must not silently fall back")
.to_string();
assert!(
message.contains("acceptance_escalation_command"),
"{message}"
);
}
#[test]
fn only_an_escalation_invocation_carries_the_alternate_reviewer_framing() {
assert!(
crate::agent::build_acceptance_escalation_context(AcceptanceCommandMode::Normal)
.is_empty(),
"an ordinary acceptance invocation must not be told it is an alternate review"
);
let context = crate::agent::build_acceptance_escalation_context(
AcceptanceCommandMode::Escalation,
);
assert!(context.contains("<acceptance_escalation>"), "{context}");
assert!(
context.contains("alternate acceptance reviewer"),
"{context}"
);
assert!(
context.contains("one fresh canonical acceptance verdict"),
"{context}"
);
}
#[test]
fn only_empty_fail_lacks_a_protocol_retry_budget() {
assert!(InvalidAcceptanceResult::MissingVerdict.has_protocol_retry_budget());
assert!(InvalidAcceptanceResult::BareBlocker.has_protocol_retry_budget());
assert!(InvalidAcceptanceResult::MalformedFinding.has_protocol_retry_budget());
assert!(!InvalidAcceptanceResult::EmptyFail.has_protocol_retry_budget());
}
}
}