From 37116758fa8a3d51295095730d06985bcf3363dd Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:18:58 -0400 Subject: [PATCH] Resume a Petri run whose store failed, as after a crash Petri now ends a run's lifetime at its first failed store write: it records nothing after it, fails no firing for it, and returns CoordinatorError::StoreFailed. The run is not over; the next lifetime resumes it from what the store holds. Fabro read that error as an unfinished run and failed it. - engine: RunError::StoreFailed, returned without reading the record back, and Conclusion::Interrupted for it. - worker: an interrupted run gets no terminal lifecycle record; the worker exits with EX_TEMPFAIL (75, the new ExitClass::Interrupted). - server: WorkerExit carries the exit code. An interrupted worker's run goes back to the scheduler in resume mode through the relaunch a restart takes (lease release, recovery, start_requested + runnable), now shared with reconcile_on_startup. The in-process path does the same. A run is resumed at most MAX_STORE_INTERRUPTIONS (3) times per server; the next interruption fails it. A pending cancel, a run that ended or was deleted, and a shutdown also end it as before. Tests: an engine run over a store whose first lease write fails is interrupted with no finish, and a resume finishes it; the server relaunches an interrupted worker in resume mode, fails the run after the bound, and fails a worker that exits 1 as before; exit code 75. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 18 +- lib/apps/fabro-server/src/petri_runs.rs | 136 ++++++++++- lib/apps/fabro-server/src/server.rs | 37 ++- .../fabro-server/src/server/petri_runs.rs | 222 ++++++++++++++---- lib/apps/fabro-server/src/worker_runtime.rs | 12 + lib/components/fabro-petri/src/engine.rs | 41 +++- lib/components/fabro-petri/src/recovery.rs | 2 +- lib/components/fabro-petri/tests/resume.rs | 161 ++++++++++++- lib/foundation/fabro-util/src/exit.rs | 26 +- 9 files changed, 568 insertions(+), 87 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 0b1e8ed13..040f33ef2 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -15,7 +15,10 @@ //! from the graphs when a crash cut its creation short. Either way the //! worker records the lifecycle transitions Fabro's read side needs //! (`starting`, `running`, then `succeeded` or `failed`) as platform -//! records through the client. +//! records through the client. A failed write to the run's store ends the +//! run's lifetime and not the run: the worker records no end for it and +//! exits with `EX_TEMPFAIL` (75, [`ExitClass::Interrupted`]), and the server +//! launches a worker to resume it, as after a crash. //! //! The server's controls arrive over the control channel and go to Petri //! through [`PetriControls`]: cancel (and `SIGTERM`/`SIGINT`) fires one @@ -94,6 +97,7 @@ use fabro_store::platform_records::{ }; use fabro_types::settings::run::{ApprovalMode, RunMode}; use fabro_types::{FailureReason, Principal, RunId, RunNoticeLevel, RunStatus, SuccessReason}; +use fabro_util::exit::{ErrorExt as _, ExitClass}; use fabro_vault::Vault; use fabro_workflow::Error as WorkflowError; use fabro_workflow::services::FabroRunToolServices; @@ -124,7 +128,9 @@ pub(super) struct PetriWorker<'a> { /// Execute the run to its end. `Ok` when the record says it succeeded; /// the failure otherwise, after the terminal event is appended, so the -/// worker exits as the legacy worker does for a failed run. +/// worker exits as the legacy worker does for a failed run. A run its +/// store interrupted gets no terminal event, and its error is classified +/// [`ExitClass::Interrupted`] for the server to resume the run. pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { let run_id = worker.run_id; let admission = worker.run_state.spec.admission.clone(); @@ -290,6 +296,14 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { "Petri run ended" ); let (record, phase, failure) = match engine::conclusion(&result) { + Conclusion::Interrupted { message } => { + // The run is not over: it continues from its records in the + // worker the server launches next. That also holds when the + // control channel was lost: a cancel the loss requested is in + // the records, or the run was not cancelled. + warn!(run_id = %run_id, error = %message, "Petri run interrupted; the server resumes it"); + return Err(anyhow!("{message}").classify(ExitClass::Interrupted)); + } Conclusion::Succeeded => { info!(run_id = %run_id, "Petri run completed"); ( diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index 1690fbe0e..c008f69d5 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -174,6 +174,7 @@ mod tests { use fabro_types::{ FailureReason, RunId, RunStatus, SuccessReason, WorkflowPath, WorkflowVersion, }; + use fabro_util::exit::ExitClass; use serde_json::json; use tokio::io::AsyncRead; use tokio::sync::Notify; @@ -181,6 +182,7 @@ mod tests { use tower::ServiceExt as _; use super::*; + use crate::server::petri_runs::MAX_STORE_INTERRUPTIONS; use crate::server::{ AppState, reconcile_incomplete_runs_on_startup, run_records, spawn_scheduler, }; @@ -209,6 +211,8 @@ mod tests { started: Notify, running: AtomicBool, exit: Arc, + /// The exit code the worker ends with; `None` for a signal. + code: Arc>>, mode: Mutex>, } @@ -220,6 +224,11 @@ mod tests { } fn end_worker(&self) { + self.end_worker_with(None); + } + + fn end_worker_with(&self, code: Option) { + *sync::lock(&self.code) = code; self.running.store(false, Ordering::SeqCst); self.exit.notify_one(); } @@ -235,15 +244,18 @@ mod tests { *sync::lock(&self.mode) = Some(spec.mode); self.running.store(true, Ordering::SeqCst); let exit = Arc::clone(&self.exit); + let code = Arc::clone(&self.code); let stderr: Pin> = Box::pin(tokio::io::empty()); let started = StartedWorker { worker_ref: WorkerRef::Local { pid: u32::MAX }, stderr, wait: Box::pin(async move { exit.notified().await; + let code = *sync::lock(&code); Ok(WorkerExit { - success: false, - detail: "test worker ended without a terminal event".to_string(), + success: code == Some(0), + code, + detail: "test worker ended without a terminal event".to_string(), }) }), }; @@ -804,4 +816,124 @@ mod tests { .expect("the run projects"); assert_eq!(run_state.status, succeeded); } + + /// Wait until the run's stored status is terminal, and return it with + /// the failure's message. + async fn terminal_status(state: &Arc, run_id: RunId) -> (RunStatus, Option) { + for _ in 0..1000 { + let run_state = run_records::projection(state, run_id) + .await + .expect("the run state loads") + .expect("the run projects"); + if run_state.status.is_terminal() { + let message = run_state + .conclusion + .as_ref() + .and_then(|conclusion| conclusion.failure.as_ref()) + .map(|failure| failure.detail.message.clone()); + return (run_state.status, message); + } + time::sleep(Duration::from_millis(10)).await; + } + panic!("the run never ended"); + } + + /// A worker whose run's store failed exits with `EX_TEMPFAIL`: the run + /// is not over, and goes back to a worker in resume mode, as after a + /// crash. The previous worker's lease is released first. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_run_its_store_interrupted_goes_back_to_a_worker_in_resume_mode() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + assert_eq!(runtime.launched_mode(), Some("start")); + + runtime.end_worker_with(Some(ExitClass::Interrupted.code())); + runtime.wait_for_start().await; + + assert_eq!(runtime.launched_mode(), Some("resume")); + assert_eq!( + state + .petri_runs + .store() + .owner(&PetriRuns::key(&run_id)) + .await + .expect("reads the lease"), + None, + "the interrupted worker's lease is released" + ); + let transitions = state + .stores + .run_summaries + .platform_records() + .read(&run_id) + .await + .expect("the records list") + .into_iter() + .filter_map(|stored| match stored.record { + PlatformRecord::RunLifecycle(record) => Some(record.transition), + _ => None, + }) + .collect::>(); + assert_eq!( + &transitions[transitions.len() - 4..], + [ + RunLifecycleKind::Starting, + RunLifecycleKind::Running, + RunLifecycleKind::StartRequested, + RunLifecycleKind::Runnable + ], + "the run was asked to start again as a resume, and did not end: {transitions:?}" + ); + runtime.end_worker(); + } + + /// A store that keeps failing is not a glitch a resume gets past: the + /// run resumes after each of its first interruptions, and the next one + /// fails it. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_run_its_store_keeps_interrupting_fails() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + + for _ in 0..MAX_STORE_INTERRUPTIONS { + runtime.end_worker_with(Some(ExitClass::Interrupted.code())); + runtime.wait_for_start().await; + assert_eq!(runtime.launched_mode(), Some("resume")); + } + runtime.end_worker_with(Some(ExitClass::Interrupted.code())); + + let (status, message) = terminal_status(&state, run_id).await; + assert_eq!(status, RunStatus::Failed { + reason: FailureReason::Terminated, + }); + let message = message.expect("the failure has a message"); + assert!( + message.contains(&format!( + "The run's store failed {} times", + MAX_STORE_INTERRUPTIONS + 1 + )), + "{message}" + ); + } + + /// Any other failed worker exit fails the run, as before. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_worker_that_fails_otherwise_fails_the_run() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + + runtime.end_worker_with(Some(1)); + + let (status, message) = terminal_status(&state, run_id).await; + assert_eq!(status, RunStatus::Failed { + reason: FailureReason::Terminated, + }); + assert!( + message.is_some_and(|message| message.contains("Worker exited before emitting")), + "the failure names the worker's exit" + ); + } } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 2a6046844..5c3933cab 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -156,9 +156,7 @@ use crate::worker_control::{ LocalWorkerControlBus, WORKER_CONTROL_ACK_WAIT, WorkerControlAcks, WorkerControlBus, WorkerControlBusError, }; -use crate::worker_runtime::{ - LocalWorkerRuntime, WorkerExit, WorkerLaunchSpec, WorkerRef, WorkerRuntime, -}; +use crate::worker_runtime::{LocalWorkerRuntime, WorkerLaunchSpec, WorkerRef, WorkerRuntime}; use crate::worker_token::{WorkerScopeSet, WorkerTokenKeys, issue_worker_token_with_scopes}; use crate::{ canonical_host, demo, diagnostics, run_manifest, security_headers, static_files, web_auth, @@ -277,6 +275,10 @@ struct ManagedRun { cancel_escalation_worker: Option, run_dir: Option, execution_mode: RunExecutionMode, + /// How many times the run's store has interrupted it in this server's + /// life. Each time, the run resumes, up to + /// `petri_runs::MAX_STORE_INTERRUPTIONS`. + store_interruptions: u32, } impl ManagedRun { @@ -3535,6 +3537,7 @@ fn managed_run( cancel_escalation_worker: None, run_dir: Some(run_dir), execution_mode, + store_interruptions: 0, } } @@ -3784,8 +3787,9 @@ async fn fail_worker_launch(state: &Arc, run_id: RunId, err: anyhow::E state.scheduler_notify.notify_one(); } -/// A worker that exited without recording the run's end left it failed. -async fn append_worker_exit_failure(state: &AppState, run_id: RunId, worker_exit: &WorkerExit) { +/// A worker that exited without recording the run's end left it failed, +/// with `failure` as the reason unless a cancel was pending. +async fn append_worker_exit_failure(state: &AppState, run_id: RunId, failure: String) { let run_state = match run_records::projection(state, run_id).await { Ok(Some(run_state)) => run_state, Ok(None) => return, @@ -3798,13 +3802,7 @@ async fn append_worker_exit_failure(state: &AppState, run_id: RunId, worker_exit return; } - let (error, reason) = failure_for_incomplete_run( - run_state.pending_control, - format!( - "Worker exited before emitting a terminal run event: {}", - worker_exit.detail - ), - ); + let (error, reason) = failure_for_incomplete_run(run_state.pending_control, failure); if let Err(err) = run_records::lifecycle( state, run_id, @@ -4288,7 +4286,20 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { // API drop here, so its lease never outlives it. state.petri_runs.worker_exited(run_id); state.petri_projector.signal(run_id); - append_worker_exit_failure(&state, run_id, &worker_exit).await; + // A worker whose run's store failed ended the run's lifetime, not the + // run: the run resumes in a new worker, as after a crash. + let failure = if worker_exit.interrupted() { + match petri_runs::resume_after_interruption(&state, run_id, &worker_exit.detail).await { + Ok(()) => return, + Err(failure) => failure, + } + } else { + format!( + "Worker exited before emitting a terminal run event: {}", + worker_exit.detail + ) + }; + append_worker_exit_failure(&state, run_id, failure).await; let final_state = match run_records::projection(&state, run_id).await { Ok(Some(final_state)) => final_state, diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 0d059cace..0be8d2675 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -30,7 +30,11 @@ //! previous server left in flight back to a worker in resume mode, once the //! recovery protocol (`fabro_petri::recovery`) has brought every live //! workspace to the snapshot its durable state names, or reports the run -//! failed when it cannot. +//! failed when it cannot. A run whose store failed under it takes the same +//! way back ([`resume_after_interruption`]): its worker exits with +//! `EX_TEMPFAIL` and records no end, since a failed store write ends the +//! run's lifetime and not the run, and the server resumes it, at most +//! [`MAX_STORE_INTERRUPTIONS`] times. use std::collections::HashMap; use std::sync::Arc; @@ -57,7 +61,9 @@ use fabro_static::EnvVars; use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord}; use fabro_types::settings::McpTransport; use fabro_types::settings::run::{ApprovalMode, McpServerSettings, RunMode}; -use fabro_types::{FailureReason, RunId, RunRunnableSource, RunStatus, RunTarget, SuccessReason}; +use fabro_types::{ + FailureReason, RunControlAction, RunId, RunRunnableSource, RunStatus, RunTarget, SuccessReason, +}; use fabro_util::error as error_util; use fabro_workflow::Error as WorkflowError; use lithos_llm::catalog::ProviderId; @@ -70,7 +76,6 @@ use super::{ stream_follower, }; use crate::petri_check; -use crate::petri_runs::PetriRuns; use crate::run_compiler::{AdmittedRun, PreparedRun, RunCompilerError}; /// The runtime Petri gets, at create and at execution: the server's run @@ -520,6 +525,14 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { "Petri run ended" ); let (status, error, record) = match engine::conclusion(&result) { + // The run's store failed: the run resumes, as after a crash, in a + // new in-process run the scheduler starts, or fails when it cannot. + Conclusion::Interrupted { message } => { + match resume_after_interruption(&state, run_id, &message).await { + Ok(()) => return, + Err(failure) => failed(FailureReason::WorkflowError, failure), + } + } Conclusion::Succeeded => { info!(run_id = %run_id, "Petri run completed"); ( @@ -553,29 +566,157 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { } } +/// How many times the run's store may interrupt a run in one server's life. +/// Each interruption resumes the run; one more fails it, since a store that +/// keeps failing is not a glitch a resume gets past. +pub(crate) const MAX_STORE_INTERRUPTIONS: u32 = 3; + /// Bring a Petri run the server left in flight back to its worker after a /// restart: the run continues from its records, as Petri's own resume does, -/// on workspaces that match them. -/// -/// The lease the previous worker held is released from outside, which -/// fences that worker should it still be alive. Then the recovery protocol -/// reads the run's durable execution state: a run with a failed checkpoint -/// is reported failed here and never resumed; otherwise every live -/// workspace on this host is verified against, reset to, or restored from -/// the snapshot its last durable finish names, and a finish with no -/// snapshot fails the run rather than resume it on stale files. The run is -/// then asked to start again as a resume (`run.start_requested` with -/// `resume`, then `run.runnable`, the same pair the API's resume appends), -/// and a managed run is registered for the scheduler in resume mode when -/// Petri's store holds the run, else in start mode: a worker that died -/// before it created the run's record left nothing to continue from, so the -/// run starts from its admitted graphs. +/// on workspaces that match them ([`relaunch`]). A run that cannot continue +/// is reported failed here and never resumed. pub(crate) async fn reconcile_on_startup( state: &Arc, run_id: RunId, run_state: &fabro_store::RunProjection, ) -> anyhow::Result<()> { - let key = PetriRuns::key(&run_id); + let mode = match relaunch(state, run_id, run_state).await? { + Relaunch::Worker(mode) => mode, + Relaunch::Failed { reason } => { + warn!( + run_id = %run_id, + error = %reason, + "Petri run left in flight by the previous server cannot resume; reporting it failed" + ); + let (_, _, record) = failed(FailureReason::WorkflowError, reason); + run_records::lifecycle(state, run_id, record).await?; + return Ok(()); + } + }; + info!( + run_id = %run_id, + mode = super::worker_mode_arg(mode), + "Petri run left in flight by the previous server; relaunching its worker" + ); + let mut runs = state.runs.lock().expect("runs lock poisoned"); + runs.insert( + run_id, + super::managed_run( + run_state.spec.graph_source.clone().unwrap_or_default(), + RunStatus::Runnable, + run_id.created_at(), + scratch_root(state, run_id), + mode, + ), + ); + Ok(()) +} + +/// Resume a run whose store failed under it, as after a crash: its worker +/// (or the in-process engine) ended the run's lifetime with nothing +/// recorded after the failure and no end for the run. The run goes back +/// to the scheduler through the same [`relaunch`] a restart takes. +/// +/// `Err` with the failure to record when the run does not resume: it was +/// deleted or ended meanwhile, a cancel is pending, the server is shutting +/// down, its store interrupted it more than [`MAX_STORE_INTERRUPTIONS`] +/// times, or it cannot continue from its records. +pub(crate) async fn resume_after_interruption( + state: &Arc, + run_id: RunId, + message: &str, +) -> Result<(), String> { + let interrupted = format!("The run's store failed: {message}"); + let interruptions = { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + let Some(managed_run) = runs.get_mut(&run_id) else { + return Err(interrupted); + }; + managed_run.store_interruptions += 1; + managed_run.store_interruptions + }; + if interruptions > MAX_STORE_INTERRUPTIONS { + warn!( + run_id = %run_id, + interruptions, + error = message, + "the run's store keeps failing; reporting the run failed" + ); + return Err(format!( + "The run's store failed {interruptions} times; the last failure: {message}" + )); + } + if state.is_shutting_down() { + return Err(interrupted); + } + let run_state = match run_records::projection(state, run_id).await { + Ok(Some(run_state)) => run_state, + Ok(None) => return Err(interrupted), + Err(err) => { + return Err(format!("{interrupted}; its state could not be read: {err}")); + } + }; + if run_state.status.is_terminal() || run_state.pending_control == Some(RunControlAction::Cancel) + { + return Err(interrupted); + } + let mode = match relaunch(state, run_id, &run_state).await { + Ok(Relaunch::Worker(mode)) => mode, + Ok(Relaunch::Failed { reason }) => return Err(reason), + Err(err) => { + return Err(format!( + "{interrupted}; the run could not be resumed: {err:#}" + )); + } + }; + warn!( + run_id = %run_id, + interruptions, + error = message, + mode = super::worker_mode_arg(mode), + "the run's store interrupted it; resuming it as after a crash" + ); + { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + let Some(managed_run) = runs.get_mut(&run_id) else { + return Err(interrupted); + }; + managed_run.status = RunStatus::Runnable; + managed_run.execution_mode = mode; + clear_live_run_state(managed_run); + } + state.scheduler_notify.notify_one(); + Ok(()) +} + +/// How a run whose lifetime ended short of its end continues. +enum Relaunch { + /// A new worker takes it, in this mode. + Worker(RunExecutionMode), + /// Its records say it cannot continue. + Failed { reason: String }, +} + +/// Ready a run whose lifetime ended short of its end for a new worker, as a +/// crash is recovered. +/// +/// The lease the previous worker held is released from outside, which +/// fences that worker should it still be alive. Then the recovery protocol +/// reads the run's durable execution state: a run with a failed checkpoint +/// cannot continue; otherwise every live workspace on this host is verified +/// against, reset to, or restored from the snapshot its last durable finish +/// names, and a finish with no snapshot cannot continue rather than resume +/// on stale files. A run that continues is asked to start again as a resume +/// (`run.start_requested` with `resume`, then `run.runnable`, the same pair +/// the API's resume appends), in resume mode when Petri's store holds the +/// run, else in start mode: a worker that died before it created the run's +/// record left nothing to continue from, so the run starts from its +/// admitted graphs. The caller registers the run with the scheduler. +async fn relaunch( + state: &Arc, + run_id: RunId, + run_state: &fabro_store::RunProjection, +) -> anyhow::Result { let held = match state.petri_runs.release_for_restart(run_id).await { Ok(()) => true, Err(StoreError::NotFound { .. }) => false, @@ -583,10 +724,6 @@ pub(crate) async fn reconcile_on_startup( return Err(anyhow::Error::new(err).context("releasing the Petri run's lease")); } }; - let run_dir = Storage::new(state.server_storage_dir()) - .run_scratch(&run_id) - .root() - .to_path_buf(); let mode = if held { let request = RecoveryRequest::for_run( run_id, @@ -608,27 +745,11 @@ pub(crate) async fn reconcile_on_startup( ); RunExecutionMode::Resume } - Recovery::Failed { reason } => { - warn!( - run_id = %run_id, - petri_key = %key, - error = %reason, - "Petri run left in flight by the previous server cannot resume; reporting it failed" - ); - let (_, _, record) = failed(FailureReason::WorkflowError, reason); - run_records::lifecycle(state, run_id, record).await?; - return Ok(()); - } + Recovery::Failed { reason } => return Ok(Relaunch::Failed { reason }), } } else { RunExecutionMode::Start }; - info!( - run_id = %run_id, - petri_key = %key, - mode = super::worker_mode_arg(mode), - "Petri run left in flight by the previous server; relaunching its worker" - ); let mut start_requested = RunLifecycleRecord::new(RunLifecycleKind::StartRequested); start_requested.source = Some("resume".to_string()); let mut runnable = run_records::transition(RunLifecycleKind::Runnable, RunStatus::Runnable); @@ -636,18 +757,15 @@ pub(crate) async fn reconcile_on_startup( for record in [start_requested, runnable] { run_records::lifecycle(state, run_id, record).await?; } - let mut runs = state.runs.lock().expect("runs lock poisoned"); - runs.insert( - run_id, - super::managed_run( - run_state.spec.graph_source.clone().unwrap_or_default(), - RunStatus::Runnable, - run_id.created_at(), - run_dir, - mode, - ), - ); - Ok(()) + Ok(Relaunch::Worker(mode)) +} + +/// The run's scratch root: its worker's run directory. +fn scratch_root(state: &AppState, run_id: RunId) -> std::path::PathBuf { + Storage::new(state.server_storage_dir()) + .run_scratch(&run_id) + .root() + .to_path_buf() } /// The failed status, its message, and the `failed` lifecycle record for it. diff --git a/lib/apps/fabro-server/src/worker_runtime.rs b/lib/apps/fabro-server/src/worker_runtime.rs index 7b4fa39b6..9d2029340 100644 --- a/lib/apps/fabro-server/src/worker_runtime.rs +++ b/lib/apps/fabro-server/src/worker_runtime.rs @@ -7,6 +7,7 @@ use async_trait::async_trait; use fabro_static::EnvVars; use fabro_types::RunId; use fabro_types::settings::server::LogDestination; +use fabro_util::exit::ExitClass; use futures_util::future::BoxFuture; use tokio::io::AsyncRead; use tokio::process::Command; @@ -62,9 +63,19 @@ pub(crate) struct StartedWorker { #[derive(Debug)] pub(crate) struct WorkerExit { pub(crate) success: bool, + /// The worker's exit code; `None` when a signal ended it. + pub(crate) code: Option, pub(crate) detail: String, } +impl WorkerExit { + /// Whether the worker ended the run's lifetime and not the run: its + /// store failed, and the run resumes ([`ExitClass::Interrupted`]). + pub(crate) fn interrupted(&self) -> bool { + self.code == Some(ExitClass::Interrupted.code()) + } +} + #[derive(Default)] pub(crate) struct LocalWorkerRuntime; @@ -137,6 +148,7 @@ impl WorkerRuntime for LocalWorkerRuntime { let status = child.wait().await.context("worker wait failed")?; Ok(WorkerExit { success: status.success(), + code: status.code(), detail: status.to_string(), }) }); diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 5a01eace9..32d707583 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -18,6 +18,13 @@ //! `inspect_run` over a read handle of the same store, so what the caller //! reports is what the durable record says. //! +//! A failed write to the run's store ends the run's lifetime, not the run: +//! Petri records nothing after it and returns `CoordinatorError::StoreFailed`, +//! and no firing fails for it. [`run`] returns [`RunError::StoreFailed`] +//! without reading the record back, and [`conclusion`] says +//! [`Conclusion::Interrupted`]: the caller records no end for the run, and +//! the run resumes from its records, as after a crash. +//! //! What the caller supplies beyond the runtime: the interviewer its //! questions go to ([`interview`](crate::interview) in the worker and the //! server), the secret provider over the vault ([`secrets`](crate::secrets)) @@ -59,8 +66,8 @@ use fabro_util::sync; use petri_execution::host::{self, HostError, HostRun}; use petri_execution::inspect::{self, InspectError, RunInspection}; use petri_execution::{ - Access, CancelReason, ExecutionObserver, InterviewDispatcher, Interviewer, RECEIPT_FILE, - RunKey, RunStore, + Access, CancelReason, CoordinatorError, ExecutionObserver, InterviewDispatcher, Interviewer, + RECEIPT_FILE, RunKey, RunStore, }; use petri_runtime::driver::ExecutionReport; use petri_runtime::driver::lifecycle::ExecutionHooks; @@ -169,6 +176,11 @@ pub enum RunError { Inspect(#[source] InspectError), #[error("the run ended without recording a status; the record says: {}", .0.join("; "))] Unfinished(Vec), + /// A write to the run's store failed. The lifetime ended there, with + /// nothing recorded after the failure; the run did not end, and resumes + /// from its records. + #[error("the run's store failed: {0}")] + StoreFailed(String), } /// How Fabro reports the run: what its read side records as the run's @@ -183,6 +195,10 @@ pub enum Conclusion { reason: FailureReason, message: String, }, + /// The run's store failed, which ended this lifetime of the run and not + /// the run: the caller records no terminal event, and the run resumes + /// from its records, as after a crash. + Interrupted { message: String }, } fn network_policy(settings: &EnvironmentNetworkSettings) -> NetworkPolicy { @@ -331,6 +347,11 @@ pub async fn run(request: RunRequest) -> Result { Ok(report) => debug!(status = %report.status, "Petri run ended"), Err(error) => warn!(error = %error, "Petri run ended with a host error"), } + // The store holds what it held at the failure, and the run is not over: + // there is no outcome to read back. + if let Err(HostError::Coordinator(CoordinatorError::StoreFailed(message))) = result { + return Err(RunError::StoreFailed(message)); + } let inspection = inspect(request.store.as_ref(), &key).await?; let mut outcome = outcome(inspection, result.err())?; // A failed checkpoint cancelled the run; what Fabro reports is the @@ -397,9 +418,9 @@ pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result) -> Conclusion { match result { @@ -419,6 +440,9 @@ pub fn conclusion(result: &Result) -> Conclusion { message: failure_message(outcome), } } + Err(error @ RunError::StoreFailed(_)) => Conclusion::Interrupted { + message: error.to_string(), + }, Err(error) => Conclusion::Failed { reason: FailureReason::WorkflowError, message: error_chain(error), @@ -641,4 +665,11 @@ mod tests { .to_string(), }); } + #[test] + fn a_failed_store_concludes_interrupted() { + let error = RunError::StoreFailed("could not append: the disk is full".to_string()); + assert_eq!(conclusion(&Err(error)), Conclusion::Interrupted { + message: "the run's store failed: could not append: the disk is full".to_string(), + }); + } } diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index 2fd05be82..e0c349cc3 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -134,7 +134,7 @@ pub async fn plan( }; // A record with no root invocation (the worker died between creating // the run and declaring it) has nothing to reconcile; the worker's - // resume reports it as such. + // resume starts it again from its admitted graphs. let coordinator = petri_execution::read_coordinator_log(&*logs) .await .map_err(RecoveryError::Log)?; diff --git a/lib/components/fabro-petri/tests/resume.rs b/lib/components/fabro-petri/tests/resume.rs index ecfbc2520..9f07cb148 100644 --- a/lib/components/fabro-petri/tests/resume.rs +++ b/lib/components/fabro-petri/tests/resume.rs @@ -5,12 +5,16 @@ mod support; use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use fabro_petri::admission::AdmittedGraphs; use fabro_petri::check::Launch; -use fabro_petri::engine::{self, Execution, RunStatus}; +use fabro_petri::engine::{self, Conclusion, Execution, RunError, RunStatus}; use fabro_petri::runtime::RuntimeSpec; -use petri_store::{Access, MemoryRunStore, OwnerId, RunKey, RunStore as _}; -use support::{SETTINGS, Silent, admit, no_questions, run_request}; +use petri_store::{ + Access, Digest, LogId, MemoryRunStore, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError, +}; +use support::{SETTINGS, Silent, admit, all_records, no_questions, run_request}; /// One command stage between start and exit. const COMMAND: &str = r#"digraph Command { @@ -41,21 +45,160 @@ async fn a_resume_of_a_run_that_never_started_starts_it_again() { Launch::default(), &runtime, ); - let mut request = run_request( + let outcome = engine::run(resume_request( "cut-short", root.path(), graphs, store, runtime, + )) + .await + .expect("the run ends"); + + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); +} + +/// A write to the run's store fails partway through the run: the lifetime +/// ends with the store's failure, which interrupts the run rather than +/// failing it, and nothing records an end. A resume over the same records +/// finishes the run. +#[tokio::test] +async fn a_failed_store_write_interrupts_the_run_and_a_resume_finishes_it() { + let root = tempfile::tempdir().expect("a temp dir"); + let memory = Arc::new(MemoryRunStore::new()); + let failing = Arc::new(FailingStore { + inner: Arc::clone(&memory), + log: LogId::Resources, + appended: Arc::new(AtomicUsize::new(0)), + }); + let runtime = RuntimeSpec::default(); + let graphs = || { + admit( + &[("workflow.fabro", COMMAND), ("workflow.toml", SETTINGS)], + Launch::default(), + &runtime, + ) + }; + + let first = engine::run(run_request( + "interrupted", + root.path(), + graphs(), + failing, + runtime.clone(), + no_questions(Arc::new(Silent)), + )) + .await; + + assert!( + matches!(&first, Err(RunError::StoreFailed(message)) if message.contains("the disk is full")), + "{first:?}" + ); + assert!( + matches!(engine::conclusion(&first), Conclusion::Interrupted { .. }), + "{first:?}" + ); + assert!(!finished(&memory).await, "nothing ended the run"); + + let outcome = engine::run(resume_request( + "interrupted", + root.path(), + graphs(), + Arc::clone(&memory) as Arc, + runtime, + )) + .await + .expect("the resumed run ends"); + + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); + assert!(finished(&memory).await, "the resume ended the run"); +} + +/// Whether the interrupted run's records hold Petri's own finish. +async fn finished(store: &MemoryRunStore) -> bool { + all_records(store, "interrupted") + .await + .iter() + .any(|record| record["body"]["event"] == "run.finished") +} + +/// A resume of the run over `store`, with the admitted graphs to start it +/// from should its creation have been cut short. +fn resume_request( + run_id: &str, + run_dir: &std::path::Path, + graphs: AdmittedGraphs, + store: Arc, + runtime: RuntimeSpec, +) -> engine::RunRequest { + let mut request = run_request( + run_id, + run_dir, + graphs, + store, + runtime, no_questions(Arc::new(Silent)), ); request.execution = match request.execution { Execution::Start(graphs) => Execution::Resume(graphs), resume @ Execution::Resume(_) => resume, }; - - let outcome = engine::run(request).await.expect("the run ends"); - - assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); - assert!(outcome.complete, "{:?}", outcome.incomplete); + request +} + +/// The in-memory store, failing its first append to `log` as a full disk +/// would, before anything is stored. +struct FailingStore { + inner: Arc, + log: LogId, + appended: Arc, +} + +#[async_trait::async_trait] +impl RunStore for FailingStore { + async fn open(&self, key: &RunKey, access: Access) -> Result, StoreError> { + Ok(Arc::new(FailingLogs { + inner: self.inner.open(key, access).await?, + log: self.log, + appended: Arc::clone(&self.appended), + })) + } +} + +struct FailingLogs { + inner: Arc, + log: LogId, + appended: Arc, +} + +#[async_trait::async_trait] +impl RunLogs for FailingLogs { + fn locator(&self) -> String { + self.inner.locator() + } + + async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> { + if *log == self.log && self.appended.fetch_add(1, Ordering::SeqCst) == 0 { + return Err(StoreError::backend( + self.locator(), + "append", + "the disk is full", + )); + } + self.inner.append(log, records).await + } + + async fn read(&self, log: &LogId) -> Result, StoreError> { + self.inner.read(log).await + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + self.inner.put_blob(bytes).await + } + + async fn get_blob(&self, digest: Digest) -> Result>, StoreError> { + self.inner.get_blob(digest).await + } } diff --git a/lib/foundation/fabro-util/src/exit.rs b/lib/foundation/fabro-util/src/exit.rs index 31b09d4c8..57fd1c7f5 100644 --- a/lib/foundation/fabro-util/src/exit.rs +++ b/lib/foundation/fabro-util/src/exit.rs @@ -3,6 +3,21 @@ use anyhow::Error; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum ExitClass { AuthRequired, + /// The command stopped short of its end for a reason a later attempt + /// can get past: a run worker whose run's store failed, which leaves + /// the run for the server to resume. `EX_TEMPFAIL` from `sysexits.h`. + Interrupted, +} + +impl ExitClass { + /// The process exit code of an error of this class. + #[must_use] + pub const fn code(self) -> i32 { + match self { + Self::AuthRequired => 4, + Self::Interrupted => 75, + } + } } // Keep the wrapper transparent so existing stderr remains unchanged while the @@ -49,9 +64,7 @@ impl ErrorExt for Error { pub fn exit_code_for(err: &Error) -> i32 { err.chain() .find_map(|cause| cause.downcast_ref::()) - .map_or(1, |classified| match classified.class() { - ExitClass::AuthRequired => 4, - }) + .map_or(1, |classified| classified.class().code()) } pub fn exit_class_for(err: &Error) -> Option { @@ -77,6 +90,13 @@ mod tests { assert_eq!(exit_code_for(&err), 4); } + #[test] + fn interrupted_errors_map_to_exit_75() { + let err = anyhow!("boom").classify(ExitClass::Interrupted); + assert_eq!(exit_code_for(&err), 75); + assert_eq!(exit_class_for(&err), Some(ExitClass::Interrupted)); + } + #[test] fn classification_keeps_display_transparent() { assert_eq!(