use crate::analyzer::{AnalysisOutcome, AnalysisResult};
use crate::config::OrchestratorConfig;
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use crate::orchestration::operator_command::{
ExecutionMarkStore, NoopQueueHooks, OperatorCommandService, OperatorMode,
};
use crate::orchestration::run_control::testing::{RecordingScheduler, SchedulerCall};
use crate::orchestration::run_control::{
ResolveReservations, RunControlOutcome, RunControlService, SchedulerEffect, StartEligibility,
};
use crate::orchestration::state::{OrchestratorState, ReducerCommand};
use crate::parallel::cleanup::WorkspaceCleanupGuard;
use crate::parallel::dynamic_queue::ReanalysisReason;
use crate::parallel::queue_state::ReanalysisDispatchContext;
use crate::parallel::{ParallelExecutor, WorkspaceResult};
use crate::tui::queue::DynamicQueue;
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;
use tempfile::TempDir;
use tokio::sync::{mpsc, RwLock};
use tokio::task::JoinSet;
const ALPHA: &str = "alpha";
const BETA: &str = "beta";
const SCHEDULER_TIMER: std::time::Duration = std::time::Duration::from_millis(500);
type AnalysisFuture<'a> = Pin<Box<dyn Future<Output = AnalysisOutcome> + Send + 'a>>;
fn test_config(workspace_base: &std::path::Path) -> OrchestratorConfig {
OrchestratorConfig {
apply_command: Some("echo apply {change_id}".to_string()),
archive_command: Some("echo archive {change_id}".to_string()),
analyze_command: Some("echo analyze".to_string()),
acceptance_command: Some("echo acceptance".to_string()),
resolve_command: Some("echo resolve".to_string()),
workspace_base_dir: Some(workspace_base.to_string_lossy().to_string()),
..Default::default()
}
}
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 git(repo_root: &std::path::Path, args: &[&str]) {
let output = Command::new("git")
.args(args)
.current_dir(repo_root)
.output()
.expect("run git command");
assert!(output.status.success(), "git {args:?} failed");
}
fn repo_with_changes(repo_root: &std::path::Path, change_ids: &[&str]) {
git(repo_root, &["init", "-b", "main"]);
git(repo_root, &["config", "user.email", "test@example.com"]);
git(repo_root, &["config", "user.name", "Test User"]);
std::fs::write(repo_root.join("README.md"), "base\n").expect("write base file");
for change_id in change_ids {
let dir = repo_root.join("openspec/changes").join(change_id);
std::fs::create_dir_all(&dir).expect("create change dir");
std::fs::write(dir.join("proposal.md"), format!("# {change_id}\n"))
.expect("write proposal");
std::fs::write(dir.join("tasks.md"), "- [ ] work\n").expect("write tasks");
}
git(repo_root, &["add", "-A"]);
git(repo_root, &["commit", "-m", "Base"]);
}
fn counting_analyzer(
invocations: Arc<AtomicUsize>,
) -> impl for<'a> Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync {
move |changes: &[Change], _in_flight: &[String], _iteration: u32| -> AnalysisFuture<'_> {
invocations.fetch_add(1, Ordering::SeqCst);
let order: Vec<String> = changes.iter().map(|change| change.id.clone()).collect();
Box::pin(async move {
AnalysisResult {
order,
dependencies: HashMap::new(),
groups: None,
}
.into()
})
}
}
#[derive(Debug, Default)]
struct ObservedEvents {
analysis_started: usize,
dispatch_started: Vec<String>,
}
struct Harness {
executor: ParallelExecutor,
queue: Arc<DynamicQueue>,
state: Arc<RwLock<OrchestratorState>>,
run_control: RunControlService,
marks: Arc<ExecutionMarkStore>,
scheduler: Arc<RecordingScheduler>,
queued: Vec<Change>,
in_flight: HashSet<String>,
join_set: JoinSet<WorkspaceResult>,
cleanup_guard: WorkspaceCleanupGuard,
reanalysis_reason: ReanalysisReason,
iteration: u32,
max_parallelism: usize,
analyses: Arc<AtomicUsize>,
events: mpsc::Receiver<ExecutionEvent>,
repo_root: std::path::PathBuf,
_repo_dir: TempDir,
_workspace_base: TempDir,
}
impl Harness {
fn new(catalog: &[&str], max_parallelism: usize) -> Self {
let repo_dir = TempDir::new().expect("create repo dir");
let workspace_base = TempDir::new().expect("create workspace base");
repo_with_changes(repo_dir.path(), catalog);
let (tx, events) = mpsc::channel(256);
let mut executor = ParallelExecutor::new(
repo_dir.path().to_path_buf(),
test_config(workspace_base.path()),
Some(tx),
);
let queue = Arc::new(DynamicQueue::new());
executor.set_dynamic_queue(queue.clone());
let state = Arc::new(RwLock::new(OrchestratorState::new(
catalog.iter().map(|id| (*id).to_string()).collect(),
10,
)));
executor.set_shared_orchestrator_state(state.clone());
let marks = Arc::new(ExecutionMarkStore::new());
let scheduler = Arc::new(RecordingScheduler::new());
scheduler.set_running(true);
let operator = Arc::new(OperatorCommandService::new(
state.clone(),
queue.clone(),
Arc::new(NoopQueueHooks),
marks.clone(),
));
let run_control = RunControlService::new(
state.clone(),
operator,
scheduler.clone(),
Arc::new(ResolveReservations::new()),
Arc::new(StartEligibility::new()),
);
Self {
executor,
queue,
state,
run_control,
marks,
scheduler,
queued: Vec::new(),
in_flight: HashSet::new(),
join_set: JoinSet::new(),
cleanup_guard: WorkspaceCleanupGuard::new(
VcsBackend::Git,
repo_dir.path().to_path_buf(),
),
reanalysis_reason: ReanalysisReason::Initial,
iteration: 2,
max_parallelism,
analyses: Arc::new(AtomicUsize::new(0)),
events,
repo_root: repo_dir.path().to_path_buf(),
_repo_dir: repo_dir,
_workspace_base: workspace_base,
}
}
async fn arm_queue_debounce(&self) {
let mut last_change = self.executor.last_queue_change_at.lock().await;
*last_change = Some(std::time::Instant::now());
}
async fn fail_with_runtime_limit(&mut self, change_id: &str) {
{
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
guard.apply_execution_event(&ExecutionEvent::ProcessingError {
id: change_id.to_string(),
error: "Apply exceeded its absolute runtime limit of 3600s".to_string(),
});
}
self.executor.failed_tracker.mark_failed(change_id);
}
fn depends_on(&mut self, edges: &[(&str, &str)]) {
let mut dependencies: HashMap<String, Vec<String>> = HashMap::new();
for (dependent, dependency) in edges {
dependencies
.entry((*dependent).to_string())
.or_default()
.push((*dependency).to_string());
}
self.executor.failed_tracker.set_dependencies(dependencies);
}
fn is_gated_by_failure(&self, dependent: &str) -> bool {
self.executor
.failed_tracker
.should_skip(dependent)
.is_some()
}
async fn consume_retry_edges(&mut self) -> bool {
self.executor
.consume_explicit_retry_edges()
.await
.bypass_armed
}
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 retry_edges = self.executor.consume_explicit_retry_edges().await;
let in_flight = self.in_flight.clone();
let dynamic_queue_added = self
.executor
.check_dynamic_queue_and_add_changes(
&mut self.queued,
&in_flight,
&mut self.reanalysis_reason,
)
.await;
let reconciliation = self
.executor
.reconcile_queued_candidates_from_shared_state(&mut self.queued, &in_flight)
.await;
self.reanalysis_reason = ParallelExecutor::derive_pass_reanalysis_reason(
self.reanalysis_reason,
retry_edges,
reconciliation,
dynamic_queue_added,
);
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 should not fail");
if let Some((_, new_iteration)) = outcome {
self.iteration = new_iteration;
}
outcome
}
async fn timer_wake(&self) {
tokio::time::sleep(SCHEDULER_TIMER).await;
}
fn analyses(&self) -> usize {
self.analyses.load(Ordering::SeqCst)
}
fn drain_events(&mut self) -> (usize, usize) {
let observed = self.drain_observed_events();
(observed.analysis_started, observed.dispatch_started.len())
}
fn catalog_candidate(&self, change_id: &str) -> Change {
crate::openspec::list_changes_native_from(&self.repo_root)
.expect("load the OpenSpec catalog")
.into_iter()
.find(|change| change.id == change_id)
.expect("the catalog contains the change")
}
async fn settle_dispatch_back_to_queued(&mut self, change_id: &str) {
self.join_set.abort_all();
while self.join_set.join_next().await.is_some() {}
self.in_flight.remove(change_id);
self.executor.lifecycle_slots.release(change_id);
if !self.queued.iter().any(|change| change.id == change_id) {
let candidate = self.catalog_candidate(change_id);
self.queued.push(candidate);
}
}
async fn stall_in_acceptance(&self, change_id: &str) {
self.state
.write()
.await
.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: change_id.to_string(),
blocker: crate::events::StalledBlocker {
category: "acceptance_finding".to_string(),
phase: "acceptance".to_string(),
gate: "acceptance".to_string(),
error_summary: "unresolved acceptance finding".to_string(),
evidence: vec!["tests/acceptance.rs:1".to_string()],
unblock_condition: None,
prerequisite_owner: None,
next_action: "resolve finding and retry".to_string(),
resumable: true,
worktree_preserved: true,
},
});
}
fn drain_observed_events(&mut self) -> ObservedEvents {
let mut observed = ObservedEvents::default();
while let Ok(event) = self.events.try_recv() {
match event {
ExecutionEvent::AnalysisStarted { .. } => observed.analysis_started += 1,
ExecutionEvent::WorkspacePreparationStarted { change_id } => {
observed.dispatch_started.push(change_id)
}
_ => {}
}
}
observed
}
async fn status(&self, change_id: &str) -> String {
self.state
.read()
.await
.display_status(change_id)
.to_string()
}
async fn shutdown(mut self) {
self.join_set.abort_all();
while self.join_set.join_next().await.is_some() {}
}
}
#[tokio::test(start_paused = true)]
async fn change_error_f5_retry_runtime_limit_failure_is_retried_only_by_explicit_intent() {
let mut harness = Harness::new(&[ALPHA, BETA], 2);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.arm_queue_debounce().await;
harness.fail_with_runtime_limit(ALPHA).await;
harness.fail_with_runtime_limit(BETA).await;
assert_eq!(harness.status(ALPHA).await, "error");
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
let (_, apply_started) = harness.drain_events();
assert_eq!(
apply_started, 0,
"a runtime-limit failure must not be redispatched without explicit intent"
);
assert!(
harness.in_flight.is_empty(),
"no terminal-error change may be admitted for execution"
);
harness.depends_on(&[("alpha-dep", ALPHA), ("beta-dep", BETA)]);
assert!(
harness.is_gated_by_failure("alpha-dep"),
"the failed classification stands until an accepted retry releases it"
);
let analyses_before_retry = harness.analyses();
harness.marks.set(ALPHA, true);
let outcome = harness
.run_control
.start(OperatorMode::Running)
.await
.expect("a marked retry-eligible error is startable while the run is live");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec![ALPHA.to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
},
"the live scheduler is woken rather than joined by a second boundary"
);
assert_eq!(
harness.scheduler.calls(),
vec![SchedulerCall::Notified],
"exactly one wake"
);
assert_eq!(
harness.status(ALPHA).await,
"queued",
"the terminal error was cleared through retry intent"
);
assert_eq!(
harness.status(BETA).await,
"error",
"an unretried failure keeps its terminal evidence"
);
assert!(
harness.consume_retry_edges().await,
"the accepted retry published an edge for this scheduler to consume"
);
harness.depends_on(&[("alpha-dep", ALPHA), ("beta-dep", BETA)]);
assert!(
!harness.is_gated_by_failure("alpha-dep"),
"the retried change's failed classification is released"
);
assert!(
harness.is_gated_by_failure("beta-dep"),
"the edge is target-specific: another failed change stays gated"
);
harness.run_loop_iteration(&analyzer).await;
let (analysis_started, _) = harness.drain_events();
assert_eq!(
harness.analyses(),
analyses_before_retry + 1,
"the consumed edge arms exactly one reevaluation"
);
assert_eq!(
analysis_started, 1,
"a distinct AnalysisStarted follows the accepted retry without \
waiting for mark settlement or debounce expiry"
);
assert!(
harness.in_flight.contains(ALPHA),
"the released change is dispatched once dependency and capacity guards allow it"
);
harness.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn change_error_f5_retry_ordinary_notification_creates_no_retry_edge() {
let mut harness = Harness::new(&[ALPHA], 1);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.arm_queue_debounce().await;
harness.fail_with_runtime_limit(ALPHA).await;
harness.queue.notify_scheduler();
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
assert!(
!harness.consume_retry_edges().await,
"a generic notification is not an explicit-retry edge"
);
harness.depends_on(&[("alpha-dep", ALPHA)]);
assert!(
harness.is_gated_by_failure("alpha-dep"),
"no failed classification may be released without accepted retry intent"
);
assert_eq!(harness.status(ALPHA).await, "error");
let (_, apply_started) = harness.drain_events();
assert_eq!(apply_started, 0);
harness.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn change_error_f5_retry_marking_alone_publishes_no_edge() {
let mut harness = Harness::new(&[ALPHA], 1);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.arm_queue_debounce().await;
harness.fail_with_runtime_limit(ALPHA).await;
harness
.run_control
.operator()
.set_execution_mark(ALPHA, true)
.await
.expect("a non-terminal row accepts a mark at any time");
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
assert!(
!harness.consume_retry_edges().await,
"an execution mark must not publish an explicit-retry edge"
);
assert_eq!(
harness.status(ALPHA).await,
"error",
"marking does not clear terminal error evidence"
);
harness.depends_on(&[("alpha-dep", ALPHA)]);
assert!(harness.is_gated_by_failure("alpha-dep"));
let (_, apply_started) = harness.drain_events();
assert_eq!(apply_started, 0);
harness.shutdown().await;
}
impl Harness {
async fn fail_with_iteration_limit(&mut self, change_id: &str, max: u32) {
for _ in 0..max {
self.executor.apply_budget.reserve(change_id, max);
}
{
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
guard.apply_execution_event(&ExecutionEvent::ProcessingError {
id: change_id.to_string(),
error: format!("reached maximum iterations ({max}/{max}) without completion"),
});
guard.record_apply_iteration_limit(change_id, max, max);
}
self.executor.failed_tracker.mark_failed(change_id);
}
fn apply_budget_exhausted(&self, change_id: &str, max: u32) -> bool {
self.executor
.apply_budget
.exhaustion(change_id, max)
.is_some()
}
}
#[tokio::test(start_paused = true)]
async fn change_error_f5_retry_iteration_limit_receives_fresh_budget_on_explicit_retry() {
const MAX: u32 = 3;
let mut harness = Harness::new(&[ALPHA, BETA], 2);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.arm_queue_debounce().await;
harness.fail_with_iteration_limit(ALPHA, MAX).await;
harness.fail_with_iteration_limit(BETA, MAX).await;
assert_eq!(harness.status(ALPHA).await, "error");
assert!(
harness.apply_budget_exhausted(ALPHA, MAX),
"the invocation really spent its ceiling"
);
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
let (_, apply_started) = harness.drain_events();
assert_eq!(
apply_started, 0,
"an iteration-limit failure must not be redispatched without explicit intent"
);
assert!(
harness.apply_budget_exhausted(ALPHA, MAX),
"and reconciliation grants no budget"
);
assert!(
harness
.state
.read()
.await
.apply_iteration_limit(ALPHA)
.is_some(),
"the diagnostic survives every automatic cycle"
);
harness.marks.set(ALPHA, true);
let outcome = harness
.run_control
.start(OperatorMode::Running)
.await
.expect("a settled iteration-limit error is startable while the run is live");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec![ALPHA.to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
},
"the live scheduler is woken rather than joined by a second boundary"
);
assert_eq!(
harness.status(ALPHA).await,
"queued",
"the terminal error was cleared through retry intent"
);
assert!(
harness
.state
.read()
.await
.apply_iteration_limit(ALPHA)
.is_none(),
"the diagnostic is consumed by the same explicit intent"
);
assert!(
harness
.state
.read()
.await
.apply_iteration_limit(BETA)
.is_some(),
"and an unretried failure keeps its own diagnostic"
);
assert!(
harness.consume_retry_edges().await,
"the accepted retry published an edge for this scheduler to consume"
);
assert!(
!harness.apply_budget_exhausted(ALPHA, MAX),
"the retried change enters its new invocation with a fresh Apply budget"
);
assert!(
harness.apply_budget_exhausted(BETA, MAX),
"release is target-specific: another exhausted change keeps its ceiling"
);
harness.run_loop_iteration(&analyzer).await;
let (analysis_started, _) = harness.drain_events();
assert!(
harness.in_flight.contains(ALPHA),
"the released change is dispatched once dependency and capacity guards allow it"
);
assert_eq!(
analysis_started, 1,
"exactly one reevaluation follows the accepted retry"
);
assert!(
!harness.in_flight.contains(BETA),
"and the unretried failure is still not admitted"
);
assert!(
matches!(
harness.executor.apply_budget.reserve(ALPHA, MAX),
crate::execution::apply::ApplyBudgetReservation::Reserved { attempt: 1, .. }
),
"the new invocation's first dispatch is admitted as attempt 1 rather than refused"
);
harness.shutdown().await;
}
#[tokio::test(start_paused = true)]
async fn change_error_f5_retry_iteration_limit_budget_survives_ordinary_notification() {
const MAX: u32 = 3;
let mut harness = Harness::new(&[ALPHA], 1);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.arm_queue_debounce().await;
harness.fail_with_iteration_limit(ALPHA, MAX).await;
harness.queue.notify_scheduler();
harness
.run_control
.operator()
.set_execution_mark(ALPHA, true)
.await
.expect("a settled limited row still accepts next-run intent");
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
assert!(
!harness.consume_retry_edges().await,
"no automatic path publishes an explicit-retry edge"
);
assert!(
harness.apply_budget_exhausted(ALPHA, MAX),
"so the exhausted budget is never released"
);
assert_eq!(harness.status(ALPHA).await, "error");
let (_, apply_started) = harness.drain_events();
assert_eq!(apply_started, 0);
harness.shutdown().await;
}
impl Harness {
fn budget_spent(&self, change_id: &str, max: u32) -> bool {
self.executor
.apply_budget
.exhaustion(change_id, max)
.is_some()
}
fn occupy_slots(&mut self, count: usize) {
for index in 0..count {
self.in_flight.insert(format!("unrelated-{index}"));
}
}
fn release_occupied_slots(&mut self) {
self.in_flight.retain(|id| !id.starts_with("unrelated-"));
}
}
#[tokio::test(start_paused = true)]
async fn retry_change_bypasses_unchanged_analysis_input_and_dispatches() {
const MAX: u32 = 3;
let mut harness = Harness::new(&[ALPHA, BETA], 1);
let analyzer = counting_analyzer(harness.analyses.clone());
harness.fail_with_runtime_limit(BETA).await;
harness.depends_on(&[("beta-dep", BETA)]);
assert!(harness.is_gated_by_failure("beta-dep"));
{
let mut guard = harness.state.write().await;
guard.apply_command(ReducerCommand::AddToQueue(ALPHA.to_string()));
}
harness.run_loop_iteration(&analyzer).await;
let baseline = harness.drain_observed_events();
assert_eq!(baseline.analysis_started, 1);
assert_eq!(
baseline.dispatch_started,
vec![ALPHA.to_string()],
"the first evaluation really dispatched the change it analysed"
);
harness.settle_dispatch_back_to_queued(ALPHA).await;
for _ in 0..MAX {
harness.executor.apply_budget.reserve(ALPHA, MAX);
}
harness.stall_in_acceptance(ALPHA).await;
assert_eq!(harness.status(ALPHA).await, "stalled");
assert!(harness.budget_spent(ALPHA, MAX));
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
assert_eq!(
harness.analyses(),
1,
"an acceptance hold is classified blocked-only; no wake analyses it"
);
assert!(
harness.queued.iter().any(|change| change.id == ALPHA),
"the held change stays reducer-visible queued work the scheduler already holds, \
so a later retry produces no scheduler-visible queue addition"
);
let plan = harness
.run_control
.operator()
.retry_change(ALPHA)
.await
.expect("a resumable acceptance stall is retryable");
assert_eq!(
plan.change_ids,
vec![ALPHA.to_string()],
"the retry is accepted through the acceptance-stall route"
);
assert_eq!(harness.status(ALPHA).await, "queued");
harness.occupy_slots(1);
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
assert_eq!(
harness.analyses(),
1,
"the pass ended at the capacity gate without evaluating anything"
);
assert_eq!(
harness.executor.pending_retry_bypass_targets(),
vec![ALPHA.to_string()],
"an abandoned pass must not discard the authority it took"
);
harness.release_occupied_slots();
harness.drain_observed_events();
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
let retried = harness.drain_observed_events();
assert_eq!(
harness.analyses(),
2,
"the retry's edge produced one fresh analyzer invocation even though the \
analysis input still matches the last completed one"
);
assert_eq!(retried.analysis_started, 1);
assert_eq!(
retried.dispatch_started,
vec![ALPHA.to_string()],
"and the retried change reached real Apply dispatch, not just re-analysis"
);
assert!(harness.in_flight.contains(ALPHA));
assert!(
harness.budget_spent(ALPHA, MAX),
"a stall-route edge must not reset the retried target's Apply budget"
);
harness.depends_on(&[("beta-dep", BETA)]);
assert!(
harness.is_gated_by_failure("beta-dep"),
"and it must not release another change's failed classification"
);
harness.settle_dispatch_back_to_queued(ALPHA).await;
assert!(
harness.executor.pending_retry_bypass_targets().is_empty(),
"the authorized evaluation spent the bypass"
);
harness.queue.notify_scheduler();
for _ in 0..3 {
harness.timer_wake().await;
harness.run_loop_iteration(&analyzer).await;
}
let after = harness.drain_observed_events();
assert_eq!(
harness.analyses(),
2,
"ordinary timer wakes and a generic notification observe the same analysis \
input and are suppressed exactly as before"
);
assert_eq!(after.analysis_started, 0);
assert!(after.dispatch_started.is_empty());
harness.shutdown().await;
}