use std::sync::{Arc, Mutex};
use crate::events::{
accepted_start_opens_idle_run_episode, all_completed_may_overwrite_mode,
graceful_stop_is_idle_origin, is_admitted_work_start, persistent_idle_may_project_ready,
EventDispatcher, ExecutionEvent, OperatorCommandEffect, OutcomeRevisions,
};
use crate::orchestration::apply_commit_evidence::ApplyCommitEvidencePort;
use crate::orchestration::mark_settlement::{
MarkSettlementAction, MarkSettlementExclusion, MarkSettlementFailure, MarkSettlementPlan,
MarkSettlementRuntime,
};
use crate::orchestration::operator_command::{
OperatorMode, OperatorOutcome, PendingForceStop, PendingTermination,
};
use crate::orchestration::run_control::{
ResolveReservation, RunControlError, RunControlOutcome, RunControlService, RunNoOpReason,
SchedulerEffect,
};
#[derive(Debug)]
pub struct CoreMode {
state: Mutex<CoreModeState>,
}
#[derive(Debug, Clone, Copy)]
struct CoreModeState {
mode: OperatorMode,
persistent_idle: bool,
}
impl Default for CoreMode {
fn default() -> Self {
Self::new()
}
}
impl CoreMode {
pub fn new() -> Self {
Self {
state: Mutex::new(CoreModeState {
mode: OperatorMode::Select,
persistent_idle: false,
}),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, CoreModeState> {
self.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn get(&self) -> OperatorMode {
self.lock().mode
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn set(&self, mode: OperatorMode) -> bool {
let mut guard = self.lock();
let changed = guard.mode != mode;
guard.mode = mode;
guard.persistent_idle &= mode == OperatorMode::Select;
changed
}
#[cfg(test)]
pub fn set_persistent_idle(&self, persistent_idle: bool) {
self.lock().persistent_idle = persistent_idle;
}
pub fn apply_event(&self, event: &ExecutionEvent) -> Option<OperatorMode> {
let mut guard = self.lock();
if is_admitted_work_start(event)
|| matches!(
event,
ExecutionEvent::Stopped | ExecutionEvent::Error { .. }
)
{
guard.persistent_idle = false;
}
let next = match event {
ExecutionEvent::ProcessingStarted(_) => OperatorMode::Running,
ExecutionEvent::Stopping => {
if graceful_stop_is_idle_origin(guard.mode.as_app_mode()) {
guard.persistent_idle = true;
}
OperatorMode::Stopping
}
ExecutionEvent::Stopped => OperatorMode::Stopped,
ExecutionEvent::Error { .. } => OperatorMode::Error,
ExecutionEvent::ProcessingError { .. } => return None,
ExecutionEvent::PersistentSchedulerIdle => {
if !persistent_idle_may_project_ready(guard.mode.as_app_mode()) {
return None;
}
guard.persistent_idle = true;
OperatorMode::Select
}
ExecutionEvent::AllCompleted => {
if !all_completed_may_overwrite_mode(guard.mode.as_app_mode()) {
return None;
}
OperatorMode::Select
}
ExecutionEvent::OperatorCommandApplied { effect } => match effect {
OperatorCommandEffect::RunDispatched {
scheduler_started: true,
..
} => OperatorMode::Running,
OperatorCommandEffect::RunDispatched {
change_ids,
scheduler_started,
..
} => {
if !accepted_start_opens_idle_run_episode(
guard.mode.as_app_mode(),
guard.persistent_idle,
*scheduler_started,
change_ids,
) {
return None;
}
guard.persistent_idle = false;
OperatorMode::Running
}
OperatorCommandEffect::StopCancelled if guard.persistent_idle => {
OperatorMode::Select
}
OperatorCommandEffect::StopCancelled => OperatorMode::Running,
OperatorCommandEffect::ForceStopAwaitingBoundary { .. } => OperatorMode::Stopping,
OperatorCommandEffect::ResolveReserved { active: true, .. } => {
OperatorMode::Running
}
OperatorCommandEffect::ResolveReserved { active: false, .. }
| OperatorCommandEffect::MarkDelta { .. }
| OperatorCommandEffect::QueueDelta { .. } => return None,
},
_ if is_admitted_work_start(event) && guard.mode == OperatorMode::Select => {
OperatorMode::Running
}
_ => return None,
};
let changed = guard.mode != next;
guard.mode = next;
changed.then_some(next)
}
}
#[cfg(test)]
mod core_mode_scope_tests {
use super::*;
fn processing_error() -> ExecutionEvent {
ExecutionEvent::ProcessingError {
id: "alpha".to_string(),
error: "acceptance command attempts exhausted".to_string(),
}
}
#[test]
fn processing_error_preserves_shared_mode() {
for mode in [
OperatorMode::Select,
OperatorMode::Running,
OperatorMode::Stopping,
OperatorMode::Stopped,
OperatorMode::Error,
] {
let core = CoreMode::new();
core.set(mode);
assert_eq!(
core.apply_event(&processing_error()),
None,
"a change-scoped failure moves no process mode (from {mode:?})"
);
assert_eq!(
core.get(),
mode,
"the process mode must be exactly the one that existed before the event"
);
}
}
#[test]
fn processing_error_preserves_shared_mode_persistent_idle_episode() {
let core = CoreMode::new();
core.set(OperatorMode::Select);
core.set_persistent_idle(true);
assert_eq!(core.apply_event(&processing_error()), None);
assert_eq!(
core.apply_event(&ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::StopCancelled,
}),
None,
"the withdrawn stop returns to the Ready the episode still describes"
);
assert_eq!(core.get(), OperatorMode::Select);
}
#[test]
fn processing_error_preserves_shared_mode_fatal_control_still_transitions() {
for mode in [
OperatorMode::Select,
OperatorMode::Running,
OperatorMode::Stopping,
OperatorMode::Stopped,
] {
let core = CoreMode::new();
core.set(mode);
assert_eq!(
core.apply_event(&ExecutionEvent::Error {
message: "the orchestrator could not start".to_string(),
}),
Some(OperatorMode::Error),
"a global Error is process-fatal (from {mode:?})"
);
assert_eq!(core.get(), OperatorMode::Error);
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperatorIntent {
Start,
Stop,
CancelStop,
ForceStop,
RetryChange {
change_id: String,
},
RetryErrors {
change_ids: Vec<String>,
},
ResolveMerge {
change_id: String,
},
SetExecutionMark {
change_id: String,
marked: bool,
},
SetQueueIntent {
change_id: String,
queued: bool,
},
SetAllExecutionMarks,
StopAndDequeue {
change_id: String,
},
ForceStopChange {
change_id: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ApplicationOutcome {
Run(RunControlOutcome),
Operator(OperatorOutcome),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApplicationResult {
pub outcome: Result<ApplicationOutcome, RunControlError>,
pub revision: Option<u64>,
}
impl ApplicationResult {
fn failed(error: RunControlError) -> Self {
Self {
outcome: Err(error),
revision: None,
}
}
fn run(outcome: RunControlOutcome, revision: Option<u64>) -> Self {
Self {
outcome: Ok(ApplicationOutcome::Run(outcome)),
revision,
}
}
fn operator(outcome: OperatorOutcome, revision: Option<u64>) -> Self {
Self {
outcome: Ok(ApplicationOutcome::Operator(outcome)),
revision,
}
}
}
pub type ApplicationGuard = tokio::sync::OwnedMutexGuard<()>;
pub struct OperatorApplication {
gate: Arc<tokio::sync::Mutex<()>>,
mode: Arc<CoreMode>,
run_control: Arc<RunControlService>,
dispatch: Arc<EventDispatcher>,
revisions: Option<Arc<dyn OutcomeRevisions>>,
apply_commit_evidence: Option<Arc<dyn ApplyCommitEvidencePort>>,
}
impl OperatorApplication {
pub fn new(
mode: Arc<CoreMode>,
run_control: Arc<RunControlService>,
dispatch: Arc<EventDispatcher>,
) -> Self {
Self {
gate: Arc::new(tokio::sync::Mutex::new(())),
mode,
run_control,
dispatch,
revisions: None,
apply_commit_evidence: None,
}
}
pub fn with_revisions(mut self, revisions: Option<Arc<dyn OutcomeRevisions>>) -> Self {
self.revisions = revisions;
self
}
pub fn with_apply_commit_evidence(
mut self,
port: Option<Arc<dyn ApplyCommitEvidencePort>>,
) -> Self {
self.apply_commit_evidence = port;
self
}
pub fn gate(&self) -> Arc<tokio::sync::Mutex<()>> {
self.gate.clone()
}
pub fn run_control(&self) -> Arc<RunControlService> {
self.run_control.clone()
}
pub async fn apply(&self, intent: OperatorIntent) -> ApplicationResult {
if let OperatorIntent::StopAndDequeue { change_id } = &intent {
return self.apply_stop_and_dequeue(change_id, None).await;
}
if let OperatorIntent::ForceStopChange { change_id } = &intent {
return self.apply_force_stop_change(change_id, None).await;
}
let guard = self.gate.clone().lock_owned().await;
self.apply_ordinary(intent, &guard).await
}
pub async fn apply_held(
&self,
intent: OperatorIntent,
guard: ApplicationGuard,
) -> ApplicationResult {
if let OperatorIntent::StopAndDequeue { change_id } = &intent {
let change_id = change_id.clone();
return self.apply_stop_and_dequeue(&change_id, Some(guard)).await;
}
if let OperatorIntent::ForceStopChange { change_id } = &intent {
let change_id = change_id.clone();
return self.apply_force_stop_change(&change_id, Some(guard)).await;
}
self.apply_ordinary(intent, &guard).await
}
pub async fn begin_stop_and_dequeue(
&self,
change_id: &str,
) -> Result<PendingTermination, RunControlError> {
self.run_control
.operator()
.begin_stop_and_dequeue(change_id)
.await
.map_err(RunControlError::Operator)
}
pub async fn settle_stop_and_dequeue(&self, pending: PendingTermination) -> ApplicationResult {
let change_id = pending.change_id().to_string();
if let Err(error) = pending.confirm_termination().await {
let _guard = self.gate.clone().lock_owned().await;
return ApplicationResult {
outcome: Err(RunControlError::Operator(error)),
revision: Some(self.current_revision()),
};
}
let apply_commit = crate::orchestration::apply_commit_evidence::observe_apply_commit(
self.run_control.operator().execution_facts().as_deref(),
self.apply_commit_evidence.as_deref(),
&change_id,
)
.await;
let _guard = self.gate.clone().lock_owned().await;
match self
.run_control
.operator()
.commit_stop_and_dequeue(&change_id, apply_commit)
.await
{
Ok(OperatorOutcome::Dequeued {
change_id,
settlement,
}) => {
let revision = self
.dispatch_outcome(ExecutionEvent::ChangeDequeued {
change_id: change_id.clone(),
})
.await;
ApplicationResult::operator(
OperatorOutcome::Dequeued {
change_id,
settlement,
},
Some(revision),
)
}
Ok(outcome) => ApplicationResult::operator(outcome, Some(self.current_revision())),
Err(error) => ApplicationResult {
outcome: Err(RunControlError::Operator(error)),
revision: Some(self.current_revision()),
},
}
}
pub async fn begin_force_stop_change(
&self,
change_id: &str,
) -> Result<PendingForceStop, RunControlError> {
self.run_control
.operator()
.begin_force_stop_change(change_id)
.await
.map_err(RunControlError::Operator)
}
pub async fn settle_force_stop_change(&self, pending: PendingForceStop) -> ApplicationResult {
let change_id = pending.change_id().to_string();
if let Err(error) = pending.confirm_termination().await {
let _guard = self.gate.clone().lock_owned().await;
return ApplicationResult {
outcome: Err(RunControlError::Operator(error)),
revision: Some(self.current_revision()),
};
}
let apply_commit = crate::orchestration::apply_commit_evidence::observe_apply_commit(
self.run_control.operator().execution_facts().as_deref(),
self.apply_commit_evidence.as_deref(),
&change_id,
)
.await;
let _guard = self.gate.clone().lock_owned().await;
match self
.run_control
.operator()
.commit_force_stop_change(&pending, apply_commit)
.await
{
Ok(outcome @ OperatorOutcome::ForceStopped { .. }) => {
let revision = self
.dispatch_outcome(ExecutionEvent::ChangeDequeued {
change_id: change_id.clone(),
})
.await;
ApplicationResult::operator(outcome, Some(revision))
}
Ok(outcome) => ApplicationResult::operator(outcome, Some(self.current_revision())),
Err(error) => ApplicationResult {
outcome: Err(RunControlError::Operator(error)),
revision: Some(self.current_revision()),
},
}
}
async fn apply_force_stop_change(
&self,
change_id: &str,
held: Option<ApplicationGuard>,
) -> ApplicationResult {
let guard = match held {
Some(guard) => guard,
None => self.gate.clone().lock_owned().await,
};
let pending = match self.begin_force_stop_change(change_id).await {
Ok(pending) => pending,
Err(error) => return ApplicationResult::failed(error),
};
drop(guard);
self.settle_force_stop_change(pending).await
}
async fn apply_stop_and_dequeue(
&self,
change_id: &str,
held: Option<ApplicationGuard>,
) -> ApplicationResult {
let guard = match held {
Some(guard) => guard,
None => self.gate.clone().lock_owned().await,
};
let pending = match self.begin_stop_and_dequeue(change_id).await {
Ok(pending) => pending,
Err(error) => return ApplicationResult::failed(error),
};
drop(guard);
self.settle_stop_and_dequeue(pending).await
}
async fn apply_ordinary(
&self,
intent: OperatorIntent,
_guard: &ApplicationGuard,
) -> ApplicationResult {
let mode = self.mode.get();
match intent {
OperatorIntent::Start => {
self.apply_prepared(self.run_control.prepare_start(mode).await)
.await
}
OperatorIntent::RetryChange { change_id } => {
self.apply_prepared(self.run_control.prepare_retry_change(&change_id).await)
.await
}
OperatorIntent::RetryErrors { change_ids } => {
self.apply_prepared(self.run_control.prepare_retry_errors(&change_ids).await)
.await
}
OperatorIntent::ResolveMerge { change_id } => self.apply_resolve(&change_id).await,
OperatorIntent::Stop => self.apply_stop(mode).await,
OperatorIntent::CancelStop => self.apply_cancel_stop(mode).await,
OperatorIntent::ForceStop => self.apply_force_stop(mode).await,
OperatorIntent::SetExecutionMark { change_id, marked } => {
self.apply_mark(&change_id, marked).await
}
OperatorIntent::SetQueueIntent { change_id, queued } => {
self.apply_queue_intent(&change_id, queued).await
}
OperatorIntent::SetAllExecutionMarks => self.apply_bulk_marks().await,
OperatorIntent::StopAndDequeue { .. } | OperatorIntent::ForceStopChange { .. } => {
unreachable!("a two-phase intent is routed before the ordinary transaction")
}
}
}
async fn apply_prepared(
&self,
prepared: Result<crate::orchestration::run_control::PreparedRunCommand, RunControlError>,
) -> ApplicationResult {
let prepared = match prepared {
Ok(prepared) => prepared,
Err(error) => return ApplicationResult::failed(error),
};
let committed = match self.run_control.commit(prepared).await {
Ok(committed) => committed,
Err(error) => return ApplicationResult::failed(error),
};
let outcome = committed.outcome.clone();
let Some(event) = run_outcome_event(&outcome) else {
committed
.activate(self.run_control.scheduler().as_ref())
.await;
return ApplicationResult::run(outcome, None);
};
let revision = self.dispatch_outcome(event).await;
committed
.activate(self.run_control.scheduler().as_ref())
.await;
ApplicationResult::run(outcome, Some(revision))
}
async fn apply_resolve(&self, change_id: &str) -> ApplicationResult {
match self.run_control.prepare_resolve(change_id).await {
Ok(Some(prepared)) => self.apply_prepared(Ok(prepared)).await,
Ok(None) => ApplicationResult::run(
RunControlOutcome::NoOp {
reason: RunNoOpReason::ResolveAlreadyReserved {
change_id: change_id.to_string(),
},
},
None,
),
Err(error) => ApplicationResult::failed(error),
}
}
async fn apply_stop(&self, mode: OperatorMode) -> ApplicationResult {
match self.run_control.stop(mode).await {
Ok(outcome @ RunControlOutcome::StopRequested) => {
let revision = self.dispatch_outcome(ExecutionEvent::Stopping).await;
ApplicationResult::run(outcome, Some(revision))
}
Ok(other) => ApplicationResult::run(other, None),
Err(error) => ApplicationResult::failed(error),
}
}
async fn apply_cancel_stop(&self, mode: OperatorMode) -> ApplicationResult {
match self.run_control.cancel_stop(mode).await {
Ok(outcome @ RunControlOutcome::StopCancelled) => {
let revision = self
.dispatch_outcome(ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::StopCancelled,
})
.await;
ApplicationResult::run(outcome, Some(revision))
}
Ok(other) => ApplicationResult::run(other, None),
Err(error) => ApplicationResult::failed(error),
}
}
async fn apply_force_stop(&self, mode: OperatorMode) -> ApplicationResult {
let outcome = match self.run_control.force_stop(mode).await {
Ok(outcome) => outcome,
Err(error) => return ApplicationResult::failed(error),
};
let RunControlOutcome::ForceStopped {
classification,
awaiting_safe_boundary,
} = &outcome
else {
return ApplicationResult::run(outcome, None);
};
let event = if *awaiting_safe_boundary {
ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::ForceStopAwaitingBoundary {
force_stop: classification.process_report.is_force_stop(),
},
}
} else {
ExecutionEvent::Stopped
};
let revision = self.dispatch_outcome(event).await;
ApplicationResult::run(outcome, Some(revision))
}
async fn apply_mark(&self, change_id: &str, marked: bool) -> ApplicationResult {
match self
.run_control
.operator()
.set_execution_mark(change_id, marked)
.await
{
Ok(outcome) => self.publish_operator_outcome(outcome).await,
Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
}
}
async fn apply_queue_intent(&self, change_id: &str, queued: bool) -> ApplicationResult {
let service = self.run_control.operator();
let result = if queued {
service.add_to_queue(change_id).await
} else {
service.remove_from_queue(change_id).await
};
match result {
Ok(queue) => {
self.publish_operator_outcome(OperatorOutcome::Queue(queue))
.await
}
Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
}
}
async fn apply_bulk_marks(&self) -> ApplicationResult {
match self.run_control.operator().set_all_execution_marks().await {
Ok(outcome) => self.publish_operator_outcome(outcome).await,
Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
}
}
async fn publish_operator_outcome(&self, outcome: OperatorOutcome) -> ApplicationResult {
let Some(event) = operator_outcome_event(&outcome) else {
return ApplicationResult::operator(outcome, None);
};
let revision = self.dispatch_outcome(event).await;
ApplicationResult::operator(outcome, Some(revision))
}
async fn dispatch_outcome(&self, event: ExecutionEvent) -> u64 {
let dispatch_id = self.dispatch.dispatch(event).await;
match &self.revisions {
Some(revisions) => revisions
.revision_for_dispatch(dispatch_id)
.unwrap_or_else(|| revisions.current_revision()),
None => 0,
}
}
fn current_revision(&self) -> u64 {
self.revisions
.as_ref()
.map_or(0, |revisions| revisions.current_revision())
}
}
pub fn bind_mark_settlement(application: &Arc<OperatorApplication>) {
let owned: std::sync::Weak<OperatorApplication> = Arc::downgrade(application);
let runtime: std::sync::Weak<dyn MarkSettlementRuntime> = owned;
application
.run_control
.operator()
.marks()
.settlement()
.bind_runtime(runtime);
}
#[async_trait::async_trait]
impl MarkSettlementRuntime for OperatorApplication {
fn admits_dynamic_queue(&self) -> bool {
self.run_control.scheduler().is_running()
}
async fn settle_marks(&self, targets: Vec<String>) -> MarkSettlementPlan {
let mut plan = self
.run_control
.operator()
.plan_mark_settlement(&targets)
.await;
let mutations: Vec<(String, MarkSettlementAction)> = plan
.additions
.iter()
.map(|change_id| (change_id.clone(), MarkSettlementAction::Add))
.chain(
plan.removals
.iter()
.map(|change_id| (change_id.clone(), MarkSettlementAction::Remove)),
)
.collect();
let mut applied_membership_change = false;
let mut guard_skipped: Vec<(String, MarkSettlementExclusion)> = Vec::new();
for (change_id, action) in mutations {
let application = {
let _guard = self.gate.clone().lock_owned().await;
self.run_control
.operator()
.apply_settlement_queue_intent(&change_id, action)
.await
};
if application.applied() {
applied_membership_change = true;
}
if let Some(reason) = application.skipped {
guard_skipped.push((change_id.clone(), reason));
}
self.publish_operator_outcome(OperatorOutcome::Queue(application.outcome))
.await;
}
for (change_id, reason) in guard_skipped {
plan.additions.retain(|planned| planned != &change_id);
plan.removals.retain(|planned| planned != &change_id);
plan.excluded.push((change_id, reason));
}
if applied_membership_change {
self.run_control
.operator()
.notify_scheduler_after_settlement()
.await;
}
if !plan.excluded.is_empty() {
let skipped = plan
.excluded
.iter()
.map(|(change_id, reason)| format!("{change_id}={}", reason.as_str()))
.collect::<Vec<_>>()
.join(", ");
tracing::debug!("Mark settlement changed no queue intent for: {skipped}");
}
plan
}
async fn report_abandoned_settlement(&self, pending: Vec<String>) {
let targets = if pending.is_empty() {
"no marked change".to_string()
} else {
pending.join(", ")
};
self.dispatch
.dispatch(ExecutionEvent::Log(crate::events::LogEntry::info(format!(
"Mark settlement abandoned because the scheduler ended: {targets}"
))))
.await;
}
async fn report_settlement_failure(
&self,
failure: MarkSettlementFailure,
targets: Vec<String>,
) {
let named = if targets.is_empty() {
"no marked change".to_string()
} else {
targets.join(", ")
};
self.dispatch
.dispatch(ExecutionEvent::Log(crate::events::LogEntry::warn(format!(
"Mark settlement did not admit (reason={}): {named}",
failure.as_str()
))))
.await;
}
}
fn run_outcome_event(outcome: &RunControlOutcome) -> Option<ExecutionEvent> {
match outcome {
RunControlOutcome::RunDispatched {
change_ids,
explicit_retry,
scheduler,
excluded: _,
} => {
debug_assert!(
scheduler.dispatched(),
"an accepted run dispatch must carry scheduler evidence"
);
Some(ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::RunDispatched {
change_ids: change_ids.clone(),
explicit_retry: *explicit_retry,
scheduler_started: matches!(scheduler, SchedulerEffect::Started),
},
})
}
RunControlOutcome::ResolveReserved {
change_id,
reservation,
..
} => Some(ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::ResolveReserved {
change_id: change_id.clone(),
active: matches!(reservation, ResolveReservation::Active),
},
}),
RunControlOutcome::StopRequested
| RunControlOutcome::StopCancelled
| RunControlOutcome::ForceStopped { .. } => None,
RunControlOutcome::NoOp { .. } => None,
}
}
pub(crate) fn operator_outcome_event(outcome: &OperatorOutcome) -> Option<ExecutionEvent> {
match outcome {
OperatorOutcome::MarkSet { change_id, marked } => {
Some(ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::MarkDelta {
change_ids: vec![change_id.clone()],
marked: *marked,
},
})
}
OperatorOutcome::BulkMarks {
marked, changed, ..
} => (!changed.is_empty()).then(|| ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::MarkDelta {
change_ids: changed.clone(),
marked: *marked,
},
}),
OperatorOutcome::Queue(queue) => (queue.reducer_changed || queue.dynamic_queue_mutated)
.then(|| ExecutionEvent::OperatorCommandApplied {
effect: OperatorCommandEffect::QueueDelta {
change_id: queue.change_id.clone(),
queued: matches!(
queue.mutation,
crate::orchestration::operator_command::QueueMutation::Added
),
},
}),
OperatorOutcome::Retry(_) => None,
OperatorOutcome::Dequeued { .. } | OperatorOutcome::ForceStopped { .. } => None,
OperatorOutcome::NoOp { .. } => None,
}
}
#[cfg(test)]
mod tests;
#[cfg(test)]
mod stop_settlement_tests;
#[cfg(all(test, feature = "web-monitoring"))]
mod mode_matrix_tests;