#![allow(dead_code)]
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::RwLock;
use tokio_util::sync::CancellationToken;
use crate::orchestration::apply_commit_evidence::ApplyCommitEvidence;
use crate::orchestration::execution_facts::{project_phase, ExecutionFactsStore, ExecutionPhase};
use crate::orchestration::mark_settlement::{
classify_mark_settlement_row, plan_mark_settlement, MarkSettlementAction,
MarkSettlementCoordinator, MarkSettlementExclusion, MarkSettlementPlan, MarkSettlementRow,
};
use crate::orchestration::state::{OrchestratorState, ReduceOutcome, ReducerCommand};
pub const DEFAULT_CANCELLATION_TIMEOUT: Duration = Duration::from_secs(30);
pub const ACTIVE_STATUSES: [&str; 6] = [
"preparing",
"applying",
"accepting",
"rejecting",
"archiving",
"resolving",
];
pub const COMPLETED_STATUSES: [&str; 3] = ["archived", "merged", "pushed"];
const FINAL_STATUSES: [&str; 4] = ["archived", "merged", "pushed", "rejected"];
const MARK_ONLY_WAIT_STATUSES: [&str; 2] = ["merge wait", "resolve pending"];
pub fn is_active_status(display_status: &str) -> bool {
ACTIVE_STATUSES.contains(&display_status)
}
pub fn is_completed_status(display_status: &str) -> bool {
COMPLETED_STATUSES.contains(&display_status)
}
pub fn is_final_status(display_status: &str) -> bool {
FINAL_STATUSES.contains(&display_status)
}
pub trait RunBoundaryLiveness: Send + Sync {
fn boundary_running(&self) -> bool;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OperatorMode {
Select,
Running,
Stopping,
Stopped,
Error,
}
impl OperatorMode {
pub fn from_app_mode(app_mode: &str) -> Self {
match app_mode {
"running" => Self::Running,
"stopping" => Self::Stopping,
"stopped" => Self::Stopped,
"error" => Self::Error,
_ => Self::Select,
}
}
pub fn as_app_mode(self) -> &'static str {
match self {
Self::Select => "select",
Self::Running => "running",
Self::Stopping => "stopping",
Self::Stopped => "stopped",
Self::Error => "error",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MarkAdmission {
Allowed,
TerminalTarget,
ArchiveComplete,
}
impl MarkAdmission {
pub fn is_allowed(self) -> bool {
matches!(self, Self::Allowed)
}
}
pub fn classify_mark_admission(display_status: &str, archive_complete: bool) -> MarkAdmission {
if is_final_status(display_status) {
MarkAdmission::TerminalTarget
} else if archive_complete {
MarkAdmission::ArchiveComplete
} else {
MarkAdmission::Allowed
}
}
pub fn is_markable_status(display_status: &str, archive_complete: bool) -> bool {
classify_mark_admission(display_status, archive_complete).is_allowed()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueIntentRoute {
NoQueueEffect,
Mutable,
RetryRequired,
Immutable,
}
pub fn classify_queue_intent_route(mode: OperatorMode, display_status: &str) -> QueueIntentRoute {
if is_final_status(display_status) {
return QueueIntentRoute::Immutable;
}
match mode {
OperatorMode::Error => QueueIntentRoute::RetryRequired,
OperatorMode::Select => QueueIntentRoute::NoQueueEffect,
OperatorMode::Stopping => QueueIntentRoute::Immutable,
OperatorMode::Stopped => {
if matches!(display_status, "not queued" | "error")
|| MARK_ONLY_WAIT_STATUSES.contains(&display_status)
{
QueueIntentRoute::NoQueueEffect
} else {
QueueIntentRoute::Immutable
}
}
OperatorMode::Running => {
if MARK_ONLY_WAIT_STATUSES.contains(&display_status) {
return QueueIntentRoute::NoQueueEffect;
}
if is_active_status(display_status) {
return QueueIntentRoute::Immutable;
}
match display_status {
"not queued" | "queued" | "error" => QueueIntentRoute::Mutable,
_ => QueueIntentRoute::Immutable,
}
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum ParallelEligibility {
#[default]
Eligible,
ProposalAbsentFromHead,
UncommittedProposalFiles,
}
impl ParallelEligibility {
pub fn observe(
change_id: &str,
committed_change_ids: &HashSet<String>,
uncommitted_file_change_ids: &HashSet<String>,
) -> Self {
if uncommitted_file_change_ids.contains(change_id) {
Self::UncommittedProposalFiles
} else if !committed_change_ids.contains(change_id) {
Self::ProposalAbsentFromHead
} else {
Self::Eligible
}
}
pub fn is_eligible(self) -> bool {
matches!(self, Self::Eligible)
}
pub fn has_uncommitted_proposal_files(self) -> bool {
matches!(self, Self::UncommittedProposalFiles)
}
pub fn queue_exclusion(self) -> Option<MarkExclusion> {
match self {
Self::Eligible => None,
Self::ProposalAbsentFromHead => Some(MarkExclusion::ParallelProposalAbsent),
Self::UncommittedProposalFiles => Some(MarkExclusion::ParallelIneligible),
}
}
}
pub const PARALLEL_INELIGIBLE_CLEANUP_REASON: &str = "not eligible for worktree execution";
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum MarkExclusion {
FinalStatus,
RetryRequired,
StopPending,
ChangeActive,
StatusImmutable,
ParallelIneligible,
ParallelProposalAbsent,
ArchiveComplete,
}
impl MarkExclusion {
pub const ALL: [MarkExclusion; 8] = [
MarkExclusion::ChangeActive,
MarkExclusion::ParallelIneligible,
MarkExclusion::ParallelProposalAbsent,
MarkExclusion::ArchiveComplete,
MarkExclusion::FinalStatus,
MarkExclusion::RetryRequired,
MarkExclusion::StopPending,
MarkExclusion::StatusImmutable,
];
pub fn as_str(self) -> &'static str {
match self {
Self::ArchiveComplete => "archive_complete",
Self::FinalStatus => "final_status",
Self::RetryRequired => "retry_required",
Self::StopPending => "stop_pending",
Self::ChangeActive => "change_active",
Self::StatusImmutable => "status_immutable",
Self::ParallelIneligible => "parallel_ineligible",
Self::ParallelProposalAbsent => "parallel_proposal_absent",
}
}
pub fn reason(self) -> &'static str {
match self {
Self::ArchiveComplete => "archive complete (no next run)",
Self::FinalStatus => "final or rejected and read-only",
Self::RetryRequired => "in error mode (use retry)",
Self::StopPending => "waiting for the pending stop",
Self::ChangeActive => "in progress (use K to stop)",
Self::StatusImmutable => "not mutable in this mode",
Self::ParallelIneligible => "uncommitted (commit first)",
Self::ParallelProposalAbsent => "not present in HEAD (cannot queue)",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MarkTargetRow<'a> {
pub change_id: &'a str,
pub display_status: &'a str,
pub archive_complete: bool,
pub marked: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BulkMarkPlan {
pub target_state: bool,
pub eligible: Vec<String>,
pub excluded: Vec<(String, MarkExclusion)>,
}
impl BulkMarkPlan {
pub fn is_empty(&self) -> bool {
self.eligible.is_empty()
}
pub fn exclusion_summary(&self) -> String {
MarkExclusion::ALL
.iter()
.filter_map(|reason| {
let count = self
.excluded
.iter()
.filter(|(_, actual)| actual == reason)
.count();
(count > 0).then(|| format!("{} {}", count, reason.reason()))
})
.collect::<Vec<_>>()
.join(", ")
}
}
pub fn classify_bulk_mark_row(
display_status: &str,
archive_complete: bool,
) -> Option<MarkExclusion> {
match classify_mark_admission(display_status, archive_complete) {
MarkAdmission::Allowed => None,
MarkAdmission::TerminalTarget => Some(MarkExclusion::FinalStatus),
MarkAdmission::ArchiveComplete => Some(MarkExclusion::ArchiveComplete),
}
}
pub fn plan_bulk_marks(rows: &[MarkTargetRow<'_>]) -> BulkMarkPlan {
let mut eligible = Vec::new();
let mut excluded = Vec::new();
let mut any_unmarked = false;
for row in rows {
match classify_bulk_mark_row(row.display_status, row.archive_complete) {
Some(reason) => excluded.push((row.change_id.to_string(), reason)),
None => {
any_unmarked |= !row.marked;
eligible.push(row.change_id.to_string());
}
}
}
BulkMarkPlan {
target_state: any_unmarked,
eligible,
excluded,
}
}
pub fn parallel_cleanup_targets(rows: &[ParallelCleanupRow<'_>]) -> Vec<String> {
rows.iter()
.filter(|row| !row.parallel_eligible && (row.marked || row.queued))
.map(|row| row.change_id.to_string())
.collect()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ParallelCleanupRow<'a> {
pub change_id: &'a str,
pub parallel_eligible: bool,
pub marked: bool,
pub queued: bool,
}
#[derive(Debug, Default)]
pub struct ParallelRuntime {
inner: Mutex<ParallelRuntimeInner>,
mutations: tokio::sync::Mutex<()>,
}
#[derive(Debug, Default)]
struct ParallelRuntimeInner {
max_concurrent: usize,
vcs_backend: String,
ineligible: HashMap<String, ParallelEligibility>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParallelRuntimeFacts {
pub max_concurrent: usize,
pub vcs_backend: String,
}
impl ParallelRuntime {
pub fn new() -> Self {
Self::default()
}
fn lock(&self) -> std::sync::MutexGuard<'_, ParallelRuntimeInner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn set_max_concurrent(&self, max_concurrent: usize) {
self.lock().max_concurrent = max_concurrent;
}
pub fn set_vcs_backend(&self, backend: impl Into<String>) {
self.lock().vcs_backend = backend.into();
}
pub fn set_parallel_ineligible(
&self,
entries: impl IntoIterator<Item = (String, ParallelEligibility)>,
) {
self.lock().ineligible = entries
.into_iter()
.filter(|(_, eligibility)| !eligibility.is_eligible())
.collect();
}
pub fn is_eligible(&self, change_id: &str) -> bool {
!self.lock().ineligible.contains_key(change_id)
}
pub fn eligibility(&self, change_id: &str) -> ParallelEligibility {
self.lock()
.ineligible
.get(change_id)
.copied()
.unwrap_or_default()
}
pub fn ineligible_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self.lock().ineligible.keys().cloned().collect();
ids.sort();
ids
}
pub fn facts(&self) -> ParallelRuntimeFacts {
let guard = self.lock();
ParallelRuntimeFacts {
max_concurrent: guard.max_concurrent,
vcs_backend: guard.vcs_backend.clone(),
}
}
pub async fn lock_mutations(&self) -> tokio::sync::MutexGuard<'_, ()> {
self.mutations.lock().await
}
pub fn rejected(&self, targets: &[String]) -> Vec<String> {
let guard = self.lock();
targets
.iter()
.filter(|id| guard.ineligible.contains_key(*id))
.cloned()
.collect()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RetryRoute {
TerminalError,
AcceptanceStall,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RetryEdgeAuthority {
AnalysisBypass,
TerminalError,
}
impl RetryEdgeAuthority {
pub fn releases_terminal_error_state(self) -> bool {
matches!(self, Self::TerminalError)
}
pub fn widen(self, other: Self) -> Self {
match (self, other) {
(Self::TerminalError, _) | (_, Self::TerminalError) => Self::TerminalError,
_ => Self::AnalysisBypass,
}
}
}
impl From<RetryRoute> for RetryEdgeAuthority {
fn from(route: RetryRoute) -> Self {
match route {
RetryRoute::TerminalError => Self::TerminalError,
RetryRoute::AcceptanceStall => Self::AnalysisBypass,
}
}
}
pub fn classify_retry_route(
display_status: &str,
blocker_kind: crate::orchestration::state::BlockerKind,
) -> Option<RetryRoute> {
use crate::orchestration::state::BlockerKind;
match (display_status, blocker_kind) {
("error", _) => Some(RetryRoute::TerminalError),
("stalled", _) => Some(RetryRoute::AcceptanceStall),
("blocked", BlockerKind::External) => Some(RetryRoute::AcceptanceStall),
_ => None,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperatorCommand {
SetExecutionMark {
change_id: String,
marked: bool,
},
AddToQueue {
change_id: String,
},
RemoveFromQueue {
change_id: String,
},
StopAndDequeue {
change_id: String,
},
ForceStopChange {
change_id: String,
},
RetryChange {
change_id: String,
},
SetAllExecutionMarks,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ForceStopExclusion {
UnknownTarget,
TerminalTarget,
NotAdmitted,
MergeWait,
ResolveWait,
NoLiveProcess,
}
impl ForceStopExclusion {
pub const ALL: [ForceStopExclusion; 6] = [
ForceStopExclusion::UnknownTarget,
ForceStopExclusion::TerminalTarget,
ForceStopExclusion::NotAdmitted,
ForceStopExclusion::MergeWait,
ForceStopExclusion::ResolveWait,
ForceStopExclusion::NoLiveProcess,
];
pub fn as_str(self) -> &'static str {
match self {
Self::UnknownTarget => "unknown_target",
Self::TerminalTarget => "terminal_target",
Self::NotAdmitted => "not_admitted",
Self::MergeWait => "merge_wait",
Self::ResolveWait => "resolve_wait",
Self::NoLiveProcess => "no_live_process",
}
}
pub fn reason(self) -> &'static str {
match self {
Self::UnknownTarget => "not tracked by this owner",
Self::TerminalTarget => "final or rejected and already stopped",
Self::NotAdmitted => "not admitted (nothing is running for it)",
Self::MergeWait => "waiting on a base merge (no process to kill)",
Self::ResolveWait => "waiting for manual resolution (no live resolver)",
Self::NoLiveProcess => "owns no live managed process (use stop instead)",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ForceStopAdmission {
KillAndDequeue,
DequeueOnly,
Refused(ForceStopExclusion),
}
impl ForceStopAdmission {
pub fn is_allowed(self) -> bool {
!matches!(self, Self::Refused(_))
}
pub fn exclusion(self) -> Option<ForceStopExclusion> {
match self {
Self::Refused(reason) => Some(reason),
_ => None,
}
}
}
pub fn classify_force_stop_change(
display_status: &str,
tracked: bool,
owns_managed_process: bool,
) -> ForceStopAdmission {
use ForceStopAdmission::{DequeueOnly, KillAndDequeue, Refused};
if !tracked {
return Refused(ForceStopExclusion::UnknownTarget);
}
if is_final_status(display_status) {
return Refused(ForceStopExclusion::TerminalTarget);
}
match display_status {
"merge wait" => Refused(ForceStopExclusion::MergeWait),
"resolve pending" => Refused(ForceStopExclusion::ResolveWait),
status if is_active_status(status) => {
if owns_managed_process {
KillAndDequeue
} else {
Refused(ForceStopExclusion::NoLiveProcess)
}
}
"queued" | "blocked" => {
if owns_managed_process {
KillAndDequeue
} else {
DequeueOnly
}
}
_ => Refused(ForceStopExclusion::NotAdmitted),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueMutation {
Added,
Removed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueueOutcome {
pub change_id: String,
pub mutation: QueueMutation,
pub reducer_changed: bool,
pub dynamic_queue_mutated: bool,
pub display_status: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SettlementApplication {
pub outcome: QueueOutcome,
pub skipped: Option<MarkSettlementExclusion>,
}
impl SettlementApplication {
pub fn applied(&self) -> bool {
self.outcome.reducer_changed || self.outcome.dynamic_queue_mutated
}
}
fn queue_mutation_for(action: MarkSettlementAction) -> QueueMutation {
match action {
MarkSettlementAction::Add => QueueMutation::Added,
MarkSettlementAction::Remove => QueueMutation::Removed,
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct StopSettlement {
pub cancelled_phase: ExecutionPhase,
pub last_completed_phase: Option<ExecutionPhase>,
pub apply_commit_present: Option<bool>,
pub apply_commit_oid: Option<String>,
}
impl StopSettlement {
pub fn none() -> Self {
Self {
cancelled_phase: ExecutionPhase::None,
last_completed_phase: None,
apply_commit_present: None,
apply_commit_oid: None,
}
}
pub fn describe(&self, change_id: &str) -> String {
let what = match self.cancelled_phase {
ExecutionPhase::None => {
format!("'{change_id}' was already terminated and was dequeued")
}
ExecutionPhase::Unknown => format!(
"'{change_id}' was dequeued; the phase it was cancelled during could not be \
determined"
),
phase => format!(
"'{change_id}' was cancelled during {} and dequeued",
phase.as_str()
),
};
let apply = match (self.apply_commit_present, self.apply_commit_oid.as_deref()) {
(Some(true), Some(oid)) => {
format!("; the final Apply commit {oid} was already created")
}
(Some(true), None) => "; the final Apply commit was already created".to_string(),
_ => "; whether the final Apply commit exists could not be proven".to_string(),
};
format!("{what}{apply}; previously completed worktree effects were not rolled back")
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NoOpReason {
MarkUnchanged,
TerminalMarkTarget,
ArchiveCompleteMarkTarget,
ReducerRejected,
BulkMarksUnchanged,
NoEligibleMarkTarget,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperatorOutcome {
MarkSet {
change_id: String,
marked: bool,
},
Queue(QueueOutcome),
Dequeued {
change_id: String,
settlement: StopSettlement,
},
ForceStopped {
change_id: String,
execution_id: Option<String>,
terminated: bool,
settlement: StopSettlement,
},
Retry(RetryPlan),
BulkMarks {
marked: bool,
changed: Vec<String>,
excluded: Vec<(String, MarkExclusion)>,
},
NoOp {
change_id: String,
reason: NoOpReason,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RetryPlan {
pub change_ids: Vec<String>,
pub routes: Vec<RetryRoute>,
pub explicit_retry: bool,
}
impl RetryPlan {
pub fn is_empty(&self) -> bool {
self.change_ids.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperatorCommandError {
MissingCancellationHandle {
change_id: String,
},
CancellationFailed {
change_id: String,
message: String,
},
TerminationTimeout {
change_id: String,
waited: Duration,
},
RetryUnsupported {
change_id: String,
display_status: String,
},
ForceStopIneligible {
change_id: String,
display_status: String,
reason: ForceStopExclusion,
},
ForceStopUnconfirmed {
change_id: String,
detail: String,
},
UnchangedAcceptanceInput {
change_id: String,
category: String,
fingerprint: String,
components: Vec<String>,
},
}
impl OperatorCommandError {
pub fn change_id(&self) -> &str {
match self {
Self::MissingCancellationHandle { change_id }
| Self::CancellationFailed { change_id, .. }
| Self::TerminationTimeout { change_id, .. }
| Self::RetryUnsupported { change_id, .. }
| Self::ForceStopIneligible { change_id, .. }
| Self::ForceStopUnconfirmed { change_id, .. }
| Self::UnchangedAcceptanceInput { change_id, .. } => change_id,
}
}
pub fn outcome_token(&self) -> Option<&'static str> {
match self {
Self::UnchangedAcceptanceInput { .. } => Some(
crate::orchestration::acceptance::execution_manifest::UNCHANGED_ACCEPTANCE_INPUT,
),
_ => None,
}
}
}
impl std::fmt::Display for OperatorCommandError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::MissingCancellationHandle { change_id } => write!(
f,
"no cancellation handle registered for active change '{change_id}'; \
the stop request is recorded and takes effect before the next operation starts"
),
Self::CancellationFailed { change_id, message } => {
write!(f, "cancellation failed for '{change_id}': {message}")
}
Self::TerminationTimeout { change_id, waited } => write!(
f,
"termination of '{change_id}' was not confirmed within {waited:?}"
),
Self::RetryUnsupported {
change_id,
display_status,
} => write!(
f,
"retry is not supported for '{change_id}' with status '{display_status}'"
),
Self::ForceStopIneligible {
change_id,
display_status,
reason,
} => write!(
f,
"'{change_id}' cannot be force-stopped with status '{display_status}': {}",
reason.reason()
),
Self::ForceStopUnconfirmed { change_id, detail } => write!(
f,
"the managed process group of '{change_id}' was not proven empty after SIGKILL, \
so no dequeue was committed: {detail}"
),
Self::UnchangedAcceptanceInput {
change_id,
category,
fingerprint,
components,
} => write!(
f,
"retry of '{change_id}' is refused as `{}`: its Acceptance input is unchanged \
since the {category} hold (fingerprint {}). No analysis, Apply, gate, or \
Acceptance work was dispatched. A repository-visible change to one of {} \
restores retry eligibility.",
crate::orchestration::acceptance::execution_manifest::UNCHANGED_ACCEPTANCE_INPUT,
&fingerprint[..16.min(fingerprint.len())],
components.join(", ")
),
}
}
}
pub type OperatorResult<T> = std::result::Result<T, OperatorCommandError>;
#[async_trait::async_trait]
pub trait AcceptanceAdmissionPort: Send + Sync {
async fn classify(
&self,
change_id: &str,
) -> crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct AlwaysAdmitAcceptance;
#[async_trait::async_trait]
impl AcceptanceAdmissionPort for AlwaysAdmitAcceptance {
async fn classify(
&self,
_change_id: &str,
) -> crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission {
crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission::Admit
}
}
#[derive(Debug, Default)]
pub struct ExecutionMarkStore {
marks: Mutex<HashSet<String>>,
settlement: Arc<MarkSettlementCoordinator>,
operator_interactions: Mutex<HashSet<String>>,
}
impl ExecutionMarkStore {
pub fn new() -> Self {
Self::default()
}
pub fn settlement(&self) -> Arc<MarkSettlementCoordinator> {
self.settlement.clone()
}
pub fn arm_settlement(&self, changed: Vec<String>) -> bool {
self.interactions().extend(changed.iter().cloned());
self.settlement.clone().notify(changed)
}
pub fn take_operator_interactions(&self) -> Vec<String> {
let mut guard = self.interactions();
if guard.is_empty() {
return Vec::new();
}
let mut targets: Vec<String> = guard.drain().collect();
targets.sort();
targets
}
fn interactions(&self) -> std::sync::MutexGuard<'_, HashSet<String>> {
self.operator_interactions
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn is_marked(&self, change_id: &str) -> bool {
self.lock().contains(change_id)
}
pub fn set(&self, change_id: &str, marked: bool) -> bool {
let mut guard = self.lock();
if marked {
guard.insert(change_id.to_string())
} else {
guard.remove(change_id)
}
}
#[cfg(test)]
pub fn replace(&self, change_ids: impl IntoIterator<Item = String>) {
*self.lock() = change_ids.into_iter().collect();
}
pub fn marked_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self.lock().iter().cloned().collect();
ids.sort();
ids
}
pub fn clear(&self) {
self.lock().clear();
}
fn lock(&self) -> std::sync::MutexGuard<'_, HashSet<String>> {
self.marks
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
}
#[derive(Clone, Debug)]
pub struct TerminationWaiter {
done: CancellationToken,
}
impl TerminationWaiter {
pub fn new(done: CancellationToken) -> Self {
Self { done }
}
pub fn already_terminated() -> Self {
let done = CancellationToken::new();
done.cancel();
Self { done }
}
pub fn never() -> Self {
Self {
done: CancellationToken::new(),
}
}
pub async fn wait(&self) {
self.done.cancelled().await;
}
}
#[derive(Clone, Debug)]
pub struct PendingTermination {
change_id: String,
waiter: TerminationWaiter,
timeout: Duration,
}
impl PendingTermination {
pub fn change_id(&self) -> &str {
&self.change_id
}
pub async fn confirm_termination(&self) -> OperatorResult<()> {
if tokio::time::timeout(self.timeout, self.waiter.wait())
.await
.is_err()
{
return Err(OperatorCommandError::TerminationTimeout {
change_id: self.change_id.clone(),
waited: self.timeout,
});
}
Ok(())
}
}
#[derive(Clone, Debug)]
pub struct PendingForceStop {
pending: PendingTermination,
execution_id: Option<String>,
terminated: bool,
}
impl PendingForceStop {
pub fn change_id(&self) -> &str {
self.pending.change_id()
}
pub fn execution_id(&self) -> Option<&str> {
self.execution_id.as_deref()
}
pub fn terminated(&self) -> bool {
self.terminated
}
pub async fn confirm_termination(&self) -> OperatorResult<()> {
self.pending.confirm_termination().await
}
}
#[async_trait]
pub trait ManagedProcessTermination: Send + Sync {
async fn owns_managed_process(&self, change_id: &str) -> bool;
async fn kill_managed_process(&self, change_id: &str) -> ImmediateKillEvidence;
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ImmediateKillEvidence {
pub identities: usize,
pub confirmed: bool,
pub detail: String,
}
impl ImmediateKillEvidence {
pub fn nothing_to_kill() -> Self {
Self {
identities: 0,
confirmed: true,
detail: String::new(),
}
}
pub fn confirmed(identities: usize) -> Self {
Self {
identities,
confirmed: true,
detail: String::new(),
}
}
pub fn unconfirmed(identities: usize, detail: impl Into<String>) -> Self {
Self {
identities,
confirmed: false,
detail: detail.into(),
}
}
pub fn signalled(&self) -> bool {
self.identities > 0
}
}
pub struct NoManagedProcesses;
#[async_trait]
impl ManagedProcessTermination for NoManagedProcesses {
async fn owns_managed_process(&self, _change_id: &str) -> bool {
false
}
async fn kill_managed_process(&self, _change_id: &str) -> ImmediateKillEvidence {
ImmediateKillEvidence::nothing_to_kill()
}
}
#[async_trait]
pub trait QueuePort: Send + Sync {
async fn add(&self, change_id: &str) -> bool;
async fn remove(&self, change_id: &str) -> bool;
async fn request_cancellation(
&self,
change_id: &str,
) -> std::result::Result<Option<TerminationWaiter>, String>;
async fn notify_scheduler(&self);
async fn publish_explicit_retry(&self, _change_id: &str, _authority: RetryEdgeAuthority) {}
}
#[async_trait]
pub trait QueueHookPort: Send + Sync {
async fn on_queue_add(&self, change_id: &str);
async fn on_queue_remove(&self, change_id: &str);
}
pub struct NoopQueueHooks;
#[async_trait]
impl QueueHookPort for NoopQueueHooks {
async fn on_queue_add(&self, _change_id: &str) {}
async fn on_queue_remove(&self, _change_id: &str) {}
}
pub struct HookRunnerQueueHooks {
runner: crate::hooks::HookRunner,
}
impl HookRunnerQueueHooks {
pub fn new(runner: crate::hooks::HookRunner) -> Self {
Self { runner }
}
async fn run(&self, hook_type: crate::hooks::HookType, change_id: &str) {
let context = crate::hooks::HookContext::new(0, 0, 0, false).with_change(change_id, 0, 0);
if let Err(error) = self.runner.run_hook(hook_type, &context).await {
tracing::warn!("{hook_type} hook failed for '{change_id}': {error}");
}
}
}
#[async_trait]
impl QueueHookPort for HookRunnerQueueHooks {
async fn on_queue_add(&self, change_id: &str) {
self.run(crate::hooks::HookType::OnQueueAdd, change_id)
.await;
}
async fn on_queue_remove(&self, change_id: &str) {
self.run(crate::hooks::HookType::OnQueueRemove, change_id)
.await;
}
}
pub struct OperatorCommandService {
state: Arc<RwLock<OrchestratorState>>,
queue: Arc<dyn QueuePort>,
hooks: Arc<dyn QueueHookPort>,
marks: Arc<ExecutionMarkStore>,
parallel: Arc<ParallelRuntime>,
cancellation_timeout: Duration,
execution_facts: Option<Arc<ExecutionFactsStore>>,
managed_termination: Arc<dyn ManagedProcessTermination>,
acceptance_admission: Arc<dyn AcceptanceAdmissionPort>,
}
impl OperatorCommandService {
pub fn new(
state: Arc<RwLock<OrchestratorState>>,
queue: Arc<dyn QueuePort>,
hooks: Arc<dyn QueueHookPort>,
marks: Arc<ExecutionMarkStore>,
) -> Self {
Self {
state,
queue,
hooks,
marks,
parallel: Arc::new(ParallelRuntime::new()),
cancellation_timeout: DEFAULT_CANCELLATION_TIMEOUT,
execution_facts: None,
managed_termination: Arc::new(NoManagedProcesses),
acceptance_admission: Arc::new(AlwaysAdmitAcceptance),
}
}
pub fn with_acceptance_admission(mut self, port: Arc<dyn AcceptanceAdmissionPort>) -> Self {
self.acceptance_admission = port;
self
}
pub fn with_parallel(mut self, parallel: Arc<ParallelRuntime>) -> Self {
self.parallel = parallel;
self
}
pub fn with_managed_termination(
mut self,
termination: Arc<dyn ManagedProcessTermination>,
) -> Self {
self.managed_termination = termination;
self
}
pub fn with_execution_facts(mut self, facts: Arc<ExecutionFactsStore>) -> Self {
self.execution_facts = Some(facts);
self
}
pub fn execution_facts(&self) -> Option<Arc<ExecutionFactsStore>> {
self.execution_facts.clone()
}
pub fn with_cancellation_timeout(mut self, timeout: Duration) -> Self {
self.cancellation_timeout = timeout;
self
}
pub fn marks(&self) -> Arc<ExecutionMarkStore> {
self.marks.clone()
}
pub fn parallel(&self) -> Arc<ParallelRuntime> {
self.parallel.clone()
}
pub async fn display_status(&self, change_id: &str) -> String {
self.state
.read()
.await
.display_status(change_id)
.to_string()
}
pub async fn execute(
&self,
_mode: OperatorMode,
command: OperatorCommand,
) -> OperatorResult<OperatorOutcome> {
match command {
OperatorCommand::SetExecutionMark { change_id, marked } => {
self.set_execution_mark(&change_id, marked).await
}
OperatorCommand::AddToQueue { change_id } => self
.add_to_queue(&change_id)
.await
.map(OperatorOutcome::Queue),
OperatorCommand::RemoveFromQueue { change_id } => self
.remove_from_queue(&change_id)
.await
.map(OperatorOutcome::Queue),
OperatorCommand::StopAndDequeue { change_id } => {
self.stop_and_dequeue(&change_id).await
}
OperatorCommand::ForceStopChange { change_id } => {
self.force_stop_change(&change_id).await
}
OperatorCommand::RetryChange { change_id } => self
.retry_change(&change_id)
.await
.map(OperatorOutcome::Retry),
OperatorCommand::SetAllExecutionMarks => self.set_all_execution_marks().await,
}
}
pub async fn apply_execution_mark(&self, change_id: &str, marked: bool) -> bool {
let _mutation = self.parallel.lock_mutations().await;
let changed = self.marks.set(change_id, marked);
if changed {
self.marks.arm_settlement(vec![change_id.to_string()]);
}
changed
}
pub async fn apply_admission_execution_mark(&self, change_id: &str, marked: bool) -> bool {
let _mutation = self.parallel.lock_mutations().await;
self.marks.set(change_id, marked)
}
pub async fn set_execution_mark(
&self,
change_id: &str,
marked: bool,
) -> OperatorResult<OperatorOutcome> {
let _mutation = self.parallel.lock_mutations().await;
let (display_status, archive_complete) = {
let guard = self.state.read().await;
(
guard.display_status(change_id).to_string(),
guard.is_archived(change_id),
)
};
match classify_mark_admission(&display_status, archive_complete) {
MarkAdmission::Allowed => {}
MarkAdmission::TerminalTarget => {
return Ok(OperatorOutcome::NoOp {
change_id: change_id.to_string(),
reason: NoOpReason::TerminalMarkTarget,
})
}
MarkAdmission::ArchiveComplete => {
return Ok(OperatorOutcome::NoOp {
change_id: change_id.to_string(),
reason: NoOpReason::ArchiveCompleteMarkTarget,
})
}
}
if self.marks.set(change_id, marked) {
self.marks.arm_settlement(vec![change_id.to_string()]);
Ok(OperatorOutcome::MarkSet {
change_id: change_id.to_string(),
marked,
})
} else {
Ok(OperatorOutcome::NoOp {
change_id: change_id.to_string(),
reason: NoOpReason::MarkUnchanged,
})
}
}
pub async fn set_all_execution_marks(&self) -> OperatorResult<OperatorOutcome> {
let _mutation = self.parallel.lock_mutations().await;
let observed: Vec<(String, String, bool, bool)> = {
let guard = self.state.read().await;
guard
.tracked_change_ids()
.into_iter()
.map(|change_id| {
let display_status = guard.display_status(&change_id).to_string();
let archive_complete = guard.is_archived(&change_id);
let marked = self.marks.is_marked(&change_id);
(change_id, display_status, archive_complete, marked)
})
.collect()
};
let rows: Vec<MarkTargetRow<'_>> = observed
.iter()
.map(
|(change_id, display_status, archive_complete, marked)| MarkTargetRow {
change_id,
display_status,
archive_complete: *archive_complete,
marked: *marked,
},
)
.collect();
let plan = plan_bulk_marks(&rows);
if plan.is_empty() {
return Ok(OperatorOutcome::NoOp {
change_id: String::new(),
reason: NoOpReason::NoEligibleMarkTarget,
});
}
let mut changed = Vec::new();
for change_id in &plan.eligible {
if self.marks.set(change_id, plan.target_state) {
changed.push(change_id.clone());
}
}
if changed.is_empty() {
return Ok(OperatorOutcome::NoOp {
change_id: String::new(),
reason: NoOpReason::BulkMarksUnchanged,
});
}
self.marks.arm_settlement(changed.clone());
Ok(OperatorOutcome::BulkMarks {
marked: plan.target_state,
changed,
excluded: plan.excluded,
})
}
pub async fn plan_mark_settlement(&self, targets: &[String]) -> MarkSettlementPlan {
let observed: Vec<(String, String, bool, bool, bool)> = {
let _mutation = self.parallel.lock_mutations().await;
let guard = self.state.read().await;
let tracked: HashSet<String> = guard.tracked_change_ids().into_iter().collect();
targets
.iter()
.map(|change_id| {
let display_status = guard.display_status(change_id).to_string();
let tracked = tracked.contains(change_id);
let eligible = self.parallel.is_eligible(change_id);
let marked = self.marks.is_marked(change_id);
(change_id.clone(), display_status, tracked, eligible, marked)
})
.collect()
};
let rows: Vec<MarkSettlementRow<'_>> = observed
.iter()
.map(
|(change_id, display_status, tracked, parallel_eligible, marked)| {
MarkSettlementRow {
change_id,
display_status,
tracked: *tracked,
parallel_eligible: *parallel_eligible,
marked: *marked,
}
},
)
.collect();
let mut plan = plan_mark_settlement(&rows);
let mut admitted = Vec::with_capacity(plan.additions.len());
for change_id in std::mem::take(&mut plan.additions) {
if self
.unchanged_acceptance_refusal(&change_id)
.await
.is_some()
{
plan.excluded
.push((change_id, MarkSettlementExclusion::UnchangedAcceptanceInput));
} else {
admitted.push(change_id);
}
}
plan.additions = admitted;
plan
}
pub async fn apply_settlement_queue_intent(
&self,
change_id: &str,
action: MarkSettlementAction,
) -> SettlementApplication {
let queued = matches!(action, MarkSettlementAction::Add);
if queued {
if let Some(refusal) = self.unchanged_acceptance_refusal(change_id).await {
tracing::warn!(
change_id = %change_id,
"Mark settlement did not queue the change: {refusal}"
);
return SettlementApplication {
outcome: QueueOutcome {
change_id: change_id.to_string(),
mutation: queue_mutation_for(action),
reducer_changed: false,
dynamic_queue_mutated: false,
display_status: self.display_status(change_id).await,
},
skipped: Some(MarkSettlementExclusion::UnchangedAcceptanceInput),
};
}
}
let guard_outcome = {
let mut guard = self.state.write().await;
let tracked = guard.change_runtime(change_id).is_some();
let display_status = guard.display_status(change_id);
let row = MarkSettlementRow {
change_id,
display_status,
tracked,
parallel_eligible: true,
marked: queued,
};
match classify_mark_settlement_row(&row) {
Ok(current) if current == action => {
let command = if queued {
ReducerCommand::AddToQueue(change_id.to_string())
} else {
ReducerCommand::RemoveFromQueue(change_id.to_string())
};
Ok(matches!(
guard.apply_command(command),
ReduceOutcome::Changed(_)
))
}
Ok(_) => Err(if queued {
MarkSettlementExclusion::AlreadyQueued
} else {
MarkSettlementExclusion::AlreadyNotQueued
}),
Err(reason) => Err(reason),
}
};
let reducer_changed = match guard_outcome {
Ok(changed) => changed,
Err(reason) => {
return SettlementApplication {
outcome: QueueOutcome {
change_id: change_id.to_string(),
mutation: queue_mutation_for(action),
reducer_changed: false,
dynamic_queue_mutated: false,
display_status: self.display_status(change_id).await,
},
skipped: Some(reason),
};
}
};
let dynamic_queue_mutated = if !reducer_changed {
false
} else if queued {
let added = self.queue.add(change_id).await;
if added {
self.hooks.on_queue_add(change_id).await;
}
added
} else {
let removed = self.queue.remove(change_id).await;
self.hooks.on_queue_remove(change_id).await;
removed
};
SettlementApplication {
outcome: QueueOutcome {
change_id: change_id.to_string(),
mutation: queue_mutation_for(action),
reducer_changed,
dynamic_queue_mutated,
display_status: self.display_status(change_id).await,
},
skipped: None,
}
}
pub async fn notify_scheduler_after_settlement(&self) {
self.queue.notify_scheduler().await;
}
pub async fn add_to_queue(&self, change_id: &str) -> OperatorResult<QueueOutcome> {
let is_retry_shaped = {
let guard = self.state.read().await;
guard.is_terminal_error_change(change_id)
|| guard.change_runtime(change_id).is_some_and(
crate::orchestration::state::ChangeRuntimeState::is_acceptance_stalled,
)
};
if is_retry_shaped {
if let Some(refusal) = self.unchanged_acceptance_refusal(change_id).await {
return Err(refusal);
}
}
let (reduce_outcome, was_error_retry) = {
let mut guard = self.state.write().await;
if guard.is_terminal_error_change(change_id) {
(
guard.apply_command(ReducerCommand::RetryError(change_id.to_string())),
true,
)
} else {
(
guard.apply_command(ReducerCommand::AddToQueue(change_id.to_string())),
false,
)
}
};
let reducer_changed = matches!(reduce_outcome, ReduceOutcome::Changed(_));
if was_error_retry && reducer_changed {
self.queue
.publish_explicit_retry(change_id, RetryEdgeAuthority::TerminalError)
.await;
}
let dynamic_queue_mutated = if reducer_changed {
self.queue.add(change_id).await
} else {
false
};
if dynamic_queue_mutated {
self.queue.notify_scheduler().await;
self.hooks.on_queue_add(change_id).await;
}
Ok(QueueOutcome {
change_id: change_id.to_string(),
mutation: QueueMutation::Added,
reducer_changed,
dynamic_queue_mutated,
display_status: self.display_status(change_id).await,
})
}
pub async fn remove_from_queue(&self, change_id: &str) -> OperatorResult<QueueOutcome> {
let reduce_outcome = {
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::RemoveFromQueue(change_id.to_string()))
};
let reducer_changed = matches!(reduce_outcome, ReduceOutcome::Changed(_));
let removed_from_dynamic_queue = self.queue.remove(change_id).await;
let dynamic_queue_mutated = reducer_changed || removed_from_dynamic_queue;
if dynamic_queue_mutated {
self.hooks.on_queue_remove(change_id).await;
}
Ok(QueueOutcome {
change_id: change_id.to_string(),
mutation: QueueMutation::Removed,
reducer_changed,
dynamic_queue_mutated,
display_status: self.display_status(change_id).await,
})
}
pub async fn stop_and_dequeue(&self, change_id: &str) -> OperatorResult<OperatorOutcome> {
let pending = self.begin_stop_and_dequeue(change_id).await?;
pending.confirm_termination().await?;
self.commit_stop_and_dequeue(change_id, ApplyCommitEvidence::unknown())
.await
}
pub fn cancellation_timeout(&self) -> Duration {
self.cancellation_timeout
}
pub async fn begin_stop_and_dequeue(
&self,
change_id: &str,
) -> OperatorResult<PendingTermination> {
let display_status = self.display_status(change_id).await;
let was_active = is_active_status(&display_status);
let waiter = match self.queue.request_cancellation(change_id).await {
Err(message) => {
return Err(OperatorCommandError::CancellationFailed {
change_id: change_id.to_string(),
message,
})
}
Ok(Some(waiter)) => waiter,
Ok(None) if was_active => {
return Err(OperatorCommandError::MissingCancellationHandle {
change_id: change_id.to_string(),
});
}
Ok(None) => TerminationWaiter::already_terminated(),
};
Ok(PendingTermination {
change_id: change_id.to_string(),
waiter,
timeout: self.cancellation_timeout,
})
}
pub async fn commit_stop_and_dequeue(
&self,
change_id: &str,
apply_commit: ApplyCommitEvidence,
) -> OperatorResult<OperatorOutcome> {
let _mutation = self.parallel.lock_mutations().await;
let (cancelled_phase, reduce_outcome) = {
let mut guard = self.state.write().await;
let phase = guard
.change_runtime(change_id)
.map(|runtime| project_phase(runtime, self.push_open(change_id)))
.unwrap_or(ExecutionPhase::Unknown);
let outcome = guard.apply_command(ReducerCommand::DequeueChange(change_id.to_string()));
(phase, outcome)
};
if matches!(reduce_outcome, ReduceOutcome::NoOp) {
return Ok(OperatorOutcome::NoOp {
change_id: change_id.to_string(),
reason: NoOpReason::ReducerRejected,
});
}
self.marks.set(change_id, false);
Ok(OperatorOutcome::Dequeued {
change_id: change_id.to_string(),
settlement: StopSettlement {
cancelled_phase,
last_completed_phase: self
.execution_facts
.as_ref()
.and_then(|facts| facts.change(change_id).last_completed_phase),
apply_commit_present: apply_commit.present,
apply_commit_oid: apply_commit.oid,
},
})
}
pub async fn force_stop_change(&self, change_id: &str) -> OperatorResult<OperatorOutcome> {
let pending = self.begin_force_stop_change(change_id).await?;
pending.confirm_termination().await?;
self.commit_force_stop_change(&pending, ApplyCommitEvidence::unknown())
.await
}
pub async fn force_stop_admission(&self, change_id: &str) -> ForceStopAdmission {
let (display_status, tracked) = self.force_stop_facts(change_id).await;
classify_force_stop_change(
&display_status,
tracked,
self.managed_termination
.owns_managed_process(change_id)
.await,
)
}
async fn force_stop_facts(&self, change_id: &str) -> (String, bool) {
let guard = self.state.read().await;
(
guard.display_status(change_id).to_string(),
guard.is_tracked_change(change_id),
)
}
pub async fn begin_force_stop_change(
&self,
change_id: &str,
) -> OperatorResult<PendingForceStop> {
let (display_status, tracked) = self.force_stop_facts(change_id).await;
let owns_managed_process = self
.managed_termination
.owns_managed_process(change_id)
.await;
let admission = classify_force_stop_change(&display_status, tracked, owns_managed_process);
if let Some(reason) = admission.exclusion() {
return Err(OperatorCommandError::ForceStopIneligible {
change_id: change_id.to_string(),
display_status,
reason,
});
}
let execution_id = self
.execution_facts
.as_ref()
.and_then(|facts| facts.execution_id(change_id));
let kill = match admission {
ForceStopAdmission::KillAndDequeue => {
self.managed_termination
.kill_managed_process(change_id)
.await
}
_ => ImmediateKillEvidence::nothing_to_kill(),
};
if !kill.confirmed {
return Err(OperatorCommandError::ForceStopUnconfirmed {
change_id: change_id.to_string(),
detail: kill.detail,
});
}
let waiter = match self.queue.request_cancellation(change_id).await {
Err(message) => {
return Err(OperatorCommandError::CancellationFailed {
change_id: change_id.to_string(),
message,
})
}
Ok(Some(waiter)) => waiter,
Ok(None) => TerminationWaiter::already_terminated(),
};
Ok(PendingForceStop {
pending: PendingTermination {
change_id: change_id.to_string(),
waiter,
timeout: self.cancellation_timeout,
},
execution_id,
terminated: kill.signalled(),
})
}
pub async fn commit_force_stop_change(
&self,
pending: &PendingForceStop,
apply_commit: ApplyCommitEvidence,
) -> OperatorResult<OperatorOutcome> {
let change_id = pending.change_id().to_string();
let _mutation = self.parallel.lock_mutations().await;
let cancelled_phase = {
let mut guard = self.state.write().await;
let phase = guard
.change_runtime(&change_id)
.map(|runtime| project_phase(runtime, self.push_open(&change_id)))
.unwrap_or(ExecutionPhase::Unknown);
guard.apply_command(ReducerCommand::StopChange(change_id.clone()));
phase
};
self.marks.set(&change_id, false);
Ok(OperatorOutcome::ForceStopped {
change_id: change_id.clone(),
execution_id: pending.execution_id.clone(),
terminated: pending.terminated,
settlement: StopSettlement {
cancelled_phase,
last_completed_phase: self
.execution_facts
.as_ref()
.and_then(|facts| facts.change(&change_id).last_completed_phase),
apply_commit_present: apply_commit.present,
apply_commit_oid: apply_commit.oid,
},
})
}
fn push_open(&self, change_id: &str) -> bool {
self.execution_facts
.as_ref()
.is_some_and(|facts| facts.change(change_id).current_phase == ExecutionPhase::Push)
}
pub async fn resolve_merge(&self, change_id: &str) -> bool {
let reduce_outcome = {
let mut guard = self.state.write().await;
guard.apply_command(ReducerCommand::ResolveMerge(change_id.to_string()))
};
matches!(reduce_outcome, ReduceOutcome::Changed(_))
}
async fn blocker_kind(&self, change_id: &str) -> crate::orchestration::state::BlockerKind {
self.state
.read()
.await
.change_runtime(change_id)
.map(crate::orchestration::state::ChangeRuntimeState::blocker_kind)
.unwrap_or_default()
}
pub async fn retry_change(&self, change_id: &str) -> OperatorResult<RetryPlan> {
let routes = match self.plan_retry_change(change_id).await? {
Some(route) => vec![(change_id.to_string(), route)],
None => Vec::new(),
};
Ok(self.commit_retry_routes(&routes).await)
}
pub async fn plan_retry_change(&self, change_id: &str) -> OperatorResult<Option<RetryRoute>> {
let display_status = self.display_status(change_id).await;
let blocker_kind = self.blocker_kind(change_id).await;
let Some(route) = classify_retry_route(&display_status, blocker_kind) else {
return Err(OperatorCommandError::RetryUnsupported {
change_id: change_id.to_string(),
display_status,
});
};
if let Some(refusal) = self.unchanged_acceptance_refusal(change_id).await {
return Err(refusal);
}
Ok(self
.route_is_committable(change_id, route, true)
.await
.then_some(route))
}
async fn unchanged_acceptance_refusal(&self, change_id: &str) -> Option<OperatorCommandError> {
use crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission;
match self.acceptance_admission.classify(change_id).await {
AcceptanceAdmission::Admit => None,
AcceptanceAdmission::Refuse {
category,
fingerprint,
components,
..
} => {
tracing::warn!(
change_id = %change_id,
category = category.as_str(),
"Retry refused: Acceptance input is unchanged since a non-resumable hold"
);
Some(OperatorCommandError::UnchangedAcceptanceInput {
change_id: change_id.to_string(),
category: category.as_str().to_string(),
fingerprint,
components: components.iter().map(|part| part.to_string()).collect(),
})
}
}
}
pub async fn plan_retry_errors(&self, change_ids: &[String]) -> Vec<(String, RetryRoute)> {
self.plan_retry_errors_with_refusals(change_ids).await.0
}
pub async fn plan_retry_errors_with_refusals(
&self,
change_ids: &[String],
) -> (Vec<(String, RetryRoute)>, Vec<OperatorCommandError>) {
let mut routes = Vec::new();
let mut refusals = Vec::new();
for change_id in change_ids {
let display_status = self.display_status(change_id).await;
let blocker_kind = self.blocker_kind(change_id).await;
let Some(route) = classify_retry_route(&display_status, blocker_kind) else {
continue;
};
if let Some(refusal) = self.unchanged_acceptance_refusal(change_id).await {
refusals.push(refusal);
continue;
}
if self.route_is_committable(change_id, route, true).await {
routes.push((change_id.clone(), route));
}
}
(routes, refusals)
}
async fn route_is_committable(
&self,
change_id: &str,
route: RetryRoute,
acceptance_admitted: bool,
) -> bool {
if !matches!(route, RetryRoute::AcceptanceStall) {
return true;
}
let (non_resumable, fingerprint_guarded) = {
let guard = self.state.read().await;
guard
.change_runtime(change_id)
.map_or((false, false), |rt| {
(
rt.is_acceptance_stalled() && !rt.is_resumable_acceptance_stall(),
rt.is_acceptance_execution_hold(),
)
})
};
if !non_resumable {
return true;
}
if fingerprint_guarded && acceptance_admitted {
tracing::info!(
change_id = %change_id,
"Explicit retry admitted: the Acceptance input fingerprint moved since the \
execution hold, so one fresh bounded attempt is permitted"
);
return true;
}
tracing::warn!(
change_id = %change_id,
fingerprint_guarded,
"Explicit retry refused: the acceptance stall is not resumable, so its \
blocker evidence is retained"
);
false
}
pub async fn commit_retry_routes(&self, routes: &[(String, RetryRoute)]) -> RetryPlan {
let mut plan = RetryPlan {
change_ids: Vec::new(),
routes: Vec::new(),
explicit_retry: false,
};
for (change_id, route) in routes {
let accepted = self.apply_retry_route(change_id, *route).await;
plan.change_ids.extend(accepted.change_ids);
plan.routes.extend(accepted.routes);
plan.explicit_retry |= accepted.explicit_retry;
}
plan
}
pub async fn retry_errors(&self, change_ids: &[String]) -> RetryPlan {
self.retry_errors_with_refusals(change_ids).await.0
}
pub async fn retry_errors_with_refusals(
&self,
change_ids: &[String],
) -> (RetryPlan, Vec<OperatorCommandError>) {
let (routes, refusals) = self.plan_retry_errors_with_refusals(change_ids).await;
(self.commit_retry_routes(&routes).await, refusals)
}
async fn apply_retry_route(&self, change_id: &str, route: RetryRoute) -> RetryPlan {
let command = match route {
RetryRoute::TerminalError => ReducerCommand::RetryError(change_id.to_string()),
RetryRoute::AcceptanceStall => ReducerCommand::AddToQueue(change_id.to_string()),
};
let _mutation = self.parallel.lock_mutations().await;
let reduce_outcome = {
let mut guard = self.state.write().await;
guard.apply_command(command)
};
if matches!(reduce_outcome, ReduceOutcome::NoOp) {
return RetryPlan {
change_ids: Vec::new(),
routes: Vec::new(),
explicit_retry: false,
};
}
self.queue
.publish_explicit_retry(change_id, RetryEdgeAuthority::from(route))
.await;
self.marks.set(change_id, true);
RetryPlan {
change_ids: vec![change_id.to_string()],
routes: vec![route],
explicit_retry: true,
}
}
}
pub struct LiveRunManagedProcesses<F>
where
F: Fn() -> Option<crate::ai_command_runner::RunCommandScope> + Send + Sync,
{
current_scope: F,
}
impl<F> LiveRunManagedProcesses<F>
where
F: Fn() -> Option<crate::ai_command_runner::RunCommandScope> + Send + Sync,
{
pub fn new(current_scope: F) -> Self {
Self { current_scope }
}
}
#[async_trait]
impl<F> ManagedProcessTermination for LiveRunManagedProcesses<F>
where
F: Fn() -> Option<crate::ai_command_runner::RunCommandScope> + Send + Sync,
{
async fn owns_managed_process(&self, change_id: &str) -> bool {
(self.current_scope)().is_some_and(|scope| scope.change_owns_managed_process(change_id))
}
async fn kill_managed_process(&self, change_id: &str) -> ImmediateKillEvidence {
let Some(scope) = (self.current_scope)() else {
return ImmediateKillEvidence::nothing_to_kill();
};
let report = scope
.force_stop_change(
change_id,
crate::ai_command_runner::FORCE_STOP_CHANGE_KILL_BUDGET,
)
.await;
if report.is_confirmed() {
ImmediateKillEvidence::confirmed(report.identities)
} else {
ImmediateKillEvidence::unconfirmed(report.identities, report.diagnostics())
}
}
}
#[async_trait]
impl QueuePort for crate::tui::queue::DynamicQueue {
async fn add(&self, change_id: &str) -> bool {
self.push(change_id.to_string()).await
}
async fn remove(&self, change_id: &str) -> bool {
let removed_from_queue = crate::tui::queue::DynamicQueue::remove(self, change_id).await;
self.mark_removed(change_id.to_string()).await;
removed_from_queue
}
async fn request_cancellation(
&self,
change_id: &str,
) -> std::result::Result<Option<TerminationWaiter>, String> {
Ok(
crate::tui::queue::DynamicQueue::request_cancellation(self, change_id)
.await
.map(TerminationWaiter::new),
)
}
async fn notify_scheduler(&self) {
crate::tui::queue::DynamicQueue::notify_scheduler(self);
}
async fn publish_explicit_retry(&self, change_id: &str, authority: RetryEdgeAuthority) {
crate::tui::queue::DynamicQueue::publish_explicit_retry(
self,
change_id.to_string(),
authority,
)
.await;
}
}
#[cfg(test)]
mod tests;
#[cfg(test)]
#[path = "operator_command/acceptance_admission_tests.rs"]
mod acceptance_admission_tests;