use super::support::create_test_config;
use crate::analyzer::{AnalysisOutcome, AnalysisProvenance, AnalysisResult};
use crate::config::OrchestratorConfig;
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use crate::parallel::analysis_signature::{AnalysisInputProbe, DEGRADED_SUPPRESSION_TTL};
use crate::parallel::cleanup::WorkspaceCleanupGuard;
use crate::parallel::dynamic_queue::ReanalysisReason;
use crate::parallel::queue_state::ReanalysisDispatchContext;
use crate::parallel::{ParallelExecutor, WorkspaceResult};
use crate::vcs::VcsBackend;
use std::collections::{HashMap, HashSet};
use std::future::Future;
use std::pin::Pin;
use std::process::Command;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex as StdMutex};
use std::time::Duration;
use tempfile::TempDir;
use tokio::sync::mpsc;
use tokio::task::JoinSet;
const SCHEDULER_TIMER: Duration = Duration::from_millis(500);
const QUEUE_DEBOUNCE: Duration = Duration::from_secs(10);
fn test_change(id: &str) -> Change {
Change {
id: id.to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: String::new(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}
}
fn init_minimal_git_repo(repo_root: &std::path::Path) {
for args in [
vec!["init", "-b", "main"],
vec!["config", "user.email", "test@example.com"],
vec!["config", "user.name", "Test User"],
] {
let output = Command::new("git")
.args(args)
.current_dir(repo_root)
.output()
.expect("run git setup command");
assert!(output.status.success(), "git setup command failed");
}
std::fs::write(repo_root.join("README.md"), "base\n").expect("write base file");
for args in [vec!["add", "-A"], vec!["commit", "-m", "Base"]] {
let output = Command::new("git")
.args(args)
.current_dir(repo_root)
.output()
.expect("run git commit command");
assert!(output.status.success(), "git commit command failed");
}
}
struct FakeAnalysisInputProbe {
revision: StdMutex<String>,
digests: StdMutex<HashMap<String, String>>,
revision_failure: StdMutex<Option<String>>,
digest_failure: StdMutex<Option<String>>,
revision_probes: AtomicUsize,
digest_probes: AtomicUsize,
}
impl FakeAnalysisInputProbe {
fn new(revision: &str) -> Arc<Self> {
Arc::new(Self {
revision: StdMutex::new(revision.to_string()),
digests: StdMutex::new(HashMap::new()),
revision_failure: StdMutex::new(None),
digest_failure: StdMutex::new(None),
revision_probes: AtomicUsize::new(0),
digest_probes: AtomicUsize::new(0),
})
}
fn set_revision(&self, revision: &str) {
*self.revision.lock().expect("revision lock") = revision.to_string();
}
fn set_effective_base(&self, base_ref: &str, revision: &str) {
self.set_revision(&format!("{base_ref}@{revision}"));
}
fn clear_revision_failure(&self) {
*self.revision_failure.lock().expect("revision failure lock") = None;
}
fn set_proposal_digest(&self, change_id: &str, digest: &str) {
self.digests
.lock()
.expect("digest lock")
.insert(change_id.to_string(), digest.to_string());
}
fn fail_revision(&self, error: &str) {
*self.revision_failure.lock().expect("revision failure lock") = Some(error.to_string());
}
fn fail_proposal_digest(&self, error: &str) {
*self.digest_failure.lock().expect("digest failure lock") = Some(error.to_string());
}
fn revision_probes(&self) -> usize {
self.revision_probes.load(Ordering::SeqCst)
}
fn digest_probes(&self) -> usize {
self.digest_probes.load(Ordering::SeqCst)
}
}
#[async_trait::async_trait]
impl AnalysisInputProbe for FakeAnalysisInputProbe {
async fn base_revision(&self) -> Result<String, String> {
self.revision_probes.fetch_add(1, Ordering::SeqCst);
if let Some(error) = self
.revision_failure
.lock()
.expect("revision failure lock")
.clone()
{
return Err(error);
}
Ok(self.revision.lock().expect("revision lock").clone())
}
fn proposal_digest(&self, change_id: &str) -> Result<String, String> {
self.digest_probes.fetch_add(1, Ordering::SeqCst);
if let Some(error) = self
.digest_failure
.lock()
.expect("digest failure lock")
.clone()
{
return Err(error);
}
Ok(self
.digests
.lock()
.expect("digest lock")
.get(change_id)
.cloned()
.unwrap_or_else(|| format!("digest-of-{change_id}")))
}
}
type AnalysisFuture<'a> = Pin<Box<dyn Future<Output = AnalysisOutcome> + Send + 'a>>;
const DISPATCH_HOLDER: &str = "dispatch-holder";
const UNRESOLVABLE_DEPENDENCY: &str = "never-resolvable-dependency";
struct AnalyzerScript {
invocations: AtomicUsize,
duration: StdMutex<Duration>,
provenance: StdMutex<AnalysisProvenance>,
dependencies: StdMutex<HashMap<String, Vec<String>>>,
blanket_dependency: StdMutex<Option<String>>,
empty_order: StdMutex<bool>,
during_analysis: StdMutex<Option<Box<dyn Fn() + Send + Sync>>>,
observed_in_flight: StdMutex<Vec<Vec<String>>>,
}
impl AnalyzerScript {
fn new() -> Arc<Self> {
Arc::new(Self {
invocations: AtomicUsize::new(0),
duration: StdMutex::new(Duration::from_secs(1)),
provenance: StdMutex::new(AnalysisProvenance::HealthyLlm),
dependencies: StdMutex::new(HashMap::new()),
blanket_dependency: StdMutex::new(None),
empty_order: StdMutex::new(false),
during_analysis: StdMutex::new(None),
observed_in_flight: StdMutex::new(Vec::new()),
})
}
fn observed_in_flight(&self) -> Vec<Vec<String>> {
self.observed_in_flight
.lock()
.expect("observed in-flight lock")
.clone()
}
fn invocations(&self) -> usize {
self.invocations.load(Ordering::SeqCst)
}
fn set_duration(&self, duration: Duration) {
*self.duration.lock().expect("duration lock") = duration;
}
fn set_provenance(&self, provenance: AnalysisProvenance) {
*self.provenance.lock().expect("provenance lock") = provenance;
}
fn set_dependencies(&self, dependencies: HashMap<String, Vec<String>>) {
*self.dependencies.lock().expect("dependencies lock") = dependencies;
}
fn set_empty_order(&self, empty_order: bool) {
*self.empty_order.lock().expect("empty order lock") = empty_order;
}
fn set_blanket_dependency(&self, dependency: Option<&str>) {
*self
.blanket_dependency
.lock()
.expect("blanket dependency lock") = dependency.map(str::to_string);
}
fn run_during_analysis(&self, effect: impl Fn() + Send + Sync + 'static) {
*self.during_analysis.lock().expect("during analysis lock") = Some(Box::new(effect));
}
}
fn scripted_analyzer(
script: Arc<AnalyzerScript>,
) -> impl for<'a> Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync {
move |changes: &[Change], in_flight: &[String], _iteration: u32| -> AnalysisFuture<'_> {
script.invocations.fetch_add(1, Ordering::SeqCst);
script
.observed_in_flight
.lock()
.expect("observed in-flight lock")
.push(in_flight.to_vec());
let order: Vec<String> = if *script.empty_order.lock().expect("empty order lock") {
Vec::new()
} else {
changes.iter().map(|change| change.id.clone()).collect()
};
let mut dependencies = script
.dependencies
.lock()
.expect("dependencies lock")
.clone();
if let Some(blanket) = script
.blanket_dependency
.lock()
.expect("blanket dependency lock")
.clone()
{
for change_id in &order {
dependencies
.entry(change_id.clone())
.or_default()
.push(blanket.clone());
}
}
let provenance = *script.provenance.lock().expect("provenance lock");
let duration = *script.duration.lock().expect("duration lock");
if let Some(effect) = script
.during_analysis
.lock()
.expect("during analysis lock")
.as_ref()
{
effect();
}
Box::pin(async move {
tokio::time::sleep(duration).await;
AnalysisOutcome::new(
AnalysisResult {
order,
dependencies,
groups: None,
},
provenance,
)
})
}
}
struct SuppressionHarness {
executor: ParallelExecutor,
queued: Vec<Change>,
in_flight: HashSet<String>,
join_set: JoinSet<WorkspaceResult>,
cleanup_guard: WorkspaceCleanupGuard,
reanalysis_reason: ReanalysisReason,
iteration: u32,
max_parallelism: usize,
probe: Arc<FakeAnalysisInputProbe>,
script: Arc<AnalyzerScript>,
events: mpsc::Receiver<ExecutionEvent>,
}
impl SuppressionHarness {
fn new(repo_root: std::path::PathBuf, queued: Vec<Change>) -> Self {
Self::with_config(repo_root, queued, create_test_config(), 1)
}
fn with_config(
repo_root: std::path::PathBuf,
queued: Vec<Change>,
config: OrchestratorConfig,
max_parallelism: usize,
) -> Self {
let (tx, events) = mpsc::channel(256);
let mut executor = ParallelExecutor::new(repo_root.clone(), config, Some(tx));
let probe = FakeAnalysisInputProbe::new("base-rev-1");
executor.set_analysis_input_probe(probe.clone());
Self {
executor,
queued,
in_flight: HashSet::new(),
join_set: JoinSet::new(),
cleanup_guard: WorkspaceCleanupGuard::new(VcsBackend::Git, repo_root),
reanalysis_reason: ReanalysisReason::Initial,
iteration: 2,
max_parallelism,
probe,
script: AnalyzerScript::new(),
events,
}
}
async fn make_queue_debounce_stale(&self) {
let mut last_change = self.executor.last_queue_change_at.lock().await;
*last_change = Some(std::time::Instant::now() - Duration::from_secs(600));
}
fn keep_dispatch_inert(&mut self) {
self.script
.set_blanket_dependency(Some(UNRESOLVABLE_DEPENDENCY));
self.in_flight.insert(DISPATCH_HOLDER.to_string());
self.grant_one_dispatch_slot();
}
fn allow_dispatch(&mut self) {
self.script.set_blanket_dependency(None);
}
fn grant_one_dispatch_slot(&mut self) {
self.max_parallelism = self.in_flight.len() + 1;
}
fn add_in_flight(&mut self, change_id: &str) {
self.in_flight.insert(change_id.to_string());
self.grant_one_dispatch_slot();
}
fn set_in_flight(&mut self, change_ids: &[&str]) {
self.in_flight = change_ids.iter().map(|id| id.to_string()).collect();
self.grant_one_dispatch_slot();
}
async fn run_loop_iteration<F>(&mut self, analyzer: &F) -> Option<(bool, u32)>
where
for<'a> F: Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync,
{
let outcome = self
.executor
.evaluate_queued_reanalysis_and_dispatch(
ReanalysisDispatchContext {
queued: &mut self.queued,
in_flight: &mut self.in_flight,
max_parallelism: self.max_parallelism,
iteration: self.iteration,
reanalysis_reason: self.reanalysis_reason,
analyzer,
join_set: &mut self.join_set,
cleanup_guard: &mut self.cleanup_guard,
work_snapshot: None,
},
&mut self.reanalysis_reason,
)
.await
.expect("scheduler re-analysis evaluation must not fail");
if let Some((_, new_iteration)) = outcome {
self.iteration = new_iteration;
}
outcome
}
async fn timer_wake(&self) {
tokio::time::sleep(SCHEDULER_TIMER).await;
}
async fn timer_wakes<F>(&mut self, analyzer: &F, wakes: usize)
where
for<'a> F: Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync,
{
for _ in 0..wakes {
self.timer_wake().await;
self.run_loop_iteration(analyzer).await;
}
}
async fn pass_probe_deadline(&self) {
tokio::time::sleep(QUEUE_DEBOUNCE + SCHEDULER_TIMER).await;
}
fn deliver_edge(&mut self, reason: ReanalysisReason) {
self.reanalysis_reason = reason;
}
fn analyses(&self) -> usize {
self.script.invocations()
}
fn drain_logs(&mut self) -> Vec<String> {
let mut messages = Vec::new();
while let Ok(event) = self.events.try_recv() {
if let ExecutionEvent::Log(entry) = event {
messages.push(entry.message);
}
}
messages
}
fn drain_attempt_counts(&mut self) -> (usize, usize) {
let (mut analysis_started, mut apply_started) = (0, 0);
while let Ok(event) = self.events.try_recv() {
match event {
ExecutionEvent::AnalysisStarted { .. } => analysis_started += 1,
ExecutionEvent::ApplyStarted { .. } => apply_started += 1,
_ => {}
}
}
(analysis_started, apply_started)
}
}
const TIMER_WAKES: usize = 24;
async fn assert_unchanged_timer_wakes_analyze_once(analysis_duration: Duration) {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness.script.set_duration(analysis_duration);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"an ordinary timer evaluation of a never-analyzed input must analyze once"
);
harness.timer_wakes(&analyzer, TIMER_WAKES).await;
assert_eq!(
harness.analyses(),
1,
"unchanged timer wakes must not invoke the analyzer again; saw {} analyses",
harness.analyses()
);
assert_eq!(
harness.queued.len(),
1,
"queued work must be retained while nothing dispatches it"
);
let (analysis_started, apply_started) = harness.drain_attempt_counts();
assert_eq!(
apply_started, 0,
"a suppressed pass must not dispatch apply work"
);
assert_eq!(
analysis_started, 1,
"a pass suppressed before analyzer invocation must not surface as a distinct analysis \
attempt in operator-visible output"
);
}
#[tokio::test(start_paused = true)]
async fn unchanged_timer_wakes_analyze_once_when_analysis_is_shorter_than_debounce() {
assert_unchanged_timer_wakes_analyze_once(Duration::from_secs(2)).await;
}
#[tokio::test(start_paused = true)]
async fn unchanged_timer_wakes_analyze_once_when_analysis_is_longer_than_debounce() {
assert_unchanged_timer_wakes_analyze_once(Duration::from_secs(25)).await;
}
#[tokio::test(start_paused = true)]
async fn stable_multi_change_queue_is_analyzed_once() {
let temp_dir = TempDir::new().unwrap();
let queued: Vec<Change> = (0..4)
.map(|index| test_change(&format!("queued-{index:02}")))
.collect();
let mut harness = SuppressionHarness::new(temp_dir.path().to_path_buf(), queued);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
1,
"a stable multi-change queue must be analyzed once, not once per wake"
);
assert_eq!(
harness.queued.len(),
4,
"every queued change must be retained while capacity is zero"
);
}
#[tokio::test(start_paused = true)]
async fn suppressed_wakes_do_not_probe_before_the_bounded_deadline() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
let probes_after_analysis = harness.probe.revision_probes();
let digest_probes_after_analysis = harness.probe.digest_probes();
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.probe.revision_probes(),
probes_after_analysis,
"the wake immediately after a completed analysis must not re-probe the VCS revision"
);
assert_eq!(
harness.probe.digest_probes(),
digest_probes_after_analysis,
"the wake immediately after a completed analysis must not re-read proposal files"
);
harness.timer_wakes(&analyzer, 18).await;
assert_eq!(
harness.probe.revision_probes(),
probes_after_analysis,
"a suppressed 500 ms wake must not spawn a VCS revision probe"
);
assert_eq!(
harness.probe.digest_probes(),
digest_probes_after_analysis,
"a suppressed 500 ms wake must not re-read proposal files"
);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.probe.revision_probes(),
probes_after_analysis + 1,
"the first evaluation after the probe deadline must re-confirm repository-visible input"
);
assert_eq!(
harness.analyses(),
1,
"re-probing an unchanged input must not invoke the analyzer"
);
let logs = harness.drain_logs();
assert!(
logs.iter()
.any(|message| message.contains("unchanged_analysis_input")),
"suppression must be observable as a deduplicated operator-visible reason; saw {logs:?}"
);
}
#[tokio::test(start_paused = true)]
async fn queue_addition_bypasses_a_matching_signature_once() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(
harness.analyses(),
1,
"the unchanged input must be suppressed before the queue event"
);
harness.deliver_edge(ReanalysisReason::QueueNotification);
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a queue notification must not be suppressed by a prior matching signature"
);
harness.deliver_edge(ReanalysisReason::Initial);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
2,
"timer wakes after the consumed event must stay quiescent"
);
}
#[tokio::test(start_paused = true)]
async fn each_explicit_edge_bypasses_a_matching_signature_once() {
for edge in [
ReanalysisReason::ResolveCompletion,
ReanalysisReason::RepairCandidate,
ReanalysisReason::SlotRecovery,
] {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(
harness.analyses(),
1,
"{edge}: the unchanged input must be suppressed before the edge"
);
harness.deliver_edge(edge);
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"{edge}: an explicit edge must analyze once even with a matching signature"
);
assert_eq!(
harness.reanalysis_reason,
ReanalysisReason::Initial,
"{edge}: the one-shot edge must still be consumed after evaluation"
);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
2,
"{edge}: timer wakes must not replay the consumed edge"
);
}
}
#[tokio::test(start_paused = true)]
async fn changed_capacity_rearms_analysis_without_a_slot_recovery_reason() {
let repo_dir = TempDir::new().unwrap();
let workspace_base = TempDir::new().unwrap();
init_minimal_git_repo(repo_dir.path());
let config = OrchestratorConfig {
workspace_base_dir: Some(workspace_base.path().to_string_lossy().to_string()),
..create_test_config()
};
let mut harness = SuppressionHarness::with_config(
repo_dir.path().to_path_buf(),
vec![test_change("queued-a")],
config,
1,
);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1, "the stable input is analyzed once");
assert!(
harness.queued.iter().any(|change| change.id == "queued-a"),
"the undispatched change stays queued"
);
harness.max_parallelism += 1;
harness.pass_probe_deadline().await;
let (should_break, _iteration) = harness
.run_loop_iteration(&analyzer)
.await
.expect("queued work must be evaluated");
assert!(!should_break, "changed capacity must resume the scheduler");
assert_eq!(
harness.reanalysis_reason,
ReanalysisReason::Initial,
"this variant must not rely on a retained SlotRecovery reason"
);
assert_eq!(
harness.analyses(),
2,
"changed capacity must re-arm dependency analysis"
);
harness.allow_dispatch();
harness.max_parallelism += 1;
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert!(
harness.in_flight.contains("queued-a"),
"eligible queued work must reach dispatch once nothing blocks it"
);
assert!(harness.queued.is_empty());
harness.join_set.abort_all();
while harness.join_set.join_next().await.is_some() {}
}
#[tokio::test(start_paused = true)]
async fn same_id_queued_proposal_change_rearms_analysis() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1);
harness.probe.set_proposal_digest("queued-a", "edited");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"an edited proposal with the same change ID must re-arm analysis"
);
}
#[tokio::test(start_paused = true)]
async fn same_id_in_flight_proposal_change_rearms_analysis() {
let temp_dir = TempDir::new().unwrap();
let mut harness = SuppressionHarness::with_config(
temp_dir.path().to_path_buf(),
vec![test_change("queued-a")],
create_test_config(),
2,
);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.add_in_flight("inflight-a");
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1);
harness.probe.set_proposal_digest("inflight-a", "edited");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"an edited in-flight proposal must re-arm analysis"
);
}
#[tokio::test(start_paused = true)]
async fn effective_base_revision_change_rearms_analysis() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1);
harness.probe.set_revision("base-rev-2");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a changed effective dependency-base revision must re-arm analysis"
);
}
#[tokio::test(start_paused = true)]
async fn input_change_during_analysis_is_visible_at_the_next_probe() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let probe = harness.probe.clone();
harness.script.run_during_analysis(move || {
probe.set_proposal_digest("queued-a", "edited-during-analysis");
});
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"recording the pre-analysis snapshot must keep a change made during analysis visible"
);
}
#[tokio::test(start_paused = true)]
async fn fresh_executor_starts_without_a_prior_analysis_signature() {
let temp_dir = TempDir::new().unwrap();
let queued = vec![test_change("queued-a")];
let mut first = SuppressionHarness::new(temp_dir.path().to_path_buf(), queued.clone());
let first_analyzer = scripted_analyzer(first.script.clone());
first.make_queue_debounce_stale().await;
first.keep_dispatch_inert();
first.run_loop_iteration(&first_analyzer).await;
first.timer_wakes(&first_analyzer, 4).await;
assert_eq!(first.analyses(), 1);
let mut restarted = SuppressionHarness::new(temp_dir.path().to_path_buf(), queued);
let restarted_analyzer = scripted_analyzer(restarted.script.clone());
restarted.make_queue_debounce_stale().await;
restarted.keep_dispatch_inert();
restarted.run_loop_iteration(&restarted_analyzer).await;
assert_eq!(
restarted.analyses(),
1,
"a fresh executor must perform its own initial analysis"
);
}
#[tokio::test(start_paused = true)]
async fn degraded_fallback_suppresses_rapid_retries_and_permits_one_retry_after_five_minutes() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness
.script
.set_provenance(AnalysisProvenance::RecoverableFailureFallback);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"a recoverable-failure fallback is still a usable degraded result"
);
for _ in 0..8 {
tokio::time::sleep(Duration::from_secs(30)).await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
1,
"a broken analyzer command must not be relaunched on every timer wake; saw {} analyses",
harness.analyses()
);
tokio::time::sleep(DEGRADED_SUPPRESSION_TTL).await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"the degraded record must expire and permit one retry"
);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
2,
"the retry must re-arm a new bounded window rather than a rapid loop"
);
}
#[tokio::test(start_paused = true)]
async fn healthy_result_after_a_degraded_record_becomes_non_expiring() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness
.script
.set_provenance(AnalysisProvenance::RecoverableFailureFallback);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
harness
.script
.set_provenance(AnalysisProvenance::HealthyLlm);
tokio::time::sleep(DEGRADED_SUPPRESSION_TTL).await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 2, "one retry must be permitted");
for _ in 0..3 {
tokio::time::sleep(DEGRADED_SUPPRESSION_TTL * 2).await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
2,
"a healthy result must record a non-expiring signature; saw {} analyses",
harness.analyses()
);
}
#[tokio::test(start_paused = true)]
async fn intentional_metadata_only_result_is_not_treated_as_degraded() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness
.script
.set_provenance(AnalysisProvenance::IntentionalMetadataOnly);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
for _ in 0..3 {
tokio::time::sleep(DEGRADED_SUPPRESSION_TTL * 2).await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
1,
"configured metadata-only analysis is the intended result and must not expire"
);
}
#[tokio::test(start_paused = true)]
async fn unusable_empty_result_establishes_no_completed_signature() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness.script.set_empty_order(true);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"an unusable empty result must not falsely establish a completed input"
);
}
#[tokio::test(start_paused = true)]
async fn zero_dispatch_with_positive_capacity_and_idle_scheduler_stays_eligible() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness.script.set_dependencies(HashMap::from([(
"queued-a".to_string(),
vec!["hallucinated-dependency".to_string()],
)]));
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
assert!(
harness.in_flight.is_empty(),
"the erroneous dependency must block dispatch"
);
assert_eq!(
harness.queued.len(),
1,
"the undispatched change stays queued"
);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a zero-dispatch idle pass must stay eligible for the next debounced analysis"
);
}
#[tokio::test(start_paused = true)]
async fn revision_probe_failure_fails_open_without_recording_suppression() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.probe.fail_revision("revision resolution failed");
for _ in 0..3 {
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
3,
"signature construction failure must fail open instead of suppressing analysis"
);
let logs = harness.drain_logs();
assert!(
logs.iter()
.any(|message| message.contains("signature unavailable")),
"the fail-open path must stay operator-visible; saw {logs:?}"
);
harness.probe.set_revision("base-rev-recovered");
harness.probe.clear_revision_failure();
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 4);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
4,
"once a signature can be built again, unchanged wakes must be suppressed"
);
}
#[tokio::test(start_paused = true)]
async fn effective_base_ref_change_rearms_analysis() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.probe.set_effective_base("main", "commit-1");
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1);
harness
.probe
.set_effective_base("integration-1", "commit-1");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a changed effective dependency-base ref must re-arm analysis"
);
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(
harness.analyses(),
2,
"the new base must then be suppressed"
);
harness
.probe
.set_effective_base("integration-1", "commit-2");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
3,
"an advancing effective dependency-base ref must re-arm analysis"
);
}
#[tokio::test(start_paused = true)]
async fn persistent_signature_failure_is_rate_limited_across_timer_wakes() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.probe.fail_revision("revision resolution failed");
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"an unavailable signature must permit analysis"
);
let probes_after_first_attempt = harness.probe.revision_probes();
harness.timer_wakes(&analyzer, 16).await;
assert_eq!(
harness.analyses(),
1,
"a persistent signature failure must not relaunch the analyzer on every 500 ms wake; saw \
{} analyses",
harness.analyses()
);
assert_eq!(
harness.probe.revision_probes(),
probes_after_first_attempt,
"a throttled wake must not re-probe repository-visible signature material"
);
let logs = harness.drain_logs();
assert!(
logs.iter()
.any(|message| message.contains("analysis_signature_unavailable_retry_pending")),
"the bounded fail-open retry must stay operator-visible; saw {logs:?}"
);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"the deadline must permit one retry of both signature construction and analysis"
);
assert_eq!(
harness.probe.revision_probes(),
probes_after_first_attempt + 1,
"the retry must probe exactly once"
);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
2,
"the retry must re-arm the bounded window rather than a rapid loop"
);
assert!(
!harness.queued.is_empty(),
"a fail-open pass must keep queued work alive"
);
}
#[tokio::test(start_paused = true)]
async fn explicit_edge_bypasses_the_signature_failure_retry_deadline_once() {
for edge in [
ReanalysisReason::QueueNotification,
ReanalysisReason::ResolveCompletion,
ReanalysisReason::RepairCandidate,
ReanalysisReason::SlotRecovery,
] {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.probe.fail_revision("revision resolution failed");
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 3).await;
assert_eq!(
harness.analyses(),
1,
"{edge}: the bounded retry deadline must hold before the edge"
);
harness.deliver_edge(edge);
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"{edge}: an explicit edge must bypass the bounded retry deadline once"
);
harness.deliver_edge(ReanalysisReason::Initial);
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(
harness.analyses(),
2,
"{edge}: timer wakes after the event must return to bounded retry"
);
}
}
#[tokio::test(start_paused = true)]
async fn signature_probe_recovery_reestablishes_suppression_without_a_queue_change() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.probe.fail_revision("revision resolution failed");
harness.run_loop_iteration(&analyzer).await;
harness.timer_wakes(&analyzer, 4).await;
assert_eq!(harness.analyses(), 1, "the failing probe fails open once");
harness.probe.clear_revision_failure();
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"the first eligible evaluation after recovery must analyze once"
);
assert_eq!(
harness.reanalysis_reason,
ReanalysisReason::Initial,
"recovery must not rely on an explicit edge"
);
harness.timer_wakes(&analyzer, 6).await;
assert_eq!(
harness.analyses(),
2,
"a recovered signature must re-establish ordinary suppression"
);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"re-probing the recovered unchanged input must not invoke the analyzer"
);
}
#[tokio::test(start_paused = true)]
async fn unusable_empty_result_retries_at_the_bounded_cadence() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness.script.set_empty_order(true);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
harness.timer_wakes(&analyzer, 18).await;
assert_eq!(
harness.analyses(),
1,
"a persistently unusable analyzer must not be relaunched on every wake; saw {} analyses",
harness.analyses()
);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"the bounded deadline must permit one retry of the unusable input"
);
}
#[tokio::test(start_paused = true)]
async fn degraded_expiry_is_not_delayed_by_the_repository_probe_cadence() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
harness
.script
.set_provenance(AnalysisProvenance::RecoverableFailureFallback);
harness.script.set_duration(Duration::ZERO);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
tokio::time::sleep(DEGRADED_SUPPRESSION_TTL - Duration::from_secs(1)).await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"the degraded record must still suppress one second before expiry"
);
tokio::time::sleep(Duration::from_secs(1)).await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a probe taken just before expiry must not delay the degraded retry"
);
}
#[tokio::test(start_paused = true)]
async fn in_flight_order_permutations_produce_one_deterministic_analyzer_input() {
let temp_dir = TempDir::new().unwrap();
let mut harness = SuppressionHarness::with_config(
temp_dir.path().to_path_buf(),
vec![test_change("queued-a")],
create_test_config(),
4,
);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
for id in ["inflight-c", "inflight-a", "inflight-b"] {
harness.add_in_flight(id);
}
harness.run_loop_iteration(&analyzer).await;
assert_eq!(harness.analyses(), 1);
harness.set_in_flight(&[DISPATCH_HOLDER, "inflight-b", "inflight-c", "inflight-a"]);
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"an order-only permutation of the same in-flight set is one semantic input"
);
harness.add_in_flight("inflight-d");
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
2,
"a real in-flight membership change must re-arm analysis"
);
for observed in harness.script.observed_in_flight() {
let mut sorted = observed.clone();
sorted.sort();
assert_eq!(
observed, sorted,
"the analyzer must receive one deterministic ordering of the in-flight set"
);
}
}
#[tokio::test(start_paused = true)]
async fn proposal_read_failure_fails_open_without_recording_suppression() {
let temp_dir = TempDir::new().unwrap();
let mut harness =
SuppressionHarness::new(temp_dir.path().to_path_buf(), vec![test_change("queued-a")]);
let analyzer = scripted_analyzer(harness.script.clone());
harness.make_queue_debounce_stale().await;
harness.keep_dispatch_inert();
harness
.probe
.fail_proposal_digest("proposal.md could not be read");
for _ in 0..3 {
harness.pass_probe_deadline().await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
3,
"a proposal read failure must fail open instead of suppressing analysis"
);
assert_eq!(
harness.queued.len(),
1,
"a fail-open pass must not lose queued work"
);
}