use std::collections::HashSet;
use std::sync::Arc;
use super::*;
use crate::events::{ExecutionEvent, RejectionOutcome, StalledBlocker};
use crate::openspec::{Change, ProposalMetadata};
use crate::orchestration::state::OrchestratorState;
fn reconciler(marked: &[&str]) -> (ExecutionMarkReconciler, Arc<ExecutionMarkStore>) {
let marks = Arc::new(ExecutionMarkStore::new());
for id in marked {
marks.set(id, true);
}
let reconciler = ExecutionMarkReconciler::new(marks.clone(), Arc::new(ParallelRuntime::new()));
(reconciler, marks)
}
fn state(change_ids: &[&str]) -> OrchestratorState {
OrchestratorState::new(change_ids.iter().map(|id| id.to_string()).collect(), 10)
}
fn apply(
state: &mut OrchestratorState,
reconciler: &ExecutionMarkReconciler,
event: &ExecutionEvent,
) -> Vec<String> {
let pre = capture_pre_state(event, state);
state.apply_execution_event(event);
match pre {
Some(pre) => reconciler.reconcile(event, &pre, state),
None => Vec::new(),
}
}
fn change(id: &str) -> Change {
Change {
id: id.to_string(),
completed_tasks: 1,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}
}
fn refresh(
active: &[&str],
rejected: &[&str],
committed: &[&str],
dirty: &[&str],
) -> ExecutionEvent {
ExecutionEvent::ChangesRefreshed {
changes: active.iter().map(|id| change(id)).collect(),
rejected_changes: rejected.iter().map(|id| change(id)).collect(),
committed_change_ids: committed.iter().map(|id| id.to_string()).collect(),
uncommitted_file_change_ids: dirty.iter().map(|id| id.to_string()).collect(),
worktree_change_ids: HashSet::new(),
worktree_paths: Default::default(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
}
}
fn on_merged_failure(change_id: &str) -> ExecutionEvent {
ExecutionEvent::HookFailed {
change_id: change_id.to_string(),
hook_type: crate::hooks::HookType::OnMerged.config_key().to_string(),
error: "publish script exited 1".to_string(),
}
}
fn stalled_blocker() -> StalledBlocker {
StalledBlocker {
category: "acceptance_finding".to_string(),
phase: "acceptance".to_string(),
gate: "acceptance".to_string(),
error_summary: "unresolved finding".to_string(),
evidence: vec!["tests/acceptance.rs:1".to_string()],
unblock_condition: None,
prerequisite_owner: None,
next_action: "resolve and retry".to_string(),
resumable: true,
worktree_preserved: true,
}
}
#[test]
fn event_mark_reconciliation_covers_failure_and_rejection_edges() {
struct Case {
name: &'static str,
arrange: Vec<ExecutionEvent>,
marked: Vec<&'static str>,
event: ExecutionEvent,
expected: Vec<&'static str>,
}
let alpha = "alpha";
let beta = "beta";
let failure_variants: Vec<(&'static str, ExecutionEvent)> = vec![
(
"processing failure",
ExecutionEvent::ProcessingError {
id: alpha.to_string(),
error: "boom".to_string(),
},
),
(
"apply failure",
ExecutionEvent::ApplyFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
},
),
(
"acceptance failure",
ExecutionEvent::AcceptanceFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
},
),
(
"archive failure",
ExecutionEvent::ArchiveFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
reason: None,
summary: None,
},
),
(
"push failure",
ExecutionEvent::PushFailed {
change_id: alpha.to_string(),
remote: "origin".to_string(),
branch: "cflx/alpha".to_string(),
error: "boom".to_string(),
},
),
(
"rejection-review failure",
ExecutionEvent::RejectionReviewFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
},
),
];
let mut cases: Vec<Case> = failure_variants
.iter()
.map(|(name, event)| Case {
name,
arrange: Vec::new(),
marked: vec![alpha, beta],
event: event.clone(),
expected: vec![beta],
})
.collect();
cases.extend(failure_variants.iter().map(|(name, event)| Case {
name,
arrange: vec![ExecutionEvent::MergeCompleted {
change_id: alpha.to_string(),
revision: "abc".to_string(),
}],
marked: vec![alpha, beta],
event: event.clone(),
expected: vec![alpha, beta],
}));
cases.extend([
Case {
name: "terminal rejection",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: ExecutionEvent::ChangeRejected {
change_id: alpha.to_string(),
reason: "blocker".to_string(),
},
expected: vec![beta],
},
Case {
name: "rejected marker row introduced by refresh",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: refresh(&[beta], &[alpha], &[alpha, beta], &[]),
expected: vec![beta],
},
Case {
name: "successful dequeue",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: ExecutionEvent::ChangeDequeued {
change_id: alpha.to_string(),
},
expected: vec![beta],
},
Case {
name: "legacy target-scoped stop",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: ExecutionEvent::ChangeStopped {
change_id: alpha.to_string(),
},
expected: vec![beta],
},
Case {
name: "duplicate dequeue after re-mark",
arrange: vec![ExecutionEvent::ChangeDequeued {
change_id: alpha.to_string(),
}],
marked: vec![alpha, beta],
event: ExecutionEvent::ChangeDequeued {
change_id: alpha.to_string(),
},
expected: vec![alpha, beta],
},
Case {
name: "first on_merged hook failure enters merge-wait recovery",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: on_merged_failure(alpha),
expected: vec![beta],
},
Case {
name: "replayed on_merged hook failure preserves a fresh re-mark",
arrange: vec![on_merged_failure(alpha)],
marked: vec![alpha, beta],
event: on_merged_failure(alpha),
expected: vec![alpha, beta],
},
Case {
name: "a non-merge hook failure is not a mark edge",
arrange: Vec::new(),
marked: vec![alpha, beta],
event: ExecutionEvent::HookFailed {
change_id: alpha.to_string(),
hook_type: "post_apply".to_string(),
error: "boom".to_string(),
},
expected: vec![alpha, beta],
},
Case {
name: "duplicate failure after re-mark keeps the fresh intent",
arrange: vec![ExecutionEvent::ApplyFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
}],
marked: vec![alpha, beta],
event: ExecutionEvent::ApplyFailed {
change_id: alpha.to_string(),
error: "boom".to_string(),
},
expected: vec![alpha, beta],
},
Case {
name: "dequeue cannot revoke a rejected row's mark a second time",
arrange: vec![ExecutionEvent::ChangeRejected {
change_id: alpha.to_string(),
reason: "blocker".to_string(),
}],
marked: vec![alpha, beta],
event: ExecutionEvent::ChangeDequeued {
change_id: alpha.to_string(),
},
expected: vec![alpha, beta],
},
]);
for case in cases {
let mut reducer = state(&[alpha, beta]);
let (reconciler, marks) = reconciler(&[]);
for event in &case.arrange {
apply(&mut reducer, &reconciler, event);
}
for id in &case.marked {
marks.set(id, true);
}
apply(&mut reducer, &reconciler, &case.event);
let expected: Vec<String> = case.expected.iter().map(|id| id.to_string()).collect();
assert_eq!(
marks.marked_ids(),
expected,
"{}: unexpected mark set after reconciliation",
case.name
);
}
}
#[test]
fn event_mark_reconciliation_is_idempotent_for_repeated_revocation() {
let mut reducer = state(&["alpha"]);
let (reconciler, marks) = reconciler(&["alpha"]);
let event = ExecutionEvent::ChangeRejected {
change_id: "alpha".to_string(),
reason: "blocker".to_string(),
};
assert_eq!(
apply(&mut reducer, &reconciler, &event),
vec!["alpha".to_string()]
);
assert!(apply(&mut reducer, &reconciler, &event).is_empty());
assert!(marks.marked_ids().is_empty());
}
#[test]
fn parallel_ineligible_refresh_revokes_target_mark() {
let mut reducer = state(&["alpha", "beta"]);
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let event = refresh(&["alpha", "beta"], &[], &["alpha", "beta"], &["alpha"]);
let revoked = apply(&mut reducer, &reconciler, &event);
assert_eq!(revoked, vec!["alpha".to_string()]);
assert_eq!(marks.marked_ids(), vec!["beta".to_string()]);
assert!(apply(&mut reducer, &reconciler, &event).is_empty());
assert_eq!(marks.marked_ids(), vec!["beta".to_string()]);
marks.set("beta", true);
let absent = refresh(&["alpha", "beta"], &[], &["alpha"], &[]);
apply(&mut reducer, &reconciler, &absent);
assert!(
marks.marked_ids().is_empty(),
"a target absent from HEAD keeps an invalid mark"
);
}
#[test]
fn eligible_refresh_preserves_every_mark() {
let mut reducer = state(&["alpha", "beta"]);
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let event = refresh(&["alpha", "beta"], &[], &["alpha", "beta"], &[]);
assert!(apply(&mut reducer, &reconciler, &event).is_empty());
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()]
);
}
#[test]
fn event_mark_reconciliation_preserves_unrelated_and_stopped_marks() {
let preserved: Vec<(&str, ExecutionEvent)> = vec![
(
"dependency block",
ExecutionEvent::DependencyBlocked {
change_id: "alpha".to_string(),
dependency_ids: vec!["beta".to_string()],
},
),
(
"acceptance-gated stall",
ExecutionEvent::AcceptanceGated {
change_id: "alpha".to_string(),
blocker: stalled_blocker(),
},
),
(
"execution hold",
ExecutionEvent::ExecutionBlocked {
change_id: "alpha".to_string(),
blocker: stalled_blocker(),
},
),
(
"skipped after a failed dependency",
ExecutionEvent::ChangeSkipped {
change_id: "alpha".to_string(),
reason: "dependency failed".to_string(),
},
),
(
"manual merge-wait deferral",
ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
},
),
(
"auto-resumable resolve wait",
ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "lane busy".to_string(),
auto_resumable: true,
},
),
(
"resolve failure returns to merge wait",
ExecutionEvent::ResolveFailed {
change_id: "alpha".to_string(),
error: "conflict".to_string(),
},
),
(
"archive success",
ExecutionEvent::ChangeArchived("alpha".to_string()),
),
(
"merge success",
ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "abc".to_string(),
},
),
(
"push success",
ExecutionEvent::PushCompleted {
change_id: "alpha".to_string(),
remote: "origin".to_string(),
branch: "cflx/alpha".to_string(),
},
),
(
"resumed rejection review",
ExecutionEvent::RejectionReviewCompleted {
change_id: "alpha".to_string(),
outcome: RejectionOutcome::Resume,
},
),
("run completion", ExecutionEvent::AllCompleted),
(
"global fatal error without a target",
ExecutionEvent::Error {
message: "scheduler died".to_string(),
},
),
("process-level stop", ExecutionEvent::Stopped),
];
for (name, event) in preserved {
let mut reducer = state(&["alpha", "beta"]);
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let revoked = apply(&mut reducer, &reconciler, &event);
assert!(revoked.is_empty(), "{name} revoked a mark it must preserve");
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()],
"{name} did not preserve the complete mark set"
);
}
}
#[test]
fn process_stop_retains_marked_resume_targets() {
let mut reducer = state(&["alpha", "beta"]);
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
apply(
&mut reducer,
&reconciler,
&ExecutionEvent::ApplyStarted {
change_id: "alpha".to_string(),
command: "apply".to_string(),
},
);
apply(&mut reducer, &reconciler, &ExecutionEvent::Stopped);
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()],
"a process stop must leave every resume target marked"
);
assert_eq!(reducer.display_status("alpha"), "not queued");
}
#[test]
fn queue_intent_never_creates_an_execution_mark() {
use crate::orchestration::state::ReducerCommand;
let mut reducer = state(&["alpha"]);
let (reconciler, marks) = reconciler(&[]);
reducer.apply_command(ReducerCommand::AddToQueue("alpha".to_string()));
assert_eq!(reducer.display_status("alpha"), "queued");
let event = refresh(&["alpha"], &[], &["alpha"], &[]);
apply(&mut reducer, &reconciler, &event);
assert!(
marks.marked_ids().is_empty(),
"queue intent must not become an execution mark"
);
}
#[test]
fn run_mark_intent_archive_preserves_target_and_unrelated_marks() {
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let mut state = state(&["alpha", "beta"]);
let revoked = apply(
&mut state,
&reconciler,
&ExecutionEvent::ChangeArchived("alpha".to_string()),
);
assert!(
revoked.is_empty(),
"archive must revoke nothing, but reported {revoked:?}"
);
assert!(
state.is_archived("alpha"),
"the reducer did record the archive"
);
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()],
"the archived target and its unrelated neighbour both stay marked"
);
}
#[test]
fn run_mark_intent_archive_repeated_delivery_preserves_marks() {
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let mut state = state(&["alpha", "beta"]);
let archived = ExecutionEvent::ChangeArchived("alpha".to_string());
assert!(apply(&mut state, &reconciler, &archived).is_empty());
assert!(apply(&mut state, &reconciler, &archived).is_empty());
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()],
"a duplicate archive delivery must not become a revocation either"
);
}
#[test]
fn run_mark_intent_archive_post_transitions_preserve_the_marks() {
let (reconciler, marks) = reconciler(&["alpha", "beta"]);
let mut state = state(&["alpha", "beta"]);
apply(
&mut state,
&reconciler,
&ExecutionEvent::ChangeArchived("alpha".to_string()),
);
for event in [
ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "rev".to_string(),
},
ExecutionEvent::PushCompleted {
change_id: "alpha".to_string(),
remote: "origin".to_string(),
branch: "cflx/alpha".to_string(),
},
] {
assert!(
revoking_edge_target(&event).is_none(),
"{event:?} carries no revoking edge of its own"
);
assert!(apply(&mut state, &reconciler, &event).is_empty());
assert_eq!(
marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()],
"{event:?} disturbed a mark it must preserve"
);
}
}
#[test]
fn run_mark_intent_archive_carries_no_revoking_edge() {
assert!(
revoking_edge_target(&ExecutionEvent::ChangeArchived("alpha".to_string())).is_none(),
"archive must not classify as a revoking edge"
);
let genuine: Vec<(ExecutionEvent, RevokingEdge)> = vec![
(
ExecutionEvent::ProcessingError {
id: "alpha".to_string(),
error: "boom".to_string(),
},
RevokingEdge::ChangeLevelFailure,
),
(
ExecutionEvent::ArchiveFailed {
change_id: "alpha".to_string(),
error: "boom".to_string(),
reason: None,
summary: None,
},
RevokingEdge::ChangeLevelFailure,
),
(
ExecutionEvent::ChangeRejected {
change_id: "alpha".to_string(),
reason: "blocker".to_string(),
},
RevokingEdge::Rejection,
),
(
ExecutionEvent::ChangeDequeued {
change_id: "alpha".to_string(),
},
RevokingEdge::Dequeue,
),
(
ExecutionEvent::ChangeStopped {
change_id: "alpha".to_string(),
},
RevokingEdge::Dequeue,
),
(on_merged_failure("alpha"), RevokingEdge::MergedHookRecovery),
];
for (event, expected) in genuine {
assert_eq!(
revoking_edge_target(&event),
Some(("alpha", expected)),
"{event:?} lost its own revocation edge"
);
}
}