Simplify the store-fault handling

- stop_lock_holder returns the stopped pid as Option<u32> instead of a
  LockHolder enum that only wrapped it
- stop_previous_worker uses with_context, and one run_scratch helper
  replaces three spellings of the run's scratch path
- relaunch builds its runnable record through run_records::runnable,
  moved out of the lifecycle handler so both callers share it
- WorkerExit derives success from its exit code instead of storing both
- the scenario tests share one worker-pid lookup and wait loop
- small readability fixes in the engine's resume arm and the resume test

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-06 16:07:36 -04:00
parent f12065e6ae
commit e889f9978d
11 changed files with 69 additions and 97 deletions

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);
}
@ -805,7 +816,7 @@ async fn a_worker_that_outlives_the_server_is_stopped_before_its_run_resumes() {
server.launch().await;
eprintln!("server restarted");
let resumed = wait_for_worker_other_than(&run_id, worker);
let resumed = wait_for_worker_except(&run_id, Some(worker));
eprintln!("worker {resumed} launched for the resume");
assert!(
!fabro_proc::process_running_strict(worker),
@ -826,30 +837,6 @@ async fn a_worker_that_outlives_the_server_is_stopped_before_its_run_resumes() {
server.shutdown();
}
/// Wait for a worker of the run other than `previous`.
fn wait_for_worker_other_than(run_id: &str, previous: u32) -> u32 {
let short_id: String = run_id.chars().take(12).collect();
let deadline = Instant::now() + RUN_TIMEOUT;
loop {
let output = Command::new("pgrep")
.args(["-f", &format!("^fabro {short_id} ")])
.output()
.expect("pgrep runs");
if let Some(pid) = String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|line| line.trim().parse::<u32>().ok())
.find(|pid| *pid != previous)
{
return pid;
}
assert!(
Instant::now() < deadline,
"no new worker process appeared for run {run_id}"
);
std::thread::sleep(POLL);
}
}
/// The run's status while the server may be down: `None` when it is.
async fn run_status_offline(server: &RunningServer) -> Option<String> {
fabro_test::test_http_client()

View file

@ -258,7 +258,6 @@ mod tests {
exit.notified().await;
let code = *sync::lock(&code);
Ok(WorkerExit {
success: code == Some(0),
code,
detail: "test worker ended without a terminal event".to_string(),
})

View file

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

@ -40,8 +40,9 @@ use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use anyhow::Context as _;
use fabro_config::{
EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, SettingsLayer, Storage,
EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, RunScratch, SettingsLayer, Storage,
};
use fabro_interview::ControlInterviewer;
use fabro_petri::artifacts::StoreArtifactWriter;
@ -605,7 +606,7 @@ pub(crate) async fn reconcile_on_startup(
run_state.spec.graph_source.clone().unwrap_or_default(),
RunStatus::Runnable,
run_id.created_at(),
scratch_root(state, run_id),
run_scratch(state, run_id).root().to_path_buf(),
mode,
),
);
@ -750,8 +751,7 @@ async fn relaunch(state: &Arc<AppState>, run_id: RunId) -> anyhow::Result<Relaun
};
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?;
}
@ -772,18 +772,11 @@ const WORKER_STOP_PATIENCE: Duration = Duration::from_secs(10);
pub(crate) async fn stop_previous_worker(state: &AppState, run_id: RunId) -> anyhow::Result<()> {
#[cfg(unix)]
{
let path = Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.worker_lock_path();
let holder = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE)
let path = run_scratch(state, run_id).worker_lock_path();
let stopped = fabro_proc::stop_lock_holder(&path, WORKER_STOP_PATIENCE)
.await
.map_err(|err| {
anyhow::Error::new(err).context(format!(
"stopping the run's previous worker ({})",
path.display()
))
})?;
if let fabro_proc::LockHolder::Stopped { pid } = holder {
.with_context(|| format!("stopping the run's previous worker ({})", path.display()))?;
if let Some(pid) = stopped {
warn!(
run_id = %run_id,
pid,
@ -796,12 +789,9 @@ pub(crate) async fn stop_previous_worker(state: &AppState, run_id: RunId) -> any
Ok(())
}
/// The run's scratch root: its worker's run directory.
fn scratch_root(state: &AppState, run_id: RunId) -> std::path::PathBuf {
Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.root()
.to_path_buf()
/// The run's scratch: its worker's run directory and worker lock.
fn run_scratch(state: &AppState, run_id: RunId) -> RunScratch {
Storage::new(state.server_storage_dir()).run_scratch(&run_id)
}
/// The failed status, its message, and the `failed` lifecycle record for it.
@ -979,10 +969,7 @@ mod tests {
async fn in_flight_run() -> (Arc<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

@ -62,13 +62,17 @@ pub(crate) struct StartedWorker {
#[derive(Debug)]
pub(crate) struct WorkerExit {
pub(crate) success: bool,
/// The worker's exit code; `None` when a signal ended it.
pub(crate) code: Option<i32>,
pub(crate) detail: String,
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 {
@ -147,9 +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(),
code: status.code(),
detail: status.to_string(),
code: status.code(),
detail: status.to_string(),
})
});

View file

@ -316,14 +316,14 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
}
Execution::Resume(graphs) => {
info!(run_id = %request.run_id, backend = %backend, "Resuming Petri run");
let resumed = Box::pin(host::resume_configured(
let outcome = Box::pin(host::resume_configured(
&runtime,
Vec::new(),
observers.clone(),
wiring(),
))
.await;
match resumed {
match outcome {
// A crash cut the run's creation short: nothing beyond its
// start is stored, so it starts again from its admitted
// graphs, and Petri takes the stored prefix over.
@ -334,7 +334,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
);
start(&runtime, graphs, observers, wiring()).await
}
resumed => resumed,
outcome => outcome,
}
}
};

View file

@ -141,10 +141,10 @@ fn resume_request(
runtime,
no_questions(Arc::new(Silent)),
);
request.execution = match request.execution {
Execution::Start(graphs) => Execution::Resume(graphs),
resume @ Execution::Resume(_) => resume,
let Execution::Start(graphs) = request.execution else {
unreachable!("run_request builds a start");
};
request.execution = Execution::Resume(graphs);
request
}

View file

@ -25,7 +25,7 @@ pub use pre_exec::pre_exec_setpgid;
#[cfg(unix)]
pub use pre_exec::pre_exec_setsid;
#[cfg(unix)]
pub use process_lock::{LockHolder, ProcessLock, stop_lock_holder};
pub use process_lock::{ProcessLock, stop_lock_holder};
pub use signal::{process_exists, process_group_alive, process_running, process_running_strict};
#[cfg(unix)]
pub use signal::{

View file

@ -69,37 +69,29 @@ impl ProcessLock {
}
}
/// What [`stop_lock_holder`] found.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LockHolder {
/// No other process held the lock.
None,
/// This process held the lock. It was killed, with its process group,
/// and the lock is free: it is gone.
Stopped { pid: u32 },
}
/// Stop the process that holds the lock on the file at `path`, if another
/// process does, and wait up to `patience` for the lock to be free. The
/// holder and its process group get `SIGKILL`: a holder that could handle
/// a signal could also keep running. A missing file has no holder.
///
/// `Err` when the lock is still held after `patience`.
pub async fn stop_lock_holder(path: &Path, patience: Duration) -> io::Result<LockHolder> {
/// `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(LockHolder::None),
Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
Err(error) => return Err(error),
};
let Some(pid) = holder(&file)? else {
return Ok(LockHolder::None);
return Ok(None);
};
signal::sigkill_process_group(pid);
signal::sigkill(pid);
let deadline = Instant::now() + patience;
loop {
match holder(&file)? {
None => return Ok(LockHolder::Stopped { pid }),
None => return Ok(Some(pid)),
Some(_) if Instant::now() >= deadline => {
return Err(io::Error::other(format!(
"process {pid} still holds {} after SIGKILL",
@ -210,7 +202,7 @@ mod tests {
stop_lock_holder(&path, Duration::from_secs(1))
.await
.expect("the check runs"),
LockHolder::None
None
);
drop(
ProcessLock::try_hold(&path)
@ -222,7 +214,7 @@ mod tests {
stop_lock_holder(&path, Duration::from_secs(1))
.await
.expect("the check runs"),
LockHolder::None
None
);
}
@ -250,7 +242,7 @@ mod tests {
.await
.expect("the holder is stopped");
assert_eq!(found, LockHolder::Stopped { pid });
assert_eq!(found, Some(pid));
let status = holder.wait().await.expect("the holder is reaped");
assert!(!status.success(), "the holder was killed: {status}");
assert!(