Settle managed runs through one rule

ManagedRun::settle owns the rule that a run's first terminal status
sticks, and the lifecycle fold, the finish settle, fail_managed_run and
the in-process finish all go through it. One release_managed_run
replaces the two helpers that released a run's live state and also
frees its scheduler slot. The engine reads the finish through the
projection's finished_status, so a cancelled finish carrying a
checkpoint failure, and a failed publication, are mapped in one place.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-09 11:59:19 -04:00
parent 04309977e5
commit b8281d6c1f
5 changed files with 57 additions and 74 deletions

View file

@ -282,6 +282,23 @@ struct ManagedRun {
}
impl ManagedRun {
/// Settle the run on a terminal `status`. The first terminal status
/// sticks: a later one that agrees only fills a missing error, and one
/// that disagrees is ignored. Returns whether the status was applied.
fn settle(&mut self, status: RunStatus, error: Option<String>) -> bool {
if self.status.is_terminal() {
if self.status == status && self.error.is_none() {
self.error = error;
}
return false;
}
self.status = status;
self.error = error;
self.active_steerable_stages.clear();
self.active_non_steerable_stages.clear();
true
}
/// True if cancellation should still escalate to `worker_ref`; clears a
/// stale escalation marker as a side effect.
fn escalation_still_current(&mut self, worker_ref: &WorkerRef) -> bool {
@ -3550,7 +3567,6 @@ pub(crate) async fn persist_run_failure(
}
}
release_managed_run(state, run_id);
state.scheduler_notify.notify_one();
}
/// How long [`persist_run_failure`] waits before each retry of a failed
@ -3649,23 +3665,22 @@ async fn durable_run_status(state: &AppState, run_id: RunId) -> anyhow::Result<O
fn fail_managed_run(state: &Arc<AppState>, run_id: RunId, reason: FailureReason, message: String) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
if !managed_run.status.is_terminal() {
managed_run.status = RunStatus::Failed { reason };
managed_run.error = Some(message);
}
managed_run.settle(RunStatus::Failed { reason }, Some(message));
}
drop(runs);
release_managed_run(state, run_id);
}
/// Drop the run's live worker state and controls, leaving its status alone.
fn release_managed_run(state: &AppState, run_id: RunId) {
/// Drop the run's live worker state and controls and free its scheduler
/// slot, leaving its status alone.
pub(in crate::server) fn release_managed_run(state: &AppState, run_id: RunId) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
clear_live_run_state(managed_run);
}
drop(runs);
cleanup_worker_control_bus_for_run(state, run_id);
state.scheduler_notify.notify_one();
}
/// Fold one lifecycle record of the run's stream into the in-memory run:
@ -3680,9 +3695,8 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
// A settled run is immutable to the lifecycle: the follower still folds
// the records before the terminal one after Petri's finish or the
// worker's terminal record settled the run, and none may reopen it.
if managed_run.status.is_terminal()
&& (!is_terminal_transition(record) || record.status != Some(managed_run.status))
{
// Terminal records go through [`ManagedRun::settle`].
if managed_run.status.is_terminal() && !is_terminal_transition(record) {
return;
}
match record.transition {
@ -3729,23 +3743,17 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
}
RunLifecycleKind::Removing => managed_run.status = RunStatus::Removing,
RunLifecycleKind::Succeeded => {
managed_run.status = record.status.unwrap_or(RunStatus::Succeeded {
let status = record.status.unwrap_or(RunStatus::Succeeded {
reason: SuccessReason::Completed,
});
managed_run.error = None;
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
managed_run.settle(status, None);
cleanup_worker_control_bus_for_run(state, run_id);
}
RunLifecycleKind::Failed | RunLifecycleKind::Dead => {
managed_run.status = record.status.unwrap_or(RunStatus::Failed {
let status = record.status.unwrap_or(RunStatus::Failed {
reason: FailureReason::WorkflowError,
});
if managed_run.error.is_none() {
managed_run.error.clone_from(&record.reason);
}
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
managed_run.settle(status, record.reason.clone());
cleanup_worker_control_bus_for_run(state, run_id);
}
RunLifecycleKind::Runnable
@ -3776,16 +3784,9 @@ pub(in crate::server) fn settle_managed_run_at_finish(
failure: Option<String>,
) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
let Some(managed_run) = runs.get_mut(&run_id) else {
return;
};
if managed_run.status.is_terminal() {
return;
if let Some(managed_run) = runs.get_mut(&run_id) {
managed_run.settle(status, failure);
}
managed_run.status = status;
managed_run.error = failure;
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
}
/// Settle the in-memory run at the terminal lifecycle record its worker
@ -4234,7 +4235,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
"Run not found at launch".to_string(),
);
state.scheduler_notify.notify_one();
return;
}
Err(err) => {
@ -4245,7 +4245,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
format!("Failed to load run state: {err}"),
);
state.scheduler_notify.notify_one();
return;
}
};
@ -4394,7 +4393,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
"The run's final state is missing from the store".to_string(),
);
state.scheduler_notify.notify_one();
return;
}
Err(err) => {
@ -4405,7 +4403,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
format!("Failed to load final run state: {err}"),
);
state.scheduler_notify.notify_one();
return;
}
};

View file

@ -824,7 +824,7 @@ async fn commit_and_finish(
Ok(_) => finish(state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
release_live_state(state, run_id);
super::release_managed_run(state, run_id);
}
}
}
@ -837,26 +837,10 @@ async fn commit_and_finish(
fn finish(state: &Arc<AppState>, run_id: RunId, status: RunStatus, error: Option<String>) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
if !managed_run.status.is_terminal() || managed_run.status == status {
managed_run.status = status;
if managed_run.error.is_none() {
managed_run.error = error;
}
}
managed_run.settle(status, error);
}
drop(runs);
release_live_state(state, run_id);
}
/// Release controls after an append failure without claiming a new terminal
/// result. A Petri finish already committed remains authoritative.
fn release_live_state(state: &Arc<AppState>, run_id: RunId) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
clear_live_run_state(managed_run);
}
drop(runs);
state.scheduler_notify.notify_one();
super::release_managed_run(state, run_id);
}
/// The run store an in-process run executes over. A coordinator finish

View file

@ -525,25 +525,30 @@ fn outcome(
inspection: RunInspection,
host_error: Option<HostError>,
) -> Result<RunOutcome, RunError> {
// A failed checkpoint cancelled the run; its finish records the
// checkpoint failure, which is what Fabro reports.
let checkpoint_failed = inspection
.finalization_failure
.as_ref()
.is_some_and(projection::is_checkpoint_failure);
let status = match inspection.status.as_deref() {
Some("success") => RunStatus::Success,
Some("failed") => RunStatus::Failed,
Some("cancelled") if checkpoint_failed => RunStatus::Failed,
Some("cancelled") => RunStatus::Cancelled,
_ => {
let mut reasons = inspection.incomplete.clone();
if let Some(error) = host_error {
reasons.push(error.to_string());
}
return Err(RunError::Unfinished(reasons));
let Some(recorded) = inspection
.status
.as_deref()
.filter(|status| matches!(*status, "success" | "failed" | "cancelled"))
else {
let mut reasons = inspection.incomplete.clone();
if let Some(error) = host_error {
reasons.push(error.to_string());
}
return Err(RunError::Unfinished(reasons));
};
// The projection's reading of the finish, so the worker's return and
// the API agree: a failed checkpoint's cancellation is a failure.
let (status, publish_failed) =
match projection::finished_status(recorded, inspection.finalization_failure.as_ref()) {
fabro_types::RunStatus::Succeeded { .. } => (RunStatus::Success, false),
fabro_types::RunStatus::Failed {
reason: FailureReason::Cancelled,
} => (RunStatus::Cancelled, false),
fabro_types::RunStatus::Failed { reason } => {
(RunStatus::Failed, reason == FailureReason::PublishFailed)
}
_ => (RunStatus::Failed, false),
};
let execution_failure = inspection
.invocations
.iter()
@ -551,10 +556,6 @@ fn outcome(
.and_then(|root| root.result.as_ref())
.and_then(|result| result.failure.as_ref())
.map(|failure| failure.message.clone());
let publish_failed = inspection
.finalization_failure
.as_ref()
.is_some_and(projection::is_publish_failure);
let failure = inspection
.finalization_failure
.map(|failure| failure.message)

View file

@ -229,7 +229,7 @@ impl RunView {
/// records (`success`, `cancelled`, or a failure). A failed checkpoint
/// cancels the run, but the run failed: its finish says so with the
/// checkpoint's finalization failure.
pub(super) fn finished_status(
pub(crate) fn finished_status(
status: &str,
finalization_failure: Option<&FinalizationFailure>,
) -> RunStatus {

View file

@ -37,6 +37,7 @@ use std::fmt;
use std::str::FromStr;
use chrono::{DateTime, TimeZone as _, Utc};
pub(crate) use coordinator::finished_status;
use fabro_store::StagePosition;
use fabro_store::platform_records::StoredPlatformRecord;
use fabro_types::{