use std::sync::Arc;
use tokio::sync::RwLock;
use super::testing::{RecordingScheduler, SchedulerCall};
use super::*;
use crate::events::{ExecutionEvent, StalledBlocker};
use crate::orchestration::operator_command::{
ExecutionMarkStore, NoopQueueHooks, OperatorCommandService, ParallelEligibility, QueuePort,
RetryEdgeAuthority, TerminationWaiter,
};
#[derive(Debug, Default)]
struct FakeQueue {
entries: Mutex<Vec<String>>,
notifications: Mutex<usize>,
explicit_retries: Mutex<Vec<(String, RetryEdgeAuthority)>>,
}
#[async_trait]
impl QueuePort for FakeQueue {
async fn add(&self, change_id: &str) -> bool {
let mut guard = self.entries.lock().unwrap();
if guard.iter().any(|id| id == change_id) {
return false;
}
guard.push(change_id.to_string());
true
}
async fn remove(&self, change_id: &str) -> bool {
let mut guard = self.entries.lock().unwrap();
let before = guard.len();
guard.retain(|id| id != change_id);
guard.len() != before
}
async fn request_cancellation(
&self,
_change_id: &str,
) -> std::result::Result<Option<TerminationWaiter>, String> {
Ok(None)
}
async fn notify_scheduler(&self) {
*self.notifications.lock().unwrap() += 1;
}
async fn publish_explicit_retry(&self, change_id: &str, authority: RetryEdgeAuthority) {
self.explicit_retries
.lock()
.unwrap()
.push((change_id.to_string(), authority));
}
}
struct Harness {
service: RunControlService,
state: Arc<RwLock<OrchestratorState>>,
scheduler: Arc<RecordingScheduler>,
marks: Arc<ExecutionMarkStore>,
queue: Arc<FakeQueue>,
operator: Arc<OperatorCommandService>,
resolves: Arc<ResolveReservations>,
eligibility: Arc<StartEligibility>,
}
struct FakeAdmission {
refused: Vec<String>,
}
#[async_trait]
impl crate::orchestration::operator_command::AcceptanceAdmissionPort for FakeAdmission {
async fn classify(
&self,
change_id: &str,
) -> crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission {
use crate::orchestration::acceptance::execution_manifest::{
AcceptanceAdmission, AcceptanceHoldCategory, FINGERPRINT_COMPONENTS,
UNCHANGED_ACCEPTANCE_INPUT,
};
if self.refused.iter().any(|id| id == change_id) {
AcceptanceAdmission::Refuse {
outcome: UNCHANGED_ACCEPTANCE_INPUT,
category: AcceptanceHoldCategory::ReviewDeadlineExhausted,
fingerprint: "a".repeat(64),
components: FINGERPRINT_COMPONENTS.to_vec(),
}
} else {
AcceptanceAdmission::Admit
}
}
}
impl Harness {
fn new(change_ids: &[&str]) -> Self {
Self::with_refused_acceptance(change_ids, &[])
}
fn with_refused_acceptance(change_ids: &[&str], refused: &[&str]) -> Self {
let admission = Arc::new(FakeAdmission {
refused: refused.iter().map(|id| id.to_string()).collect(),
});
Self::build(change_ids, Some(admission))
}
fn build(
change_ids: &[&str],
admission: Option<Arc<dyn crate::orchestration::operator_command::AcceptanceAdmissionPort>>,
) -> Self {
let state = Arc::new(RwLock::new(OrchestratorState::new(
change_ids.iter().map(|id| id.to_string()).collect(),
10,
)));
let marks = Arc::new(ExecutionMarkStore::new());
let scheduler = Arc::new(RecordingScheduler::new());
let queue = Arc::new(FakeQueue::default());
let operator = {
let service = OperatorCommandService::new(
state.clone(),
queue.clone(),
Arc::new(NoopQueueHooks),
marks.clone(),
);
Arc::new(match admission {
Some(admission) => service.with_acceptance_admission(admission),
None => service,
})
};
let resolves = Arc::new(ResolveReservations::new());
let eligibility = Arc::new(StartEligibility::new());
Self {
service: RunControlService::new(
state.clone(),
operator.clone(),
scheduler.clone(),
resolves.clone(),
eligibility.clone(),
),
state,
scheduler,
marks,
queue,
operator,
resolves,
eligibility,
}
}
fn mark(&self, change_ids: &[&str]) {
self.marks
.replace(change_ids.iter().map(|id| id.to_string()));
}
async fn effects(&self) -> RunEffects {
let queue = self.queue.entries.lock().unwrap().clone();
let notifications = *self.queue.notifications.lock().unwrap();
let explicit_retries = self.queue.explicit_retries.lock().unwrap().clone();
let statuses = {
let guard = self.state.read().await;
let mut statuses: Vec<(String, String)> = guard
.tracked_change_ids()
.into_iter()
.map(|id| {
let status = guard.display_status(&id).to_string();
(id, status)
})
.collect();
statuses.sort();
statuses
};
RunEffects {
scheduler: self.scheduler.calls(),
queue,
notifications,
explicit_retries,
statuses,
}
}
async fn apply(&self, event: ExecutionEvent) {
self.state.write().await.apply_execution_event(&event);
}
async fn status(&self, change_id: &str) -> String {
self.state
.read()
.await
.display_status(change_id)
.to_string()
}
async fn to_merge_wait(&self, change_id: &str) {
self.apply(ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "manual resolution required".to_string(),
auto_resumable: false,
})
.await;
}
async fn to_startup_merge_wait(&self, change_id: &str) {
use std::collections::{HashMap, HashSet};
self.apply(ExecutionEvent::ChangesRefreshed {
changes: Vec::new(),
rejected_changes: Vec::new(),
committed_change_ids: HashSet::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::from([change_id.to_string()]),
})
.await;
}
async fn to_error(&self, change_id: &str) {
self.apply(ExecutionEvent::ProcessingError {
id: change_id.to_string(),
error: "boom".to_string(),
})
.await;
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct RunEffects {
scheduler: Vec<SchedulerCall>,
queue: Vec<String>,
notifications: usize,
explicit_retries: Vec<(String, RetryEdgeAuthority)>,
statuses: Vec<(String, String)>,
}
fn external_blocker() -> StalledBlocker {
StalledBlocker {
category: "external_service".to_string(),
phase: "acceptance".to_string(),
gate: "prerequisite".to_string(),
error_summary: "registry unreachable".to_string(),
evidence: vec!["curl: (6) could not resolve host".to_string()],
unblock_condition: Some("registry responds to a health check".to_string()),
prerequisite_owner: Some("platform".to_string()),
next_action: "retry once the registry is reachable".to_string(),
resumable: true,
worktree_preserved: true,
}
}
#[tokio::test]
async fn start_consumes_the_authoritative_marked_target_set() {
let harness = Harness::new(&["a", "b", "c"]);
harness.mark(&["a", "c"]);
let outcome = harness
.service
.start(OperatorMode::Select)
.await
.expect("marked changes are startable");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["a".to_string(), "c".to_string()],
explicit_retry: false,
scheduler: SchedulerEffect::Started,
excluded: Vec::new(),
}
);
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["a".to_string(), "c".to_string()]],
"the run must be dispatched for exactly the marked target set"
);
assert_eq!(harness.status("a").await, "queued");
assert_eq!(harness.status("b").await, "not queued");
}
#[tokio::test]
async fn start_without_any_mark_fails_and_starts_nothing() {
let harness = Harness::new(&["a"]);
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("an empty target set is not a successful start");
assert!(matches!(
error,
RunControlError::NoEligibleTarget {
command: RunCommandKind::Start,
..
}
));
assert!(
harness.scheduler.calls().is_empty(),
"a failed start must not touch the scheduler"
);
}
#[tokio::test]
async fn start_with_only_ineligible_marks_fails_and_starts_nothing() {
let harness = Harness::new(&["a"]);
harness.mark(&["a"]);
harness.to_merge_wait("a").await;
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("a merge-wait row is not startable");
match error {
RunControlError::NoEligibleTarget { detail, .. } => {
assert!(detail.contains("not queued"), "detail must be actionable");
}
other => panic!("unexpected error: {other:?}"),
}
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn start_refuses_parallel_ineligible_targets() {
let harness = Harness::new(&["a"]);
harness.mark(&["a"]);
harness.eligibility.set_parallel_ineligible([(
"a".to_string(),
ParallelEligibility::UncommittedProposalFiles,
)]);
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("worktree execution refuses uncommitted changes");
assert!(matches!(error, RunControlError::NoEligibleTarget { .. }));
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn one_ineligible_mark_refuses_the_whole_parallel_start() {
for arrange_ineligible_as_unstartable in [false, true] {
let harness = Harness::new(&["eligible", "ineligible"]);
harness.mark(&["eligible", "ineligible"]);
harness.eligibility.set_parallel_ineligible([(
"ineligible".to_string(),
ParallelEligibility::UncommittedProposalFiles,
)]);
if arrange_ineligible_as_unstartable {
harness.to_merge_wait("ineligible").await;
}
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("one ineligible marked target refuses the whole start");
match &error {
RunControlError::NoEligibleTarget { detail, .. } => assert!(
detail.contains("ineligible") && !detail.contains("eligible,"),
"the refusal must name the ineligible target: {detail}"
),
other => panic!("unexpected error: {other:?}"),
}
assert!(
harness.scheduler.calls().is_empty(),
"no scheduler may be spawned for a partially eligible target set"
);
assert_eq!(
harness.status("eligible").await,
"not queued",
"the eligible target must not be left queued by a refused start"
);
assert_eq!(
harness.marks.marked_ids(),
vec!["eligible".to_string(), "ineligible".to_string()],
"a refused start leaves marks coherent"
);
}
}
#[tokio::test]
async fn start_wakes_a_live_scheduler_instead_of_spawning_a_second_run() {
let harness = Harness::new(&["a"]);
harness.mark(&["a"]);
harness.scheduler.set_running(true);
let outcome = harness.service.start(OperatorMode::Stopped).await.unwrap();
assert!(matches!(
outcome,
RunControlOutcome::RunDispatched {
scheduler: SchedulerEffect::Notified,
..
}
));
assert_eq!(harness.scheduler.calls(), vec![SchedulerCall::Notified]);
}
#[tokio::test]
async fn start_is_refused_while_a_run_owns_the_lifecycle() {
for mode in [OperatorMode::Running, OperatorMode::Stopping] {
let harness = Harness::new(&["a"]);
harness.mark(&["a"]);
let error = harness.service.start(mode).await.expect_err("mode refuses");
match (mode, &error) {
(
OperatorMode::Stopping,
RunControlError::InvalidMode {
command: RunCommandKind::Start,
..
},
) => {}
(OperatorMode::Running, RunControlError::NoEligibleTarget { detail, .. }) => {
assert!(
detail.contains("a (not queued)"),
"the refusal must name the mark it could not route: {detail}"
);
}
_ => panic!("{mode:?}: unexpected refusal: {error:?}"),
}
assert!(harness.scheduler.calls().is_empty());
}
}
#[tokio::test]
async fn start_reports_a_runtime_launch_failure_instead_of_claiming_success() {
let harness = Harness::new(&["a"]);
harness.mark(&["a"]);
harness
.scheduler
.fail_launch("the scheduler refused this launch");
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("a refused launch is not a successful start");
assert!(matches!(
error,
RunControlError::DispatchFailed {
command: RunCommandKind::Start,
..
}
));
}
#[tokio::test]
async fn retry_routes_a_terminal_error_and_dispatches_the_scheduler() {
let harness = Harness::new(&["a"]);
harness.to_error("a").await;
let outcome = harness.service.retry_change("a").await.unwrap();
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["a".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Started,
excluded: Vec::new(),
}
);
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["a".to_string()]]
);
}
#[tokio::test]
async fn retry_resumes_a_resumable_external_hold() {
let harness = Harness::new(&["a"]);
harness
.apply(ExecutionEvent::ExecutionBlocked {
change_id: "a".to_string(),
blocker: external_blocker(),
})
.await;
assert_eq!(harness.status("a").await, "blocked");
let outcome = harness.service.retry_change("a").await.unwrap();
assert!(matches!(
outcome,
RunControlOutcome::RunDispatched {
explicit_retry: true,
scheduler: SchedulerEffect::Started,
..
}
));
}
#[tokio::test]
async fn retry_of_an_unsupported_target_changes_nothing() {
let harness = Harness::new(&["a"]);
let error = harness
.service
.retry_change("a")
.await
.expect_err("a not-queued row carries no retryable evidence");
assert!(matches!(error, RunControlError::Operator(_)));
assert!(
harness.scheduler.calls().is_empty(),
"an unsupported retry must not dispatch work"
);
}
#[tokio::test]
async fn bulk_retry_without_retryable_evidence_is_a_no_op() {
let harness = Harness::new(&["a", "b"]);
let outcome = harness
.service
.retry_errors(&["a".to_string(), "b".to_string()])
.await
.unwrap();
assert_eq!(
outcome,
RunControlOutcome::NoOp {
reason: RunNoOpReason::NoRetryableTarget
}
);
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn start_in_error_mode_retries_the_marked_error_rows() {
let harness = Harness::new(&["a", "b"]);
harness.to_error("a").await;
harness.mark(&["a", "b"]);
let outcome = harness.service.start(OperatorMode::Error).await.unwrap();
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["a".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Started,
excluded: vec![ExcludedTarget::new("b", "not queued")],
},
"only the row that carries retryable evidence is retried"
);
}
#[tokio::test]
async fn stop_and_cancel_stop_enforce_the_mode_matrix() {
let harness = Harness::new(&["a"]);
assert_eq!(
harness.service.stop(OperatorMode::Running).await.unwrap(),
RunControlOutcome::StopRequested
);
assert_eq!(
harness.scheduler.calls(),
vec![SchedulerCall::GracefulStop(true), SchedulerCall::Notified],
"the request is recorded, then the idle waiter is woken so a parked \
scheduler can reach its stop boundary without an unrelated event"
);
for refused in [
OperatorMode::Select,
OperatorMode::Stopped,
OperatorMode::Stopping,
OperatorMode::Error,
] {
let error = harness.service.stop(refused).await.expect_err("refused");
assert!(matches!(
error,
RunControlError::InvalidMode {
command: RunCommandKind::Stop,
..
}
));
}
assert_eq!(
harness
.service
.cancel_stop(OperatorMode::Stopping)
.await
.unwrap(),
RunControlOutcome::StopCancelled
);
let error = harness
.service
.cancel_stop(OperatorMode::Running)
.await
.expect_err("cancel stop needs a pending stop");
assert!(matches!(
error,
RunControlError::InvalidMode {
command: RunCommandKind::CancelStop,
..
}
));
}
#[tokio::test]
async fn force_stop_reports_classification_truthfully_and_always_cancels() {
use crate::tui::stop_classification::{ExecutionEvidence, ProcessReport, ShutdownWorkEvidence};
let harness = Harness::new(&["a"]);
harness.scheduler.set_running(true);
harness.scheduler.set_activity(StopActivitySnapshot {
execution_handles: ExecutionEvidence::Known { registered: 1 },
reducer_agent_execution_active: false,
shutdown_work: ShutdownWorkEvidence::Known { pending: false },
});
let outcome = harness
.service
.force_stop(OperatorMode::Running)
.await
.unwrap();
match outcome {
RunControlOutcome::ForceStopped {
classification,
awaiting_safe_boundary,
} => {
assert_eq!(classification.process_report, ProcessReport::ForceStopped);
assert!(awaiting_safe_boundary);
}
other => panic!("unexpected outcome: {other:?}"),
}
assert!(harness
.scheduler
.calls()
.contains(&SchedulerCall::Cancelled));
let idle = Harness::new(&["a"]);
let outcome = idle
.service
.force_stop(OperatorMode::Stopping)
.await
.unwrap();
match outcome {
RunControlOutcome::ForceStopped {
classification,
awaiting_safe_boundary,
} => {
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert!(!awaiting_safe_boundary);
}
other => panic!("unexpected outcome: {other:?}"),
}
assert!(idle.scheduler.calls().contains(&SchedulerCall::Cancelled));
}
#[tokio::test]
async fn force_stop_is_refused_outside_running_and_stopping() {
let harness = Harness::new(&["a"]);
for refused in [
OperatorMode::Select,
OperatorMode::Stopped,
OperatorMode::Error,
] {
let error = harness
.service
.force_stop(refused)
.await
.expect_err("refused");
assert!(matches!(
error,
RunControlError::InvalidMode {
command: RunCommandKind::ForceStop,
..
}
));
}
assert!(
harness.scheduler.calls().is_empty(),
"a refused force stop must not cancel the run"
);
}
#[tokio::test]
async fn resolve_reserves_one_active_resolver_and_queues_the_rest_in_fifo_order() {
let harness = Harness::new(&["a", "b", "c"]);
for id in ["a", "b", "c"] {
harness.to_merge_wait(id).await;
}
let first = harness.service.resolve_merge("a").await.unwrap();
assert_eq!(
first,
RunControlOutcome::ResolveReserved {
change_id: "a".to_string(),
reservation: ResolveReservation::Active,
scheduler: SchedulerEffect::Started,
}
);
let second = harness.service.resolve_merge("b").await.unwrap();
let third = harness.service.resolve_merge("c").await.unwrap();
assert_eq!(
second,
RunControlOutcome::ResolveReserved {
change_id: "b".to_string(),
reservation: ResolveReservation::Queued { position: 1 },
scheduler: SchedulerEffect::None,
}
);
assert_eq!(
third,
RunControlOutcome::ResolveReserved {
change_id: "c".to_string(),
reservation: ResolveReservation::Queued { position: 2 },
scheduler: SchedulerEffect::None,
}
);
assert_eq!(
harness.resolves.waiting(),
vec!["b".to_string(), "c".to_string()]
);
assert_eq!(
harness.scheduler.started_targets().len(),
1,
"only the active resolver dispatches scheduler work"
);
}
#[tokio::test]
async fn duplicate_resolve_submission_does_not_create_a_second_queue_entry() {
let harness = Harness::new(&["a", "b"]);
harness.to_merge_wait("a").await;
harness.to_merge_wait("b").await;
harness.service.resolve_merge("a").await.unwrap();
harness.service.resolve_merge("b").await.unwrap();
let duplicate = harness.service.resolve_merge("b").await.unwrap();
assert_eq!(
duplicate,
RunControlOutcome::NoOp {
reason: RunNoOpReason::ResolveAlreadyReserved {
change_id: "b".to_string()
}
}
);
assert_eq!(harness.resolves.waiting(), vec!["b".to_string()]);
}
#[tokio::test]
async fn resolve_from_startup_reconstructed_merge_wait_reserves_and_dispatches() {
let harness = Harness::new(&["a"]);
harness.to_startup_merge_wait("a").await;
assert_eq!(
harness.status("a").await,
"merge wait",
"workspace evidence alone must make the reducer agree with the row"
);
let outcome = harness.service.resolve_merge("a").await.unwrap();
assert_eq!(
outcome,
RunControlOutcome::ResolveReserved {
change_id: "a".to_string(),
reservation: ResolveReservation::Active,
scheduler: SchedulerEffect::Started,
}
);
assert_eq!(harness.status("a").await, "resolve pending");
assert!(harness.resolves.is_active());
assert_eq!(
harness.state.read().await.resolve_wait_change_ids(),
vec!["a".to_string()],
"the accepted intent must be scheduler-consumable retry membership"
);
assert_eq!(
harness.scheduler.started_targets(),
vec![Vec::<String>::new()]
);
}
#[tokio::test]
async fn resolve_from_startup_merge_wait_notifies_a_live_scheduler() {
let harness = Harness::new(&["a"]);
harness.to_startup_merge_wait("a").await;
harness.scheduler.set_running(true);
let outcome = harness.service.resolve_merge("a").await.unwrap();
assert_eq!(
outcome,
RunControlOutcome::ResolveReserved {
change_id: "a".to_string(),
reservation: ResolveReservation::Active,
scheduler: SchedulerEffect::Notified,
}
);
assert_eq!(harness.scheduler.calls(), vec![SchedulerCall::Notified]);
}
#[tokio::test]
async fn resolve_of_a_not_queued_target_without_workspace_evidence_is_refused() {
let harness = Harness::new(&["a"]);
assert_eq!(harness.status("a").await, "not queued");
let error = harness
.service
.resolve_merge("a")
.await
.expect_err("an idle, not-queued change is not waiting on a merge");
assert!(matches!(
error,
RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
..
}
));
assert_eq!(harness.status("a").await, "not queued");
assert!(!harness.resolves.is_active());
assert!(harness
.state
.read()
.await
.resolve_wait_change_ids()
.is_empty());
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn refresh_evidence_for_another_change_leaves_a_stale_target_ineligible() {
let harness = Harness::new(&["a", "b"]);
harness.to_startup_merge_wait("b").await;
let error = harness
.service
.resolve_merge("a")
.await
.expect_err("evidence for 'b' must not make 'a' resolve-eligible");
assert!(matches!(
error,
RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
..
}
));
assert!(!harness.resolves.is_reserved("a"));
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn resolve_of_a_stale_target_is_refused_without_a_reservation() {
let harness = Harness::new(&["a"]);
harness
.apply(ExecutionEvent::MergeCompleted {
change_id: "a".to_string(),
revision: "rev-a".to_string(),
})
.await;
let error = harness
.service
.resolve_merge("a")
.await
.expect_err("a merged change is not waiting on a merge");
assert!(matches!(
error,
RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
..
}
));
assert!(!harness.resolves.is_active());
assert!(harness.scheduler.calls().is_empty());
}
#[tokio::test]
async fn resolve_wakes_a_live_scheduler_instead_of_starting_a_second_one() {
let harness = Harness::new(&["a"]);
harness.to_merge_wait("a").await;
harness.scheduler.set_running(true);
let outcome = harness.service.resolve_merge("a").await.unwrap();
assert_eq!(
outcome,
RunControlOutcome::ResolveReserved {
change_id: "a".to_string(),
reservation: ResolveReservation::Active,
scheduler: SchedulerEffect::Notified,
}
);
assert_eq!(harness.scheduler.calls(), vec![SchedulerCall::Notified]);
}
#[tokio::test]
async fn finishing_the_active_resolve_promotes_the_next_waiting_change() {
let harness = Harness::new(&["a", "b"]);
harness.to_merge_wait("a").await;
harness.to_merge_wait("b").await;
harness.service.resolve_merge("a").await.unwrap();
harness.service.resolve_merge("b").await.unwrap();
assert_eq!(harness.resolves.finish_active(), Some("b".to_string()));
assert!(!harness.resolves.is_active());
assert!(
!harness.resolves.is_reserved("b"),
"a promoted change releases its reservation so it can reserve again"
);
}
#[test]
fn cancelling_a_queued_reservation_preserves_fifo_order() {
let ledger = ResolveReservations::new();
assert_eq!(ledger.reserve("a"), Some(ResolveReservation::Active));
assert_eq!(
ledger.reserve("b"),
Some(ResolveReservation::Queued { position: 1 })
);
assert_eq!(
ledger.reserve("c"),
Some(ResolveReservation::Queued { position: 2 })
);
assert!(ledger.cancel("b"));
assert!(!ledger.cancel("b"), "cancelling twice reports no change");
assert_eq!(ledger.waiting(), vec!["c".to_string()]);
assert_eq!(ledger.finish_active(), Some("c".to_string()));
}
#[test]
fn marking_an_active_resolver_removes_it_from_the_waiting_queue() {
let ledger = ResolveReservations::new();
ledger.reserve("a");
ledger.reserve("b");
ledger.mark_active("b");
assert_eq!(ledger.active(), Some("b".to_string()));
assert!(ledger.waiting().is_empty());
}
#[test]
fn start_eligibility_always_rejects_an_ineligible_target() {
let eligibility = StartEligibility::new();
let targets = vec!["a".to_string(), "b".to_string()];
assert!(
eligibility.rejected(&targets).is_empty(),
"nothing is rejected before an ineligible observation exists"
);
eligibility.set_parallel_ineligible([(
"a".to_string(),
ParallelEligibility::UncommittedProposalFiles,
)]);
assert_eq!(eligibility.rejected(&targets), vec!["a".to_string()]);
}
#[tokio::test]
async fn start_targets_only_returns_marked_rows_the_reducer_calls_not_queued() {
let harness = Harness::new(&["a", "b", "c"]);
harness.mark(&["a", "b", "c"]);
harness.to_merge_wait("b").await;
harness
.service
.operator()
.add_to_queue("c")
.await
.expect("queueing a change is allowed");
assert_eq!(harness.status("b").await, "merge wait");
assert_eq!(harness.status("c").await, "queued");
assert_eq!(harness.service.start_targets().await, vec!["a".to_string()]);
}
#[tokio::test]
async fn tui_and_remote_start_produce_identical_explicit_target_eligibility() {
let tui = Harness::new(&["marked-a", "marked-b", "unmarked-residue"]);
let remote = Harness::new(&["marked-a", "marked-b", "unmarked-residue"]);
tui.mark(&["marked-a", "marked-b"]);
remote.mark(&["marked-a", "marked-b"]);
let tui_outcome = tui
.service
.start(OperatorMode::Select)
.await
.expect("TUI start");
let remote_outcome = remote
.service
.start(OperatorMode::Select)
.await
.expect("remote start");
assert_eq!(tui_outcome, remote_outcome);
assert_eq!(
tui.scheduler.started_targets(),
remote.scheduler.started_targets()
);
let mut tui_queued = tui.state.read().await.queued_change_ids();
let mut remote_queued = remote.state.read().await.queued_change_ids();
tui_queued.sort();
remote_queued.sort();
assert_eq!(
tui_queued,
vec!["marked-a".to_string(), "marked-b".to_string()]
);
assert_eq!(
tui_queued, remote_queued,
"equivalent accepted intent must produce identical scheduler eligibility"
);
for harness in [&tui, &remote] {
assert!(
!harness
.state
.read()
.await
.is_ordinary_queue_eligible("unmarked-residue"),
"an unmarked catalog or worktree entry must stay ineligible"
);
}
}
#[tokio::test]
async fn queue_removal_and_dequeue_revoke_eligibility_until_explicit_requeue() {
let harness = Harness::new(&["alpha"]);
harness.mark(&["alpha"]);
harness
.service
.start(OperatorMode::Select)
.await
.expect("start");
assert!(harness
.state
.read()
.await
.is_ordinary_queue_eligible("alpha"));
harness
.service
.operator()
.remove_from_queue("alpha")
.await
.expect("queue removal");
assert!(
!harness
.state
.read()
.await
.is_ordinary_queue_eligible("alpha"),
"queue removal revokes ordinary execution eligibility immediately"
);
harness
.service
.operator()
.add_to_queue("alpha")
.await
.expect("requeue");
harness
.service
.operator()
.stop_and_dequeue("alpha")
.await
.expect("stop and dequeue");
assert!(
!harness
.state
.read()
.await
.is_ordinary_queue_eligible("alpha"),
"a dequeued change must not be reacquired"
);
harness
.service
.operator()
.add_to_queue("alpha")
.await
.expect("explicit requeue");
assert!(
harness
.state
.read()
.await
.is_ordinary_queue_eligible("alpha"),
"explicit requeue is the ordinary way back"
);
}
impl Harness {
async fn to_iteration_limit(&self, change_id: &str, attempts: u32, max: u32) {
self.to_error(change_id).await;
self.state
.write()
.await
.record_apply_iteration_limit(change_id, attempts, max);
}
async fn retains_iteration_limit(&self, change_id: &str) -> bool {
self.state
.read()
.await
.apply_iteration_limit(change_id)
.is_some()
}
}
#[tokio::test]
async fn settled_apply_limit_retry_is_admitted_while_the_scheduler_task_is_live() {
let harness = Harness::new(&["limited"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.scheduler.set_running(true);
assert!(harness.retains_iteration_limit("limited").await);
let outcome = harness
.service
.retry_change("limited")
.await
.expect("a settled terminal error is retryable on its own evidence");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["limited".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
},
"the live boundary is woken rather than joined by a second one"
);
assert_ne!(
harness.status("limited").await,
"error",
"the terminal error was consumed by the explicit intent"
);
assert!(
!harness.retains_iteration_limit("limited").await,
"and the diagnostic it explained is consumed by the same intent, so \
the later invocation starts from a fresh ceiling"
);
}
#[tokio::test]
async fn settled_apply_limit_retry_after_task_exit_starts_a_later_boundary() {
let harness = Harness::new(&["limited"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.scheduler.set_running(false);
let outcome = harness
.service
.retry_change("limited")
.await
.expect("a closed boundary admits the same ordinary retry route");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["limited".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Started,
excluded: Vec::new(),
},
"a later boundary is started, never a wake-up of the exited scheduler"
);
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["limited".to_string()]]
);
assert!(
!harness
.scheduler
.calls()
.iter()
.any(|call| matches!(call, SchedulerCall::Notified)),
"the exited scheduler is never notified: {:?}",
harness.scheduler.calls()
);
}
#[tokio::test]
async fn settled_apply_limit_later_state_starts_with_a_fresh_budget() {
let harness = Harness::new(&["limited"]);
harness.to_iteration_limit("limited", 50, 50).await;
let later = OrchestratorState::new(vec!["limited".to_string()], 50);
assert!(
later.apply_iteration_limit("limited").is_none(),
"the ephemeral diagnostic never crosses into a later boundary"
);
assert_eq!(later.apply_count("limited"), 0, "the budget starts fresh");
assert_eq!(later.parallel_finish_report(), ("completed", 0));
}
#[tokio::test]
async fn settled_apply_limit_marking_alone_produces_no_scheduler_effect() {
let harness = Harness::new(&["limited"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.scheduler.set_running(true);
let before = harness.effects().await;
harness
.service
.operator()
.set_execution_mark("limited", true)
.await
.expect("a settled limited row still accepts next-run intent");
assert_eq!(
harness.effects().await,
before,
"a mark leaves no queue, scheduler, retry-edge, or reducer effect"
);
assert!(
harness.retains_iteration_limit("limited").await,
"and the diagnostic survives until explicit retry consumes it"
);
}
#[tokio::test]
async fn settled_apply_limit_bulk_retry_dispatches_every_admitted_target() {
let harness = Harness::new(&["limited", "ordinary"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.to_error("ordinary").await;
harness.scheduler.set_running(true);
harness.mark(&["limited", "ordinary"]);
let outcome = harness
.service
.retry_errors(&["limited".to_string(), "ordinary".to_string()])
.await
.expect("both targets carry ordinary terminal-error evidence");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["limited".to_string(), "ordinary".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
},
"the settled limited target is accepted alongside the unrelated one"
);
assert_ne!(harness.status("limited").await, "error");
assert_ne!(harness.status("ordinary").await, "error");
}
#[tokio::test]
async fn settled_apply_limit_bulk_retry_keeps_unsupported_evidence_intact() {
let harness = Harness::new(&["limited", "waiting"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.to_merge_wait("waiting").await;
harness.scheduler.set_running(true);
let outcome = harness
.service
.retry_errors(&["limited".to_string(), "waiting".to_string()])
.await
.expect("one supported target keeps the bulk request useful");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["limited".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
},
"a merge-wait row carries no retryable evidence and is not claimed"
);
assert_eq!(harness.status("waiting").await, "merge wait");
}
#[tokio::test]
async fn settled_apply_limit_error_mode_start_retries_the_marked_row() {
let harness = Harness::new(&["limited"]);
harness.to_iteration_limit("limited", 50, 50).await;
harness.scheduler.set_running(true);
harness.mark(&["limited"]);
let outcome = harness
.service
.start(OperatorMode::Error)
.await
.expect("the marked settled error is a retry-class Start target");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["limited".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Notified,
excluded: Vec::new(),
}
);
assert_eq!(
harness.effects().await.explicit_retries,
vec![("limited".to_string(), RetryEdgeAuthority::TerminalError)],
"exactly one target-specific explicit-retry edge reaches the scheduler"
);
assert!(!harness.retains_iteration_limit("limited").await);
}
#[tokio::test]
async fn run_mark_intent_start_admission_reads_current_status_not_mark_time_status() {
let harness = Harness::new(&["a"]);
harness
.state
.write()
.await
.apply_command(ReducerCommand::AddToQueue("a".to_string()));
harness.mark(&["a"]);
let refused = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("an already-queued row is not a start target");
assert!(matches!(refused, RunControlError::NoEligibleTarget { .. }));
harness
.state
.write()
.await
.apply_command(ReducerCommand::RemoveFromQueue("a".to_string()));
assert_eq!(harness.status("a").await, "not queued");
let outcome = harness
.service
.start(OperatorMode::Select)
.await
.expect("the unchanged mark is admitted once the status allows it");
assert!(matches!(outcome, RunControlOutcome::RunDispatched { .. }));
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["a".to_string()]]
);
}
#[tokio::test]
async fn run_mark_intent_start_admission_excludes_status_without_blocking_runnable_targets() {
let harness = Harness::new(&["runnable", "waiting"]);
harness.to_merge_wait("waiting").await;
harness.mark(&["runnable", "waiting"]);
let outcome = harness
.service
.start(OperatorMode::Select)
.await
.expect("one non-startable mark does not refuse the runnable subset");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["runnable".to_string()],
explicit_retry: false,
scheduler: SchedulerEffect::Started,
excluded: vec![ExcludedTarget::new("waiting", "merge wait")],
},
"the admitted run must name the marked target it left out, and why"
);
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["runnable".to_string()]]
);
assert_eq!(
harness.marks.marked_ids(),
vec!["runnable".to_string(), "waiting".to_string()],
"admission consumes marks; it does not revoke them"
);
}
#[tokio::test]
async fn run_mark_intent_start_admission_rejects_when_no_runnable_target_remains() {
let harness = Harness::new(&["waiting", "queued"]);
harness.to_merge_wait("waiting").await;
harness
.state
.write()
.await
.apply_command(ReducerCommand::AddToQueue("queued".to_string()));
harness.mark(&["waiting", "queued"]);
let before = harness.effects().await;
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("no marked target is startable");
match &error {
RunControlError::NoEligibleTarget { detail, .. } => {
assert!(
detail.contains("waiting (merge wait)") && detail.contains("queued (queued)"),
"every exclusion must be named with its status: {detail}"
);
}
other => panic!("unexpected error: {other:?}"),
}
assert_eq!(
harness.effects().await,
before,
"a rejected admission leaves no queue, scheduler, or reducer effect"
);
}
#[tokio::test]
async fn run_mark_intent_start_admission_worktree_fence_rejects_the_whole_request() {
let harness = Harness::new(&["runnable", "ineligible"]);
harness.mark(&["runnable", "ineligible"]);
harness.eligibility.set_parallel_ineligible([(
"ineligible".to_string(),
ParallelEligibility::UncommittedProposalFiles,
)]);
let before = harness.effects().await;
let error = harness
.service
.start(OperatorMode::Select)
.await
.expect_err("one worktree-ineligible mark refuses the whole request");
match &error {
RunControlError::NoEligibleTarget { detail, .. } => assert!(
detail.contains("ineligible"),
"the refusal must name the fenced target: {detail}"
),
other => panic!("unexpected error: {other:?}"),
}
assert_eq!(
harness.effects().await,
before,
"the fence refuses before any partial effect exists"
);
}
#[tokio::test]
async fn run_mark_intent_start_admission_error_mode_routes_only_retryable_marks() {
let harness = Harness::new(&["failed", "idle"]);
harness.to_error("failed").await;
harness.mark(&["failed", "idle"]);
let outcome = harness
.service
.start(OperatorMode::Error)
.await
.expect("the retryable mark is routed");
assert_eq!(
outcome,
RunControlOutcome::RunDispatched {
change_ids: vec!["failed".to_string()],
explicit_retry: true,
scheduler: SchedulerEffect::Started,
excluded: vec![ExcludedTarget::new("idle", "not queued")],
},
"a marked row without retryable evidence is reported, not retried"
);
assert_eq!(
harness.scheduler.started_targets(),
vec![vec!["failed".to_string()]]
);
}
#[tokio::test]
async fn run_mark_intent_start_admission_error_mode_rejects_without_retryable_marks() {
let harness = Harness::new(&["idle"]);
harness.mark(&["idle"]);
let before = harness.effects().await;
let error = harness
.service
.start(OperatorMode::Error)
.await
.expect_err("no marked row carries retryable evidence");
match &error {
RunControlError::NoEligibleTarget { detail, .. } => assert!(
detail.contains("idle (not queued)"),
"the refusal must name the exclusion: {detail}"
),
other => panic!("unexpected error: {other:?}"),
}
assert_eq!(harness.effects().await, before);
}
#[tokio::test]
async fn run_mark_intent_start_admission_unmark_does_not_cancel_or_dequeue_admitted_work() {
let harness = Harness::new(&["a", "b"]);
harness.mark(&["a", "b"]);
harness
.service
.start(OperatorMode::Select)
.await
.expect("both marked targets are startable");
let admitted = harness.effects().await;
let changed = harness
.operator
.set_execution_mark("a", false)
.await
.expect("a non-terminal row accepts an unmark at any time");
assert!(matches!(
changed,
crate::orchestration::operator_command::OperatorOutcome::MarkSet { marked: false, .. }
));
assert_eq!(harness.marks.marked_ids(), vec!["b".to_string()]);
assert_eq!(
harness.effects().await,
admitted,
"unmarking must not dequeue, cancel, unschedule, or restatus admitted work"
);
}
#[tokio::test]
async fn bulk_retry_reports_the_unchanged_acceptance_input_refusal_per_target() {
use crate::orchestration::acceptance::execution_manifest::UNCHANGED_ACCEPTANCE_INPUT;
let harness = Harness::with_refused_acceptance(&["alpha", "beta"], &["alpha"]);
harness.to_error("alpha").await;
harness.to_error("beta").await;
let prepared = harness
.service
.prepare_retry_errors(&["alpha".to_string(), "beta".to_string()])
.await
.expect("one refused target must not refuse the whole request");
let PreparedIntent::Retry { routes, excluded } = &prepared.intent else {
panic!(
"a bulk retry must produce retry intent, got {:?}",
prepared.intent
);
};
assert_eq!(
routes.iter().map(|(id, _)| id.as_str()).collect::<Vec<_>>(),
vec!["beta"],
"only the admitted sibling may be routed"
);
assert_eq!(excluded.len(), 1, "{excluded:?}");
assert_eq!(excluded[0].change_id, "alpha");
assert_eq!(
excluded[0].status, UNCHANGED_ACCEPTANCE_INPUT,
"the exclusion must carry the stable outcome token, not a lifecycle status"
);
assert!(
excluded[0]
.detail
.as_deref()
.is_some_and(|detail| detail.contains("restores retry eligibility")),
"the exclusion must say what to change: {:?}",
excluded[0].detail
);
}
#[tokio::test]
async fn start_retry_admission_reports_the_unchanged_acceptance_input_refusal() {
use crate::orchestration::acceptance::execution_manifest::UNCHANGED_ACCEPTANCE_INPUT;
let harness = Harness::with_refused_acceptance(&["alpha", "beta"], &["alpha"]);
harness.to_error("alpha").await;
harness.to_error("beta").await;
harness.mark(&["alpha", "beta"]);
let outcome = harness
.service
.start(OperatorMode::Select)
.await
.expect("a retry-eligible sibling keeps Start runnable");
let RunControlOutcome::RunDispatched {
change_ids,
explicit_retry,
excluded,
..
} = outcome
else {
panic!("Start must dispatch the admitted retry, got {outcome:?}");
};
assert_eq!(change_ids, vec!["beta".to_string()]);
assert!(explicit_retry);
assert_eq!(excluded.len(), 1, "{excluded:?}");
assert_eq!(excluded[0].change_id, "alpha");
assert_eq!(excluded[0].status, UNCHANGED_ACCEPTANCE_INPUT);
assert!(
excluded[0].describe().contains(UNCHANGED_ACCEPTANCE_INPUT),
"the one operator-facing spelling must name the refusal: {}",
excluded[0].describe()
);
assert_eq!(
harness.status("alpha").await,
"error",
"a refused target keeps its terminal evidence"
);
}
mod change_error_f5_retry;
mod stopped_marked_resume;