use crate::analyzer::{AnalysisOutcome, AnalysisProvenance, AnalysisResult};
use crate::config::OrchestratorConfig;
use crate::events::{ExecutionEvent, StalledBlocker};
use crate::openspec::{Change, ProposalMetadata};
use crate::orchestration::state::{OrchestratorState, ReducerCommand};
use crate::parallel::work_snapshot::ReducerWorkSnapshot;
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, Mutex};
use std::time::Duration;
use tempfile::TempDir;
use tokio::sync::{mpsc, Notify, RwLock};
use tokio_util::sync::CancellationToken;
const MUST_HAPPEN: Duration = Duration::from_secs(5);
const NO_PROGRESS_WINDOW: Duration = Duration::from_millis(120);
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) {
let change_dir = repo_root.join("openspec/changes").join(change_id);
std::fs::create_dir_all(&change_dir).expect("create change directory");
std::fs::write(
change_dir.join("proposal.md"),
format!("---\ndependencies:\n---\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(change_ids: &[&str]) -> 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");
for change_id in change_ids {
write_change(root, change_id);
}
git(root, &["add", "-A"]);
git(root, &["commit", "-m", "Base"]);
repo
}
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 acceptance_hold(change_id: &str) -> ExecutionEvent {
ExecutionEvent::AcceptanceGated {
change_id: change_id.to_string(),
blocker: StalledBlocker {
category: "external_service".to_string(),
phase: "acceptance".to_string(),
gate: "acceptance".to_string(),
error_summary: "waiting on an operator-owned prerequisite".to_string(),
evidence: vec!["tasks.md: Implementation Blocker #1".to_string()],
unblock_condition: None,
prerequisite_owner: None,
next_action: "operator clears the prerequisite".to_string(),
resumable: true,
worktree_preserved: true,
},
}
}
#[derive(Default)]
struct AnalyzerProbe {
invocations: AtomicUsize,
analyzed_ids: Mutex<Vec<Vec<String>>>,
reducer_writable_during_analysis: AtomicUsize,
started: Notify,
}
impl AnalyzerProbe {
fn invocations(&self) -> usize {
self.invocations.load(Ordering::SeqCst)
}
fn analyzed_ids(&self) -> Vec<Vec<String>> {
self.analyzed_ids.lock().expect("probe mutex").clone()
}
}
fn probing_analyzer(
probe: Arc<AnalyzerProbe>,
shared: Arc<RwLock<OrchestratorState>>,
) -> impl for<'a> Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync {
move |changes: &[Change], _in_flight: &[String], _iteration: u32| -> AnalysisFuture<'_> {
let probe = probe.clone();
let shared = shared.clone();
let ids: Vec<String> = changes.iter().map(|change| change.id.clone()).collect();
Box::pin(async move {
probe.invocations.fetch_add(1, Ordering::SeqCst);
probe
.analyzed_ids
.lock()
.expect("probe mutex")
.push(ids.clone());
if shared.try_write().is_ok() {
probe
.reducer_writable_during_analysis
.fetch_add(1, Ordering::SeqCst);
}
probe.started.notify_waiters();
AnalysisOutcome::new(
AnalysisResult {
order: Vec::new(),
dependencies: HashMap::new(),
groups: None,
},
AnalysisProvenance::HealthyLlm,
)
})
}
}
#[derive(Default)]
struct TerminalEvents {
all_completed: AtomicUsize,
stopped: AtomicUsize,
}
fn spawn_event_recorder(
mut events: mpsc::Receiver<ParallelEvent>,
recorded: Arc<TerminalEvents>,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
while let Some(event) = events.recv().await {
match event {
ParallelEvent::AllCompleted => {
recorded.all_completed.fetch_add(1, Ordering::SeqCst);
}
ParallelEvent::Stopped => {
recorded.stopped.fetch_add(1, Ordering::SeqCst);
}
_ => {}
}
}
})
}
struct ContendedRun {
executor: ParallelExecutor,
shared: Arc<RwLock<OrchestratorState>>,
queue: Arc<DynamicQueue>,
cancel: CancellationToken,
probe: Arc<AnalyzerProbe>,
terminal: Arc<TerminalEvents>,
_recorder: tokio::task::JoinHandle<()>,
_workspace_base: TempDir,
}
fn wire_run(
repo_root: &std::path::Path,
shared: Arc<RwLock<OrchestratorState>>,
lifetime: SchedulerLifetime,
) -> ContendedRun {
let workspace_base = TempDir::new().expect("create workspace base");
let (event_tx, events) = mpsc::channel(512);
let mut executor = ParallelExecutor::new(
repo_root.to_path_buf(),
test_config(workspace_base.path()),
Some(event_tx),
);
let queue = Arc::new(DynamicQueue::new());
let cancel = CancellationToken::new();
executor.set_shared_orchestrator_state(shared.clone());
executor.set_dynamic_queue(queue.clone());
executor.set_scheduler_lifetime(lifetime);
executor.set_cancel_token(cancel.clone());
let terminal = Arc::new(TerminalEvents::default());
let recorder = spawn_event_recorder(events, terminal.clone());
ContendedRun {
executor,
probe: Arc::new(AnalyzerProbe::default()),
shared,
queue,
cancel,
terminal,
_recorder: recorder,
_workspace_base: workspace_base,
}
}
#[tokio::test]
async fn reducer_snapshot_contention_does_not_drain_or_strand_a_finite_scheduler() {
let repo = init_repo(&["queued-a"]);
let shared = reducer_state(&["queued-a"], &["queued-a"]).await;
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Finite);
let ContendedRun {
mut executor,
queue,
cancel,
probe,
terminal,
..
} = run;
assert!(queue.push("queued-a".to_string()).await);
let contended = shared.write().await;
let analyzer = probing_analyzer(probe.clone(), shared.clone());
let mut scheduler = tokio::spawn(async move {
executor
.execute_with_order_based_reanalysis(Vec::new(), analyzer)
.await
});
assert!(
tokio::time::timeout(NO_PROGRESS_WINDOW, &mut scheduler)
.await
.is_err(),
"a finite scheduler must not return while reducer evidence is unavailable"
);
assert_eq!(
probe.invocations(),
0,
"no dependency analysis may start from incomplete reducer evidence"
);
assert_eq!(
terminal.all_completed.load(Ordering::SeqCst),
0,
"contention must never be announced as a completed run"
);
assert!(
queue.contains("queued-a").await,
"a hint the pass cannot yet judge keeps its wake edge instead of being consumed"
);
drop(contended);
tokio::time::timeout(MUST_HAPPEN, probe.started.notified())
.await
.expect("the same evaluation must continue once the writer releases");
assert_eq!(
probe.analyzed_ids().first().cloned().unwrap_or_default(),
vec!["queued-a".to_string()],
"the reducer-queued candidate is reconciled into the scheduler-local queue and analyzed"
);
assert_eq!(
probe
.reducer_writable_during_analysis
.load(Ordering::SeqCst),
1,
"the reducer read guard must be released before dependency analysis begins"
);
cancel.cancel();
let report = tokio::time::timeout(MUST_HAPPEN, scheduler)
.await
.expect("cancelled scheduler must return")
.expect("scheduler task must not panic")
.expect("scheduler loop must not fail");
assert_eq!(
report,
SchedulerRunReport::Stopped,
"the run ended by cancellation, not by a contention-derived drain"
);
assert_eq!(
terminal.all_completed.load(Ordering::SeqCst),
0,
"queued work was never drained, so completion must never be announced"
);
}
#[tokio::test]
async fn reducer_snapshot_contention_release_preserves_a_real_acceptance_hold() {
let repo = init_repo(&["held-a"]);
let shared = reducer_state(&["held-a"], &["held-a"]).await;
{
let mut guard = shared.write().await;
guard.apply_execution_event(&acceptance_hold("held-a"));
assert!(
guard.acceptance_stalled_change_ids().contains("held-a"),
"fixture must install a reducer-owned Acceptance hold"
);
assert!(
guard.queued_change_ids().contains(&"held-a".to_string()),
"the held change must still carry queue intent, or the test proves nothing"
);
}
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Finite);
let ContendedRun {
mut executor,
queue,
probe,
terminal,
..
} = run;
assert!(queue.push("held-a".to_string()).await);
let contended = shared.write().await;
let analyzer = probing_analyzer(probe.clone(), shared.clone());
let mut scheduler = tokio::spawn(async move {
executor
.execute_with_order_based_reanalysis(Vec::new(), analyzer)
.await
});
assert!(
tokio::time::timeout(NO_PROGRESS_WINDOW, &mut scheduler)
.await
.is_err(),
"even a run that will end blocked must not report that while evidence is unavailable"
);
drop(contended);
let report = tokio::time::timeout(MUST_HAPPEN, scheduler)
.await
.expect("the released writer must let the blocked-only decision complete")
.expect("scheduler task must not panic")
.expect("scheduler loop must not fail");
assert_eq!(
report,
SchedulerRunReport::BlockedOrStalled,
"a real hold is stable blocked-only work once evidence is coherent"
);
assert_eq!(
probe.invocations(),
0,
"a held candidate must not reach dependency analysis or ordinary dispatch"
);
assert_eq!(
terminal.all_completed.load(Ordering::SeqCst),
0,
"blocked-only work is not completion"
);
assert!(
shared
.read()
.await
.acceptance_stalled_change_ids()
.contains("held-a"),
"the hold is still reducer-owned; the resumed pass did not repeat acceptance"
);
}
#[tokio::test]
async fn reducer_snapshot_contention_release_admits_the_candidate_to_dispatch() {
let repo = init_repo(&["queued-a"]);
let shared = reducer_state(&["queued-a"], &["queued-a"]).await;
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Finite);
let ContendedRun { mut executor, .. } = run;
let contended = shared.write().await;
let mut selection = tokio::spawn(async move {
let mut queued: Vec<Change> = Vec::new();
let in_flight = HashSet::new();
let reconciled = executor
.reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
.await;
let analysis = AnalysisResult {
order: queued.iter().map(|change| change.id.clone()).collect(),
dependencies: HashMap::new(),
groups: None,
};
let selected = executor
.select_changes_for_dispatch(&analysis, 1, &in_flight)
.await;
(reconciled.queued_added, selected)
});
assert!(
tokio::time::timeout(NO_PROGRESS_WINDOW, &mut selection)
.await
.is_err(),
"reconciliation and dispatch selection must both wait for coherent evidence"
);
drop(contended);
let (queued_added, selected) = tokio::time::timeout(MUST_HAPPEN, selection)
.await
.expect("the released writer must let reconciliation and selection finish")
.expect("selection task must not panic");
assert_eq!(
queued_added, 1,
"reducer-visible queue intent is reconciled into the scheduler-local queue"
);
assert_eq!(
selected,
vec!["queued-a".to_string()],
"the resumed snapshot admits the ordinary candidate to dispatch"
);
}
#[tokio::test]
async fn reducer_snapshot_contention_never_authorizes_termination_or_idle() {
let repo = init_repo(&[]);
let shared = reducer_state(&[], &[]).await;
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Finite);
let ContendedRun { mut executor, .. } = run;
let incomplete = ReducerWorkSnapshot::incomplete();
let complete = executor.capture_reducer_work_snapshot().await;
let empty_queue: Vec<Change> = Vec::new();
let in_flight = HashSet::new();
assert!(
!executor
.should_exit_when_idle(true, &empty_queue, &in_flight, Some(&incomplete))
.await,
"a finite run must not report DrainedSuccessfully or BlockedOrStalled on incomplete evidence"
);
assert!(
executor
.should_exit_when_idle(true, &empty_queue, &in_flight, Some(&complete))
.await,
"the same state does drain once the reducer view is coherent"
);
executor.set_persistent_lifetime();
assert!(
!executor
.should_enter_persistent_idle_wait(true, &empty_queue, &in_flight, Some(&incomplete))
.await,
"a persistent run must not park in the timer-free idle wait on incomplete evidence"
);
assert!(
executor
.should_enter_persistent_idle_wait(true, &empty_queue, &in_flight, Some(&complete))
.await,
"stable, genuinely drained state still enters event-driven persistent idle"
);
}
#[tokio::test]
async fn reducer_snapshot_contention_retains_an_unjudged_dynamic_hint() {
let repo = init_repo(&["queued-a"]);
let shared = reducer_state(&["queued-a"], &["queued-a"]).await;
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Finite);
let ContendedRun {
mut executor,
queue,
..
} = run;
assert!(queue.push("queued-a".to_string()).await);
let mut queued: Vec<Change> = Vec::new();
let in_flight = HashSet::new();
let mut reason = crate::parallel::dynamic_queue::ReanalysisReason::Initial;
let ingested = executor
.check_dynamic_queue_and_add_changes_with_snapshot(
&mut queued,
&in_flight,
&mut reason,
&ReducerWorkSnapshot::incomplete(),
)
.await;
assert!(!ingested, "incomplete evidence admits nothing");
assert!(queued.is_empty(), "and adds no scheduler-local candidate");
assert!(
queue.contains("queued-a").await,
"the hint must survive so its wake edge is not lost"
);
let snapshot = executor.capture_reducer_work_snapshot().await;
assert!(
executor
.check_dynamic_queue_and_add_changes_with_snapshot(
&mut queued,
&in_flight,
&mut reason,
&snapshot,
)
.await,
"the retained hint is ingested once the reducer view is coherent"
);
assert_eq!(
queued
.iter()
.map(|change| change.id.as_str())
.collect::<Vec<_>>(),
vec!["queued-a"]
);
}
#[tokio::test]
async fn reducer_snapshot_contention_stays_cancellable_while_acquisition_is_pending() {
let repo = init_repo(&["queued-a"]);
let shared = reducer_state(&["queued-a"], &["queued-a"]).await;
let run = wire_run(repo.path(), shared.clone(), SchedulerLifetime::Persistent);
let ContendedRun {
mut executor,
queue,
cancel,
probe,
..
} = run;
assert!(queue.push("queued-a".to_string()).await);
let contended = shared.write().await;
let analyzer = probing_analyzer(probe.clone(), shared.clone());
let mut scheduler = tokio::spawn(async move {
executor
.execute_with_order_based_reanalysis(Vec::new(), analyzer)
.await
});
assert!(
tokio::time::timeout(NO_PROGRESS_WINDOW, &mut scheduler)
.await
.is_err(),
"the scheduler must still be suspended on snapshot acquisition"
);
cancel.cancel();
let report = tokio::time::timeout(MUST_HAPPEN, scheduler)
.await
.expect("cancellation must terminate a scheduler that is waiting for a reducer snapshot")
.expect("scheduler task must not panic")
.expect("scheduler loop must not fail");
assert_eq!(report, SchedulerRunReport::Stopped);
assert_eq!(
probe.invocations(),
0,
"a cancelled acquisition must not fall through into analysis"
);
drop(contended);
}