use std::sync::Arc;
use async_trait::async_trait;
use crate::orchestration::operator_command::{
MarkExclusion, NoOpReason, OperatorCommandError, OperatorOutcome, StopSettlement,
};
use crate::orchestration::operator_coordinator::{
ApplicationOutcome, ApplicationResult, OperatorApplication, OperatorIntent,
};
use crate::orchestration::run_control::{
ExcludedTarget, ResolveReservation, RunControlError, RunControlOutcome, RunNoOpReason,
SchedulerEffect,
};
use crate::web::state::WebState;
use super::dto::{
ApplyCommitEvidence, CommandResult, CommandSpec, ErrorCode, ExecutionPhase as DtoPhase,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ExecutionSummary {
pub changed: bool,
pub detail: Option<String>,
pub result_revision: Option<u64>,
pub result: Option<CommandResult>,
}
impl ExecutionSummary {
pub fn changed(detail: impl Into<String>) -> Self {
Self {
changed: true,
detail: Some(detail.into()),
result_revision: None,
result: None,
}
}
pub fn no_op(detail: impl Into<String>) -> Self {
Self {
changed: false,
detail: Some(detail.into()),
result_revision: None,
result: None,
}
}
pub fn at_revision(mut self, revision: Option<u64>) -> Self {
self.result_revision = revision;
self
}
pub fn with_result(mut self, result: CommandResult) -> Self {
self.result = Some(result);
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommandFailure {
pub error_code: ErrorCode,
pub message: String,
pub result_revision: Option<u64>,
}
impl CommandFailure {
pub fn new(error_code: ErrorCode, message: impl Into<String>) -> Self {
Self {
error_code,
message: message.into(),
result_revision: None,
}
}
pub fn at_revision(mut self, revision: Option<u64>) -> Self {
self.result_revision = revision;
self
}
}
pub type PendingCommand = std::pin::Pin<
Box<dyn std::future::Future<Output = Result<ExecutionSummary, CommandFailure>> + Send>,
>;
pub enum Applied {
Ordinary(Option<GateGuard>),
Pending(PendingCommand),
Settled(Result<ExecutionSummary, CommandFailure>),
}
pub type GateGuard = tokio::sync::OwnedMutexGuard<()>;
#[async_trait]
pub trait RemoteControlExecutor: Send + Sync {
async fn execute(&self, command: &CommandSpec) -> Result<ExecutionSummary, CommandFailure>;
async fn is_command_capable(&self) -> bool {
true
}
async fn begin(&self, _command: &CommandSpec, gate: Option<GateGuard>) -> Applied {
Applied::Ordinary(gate)
}
async fn execute_held(
&self,
command: &CommandSpec,
gate: Option<GateGuard>,
) -> Result<ExecutionSummary, CommandFailure> {
drop(gate);
self.execute(command).await
}
}
pub fn map_operator_error(error: &OperatorCommandError) -> CommandFailure {
match error {
OperatorCommandError::MissingCancellationHandle { .. }
| OperatorCommandError::RetryUnsupported { .. }
| OperatorCommandError::ForceStopIneligible { .. }
| OperatorCommandError::UnchangedAcceptanceInput { .. } => {
CommandFailure::new(ErrorCode::TargetIneligible, error.to_string())
}
OperatorCommandError::TerminationTimeout { .. }
| OperatorCommandError::ForceStopUnconfirmed { .. } => {
CommandFailure::new(ErrorCode::RootBusy, error.to_string())
}
OperatorCommandError::CancellationFailed { .. } => {
CommandFailure::new(ErrorCode::InternalError, error.to_string())
}
}
}
pub fn summarize_outcome(outcome: &OperatorOutcome) -> ExecutionSummary {
match outcome {
OperatorOutcome::MarkSet { change_id, marked } => {
ExecutionSummary::changed(format!("execution mark for '{change_id}' set to {marked}"))
}
OperatorOutcome::Queue(queue) => {
if queue.reducer_changed || queue.dynamic_queue_mutated {
ExecutionSummary::changed(format!(
"queue intent for '{}' is now '{}'",
queue.change_id, queue.display_status
))
} else {
ExecutionSummary::no_op(format!(
"queue intent for '{}' already '{}'",
queue.change_id, queue.display_status
))
}
}
OperatorOutcome::Dequeued {
change_id,
settlement,
} => ExecutionSummary::changed(settlement.describe(change_id))
.with_result(stop_result(settlement)),
OperatorOutcome::ForceStopped {
change_id,
execution_id,
terminated,
settlement,
} => ExecutionSummary::changed(describe_force_stop(change_id, *terminated, settlement))
.with_result(force_stop_result(
change_id,
execution_id.as_deref(),
*terminated,
settlement,
)),
OperatorOutcome::Retry(plan) => {
if plan.is_empty() {
ExecutionSummary::no_op("no change carried retryable evidence")
} else {
ExecutionSummary::changed(format!("retry accepted for {:?}", plan.change_ids))
}
}
OperatorOutcome::BulkMarks {
marked,
changed,
excluded,
} => {
let action = if *marked { "marked" } else { "unmarked" };
let mut detail = format!("{} change(s) {action}: {changed:?}", changed.len());
if !excluded.is_empty() {
detail.push_str(&format!(
", {} excluded ({})",
excluded.len(),
summarize_exclusions(excluded)
));
}
ExecutionSummary::changed(detail)
}
OperatorOutcome::NoOp { change_id, reason } => {
let why = match reason {
NoOpReason::MarkUnchanged => "execution mark already had the requested value",
NoOpReason::TerminalMarkTarget => {
"the target is terminal and carries no next-run intent"
}
NoOpReason::ArchiveCompleteMarkTarget => {
"the target has completed archive and carries no next-run intent"
}
NoOpReason::ReducerRejected => "the reducer produced no state change",
NoOpReason::BulkMarksUnchanged => {
"every eligible change already carried the derived mark"
}
NoOpReason::NoEligibleMarkTarget => {
"no change is eligible for a bulk execution-mark mutation"
}
};
if change_id.is_empty() {
ExecutionSummary::no_op(why)
} else {
ExecutionSummary::no_op(format!("'{change_id}': {why}"))
}
}
}
}
pub fn stop_result(settlement: &StopSettlement) -> CommandResult {
CommandResult::StopAndDequeue {
cancelled_phase: DtoPhase::from_shared(settlement.cancelled_phase),
last_completed_phase: settlement.last_completed_phase.map(DtoPhase::from_shared),
apply_commit: ApplyCommitEvidence {
present: settlement.apply_commit_present,
oid: settlement.apply_commit_oid.clone(),
},
effects_rolled_back: false,
}
}
pub fn force_stop_result(
change_id: &str,
execution_id: Option<&str>,
terminated: bool,
settlement: &StopSettlement,
) -> CommandResult {
CommandResult::ForceStopChange {
change_id: change_id.to_string(),
execution_id: execution_id.map(str::to_string),
cancelled_phase: DtoPhase::from_shared(settlement.cancelled_phase),
last_completed_phase: settlement.last_completed_phase.map(DtoPhase::from_shared),
terminated,
apply_commit: ApplyCommitEvidence {
present: settlement.apply_commit_present,
oid: settlement.apply_commit_oid.clone(),
},
effects_rolled_back: false,
}
}
fn describe_force_stop(change_id: &str, terminated: bool, settlement: &StopSettlement) -> String {
let how = if terminated {
"its managed process group was killed immediately and confirmed reaped"
} else {
"it owned no managed process, so it was dequeued without signalling anything"
};
format!("{}; {how}", settlement.describe(change_id))
}
fn summarize_exclusions(excluded: &[(String, MarkExclusion)]) -> String {
MarkExclusion::ALL
.iter()
.filter_map(|reason| {
let count = excluded
.iter()
.filter(|(_, actual)| actual == reason)
.count();
(count > 0).then(|| format!("{count} {}", reason.as_str()))
})
.collect::<Vec<_>>()
.join(", ")
}
pub fn map_run_control_error(error: &RunControlError) -> CommandFailure {
match error {
RunControlError::InvalidMode { .. } => {
CommandFailure::new(ErrorCode::LifecycleConflict, error.to_string())
}
RunControlError::NoEligibleTarget { .. } | RunControlError::TargetIneligible { .. } => {
CommandFailure::new(ErrorCode::TargetIneligible, error.to_string())
}
RunControlError::DispatchFailed { .. } => {
CommandFailure::new(ErrorCode::InternalError, error.to_string())
}
RunControlError::Operator(inner) => map_operator_error(inner),
}
}
pub fn summarize_run_outcome(outcome: &RunControlOutcome) -> ExecutionSummary {
match outcome {
RunControlOutcome::RunDispatched {
change_ids,
explicit_retry,
scheduler,
excluded,
} => {
let how = match scheduler {
SchedulerEffect::Started => "started the scheduler",
SchedulerEffect::Notified => "woke the running scheduler",
SchedulerEffect::None => "dispatched no scheduler work",
};
let kind = if *explicit_retry { "retry" } else { "run" };
let left_behind = if excluded.is_empty() {
String::new()
} else {
let detail = excluded
.iter()
.map(ExcludedTarget::describe)
.collect::<Vec<_>>()
.join(", ");
format!("; excluded: {detail}")
};
ExecutionSummary::changed(format!("{kind} for {change_ids:?} {how}{left_behind}"))
}
RunControlOutcome::StopRequested => {
ExecutionSummary::changed("graceful stop requested; the run stops at its next boundary")
}
RunControlOutcome::StopCancelled => {
ExecutionSummary::changed("pending graceful stop was withdrawn")
}
RunControlOutcome::ForceStopped {
classification,
awaiting_safe_boundary,
} => {
let report = if classification.process_report.is_force_stop() {
"force stop applied to active execution"
} else {
"run cancelled; no agent execution was active"
};
let boundary = if *awaiting_safe_boundary {
"; waiting for in-flight work to reach a safe stop boundary"
} else {
""
};
ExecutionSummary::changed(format!("{report}{boundary}"))
}
RunControlOutcome::ResolveReserved {
change_id,
reservation,
..
} => match reservation {
ResolveReservation::Active => ExecutionSummary::changed(format!(
"merge resolution for '{change_id}' is the active resolve"
)),
ResolveReservation::Queued { position } => ExecutionSummary::changed(format!(
"merge resolution for '{change_id}' is queued at position {position}"
)),
},
RunControlOutcome::NoOp { reason } => match reason {
RunNoOpReason::ResolveAlreadyReserved { change_id } => ExecutionSummary::no_op(
format!("'{change_id}' already holds a resolve reservation"),
),
RunNoOpReason::NoRetryableTarget => {
ExecutionSummary::no_op("no marked change carried retryable evidence")
}
},
}
}
pub struct SharedServiceExecutor {
application: Arc<OperatorApplication>,
web_state: Arc<WebState>,
worktrees: Arc<dyn super::worktrees::WorktreeOperations>,
}
impl SharedServiceExecutor {
pub fn new(application: Arc<OperatorApplication>, web_state: Arc<WebState>) -> Self {
Self {
application,
web_state,
worktrees: Arc::new(super::worktrees::UnboundWorktreeOperations),
}
}
pub fn with_worktrees(
mut self,
worktrees: Arc<dyn super::worktrees::WorktreeOperations>,
) -> Self {
self.worktrees = worktrees;
self
}
fn intent(command: &CommandSpec) -> Option<OperatorIntent> {
Some(match command {
CommandSpec::Start => OperatorIntent::Start,
CommandSpec::Stop => OperatorIntent::Stop,
CommandSpec::CancelStop => OperatorIntent::CancelStop,
CommandSpec::ForceStop => OperatorIntent::ForceStop,
CommandSpec::SetExecutionMark { change_id, marked } => {
OperatorIntent::SetExecutionMark {
change_id: change_id.clone(),
marked: *marked,
}
}
CommandSpec::SetQueueIntent { change_id, queued } => OperatorIntent::SetQueueIntent {
change_id: change_id.clone(),
queued: *queued,
},
CommandSpec::RetryChange { change_id } => OperatorIntent::RetryChange {
change_id: change_id.clone(),
},
CommandSpec::RetryErrors { change_ids } => OperatorIntent::RetryErrors {
change_ids: change_ids.clone(),
},
CommandSpec::StopAndDequeue { change_id } => OperatorIntent::StopAndDequeue {
change_id: change_id.clone(),
},
CommandSpec::ForceStopChange { change_id } => OperatorIntent::ForceStopChange {
change_id: change_id.clone(),
},
CommandSpec::ResolveMerge { change_id } => OperatorIntent::ResolveMerge {
change_id: change_id.clone(),
},
CommandSpec::SetAllExecutionMarks {} => OperatorIntent::SetAllExecutionMarks,
CommandSpec::CreateWorktree { .. }
| CommandSpec::DeleteWorktree { .. }
| CommandSpec::MergeWorktree { .. } => return None,
})
}
async fn worktree(&self, command: &CommandSpec) -> Result<ExecutionSummary, CommandFailure> {
let outcome = match command {
CommandSpec::CreateWorktree { target, .. } => {
self.worktrees.create(&target.change_id).await
}
CommandSpec::DeleteWorktree { target, .. } => {
self.worktrees.delete(&target.worktree_id).await
}
CommandSpec::MergeWorktree { target, .. } => {
self.worktrees.merge(&target.worktree_id).await
}
other => {
debug_assert!(false, "not a worktree command: {other:?}");
return Err(CommandFailure::new(
ErrorCode::InternalError,
"command is not a worktree operation",
));
}
};
if matches!(&outcome, Ok(summary) if summary.changed) {
self.web_state.sync_remote_control_projection().await;
}
outcome
}
}
#[doc(hidden)]
#[allow(dead_code)] pub fn wired_for_test(
reducer: Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
run_control: Arc<crate::orchestration::run_control::RunControlService>,
web_state: Arc<WebState>,
core_mode: Arc<crate::orchestration::operator_coordinator::CoreMode>,
) -> (SharedServiceExecutor, Arc<OperatorApplication>) {
let dispatcher = Arc::new(
crate::events::EventDispatcher::new(
reducer,
vec![Arc::new(crate::web::state::WebEventSink::new(
web_state.clone(),
))],
)
.with_core_mode(Some(core_mode.clone())),
);
let revisions: Arc<dyn crate::events::OutcomeRevisions> = web_state.clone();
let application = Arc::new(
OperatorApplication::new(core_mode.clone(), run_control, dispatcher)
.with_revisions(Some(revisions)),
);
(
SharedServiceExecutor::new(application.clone(), web_state),
application,
)
}
pub fn summarize_application(
result: ApplicationResult,
) -> Result<ExecutionSummary, CommandFailure> {
let ApplicationResult { outcome, revision } = result;
match outcome {
Ok(ApplicationOutcome::Run(outcome)) => {
Ok(summarize_run_outcome(&outcome).at_revision(revision))
}
Ok(ApplicationOutcome::Operator(outcome)) => {
Ok(summarize_outcome(&outcome).at_revision(revision))
}
Err(error) => Err(map_run_control_error(&error).at_revision(revision)),
}
}
#[async_trait]
impl RemoteControlExecutor for SharedServiceExecutor {
async fn execute(&self, command: &CommandSpec) -> Result<ExecutionSummary, CommandFailure> {
match Self::intent(command) {
Some(intent) => summarize_application(self.application.apply(intent).await),
None => self.worktree(command).await,
}
}
async fn execute_held(
&self,
command: &CommandSpec,
gate: Option<GateGuard>,
) -> Result<ExecutionSummary, CommandFailure> {
match (Self::intent(command), gate) {
(Some(intent), Some(gate)) => {
summarize_application(self.application.apply_held(intent, gate).await)
}
(Some(intent), None) => summarize_application(self.application.apply(intent).await),
(None, gate) => {
drop(gate);
self.worktree(command).await
}
}
}
async fn begin(&self, command: &CommandSpec, gate: Option<GateGuard>) -> Applied {
match command {
CommandSpec::StopAndDequeue { change_id } => {
let pending = match self.application.begin_stop_and_dequeue(change_id).await {
Ok(pending) => pending,
Err(error) => {
return Applied::Settled(Err(map_run_control_error(&error)));
}
};
drop(gate);
let application = self.application.clone();
Applied::Pending(Box::pin(async move {
summarize_application(application.settle_stop_and_dequeue(pending).await)
}))
}
CommandSpec::ForceStopChange { change_id } => {
let pending = match self.application.begin_force_stop_change(change_id).await {
Ok(pending) => pending,
Err(error) => return Applied::Settled(Err(map_run_control_error(&error))),
};
drop(gate);
let application = self.application.clone();
Applied::Pending(Box::pin(async move {
summarize_application(application.settle_force_stop_change(pending).await)
}))
}
_ => Applied::Ordinary(gate),
}
}
}