use std::sync::Arc;
use crate::orchestration::state::OrchestratorState;
use crate::tui::queue::DynamicQueue;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ExecutionEvidence {
Known { registered: usize },
Unavailable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ShutdownWorkEvidence {
Known { pending: bool },
Unavailable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct StopActivitySnapshot {
pub execution_handles: ExecutionEvidence,
pub reducer_agent_execution_active: bool,
pub shutdown_work: ShutdownWorkEvidence,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProcessReport {
ForceStopped,
OrdinaryStop,
}
impl ProcessReport {
pub fn is_force_stop(self) -> bool {
matches!(self, ProcessReport::ForceStopped)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ShutdownBarrier {
Required,
NotRequired,
}
impl ShutdownBarrier {
pub fn is_required(self) -> bool {
matches!(self, ShutdownBarrier::Required)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct StopClassification {
pub process_report: ProcessReport,
pub shutdown_barrier: ShutdownBarrier,
}
impl StopActivitySnapshot {
pub fn scheduler_owns_cleanup(&self) -> bool {
matches!(self.execution_handles, ExecutionEvidence::Known { .. })
}
pub fn classify(&self) -> StopClassification {
let execution_active = match self.execution_handles {
ExecutionEvidence::Known { registered } => registered > 0,
ExecutionEvidence::Unavailable => true,
} || self.reducer_agent_execution_active;
let shutdown_work_pending = match self.shutdown_work {
ShutdownWorkEvidence::Known { pending } => pending,
ShutdownWorkEvidence::Unavailable => true,
};
StopClassification {
process_report: if execution_active {
ProcessReport::ForceStopped
} else {
ProcessReport::OrdinaryStop
},
shutdown_barrier: if execution_active || shutdown_work_pending {
ShutdownBarrier::Required
} else {
ShutdownBarrier::NotRequired
},
}
}
}
const REDUCER_SNAPSHOT_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(200);
pub async fn collect_stop_activity_snapshot(
dynamic_queue: &DynamicQueue,
shared_state: &Arc<tokio::sync::RwLock<OrchestratorState>>,
) -> StopActivitySnapshot {
collect_stop_activity_snapshot_within(dynamic_queue, shared_state, REDUCER_SNAPSHOT_TIMEOUT)
.await
}
pub(crate) async fn collect_stop_activity_snapshot_within(
dynamic_queue: &DynamicQueue,
shared_state: &Arc<tokio::sync::RwLock<OrchestratorState>>,
reducer_timeout: std::time::Duration,
) -> StopActivitySnapshot {
let Ok(state) = tokio::time::timeout(reducer_timeout, shared_state.read()).await else {
return StopActivitySnapshot {
execution_handles: ExecutionEvidence::Unavailable,
reducer_agent_execution_active: false,
shutdown_work: ShutdownWorkEvidence::Unavailable,
};
};
let reducer_agent_execution_active = state.is_agent_execution_active();
let shutdown_work = ShutdownWorkEvidence::Known {
pending: state.is_base_mutating_lane_occupied(),
};
drop(state);
let execution_handles = ExecutionEvidence::Known {
registered: dynamic_queue.registered_execution_count().await,
};
StopActivitySnapshot {
execution_handles,
reducer_agent_execution_active,
shutdown_work,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::orchestration::state::OrchestratorState;
use tokio_util::sync::CancellationToken;
fn snapshot(
execution_handles: ExecutionEvidence,
reducer_agent_execution_active: bool,
shutdown_work: ShutdownWorkEvidence,
) -> StopActivitySnapshot {
StopActivitySnapshot {
execution_handles,
reducer_agent_execution_active,
shutdown_work,
}
}
#[test]
fn idle_parallel_stop_registered_handle_selects_force_stop_and_barrier() {
let classification = snapshot(
ExecutionEvidence::Known { registered: 1 },
false,
ShutdownWorkEvidence::Known { pending: false },
)
.classify();
assert_eq!(classification.process_report, ProcessReport::ForceStopped);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[test]
fn idle_parallel_stop_reducer_activity_selects_force_stop_without_handles() {
let classification = snapshot(
ExecutionEvidence::Known { registered: 0 },
true,
ShutdownWorkEvidence::Known { pending: false },
)
.classify();
assert_eq!(classification.process_report, ProcessReport::ForceStopped);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[test]
fn idle_parallel_stop_empty_execution_set_selects_ordinary_stop() {
let classification = snapshot(
ExecutionEvidence::Known { registered: 0 },
false,
ShutdownWorkEvidence::Known { pending: false },
)
.classify();
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert_eq!(
classification.shutdown_barrier,
ShutdownBarrier::NotRequired
);
}
#[test]
fn idle_parallel_stop_pending_merge_keeps_barrier_without_force_stop() {
let classification = snapshot(
ExecutionEvidence::Known { registered: 0 },
false,
ShutdownWorkEvidence::Known { pending: true },
)
.classify();
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[test]
fn idle_parallel_stop_unavailable_execution_evidence_fails_safe() {
let classification = snapshot(
ExecutionEvidence::Unavailable,
false,
ShutdownWorkEvidence::Known { pending: false },
)
.classify();
assert_eq!(classification.process_report, ProcessReport::ForceStopped);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[test]
fn idle_parallel_stop_unavailable_shutdown_evidence_keeps_barrier() {
let classification = snapshot(
ExecutionEvidence::Known { registered: 0 },
false,
ShutdownWorkEvidence::Unavailable,
)
.classify();
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
fn parallel_state(change_ids: Vec<String>) -> Arc<tokio::sync::RwLock<OrchestratorState>> {
Arc::new(tokio::sync::RwLock::new(OrchestratorState::new(
change_ids, 5,
)))
}
#[tokio::test]
async fn idle_parallel_stop_snapshot_reports_registered_execution_handles() {
let queue = DynamicQueue::new();
queue
.register_kill_token("change-a".to_string(), CancellationToken::new())
.await;
let state = parallel_state(vec!["change-a".to_string()]);
let snapshot = collect_stop_activity_snapshot(&queue, &state).await;
assert_eq!(
snapshot.execution_handles,
ExecutionEvidence::Known { registered: 1 }
);
assert_eq!(
snapshot.classify().process_report,
ProcessReport::ForceStopped
);
}
#[tokio::test]
async fn idle_parallel_stop_snapshot_reports_merge_wait_as_ordinary_stop() {
let queue = DynamicQueue::new();
let state = parallel_state(vec!["change-a".to_string()]);
{
let mut guard = state.write().await;
guard.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "change-a".to_string(),
reason: "manual".to_string(),
auto_resumable: false,
});
}
let snapshot = collect_stop_activity_snapshot(&queue, &state).await;
assert_eq!(
snapshot.execution_handles,
ExecutionEvidence::Known { registered: 0 }
);
assert!(!snapshot.reducer_agent_execution_active);
assert_eq!(
snapshot.shutdown_work,
ShutdownWorkEvidence::Known { pending: false }
);
let classification = snapshot.classify();
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert_eq!(
classification.shutdown_barrier,
ShutdownBarrier::NotRequired
);
}
#[tokio::test]
async fn idle_parallel_stop_snapshot_reports_background_merge_as_shutdown_work() {
let queue = DynamicQueue::new();
let state = parallel_state(vec!["change-a".to_string()]);
{
let mut guard = state.write().await;
guard.apply_execution_event(&crate::events::ExecutionEvent::ResolveStarted {
change_id: "change-a".to_string(),
command: "merge".to_string(),
});
}
let snapshot = collect_stop_activity_snapshot(&queue, &state).await;
assert!(!snapshot.reducer_agent_execution_active);
assert_eq!(
snapshot.shutdown_work,
ShutdownWorkEvidence::Known { pending: true }
);
let classification = snapshot.classify();
assert_eq!(classification.process_report, ProcessReport::OrdinaryStop);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[tokio::test]
async fn idle_parallel_stop_snapshot_fails_safe_when_reducer_state_is_unreadable() {
let queue = DynamicQueue::new();
let state = parallel_state(vec!["change-a".to_string()]);
let blocker = state.clone();
let held = blocker.write_owned().await;
let snapshot = collect_stop_activity_snapshot_within(
&queue,
&state,
std::time::Duration::from_millis(5),
)
.await;
drop(held);
assert_eq!(snapshot.execution_handles, ExecutionEvidence::Unavailable);
assert_eq!(snapshot.shutdown_work, ShutdownWorkEvidence::Unavailable);
let classification = snapshot.classify();
assert_eq!(classification.process_report, ProcessReport::ForceStopped);
assert_eq!(classification.shutdown_barrier, ShutdownBarrier::Required);
}
#[tokio::test]
async fn an_idle_scheduler_reports_a_known_zero_rather_than_unavailable() {
let queue = DynamicQueue::new();
let state = Arc::new(tokio::sync::RwLock::new(OrchestratorState::new(
vec!["change-a".to_string()],
5,
)));
let snapshot = collect_stop_activity_snapshot(&queue, &state).await;
assert_eq!(
snapshot.execution_handles,
ExecutionEvidence::Known { registered: 0 }
);
assert_eq!(
snapshot.classify().process_report,
ProcessReport::OrdinaryStop
);
}
}