From b8281d6c1f62a3a1b9c8993da654b90171e32070 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Fri, 9 Oct 2026 11:59:19 -0400 Subject: [PATCH] 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 --- lib/apps/fabro-server/src/server.rs | 63 +++++++++---------- .../fabro-server/src/server/petri_runs.rs | 22 +------ lib/components/fabro-petri/src/engine.rs | 43 ++++++------- .../fabro-petri/src/projection/coordinator.rs | 2 +- .../fabro-petri/src/projection/mod.rs | 1 + 5 files changed, 57 insertions(+), 74 deletions(-) diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 8a5aa2b9c..c0069200d 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -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) -> 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, 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, ) { 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, 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, 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, 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, run_id: RunId) { FailureReason::WorkflowError, format!("Failed to load final run state: {err}"), ); - state.scheduler_notify.notify_one(); return; } }; diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 176fd5232..7a09f4743 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, run_id: RunId, status: RunStatus, error: Option) { 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, 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 diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 2f5ac7fc4..f6ed2d194 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -525,25 +525,30 @@ fn outcome( inspection: RunInspection, host_error: Option, ) -> Result { - // 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) diff --git a/lib/components/fabro-petri/src/projection/coordinator.rs b/lib/components/fabro-petri/src/projection/coordinator.rs index 5f624db80..7225dd007 100644 --- a/lib/components/fabro-petri/src/projection/coordinator.rs +++ b/lib/components/fabro-petri/src/projection/coordinator.rs @@ -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 { diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index ac5bea3b1..d65b65341 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -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::{