use std::sync::Arc;
use std::time::Duration;
use super::tests::AdapterHarness;
use super::*;
use crate::orchestration::mark_settlement::{MarkSettlementCoordinator, MARK_STABILITY_WINDOW};
use crate::tui::state::AppState;
use crate::tui::types::AppExecutionMode;
struct PendingPass {
passes: tokio::sync::watch::Receiver<u64>,
target: u64,
}
fn pending_pass(coordinator: &Arc<MarkSettlementCoordinator>) -> PendingPass {
let passes = coordinator.passes();
let target = *passes.borrow() + 1;
PendingPass { passes, target }
}
impl PendingPass {
async fn wait(mut self) {
self.passes
.wait_for(|observed| *observed >= self.target)
.await
.expect("the coordinator outlives its watchers");
}
}
async fn settle(coordinator: &Arc<MarkSettlementCoordinator>) {
let pending = pending_pass(coordinator);
tokio::time::advance(MARK_STABILITY_WINDOW + Duration::from_millis(1)).await;
pending.wait().await;
}
async fn expect_no_settlement(coordinator: &Arc<MarkSettlementCoordinator>) {
let before = *coordinator.passes().borrow();
tokio::time::advance(MARK_STABILITY_WINDOW * 2).await;
tokio::task::yield_now().await;
assert_eq!(
*coordinator.passes().borrow(),
before,
"no settlement pass may run"
);
}
async fn drain_marks(harness: &AdapterHarness, app: &mut AppState) {
let service = harness.application.run_control().operator();
crate::tui::runner::apply_pending_mark_writes(app, &service).await;
}
fn coordinator(harness: &AdapterHarness) -> Arc<MarkSettlementCoordinator> {
harness.marks.settlement()
}
fn running_harness(change_ids: &[&str]) -> AdapterHarness {
let harness = AdapterHarness::new(change_ids);
harness.scheduler.set_running(true);
harness
}
async fn refresh_catalog(harness: &AdapterHarness, change_ids: &[&str]) {
harness
.dispatcher
.dispatch(crate::events::ExecutionEvent::ChangesRefreshed {
changes: change_ids
.iter()
.map(|id| super::tests::create_test_change(id))
.collect(),
rejected_changes: Vec::new(),
committed_change_ids: change_ids.iter().map(|id| (*id).to_string()).collect(),
uncommitted_file_change_ids: std::collections::HashSet::new(),
worktree_change_ids: std::collections::HashSet::new(),
worktree_paths: std::collections::HashMap::new(),
worktree_not_ahead_ids: std::collections::HashSet::new(),
merge_wait_ids: std::collections::HashSet::new(),
})
.await;
}
async fn start_executing(harness: &AdapterHarness, change_id: &str) {
harness
.application
.apply(OperatorIntent::SetQueueIntent {
change_id: change_id.to_string(),
queued: true,
})
.await;
harness
.dispatcher
.dispatch(crate::events::ExecutionEvent::ApplyStarted {
change_id: change_id.to_string(),
command: "echo".to_string(),
})
.await;
assert_eq!(harness.status(change_id).await, "applying");
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_running_owner_admits_a_late_catalog_target() {
let harness = running_harness(&["active"]);
start_executing(&harness, "active").await;
let wakes_before = harness.queue_port.wakes();
assert!(
harness.state.read().await.change_runtime("alpha").is_none(),
"the reducer must not yet track the freshly created proposal"
);
harness
.application
.apply(OperatorIntent::SetExecutionMark {
change_id: "alpha".to_string(),
marked: true,
})
.await;
assert!(coordinator(&harness).is_armed());
settle(&coordinator(&harness)).await;
assert_eq!(
coordinator(&harness).pending_snapshot(),
Some(vec!["alpha".to_string()]),
"an unloadable target keeps its deadline instead of losing the mark"
);
assert_eq!(
harness.queue_port.wakes(),
wakes_before,
"nothing was applied"
);
refresh_catalog(&harness, &["active", "alpha"]).await;
assert!(
harness.marks.is_marked("alpha"),
"the mark survives the refresh"
);
settle(&coordinator(&harness)).await;
assert_eq!(
harness.status("alpha").await,
"queued",
"a marked, eligible, ordinary target must gain queue intent without another Start; plan={:?}",
coordinator(&harness).last_plan()
);
assert_eq!(
harness.queue_port.wakes(),
wakes_before + 1,
"one applied membership change wakes scheduler analysis exactly once"
);
assert!(
!coordinator(&harness).is_armed(),
"a reconciled batch releases its deadline"
);
assert!(
coordinator(&harness).last_failure().is_none(),
"a batch that reconciled reports no lifecycle failure"
);
assert_eq!(harness.status("active").await, "applying");
assert!(harness.queue.contains("alpha").await);
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_unloadable_target_reports_a_stable_reason() {
let harness = running_harness(&["active"]);
start_executing(&harness, "active").await;
harness
.application
.apply(OperatorIntent::SetExecutionMark {
change_id: "ghost".to_string(),
marked: true,
})
.await;
for _ in 0..crate::orchestration::mark_settlement::MARK_SETTLEMENT_ATTEMPTS {
settle(&coordinator(&harness)).await;
}
assert!(
!coordinator(&harness).is_armed(),
"the retry budget is finite"
);
assert_eq!(
coordinator(&harness).last_failure(),
Some(crate::orchestration::mark_settlement::MarkSettlementFailure::UnreconciledBatch),
);
expect_no_settlement(&coordinator(&harness)).await;
assert_eq!(harness.status("active").await, "applying");
}
#[derive(Debug, Clone, Copy)]
enum MarkEntryPoint {
Space,
BulkX,
Api,
}
impl MarkEntryPoint {
async fn mark(self, harness: &AdapterHarness, change_id: &str) {
match self {
Self::Space => {
let mut app = harness.app(&[change_id]);
app.execution_mode = AppExecutionMode::Running;
app.cursor_index = 0;
app.toggle_selection();
drain_marks(harness, &mut app).await;
}
Self::BulkX => {
let mut app = harness.app(&[change_id]);
app.execution_mode = AppExecutionMode::Running;
assert!(
crate::tui::key_handlers::handle_bulk_toggle_key(&mut app).is_empty(),
"a bulk mark emits no TUI command"
);
drain_marks(harness, &mut app).await;
}
Self::Api => {
harness
.application
.apply(OperatorIntent::SetExecutionMark {
change_id: change_id.to_string(),
marked: true,
})
.await;
}
}
}
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_every_mark_adapter_settles_a_late_catalog_target() {
for entry in [
MarkEntryPoint::Space,
MarkEntryPoint::BulkX,
MarkEntryPoint::Api,
] {
let harness = running_harness(&["active", "queued_only", "marked_but_removed"]);
start_executing(&harness, "active").await;
harness
.application
.apply(OperatorIntent::SetQueueIntent {
change_id: "queued_only".to_string(),
queued: true,
})
.await;
harness
.application
.apply(OperatorIntent::SetExecutionMark {
change_id: "marked_but_removed".to_string(),
marked: true,
})
.await;
settle(&coordinator(&harness)).await;
harness
.application
.apply(OperatorIntent::SetQueueIntent {
change_id: "marked_but_removed".to_string(),
queued: false,
})
.await;
let wakes_before = harness.queue_port.wakes();
entry.mark(&harness, "alpha").await;
assert!(
coordinator(&harness).is_armed(),
"{entry:?} must reach the shared coordinator"
);
settle(&coordinator(&harness)).await;
refresh_catalog(
&harness,
&["active", "queued_only", "marked_but_removed", "alpha"],
)
.await;
settle(&coordinator(&harness)).await;
assert_eq!(
harness.status("alpha").await,
"queued",
"{entry:?} did not admit the late-catalog target"
);
assert_eq!(
harness.queue_port.wakes(),
wakes_before + 1,
"{entry:?} woke scheduler analysis more than once"
);
assert_eq!(harness.status("active").await, "applying");
assert_eq!(
harness.status("queued_only").await,
"queued",
"{entry:?} disturbed an explicit queue addition"
);
assert_eq!(
harness.status("marked_but_removed").await,
"not queued",
"{entry:?} undid an explicit queue removal"
);
}
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_space_and_bulk_x_reach_the_same_settlement() {
let space = running_harness(&["alpha"]);
let mut app = space.app(&["alpha"]);
app.execution_mode = AppExecutionMode::Running;
app.cursor_index = 0;
app.toggle_selection();
drain_marks(&space, &mut app).await;
assert!(
coordinator(&space).is_armed(),
"Space must reach the shared coordinator"
);
let bulk = running_harness(&["alpha"]);
let mut bulk_app = bulk.app(&["alpha"]);
bulk_app.execution_mode = AppExecutionMode::Running;
assert!(
crate::tui::key_handlers::handle_bulk_toggle_key(&mut bulk_app).is_empty(),
"a bulk mark emits no TUI command"
);
drain_marks(&bulk, &mut bulk_app).await;
assert!(
coordinator(&bulk).is_armed(),
"bulk `x` must reach the shared coordinator"
);
let space_pass = pending_pass(&coordinator(&space));
let bulk_pass = pending_pass(&coordinator(&bulk));
tokio::time::advance(MARK_STABILITY_WINDOW + Duration::from_millis(1)).await;
space_pass.wait().await;
bulk_pass.wait().await;
assert_eq!(space.status("alpha").await, "queued");
assert_eq!(
space.status("alpha").await,
bulk.status("alpha").await,
"both interactions settle into the identical queue state"
);
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_shared_operator_mark_command_schedules_the_same_settlement() {
let harness = running_harness(&["alpha"]);
harness
.application
.apply(OperatorIntent::SetExecutionMark {
change_id: "alpha".to_string(),
marked: true,
})
.await;
assert!(coordinator(&harness).is_armed());
settle(&coordinator(&harness)).await;
assert_eq!(harness.status("alpha").await, "queued");
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_overlay_input_ownership_arms_nothing() {
let harness = running_harness(&["alpha"]);
let mut app = harness.app(&["alpha"]);
app.execution_mode = AppExecutionMode::Running;
app.show_warning_popup("blocked".to_string(), "an overlay owns input".to_string());
assert!(
crate::tui::key_handlers::handle_warning_popup_key(
&mut app,
crossterm::event::KeyEvent::new(
crossterm::event::KeyCode::Char('x'),
crossterm::event::KeyModifiers::NONE,
),
),
"the warning popup must own the keypress"
);
drain_marks(&harness, &mut app).await;
assert!(harness.marks.marked_ids().is_empty());
assert!(!coordinator(&harness).is_armed());
expect_no_settlement(&coordinator(&harness)).await;
assert_eq!(harness.status("alpha").await, "not queued");
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_parked_persistent_scheduler_remains_eligible() {
let harness = running_harness(&["alpha"]);
let mut app = harness.app(&["alpha"]);
app.execution_mode = AppExecutionMode::Select;
app.persistent_scheduler_idle = true;
app.cursor_index = 0;
app.toggle_selection();
drain_marks(&harness, &mut app).await;
assert!(coordinator(&harness).is_armed());
settle(&coordinator(&harness)).await;
assert_eq!(harness.status("alpha").await, "queued");
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_process_without_live_scheduler_stays_mark_only() {
let harness = AdapterHarness::new(&["alpha"]);
for mode in [
AppExecutionMode::Select,
AppExecutionMode::Stopping,
AppExecutionMode::Stopped,
AppExecutionMode::Error,
] {
let mut app = harness.app(&["alpha"]);
app.execution_mode = mode;
app.cursor_index = 0;
app.toggle_selection();
drain_marks(&harness, &mut app).await;
assert_eq!(
harness.marks.is_marked("alpha"),
app.changes[0].selected,
"{mode:?} still applies the mark itself"
);
assert!(
!coordinator(&harness).is_armed(),
"{mode:?} must arm no stability deadline without a live scheduler"
);
}
expect_no_settlement(&coordinator(&harness)).await;
assert_eq!(harness.status("alpha").await, "not queued");
assert!(
harness.queue.pop().await.is_none(),
"DynamicQueue is untouched"
);
}
#[tokio::test(start_paused = true)]
async fn running_mark_reanalysis_rejected_start_leaves_no_delayed_queue_effect() {
let harness = running_harness(&["alpha"]);
let mut app = harness.app(&["alpha"]);
app.execution_mode = AppExecutionMode::Select;
app.changes[0].parallel_eligibility =
crate::orchestration::operator_command::ParallelEligibility::UncommittedProposalFiles;
app.publish_parallel_runtime();
assert_eq!(harness.parallel.ineligible_ids(), vec!["alpha".to_string()]);
harness
.run(
&mut app,
TuiCommand::StartProcessing(vec!["alpha".to_string()]),
)
.await;
assert!(
harness.marks.is_marked("alpha"),
"the admission mark write still happened"
);
assert!(
!coordinator(&harness).is_armed(),
"a Start-admission mark write must arm no delayed settlement"
);
expect_no_settlement(&coordinator(&harness)).await;
assert_eq!(
harness.status("alpha").await,
"not queued",
"a rejected Start leaves no queue effect, immediate or delayed"
);
assert!(harness.queue.pop().await.is_none());
}