use std::sync::Arc;
use super::tests::{create_test_change, AdapterHarness};
use super::*;
use crate::events::ExecutionEvent;
use crate::orchestration::operator_command::OperatorMode;
use crate::tui::state::AppState;
use crate::tui::types::AppExecutionMode;
use crate::web::state::{WebEventSink, WebState};
struct Converged {
harness: AdapterHarness,
app: AppState,
web: Arc<WebState>,
}
impl Converged {
async fn new(change_ids: &[&str]) -> Self {
let harness = AdapterHarness::new(change_ids);
let mut app = harness.app(change_ids);
app.apply_display_statuses_from_reducer(&harness.state.read().await.all_display_statuses());
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;
let changes: Vec<_> = change_ids.iter().map(|id| create_test_change(id)).collect();
web.seed_workspace_observation_for_tests(&changes, "select")
.await;
web.sync_remote_control_projection().await;
harness.attach(Arc::new(WebEventSink::new(web.clone())));
Self { harness, app, web }
}
async fn remote(&mut self, intent: OperatorIntent) -> ApplicationResult {
let result = self.harness.application.apply(intent).await;
self.harness.deliver(&mut self.app).await;
result
}
async fn local(&mut self, command: TuiCommand) {
self.harness.run(&mut self.app, command).await;
}
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 web_mark(&self, change_id: &str) -> bool {
let (snapshot, _, _) = self.web.remote_control().projection().snapshot();
snapshot
.changes
.iter()
.find(|change| change.id == change_id)
.map(|change| change.execution_marked)
.unwrap_or_default()
}
async fn web_mode(&self) -> String {
self.web.get_state().await.app_mode
}
async fn arrange_running(&mut self) {
self.harness.scheduler.set_running(true);
self.harness.core_mode.set(OperatorMode::Running);
self.app.execution_mode = AppExecutionMode::Running;
let changes: Vec<_> = self
.app
.changes
.iter()
.map(|change| create_test_change(&change.id))
.collect();
self.web
.seed_workspace_observation_for_tests(&changes, "running")
.await;
self.web.sync_remote_control_projection().await;
}
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_remote_mark_reaches_the_next_frame() {
let mut converged = Converged::new(&["alpha", "beta"]).await;
assert!(!converged.row("alpha").selected);
converged
.remote(OperatorIntent::SetExecutionMark {
change_id: "alpha".to_string(),
marked: true,
})
.await;
assert!(
converged.row("alpha").selected,
"a remote mark must be visible on the very next TUI pass"
);
assert!(converged.web_mark("alpha").await);
assert!(
!converged.row("beta").selected,
"only the named target may be written"
);
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_local_run_reaches_web() {
let mut converged = Converged::new(&["alpha"]).await;
converged.harness.marks.set("alpha", true);
converged
.local(TuiCommand::StartProcessing(Vec::new()))
.await;
assert_eq!(
converged.app.execution_mode,
AppExecutionMode::Running,
"the accepted run must move this frontend's mode"
);
assert_eq!(
converged.web_mode().await,
"running",
"and the other frontend's, at the same dispatch"
);
assert_eq!(
converged.harness.core_mode.get(),
OperatorMode::Running,
"both are projections of one Core value"
);
assert_eq!(converged.harness.status("alpha").await, "queued");
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_local_command_preserves_unrelated_remote_marks()
{
let mut converged = Converged::new(&["alpha", "beta"]).await;
converged.harness.marks.set("beta", true);
converged
.remote(OperatorIntent::SetExecutionMark {
change_id: "alpha".to_string(),
marked: true,
})
.await;
assert!(
converged.harness.marks.is_marked("beta"),
"the shared store must retain the API-provided mark"
);
assert!(
converged.row("beta").selected,
"and the TUI must render it rather than the stale row it held"
);
assert!(converged.web_mark("beta").await);
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_queue_presentation_is_not_mark_authority() {
let mut converged = Converged::new(&["alpha", "beta"]).await;
converged
.harness
.state
.write()
.await
.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "beta".to_string(),
error: "boom".to_string(),
});
converged.harness.core_mode.set(OperatorMode::Running);
converged.app.execution_mode = AppExecutionMode::Running;
converged
.remote(OperatorIntent::SetQueueIntent {
change_id: "alpha".to_string(),
queued: true,
})
.await;
assert_eq!(
converged.harness.status("alpha").await,
"queued",
"the queue intent really committed"
);
assert!(
!converged.row("beta").selected,
"an Error row's hidden retry intent must not become a checked row"
);
assert!(!converged.harness.marks.is_marked("beta"));
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_dequeue_clears_only_its_target() {
let mut converged = Converged::new(&["alpha", "beta"]).await;
converged.harness.marks.set("alpha", true);
converged.harness.marks.set("beta", true);
{
let mut guard = converged.harness.state.write().await;
for id in ["alpha", "beta"] {
guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
id.to_string(),
));
}
}
converged.harness.core_mode.set(OperatorMode::Running);
converged.app.execution_mode = AppExecutionMode::Running;
let result = converged
.remote(OperatorIntent::StopAndDequeue {
change_id: "alpha".to_string(),
})
.await;
assert!(
matches!(
&result.outcome,
Ok(ApplicationOutcome::Operator(
OperatorOutcome::Dequeued { .. }
))
),
"an idle queued row dequeues once cancellation confirms: {result:?}"
);
assert!(
!converged.harness.marks.is_marked("alpha"),
"the dequeued target loses its mark"
);
assert!(
converged.harness.marks.is_marked("beta"),
"an unrelated target keeps its mark"
);
assert!(!converged.row("alpha").selected);
assert!(converged.row("beta").selected);
assert!(converged.web_mark("beta").await);
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_process_stop_retains_marks() {
let mut converged = Converged::new(&["alpha", "beta"]).await;
converged.harness.marks.set("alpha", true);
converged.harness.core_mode.set(OperatorMode::Running);
converged.app.execution_mode = AppExecutionMode::Running;
converged
.harness
.dispatcher
.dispatch(ExecutionEvent::Stopped)
.await;
converged.harness.deliver(&mut converged.app).await;
assert_eq!(converged.harness.core_mode.get(), OperatorMode::Stopped);
assert_eq!(converged.app.execution_mode, AppExecutionMode::Stopped);
assert_eq!(converged.web_mode().await, "stopped");
assert!(
converged.harness.marks.is_marked("alpha"),
"a process-level stop must not revoke the operator's next-run intent"
);
assert!(converged.row("alpha").selected);
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_remote_stop_family_reaches_both_frontends() {
let mut converged = Converged::new(&["alpha"]).await;
converged.arrange_running().await;
converged.remote(OperatorIntent::Stop).await;
assert_eq!(converged.app.execution_mode, AppExecutionMode::Stopping);
assert_eq!(converged.web_mode().await, "stopping");
converged.remote(OperatorIntent::CancelStop).await;
assert_eq!(converged.app.execution_mode, AppExecutionMode::Running);
assert_eq!(converged.web_mode().await, "running");
converged.remote(OperatorIntent::ForceStop).await;
assert_eq!(converged.app.execution_mode, AppExecutionMode::Stopped);
assert_eq!(converged.web_mode().await, "stopped");
}
#[tokio::test]
async fn accepted_operator_command_tui_convergence_remote_resolve_reaches_the_next_frame() {
let mut converged = Converged::new(&["alpha"]).await;
converged
.harness
.state
.write()
.await
.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "manual resolution required".to_string(),
auto_resumable: false,
});
converged
.remote(OperatorIntent::ResolveMerge {
change_id: "alpha".to_string(),
})
.await;
assert_eq!(
converged.harness.resolves.active().as_deref(),
Some("alpha"),
"the shared ledger owns the single resolver"
);
assert!(
converged.app.is_resolving(),
"the TUI reads that ledger rather than a cache of its own"
);
assert_eq!(converged.app.execution_mode, AppExecutionMode::Running);
assert!(converged.web.get_state().await.is_resolving);
}