From 3c395f9e6ea3047f649f0ace561f3fb333562c41 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 28 Sep 2026 10:05:40 -0400 Subject: [PATCH] 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); +}