use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use crate::events::{dispatch_event_with_marks, EventSink, ExecutionEvent};
use crate::openspec::{Change, ProposalMetadata};
use crate::orchestration::mark_reconciliation::ExecutionMarkReconciler;
use crate::orchestration::operator_command::{ExecutionMarkStore, ParallelRuntime};
use crate::orchestration::state::OrchestratorState;
use crate::web::remote_control_api::dto::InstanceSnapshot;
use crate::web::remote_control_api::projection::EventsSince;
use crate::web::state::{WebEventSink, WebState};
struct Wired {
web_state: Arc<WebState>,
reducer: Arc<tokio::sync::RwLock<OrchestratorState>>,
marks: Arc<ExecutionMarkStore>,
reconciler: ExecutionMarkReconciler,
sinks: Vec<Arc<dyn EventSink>>,
}
impl Wired {
async fn new(change_ids: &[&str]) -> Self {
let reducer = Arc::new(tokio::sync::RwLock::new(OrchestratorState::new(
change_ids.iter().map(|id| id.to_string()).collect(),
10,
)));
let marks = Arc::new(ExecutionMarkStore::new());
let parallel = Arc::new(ParallelRuntime::new());
let web_state = Arc::new(WebState::new(
&change_ids.iter().map(|id| change(id)).collect::<Vec<_>>(),
));
web_state.set_shared_state(reducer.clone()).await;
web_state.set_execution_marks(marks.clone()).await;
web_state.set_parallel_runtime(parallel.clone()).await;
Self {
sinks: vec![Arc::new(WebEventSink::new(web_state.clone()))],
reconciler: ExecutionMarkReconciler::new(marks.clone(), parallel),
web_state,
reducer,
marks,
}
}
async fn dispatch(&self, event: ExecutionEvent) {
dispatch_event_with_marks(&self.reducer, &self.sinks, event, Some(&self.reconciler)).await;
}
fn published(&self) -> (InstanceSnapshot, u64) {
let (snapshot, revision, _) = self.web_state.remote_control().projection().snapshot();
(snapshot, revision)
}
fn envelope_revision(&self, event_type: &str) -> u64 {
let EventsSince::Replay(events) =
self.web_state.remote_control().projection().events_after(0)
else {
panic!("the replay window must still hold every event of this test");
};
let envelope = events
.iter()
.rev()
.find(|event| event.event_type == event_type)
.unwrap_or_else(|| panic!("no `{event_type}` event reached the v2 stream"));
envelope.state_revision
}
}
fn change(id: &str) -> Change {
Change {
id: id.to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}
}
fn marked(snapshot: &InstanceSnapshot, change_id: &str) -> bool {
snapshot
.changes
.iter()
.find(|change| change.id == change_id)
.unwrap_or_else(|| panic!("change '{change_id}' must be projected"))
.execution_marked
}
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: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
}
}
#[tokio::test]
async fn execution_event_clears_authoritative_mark_before_projection() {
let wired = Wired::new(&["alpha", "beta"]).await;
wired.marks.set("alpha", true);
wired.marks.set("beta", true);
wired
.dispatch(ExecutionEvent::ApplyStarted {
change_id: "alpha".to_string(),
command: "apply".to_string(),
})
.await;
let (before, before_revision) = wired.published();
assert!(marked(&before, "alpha") && marked(&before, "beta"));
wired
.dispatch(ExecutionEvent::ApplyFailed {
change_id: "alpha".to_string(),
error: "boom".to_string(),
})
.await;
let (after, after_revision) = wired.published();
assert!(
after_revision > before_revision,
"the failure must advance the published revision"
);
assert_eq!(
wired.envelope_revision("apply_failed"),
after_revision,
"the failure event must name the revision that already reports the cleared mark"
);
assert!(
!marked(&after, "alpha"),
"the failure revision still exposes the stale mark"
);
assert!(
marked(&after, "beta"),
"an unrelated target lost its mark in the same snapshot"
);
assert_eq!(
after
.changes
.iter()
.find(|change| change.id == "alpha")
.map(|change| change.display_status.as_str()),
Some("error"),
"the same revision must carry the reducer transition"
);
wired
.dispatch(ExecutionEvent::ApplyFailed {
change_id: "alpha".to_string(),
error: "boom".to_string(),
})
.await;
let (_, duplicate_revision) = wired.published();
assert_eq!(
duplicate_revision, after_revision,
"a duplicate revocation must be revision-idempotent"
);
assert_eq!(wired.marks.marked_ids(), vec!["beta".to_string()]);
}
#[tokio::test]
async fn refresh_publishes_rejected_and_ineligible_revocations_in_one_revision() {
let wired = Wired::new(&["alpha", "beta", "gamma"]).await;
for id in ["alpha", "beta", "gamma"] {
wired.marks.set(id, true);
}
wired
.dispatch(refresh(
&["beta", "gamma"],
&["alpha"],
&["beta", "gamma"],
&["gamma"],
))
.await;
let (snapshot, revision) = wired.published();
assert_eq!(
wired.envelope_revision("changes_refreshed"),
revision,
"the refresh event must name the revision its own decision produced"
);
assert!(
!marked(&snapshot, "gamma"),
"a parallel-ineligible target stayed marked"
);
assert!(
marked(&snapshot, "beta"),
"an eligible target lost its mark to another row's cleanup"
);
assert_eq!(
wired.marks.marked_ids(),
vec!["beta".to_string()],
"the rejected marker row kept its shared mark"
);
wired
.dispatch(refresh(
&["beta", "gamma"],
&["alpha"],
&["beta", "gamma"],
&["gamma"],
))
.await;
let (_, repeated) = wired.published();
assert_eq!(
repeated, revision,
"a repeated refresh cleanup must be idempotent"
);
assert_eq!(wired.marks.marked_ids(), vec!["beta".to_string()]);
}
#[tokio::test]
async fn on_merged_recovery_and_cleared_mark_are_coherent() {
let wired = Wired::new(&["alpha"]).await;
wired.marks.set("alpha", true);
let hook_failure = ExecutionEvent::HookFailed {
change_id: "alpha".to_string(),
hook_type: crate::hooks::HookType::OnMerged.config_key().to_string(),
error: "publish script exited 1".to_string(),
};
wired.dispatch(hook_failure.clone()).await;
let (snapshot, revision) = wired.published();
assert_eq!(wired.envelope_revision("hook_failed"), revision);
assert!(!marked(&snapshot, "alpha"));
assert_eq!(
snapshot
.changes
.iter()
.find(|change| change.id == "alpha")
.map(|change| change.display_status.as_str()),
Some("merge wait"),
"the recovery row and the cleared mark must land together"
);
wired.marks.set("alpha", true);
wired.dispatch(hook_failure).await;
let (after, _) = wired.published();
assert!(
marked(&after, "alpha"),
"a duplicate hook failure discarded a fresh re-mark"
);
}
#[tokio::test]
async fn run_mark_intent_archive_revision_preserves_target_and_unrelated_marks() {
let wired = Wired::new(&["alpha", "beta"]).await;
wired.marks.set("alpha", true);
wired.marks.set("beta", true);
wired
.dispatch(ExecutionEvent::ApplyStarted {
change_id: "alpha".to_string(),
command: "apply".to_string(),
})
.await;
let (before, before_revision) = wired.published();
assert!(marked(&before, "alpha") && marked(&before, "beta"));
wired
.dispatch(ExecutionEvent::ChangeArchived("alpha".to_string()))
.await;
let (after, after_revision) = wired.published();
assert!(
after_revision > before_revision,
"the archive must advance the published revision"
);
assert_eq!(
wired.envelope_revision("change_archived"),
after_revision,
"the archive event must name the revision that reports the preserved marks"
);
assert!(
marked(&after, "alpha"),
"the archive revision revoked the mark on its own target"
);
assert!(
marked(&after, "beta"),
"an unrelated target lost its mark to another row's archive"
);
assert_eq!(
after
.changes
.iter()
.find(|change| change.id == "alpha")
.and_then(|change| change.latest_activity.as_ref())
.map(|activity| activity.event_type.as_str()),
Some("change_archived"),
"the same revision must carry the reducer transition"
);
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(),
},
] {
wired.dispatch(event.clone()).await;
let (snapshot, _) = wired.published();
assert!(
marked(&snapshot, "alpha"),
"{event:?} revoked the mark the archive preserved"
);
assert!(
marked(&snapshot, "beta"),
"{event:?} disturbed an unrelated mark"
);
}
assert_eq!(
wired.marks.marked_ids(),
vec!["alpha".to_string(), "beta".to_string()]
);
}
#[tokio::test]
async fn process_stop_revision_retains_marked_resume_targets() {
let wired = Wired::new(&["alpha", "beta"]).await;
wired.marks.set("alpha", true);
wired.marks.set("beta", true);
wired
.dispatch(ExecutionEvent::ApplyStarted {
change_id: "alpha".to_string(),
command: "apply".to_string(),
})
.await;
wired.dispatch(ExecutionEvent::Stopped).await;
let (snapshot, _) = wired.published();
assert!(
marked(&snapshot, "alpha") && marked(&snapshot, "beta"),
"a process stop must publish the same marked resume target set"
);
}