use std::sync::Arc;
use super::tests::{create_test_change, AdapterHarness};
use super::*;
use crate::events::ExecutionEvent;
use crate::orchestration::operator_command::ParallelEligibility;
use crate::orchestration::run_control::testing::SchedulerCall;
use crate::tui::state::AppState;
use crate::tui::types::AppExecutionMode;
use crate::web::remote_control_api::dto::{CommandSpec, ErrorCode};
use crate::web::remote_control_api::executor::{RemoteControlExecutor, SharedServiceExecutor};
use crate::web::state::{WebEventSink, WebState};
const CHANGES: [&str; 2] = ["c1", "c2"];
struct Wired {
harness: AdapterHarness,
app: AppState,
web: Arc<WebState>,
change_ids: Vec<String>,
baseline_revision: u64,
}
async fn wired(
change_ids: &[&str],
mode: AppExecutionMode,
harness: AdapterHarness,
mut app: AppState,
) -> Wired {
app.execution_mode = mode;
app.warning_message = None;
let web = Arc::new(WebState::new(&[]));
web.set_shared_state(harness.state.clone()).await;
web.set_execution_marks(harness.marks.clone()).await;
web.set_parallel_runtime(harness.parallel.clone()).await;
web.set_repo_root(std::path::PathBuf::from("/repo")).await;
let changes: Vec<_> = change_ids.iter().map(|id| create_test_change(id)).collect();
web.seed_workspace_observation_for_tests(&changes, app_mode_string(&mode))
.await;
web.sync_remote_control_projection().await;
harness.attach(Arc::new(WebEventSink::new(web.clone())));
harness.attach_revisions(web.clone());
harness.core_mode.set(mode.operator_mode());
let baseline_revision = web.remote_control().projection().revision();
Wired {
harness,
app,
web,
change_ids: change_ids.iter().map(|id| (*id).to_string()).collect(),
baseline_revision,
}
}
#[derive(Debug, PartialEq, Eq)]
struct SharedEffects {
scheduler: Vec<SchedulerCall>,
statuses: Vec<(String, String)>,
reducer_queue: Vec<String>,
explicit_retries: Vec<String>,
active_resolver: Option<String>,
queued_resolves: Vec<String>,
marks: Vec<String>,
mode: crate::orchestration::operator_command::OperatorMode,
}
#[derive(Debug, PartialEq, Eq)]
struct Effects {
shared: SharedEffects,
tui_mode: AppExecutionMode,
tui_marked_rows: Vec<String>,
tui_resolving: bool,
web_mode: String,
web_marks: Vec<String>,
web_resolving: bool,
dispatches: usize,
revisions: u64,
}
impl Wired {
fn row(&self, change_id: &str) -> &crate::tui::state::ChangeState {
self.app
.changes
.iter()
.find(|change| change.id == change_id)
.expect("the arranged change exists")
}
async fn shared_effects(&self) -> SharedEffects {
let (statuses, reducer_queue) = {
let guard = self.harness.state.read().await;
let statuses = self
.change_ids
.iter()
.map(|id| (id.clone(), guard.display_status(id).to_string()))
.collect();
let mut queued = guard.queued_change_ids();
queued.sort();
(statuses, queued)
};
SharedEffects {
scheduler: self.harness.scheduler.calls(),
statuses,
reducer_queue,
explicit_retries: self
.harness
.queue
.drain_explicit_retries()
.await
.into_iter()
.map(|edge| edge.change_id)
.collect(),
active_resolver: self.harness.resolves.active(),
queued_resolves: self.harness.resolves.waiting(),
marks: self.harness.marks.marked_ids(),
mode: self.harness.core_mode.get(),
}
}
async fn effects(&self) -> Effects {
let shared = self.shared_effects().await;
let (snapshot, _, _) = self.web.remote_control().projection().snapshot();
let mut web_marks: Vec<String> = snapshot
.changes
.iter()
.filter(|change| change.execution_marked)
.map(|change| change.id.clone())
.collect();
web_marks.sort();
let web = self.web.get_state().await;
Effects {
shared,
tui_mode: self.app.execution_mode,
tui_marked_rows: self
.app
.changes
.iter()
.filter(|change| change.selected)
.map(|change| change.id.clone())
.collect(),
tui_resolving: self.app.is_resolving(),
web_mode: web.app_mode,
web_marks,
web_resolving: web.is_resolving,
dispatches: self.harness.dispatch_count(),
revisions: self.web.remote_control().projection().revision() - self.baseline_revision,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Settlement {
Changed,
NoOp,
Failed(ErrorCode),
}
impl Settlement {
fn is_reported_to_the_operator(self) -> bool {
!matches!(self, Self::Changed)
}
}
#[derive(Debug, Clone, Copy)]
enum Setup {
Bare,
Marked,
MarkedWithLiveScheduler,
MarkedWithFailingLaunch,
MarkedError,
MarkedStopped,
LiveScheduler,
MergeWait,
MergeWaitAlreadyReserved,
BothMergeWaitingFirstReserved,
MarkedWithIneligible,
LiveSchedulerWithInFlightExecution,
QueuedInLiveRun,
ActiveWithoutCancellationHandle,
}
async fn arrange(harness: &AdapterHarness, setup: Setup) {
let mark_all = || {
harness
.marks
.replace(CHANGES.iter().map(|id| (*id).to_string()))
};
match setup {
Setup::Bare => {}
Setup::Marked => mark_all(),
Setup::MarkedWithLiveScheduler => {
mark_all();
harness.scheduler.set_running(true);
}
Setup::MarkedWithFailingLaunch => {
mark_all();
harness.scheduler.fail_launch("runtime refused the launch");
}
Setup::MarkedError => {
harness
.state
.write()
.await
.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c1".to_string(),
error: "boom".to_string(),
});
harness.marks.replace(["c1".to_string()]);
}
Setup::MarkedStopped => {
harness.state.write().await.apply_command(
crate::orchestration::state::ReducerCommand::StopChange("c1".to_string()),
);
harness.marks.replace(["c1".to_string()]);
}
Setup::LiveScheduler => harness.scheduler.set_running(true),
Setup::MergeWait => merge_wait(harness, "c1").await,
Setup::MergeWaitAlreadyReserved => {
merge_wait(harness, "c1").await;
harness
.run_control
.resolve_merge("c1")
.await
.expect("the first resolve of a merge-wait change is accepted");
}
Setup::BothMergeWaitingFirstReserved => {
merge_wait(harness, "c1").await;
merge_wait(harness, "c2").await;
harness
.run_control
.resolve_merge("c1")
.await
.expect("the first resolve of a merge-wait change is accepted");
}
Setup::ActiveWithoutCancellationHandle => {
harness.scheduler.set_running(true);
harness
.state
.write()
.await
.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c1".to_string(),
command: "apply".to_string(),
});
}
Setup::LiveSchedulerWithInFlightExecution => {
use crate::tui::stop_classification::{ExecutionEvidence, ShutdownWorkEvidence};
harness.scheduler.set_running(true);
harness
.scheduler
.set_activity(crate::tui::stop_classification::StopActivitySnapshot {
execution_handles: ExecutionEvidence::Known { registered: 1 },
reducer_agent_execution_active: true,
shutdown_work: ShutdownWorkEvidence::Known { pending: true },
});
}
Setup::QueuedInLiveRun => {
harness.scheduler.set_running(true);
let mut guard = harness.state.write().await;
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
"c1".to_string(),
));
}
Setup::MarkedWithIneligible => {
mark_all();
{
let mut guard = harness.state.write().await;
for id in CHANGES {
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
id.to_string(),
));
}
}
harness.parallel.set_parallel_ineligible([(
"c2".to_string(),
ParallelEligibility::UncommittedProposalFiles,
)]);
}
}
}
async fn merge_wait(harness: &AdapterHarness, change_id: &str) {
harness
.state
.write()
.await
.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "manual resolution required".to_string(),
auto_resumable: false,
});
}
fn app_mode_string(mode: &AppExecutionMode) -> &'static str {
mode.app_mode_token()
}
async fn arranged(setup: Setup, mode: AppExecutionMode) -> Wired {
let harness = AdapterHarness::new(&CHANGES);
arrange(&harness, setup).await;
let ineligible = harness.parallel.ineligible_ids();
let mut app = harness.app(&CHANGES);
for change in &mut app.changes {
change.parallel_eligibility = if ineligible.contains(&change.id) {
ParallelEligibility::UncommittedProposalFiles
} else {
ParallelEligibility::Eligible
};
}
app.publish_parallel_runtime();
app.apply_display_statuses_from_reducer(&harness.state.read().await.all_display_statuses());
app.sync_execution_marks_from_store();
wired(&CHANGES, mode, harness, app).await
}
struct Run {
before: Effects,
after: Effects,
report: Option<String>,
}
async fn through_tui(setup: Setup, mode: AppExecutionMode, command: TuiCommand) -> Run {
let mut wired = arranged(setup, mode).await;
let before = wired.effects().await;
let is_two_phase = matches!(command, TuiCommand::DequeueChange(_));
wired.harness.run(&mut wired.app, command).await;
if is_two_phase {
wired
.harness
.await_feedback(&mut wired.app, std::time::Duration::from_secs(5))
.await;
}
let report = wired.app.warning_message.clone();
Run {
before,
after: wired.effects().await,
report,
}
}
async fn through_v2(
setup: Setup,
mode: AppExecutionMode,
command: CommandSpec,
) -> (Run, Settlement) {
let mut wired = arranged(setup, mode).await;
let before = wired.effects().await;
let executor = SharedServiceExecutor::new(wired.harness.application.clone(), wired.web.clone());
let (settlement, detail) = match executor.execute(&command).await {
Ok(summary) if summary.changed => (Settlement::Changed, summary.detail),
Ok(summary) => (Settlement::NoOp, summary.detail),
Err(failure) => (
Settlement::Failed(failure.error_code),
Some(failure.message),
),
};
wired.harness.deliver(&mut wired.app).await;
(
Run {
before,
after: wired.effects().await,
report: detail,
},
settlement,
)
}
struct Row {
name: &'static str,
setup: Setup,
mode: AppExecutionMode,
tui: TuiCommand,
v2: CommandSpec,
expect: Settlement,
notice: Option<&'static str>,
}
fn rows() -> Vec<Row> {
vec![
Row {
name: "start with an idle scheduler spawns one run over the marked set",
setup: Setup::Marked,
mode: AppExecutionMode::Select,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "start with a live scheduler wakes it instead of spawning a second run",
setup: Setup::MarkedWithLiveScheduler,
mode: AppExecutionMode::Select,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "start from a stopped run resumes the marked set",
setup: Setup::Marked,
mode: AppExecutionMode::Stopped,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "start from a stopped run resumes a preserved stopped mark",
setup: Setup::MarkedStopped,
mode: AppExecutionMode::Stopped,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "a preserved stopped mark is not resumed outside Stopped",
setup: Setup::MarkedStopped,
mode: AppExecutionMode::Select,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
Row {
name: "start under a live run is refused when no mark is retryable",
setup: Setup::MarkedWithLiveScheduler,
mode: AppExecutionMode::Running,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
Row {
name: "start with an empty target set is not a success",
setup: Setup::Bare,
mode: AppExecutionMode::Select,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
Row {
name: "a runtime launch failure is reported, not claimed as started",
setup: Setup::MarkedWithFailingLaunch,
mode: AppExecutionMode::Select,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Failed(ErrorCode::InternalError),
notice: None,
},
Row {
name: "retry routes a marked error row and dispatches the scheduler",
setup: Setup::MarkedError,
mode: AppExecutionMode::Error,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "retry without retryable evidence is rejected without effects",
setup: Setup::Marked,
mode: AppExecutionMode::Error,
tui: TuiCommand::StartProcessing(Vec::new()),
v2: CommandSpec::Start,
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
Row {
name: "graceful stop while running sets the stop request",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Running,
tui: TuiCommand::Stop,
v2: CommandSpec::Stop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "graceful stop outside running is refused",
setup: Setup::Bare,
mode: AppExecutionMode::Select,
tui: TuiCommand::Stop,
v2: CommandSpec::Stop,
expect: Settlement::Failed(ErrorCode::LifecycleConflict),
notice: None,
},
Row {
name: "graceful stop from persistent-idle Ready addresses the live scheduler",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Select,
tui: TuiCommand::Stop,
v2: CommandSpec::Stop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "force stop from persistent-idle Ready cancels the live scheduler",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Select,
tui: TuiCommand::ForceStop,
v2: CommandSpec::ForceStop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "cancel stop while stopping withdraws the request",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Stopping,
tui: TuiCommand::CancelStop,
v2: CommandSpec::CancelStop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "cancel stop outside stopping is refused",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Running,
tui: TuiCommand::CancelStop,
v2: CommandSpec::CancelStop,
expect: Settlement::Failed(ErrorCode::LifecycleConflict),
notice: None,
},
Row {
name: "force stop while running cancels the live run",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Running,
tui: TuiCommand::ForceStop,
v2: CommandSpec::ForceStop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "force stop escalates a pending graceful stop",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Stopping,
tui: TuiCommand::ForceStop,
v2: CommandSpec::ForceStop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "force stop outside running and stopping is refused",
setup: Setup::Bare,
mode: AppExecutionMode::Select,
tui: TuiCommand::ForceStop,
v2: CommandSpec::ForceStop,
expect: Settlement::Failed(ErrorCode::LifecycleConflict),
notice: None,
},
Row {
name: "resolve of a merge-wait change takes the single resolver slot",
setup: Setup::MergeWait,
mode: AppExecutionMode::Select,
tui: TuiCommand::ResolveMerge("c1".to_string()),
v2: CommandSpec::ResolveMerge {
change_id: "c1".to_string(),
},
expect: Settlement::Changed,
notice: None,
},
Row {
name: "a duplicate resolve submission does not reserve twice",
setup: Setup::MergeWaitAlreadyReserved,
mode: AppExecutionMode::Select,
tui: TuiCommand::ResolveMerge("c1".to_string()),
v2: CommandSpec::ResolveMerge {
change_id: "c1".to_string(),
},
expect: Settlement::NoOp,
notice: None,
},
Row {
name: "a second resolve target takes a FIFO position, not a second resolver",
setup: Setup::BothMergeWaitingFirstReserved,
mode: AppExecutionMode::Select,
tui: TuiCommand::ResolveMerge("c2".to_string()),
v2: CommandSpec::ResolveMerge {
change_id: "c2".to_string(),
},
expect: Settlement::Changed,
notice: None,
},
Row {
name: "force stop with in-flight execution waits for the safe boundary",
setup: Setup::LiveSchedulerWithInFlightExecution,
mode: AppExecutionMode::Running,
tui: TuiCommand::ForceStop,
v2: CommandSpec::ForceStop,
expect: Settlement::Changed,
notice: None,
},
Row {
name: "queue add in a live run commits the same intent on both adapters",
setup: Setup::LiveScheduler,
mode: AppExecutionMode::Running,
tui: TuiCommand::AddToQueue("c1".to_string()),
v2: CommandSpec::SetQueueIntent {
change_id: "c1".to_string(),
queued: true,
},
expect: Settlement::Changed,
notice: None,
},
Row {
name: "queue remove takes the same intent back out on both adapters",
setup: Setup::QueuedInLiveRun,
mode: AppExecutionMode::Running,
tui: TuiCommand::RemoveFromQueue("c1".to_string()),
v2: CommandSpec::SetQueueIntent {
change_id: "c1".to_string(),
queued: false,
},
expect: Settlement::Changed,
notice: None,
},
Row {
name: "queue add for an already queued row changes nothing on either adapter",
setup: Setup::QueuedInLiveRun,
mode: AppExecutionMode::Running,
tui: TuiCommand::AddToQueue("c1".to_string()),
v2: CommandSpec::SetQueueIntent {
change_id: "c1".to_string(),
queued: true,
},
expect: Settlement::NoOp,
notice: None,
},
Row {
name: "stop-and-dequeue clears exactly its target on both adapters",
setup: Setup::QueuedInLiveRun,
mode: AppExecutionMode::Running,
tui: TuiCommand::DequeueChange("c1".to_string()),
v2: CommandSpec::StopAndDequeue {
change_id: "c1".to_string(),
},
expect: Settlement::Changed,
notice: None,
},
Row {
name: "stop-and-dequeue of an unprovable termination commits no dequeue",
setup: Setup::ActiveWithoutCancellationHandle,
mode: AppExecutionMode::Running,
tui: TuiCommand::DequeueChange("c1".to_string()),
v2: CommandSpec::StopAndDequeue {
change_id: "c1".to_string(),
},
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
Row {
name: "resolve of a stale target is refused without a reservation",
setup: Setup::Bare,
mode: AppExecutionMode::Select,
tui: TuiCommand::ResolveMerge("c1".to_string()),
v2: CommandSpec::ResolveMerge {
change_id: "c1".to_string(),
},
expect: Settlement::Failed(ErrorCode::TargetIneligible),
notice: None,
},
]
}
#[tokio::test]
async fn tui_and_v2_settle_every_lifecycle_intent_identically() {
for row in rows() {
let tui = through_tui(row.setup, row.mode, row.tui.clone()).await;
let (v2, v2_settlement) = through_v2(row.setup, row.mode, row.v2).await;
assert_eq!(
v2_settlement, row.expect,
"{}: /api/v2 settlement must match the declared outcome",
row.name
);
assert_eq!(
tui.after, v2.after,
"{}: the TUI and /api/v2 must produce the same reducer, scheduler, mark, resolver, \
frontend, and event effects",
row.name
);
for (adapter, run) in [("the TUI", &tui), ("/api/v2", &v2)] {
match row.expect {
Settlement::Changed => assert_ne!(
run.after, run.before,
"{}: {adapter} declared a change and produced none",
row.name
),
Settlement::NoOp | Settlement::Failed(_) => assert_eq!(
run.after, run.before,
"{}: {adapter} must leave the process exactly as it found it",
row.name
),
}
}
let must_report = row.expect.is_reported_to_the_operator() || row.notice.is_some();
assert_eq!(
tui.report.is_some(),
must_report,
"{}: the TUI must surface exactly what /api/v2 reports, got {:?}",
row.name,
tui.report
);
if let Some(cleared) = row.notice {
let tui_message = tui
.report
.expect("a consequence row surfaces a TUI message");
assert!(
tui_message.contains(cleared),
"{}: the TUI must name '{cleared}', got {tui_message:?}",
row.name
);
let v2_detail = v2.report.expect("a consequence row carries a v2 detail");
assert!(
v2_detail.contains(cleared),
"{}: /api/v2 must name '{cleared}', got {v2_detail:?}",
row.name
);
}
}
}
#[tokio::test]
async fn persistent_idle_commands_use_live_scheduler() {
use crate::events::persistent_idle_may_project_ready;
use crate::tui::key_handlers::{esc_stop_action, EscStopAction};
use crate::tui::types::StopMode;
let harness = AdapterHarness::new(&CHANGES);
harness
.marks
.replace(CHANGES.iter().map(|id| id.to_string()));
harness.scheduler.set_running(true);
let mut app = harness.app(&CHANGES);
app.execution_mode = AppExecutionMode::Select;
app.persistent_scheduler_idle = true;
let statuses_before = harness.state.read().await.all_display_statuses();
harness
.run_control
.operator()
.set_execution_mark("c1", true)
.await
.expect("marking in Select is accepted");
assert_eq!(
harness.state.read().await.all_display_statuses(),
statuses_before,
"a mark made in idle Ready must not synthesize queue intent"
);
harness
.run(&mut app, TuiCommand::StartProcessing(Vec::new()))
.await;
assert!(
harness.scheduler.started_targets().is_empty(),
"Start against a live idle scheduler must not spawn a second run"
);
assert!(
harness.scheduler.calls().contains(&SchedulerCall::Notified),
"Start must wake the scheduler that is already alive"
);
for id in CHANGES {
assert_eq!(
harness.state.read().await.display_status(id),
"queued",
"Start applies the existing reducer queue intent for '{id}'"
);
}
assert_eq!(
app.execution_mode,
AppExecutionMode::Running,
"an accepted Start opens the run episode the operator asked for"
);
assert!(
!app.persistent_scheduler_idle,
"the accepted Start closes the idle presentation episode"
);
assert_eq!(
esc_stop_action(&AppExecutionMode::Select, &StopMode::None, true),
EscStopAction::RequestGracefulStop,
"idle Ready keeps the first-Esc graceful stop"
);
assert_eq!(
esc_stop_action(&AppExecutionMode::Select, &StopMode::None, false),
EscStopAction::None,
"pre-run Select has no scheduler to stop"
);
let harness = AdapterHarness::new(&CHANGES);
harness.scheduler.set_running(true);
let mut app = harness.app(&CHANGES);
app.execution_mode = AppExecutionMode::Select;
app.persistent_scheduler_idle = true;
harness.run(&mut app, TuiCommand::Stop).await;
assert_eq!(
harness.scheduler.calls(),
vec![SchedulerCall::GracefulStop(true), SchedulerCall::Notified],
"the stop request is recorded before the idle waiter is woken"
);
assert_eq!(app.execution_mode, AppExecutionMode::Stopping);
assert_eq!(app.stop_mode, StopMode::GracefulPending);
assert!(
app.persistent_scheduler_idle,
"an idle-origin stop keeps its episode identity"
);
assert_eq!(
esc_stop_action(&app.execution_mode, &app.stop_mode, true),
EscStopAction::RequestImmediateStop
);
harness.run(&mut app, TuiCommand::CancelStop).await;
assert_eq!(
app.execution_mode,
AppExecutionMode::Select,
"withdrawing an idle-origin stop restores Ready"
);
assert_eq!(app.stop_mode, StopMode::None);
assert!(app.persistent_scheduler_idle);
app.handle_orchestrator_event(ExecutionEvent::WorkspacePreparationStarted {
change_id: "c1".to_string(),
});
assert_eq!(app.execution_mode, AppExecutionMode::Running);
assert!(!app.persistent_scheduler_idle);
harness.run(&mut app, TuiCommand::Stop).await;
harness.run(&mut app, TuiCommand::CancelStop).await;
assert_eq!(app.execution_mode, AppExecutionMode::Running);
let harness = AdapterHarness::new(&CHANGES);
harness.scheduler.set_running(true);
let mut app = harness.app(&CHANGES);
app.execution_mode = AppExecutionMode::Select;
app.persistent_scheduler_idle = true;
harness.run(&mut app, TuiCommand::ForceStop).await;
assert!(
harness
.scheduler
.calls()
.contains(&SchedulerCall::Cancelled),
"force stop must cancel the live scheduler behind idle Ready"
);
assert!(
harness.scheduler.started_targets().is_empty(),
"force stop never spawns anything"
);
let harness = AdapterHarness::new(&CHANGES);
let mut app = harness.app(&CHANGES);
app.execution_mode = AppExecutionMode::Select;
app.persistent_scheduler_idle = true;
harness.run(&mut app, TuiCommand::Stop).await;
assert!(
harness.scheduler.calls().is_empty(),
"a stale idle fact must not authorize a stop against an exited scheduler"
);
assert_eq!(app.execution_mode, AppExecutionMode::Select);
assert!(app
.warning_message
.as_deref()
.is_some_and(|message| message.contains("stop is not available")));
assert!(persistent_idle_may_project_ready(
AppExecutionMode::Running.app_mode_token()
));
for retained in [
AppExecutionMode::Select,
AppExecutionMode::Stopping,
AppExecutionMode::Stopped,
AppExecutionMode::Error,
] {
assert!(
!persistent_idle_may_project_ready(retained.app_mode_token()),
"{retained:?} must not be turned into persistent-idle Ready"
);
}
}
#[tokio::test]
async fn admitted_work_restores_running_after_idle() {
use crate::events::EventSink;
use crate::orchestration::state::OrchestratorState;
use crate::web::state::WebEventSink;
struct Step {
what: &'static str,
event: ExecutionEvent,
mode: &'static str,
idle: bool,
}
let prepare = |id: &str| ExecutionEvent::WorkspacePreparationStarted {
change_id: id.to_string(),
};
let steps = vec![
Step {
what: "ordinary workspace preparation starts the run",
event: prepare("change-a"),
mode: "running",
idle: false,
},
Step {
what: "the scheduler parks with nothing to execute",
event: ExecutionEvent::PersistentSchedulerIdle,
mode: "select",
idle: true,
},
Step {
what: "a no-op wake analyses and admits nothing",
event: ExecutionEvent::AnalysisStarted {
remaining_changes: 1,
attempt_id: "attempt-1".to_string(),
},
mode: "select",
idle: true,
},
Step {
what: "a catalog refresh is not execution evidence either",
event: ExecutionEvent::WorktreesRefreshed {
worktrees: Vec::new(),
},
mode: "select",
idle: true,
},
Step {
what: "actual admitted work resumes Running and closes the episode",
event: prepare("change-a"),
mode: "running",
idle: false,
},
Step {
what: "the scheduler parks again: a second idle edge",
event: ExecutionEvent::PersistentSchedulerIdle,
mode: "select",
idle: true,
},
Step {
what: "scheduler-owned resolve work resumes Running",
event: ExecutionEvent::ResolveStarted {
change_id: "change-a".to_string(),
command: "resolve".to_string(),
},
mode: "running",
idle: false,
},
Step {
what: "and parks once more",
event: ExecutionEvent::PersistentSchedulerIdle,
mode: "select",
idle: true,
},
Step {
what: "scheduler-owned base-lane rejection review resumes Running",
event: ExecutionEvent::WorkspaceStatusUpdated {
change_id: "change-a".to_string(),
workspace_name: "ws-a".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
},
mode: "running",
idle: false,
},
];
let reducer = Arc::new(tokio::sync::RwLock::new(OrchestratorState::new(
vec!["change-a".to_string()],
10,
)));
let web_state = Arc::new(WebState::new(&[]));
web_state.set_shared_state(reducer.clone()).await;
let sinks: Vec<Arc<dyn EventSink>> = vec![Arc::new(WebEventSink::new(web_state.clone()))];
let mut app = AppState::new(Vec::new());
for step in steps {
crate::events::dispatch_event(reducer.as_ref(), &sinks, step.event.clone()).await;
app.handle_orchestrator_event(step.event);
let web = web_state.get_state().await;
assert_eq!(web.app_mode, step.mode, "web: {}", step.what);
assert_eq!(
web.persistent_scheduler_idle, step.idle,
"web idle episode: {}",
step.what
);
assert_eq!(
app.execution_mode.app_mode_token(),
step.mode,
"tui: {}",
step.what
);
assert_eq!(
app.persistent_scheduler_idle, step.idle,
"tui idle episode: {}",
step.what
);
assert_eq!(
(app.execution_mode == AppExecutionMode::Running),
(web.app_mode == "running"),
"tui/web divergence: {}",
step.what
);
}
}
type BulkRow = (&'static str, &'static str, bool);
async fn arrange_bulk(harness: &AdapterHarness, rows: &[BulkRow], marked: &[&str]) {
harness
.parallel
.set_parallel_ineligible(rows.iter().filter(|(_, _, ok)| !ok).map(|(id, ..)| {
(
id.to_string(),
ParallelEligibility::UncommittedProposalFiles,
)
}));
harness
.marks
.replace(marked.iter().map(|id| id.to_string()));
let mut guard = harness.state.write().await;
for (id, status, _) in rows {
match *status {
"queued" => {
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
(*id).to_string(),
));
}
"applying" => guard.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: (*id).to_string(),
command: "apply".to_string(),
}),
"rejected" => guard.apply_execution_event(&ExecutionEvent::ChangeRejected {
change_id: (*id).to_string(),
reason: "acceptance refused the proposal".to_string(),
}),
"merged" => {
guard.apply_execution_event(&ExecutionEvent::ChangeArchived((*id).to_string()));
guard.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: (*id).to_string(),
revision: "deadbeef".to_string(),
});
}
"not queued" => {}
other => panic!("unsupported bulk-mark arrangement status '{other}'"),
}
}
}
async fn arranged_bulk(rows: &[BulkRow], marked: &[&str], mode: AppExecutionMode) -> Wired {
let ids: Vec<&str> = rows.iter().map(|(id, ..)| *id).collect();
let harness = AdapterHarness::new(&ids);
arrange_bulk(&harness, rows, marked).await;
let mut app = harness.app(&ids);
app.apply_display_statuses_from_reducer(&harness.state.read().await.all_display_statuses());
for (index, (_, _, eligible)) in rows.iter().enumerate() {
app.changes[index].parallel_eligibility = if *eligible {
ParallelEligibility::Eligible
} else {
ParallelEligibility::UncommittedProposalFiles
};
app.changes[index].selected = marked.contains(&app.changes[index].id.as_str());
}
app.publish_parallel_runtime();
wired(&ids, mode, harness, app).await
}
async fn bulk_through_tui(
rows: &[BulkRow],
marked: &[&str],
mode: AppExecutionMode,
) -> SharedEffects {
let mut wired = arranged_bulk(rows, marked, mode).await;
let commands = wired.app.toggle_all_marks();
for (change_id, marked) in wired.app.take_pending_mark_writes() {
wired
.harness
.run_control
.operator()
.apply_execution_mark(&change_id, marked)
.await;
}
wired.app.sync_execution_marks_from_store();
for command in commands {
wired.harness.run(&mut wired.app, command).await;
}
wired.shared_effects().await
}
async fn bulk_through_v2(
rows: &[BulkRow],
marked: &[&str],
mode: AppExecutionMode,
) -> (SharedEffects, Settlement) {
let mut wired = arranged_bulk(rows, marked, mode).await;
let executor = SharedServiceExecutor::new(wired.harness.application.clone(), wired.web.clone());
let settlement = match executor
.execute(&CommandSpec::SetAllExecutionMarks {})
.await
{
Ok(summary) if summary.changed => Settlement::Changed,
Ok(_) => Settlement::NoOp,
Err(failure) => Settlement::Failed(failure.error_code),
};
wired.harness.deliver(&mut wired.app).await;
(wired.shared_effects().await, settlement)
}
#[tokio::test]
async fn tui_and_v2_derive_the_same_bulk_mark_target_set_and_exclusions() {
struct Case {
name: &'static str,
rows: Vec<BulkRow>,
marked: Vec<&'static str>,
mode: AppExecutionMode,
expect: Settlement,
}
let cases = vec![
Case {
name: "select mode marks every eligible row and skips a final one",
rows: vec![("c1", "not queued", true), ("c2", "rejected", true)],
marked: vec![],
mode: AppExecutionMode::Select,
expect: Settlement::Changed,
},
Case {
name: "a fully marked eligible set unmarks, ignoring the excluded row's mark state",
rows: vec![("c1", "not queued", true), ("c2", "rejected", true)],
marked: vec!["c1"],
mode: AppExecutionMode::Select,
expect: Settlement::Changed,
},
Case {
name: "an uncommitted row joins the target set: eligibility is a start-time fact",
rows: vec![("c1", "not queued", true), ("c2", "not queued", false)],
marked: vec![],
mode: AppExecutionMode::Select,
expect: Settlement::Changed,
},
Case {
name: "running mode marks the active row too, and writes no queue intent",
rows: vec![("c1", "not queued", true), ("c2", "applying", true)],
marked: vec![],
mode: AppExecutionMode::Running,
expect: Settlement::Changed,
},
Case {
name: "an already-queued row is marked without its queue membership moving",
rows: vec![("c1", "queued", true), ("c2", "applying", true)],
marked: vec!["c1"],
mode: AppExecutionMode::Running,
expect: Settlement::Changed,
},
Case {
name: "a fully marked non-terminal set unmarks without dequeuing anything",
rows: vec![("c1", "queued", true), ("c2", "applying", true)],
marked: vec!["c1", "c2"],
mode: AppExecutionMode::Running,
expect: Settlement::Changed,
},
Case {
name: "a terminal-only target set changes nothing on either side",
rows: vec![("c1", "rejected", true), ("c2", "merged", true)],
marked: vec![],
mode: AppExecutionMode::Running,
expect: Settlement::NoOp,
},
];
for case in cases {
let tui = bulk_through_tui(&case.rows, &case.marked, case.mode).await;
let (v2, settlement) = bulk_through_v2(&case.rows, &case.marked, case.mode).await;
assert_eq!(
settlement, case.expect,
"{}: /api/v2 settlement must match the declared outcome",
case.name
);
assert_eq!(
tui, v2,
"{}: the TUI and /api/v2 must produce the same marks, queue intent, and scheduler effects",
case.name
);
}
}
#[tokio::test]
async fn processing_error_keeps_bulk_mark_available() {
const RUNNING_RUN: [&str; 3] = ["alpha", "beta", "gamma"];
let harness = AdapterHarness::new(&RUNNING_RUN);
{
let mut guard = harness.state.write().await;
for id in ["alpha", "beta"] {
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
id.to_string(),
));
}
guard.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "alpha".to_string(),
command: "apply".to_string(),
});
}
harness
.marks
.replace(["alpha".to_string(), "beta".to_string()]);
let mut app = harness.app(&RUNNING_RUN);
app.apply_display_statuses_from_reducer(&harness.state.read().await.all_display_statuses());
app.sync_execution_marks_from_store();
let mut wired = wired(&RUNNING_RUN, AppExecutionMode::Running, harness, app).await;
wired
.harness
.dispatcher
.dispatch(ExecutionEvent::ProcessingError {
id: "alpha".to_string(),
error: "acceptance command attempts exhausted".to_string(),
})
.await;
wired.harness.deliver(&mut wired.app).await;
assert_eq!(
wired.harness.core_mode.get(),
crate::orchestration::operator_command::OperatorMode::Running,
"a change-scoped failure leaves the one process mode alone"
);
assert_eq!(
wired.app.execution_mode,
AppExecutionMode::Running,
"so the frame the TUI adopts stays Running"
);
assert_eq!(
wired.web.get_state().await.app_mode,
"running",
"and the other frontend agrees at the same dispatch"
);
assert_eq!(
wired.harness.status("alpha").await,
"error",
"the failed change still receives its change-level Error transition"
);
assert!(
!wired.harness.marks.is_marked("alpha"),
"and its stale execution intent is revoked by the same dispatch"
);
assert!(
!wired.row("alpha").selected,
"the row follows the reconciled store rather than its stale cache"
);
assert!(
wired.harness.marks.is_marked("beta"),
"an unrelated change's mark is untouched by another change's failure"
);
let commands = wired.app.toggle_all_marks();
for (change_id, marked) in wired.app.take_pending_mark_writes() {
wired
.harness
.run_control
.operator()
.apply_execution_mark(&change_id, marked)
.await;
}
wired.app.sync_execution_marks_from_store();
for command in commands {
wired.harness.run(&mut wired.app, command).await;
}
let report = wired.app.warning_message.clone().unwrap_or_default();
assert!(
!report.contains("recovery is owned by retry"),
"the Error-mode bulk rejection must not fire for a change-scoped failure: {report}"
);
assert!(
!report.contains("Error mode"),
"and no Error-mode block of any wording may be reported: {report}"
);
assert!(
wired.harness.marks.is_marked("gamma"),
"an unrelated eligible row is still bulk-markable"
);
assert!(wired.row("gamma").selected);
assert_eq!(
wired.harness.status("gamma").await,
"not queued",
"and the mark wrote no Running-mode queue intent"
);
assert!(
wired.harness.marks.is_marked("beta"),
"an already-marked eligible row keeps its mark"
);
assert_eq!(
wired.harness.status("alpha").await,
"error",
"the failed change stays in change-level Error"
);
}