mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Resume a Petri run whose store failed, as after a crash
Petri now ends a run's lifetime at its first failed store write: it records nothing after it, fails no firing for it, and returns CoordinatorError::StoreFailed. The run is not over; the next lifetime resumes it from what the store holds. Fabro read that error as an unfinished run and failed it. - engine: RunError::StoreFailed, returned without reading the record back, and Conclusion::Interrupted for it. - worker: an interrupted run gets no terminal lifecycle record; the worker exits with EX_TEMPFAIL (75, the new ExitClass::Interrupted). - server: WorkerExit carries the exit code. An interrupted worker's run goes back to the scheduler in resume mode through the relaunch a restart takes (lease release, recovery, start_requested + runnable), now shared with reconcile_on_startup. The in-process path does the same. A run is resumed at most MAX_STORE_INTERRUPTIONS (3) times per server; the next interruption fails it. A pending cancel, a run that ended or was deleted, and a shutdown also end it as before. Tests: an engine run over a store whose first lease write fails is interrupted with no finish, and a resume finishes it; the server relaunches an interrupted worker in resume mode, fails the run after the bound, and fails a worker that exits 1 as before; exit code 75. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
parent
3c395f9e6e
commit
37116758fa
9 changed files with 568 additions and 87 deletions
|
|
@ -15,7 +15,10 @@
|
|||
//! from the graphs when a crash cut its creation short. Either way the
|
||||
//! worker records the lifecycle transitions Fabro's read side needs
|
||||
//! (`starting`, `running`, then `succeeded` or `failed`) as platform
|
||||
//! records through the client.
|
||||
//! records through the client. A failed write to the run's store ends the
|
||||
//! run's lifetime and not the run: the worker records no end for it and
|
||||
//! exits with `EX_TEMPFAIL` (75, [`ExitClass::Interrupted`]), and the server
|
||||
//! launches a worker to resume it, as after a crash.
|
||||
//!
|
||||
//! The server's controls arrive over the control channel and go to Petri
|
||||
//! through [`PetriControls`]: cancel (and `SIGTERM`/`SIGINT`) fires one
|
||||
|
|
@ -94,6 +97,7 @@ use fabro_store::platform_records::{
|
|||
};
|
||||
use fabro_types::settings::run::{ApprovalMode, RunMode};
|
||||
use fabro_types::{FailureReason, Principal, RunId, RunNoticeLevel, RunStatus, SuccessReason};
|
||||
use fabro_util::exit::{ErrorExt as _, ExitClass};
|
||||
use fabro_vault::Vault;
|
||||
use fabro_workflow::Error as WorkflowError;
|
||||
use fabro_workflow::services::FabroRunToolServices;
|
||||
|
|
@ -124,7 +128,9 @@ pub(super) struct PetriWorker<'a> {
|
|||
|
||||
/// Execute the run to its end. `Ok` when the record says it succeeded;
|
||||
/// the failure otherwise, after the terminal event is appended, so the
|
||||
/// worker exits as the legacy worker does for a failed run.
|
||||
/// worker exits as the legacy worker does for a failed run. A run its
|
||||
/// store interrupted gets no terminal event, and its error is classified
|
||||
/// [`ExitClass::Interrupted`] for the server to resume the run.
|
||||
pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
|
||||
let run_id = worker.run_id;
|
||||
let admission = worker.run_state.spec.admission.clone();
|
||||
|
|
@ -290,6 +296,14 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
|
|||
"Petri run ended"
|
||||
);
|
||||
let (record, phase, failure) = match engine::conclusion(&result) {
|
||||
Conclusion::Interrupted { message } => {
|
||||
// The run is not over: it continues from its records in the
|
||||
// worker the server launches next. That also holds when the
|
||||
// control channel was lost: a cancel the loss requested is in
|
||||
// the records, or the run was not cancelled.
|
||||
warn!(run_id = %run_id, error = %message, "Petri run interrupted; the server resumes it");
|
||||
return Err(anyhow!("{message}").classify(ExitClass::Interrupted));
|
||||
}
|
||||
Conclusion::Succeeded => {
|
||||
info!(run_id = %run_id, "Petri run completed");
|
||||
(
|
||||
|
|
|
|||
|
|
@ -174,6 +174,7 @@ mod tests {
|
|||
use fabro_types::{
|
||||
FailureReason, RunId, RunStatus, SuccessReason, WorkflowPath, WorkflowVersion,
|
||||
};
|
||||
use fabro_util::exit::ExitClass;
|
||||
use serde_json::json;
|
||||
use tokio::io::AsyncRead;
|
||||
use tokio::sync::Notify;
|
||||
|
|
@ -181,6 +182,7 @@ mod tests {
|
|||
use tower::ServiceExt as _;
|
||||
|
||||
use super::*;
|
||||
use crate::server::petri_runs::MAX_STORE_INTERRUPTIONS;
|
||||
use crate::server::{
|
||||
AppState, reconcile_incomplete_runs_on_startup, run_records, spawn_scheduler,
|
||||
};
|
||||
|
|
@ -209,6 +211,8 @@ mod tests {
|
|||
started: Notify,
|
||||
running: AtomicBool,
|
||||
exit: Arc<Notify>,
|
||||
/// The exit code the worker ends with; `None` for a signal.
|
||||
code: Arc<Mutex<Option<i32>>>,
|
||||
mode: Mutex<Option<&'static str>>,
|
||||
}
|
||||
|
||||
|
|
@ -220,6 +224,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 +244,18 @@ mod tests {
|
|||
*sync::lock(&self.mode) = Some(spec.mode);
|
||||
self.running.store(true, Ordering::SeqCst);
|
||||
let exit = Arc::clone(&self.exit);
|
||||
let code = Arc::clone(&self.code);
|
||||
let stderr: Pin<Box<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(),
|
||||
success: code == Some(0),
|
||||
code,
|
||||
detail: "test worker ended without a terminal event".to_string(),
|
||||
})
|
||||
}),
|
||||
};
|
||||
|
|
@ -804,4 +816,124 @@ mod tests {
|
|||
.expect("the run projects");
|
||||
assert_eq!(run_state.status, succeeded);
|
||||
}
|
||||
|
||||
/// Wait until the run's stored status is terminal, and return it with
|
||||
/// the failure's message.
|
||||
async fn terminal_status(state: &Arc<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"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
@ -3535,6 +3537,7 @@ fn managed_run(
|
|||
cancel_escalation_worker: None,
|
||||
run_dir: Some(run_dir),
|
||||
execution_mode,
|
||||
store_interruptions: 0,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -3784,8 +3787,9 @@ async fn fail_worker_launch(state: &Arc<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 +3802,7 @@ async fn append_worker_exit_failure(state: &AppState, run_id: RunId, worker_exit
|
|||
return;
|
||||
}
|
||||
|
||||
let (error, reason) = failure_for_incomplete_run(
|
||||
run_state.pending_control,
|
||||
format!(
|
||||
"Worker exited before emitting a terminal run event: {}",
|
||||
worker_exit.detail
|
||||
),
|
||||
);
|
||||
let (error, reason) = failure_for_incomplete_run(run_state.pending_control, failure);
|
||||
if let Err(err) = run_records::lifecycle(
|
||||
state,
|
||||
run_id,
|
||||
|
|
@ -4288,7 +4286,20 @@ async fn execute_run_subprocess(state: Arc<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,
|
||||
|
|
|
|||
|
|
@ -30,7 +30,11 @@
|
|||
//! previous server left in flight back to a worker in resume mode, once the
|
||||
//! recovery protocol (`fabro_petri::recovery`) has brought every live
|
||||
//! workspace to the snapshot its durable state names, or reports the run
|
||||
//! failed when it cannot.
|
||||
//! failed when it cannot. A run whose store failed under it takes the same
|
||||
//! way back ([`resume_after_interruption`]): its worker exits with
|
||||
//! `EX_TEMPFAIL` and records no end, since a failed store write ends the
|
||||
//! run's lifetime and not the run, and the server resumes it, at most
|
||||
//! [`MAX_STORE_INTERRUPTIONS`] times.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
|
@ -57,7 +61,9 @@ use fabro_static::EnvVars;
|
|||
use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord};
|
||||
use fabro_types::settings::McpTransport;
|
||||
use fabro_types::settings::run::{ApprovalMode, McpServerSettings, RunMode};
|
||||
use fabro_types::{FailureReason, RunId, RunRunnableSource, RunStatus, RunTarget, SuccessReason};
|
||||
use fabro_types::{
|
||||
FailureReason, RunControlAction, RunId, RunRunnableSource, RunStatus, RunTarget, SuccessReason,
|
||||
};
|
||||
use fabro_util::error as error_util;
|
||||
use fabro_workflow::Error as WorkflowError;
|
||||
use lithos_llm::catalog::ProviderId;
|
||||
|
|
@ -70,7 +76,6 @@ use super::{
|
|||
stream_follower,
|
||||
};
|
||||
use crate::petri_check;
|
||||
use crate::petri_runs::PetriRuns;
|
||||
use crate::run_compiler::{AdmittedRun, PreparedRun, RunCompilerError};
|
||||
|
||||
/// The runtime Petri gets, at create and at execution: the server's run
|
||||
|
|
@ -520,6 +525,14 @@ pub(crate) async fn execute(state: Arc<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");
|
||||
(
|
||||
|
|
@ -553,29 +566,157 @@ 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, run_state).await? {
|
||||
Relaunch::Worker(mode) => mode,
|
||||
Relaunch::Failed { reason } => {
|
||||
warn!(
|
||||
run_id = %run_id,
|
||||
error = %reason,
|
||||
"Petri run left in flight by the previous server cannot resume; reporting it failed"
|
||||
);
|
||||
let (_, _, record) = failed(FailureReason::WorkflowError, reason);
|
||||
run_records::lifecycle(state, run_id, record).await?;
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
info!(
|
||||
run_id = %run_id,
|
||||
mode = super::worker_mode_arg(mode),
|
||||
"Petri run left in flight by the previous server; relaunching its worker"
|
||||
);
|
||||
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
||||
runs.insert(
|
||||
run_id,
|
||||
super::managed_run(
|
||||
run_state.spec.graph_source.clone().unwrap_or_default(),
|
||||
RunStatus::Runnable,
|
||||
run_id.created_at(),
|
||||
scratch_root(state, run_id),
|
||||
mode,
|
||||
),
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Resume a run whose store failed under it, as after a crash: its worker
|
||||
/// (or the in-process engine) ended the run's lifetime with nothing
|
||||
/// recorded after the failure and no end for the run. The run goes back
|
||||
/// to the scheduler through the same [`relaunch`] a restart takes.
|
||||
///
|
||||
/// `Err` with the failure to record when the run does not resume: it was
|
||||
/// deleted or ended meanwhile, a cancel is pending, the server is shutting
|
||||
/// down, its store interrupted it more than [`MAX_STORE_INTERRUPTIONS`]
|
||||
/// times, or it cannot continue from its records.
|
||||
pub(crate) async fn resume_after_interruption(
|
||||
state: &Arc<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, &run_state).await {
|
||||
Ok(Relaunch::Worker(mode)) => mode,
|
||||
Ok(Relaunch::Failed { reason }) => return Err(reason),
|
||||
Err(err) => {
|
||||
return Err(format!(
|
||||
"{interrupted}; the run could not be resumed: {err:#}"
|
||||
));
|
||||
}
|
||||
};
|
||||
warn!(
|
||||
run_id = %run_id,
|
||||
interruptions,
|
||||
error = message,
|
||||
mode = super::worker_mode_arg(mode),
|
||||
"the run's store interrupted it; resuming it as after a crash"
|
||||
);
|
||||
{
|
||||
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
||||
let Some(managed_run) = runs.get_mut(&run_id) else {
|
||||
return Err(interrupted);
|
||||
};
|
||||
managed_run.status = RunStatus::Runnable;
|
||||
managed_run.execution_mode = mode;
|
||||
clear_live_run_state(managed_run);
|
||||
}
|
||||
state.scheduler_notify.notify_one();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// How a run whose lifetime ended short of its end continues.
|
||||
enum Relaunch {
|
||||
/// A new worker takes it, in this mode.
|
||||
Worker(RunExecutionMode),
|
||||
/// Its records say it cannot continue.
|
||||
Failed { reason: String },
|
||||
}
|
||||
|
||||
/// Ready a run whose lifetime ended short of its end for a new worker, as a
|
||||
/// crash is recovered.
|
||||
///
|
||||
/// The lease the previous worker held is released from outside, which
|
||||
/// fences that worker should it still be alive. Then the recovery protocol
|
||||
/// reads the run's durable execution state: a run with a failed checkpoint
|
||||
/// cannot continue; otherwise every live workspace on this host is verified
|
||||
/// against, reset to, or restored from the snapshot its last durable finish
|
||||
/// names, and a finish with no snapshot cannot continue rather than resume
|
||||
/// on stale files. A run that continues is asked to start again as a resume
|
||||
/// (`run.start_requested` with `resume`, then `run.runnable`, the same pair
|
||||
/// the API's resume appends), in resume mode when Petri's store holds the
|
||||
/// run, else in start mode: a worker that died before it created the run's
|
||||
/// record left nothing to continue from, so the run starts from its
|
||||
/// admitted graphs. The caller registers the run with the scheduler.
|
||||
async fn relaunch(
|
||||
state: &Arc<AppState>,
|
||||
run_id: RunId,
|
||||
run_state: &fabro_store::RunProjection,
|
||||
) -> anyhow::Result<Relaunch> {
|
||||
let held = match state.petri_runs.release_for_restart(run_id).await {
|
||||
Ok(()) => true,
|
||||
Err(StoreError::NotFound { .. }) => false,
|
||||
|
|
@ -583,10 +724,6 @@ pub(crate) async fn reconcile_on_startup(
|
|||
return Err(anyhow::Error::new(err).context("releasing the Petri run's lease"));
|
||||
}
|
||||
};
|
||||
let run_dir = Storage::new(state.server_storage_dir())
|
||||
.run_scratch(&run_id)
|
||||
.root()
|
||||
.to_path_buf();
|
||||
let mode = if held {
|
||||
let request = RecoveryRequest::for_run(
|
||||
run_id,
|
||||
|
|
@ -608,27 +745,11 @@ pub(crate) async fn reconcile_on_startup(
|
|||
);
|
||||
RunExecutionMode::Resume
|
||||
}
|
||||
Recovery::Failed { reason } => {
|
||||
warn!(
|
||||
run_id = %run_id,
|
||||
petri_key = %key,
|
||||
error = %reason,
|
||||
"Petri run left in flight by the previous server cannot resume; reporting it failed"
|
||||
);
|
||||
let (_, _, record) = failed(FailureReason::WorkflowError, reason);
|
||||
run_records::lifecycle(state, run_id, record).await?;
|
||||
return Ok(());
|
||||
}
|
||||
Recovery::Failed { reason } => return Ok(Relaunch::Failed { reason }),
|
||||
}
|
||||
} else {
|
||||
RunExecutionMode::Start
|
||||
};
|
||||
info!(
|
||||
run_id = %run_id,
|
||||
petri_key = %key,
|
||||
mode = super::worker_mode_arg(mode),
|
||||
"Petri run left in flight by the previous server; relaunching its worker"
|
||||
);
|
||||
let mut start_requested = RunLifecycleRecord::new(RunLifecycleKind::StartRequested);
|
||||
start_requested.source = Some("resume".to_string());
|
||||
let mut runnable = run_records::transition(RunLifecycleKind::Runnable, RunStatus::Runnable);
|
||||
|
|
@ -636,18 +757,15 @@ pub(crate) async fn reconcile_on_startup(
|
|||
for record in [start_requested, runnable] {
|
||||
run_records::lifecycle(state, run_id, record).await?;
|
||||
}
|
||||
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
||||
runs.insert(
|
||||
run_id,
|
||||
super::managed_run(
|
||||
run_state.spec.graph_source.clone().unwrap_or_default(),
|
||||
RunStatus::Runnable,
|
||||
run_id.created_at(),
|
||||
run_dir,
|
||||
mode,
|
||||
),
|
||||
);
|
||||
Ok(())
|
||||
Ok(Relaunch::Worker(mode))
|
||||
}
|
||||
|
||||
/// The run's scratch root: its worker's run directory.
|
||||
fn scratch_root(state: &AppState, run_id: RunId) -> std::path::PathBuf {
|
||||
Storage::new(state.server_storage_dir())
|
||||
.run_scratch(&run_id)
|
||||
.root()
|
||||
.to_path_buf()
|
||||
}
|
||||
|
||||
/// The failed status, its message, and the `failed` lifecycle record for it.
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ use async_trait::async_trait;
|
|||
use fabro_static::EnvVars;
|
||||
use fabro_types::RunId;
|
||||
use fabro_types::settings::server::LogDestination;
|
||||
use fabro_util::exit::ExitClass;
|
||||
use futures_util::future::BoxFuture;
|
||||
use tokio::io::AsyncRead;
|
||||
use tokio::process::Command;
|
||||
|
|
@ -62,9 +63,19 @@ pub(crate) struct StartedWorker {
|
|||
#[derive(Debug)]
|
||||
pub(crate) struct WorkerExit {
|
||||
pub(crate) success: bool,
|
||||
/// The worker's exit code; `None` when a signal ended it.
|
||||
pub(crate) code: Option<i32>,
|
||||
pub(crate) detail: String,
|
||||
}
|
||||
|
||||
impl WorkerExit {
|
||||
/// Whether the worker ended the run's lifetime and not the run: its
|
||||
/// store failed, and the run resumes ([`ExitClass::Interrupted`]).
|
||||
pub(crate) fn interrupted(&self) -> bool {
|
||||
self.code == Some(ExitClass::Interrupted.code())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
pub(crate) struct LocalWorkerRuntime;
|
||||
|
||||
|
|
@ -137,6 +148,7 @@ impl WorkerRuntime for LocalWorkerRuntime {
|
|||
let status = child.wait().await.context("worker wait failed")?;
|
||||
Ok(WorkerExit {
|
||||
success: status.success(),
|
||||
code: status.code(),
|
||||
detail: status.to_string(),
|
||||
})
|
||||
});
|
||||
|
|
|
|||
|
|
@ -18,6 +18,13 @@
|
|||
//! `inspect_run` over a read handle of the same store, so what the caller
|
||||
//! reports is what the durable record says.
|
||||
//!
|
||||
//! A failed write to the run's store ends the run's lifetime, not the run:
|
||||
//! Petri records nothing after it and returns `CoordinatorError::StoreFailed`,
|
||||
//! and no firing fails for it. [`run`] returns [`RunError::StoreFailed`]
|
||||
//! without reading the record back, and [`conclusion`] says
|
||||
//! [`Conclusion::Interrupted`]: the caller records no end for the run, and
|
||||
//! the run resumes from its records, as after a crash.
|
||||
//!
|
||||
//! What the caller supplies beyond the runtime: the interviewer its
|
||||
//! questions go to ([`interview`](crate::interview) in the worker and the
|
||||
//! server), the secret provider over the vault ([`secrets`](crate::secrets))
|
||||
|
|
@ -59,8 +66,8 @@ use fabro_util::sync;
|
|||
use petri_execution::host::{self, HostError, HostRun};
|
||||
use petri_execution::inspect::{self, InspectError, RunInspection};
|
||||
use petri_execution::{
|
||||
Access, CancelReason, ExecutionObserver, InterviewDispatcher, Interviewer, RECEIPT_FILE,
|
||||
RunKey, RunStore,
|
||||
Access, CancelReason, CoordinatorError, ExecutionObserver, InterviewDispatcher, Interviewer,
|
||||
RECEIPT_FILE, RunKey, RunStore,
|
||||
};
|
||||
use petri_runtime::driver::ExecutionReport;
|
||||
use petri_runtime::driver::lifecycle::ExecutionHooks;
|
||||
|
|
@ -169,6 +176,11 @@ pub enum RunError {
|
|||
Inspect(#[source] InspectError),
|
||||
#[error("the run ended without recording a status; the record says: {}", .0.join("; "))]
|
||||
Unfinished(Vec<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
|
||||
|
|
@ -183,6 +195,10 @@ pub enum Conclusion {
|
|||
reason: FailureReason,
|
||||
message: String,
|
||||
},
|
||||
/// The run's store failed, which ended this lifetime of the run and not
|
||||
/// the run: the caller records no terminal event, and the run resumes
|
||||
/// from its records, as after a crash.
|
||||
Interrupted { message: String },
|
||||
}
|
||||
|
||||
fn network_policy(settings: &EnvironmentNetworkSettings) -> NetworkPolicy {
|
||||
|
|
@ -331,6 +347,11 @@ pub async fn run(request: RunRequest) -> Result<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
|
||||
|
|
@ -397,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 {
|
||||
|
|
@ -419,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),
|
||||
|
|
@ -641,4 +665,11 @@ mod tests {
|
|||
.to_string(),
|
||||
});
|
||||
}
|
||||
#[test]
|
||||
fn a_failed_store_concludes_interrupted() {
|
||||
let error = RunError::StoreFailed("could not append: the disk is full".to_string());
|
||||
assert_eq!(conclusion(&Err(error)), Conclusion::Interrupted {
|
||||
message: "the run's store failed: could not append: the disk is full".to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)?;
|
||||
|
|
|
|||
|
|
@ -5,12 +5,16 @@
|
|||
mod support;
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
use fabro_petri::admission::AdmittedGraphs;
|
||||
use fabro_petri::check::Launch;
|
||||
use fabro_petri::engine::{self, Execution, RunStatus};
|
||||
use fabro_petri::engine::{self, Conclusion, Execution, RunError, RunStatus};
|
||||
use fabro_petri::runtime::RuntimeSpec;
|
||||
use petri_store::{Access, MemoryRunStore, OwnerId, RunKey, RunStore as _};
|
||||
use support::{SETTINGS, Silent, admit, no_questions, run_request};
|
||||
use petri_store::{
|
||||
Access, Digest, LogId, MemoryRunStore, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError,
|
||||
};
|
||||
use support::{SETTINGS, Silent, admit, all_records, no_questions, run_request};
|
||||
|
||||
/// One command stage between start and exit.
|
||||
const COMMAND: &str = r#"digraph Command {
|
||||
|
|
@ -41,21 +45,160 @@ async fn a_resume_of_a_run_that_never_started_starts_it_again() {
|
|||
Launch::default(),
|
||||
&runtime,
|
||||
);
|
||||
let mut request = run_request(
|
||||
let outcome = engine::run(resume_request(
|
||||
"cut-short",
|
||||
root.path(),
|
||||
graphs,
|
||||
store,
|
||||
runtime,
|
||||
))
|
||||
.await
|
||||
.expect("the run ends");
|
||||
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
}
|
||||
|
||||
/// A write to the run's store fails partway through the run: the lifetime
|
||||
/// ends with the store's failure, which interrupts the run rather than
|
||||
/// failing it, and nothing records an end. A resume over the same records
|
||||
/// finishes the run.
|
||||
#[tokio::test]
|
||||
async fn a_failed_store_write_interrupts_the_run_and_a_resume_finishes_it() {
|
||||
let root = tempfile::tempdir().expect("a temp dir");
|
||||
let memory = Arc::new(MemoryRunStore::new());
|
||||
let failing = Arc::new(FailingStore {
|
||||
inner: Arc::clone(&memory),
|
||||
log: LogId::Resources,
|
||||
appended: Arc::new(AtomicUsize::new(0)),
|
||||
});
|
||||
let runtime = RuntimeSpec::default();
|
||||
let graphs = || {
|
||||
admit(
|
||||
&[("workflow.fabro", COMMAND), ("workflow.toml", SETTINGS)],
|
||||
Launch::default(),
|
||||
&runtime,
|
||||
)
|
||||
};
|
||||
|
||||
let first = engine::run(run_request(
|
||||
"interrupted",
|
||||
root.path(),
|
||||
graphs(),
|
||||
failing,
|
||||
runtime.clone(),
|
||||
no_questions(Arc::new(Silent)),
|
||||
))
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
matches!(&first, Err(RunError::StoreFailed(message)) if message.contains("the disk is full")),
|
||||
"{first:?}"
|
||||
);
|
||||
assert!(
|
||||
matches!(engine::conclusion(&first), Conclusion::Interrupted { .. }),
|
||||
"{first:?}"
|
||||
);
|
||||
assert!(!finished(&memory).await, "nothing ended the run");
|
||||
|
||||
let outcome = engine::run(resume_request(
|
||||
"interrupted",
|
||||
root.path(),
|
||||
graphs(),
|
||||
Arc::clone(&memory) as Arc<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)),
|
||||
);
|
||||
request.execution = match request.execution {
|
||||
Execution::Start(graphs) => Execution::Resume(graphs),
|
||||
resume @ Execution::Resume(_) => resume,
|
||||
};
|
||||
|
||||
let outcome = engine::run(request).await.expect("the run ends");
|
||||
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
request
|
||||
}
|
||||
|
||||
/// The in-memory store, failing its first append to `log` as a full disk
|
||||
/// would, before anything is stored.
|
||||
struct FailingStore {
|
||||
inner: Arc<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
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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!(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue