use crate::analyzer::{AnalysisOutcome, AnalysisProvenance, AnalysisResult};
use crate::config::OrchestratorConfig;
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use crate::orchestration::operator_command::RetryEdgeAuthority;
use crate::orchestration::state::{OrchestratorState, ReducerCommand};
use crate::parallel::lifecycle_slots::SlotPhase;
use crate::parallel::queue_state::{QueuedWorkClass, RetryEdgeConsumption};
use crate::parallel::{ParallelEvent, ParallelExecutor, SchedulerLifetime, SchedulerRunReport};
use crate::tui::queue::DynamicQueue;
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 std::time::Duration;
use tempfile::TempDir;
use tokio::sync::{mpsc, RwLock};
use tokio_util::sync::CancellationToken;
const DEBOUNCE_ELIGIBLE_WAKE: Duration = Duration::from_secs(30);
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 {:?} failed", args);
}
fn write_change(repo_root: &std::path::Path, change_id: &str, dependencies: &[&str]) {
let change_dir = repo_root.join("openspec/changes").join(change_id);
std::fs::create_dir_all(&change_dir).expect("create change directory");
let declared = dependencies
.iter()
.map(|dependency| format!(" - {dependency}\n"))
.collect::<String>();
std::fs::write(
change_dir.join("proposal.md"),
format!("---\ndependencies:\n{declared}---\n# {change_id}\n"),
)
.expect("write proposal");
std::fs::write(
change_dir.join("tasks.md"),
"## Implementation Tasks\n- [ ] apply\n",
)
.expect("write tasks");
}
fn init_repo() -> TempDir {
let repo = TempDir::new().expect("create temp repo");
let root = repo.path();
git(root, &["init", "-b", "main"]);
git(root, &["config", "user.email", "test@example.com"]);
git(root, &["config", "user.name", "Test User"]);
std::fs::write(root.join("README.md"), "base\n").expect("write base file");
write_change(root, "a", &[]);
write_change(root, "b", &["a"]);
write_change(root, "c", &[]);
git(root, &["add", "-A"]);
git(root, &["commit", "-m", "Base"]);
repo
}
fn archive_to_base(repo_root: &std::path::Path, archive_leaf: &str, change_id: &str) {
let archive_dir = repo_root
.join("openspec/changes/archive")
.join(archive_leaf);
std::fs::create_dir_all(&archive_dir).expect("create archive directory");
std::fs::write(
archive_dir.join("proposal.md"),
format!("# Archived {change_id}\n"),
)
.expect("write archived proposal");
let active_dir = repo_root.join("openspec/changes").join(change_id);
if active_dir.exists() {
std::fs::remove_dir_all(&active_dir).expect("remove active change directory");
}
git(repo_root, &["add", "-A"]);
git(
repo_root,
&["commit", "-m", &format!("Archive {change_id}")],
);
}
async fn reducer_state(known: &[&str], queued_ids: &[&str]) -> Arc<RwLock<OrchestratorState>> {
let state = Arc::new(RwLock::new(OrchestratorState::new(
known.iter().map(|id| id.to_string()).collect(),
10,
)));
{
let mut guard = state.write().await;
for id in queued_ids {
guard.apply_command(ReducerCommand::AddToQueue(id.to_string()));
}
}
state
}
fn seed_failed_blocker(executor: &mut ParallelExecutor, dependent: &str, blocker: &str) {
let mut dependencies = HashMap::new();
dependencies.insert(dependent.to_string(), vec![blocker.to_string()]);
executor.failed_tracker.set_dependencies(dependencies);
executor.failed_tracker.mark_failed(blocker);
}
fn dependency_analysis(
order: &[&str],
dependencies: HashMap<String, Vec<String>>,
) -> AnalysisResult {
AnalysisResult {
order: order.iter().map(|id| id.to_string()).collect(),
dependencies,
groups: None,
}
}
fn b_depends_on_a() -> HashMap<String, Vec<String>> {
let mut dependencies = HashMap::new();
dependencies.insert("b".to_string(), vec!["a".to_string()]);
dependencies
}
fn counting_analyzer(
analyses: Arc<AtomicUsize>,
) -> impl for<'a> Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync {
move |changes: &[Change], _in_flight: &[String], _iteration: u32| -> AnalysisFuture<'_> {
analyses.fetch_add(1, Ordering::SeqCst);
let order: Vec<String> = changes.iter().map(|change| change.id.clone()).collect();
Box::pin(async move {
AnalysisOutcome::new(
dependency_analysis(
&order.iter().map(String::as_str).collect::<Vec<_>>(),
b_depends_on_a(),
),
AnalysisProvenance::HealthyLlm,
)
})
}
}
#[tokio::test]
async fn failed_dependency_keeps_dependent_locally_queued_and_blocked() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["b", "c"]).await;
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
executor.set_shared_orchestrator_state(state);
seed_failed_blocker(&mut executor, "b", "a");
let mut queued = vec![test_change("b"), test_change("c")];
let in_flight = HashSet::new();
let classification = executor.classify_queued_work(&queued, &in_flight).await;
assert_eq!(
classification.class_for("b"),
Some(QueuedWorkClass::DependencyBlocked),
"a failed dependency must classify its dependent as blocked queued work"
);
assert_eq!(
classification.class_for("c"),
Some(QueuedWorkClass::DispatchableApply),
"independent queued work must stay dispatchable"
);
for pass in 0..5 {
let outcome = executor
.reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
.await;
assert_eq!(
outcome.queued_added, 0,
"pass {pass}: rediscovering an already represented blocked candidate is not a queue addition"
);
assert_eq!(
outcome.revoked_removed, 0,
"pass {pass}: accepted queue intent must not be revoked by dependency blocking"
);
assert!(
queued.iter().any(|change| change.id == "b"),
"pass {pass}: the blocked dependent must remain locally represented"
);
}
}
#[tokio::test]
async fn failed_dependency_blocks_dispatch_selection_but_not_independent_work() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["b", "c"]).await;
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
executor.set_shared_orchestrator_state(state);
seed_failed_blocker(&mut executor, "b", "a");
let analysis = dependency_analysis(&["b", "c"], b_depends_on_a());
let selected = executor
.select_changes_for_dispatch(&analysis, 4, &HashSet::new())
.await;
assert_eq!(
selected,
vec!["c".to_string()],
"the failed-dependent candidate must be excluded from dispatch while the \
independent candidate still reaches it"
);
}
struct EpochProbe {
executor: ParallelExecutor,
events: mpsc::Receiver<ExecutionEvent>,
queued: Vec<Change>,
in_flight: HashSet<String>,
analyses: Arc<AtomicUsize>,
_repo: TempDir,
_workspace_base: TempDir,
}
#[derive(Debug, Default, PartialEq, Eq)]
struct EmittedEvents {
change_skipped: Vec<String>,
dependency_blocked: Vec<String>,
}
impl EpochProbe {
async fn new(queued_ids: &[&str]) -> Self {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], queued_ids).await;
let (tx, events) = mpsc::channel(256);
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
Some(tx),
);
executor.set_shared_orchestrator_state(state);
seed_failed_blocker(&mut executor, "b", "a");
Self {
executor,
events,
queued: queued_ids.iter().map(|id| test_change(id)).collect(),
in_flight: HashSet::new(),
analyses: Arc::new(AtomicUsize::new(0)),
_repo: repo,
_workspace_base: workspace_base,
}
}
async fn wake(&mut self, iteration: u32) -> usize {
let reconciliation = self
.executor
.reconcile_queued_candidates_from_shared_state(&mut self.queued, &self.in_flight)
.await;
let analyzer = counting_analyzer(self.analyses.clone());
let mut reanalysis_reason = if reconciliation.has_queued_additions() {
crate::parallel::dynamic_queue::ReanalysisReason::QueueNotification
} else {
crate::parallel::dynamic_queue::ReanalysisReason::Initial
};
let mut join_set = tokio::task::JoinSet::new();
let mut cleanup_guard = crate::parallel::cleanup::WorkspaceCleanupGuard::new(
crate::vcs::VcsBackend::Git,
self.executor.repo_root.clone(),
);
self.executor
.evaluate_queued_reanalysis_and_dispatch(
crate::parallel::queue_state::ReanalysisDispatchContext {
queued: &mut self.queued,
in_flight: &mut self.in_flight,
max_parallelism: 0,
iteration,
reanalysis_reason,
analyzer: &analyzer,
join_set: &mut join_set,
cleanup_guard: &mut cleanup_guard,
work_snapshot: None,
},
&mut reanalysis_reason,
)
.await
.expect("scheduler evaluation must not fail");
reconciliation.queued_added
}
fn drain(&mut self) -> EmittedEvents {
let mut emitted = EmittedEvents::default();
while let Ok(event) = self.events.try_recv() {
match event {
ExecutionEvent::ChangeSkipped { change_id, reason } => {
assert!(
reason.contains("Dependency 'a' failed"),
"the compatibility reason must keep its existing wording: {reason}"
);
emitted.change_skipped.push(change_id);
}
ExecutionEvent::DependencyBlocked {
change_id,
dependency_ids,
} if dependency_ids.iter().any(|id| id == "a") => {
emitted.dependency_blocked.push(change_id);
}
_ => {}
}
}
emitted
}
}
#[tokio::test]
async fn failed_dependency_events_are_emitted_once_per_blocker_epoch() {
let mut probe = EpochProbe::new(&["b"]).await;
for iteration in 2..8 {
let added = probe.wake(iteration).await;
assert_eq!(
added, 0,
"iteration {iteration}: an already represented blocked candidate is not a new queue edge"
);
}
let emitted = probe.drain();
assert_eq!(
emitted.change_skipped,
vec!["b".to_string()],
"exactly one compatibility ChangeSkipped per blocker epoch"
);
assert_eq!(
emitted.dependency_blocked,
vec!["b".to_string()],
"exactly one authoritative DependencyBlocked per blocker epoch"
);
assert_eq!(
probe.analyses.load(Ordering::SeqCst),
0,
"rediscovering the same blocked candidate must not invoke the analyzer"
);
assert!(
probe.queued.iter().any(|change| change.id == "b"),
"the blocked dependent must still be locally represented"
);
}
#[tokio::test]
async fn failed_dependency_events_reopen_when_the_blocker_set_changes() {
let mut probe = EpochProbe::new(&["b"]).await;
probe.wake(2).await;
assert_eq!(probe.drain().change_skipped, vec!["b".to_string()]);
let mut dependencies = HashMap::new();
dependencies.insert("b".to_string(), vec!["a".to_string(), "c".to_string()]);
probe.executor.failed_tracker.set_dependencies(dependencies);
probe.executor.failed_tracker.mark_failed("c");
probe.wake(3).await;
probe.wake(4).await;
let emitted = probe.drain();
assert_eq!(
emitted.change_skipped,
vec!["b".to_string()],
"a changed blocker set opens exactly one new epoch"
);
assert_eq!(emitted.dependency_blocked, vec!["b".to_string()]);
}
#[tokio::test]
async fn failed_dependency_dequeue_readd_is_a_genuine_state_change() {
let mut probe = EpochProbe::new(&["b"]).await;
probe.wake(2).await;
assert_eq!(probe.drain().change_skipped, vec!["b".to_string()]);
{
let state = probe
.executor
.shared_orchestrator_state
.clone()
.expect("reducer state");
let mut guard = state.write().await;
guard.apply_command(ReducerCommand::RemoveFromQueue("b".to_string()));
}
probe.wake(3).await;
assert!(
!probe.queued.iter().any(|change| change.id == "b"),
"revocation must drop the blocked candidate from the local queue"
);
assert!(
!probe.executor.failed_tracker.has_blocker_epoch("b"),
"revocation must clear the blocker notification epoch"
);
{
let state = probe
.executor
.shared_orchestrator_state
.clone()
.expect("reducer state");
let mut guard = state.write().await;
guard.apply_command(ReducerCommand::AddToQueue("b".to_string()));
}
let added = probe.wake(4).await;
assert_eq!(added, 1, "an explicit re-add is a genuine queue addition");
probe.wake(5).await;
let emitted = probe.drain();
assert_eq!(
emitted.change_skipped,
vec!["b".to_string()],
"re-add announces exactly one new bounded blocker transition"
);
assert_eq!(emitted.dependency_blocked, vec!["b".to_string()]);
}
#[tokio::test]
async fn failed_dependency_retry_edge_clears_only_the_retried_failure() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["b"]).await;
let queue = Arc::new(DynamicQueue::new());
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
executor.set_shared_orchestrator_state(state);
executor.set_dynamic_queue(queue.clone());
let mut dependencies = HashMap::new();
dependencies.insert("b".to_string(), vec!["a".to_string()]);
dependencies.insert("other".to_string(), vec!["unrelated".to_string()]);
executor.failed_tracker.set_dependencies(dependencies);
executor.failed_tracker.mark_failed("a");
executor.failed_tracker.mark_failed("unrelated");
let blockers = executor.failed_tracker.failed_blockers("b");
executor.failed_tracker.begin_blocker_epoch("b", &blockers);
let other_blockers = executor.failed_tracker.failed_blockers("other");
executor
.failed_tracker
.begin_blocker_epoch("other", &other_blockers);
assert_eq!(
executor.consume_explicit_retry_edges().await,
RetryEdgeConsumption::default(),
"no pending edge means no armed reevaluation"
);
assert_eq!(
executor.failed_tracker.should_skip("b"),
Some("a".to_string())
);
queue
.publish_explicit_retry("a".to_string(), RetryEdgeAuthority::TerminalError)
.await;
assert_eq!(
executor.consume_explicit_retry_edges().await,
RetryEdgeConsumption {
newly_drained: 1,
bypass_armed: true,
},
"a drained edge must arm a reevaluation"
);
assert!(
executor.failed_tracker.should_skip("b").is_none(),
"the retried change's ephemeral failed classification must be cleared"
);
assert!(
!executor.failed_tracker.has_blocker_epoch("b"),
"the dependent's blocker notification epoch must be cleared"
);
assert_eq!(
executor.failed_tracker.should_skip("other"),
Some("unrelated".to_string()),
"an unrelated failure must survive"
);
assert!(
executor.failed_tracker.has_blocker_epoch("other"),
"an unrelated notification epoch must survive"
);
executor.failed_tracker.mark_failed("a");
assert_eq!(
executor.consume_explicit_retry_edges().await,
RetryEdgeConsumption {
newly_drained: 0,
bypass_armed: true,
},
"the drained edge is not replayable, and its unspent bypass survives the pass"
);
assert_eq!(
executor.failed_tracker.should_skip("b"),
Some("a".to_string()),
"a retained bypass releases no second failed classification"
);
assert_eq!(
executor.pending_retry_bypass_targets(),
vec!["a".to_string()]
);
}
#[tokio::test]
async fn failed_dependency_retry_does_not_prove_dependency_resolution() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["a", "b"]).await;
let queue = Arc::new(DynamicQueue::new());
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
executor.set_shared_orchestrator_state(state);
executor.set_dynamic_queue(queue.clone());
seed_failed_blocker(&mut executor, "b", "a");
queue
.publish_explicit_retry("a".to_string(), RetryEdgeAuthority::TerminalError)
.await;
assert!(executor.consume_explicit_retry_edges().await.bypass_armed);
let queued = vec![test_change("a"), test_change("b")];
let classification = executor
.classify_queued_work(&queued, &HashSet::new())
.await;
assert_eq!(
classification.class_for("b"),
Some(QueuedWorkClass::DependencyBlocked),
"clearing the fast failure gate must not make an unmerged dependency look resolved"
);
executor.failed_tracker.mark_failed("a");
let blockers = executor.failed_tracker.failed_blockers("b");
assert!(
executor.failed_tracker.begin_blocker_epoch("b", &blockers),
"refailure after an accepted retry is a new blocker epoch"
);
}
#[tokio::test]
async fn failed_dependency_retry_authoritative_resolution_unblocks_dependent() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["b"]).await;
let queue = Arc::new(DynamicQueue::new());
let mut executor = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
executor.set_shared_orchestrator_state(state);
executor.set_dynamic_queue(queue.clone());
seed_failed_blocker(&mut executor, "b", "a");
queue
.publish_explicit_retry("a".to_string(), RetryEdgeAuthority::TerminalError)
.await;
assert!(executor.consume_explicit_retry_edges().await.bypass_armed);
let analysis = dependency_analysis(&["b"], b_depends_on_a());
assert!(
executor
.select_changes_for_dispatch(&analysis, 2, &HashSet::new())
.await
.is_empty(),
"a retried but unmerged dependency must keep its dependent blocked"
);
archive_to_base(repo.path(), "2026-05-12-a", "a");
assert_eq!(
executor
.select_changes_for_dispatch(&analysis, 2, &HashSet::new())
.await,
vec!["b".to_string()],
"authoritative repository evidence is what unblocks the dependent"
);
}
#[tokio::test]
async fn failed_dependency_restart_discards_ephemeral_failure_tracking() {
let repo = init_repo();
let workspace_base = TempDir::new().expect("workspace base");
let state = reducer_state(&["a", "b", "c"], &["b"]).await;
let mut first = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
first.set_shared_orchestrator_state(state.clone());
seed_failed_blocker(&mut first, "b", "a");
let blockers = first.failed_tracker.failed_blockers("b");
first.failed_tracker.begin_blocker_epoch("b", &blockers);
assert!(first.failed_tracker.has_blocker_epoch("b"));
drop(first);
let restarted = ParallelExecutor::new(
repo.path().to_path_buf(),
test_config(workspace_base.path()),
None,
);
assert!(
restarted.failed_tracker.failed_changes().is_empty(),
"a restarted process must begin with an empty failure tracker"
);
assert!(
!restarted.failed_tracker.has_blocker_epoch("b"),
"a restarted process must begin with no blocker notification epochs"
);
assert!(
restarted.failed_tracker.should_skip("b").is_none(),
"routing after restart is recomputed from workspace and Git evidence"
);
}
#[derive(Default)]
struct ObservedEvents {
skips: AtomicUsize,
blocked: AtomicUsize,
capacity_gated: AtomicUsize,
}
struct LoopHandles {
queue: Arc<DynamicQueue>,
state: Arc<RwLock<OrchestratorState>>,
cancel: CancellationToken,
observed: Arc<ObservedEvents>,
analyses: Arc<AtomicUsize>,
}
async fn await_counter(counter: &AtomicUsize, target: usize, what: &str) {
for _ in 0..1_000_000 {
if counter.load(Ordering::SeqCst) >= target {
return;
}
tokio::task::yield_now().await;
}
panic!("scheduler loop never reached {target} {what}");
}
impl LoopHandles {
async fn await_skips(&self, target: usize) {
await_counter(&self.observed.skips, target, "ChangeSkipped events").await;
await_counter(&self.observed.blocked, target, "DependencyBlocked events").await;
}
async fn await_analyses(&self, target: usize) {
await_counter(&self.analyses, target, "analyzer invocations").await;
}
async fn await_capacity_gated(&self, target: usize) {
await_counter(
&self.observed.capacity_gated,
target,
"zero-capacity no-analysis diagnostics",
)
.await;
}
}
struct LoopRun {
report: SchedulerRunReport,
analyses: usize,
change_skipped: Vec<String>,
dependency_blocked: Vec<String>,
all_completed: bool,
stopped: bool,
}
async fn run_loop<D, Fut>(
repo_root: &std::path::Path,
lifetime: SchedulerLifetime,
queued_ids: &[&str],
occupy_dispatch_capacity: bool,
drive: D,
) -> LoopRun
where
D: FnOnce(LoopHandles) -> Fut,
Fut: Future<Output = ()> + Send + 'static,
{
let workspace_base = TempDir::new().expect("workspace base");
let (event_tx, mut events) = mpsc::channel(512);
let state = reducer_state(&["a", "b", "c"], queued_ids).await;
let queue = Arc::new(DynamicQueue::new());
let cancel_token = CancellationToken::new();
let mut executor = ParallelExecutor::new(
repo_root.to_path_buf(),
test_config(workspace_base.path()),
Some(event_tx),
);
executor.set_shared_orchestrator_state(state.clone());
executor.set_dynamic_queue(queue.clone());
executor.set_scheduler_lifetime(lifetime);
executor.set_cancel_token(cancel_token.clone());
if occupy_dispatch_capacity {
executor.auto_resolve_count.store(64, Ordering::SeqCst);
for index in 0..executor.configured_max_concurrent() {
executor
.lifecycle_slots
.occupy_now(&format!("resolving-{index}"), SlotPhase::Merge)
.await;
}
}
seed_failed_blocker(&mut executor, "b", "a");
let analyses = Arc::new(AtomicUsize::new(0));
let analyzer = counting_analyzer(analyses.clone());
let observed = Arc::new(ObservedEvents::default());
let forwarder_observed = observed.clone();
let forwarder = tokio::spawn(async move {
let mut collected = Vec::new();
while let Some(event) = events.recv().await {
match &event {
ParallelEvent::ChangeSkipped { .. } => {
forwarder_observed.skips.fetch_add(1, Ordering::SeqCst);
}
ParallelEvent::DependencyBlocked { .. } => {
forwarder_observed.blocked.fetch_add(1, Ordering::SeqCst);
}
ParallelEvent::Log(entry) if entry.message.contains("analysis_capacity_zero") => {
forwarder_observed
.capacity_gated
.fetch_add(1, Ordering::SeqCst);
}
_ => {}
}
collected.push(event);
}
collected
});
let driver = tokio::spawn(drive(LoopHandles {
queue,
state,
cancel: cancel_token,
observed,
analyses: analyses.clone(),
}));
let report = executor
.execute_with_order_based_reanalysis(
queued_ids.iter().map(|id| test_change(id)).collect(),
analyzer,
)
.await
.expect("scheduler loop must not fail");
driver.abort();
drop(executor);
let collected = forwarder.await.expect("event forwarder must not panic");
let mut run = LoopRun {
report,
analyses: analyses.load(Ordering::SeqCst),
change_skipped: Vec::new(),
dependency_blocked: Vec::new(),
all_completed: false,
stopped: false,
};
for event in collected {
match event {
ParallelEvent::AllCompleted => run.all_completed = true,
ParallelEvent::Stopped => run.stopped = true,
ParallelEvent::ChangeSkipped { change_id, .. } => run.change_skipped.push(change_id),
ParallelEvent::DependencyBlocked {
change_id,
dependency_ids,
} if dependency_ids.iter().any(|id| id == "a") => {
run.dependency_blocked.push(change_id);
}
_ => {}
}
}
run
}
#[tokio::test(start_paused = true)]
async fn failed_dependency_lifetime_finite_reports_blocked_rather_than_completed() {
let repo = init_repo();
let run = run_loop(
repo.path(),
SchedulerLifetime::Finite,
&["b"],
false,
|_handles| async move {
std::future::pending::<()>().await;
},
)
.await;
assert_eq!(
run.report,
SchedulerRunReport::BlockedOrStalled,
"blocked-only finite scheduling must not report completion"
);
assert!(
!run.all_completed,
"a run with blocked work remaining must not emit AllCompleted"
);
assert!(!run.stopped, "nothing cancelled this run");
assert_eq!(
run.analyses, 0,
"blocked-only work must not invoke the analyzer"
);
assert_eq!(run.change_skipped, vec!["b".to_string()]);
assert_eq!(run.dependency_blocked, vec!["b".to_string()]);
}
#[tokio::test(start_paused = true)]
async fn failed_dependency_loop_converges_across_repeated_notifications() {
let repo = init_repo();
let run = run_loop(
repo.path(),
SchedulerLifetime::Persistent,
&["b"],
false,
|handles| async move {
handles.await_skips(1).await;
for _ in 0..6 {
tokio::time::sleep(DEBOUNCE_ELIGIBLE_WAKE).await;
handles.queue.notify_scheduler();
}
tokio::time::sleep(DEBOUNCE_ELIGIBLE_WAKE).await;
handles.cancel.cancel();
},
)
.await;
assert_eq!(run.report, SchedulerRunReport::Stopped);
assert!(
run.stopped,
"the loop must still have been alive when the driver cancelled it"
);
assert!(!run.all_completed, "blocked work never completes");
assert_eq!(
run.analyses, 0,
"unchanged wakes must not invoke the analyzer even once"
);
assert_eq!(
run.change_skipped,
vec!["b".to_string()],
"repeated wakes must not repeat the compatibility observation"
);
assert_eq!(
run.dependency_blocked,
vec!["b".to_string()],
"repeated wakes must not repeat the blocked transition"
);
}
#[tokio::test(start_paused = true)]
async fn failed_dependency_lifetime_admits_genuine_dynamic_additions() {
let repo = init_repo();
let run = run_loop(
repo.path(),
SchedulerLifetime::Persistent,
&["b"],
true,
|handles| async move {
handles.await_skips(1).await;
{
let mut guard = handles.state.write().await;
guard.apply_command(ReducerCommand::AddToQueue("c".to_string()));
}
handles.queue.push("c".to_string()).await;
handles.await_capacity_gated(1).await;
handles.cancel.cancel();
},
)
.await;
assert!(run.stopped);
assert_eq!(
run.analyses, 0,
"a genuine queue addition is still capacity-gated; the analyzer waits for a slot"
);
assert_eq!(
run.change_skipped,
vec!["b".to_string()],
"the still-blocked dependent must not be re-announced by unrelated new work"
);
assert_eq!(run.dependency_blocked, vec!["b".to_string()]);
}