Merge pull request #907 from fabro-sh/petri-store-faults
Some checks are pending
Rust / Test (macOS) (push) Waiting to run
Rust / Process titles (musl) (push) Waiting to run
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run

Handle Petri store faults: start again, resume after a failed write, stop surviving workers
This commit is contained in:
Scott Werner 2026-10-06 16:43:06 -04:00 • committed by GitHub
commit eac82ee71c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
19 changed files with 1172 additions and 193 deletions

View file

@ -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.

View file

@ -9,12 +9,16 @@
//!
//! 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.
//! 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
@ -93,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;
@ -123,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();
@ -180,21 +187,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();
@ -291,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");
(

View file

@ -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::<ServerTarget>()?;
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<fabro_proc::ProcessLock> {
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";

View file

@ -595,6 +595,11 @@ pub(super) async fn settled_stream(server: &RunningServer, run_id: &str) -> Vec<
/// worker retitles itself `fabro <first 12 of the run id> <phase>`, so that
/// is what the process table shows.
fn worker_pid(run_id: &str) -> Option<u32> {
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<u32>) -> Option<u32> {
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<u32> {
.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>) -> 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);
}
@ -779,6 +790,53 @@ 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_except(&run_id, Some(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();
}
/// The run's status while the server may be down: `None` when it is.
async fn run_status_offline(server: &RunningServer) -> Option<String> {
fabro_test::test_http_client()

View file

@ -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
@ -174,6 +179,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 +187,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 +216,8 @@ mod tests {
started: Notify,
running: AtomicBool,
exit: Arc<Notify>,
/// The exit code the worker ends with; `None` for a signal.
code: Arc<Mutex<Option<i32>>>,
mode: Mutex<Option<&'static str>>,
}
@ -220,6 +229,11 @@ mod tests {
}
fn end_worker(&self) {
self.end_worker_with(None);
}
fn end_worker_with(&self, code: Option<i32>) {
*sync::lock(&self.code) = code;
self.running.store(false, Ordering::SeqCst);
self.exit.notify_one();
}
@ -235,15 +249,17 @@ 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<dyn AsyncRead + Send + 'static>> = 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(),
code,
detail: "test worker ended without a terminal event".to_string(),
})
}),
};
@ -804,4 +820,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<AppState>, run_id: RunId) -> (RunStatus, Option<String>) {
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::<Vec<_>>();
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"
);
}
}

View file

@ -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<WorkerRef>,
run_dir: Option<std::path::PathBuf>,
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 {
@ -2786,8 +2788,18 @@ 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| {
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?;
@ -3535,6 +3547,7 @@ fn managed_run(
cancel_escalation_worker: None,
run_dir: Some(run_dir),
execution_mode,
store_interruptions: 0,
}
}
@ -3784,8 +3797,9 @@ async fn fail_worker_launch(state: &Arc<AppState>, 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 +3812,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 +4296,20 @@ async fn execute_run_subprocess(state: Arc<AppState>, 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,
@ -4320,8 +4341,11 @@ async fn execute_run_subprocess(state: Arc<AppState>, 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()

View file

@ -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.

View file

@ -30,14 +30,19 @@
//! 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;
use std::time::Instant;
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;
@ -57,7 +62,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 +77,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
@ -391,18 +397,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.
@ -519,6 +526,14 @@ pub(crate) async fn execute(state: Arc<AppState>, 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");
(
@ -552,29 +567,155 @@ pub(crate) async fn execute(state: Arc<AppState>, 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<AppState>,
run_id: RunId,
run_state: &fabro_store::RunProjection,
) -> anyhow::Result<()> {
let key = PetriRuns::key(&run_id);
let mode = match relaunch(state, run_id).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(),
run_scratch(state, run_id).root().to_path_buf(),
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<AppState>,
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).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 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
/// 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<AppState>, run_id: RunId) -> anyhow::Result<Relaunch> {
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,
@ -582,10 +723,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,
@ -607,48 +744,56 @@ 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);
runnable.source = Some(<&'static str>::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?;
}
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(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 = run_scratch(state, run_id).worker_lock_path();
let stopped = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE)
.await
.with_context(|| format!("stopping the run's previous worker ({})", path.display()))?;
if let Some(pid) = stopped {
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: 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.
fn failed(
reason: FailureReason,
@ -824,10 +969,7 @@ mod tests {
async fn in_flight_run() -> (Arc<AppState>, RunId, SettlingLogs, Arc<RecordingLogs>) {
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(

View file

@ -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<String>) -> RunLifecycleRecord {

View file

@ -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;
@ -61,8 +62,22 @@ pub(crate) struct StartedWorker {
#[derive(Debug)]
pub(crate) struct WorkerExit {
pub(crate) success: bool,
pub(crate) detail: String,
/// The worker's exit code; `None` when a signal ended it.
pub(crate) code: Option<i32>,
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 {
self.code == Some(ExitClass::Interrupted.code())
}
}
#[derive(Default)]
@ -136,8 +151,8 @@ impl WorkerRuntime for LocalWorkerRuntime {
let wait: BoxFuture<'static, Result<WorkerExit>> = Box::pin(async move {
let status = child.wait().await.context("worker wait failed")?;
Ok(WorkerExit {
success: status.success(),
detail: status.to_string(),
code: status.code(),
detail: status.to_string(),
})
});

View file

@ -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

View file

@ -11,10 +11,20 @@
//! 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.
//!
//! 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))
@ -45,25 +55,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,
Access, CancelReason, CoordinatorError, 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 +93,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,12 +172,15 @@ 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("; "))]
Unfinished(Vec<String>),
/// 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
@ -179,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 {
@ -217,7 +237,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 +286,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 +312,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 outcome = Box::pin(host::resume_configured(
&runtime,
Vec::new(),
observers,
with_handle,
observers.clone(),
wiring(),
))
.await
.await;
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.
Err(HostError::NotStarted) => {
info!(
run_id = %request.run_id,
"The Petri run never started; starting it again"
);
start(&runtime, graphs, observers, wiring()).await
}
outcome => outcome,
}
}
};
if let Some(task) = cancel_task {
if let Some(task) = sync::lock(&cancel_task).take() {
task.abort();
}
let receipt = dispatcher.shutdown().await;
@ -313,6 +347,11 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
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
@ -379,9 +418,9 @@ pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result<RunOutcome
}
/// How Fabro reports what [`run`] returned. A cancelled run is a failure
/// with the cancelled reason, as the legacy executor reports one; every
/// other shortfall is a workflow error whose message says what the record,
/// or the host, said.
/// with the cancelled reason, as the legacy executor reports one; a failed
/// store interrupted the run; every other shortfall is a workflow error
/// whose message says what the record, or the host, said.
#[must_use]
pub fn conclusion(result: &Result<RunOutcome, RunError>) -> Conclusion {
match result {
@ -401,6 +440,9 @@ pub fn conclusion(result: &Result<RunOutcome, RunError>) -> 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),
@ -469,21 +511,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.
@ -625,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(),
});
}
}

View file

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

View file

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

View file

@ -0,0 +1,204 @@
//! 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 std::sync::atomic::{AtomicUsize, Ordering};
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, Conclusion, Execution, RunError, RunStatus};
use fabro_petri::runtime::RuntimeSpec;
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 {
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 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<dyn RunStore>,
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<dyn RunStore>,
runtime: RuntimeSpec,
) -> engine::RunRequest {
let mut request = run_request(
run_id,
run_dir,
graphs,
store,
runtime,
no_questions(Arc::new(Silent)),
);
let Execution::Start(graphs) = request.execution else {
unreachable!("run_request builds a start");
};
request.execution = Execution::Resume(graphs);
request
}
/// The in-memory store, failing its first append to `log` as a full disk
/// would, before anything is stored.
struct FailingStore {
inner: Arc<MemoryRunStore>,
log: LogId,
appended: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl RunStore for FailingStore {
async fn open(&self, key: &RunKey, access: Access) -> Result<Arc<dyn RunLogs>, StoreError> {
Ok(Arc::new(FailingLogs {
inner: self.inner.open(key, access).await?,
log: self.log,
appended: Arc::clone(&self.appended),
}))
}
}
struct FailingLogs {
inner: Arc<dyn RunLogs>,
log: LogId,
appended: Arc<AtomicUsize>,
}
#[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<Vec<Record>, StoreError> {
self.inner.read(log).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<Digest, StoreError> {
self.inner.put_blob(bytes).await
}
async fn get_blob(&self, digest: Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}

View file

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

View file

@ -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::{ProcessLock, stop_lock_holder};
pub use signal::{process_exists, process_group_alive, process_running, process_running_strict};
#[cfg(unix)]
pub use signal::{

View file

@ -0,0 +1,256 @@
//! 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;
/// 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<Option<Self>> {
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),
}
}
}
/// 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.
///
/// `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<Option<u32>> {
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(None),
Err(error) => return Err(error),
};
let Some(pid) = holder(&file)? else {
return Ok(None);
};
signal::sigkill_process_group(pid);
signal::sigkill(pid);
let deadline = Instant::now() + patience;
loop {
match holder(&file)? {
None => return Ok(Some(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<Option<u32>> {
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"),
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"),
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, Some(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"
);
}
}

View file

@ -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::<Classified>())
.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<ExitClass> {
@ -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!(