use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use tokio::sync::RwLock;
use crate::orchestration::operator_command::{
OperatorCommandError, OperatorCommandService, OperatorMode, RetryPlan, RetryRoute,
RunBoundaryLiveness,
};
use crate::orchestration::state::{OrchestratorState, ReduceOutcome, ReducerCommand};
use crate::tui::stop_classification::{StopActivitySnapshot, StopClassification};
const MERGE_WAIT_STATUS: &str = "merge wait";
const NOT_QUEUED_STATUS: &str = "not queued";
const STOPPED_STATUS: &str = "stopped";
const RETRY_DEFERRED_TO_ORDINARY_START: &str =
"retry-class Start selects it only when no ordinary marked change is startable; \
remove the ordinary marks first";
const ORDINARY_DEFERRED_TO_MARK_SETTLEMENT: &str =
"a live run admits ordinary marks through mark settlement rather than Start";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ExcludedTarget {
pub change_id: String,
pub status: String,
pub detail: Option<String>,
}
impl ExcludedTarget {
pub fn new(change_id: impl Into<String>, status: impl Into<String>) -> Self {
Self {
change_id: change_id.into(),
status: status.into(),
detail: None,
}
}
pub fn with_detail(mut self, detail: impl Into<String>) -> Self {
self.detail = Some(detail.into());
self
}
pub fn describe(&self) -> String {
match &self.detail {
Some(detail) => format!("{} ({}): {detail}", self.change_id, self.status),
None => format!("{} ({})", self.change_id, self.status),
}
}
}
fn describe_exclusions(excluded: &[ExcludedTarget]) -> String {
excluded
.iter()
.map(ExcludedTarget::describe)
.collect::<Vec<_>>()
.join(", ")
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RunCommandKind {
Start,
Stop,
CancelStop,
ForceStop,
Retry,
Resolve,
}
impl RunCommandKind {
pub fn as_str(self) -> &'static str {
match self {
Self::Start => "start",
Self::Stop => "stop",
Self::CancelStop => "cancel stop",
Self::ForceStop => "force stop",
Self::Retry => "retry",
Self::Resolve => "resolve",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SchedulerEffect {
Started,
Notified,
None,
}
impl SchedulerEffect {
#[cfg_attr(not(test), allow(dead_code))]
pub fn dispatched(self) -> bool {
matches!(self, Self::Started | Self::Notified)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunNoOpReason {
ResolveAlreadyReserved {
change_id: String,
},
NoRetryableTarget,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunControlOutcome {
RunDispatched {
change_ids: Vec<String>,
explicit_retry: bool,
scheduler: SchedulerEffect,
excluded: Vec<ExcludedTarget>,
},
StopRequested,
StopCancelled,
ForceStopped {
classification: StopClassification,
awaiting_safe_boundary: bool,
},
ResolveReserved {
change_id: String,
reservation: ResolveReservation,
scheduler: SchedulerEffect,
},
NoOp {
reason: RunNoOpReason,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunControlError {
InvalidMode {
command: RunCommandKind,
mode: OperatorMode,
},
NoEligibleTarget {
command: RunCommandKind,
detail: String,
},
TargetIneligible {
command: RunCommandKind,
change_id: String,
display_status: String,
},
DispatchFailed {
command: RunCommandKind,
message: String,
},
Operator(OperatorCommandError),
}
impl std::fmt::Display for RunControlError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::InvalidMode { command, mode } => {
write!(f, "{} is not available in {mode:?} mode", command.as_str())
}
Self::NoEligibleTarget { command, detail } => {
write!(f, "{} has no eligible target: {detail}", command.as_str())
}
Self::TargetIneligible {
command,
change_id,
display_status,
} => write!(
f,
"{} is not available for '{change_id}' with status '{display_status}'",
command.as_str()
),
Self::DispatchFailed { command, message } => {
write!(
f,
"{} could not dispatch the run: {message}",
command.as_str()
)
}
Self::Operator(error) => write!(f, "{error}"),
}
}
}
impl From<OperatorCommandError> for RunControlError {
fn from(error: OperatorCommandError) -> Self {
Self::Operator(error)
}
}
pub type RunControlResult<T> = std::result::Result<T, RunControlError>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ResolveReservation {
Active,
Queued {
position: usize,
},
}
#[derive(Debug, Default)]
pub struct ResolveReservations {
inner: Mutex<ResolveInner>,
}
#[derive(Debug, Default)]
struct ResolveInner {
active: Option<String>,
waiting: VecDeque<String>,
}
impl ResolveReservations {
pub fn new() -> Self {
Self::default()
}
fn lock(&self) -> std::sync::MutexGuard<'_, ResolveInner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn reserve(&self, change_id: &str) -> Option<ResolveReservation> {
let mut guard = self.lock();
if guard.active.as_deref() == Some(change_id)
|| guard.waiting.iter().any(|q| q == change_id)
{
return None;
}
if guard.active.is_none() {
guard.active = Some(change_id.to_string());
return Some(ResolveReservation::Active);
}
guard.waiting.push_back(change_id.to_string());
Some(ResolveReservation::Queued {
position: guard.waiting.len(),
})
}
pub fn is_active(&self) -> bool {
self.lock().active.is_some()
}
pub fn active(&self) -> Option<String> {
self.lock().active.clone()
}
pub fn is_reserved(&self, change_id: &str) -> bool {
let guard = self.lock();
guard.active.as_deref() == Some(change_id)
|| guard.waiting.iter().any(|queued| queued == change_id)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn waiting(&self) -> Vec<String> {
self.lock().waiting.iter().cloned().collect()
}
pub fn has_waiting(&self) -> bool {
!self.lock().waiting.is_empty()
}
pub fn mark_active(&self, change_id: &str) {
let mut guard = self.lock();
guard.waiting.retain(|queued| queued != change_id);
guard.active = Some(change_id.to_string());
}
pub fn finish_active(&self) -> Option<String> {
let mut guard = self.lock();
guard.active = None;
guard.waiting.pop_front()
}
pub fn cancel(&self, change_id: &str) -> bool {
let mut guard = self.lock();
let was_active = guard.active.as_deref() == Some(change_id);
if was_active {
guard.active = None;
}
let before = guard.waiting.len();
guard.waiting.retain(|queued| queued != change_id);
was_active || guard.waiting.len() != before
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn clear(&self) {
let mut guard = self.lock();
guard.active = None;
guard.waiting.clear();
}
}
pub use crate::orchestration::operator_command::ParallelRuntime as StartEligibility;
pub struct RunPermit {
activate: Box<dyn FnOnce() + Send>,
}
impl std::fmt::Debug for RunPermit {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("RunPermit")
}
}
impl RunPermit {
pub fn new(activate: impl FnOnce() + Send + 'static) -> Self {
Self {
activate: Box::new(activate),
}
}
pub fn activate(self) {
(self.activate)();
}
}
#[async_trait]
pub trait RunSchedulerPort: Send + Sync {
fn is_running(&self) -> bool;
async fn prepare_run(
&self,
targets: Vec<String>,
explicit_retry: bool,
) -> Result<RunPermit, String>;
#[cfg_attr(not(test), allow(dead_code))]
async fn start_run(&self, targets: Vec<String>, explicit_retry: bool) -> Result<(), String> {
self.prepare_run(targets, explicit_retry)
.await
.map(RunPermit::activate)
}
async fn notify_scheduler(&self);
async fn cancel_run(&self);
fn set_graceful_stop(&self, requested: bool);
async fn stop_activity(&self) -> StopActivitySnapshot;
}
impl<T: RunSchedulerPort + ?Sized> RunBoundaryLiveness for T {
fn boundary_running(&self) -> bool {
RunSchedulerPort::is_running(self)
}
}
#[derive(Debug)]
pub enum PreparedDispatch {
Start(RunPermit),
Wake,
None,
}
impl PreparedDispatch {
pub fn effect(&self) -> SchedulerEffect {
match self {
Self::Start(_) => SchedulerEffect::Started,
Self::Wake => SchedulerEffect::Notified,
Self::None => SchedulerEffect::None,
}
}
pub async fn activate(self, scheduler: &dyn RunSchedulerPort) {
match self {
Self::Start(permit) => permit.activate(),
Self::Wake => scheduler.notify_scheduler().await,
Self::None => {}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OrdinaryAdmission {
Queue,
ResumeStopped,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct OrdinaryTarget {
change_id: String,
admission: OrdinaryAdmission,
}
impl OrdinaryTarget {
fn command(&self) -> ReducerCommand {
match self.admission {
OrdinaryAdmission::Queue => ReducerCommand::AddToQueue(self.change_id.clone()),
OrdinaryAdmission::ResumeStopped => {
ReducerCommand::ResumeStopped(self.change_id.clone())
}
}
}
}
#[derive(Debug)]
enum StartAdmission {
Ordinary {
targets: Vec<OrdinaryTarget>,
excluded: Vec<ExcludedTarget>,
},
Retry {
routes: Vec<(String, RetryRoute)>,
excluded: Vec<ExcludedTarget>,
},
}
fn exhausted_start_detail(
mode: OperatorMode,
marked: usize,
excluded: &[ExcludedTarget],
) -> String {
let cause = match mode {
OperatorMode::Select => format!(
"no marked change is startable ({marked} marked, none with status \
'{NOT_QUEUED_STATUS}' and none carrying retryable evidence)"
),
OperatorMode::Stopped => format!(
"no marked change is startable ({marked} marked, none with status \
'{NOT_QUEUED_STATUS}', none resumable from an operator '{STOPPED_STATUS}' \
outcome, and none carrying retryable evidence)"
),
OperatorMode::Running => format!(
"a live run admits only marked retry-eligible changes \
({marked} marked, none carries retryable evidence)"
),
_ => format!("no marked change carries retryable evidence ({marked} marked)"),
};
format!("{cause}; excluded: {}", describe_exclusions(excluded))
}
#[derive(Debug)]
enum PreparedIntent {
Start {
targets: Vec<OrdinaryTarget>,
excluded: Vec<ExcludedTarget>,
},
Retry {
routes: Vec<(String, RetryRoute)>,
excluded: Vec<ExcludedTarget>,
},
Resolve { change_id: String },
}
#[derive(Debug)]
pub struct PreparedRunCommand {
intent: PreparedIntent,
dispatch: PreparedDispatch,
}
#[derive(Debug)]
#[must_use = "a committed command holds a scheduler dispatch that must be activated"]
pub struct CommittedRunCommand {
pub outcome: RunControlOutcome,
dispatch: PreparedDispatch,
}
impl CommittedRunCommand {
pub async fn activate(self, scheduler: &dyn RunSchedulerPort) {
self.dispatch.activate(scheduler).await;
}
}
pub struct RunControlService {
state: Arc<RwLock<OrchestratorState>>,
operator: Arc<OperatorCommandService>,
scheduler: Arc<dyn RunSchedulerPort>,
resolves: Arc<ResolveReservations>,
eligibility: Arc<StartEligibility>,
}
impl RunControlService {
pub fn new(
state: Arc<RwLock<OrchestratorState>>,
operator: Arc<OperatorCommandService>,
scheduler: Arc<dyn RunSchedulerPort>,
resolves: Arc<ResolveReservations>,
eligibility: Arc<StartEligibility>,
) -> Self {
Self {
state,
operator,
scheduler,
resolves,
eligibility,
}
}
pub fn operator(&self) -> Arc<OperatorCommandService> {
self.operator.clone()
}
async fn display_status(&self, change_id: &str) -> String {
self.state
.read()
.await
.display_status(change_id)
.to_string()
}
pub fn scheduler(&self) -> Arc<dyn RunSchedulerPort> {
self.scheduler.clone()
}
async fn prepare_dispatch(
&self,
command: RunCommandKind,
targets: Vec<String>,
explicit_retry: bool,
) -> RunControlResult<PreparedDispatch> {
if self.scheduler.is_running() {
return Ok(PreparedDispatch::Wake);
}
self.scheduler
.prepare_run(targets, explicit_retry)
.await
.map(PreparedDispatch::Start)
.map_err(|message| RunControlError::DispatchFailed { command, message })
}
#[cfg_attr(not(test), allow(dead_code))]
async fn dispatch(
&self,
command: RunCommandKind,
targets: Vec<String>,
explicit_retry: bool,
) -> RunControlResult<SchedulerEffect> {
let prepared = self
.prepare_dispatch(command, targets, explicit_retry)
.await?;
let effect = prepared.effect();
prepared.activate(self.scheduler.as_ref()).await;
Ok(effect)
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn start(&self, mode: OperatorMode) -> RunControlResult<RunControlOutcome> {
let prepared = self.prepare_start(mode).await?;
let committed = self.commit(prepared).await?;
let outcome = committed.outcome.clone();
committed.activate(self.scheduler.as_ref()).await;
Ok(outcome)
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn start_targets(&self) -> Vec<String> {
self.classify_start_targets(false)
.await
.0
.into_iter()
.map(|target| target.change_id)
.collect()
}
async fn classify_start_targets(
&self,
resume_stopped: bool,
) -> (Vec<OrdinaryTarget>, Vec<ExcludedTarget>) {
let marked = self.operator.marks().marked_ids();
let guard = self.state.read().await;
let mut startable = Vec::new();
let mut excluded = Vec::new();
for id in marked {
let status = guard.display_status(&id);
let admission = if status == NOT_QUEUED_STATUS {
Some(OrdinaryAdmission::Queue)
} else if resume_stopped && status == STOPPED_STATUS && guard.is_resumable_stopped(&id)
{
Some(OrdinaryAdmission::ResumeStopped)
} else {
None
};
match admission {
Some(admission) => startable.push(OrdinaryTarget {
change_id: id,
admission,
}),
None => {
let status = status.to_string();
excluded.push(ExcludedTarget::new(id, status));
}
}
}
(startable, excluded)
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn retry_change(&self, change_id: &str) -> RunControlResult<RunControlOutcome> {
let plan = self.operator.retry_change(change_id).await?;
self.dispatch_retry(plan, Vec::new()).await
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn retry_errors(&self, change_ids: &[String]) -> RunControlResult<RunControlOutcome> {
let plan = self.operator.retry_errors(change_ids).await;
self.dispatch_retry(plan, Vec::new()).await
}
#[cfg_attr(not(test), allow(dead_code))]
async fn dispatch_retry(
&self,
plan: RetryPlan,
excluded: Vec<ExcludedTarget>,
) -> RunControlResult<RunControlOutcome> {
if plan.is_empty() {
return Ok(RunControlOutcome::NoOp {
reason: RunNoOpReason::NoRetryableTarget,
});
}
let scheduler = self
.dispatch(
RunCommandKind::Retry,
plan.change_ids.clone(),
plan.explicit_retry,
)
.await?;
Ok(RunControlOutcome::RunDispatched {
change_ids: plan.change_ids,
explicit_retry: plan.explicit_retry,
scheduler,
excluded,
})
}
pub async fn stop(&self, mode: OperatorMode) -> RunControlResult<RunControlOutcome> {
if mode != OperatorMode::Running && !self.is_persistent_idle_ready(mode) {
return Err(RunControlError::InvalidMode {
command: RunCommandKind::Stop,
mode,
});
}
self.scheduler.set_graceful_stop(true);
self.scheduler.notify_scheduler().await;
Ok(RunControlOutcome::StopRequested)
}
fn is_persistent_idle_ready(&self, mode: OperatorMode) -> bool {
mode == OperatorMode::Select && self.scheduler.is_running()
}
pub async fn cancel_stop(&self, mode: OperatorMode) -> RunControlResult<RunControlOutcome> {
if mode != OperatorMode::Stopping {
return Err(RunControlError::InvalidMode {
command: RunCommandKind::CancelStop,
mode,
});
}
self.scheduler.set_graceful_stop(false);
Ok(RunControlOutcome::StopCancelled)
}
pub async fn force_stop(&self, mode: OperatorMode) -> RunControlResult<RunControlOutcome> {
if !matches!(mode, OperatorMode::Running | OperatorMode::Stopping)
&& !self.is_persistent_idle_ready(mode)
{
return Err(RunControlError::InvalidMode {
command: RunCommandKind::ForceStop,
mode,
});
}
let snapshot = self.scheduler.stop_activity().await;
let classification = snapshot.classify();
let scheduler_running = self.scheduler.is_running();
self.scheduler.cancel_run().await;
Ok(RunControlOutcome::ForceStopped {
classification,
awaiting_safe_boundary: scheduler_running
&& snapshot.scheduler_owns_cleanup()
&& classification.shutdown_barrier.is_required(),
})
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn resolve_merge(&self, change_id: &str) -> RunControlResult<RunControlOutcome> {
if self.resolves.is_reserved(change_id) {
return Ok(RunControlOutcome::NoOp {
reason: RunNoOpReason::ResolveAlreadyReserved {
change_id: change_id.to_string(),
},
});
}
let display_status = self.display_status(change_id).await;
if display_status != MERGE_WAIT_STATUS {
return Err(RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
change_id: change_id.to_string(),
display_status,
});
}
let reduce_outcome = {
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::ResolveMerge(change_id.to_string()))
};
if matches!(reduce_outcome, ReduceOutcome::NoOp) {
return Err(RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
change_id: change_id.to_string(),
display_status,
});
}
let Some(reservation) = self.resolves.reserve(change_id) else {
return Ok(RunControlOutcome::NoOp {
reason: RunNoOpReason::ResolveAlreadyReserved {
change_id: change_id.to_string(),
},
});
};
let scheduler = match reservation {
ResolveReservation::Active => {
self.dispatch(RunCommandKind::Resolve, Vec::new(), false)
.await?
}
ResolveReservation::Queued { .. } => SchedulerEffect::None,
};
Ok(RunControlOutcome::ResolveReserved {
change_id: change_id.to_string(),
reservation,
scheduler,
})
}
pub async fn prepare_start(&self, mode: OperatorMode) -> RunControlResult<PreparedRunCommand> {
match self.classify_start_admission(mode).await? {
StartAdmission::Ordinary { targets, excluded } => {
let change_ids: Vec<String> = targets
.iter()
.map(|target| target.change_id.clone())
.collect();
let dispatch = self
.prepare_dispatch(RunCommandKind::Start, change_ids, false)
.await?;
Ok(PreparedRunCommand {
intent: PreparedIntent::Start { targets, excluded },
dispatch,
})
}
StartAdmission::Retry { routes, excluded } => {
self.prepare_routes(RunCommandKind::Start, routes, excluded)
.await
}
}
}
async fn classify_start_admission(
&self,
mode: OperatorMode,
) -> RunControlResult<StartAdmission> {
if mode == OperatorMode::Stopping {
return Err(RunControlError::InvalidMode {
command: RunCommandKind::Start,
mode,
});
}
let marked = self.operator.marks().marked_ids();
if marked.is_empty() {
return Err(RunControlError::NoEligibleTarget {
command: RunCommandKind::Start,
detail: "no change carries an execution mark".to_string(),
});
}
let rejected = self.eligibility.rejected(&marked);
if !rejected.is_empty() {
return Err(RunControlError::NoEligibleTarget {
command: RunCommandKind::Start,
detail: format!(
"worktree execution requires committed changes with no uncommitted \
files; ineligible marked targets: {}",
rejected.join(", ")
),
});
}
let admits_ordinary = matches!(mode, OperatorMode::Select | OperatorMode::Stopped);
let resume_stopped = mode == OperatorMode::Stopped;
let (ordinary, ordinary_excluded) = if admits_ordinary {
self.classify_start_targets(resume_stopped).await
} else {
(Vec::new(), Vec::new())
};
if !ordinary.is_empty() {
let excluded = self.defer_retry_only(ordinary_excluded).await;
return Ok(StartAdmission::Ordinary {
targets: ordinary,
excluded,
});
}
let routes = self.operator.plan_retry_errors(&marked).await;
let excluded = self.describe_non_retryable(&marked, &routes, mode).await;
if routes.is_empty() {
return Err(RunControlError::NoEligibleTarget {
command: RunCommandKind::Start,
detail: exhausted_start_detail(mode, marked.len(), &excluded),
});
}
Ok(StartAdmission::Retry { routes, excluded })
}
async fn defer_retry_only(&self, excluded: Vec<ExcludedTarget>) -> Vec<ExcludedTarget> {
let ids: Vec<String> = excluded
.iter()
.map(|target| target.change_id.clone())
.collect();
if ids.is_empty() {
return excluded;
}
let routes = self.operator.plan_retry_errors(&ids).await;
excluded
.into_iter()
.map(|target| {
if routes.iter().any(|(id, _)| *id == target.change_id) {
target.with_detail(RETRY_DEFERRED_TO_ORDINARY_START)
} else {
target
}
})
.collect()
}
async fn describe_non_retryable(
&self,
marked: &[String],
routes: &[(String, RetryRoute)],
mode: OperatorMode,
) -> Vec<ExcludedTarget> {
let guard = self.state.read().await;
marked
.iter()
.filter(|id| !routes.iter().any(|(routed, _)| routed == *id))
.map(|id| {
let status = guard.display_status(id).to_string();
let deferred = mode == OperatorMode::Running && status == NOT_QUEUED_STATUS;
let target = ExcludedTarget::new(id.clone(), status);
if deferred {
target.with_detail(ORDINARY_DEFERRED_TO_MARK_SETTLEMENT)
} else {
target
}
})
.collect()
}
pub async fn prepare_retry_change(
&self,
change_id: &str,
) -> RunControlResult<PreparedRunCommand> {
let routes = match self.operator.plan_retry_change(change_id).await? {
Some(route) => vec![(change_id.to_string(), route)],
None => Vec::new(),
};
self.prepare_routes(RunCommandKind::Retry, routes, Vec::new())
.await
}
pub async fn prepare_retry_errors(
&self,
change_ids: &[String],
) -> RunControlResult<PreparedRunCommand> {
let routes = self.operator.plan_retry_errors(change_ids).await;
self.prepare_routes(RunCommandKind::Retry, routes, Vec::new())
.await
}
async fn prepare_routes(
&self,
command: RunCommandKind,
routes: Vec<(String, RetryRoute)>,
excluded: Vec<ExcludedTarget>,
) -> RunControlResult<PreparedRunCommand> {
if routes.is_empty() {
return Ok(PreparedRunCommand {
intent: PreparedIntent::Retry { routes, excluded },
dispatch: PreparedDispatch::None,
});
}
let targets: Vec<String> = routes.iter().map(|(id, _)| id.clone()).collect();
let dispatch = self.prepare_dispatch(command, targets, true).await?;
Ok(PreparedRunCommand {
intent: PreparedIntent::Retry { routes, excluded },
dispatch,
})
}
pub async fn prepare_resolve(
&self,
change_id: &str,
) -> RunControlResult<Option<PreparedRunCommand>> {
if self.resolves.is_reserved(change_id) {
return Ok(None);
}
let display_status = self.display_status(change_id).await;
if display_status != MERGE_WAIT_STATUS {
return Err(RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
change_id: change_id.to_string(),
display_status,
});
}
let dispatch = if self.resolves.is_active() {
PreparedDispatch::None
} else {
self.prepare_dispatch(RunCommandKind::Resolve, Vec::new(), false)
.await?
};
Ok(Some(PreparedRunCommand {
intent: PreparedIntent::Resolve {
change_id: change_id.to_string(),
},
dispatch,
}))
}
pub async fn commit(
&self,
prepared: PreparedRunCommand,
) -> RunControlResult<CommittedRunCommand> {
let PreparedRunCommand { intent, dispatch } = prepared;
let scheduler = dispatch.effect();
match intent {
PreparedIntent::Start { targets, excluded } => {
{
let mut guard = self.state.write().await;
for target in &targets {
guard.apply_command(target.command());
}
}
Ok(CommittedRunCommand {
outcome: RunControlOutcome::RunDispatched {
change_ids: targets.into_iter().map(|target| target.change_id).collect(),
explicit_retry: false,
scheduler,
excluded,
},
dispatch,
})
}
PreparedIntent::Retry { routes, excluded } => {
let plan = self.operator.commit_retry_routes(&routes).await;
if plan.is_empty() {
return Ok(CommittedRunCommand {
outcome: RunControlOutcome::NoOp {
reason: RunNoOpReason::NoRetryableTarget,
},
dispatch: PreparedDispatch::None,
});
}
Ok(CommittedRunCommand {
outcome: RunControlOutcome::RunDispatched {
change_ids: plan.change_ids,
explicit_retry: plan.explicit_retry,
scheduler,
excluded,
},
dispatch,
})
}
PreparedIntent::Resolve { change_id } => {
let reduce_outcome = {
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::ResolveMerge(change_id.clone()))
};
if matches!(reduce_outcome, ReduceOutcome::NoOp) {
return Err(RunControlError::TargetIneligible {
command: RunCommandKind::Resolve,
change_id: change_id.clone(),
display_status: self.display_status(&change_id).await,
});
}
let Some(reservation) = self.resolves.reserve(&change_id) else {
return Ok(CommittedRunCommand {
outcome: RunControlOutcome::NoOp {
reason: RunNoOpReason::ResolveAlreadyReserved { change_id },
},
dispatch: PreparedDispatch::None,
});
};
let (dispatch, scheduler) = match reservation {
ResolveReservation::Active => (dispatch, scheduler),
ResolveReservation::Queued { .. } => {
(PreparedDispatch::None, SchedulerEffect::None)
}
};
Ok(CommittedRunCommand {
outcome: RunControlOutcome::ResolveReserved {
change_id,
reservation,
scheduler,
},
dispatch,
})
}
}
}
}
#[cfg(test)]
pub(crate) mod testing {
use super::*;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum SchedulerCall {
Started {
targets: Vec<String>,
explicit_retry: bool,
},
Notified,
Cancelled,
GracefulStop(bool),
}
type ActivationHook = Arc<dyn Fn(Vec<String>, bool) + Send + Sync>;
#[derive(Default)]
struct SchedulerRecorder {
calls: Mutex<Vec<SchedulerCall>>,
running: std::sync::atomic::AtomicBool,
graceful_stop: Arc<std::sync::atomic::AtomicBool>,
on_activate: Mutex<Option<ActivationHook>>,
}
impl std::fmt::Debug for SchedulerRecorder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SchedulerRecorder")
.field("calls", &self.calls)
.field("running", &self.running)
.finish_non_exhaustive()
}
}
impl SchedulerRecorder {
fn record(&self, call: SchedulerCall) {
self.calls.lock().unwrap().push(call);
}
}
#[derive(Debug)]
pub(crate) struct RecordingScheduler {
recorder: Arc<SchedulerRecorder>,
activity: Mutex<StopActivitySnapshot>,
launch_failure: Mutex<Option<String>>,
}
impl Default for RecordingScheduler {
fn default() -> Self {
Self::new()
}
}
impl RecordingScheduler {
pub(crate) fn new() -> Self {
use crate::tui::stop_classification::{ExecutionEvidence, ShutdownWorkEvidence};
Self {
recorder: Arc::new(SchedulerRecorder::default()),
activity: Mutex::new(StopActivitySnapshot {
execution_handles: ExecutionEvidence::Known { registered: 0 },
reducer_agent_execution_active: false,
shutdown_work: ShutdownWorkEvidence::Known { pending: false },
}),
launch_failure: Mutex::new(None),
}
}
pub(crate) fn set_running(&self, running: bool) {
self.recorder
.running
.store(running, std::sync::atomic::Ordering::SeqCst);
}
pub(crate) fn set_activity(&self, activity: StopActivitySnapshot) {
*self.activity.lock().unwrap() = activity;
}
pub(crate) fn fail_launch(&self, message: &str) {
*self.launch_failure.lock().unwrap() = Some(message.to_string());
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn on_activate(&self, hook: ActivationHook) {
*self.recorder.on_activate.lock().unwrap() = Some(hook);
}
pub(crate) fn graceful_stop_flag(&self) -> Arc<std::sync::atomic::AtomicBool> {
self.recorder.graceful_stop.clone()
}
pub(crate) fn calls(&self) -> Vec<SchedulerCall> {
self.recorder.calls.lock().unwrap().clone()
}
pub(crate) fn started_targets(&self) -> Vec<Vec<String>> {
self.calls()
.into_iter()
.filter_map(|call| match call {
SchedulerCall::Started { targets, .. } => Some(targets),
_ => None,
})
.collect()
}
}
#[async_trait]
impl RunSchedulerPort for RecordingScheduler {
fn is_running(&self) -> bool {
self.recorder
.running
.load(std::sync::atomic::Ordering::SeqCst)
}
async fn prepare_run(
&self,
targets: Vec<String>,
explicit_retry: bool,
) -> std::result::Result<RunPermit, String> {
if let Some(message) = self.launch_failure.lock().unwrap().take() {
return Err(message);
}
let recorder = self.recorder.clone();
Ok(RunPermit::new(move || {
recorder.record(SchedulerCall::Started {
targets: targets.clone(),
explicit_retry,
});
recorder
.running
.store(true, std::sync::atomic::Ordering::SeqCst);
let hook = recorder.on_activate.lock().unwrap().clone();
if let Some(hook) = hook {
hook(targets, explicit_retry);
}
}))
}
async fn notify_scheduler(&self) {
self.recorder.record(SchedulerCall::Notified);
}
async fn cancel_run(&self) {
self.recorder.record(SchedulerCall::Cancelled);
}
fn set_graceful_stop(&self, requested: bool) {
self.recorder.record(SchedulerCall::GracefulStop(requested));
self.recorder
.graceful_stop
.store(requested, std::sync::atomic::Ordering::SeqCst);
}
async fn stop_activity(&self) -> StopActivitySnapshot {
*self.activity.lock().unwrap()
}
}
}
#[cfg(test)]
mod tests;