From 3c395f9e6ea3047f649f0ace561f3fb333562c41 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:05:40 -0400 Subject: [PATCH 1/6] Start a Petri run again when its creation was cut short Petri now refuses to resume a run whose creation a crash cut short (the key is stored, the root invocation is not) with HostError::NotStarted, and starts it again when the host runs it under the same key. The engine used its own guard, check_resumable, which failed the run with NothingToResume instead. Execution::Resume now carries the admitted graphs, and a resume Petri answers with NotStarted starts the run from them. The worker loads the graphs in resume mode too, as does the server's in-process path. The guard and RunError::NothingToResume are gone. New test: a_resume_of_a_run_that_never_started_starts_it_again. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 33 +++-- .../fabro-server/src/server/petri_runs.rs | 23 ++-- lib/components/fabro-petri/src/engine.rs | 114 ++++++++++-------- lib/components/fabro-petri/tests/resume.rs | 61 ++++++++++ 4 files changed, 154 insertions(+), 77 deletions(-) create mode 100644 lib/components/fabro-petri/tests/resume.rs 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 f59b415ed..0b1e8ed13 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -9,9 +9,10 @@ //! //! The run's record is [`HttpRunStore`] over the worker's client, leased //! for this launch: the worker mints one owner id at start, logs it, and -//! every lease the run takes over the API names it. `--mode start` loads -//! the admitted graphs through the client's blob read and runs them; -//! `--mode resume` continues the run from its records. Either way the +//! every lease the run takes over the API names it. Both modes load the +//! admitted graphs through the client's blob read: `--mode start` runs them; +//! `--mode resume` continues the run from its records, and starts it again +//! 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. @@ -180,21 +181,19 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { run_tools, ) .await?; + let client = worker.client.clone_for_reuse(); + let graphs = admission::load_with( + |blob| { + let client = client.clone_for_reuse(); + async move { client.read_run_blob(&run_id, &blob).await } + }, + &admission, + ) + .await + .context("loading the admitted graphs")?; let execution = match worker.mode { - RunWorkerMode::Start => { - let client = worker.client.clone_for_reuse(); - let graphs = admission::load_with( - |blob| { - let client = client.clone_for_reuse(); - async move { client.read_run_blob(&run_id, &blob).await } - }, - &admission, - ) - .await - .context("loading the admitted graphs")?; - Execution::Start(graphs) - } - RunWorkerMode::Resume => Execution::Resume, + RunWorkerMode::Start => Execution::Start(graphs), + RunWorkerMode::Resume => Execution::Resume(graphs), }; let started = Instant::now(); diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 2c7c2d815..0d059cace 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -391,18 +391,19 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { { return; } - let execution = match mode { - RunExecutionMode::Start => { - match admission::load(&state.store_ref().blobs(), &admission).await { - Ok(graphs) => Execution::Start(graphs), - Err(err) => { - let message = error_util::collect_chain(&err).join(": "); - fail_before_execution(&state, run_id, &message).await; - return; - } - } + // A resume loads the graphs too: they start the run again when a crash + // cut its creation short. + let graphs = match admission::load(&state.store_ref().blobs(), &admission).await { + Ok(graphs) => graphs, + Err(err) => { + let message = error_util::collect_chain(&err).join(": "); + fail_before_execution(&state, run_id, &message).await; + return; } - RunExecutionMode::Resume => Execution::Resume, + }; + let execution = match mode { + RunExecutionMode::Start => Execution::Start(graphs), + RunExecutionMode::Resume => Execution::Resume(graphs), }; // The run's secrets: a snapshot of the server's vault, as a worker // takes one at launch. diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index b833bc12d..5a01eace9 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -11,7 +11,10 @@ //! to `execution::host`: [`Execution::Start`] runs the admitted graphs //! through `run_configured`; [`Execution::Resume`] continues the run from //! its records through `resume_configured`, with the same observers a start -//! installs, as the host's docs require. The outcome is then derived from +//! installs, as the host's docs require. A run whose creation a crash cut +//! short (Petri's `HostError::NotStarted`: the key is stored, the root +//! invocation is not) starts again from its admitted graphs under the same +//! key, which Petri takes over. The outcome is then derived from //! `inspect_run` over a read handle of the same store, so what the caller //! reports is what the durable record says. //! @@ -45,25 +48,28 @@ //! projection over Petri's records is the read-side item that follows. use std::path::PathBuf; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use fabro_types::settings::run::{ EnvironmentNetworkMode, EnvironmentNetworkSettings, EnvironmentResourcesSettings, }; use fabro_types::settings::size::Size; use fabro_types::{FailureReason, RunId, SandboxProviderKind}; +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, InvocationId, - RECEIPT_FILE, RunKey, RunStore, + Access, CancelReason, ExecutionObserver, InterviewDispatcher, Interviewer, RECEIPT_FILE, + RunKey, RunStore, }; +use petri_runtime::driver::ExecutionReport; use petri_runtime::driver::lifecycle::ExecutionHooks; pub use petri_runtime::executor::Retention; use petri_runtime::executor::SecretProvider; -use petri_runtime::{DaytonaResources, LostSandbox, RunOptions, SandboxBackend}; +use petri_runtime::{DaytonaResources, LostSandbox, RunOptions, Runtime, SandboxBackend}; use sandbox_driver::NetworkPolicy; use tokio::fs; +use tokio::task::JoinHandle; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; @@ -80,9 +86,9 @@ pub enum Execution { /// Run the admitted graphs from the start; the run must not exist in /// the store yet. Start(AdmittedGraphs), - /// Continue the run from its records; the run must exist in the store - /// with its root invocation declared. - Resume, + /// Continue the run from its records. The admitted graphs start the + /// run again when a crash cut its creation short. + Resume(AdmittedGraphs), } /// One run to execute. @@ -159,8 +165,6 @@ pub enum RunError { Open(#[source] petri_store::StoreError), #[error("the run's record could not be read")] Read(#[source] HostError), - #[error("the run's record has no root invocation, so there is nothing to resume")] - NothingToResume, #[error("the run's record could not be inspected")] Inspect(#[source] InspectError), #[error("the run ended without recording a status; the record says: {}", .0.join("; "))] @@ -217,7 +221,7 @@ pub async fn run(request: RunRequest) -> Result { } // A normal resume requires its original sandbox to survive. options.sandbox.lost_sandbox = LostSandbox::Refuse; - let resumed = matches!(request.execution, Execution::Resume); + let resumed = matches!(request.execution, Execution::Resume(_)); let mut runtime = request .runtime .runtime(true) @@ -266,19 +270,25 @@ pub async fn run(request: RunRequest) -> Result { .capability(controls.turns()); let dispatcher = InterviewDispatcher::new(request.interviewer); - let cancel = request.cancel.clone(); - let mut cancel_task = None; - let with_handle = |handle: petri_execution::CoordinatorHandle, secrets| { - dispatcher.wire(handle.clone(), secrets); - if let Some(hooks) = &fabro_hooks { - hooks.attach(handle.clone()); + let cancel_task: Mutex>> = Mutex::new(None); + // What the coordinator's handle is wired to, on a start and a resume + // alike: built again when a resume starts the run over. + let wiring = || { + let cancel = request.cancel.clone(); + let (dispatcher, fabro_hooks, controls, cancel_task) = + (&dispatcher, &fabro_hooks, &controls, &cancel_task); + move |handle: petri_execution::CoordinatorHandle, secrets| { + dispatcher.wire(handle.clone(), secrets); + if let Some(hooks) = fabro_hooks { + hooks.attach(handle.clone()); + } + controls.wire(handle.clone()); + *sync::lock(cancel_task) = Some(tokio::spawn(async move { + cancel.cancelled().await; + info!("cancelling the Petri run"); + handle.cancel_root_for(CancelReason::Control); + })); } - controls.wire(handle.clone()); - cancel_task = Some(tokio::spawn(async move { - cancel.cancelled().await; - info!("cancelling the Petri run"); - handle.cancel_root_for(CancelReason::Control); - })); }; let mut observers = request.observers; observers.push(controls.observer()); @@ -286,25 +296,33 @@ pub async fn run(request: RunRequest) -> Result { let result = match request.execution { Execution::Start(graphs) => { info!(run_id = %request.run_id, backend = %backend, "Starting Petri run"); - let mut host_run = HostRun::new(graphs.graph).with_children(graphs.children); - for observer in observers { - host_run = host_run.observe(observer); - } - Box::pin(host::run_configured(&runtime, host_run, with_handle)).await + start(&runtime, graphs, observers, wiring()).await } - Execution::Resume => { - check_resumable(request.store.as_ref(), &key).await?; + Execution::Resume(graphs) => { info!(run_id = %request.run_id, backend = %backend, "Resuming Petri run"); - Box::pin(host::resume_configured( + let resumed = Box::pin(host::resume_configured( &runtime, Vec::new(), - observers, - with_handle, + observers.clone(), + wiring(), )) - .await + .await; + match resumed { + // A crash cut the run's creation short: nothing beyond its + // start is stored, so it starts again from its admitted + // graphs, and Petri takes the stored prefix over. + Err(HostError::NotStarted) => { + info!( + run_id = %request.run_id, + "The Petri run never started; starting it again" + ); + start(&runtime, graphs, observers, wiring()).await + } + resumed => resumed, + } } }; - if let Some(task) = cancel_task { + if let Some(task) = sync::lock(&cancel_task).take() { task.abort(); } let receipt = dispatcher.shutdown().await; @@ -469,21 +487,19 @@ pub(crate) fn backend(provider: &SandboxProviderKind) -> Option } } -/// Refuse a resume the host would not survive: `resume_configured` indexes -/// the root invocation of the stored state, so a record with none (the run -/// was created in the store and nothing more) is refused here with a named -/// error instead. -async fn check_resumable(store: &dyn RunStore, key: &RunKey) -> Result<(), RunError> { - let logs = store - .open(key, Access::Read) - .await - .map_err(RunError::Open)?; - let state = host::stored_state(&*logs).await.map_err(RunError::Read)?; - if state.invocations.contains_key(&InvocationId::ROOT) { - Ok(()) - } else { - Err(RunError::NothingToResume) +/// Run the admitted graphs under the run's key: a fresh run, or one whose +/// creation a crash cut short, which Petri takes over. +async fn start( + runtime: &Runtime, + graphs: AdmittedGraphs, + observers: Vec>, + with_handle: impl FnOnce(petri_execution::CoordinatorHandle, Arc), +) -> Result { + let mut host_run = HostRun::new(graphs.graph).with_children(graphs.children); + for observer in observers { + host_run = host_run.observe(observer); } + Box::pin(host::run_configured(runtime, host_run, with_handle)).await } /// Read the run back through a handle that holds no lease. diff --git a/lib/components/fabro-petri/tests/resume.rs b/lib/components/fabro-petri/tests/resume.rs new file mode 100644 index 000000000..ecfbc2520 --- /dev/null +++ b/lib/components/fabro-petri/tests/resume.rs @@ -0,0 +1,61 @@ +//! A resume through the engine assembly: a run whose creation a crash cut +//! short starts again from its admitted graphs, and a run whose store +//! failed ends its lifetime without an end of its own. + +mod support; + +use std::sync::Arc; + +use fabro_petri::check::Launch; +use fabro_petri::engine::{self, Execution, RunStatus}; +use fabro_petri::runtime::RuntimeSpec; +use petri_store::{Access, MemoryRunStore, OwnerId, RunKey, RunStore as _}; +use support::{SETTINGS, Silent, admit, no_questions, run_request}; + +/// One command stage between start and exit. +const COMMAND: &str = r#"digraph Command { + start [shape=Mdiamond] + exit [shape=Msquare] + say [shape=parallelogram, script="true"] + start -> say -> exit +}"#; + +/// A crash cut the run's creation short: its key is stored, and nothing +/// else. The resume Petri refuses as never started becomes a start from +/// the admitted graphs, under the same key. +#[tokio::test] +async fn a_resume_of_a_run_that_never_started_starts_it_again() { + let root = tempfile::tempdir().expect("a temp dir"); + let store = Arc::new(MemoryRunStore::new()); + drop( + store + .open(&RunKey::new("cut-short"), Access::Create { + owner: OwnerId::new("crashed"), + }) + .await + .expect("the key is stored"), + ); + let runtime = RuntimeSpec::default(); + let graphs = admit( + &[("workflow.fabro", COMMAND), ("workflow.toml", SETTINGS)], + Launch::default(), + &runtime, + ); + let mut request = run_request( + "cut-short", + root.path(), + 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); +} From 37116758fa8a3d51295095730d06985bcf3363dd Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:18:58 -0400 Subject: [PATCH 2/6] 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!( From cb54a0df90ed39b636f56b2a9b124ffea332dd14 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:33:26 -0400 Subject: [PATCH 3/6] Stop a surviving worker before ending its lease from outside A worker leads a process group of its own, so it outlives a server crash. The restarted server released the run's lease from outside and launched a resume while that worker could still be running: a worker whose lease is released keeps acting until its next write, beside its successor. A delete likewise dropped the lease without knowing the worker was gone. Now each worker holds a lock on `worker.lock` in its run's scratch directory for its whole life, taken before anything else. It is a POSIX record lock: the kernel frees it only when the worker exits, and names the process that holds it. Before the server ends a lease from outside (the relaunch after a restart or a store interruption, and a delete), it kills whatever process still holds the lock, with its process group, and waits until the lock is free. A worker that finds the lock held does not start. - fabro-proc: ProcessLock::try_hold and stop_lock_holder, tested with this test binary as the holding process. - fabro-config: RunScratch::worker_lock_path. - New scenario: a worker that outlives the server is gone before the resume's worker launches, and the run succeeds once. It fails without the server-side stop. A host stage process runs in a process group of its own and still outlives its killed worker, as it did at a worker crash. Co-Authored-By: Claude Opus 5.5 --- lib/apps/fabro-cli/src/commands/run/runner.rs | 31 ++ lib/apps/fabro-cli/tests/it/scenario/petri.rs | 71 +++++ lib/apps/fabro-server/src/petri_runs.rs | 15 +- lib/apps/fabro-server/src/server.rs | 8 +- .../fabro-server/src/server/petri_runs.rs | 46 ++- lib/foundation/fabro-config/src/storage.rs | 7 + lib/foundation/fabro-proc/src/lib.rs | 4 + lib/foundation/fabro-proc/src/process_lock.rs | 264 ++++++++++++++++++ 8 files changed, 436 insertions(+), 10 deletions(-) create mode 100644 lib/foundation/fabro-proc/src/process_lock.rs diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index e8120b462..3aceaa24a 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -5,6 +5,8 @@ use std::time::Duration; use anyhow::{Context, Result, anyhow}; use fabro_client::ServerTarget; +#[cfg(unix)] +use fabro_config::RunScratch; use fabro_config::Storage; use fabro_interview::{ AnswerSubmission, ControlInterviewer, WORKER_CONTROL_INVALID_CURSOR_REASON, @@ -22,6 +24,8 @@ use futures::{SinkExt, StreamExt}; use jsonwebtoken::dangerous::insecure_decode; #[cfg(unix)] use nix::unistd; +#[cfg(unix)] +use tokio::fs; #[cfg(test)] use tokio::io::DuplexStream; use tokio::net::TcpStream; @@ -65,6 +69,10 @@ pub(crate) async fn execute( ) -> Result<()> { let _ = fabro_proc::title_init(); set_worker_title(&run_id, initial_worker_title_phase(mode)); + // Held until this process exits: while it is, the server ends no lease + // of this run from outside. + #[cfg(unix)] + let _running = hold_worker_lock(&run_dir, &run_id).await?; let target = server.parse::()?; let client = server_client::connect_server_target_with_bearer(&target, worker_token).await?; @@ -93,6 +101,29 @@ pub(crate) async fn execute( .await } +/// Take the run's worker lock for this process's whole life, before +/// anything else. The server ends a worker's lease from outside only once +/// the lock is free, which the kernel makes it only when the process that +/// held it is gone: a worker that outlived a server crash is stopped +/// before its run resumes. Another process holding the lock is a worker of +/// the same run still running, and this one does not start beside it. +#[cfg(unix)] +async fn hold_worker_lock(run_dir: &Path, run_id: &RunId) -> Result { + fs::create_dir_all(run_dir) + .await + .with_context(|| format!("creating the run directory {}", run_dir.display()))?; + let path = RunScratch::new(run_dir).worker_lock_path(); + fabro_proc::ProcessLock::try_hold(&path) + .await + .with_context(|| format!("taking the worker lock {}", path.display()))? + .ok_or_else(|| { + anyhow!( + "another worker of run {run_id} is still running: it holds {}", + path.display() + ) + }) +} + const WORKER_TOKEN_SCOPE: &str = "run:worker"; const WORKER_RUN_TOOLS_SCOPE: &str = "agent:run_tools"; diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index 7beb2476c..8a3482286 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -779,6 +779,77 @@ async fn a_petri_run_resumes_in_a_new_worker_after_the_server_restarts() { server.shutdown(); } +/// A worker outlives a server crash: it leads a process group of its own. +/// The restarted server stops it before it ends the worker's lease and +/// launches the resume, so no two workers of the run ever run side by side. +#[tokio::test(flavor = "multi_thread")] +async fn a_worker_that_outlives_the_server_is_stopped_before_its_run_resumes() { + let context = test_context!(); + let mut server = RunningServer::start().await; + let gate = context.temp_dir.join("survivor.gate"); + let script = format!("while [ ! -f {} ]; do sleep 0.05; done", gate.display()); + let workspace = write_petri_workspace(&context, &script); + let run_id = run_detached(&context, &server, &workspace); + + wait_for_status(&server, &run_id, &["running"]).await; + let worker = wait_for_worker(&run_id); + eprintln!("worker {worker} launched"); + wait_until_gate_is_polled(&gate); + + // The server alone dies: its worker keeps running. + server.kill(); + assert!( + fabro_proc::process_running_strict(worker), + "the worker outlived the server" + ); + + server.launch().await; + eprintln!("server restarted"); + let resumed = wait_for_worker_other_than(&run_id, worker); + eprintln!("worker {resumed} launched for the resume"); + assert!( + !fabro_proc::process_running_strict(worker), + "the surviving worker {worker} was stopped before the resume's worker launched" + ); + std::fs::write(&gate, "go").expect("the gate opens"); + + let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await; + let names = stream_names(&settled_stream(&server, &run_id).await); + assert_eq!( + status, + "succeeded", + "stream: {names:?}\nserver stderr:\n{}", + server.stderr_text() + ); + assert_eq!(count_of(&names, "lifecycle:succeeded"), 1, "{names:?}"); + assert_eq!(count_of(&names, "run.finished"), 1, "{names:?}"); + server.shutdown(); +} + +/// Wait for a worker of the run other than `previous`. +fn wait_for_worker_other_than(run_id: &str, previous: u32) -> u32 { + let short_id: String = run_id.chars().take(12).collect(); + let deadline = Instant::now() + RUN_TIMEOUT; + loop { + let output = Command::new("pgrep") + .args(["-f", &format!("^fabro {short_id} ")]) + .output() + .expect("pgrep runs"); + if let Some(pid) = String::from_utf8_lossy(&output.stdout) + .lines() + .filter_map(|line| line.trim().parse::().ok()) + .find(|pid| *pid != previous) + { + return pid; + } + assert!( + Instant::now() < deadline, + "no new worker process appeared for run {run_id}" + ); + std::thread::sleep(POLL); + } +} + /// The run's status while the server may be down: `None` when it is. async fn run_status_offline(server: &RunningServer) -> Option { fabro_test::test_http_client() diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index c008f69d5..b977b1d7f 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -8,7 +8,10 @@ //! should last: until the worker releases it, or until the server observes //! the worker exit. That is the integration plan's rule for a lease: it ends //! when the handle drops, when the server observes the worker exit, or by -//! operator release, never by timeout. +//! operator release, never by timeout. A release from outside comes only +//! once the worker is gone: a worker whose lease was released keeps running +//! until its next write, so the server first stops a worker that still +//! holds the run's worker lock (a worker outlives a server crash). //! //! The Petri run key of a Fabro run is the run id's text, as the plan sets //! `RunOptions::run_key`. @@ -119,10 +122,12 @@ impl PetriRuns { } /// End whatever lease the run's previous worker held, from outside: - /// what the server does for a run it finds in flight at startup, before - /// it launches a new worker for it. The previous worker, should it still - /// be alive, finds its handles stale on its next write. `NotFound` when - /// the store never held the run. + /// what the server does before it launches a new worker for a run whose + /// lifetime ended short of its end. The caller first makes sure the + /// previous worker is gone (`server::petri_runs::stop_previous_worker`): + /// a live worker whose lease is released keeps running until its next + /// write finds its handles stale. `NotFound` when the store never held + /// the run. pub(crate) async fn release_for_restart(&self, run_id: RunId) -> Result<(), StoreError> { self.worker_exited(run_id); self.store.release_lease(&Self::key(&run_id)).await diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 5c3933cab..a21ce3692 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -2788,8 +2788,12 @@ async fn delete_run_internal( // Whatever Petri run handles the run's worker held open over the API // drop here, before its sandboxes are pruned through the lease ledger: - // the worker is gone or was told to stop above, and a lease it still - // held would refuse the prune. + // a lease the worker still held would refuse the prune. The worker was + // told to stop above; a worker that is still running, or one that + // outlived a server crash, is stopped first. + petri_runs::stop_previous_worker(state, id) + .await + .map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, format!("{err:#}")))?; state.petri_runs.worker_exited(id); let delete_outcome = delete_run_sandbox_resource(state, id, force).await?; diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 0be8d2675..2d3201c35 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -38,7 +38,7 @@ use std::collections::HashMap; use std::sync::Arc; -use std::time::Instant; +use std::time::{Duration, Instant}; use fabro_config::{ EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, SettingsLayer, Storage, @@ -700,8 +700,9 @@ enum Relaunch { /// 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 +/// The previous worker is stopped should it still be running +/// ([`stop_previous_worker`]), and only then is the lease it held released +/// from outside. 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 @@ -717,6 +718,7 @@ async fn relaunch( run_id: RunId, run_state: &fabro_store::RunProjection, ) -> anyhow::Result { + stop_previous_worker(state, run_id).await?; let held = match state.petri_runs.release_for_restart(run_id).await { Ok(()) => true, Err(StoreError::NotFound { .. }) => false, @@ -760,6 +762,44 @@ async fn relaunch( Ok(Relaunch::Worker(mode)) } +/// How long the server waits for a worker it killed to be gone. +const WORKER_STOP_PATIENCE: Duration = Duration::from_secs(10); + +/// Make sure no worker of the run is still running, before its lease is +/// ended from outside. A worker whose lease was released keeps running +/// until its next write, beside any successor; and a worker outlives a +/// server crash, since it leads a process group of its own. So a worker +/// that still holds the run's worker lock is killed, with its process +/// group, and the lock is waited on until the kernel frees it at the +/// worker's exit. `Err` when it is still held after +/// [`WORKER_STOP_PATIENCE`]. +pub(crate) async fn stop_previous_worker(state: &AppState, run_id: RunId) -> anyhow::Result<()> { + #[cfg(unix)] + { + let path = Storage::new(state.server_storage_dir()) + .run_scratch(&run_id) + .worker_lock_path(); + let holder = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE) + .await + .map_err(|err| { + anyhow::Error::new(err).context(format!( + "stopping the run's previous worker ({})", + path.display() + )) + })?; + if let fabro_proc::LockHolder::Stopped { pid } = holder { + warn!( + run_id = %run_id, + pid, + "the run's previous worker was still running; stopped it before ending its lease" + ); + } + } + #[cfg(not(unix))] + let _ = (state, run_id); + Ok(()) +} + /// 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()) diff --git a/lib/foundation/fabro-config/src/storage.rs b/lib/foundation/fabro-config/src/storage.rs index e7a249798..8a76c9d73 100644 --- a/lib/foundation/fabro-config/src/storage.rs +++ b/lib/foundation/fabro-config/src/storage.rs @@ -143,6 +143,13 @@ impl RunScratch { self.root.join("runtime") } + /// The lock the run's worker holds for its whole life: while it is + /// held, the worker is running. + #[must_use] + pub fn worker_lock_path(&self) -> PathBuf { + self.root.join("worker.lock") + } + pub fn create(&self) -> std::io::Result<()> { std::fs::create_dir_all(self.worktree_dir())?; std::fs::create_dir_all(self.runtime_dir())?; diff --git a/lib/foundation/fabro-proc/src/lib.rs b/lib/foundation/fabro-proc/src/lib.rs index 7dcc9fe97..ff3103f5f 100644 --- a/lib/foundation/fabro-proc/src/lib.rs +++ b/lib/foundation/fabro-proc/src/lib.rs @@ -11,6 +11,8 @@ mod flock; #[cfg(unix)] mod pre_exec; +#[cfg(unix)] +mod process_lock; mod signal; mod title; @@ -22,6 +24,8 @@ pub use pre_exec::pre_exec_pdeathsig; pub use pre_exec::pre_exec_setpgid; #[cfg(unix)] pub use pre_exec::pre_exec_setsid; +#[cfg(unix)] +pub use process_lock::{LockHolder, ProcessLock, stop_lock_holder}; pub use signal::{process_exists, process_group_alive, process_running, process_running_strict}; #[cfg(unix)] pub use signal::{ diff --git a/lib/foundation/fabro-proc/src/process_lock.rs b/lib/foundation/fabro-proc/src/process_lock.rs new file mode 100644 index 000000000..ceaadc8d0 --- /dev/null +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -0,0 +1,264 @@ +//! A lock one process holds on a file for its whole life. +//! +//! The kernel ends the lock when its process exits, however it exits, so a +//! free lock proves the process that held it is gone, and the kernel names +//! the process that holds a lock to any other process that asks. The lock +//! is a POSIX record lock (`fcntl`), not an `flock`, because only a record +//! lock names its holder. Two rules follow from record-lock semantics: the +//! holder must not open the file a second time, since closing any +//! descriptor of the file ends the process's record locks on it; and a +//! process never sees its own lock as held. + +use std::fs::File; +use std::io; +use std::os::unix::io::AsRawFd; +use std::path::Path; +use std::time::{Duration, Instant}; + +use tokio::fs::OpenOptions; +use tokio::time; + +use crate::signal::{sigkill, sigkill_process_group}; + +/// How often [`stop_lock_holder`] checks whether the lock is free. +const POLL: Duration = Duration::from_millis(20); + +#[allow( + clippy::cast_possible_truncation, + clippy::unnecessary_cast, + reason = "the lock constants are c_int on Linux and c_short on macOS, where flock's fields \ + are c_short on both; every value is small" +)] +mod consts { + pub(super) const WRITE_LOCK: libc::c_short = libc::F_WRLCK as libc::c_short; + pub(super) const UNLOCKED: libc::c_short = libc::F_UNLCK as libc::c_short; + pub(super) const FROM_START: libc::c_short = libc::SEEK_SET as libc::c_short; +} + +/// The lock this process holds on a file, until the lock is dropped or the +/// process exits. +#[derive(Debug)] +pub struct ProcessLock { + _file: File, +} + +impl ProcessLock { + /// Take the lock on the file at `path`, created when missing, for this + /// process, without waiting. `Ok(None)` when another process holds it. + pub async fn try_hold(path: &Path) -> io::Result> { + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(path) + .await? + .into_std() + .await; + let mut lock = whole_file(consts::WRITE_LOCK); + // SAFETY: fcntl(F_SETLK) on a valid descriptor with a valid flock + // struct; F_SETLK does not wait. + if unsafe { libc::fcntl(file.as_raw_fd(), libc::F_SETLK, &raw mut lock) } == 0 { + return Ok(Some(Self { _file: file })); + } + let error = io::Error::last_os_error(); + match error.raw_os_error() { + Some(libc::EACCES | libc::EAGAIN) => Ok(None), + _ => Err(error), + } + } +} + +/// What [`stop_lock_holder`] found. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum LockHolder { + /// No other process held the lock. + None, + /// This process held the lock. It was killed, with its process group, + /// and the lock is free: it is gone. + Stopped { pid: u32 }, +} + +/// Stop the process that holds the lock on the file at `path`, if another +/// process does, and wait up to `patience` for the lock to be free. The +/// holder and its process group get `SIGKILL`: a holder that could handle +/// a signal could also keep running. A missing file has no holder. +/// +/// `Err` when the lock is still held after `patience`. +pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result { + let file = match OpenOptions::new().read(true).write(true).open(path).await { + Ok(file) => file.into_std().await, + Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(LockHolder::None), + Err(error) => return Err(error), + }; + let Some(pid) = holder(&file)? else { + return Ok(LockHolder::None); + }; + sigkill_process_group(pid); + sigkill(pid); + let deadline = Instant::now() + patience; + loop { + match holder(&file)? { + None => return Ok(LockHolder::Stopped { pid }), + Some(_) if Instant::now() >= deadline => { + return Err(io::Error::other(format!( + "process {pid} still holds {} after SIGKILL", + path.display() + ))); + } + Some(_) => time::sleep(POLL).await, + } + } +} + +/// The process that holds the lock on `file`, if another process does. +fn holder(file: &File) -> io::Result> { + let mut lock = whole_file(consts::WRITE_LOCK); + // SAFETY: fcntl(F_GETLK) on a valid descriptor with a valid flock struct + // only reads the lock table. + if unsafe { libc::fcntl(file.as_raw_fd(), libc::F_GETLK, &raw mut lock) } != 0 { + return Err(io::Error::last_os_error()); + } + if lock.l_type == consts::UNLOCKED { + return Ok(None); + } + u32::try_from(lock.l_pid).map(Some).map_err(|_| { + io::Error::other(format!( + "the lock on the file is held, by no process id ({})", + lock.l_pid + )) + }) +} + +/// A lock request covering the whole file, however long it grows. +fn whole_file(kind: libc::c_short) -> libc::flock { + // SAFETY: flock is a plain C struct, for which all zeros is valid. + let mut lock: libc::flock = unsafe { std::mem::zeroed() }; + lock.l_type = kind; + lock.l_whence = consts::FROM_START; + lock.l_start = 0; + lock.l_len = 0; + lock +} + +#[cfg(test)] +#[expect( + clippy::disallowed_methods, + clippy::print_stdout, + reason = "the tests start this test binary as the process that holds a lock, and it says so \ + on its stdout" +)] +mod tests { + use std::process::Stdio; + + use tokio::io::{AsyncBufReadExt as _, BufReader}; + use tokio::process::{Child, Command}; + + use super::*; + + /// Set for the test binary started as a lock's holder: the lock's path. + const HOLD: &str = "FABRO_PROC_TEST_HOLD_LOCK"; + + /// Not a test: the process the other tests start to hold a lock. It + /// holds the lock at `$FABRO_PROC_TEST_HOLD_LOCK`, says so, and waits + /// to be killed. Without the variable it does nothing. + #[tokio::test] + async fn hold_the_lock_until_killed() { + let Some(path) = std::env::var_os(HOLD) else { + return; + }; + let lock = ProcessLock::try_hold(Path::new(&path)) + .await + .expect("the lock file opens") + .expect("the lock is free"); + println!("held"); + time::sleep(Duration::from_mins(1)).await; + drop(lock); + } + + /// This test binary, holding the lock at `path` in a process group of + /// its own, once it says it holds it. + async fn holder_process(path: &Path) -> Child { + let mut command = Command::new(std::env::current_exe().expect("the test binary")); + command + .args([ + "--exact", + "process_lock::tests::hold_the_lock_until_killed", + "--nocapture", + ]) + .env(HOLD, path) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .kill_on_drop(true); + crate::pre_exec_setpgid(command.as_std_mut()); + let mut child = command.spawn().expect("the holder starts"); + let stdout = child.stdout.take().expect("the holder's stdout"); + let mut lines = BufReader::new(stdout).lines(); + while let Some(line) = lines.next_line().await.expect("the holder's stdout reads") { + if line.trim() == "held" { + return child; + } + } + panic!("the holder exited before it held the lock"); + } + + #[tokio::test] + async fn a_missing_or_free_lock_has_no_holder() { + let dir = tempfile::tempdir().expect("a temp dir"); + let path = dir.path().join("worker.lock"); + assert_eq!( + stop_lock_holder(&path, Duration::from_secs(1)) + .await + .expect("the check runs"), + LockHolder::None + ); + drop( + ProcessLock::try_hold(&path) + .await + .expect("the file opens") + .expect("the lock is free"), + ); + assert_eq!( + stop_lock_holder(&path, Duration::from_secs(1)) + .await + .expect("the check runs"), + LockHolder::None + ); + } + + #[tokio::test] + async fn a_lock_another_process_holds_cannot_be_taken() { + let dir = tempfile::tempdir().expect("a temp dir"); + let path = dir.path().join("worker.lock"); + let _holder = holder_process(&path).await; + assert!( + ProcessLock::try_hold(&path) + .await + .expect("the file opens") + .is_none() + ); + } + + #[tokio::test] + async fn the_holder_is_killed_and_the_lock_is_free_once_it_is_gone() { + let dir = tempfile::tempdir().expect("a temp dir"); + let path = dir.path().join("worker.lock"); + let mut holder = holder_process(&path).await; + let pid = holder.id().expect("the holder runs"); + + let found = stop_lock_holder(&path, Duration::from_secs(10)) + .await + .expect("the holder is stopped"); + + assert_eq!(found, LockHolder::Stopped { pid }); + let status = holder.wait().await.expect("the holder is reaped"); + assert!(!status.success(), "the holder was killed: {status}"); + assert!( + ProcessLock::try_hold(&path) + .await + .expect("the file opens") + .is_some(), + "the lock is free" + ); + } +} From eb39e157811d7abae76794524b9fd14b6968273e Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:33:53 -0400 Subject: [PATCH 4/6] Describe store interruptions and the worker lock in the Petri README Co-Authored-By: Claude Opus 5.5 --- lib/components/fabro-petri/README.md | 30 ++++++++++++++++++++++------ 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index ba2f938c4..768ce6007 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -33,10 +33,13 @@ Every adapter the integration plan describes lands here. - `engine`: a run executed by Petri, started from its admitted graphs or resumed from its records, with the outcome read from the run's record through `inspect_run` and mapped to the conclusion Fabro's read side - records. The run's worker process runs it over `HttpRunStore`; the server - runs it in its own process only under its test override, over - `SqliteRunStore`. The caller supplies the interviewer, and the secret - provider and blob table when it has them. + records. A resume of a run whose creation a crash cut short starts it + again from its admitted graphs. A failed store write ends the run's + lifetime, not the run: the conclusion is `Interrupted`, and the run + resumes from its records. The run's worker process runs it over + `HttpRunStore`; the server runs it in its own process only under its + test override, over `SqliteRunStore`. The caller supplies the + interviewer, and the secret provider and blob table when it has them. - `interview`: Petri's `Interviewer` over Fabro's questions API and the worker's control channel. A question has one id in Fabro, Petri's own (`gate#2`): the projection lists it pending from the `question` record, @@ -165,6 +168,13 @@ Every run executes on Petri. The server side is `fabro-server`'s a server restart, a run left in flight goes back to a worker in `--mode resume`: the run continues from its records, as Petri's own resume does, on workspaces the recovery protocol brought to their durable snapshots. +A run whose store failed under it takes the same way back: its worker +records no end and exits with `EX_TEMPFAIL` (75), and the server resumes +the run, at most three times. A worker holds a lock on `worker.lock` in the +run's scratch directory for its whole life; before the server ends a +lease from outside (at that relaunch, or at a delete), it kills whatever +process still holds the lock and waits until it is gone, since a worker +outlives a server crash and keeps acting until its next write. ## How it is tested @@ -175,6 +185,10 @@ Integration tests live under `tests/`: command-only workflow on the host sandbox through the real step registry. Both acquire real Host scopes through the built-in in-process provider; no plugin executable or checksum is required. +- `resume.rs` resumes a run whose creation a crash cut short, which starts + it again, and runs a command workflow over a store whose first lease + write fails: the run is interrupted with no finish recorded, and a + resume finishes it. - `check.rs` admits the `hello` bundle and round-trips its graph through the blob store, binds the launch, admits a version whose `workflow.toml` names `engine = "petri"`, reads the project settings from the map, and @@ -237,8 +251,10 @@ the scheduler, in the server process under its test override, with `GET /runs/{id}/state` serving the projection over Petri's records; a human gate is answered through the questions API; and Petri's diagnostics refuse a run at create. -The server's `petri_runs` unit tests cover the lease ending at worker exit -and the restart reconcile that relaunches a worker in resume mode. +The server's `petri_runs` unit tests cover the lease ending at worker exit, +the restart reconcile that relaunches a worker in resume mode, and a worker +its store interrupted: relaunched in resume mode, and failed after the +bound. `lib/apps/fabro-server/tests/it/scenario/petri_stream.rs` covers the stream: a client attached to a two-branch parallel run disconnects once both branches started, a platform notice is recorded while both branch scripts @@ -255,6 +271,8 @@ executes in the worker a foreground server launched, its records reach `petri_records` over the HTTP store and its lease ends with the worker; and a run whose server and worker are both killed mid-stage resumes in a new worker after the server restarts, with one terminal lifecycle record; a +worker that outlives its server is stopped before the resume's worker +launches; a human gate in the worker is answered through the questions API over the control channel; two parallel gates each bind their own answer; and an unanswered gate expires with its default. The same file reads a finished From f12065e6ae8df96c7db7228b6aab1f8ec8f7870f Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Tue, 6 Oct 2026 12:01:01 -0400 Subject: [PATCH 5/6] Adapt recovery to current resume callers and API errors --- docs/internal/events-strategy.md | 5 ++++- lib/apps/fabro-server/src/server.rs | 8 +++++++- lib/apps/fabro-server/src/server/petri_runs.rs | 10 +++------- lib/components/fabro-petri/tests/hooks.rs | 2 +- lib/foundation/fabro-proc/src/process_lock.rs | 6 +++--- 5 files changed, 18 insertions(+), 13 deletions(-) diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index 4be756797..f6c404911 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -104,4 +104,7 @@ same transaction as the projection that consumed the record. A client that resumes from its last `stream_seq` sees every item exactly once. A worker cannot continue past a record it failed to append: the store's -error reaches the engine and fails the run. +error reaches the engine and interrupts the run's lifetime. The worker +records no terminal lifecycle transition for that interruption. The server +resumes the run from its durable records, up to three times per server +process; a fourth interruption fails the run. diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index a21ce3692..9f0024990 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -2793,7 +2793,13 @@ async fn delete_run_internal( // outlived a server crash, is stopped first. petri_runs::stop_previous_worker(state, id) .await - .map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, format!("{err:#}")))?; + .map_err(|err| { + error!(run_id = %id, error = %format!("{err:#}"), "Stopping the run's previous worker failed"); + ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + "failed to stop the run's previous worker", + ) + })?; state.petri_runs.worker_exited(id); let delete_outcome = delete_run_sandbox_resource(state, id, force).await?; diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 2d3201c35..31f4980ca 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -580,7 +580,7 @@ pub(crate) async fn reconcile_on_startup( run_id: RunId, run_state: &fabro_store::RunProjection, ) -> anyhow::Result<()> { - let mode = match relaunch(state, run_id, run_state).await? { + let mode = match relaunch(state, run_id).await? { Relaunch::Worker(mode) => mode, Relaunch::Failed { reason } => { warn!( @@ -660,7 +660,7 @@ pub(crate) async fn resume_after_interruption( { return Err(interrupted); } - let mode = match relaunch(state, run_id, &run_state).await { + let mode = match relaunch(state, run_id).await { Ok(Relaunch::Worker(mode)) => mode, Ok(Relaunch::Failed { reason }) => return Err(reason), Err(err) => { @@ -713,11 +713,7 @@ enum Relaunch { /// 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 { +async fn relaunch(state: &Arc, run_id: RunId) -> anyhow::Result { stop_previous_worker(state, run_id).await?; let held = match state.petri_runs.release_for_restart(run_id).await { Ok(()) => true, diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 0a118a3a9..d548f3a79 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -191,7 +191,7 @@ impl Harness { run_id: self.run_id.to_string(), run_dir: self.run_dir.clone(), execution: if resumed { - Execution::Resume + Execution::Resume(admit(workflow, settings)) } else { Execution::Start(admit(workflow, settings)) }, diff --git a/lib/foundation/fabro-proc/src/process_lock.rs b/lib/foundation/fabro-proc/src/process_lock.rs index ceaadc8d0..3b1873dc6 100644 --- a/lib/foundation/fabro-proc/src/process_lock.rs +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -18,7 +18,7 @@ use std::time::{Duration, Instant}; use tokio::fs::OpenOptions; use tokio::time; -use crate::signal::{sigkill, sigkill_process_group}; +use crate::signal; /// How often [`stop_lock_holder`] checks whether the lock is free. const POLL: Duration = Duration::from_millis(20); @@ -94,8 +94,8 @@ pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result Date: Tue, 6 Oct 2026 16:07:36 -0400 Subject: [PATCH 6/6] Simplify the store-fault handling - stop_lock_holder returns the stopped pid as Option instead of a LockHolder enum that only wrapped it - stop_previous_worker uses with_context, and one run_scratch helper replaces three spellings of the run's scratch path - relaunch builds its runnable record through run_records::runnable, moved out of the lifecycle handler so both callers share it - WorkerExit derives success from its exit code instead of storing both - the scenario tests share one worker-pid lookup and wait loop - small readability fixes in the engine's resume arm and the resume test Co-Authored-By: Claude Opus 5.5 --- lib/apps/fabro-cli/tests/it/scenario/petri.rs | 43 +++++++------------ lib/apps/fabro-server/src/petri_runs.rs | 1 - lib/apps/fabro-server/src/server.rs | 7 ++- .../src/server/handler/lifecycle.rs | 11 +---- .../fabro-server/src/server/petri_runs.rs | 37 ++++++---------- .../fabro-server/src/server/run_records.rs | 10 ++++- lib/apps/fabro-server/src/worker_runtime.rs | 15 ++++--- lib/components/fabro-petri/src/engine.rs | 6 +-- lib/components/fabro-petri/tests/resume.rs | 6 +-- lib/foundation/fabro-proc/src/lib.rs | 2 +- lib/foundation/fabro-proc/src/process_lock.rs | 28 +++++------- 11 files changed, 69 insertions(+), 97 deletions(-) diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index 8a3482286..c641054fa 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -595,6 +595,11 @@ pub(super) async fn settled_stream(server: &RunningServer, run_id: &str) -> Vec< /// worker retitles itself `fabro `, so that /// is what the process table shows. fn worker_pid(run_id: &str) -> Option { + worker_pid_other_than(run_id, None) +} + +/// [`worker_pid`], skipping the worker `previous` when given. +fn worker_pid_other_than(run_id: &str, previous: Option) -> Option { let short_id: String = run_id.chars().take(12).collect(); let output = Command::new("pgrep") .args(["-f", &format!("^fabro {short_id} ")]) @@ -602,18 +607,24 @@ fn worker_pid(run_id: &str) -> Option { .expect("pgrep runs"); String::from_utf8_lossy(&output.stdout) .lines() - .find_map(|line| line.trim().parse().ok()) + .filter_map(|line| line.trim().parse().ok()) + .find(|pid| Some(*pid) != previous) } pub(super) fn wait_for_worker(run_id: &str) -> u32 { + wait_for_worker_except(run_id, None) +} + +/// Wait for a worker of the run other than `previous`, when given. +fn wait_for_worker_except(run_id: &str, previous: Option) -> u32 { let deadline = Instant::now() + RUN_TIMEOUT; loop { - if let Some(pid) = worker_pid(run_id) { + if let Some(pid) = worker_pid_other_than(run_id, previous) { return pid; } assert!( Instant::now() < deadline, - "no worker process appeared for run {run_id}" + "no worker process other than {previous:?} appeared for run {run_id}" ); std::thread::sleep(POLL); } @@ -805,7 +816,7 @@ async fn a_worker_that_outlives_the_server_is_stopped_before_its_run_resumes() { server.launch().await; eprintln!("server restarted"); - let resumed = wait_for_worker_other_than(&run_id, worker); + let resumed = wait_for_worker_except(&run_id, Some(worker)); eprintln!("worker {resumed} launched for the resume"); assert!( !fabro_proc::process_running_strict(worker), @@ -826,30 +837,6 @@ async fn a_worker_that_outlives_the_server_is_stopped_before_its_run_resumes() { server.shutdown(); } -/// Wait for a worker of the run other than `previous`. -fn wait_for_worker_other_than(run_id: &str, previous: u32) -> u32 { - let short_id: String = run_id.chars().take(12).collect(); - let deadline = Instant::now() + RUN_TIMEOUT; - loop { - let output = Command::new("pgrep") - .args(["-f", &format!("^fabro {short_id} ")]) - .output() - .expect("pgrep runs"); - if let Some(pid) = String::from_utf8_lossy(&output.stdout) - .lines() - .filter_map(|line| line.trim().parse::().ok()) - .find(|pid| *pid != previous) - { - return pid; - } - assert!( - Instant::now() < deadline, - "no new worker process appeared for run {run_id}" - ); - std::thread::sleep(POLL); - } -} - /// The run's status while the server may be down: `None` when it is. async fn run_status_offline(server: &RunningServer) -> Option { fabro_test::test_http_client() diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index b977b1d7f..28380e8f2 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -258,7 +258,6 @@ mod tests { exit.notified().await; let code = *sync::lock(&code); Ok(WorkerExit { - success: code == Some(0), code, detail: "test worker ended without a terminal event".to_string(), }) diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 9f0024990..2a2ae3079 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -4341,8 +4341,11 @@ async fn execute_run_subprocess(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) { - managed_run.status = - status_after_worker_exit(managed_run.status, final_state.status, worker_exit.success); + managed_run.status = status_after_worker_exit( + managed_run.status, + final_state.status, + worker_exit.succeeded(), + ); managed_run.error = final_state .conclusion .as_ref() diff --git a/lib/apps/fabro-server/src/server/handler/lifecycle.rs b/lib/apps/fabro-server/src/server/handler/lifecycle.rs index 1d53aec93..b5b535c93 100644 --- a/lib/apps/fabro-server/src/server/handler/lifecycle.rs +++ b/lib/apps/fabro-server/src/server/handler/lifecycle.rs @@ -153,7 +153,7 @@ pub(super) async fn queue_run( let next = if approval_required { run_records::transition(RunLifecycleKind::Pending, next_status) } else { - runnable(RunRunnableSource::StartRequested) + run_records::runnable(RunRunnableSource::StartRequested) }; for record in [start_requested, next] { if let Err(err) = run_records::lifecycle(state, id, record).await { @@ -213,7 +213,7 @@ async fn approve_run( for record in [ RunLifecycleRecord::new(RunLifecycleKind::Approved), - runnable(RunRunnableSource::Approved), + run_records::runnable(RunRunnableSource::Approved), ] { if let Err(err) = run_records::lifecycle(state.as_ref(), id, record).await { return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) @@ -1072,13 +1072,6 @@ async fn archive_status_response(state: &AppState, id: RunId) -> Response { run_response(state, id, StatusCode::OK).await } -/// The runnable transition, with what made the run runnable. -fn runnable(source: RunRunnableSource) -> RunLifecycleRecord { - let mut record = run_records::transition(RunLifecycleKind::Runnable, RunStatus::Runnable); - record.source = Some(<&'static str>::from(source).to_string()); - record -} - /// Persist a synchronous pause/unpause transition: record it and mirror the /// new status in the in-memory run map. Returns `Some(Response)` on error, /// `None` on success. diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 31f4980ca..356054d94 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -40,8 +40,9 @@ use std::collections::HashMap; use std::sync::Arc; use std::time::{Duration, Instant}; +use anyhow::Context as _; use fabro_config::{ - EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, SettingsLayer, Storage, + EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, RunScratch, SettingsLayer, Storage, }; use fabro_interview::ControlInterviewer; use fabro_petri::artifacts::StoreArtifactWriter; @@ -605,7 +606,7 @@ pub(crate) async fn reconcile_on_startup( run_state.spec.graph_source.clone().unwrap_or_default(), RunStatus::Runnable, run_id.created_at(), - scratch_root(state, run_id), + run_scratch(state, run_id).root().to_path_buf(), mode, ), ); @@ -750,8 +751,7 @@ async fn relaunch(state: &Arc, run_id: RunId) -> anyhow::Result::from(RunRunnableSource::StartRequested).to_string()); + let runnable = run_records::runnable(RunRunnableSource::StartRequested); for record in [start_requested, runnable] { run_records::lifecycle(state, run_id, record).await?; } @@ -772,18 +772,11 @@ const WORKER_STOP_PATIENCE: Duration = Duration::from_secs(10); pub(crate) async fn stop_previous_worker(state: &AppState, run_id: RunId) -> anyhow::Result<()> { #[cfg(unix)] { - let path = Storage::new(state.server_storage_dir()) - .run_scratch(&run_id) - .worker_lock_path(); - let holder = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE) + let path = run_scratch(state, run_id).worker_lock_path(); + let stopped = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE) .await - .map_err(|err| { - anyhow::Error::new(err).context(format!( - "stopping the run's previous worker ({})", - path.display() - )) - })?; - if let fabro_proc::LockHolder::Stopped { pid } = holder { + .with_context(|| format!("stopping the run's previous worker ({})", path.display()))?; + if let Some(pid) = stopped { warn!( run_id = %run_id, pid, @@ -796,12 +789,9 @@ pub(crate) async fn stop_previous_worker(state: &AppState, run_id: RunId) -> any Ok(()) } -/// 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 run's scratch: its worker's run directory and worker lock. +fn run_scratch(state: &AppState, run_id: RunId) -> RunScratch { + Storage::new(state.server_storage_dir()).run_scratch(&run_id) } /// The failed status, its message, and the `failed` lifecycle record for it. @@ -979,10 +969,7 @@ mod tests { async fn in_flight_run() -> (Arc, RunId, SettlingLogs, Arc) { let state = TestAppStateBuilder::new().in_process_execution().build(); let run_id = RunId::new(); - let run_dir = Storage::new(state.server_storage_dir()) - .run_scratch(&run_id) - .root() - .to_path_buf(); + let run_dir = super::run_scratch(&state, run_id).root().to_path_buf(); state.runs.lock().expect("runs lock poisoned").insert( run_id, super::super::managed_run( diff --git a/lib/apps/fabro-server/src/server/run_records.rs b/lib/apps/fabro-server/src/server/run_records.rs index ad46f02aa..a32c4b93c 100644 --- a/lib/apps/fabro-server/src/server/run_records.rs +++ b/lib/apps/fabro-server/src/server/run_records.rs @@ -16,7 +16,7 @@ use fabro_store::RunProjection; use fabro_store::platform_records::{ PlatformRecord, RunLifecycleKind, RunLifecycleRecord, StoredPlatformRecord, }; -use fabro_types::{FailureReason, RunId, RunStatus, SuccessReason}; +use fabro_types::{FailureReason, RunId, RunRunnableSource, RunStatus, SuccessReason}; use super::AppState; use crate::error::ApiError; @@ -54,6 +54,14 @@ pub(crate) fn transition(kind: RunLifecycleKind, status: RunStatus) -> RunLifecy RunLifecycleRecord::new(kind).with_status(status) } +/// The runnable transition, with what made the run runnable. +#[must_use] +pub(crate) fn runnable(source: RunRunnableSource) -> RunLifecycleRecord { + let mut record = transition(RunLifecycleKind::Runnable, RunStatus::Runnable); + record.source = Some(<&'static str>::from(source).to_string()); + record +} + /// The run failed for `reason`, with `message` as the failure's detail. #[must_use] pub(crate) fn failed(reason: FailureReason, message: impl Into) -> RunLifecycleRecord { diff --git a/lib/apps/fabro-server/src/worker_runtime.rs b/lib/apps/fabro-server/src/worker_runtime.rs index 9d2029340..1171e7153 100644 --- a/lib/apps/fabro-server/src/worker_runtime.rs +++ b/lib/apps/fabro-server/src/worker_runtime.rs @@ -62,13 +62,17 @@ 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, + pub(crate) code: Option, + pub(crate) detail: String, } impl WorkerExit { + /// Whether the worker exited with status zero. + pub(crate) fn succeeded(&self) -> bool { + self.code == Some(0) + } + /// 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 { @@ -147,9 +151,8 @@ impl WorkerRuntime for LocalWorkerRuntime { let wait: BoxFuture<'static, Result> = Box::pin(async move { let status = child.wait().await.context("worker wait failed")?; Ok(WorkerExit { - success: status.success(), - code: status.code(), - detail: status.to_string(), + 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 32d707583..f9f2380ed 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -316,14 +316,14 @@ pub async fn run(request: RunRequest) -> Result { } Execution::Resume(graphs) => { info!(run_id = %request.run_id, backend = %backend, "Resuming Petri run"); - let resumed = Box::pin(host::resume_configured( + let outcome = Box::pin(host::resume_configured( &runtime, Vec::new(), observers.clone(), wiring(), )) .await; - match resumed { + match outcome { // A crash cut the run's creation short: nothing beyond its // start is stored, so it starts again from its admitted // graphs, and Petri takes the stored prefix over. @@ -334,7 +334,7 @@ pub async fn run(request: RunRequest) -> Result { ); start(&runtime, graphs, observers, wiring()).await } - resumed => resumed, + outcome => outcome, } } }; diff --git a/lib/components/fabro-petri/tests/resume.rs b/lib/components/fabro-petri/tests/resume.rs index 9f07cb148..bf607dfed 100644 --- a/lib/components/fabro-petri/tests/resume.rs +++ b/lib/components/fabro-petri/tests/resume.rs @@ -141,10 +141,10 @@ fn resume_request( runtime, no_questions(Arc::new(Silent)), ); - request.execution = match request.execution { - Execution::Start(graphs) => Execution::Resume(graphs), - resume @ Execution::Resume(_) => resume, + let Execution::Start(graphs) = request.execution else { + unreachable!("run_request builds a start"); }; + request.execution = Execution::Resume(graphs); request } diff --git a/lib/foundation/fabro-proc/src/lib.rs b/lib/foundation/fabro-proc/src/lib.rs index ff3103f5f..6ca06b1af 100644 --- a/lib/foundation/fabro-proc/src/lib.rs +++ b/lib/foundation/fabro-proc/src/lib.rs @@ -25,7 +25,7 @@ pub use pre_exec::pre_exec_setpgid; #[cfg(unix)] pub use pre_exec::pre_exec_setsid; #[cfg(unix)] -pub use process_lock::{LockHolder, ProcessLock, stop_lock_holder}; +pub use process_lock::{ProcessLock, stop_lock_holder}; pub use signal::{process_exists, process_group_alive, process_running, process_running_strict}; #[cfg(unix)] pub use signal::{ diff --git a/lib/foundation/fabro-proc/src/process_lock.rs b/lib/foundation/fabro-proc/src/process_lock.rs index 3b1873dc6..90c7f7df6 100644 --- a/lib/foundation/fabro-proc/src/process_lock.rs +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -69,37 +69,29 @@ impl ProcessLock { } } -/// What [`stop_lock_holder`] found. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum LockHolder { - /// No other process held the lock. - None, - /// This process held the lock. It was killed, with its process group, - /// and the lock is free: it is gone. - Stopped { pid: u32 }, -} - /// Stop the process that holds the lock on the file at `path`, if another /// process does, and wait up to `patience` for the lock to be free. The /// holder and its process group get `SIGKILL`: a holder that could handle /// a signal could also keep running. A missing file has no holder. /// -/// `Err` when the lock is still held after `patience`. -pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result { +/// `Ok(Some(pid))` names the holder that was stopped: it is gone and the +/// lock is free. `Ok(None)` when no other process held the lock. `Err` +/// when the lock is still held after `patience`. +pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result> { let file = match OpenOptions::new().read(true).write(true).open(path).await { Ok(file) => file.into_std().await, - Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(LockHolder::None), + Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None), Err(error) => return Err(error), }; let Some(pid) = holder(&file)? else { - return Ok(LockHolder::None); + return Ok(None); }; signal::sigkill_process_group(pid); signal::sigkill(pid); let deadline = Instant::now() + patience; loop { match holder(&file)? { - None => return Ok(LockHolder::Stopped { pid }), + None => return Ok(Some(pid)), Some(_) if Instant::now() >= deadline => { return Err(io::Error::other(format!( "process {pid} still holds {} after SIGKILL", @@ -210,7 +202,7 @@ mod tests { stop_lock_holder(&path, Duration::from_secs(1)) .await .expect("the check runs"), - LockHolder::None + None ); drop( ProcessLock::try_hold(&path) @@ -222,7 +214,7 @@ mod tests { stop_lock_holder(&path, Duration::from_secs(1)) .await .expect("the check runs"), - LockHolder::None + None ); } @@ -250,7 +242,7 @@ mod tests { .await .expect("the holder is stopped"); - assert_eq!(found, LockHolder::Stopped { pid }); + assert_eq!(found, Some(pid)); let status = holder.wait().await.expect("the holder is reaped"); assert!(!status.success(), "the holder was killed: {status}"); assert!(