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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-28 10:05:40 -04:00 • committed by Scott Werner
parent d89d0c577b
commit 3c395f9e6e
4 changed files with 154 additions and 77 deletions

View file

@ -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();

View file

@ -391,18 +391,19 @@ pub(crate) async fn execute(state: Arc<AppState>, 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.

View file

@ -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<RunOutcome, RunError> {
}
// 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<RunOutcome, RunError> {
.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<Option<JoinHandle<()>>> = 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<RunOutcome, RunError> {
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<SandboxBackend>
}
}
/// 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<Arc<dyn ExecutionObserver>>,
with_handle: impl FnOnce(petri_execution::CoordinatorHandle, Arc<dyn SecretProvider>),
) -> Result<ExecutionReport, HostError> {
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.

View file

@ -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);
}