From e889f9978d52fa4f9bc763efc890a37a676c2953 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Tue, 6 Oct 2026 16:07:36 -0400 Subject: [PATCH] Simplify the store-fault handling - stop_lock_holder returns the stopped pid as Option 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 --- lib/apps/fabro-cli/tests/it/scenario/petri.rs | 43 +++++++------------ lib/apps/fabro-server/src/petri_runs.rs | 1 - lib/apps/fabro-server/src/server.rs | 7 ++- .../src/server/handler/lifecycle.rs | 11 +---- .../fabro-server/src/server/petri_runs.rs | 37 ++++++---------- .../fabro-server/src/server/run_records.rs | 10 ++++- lib/apps/fabro-server/src/worker_runtime.rs | 15 ++++--- lib/components/fabro-petri/src/engine.rs | 6 +-- lib/components/fabro-petri/tests/resume.rs | 6 +-- lib/foundation/fabro-proc/src/lib.rs | 2 +- lib/foundation/fabro-proc/src/process_lock.rs | 28 +++++------- 11 files changed, 69 insertions(+), 97 deletions(-) diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index 8a3482286..c641054fa 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -595,6 +595,11 @@ pub(super) async fn settled_stream(server: &RunningServer, run_id: &str) -> Vec< /// worker retitles itself `fabro `, so that /// is what the process table shows. fn worker_pid(run_id: &str) -> Option { + 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) -> Option { 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 { .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 { 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::().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 { fabro_test::test_http_client() diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index b977b1d7f..28380e8f2 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -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(), }) diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 9f0024990..2a2ae3079 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -4341,8 +4341,11 @@ async fn execute_run_subprocess(state: Arc, 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() diff --git a/lib/apps/fabro-server/src/server/handler/lifecycle.rs b/lib/apps/fabro-server/src/server/handler/lifecycle.rs index 1d53aec93..b5b535c93 100644 --- a/lib/apps/fabro-server/src/server/handler/lifecycle.rs +++ b/lib/apps/fabro-server/src/server/handler/lifecycle.rs @@ -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. diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 31f4980ca..356054d94 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, run_id: RunId) -> anyhow::Result::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, RunId, SettlingLogs, Arc) { 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( diff --git a/lib/apps/fabro-server/src/server/run_records.rs b/lib/apps/fabro-server/src/server/run_records.rs index ad46f02aa..a32c4b93c 100644 --- a/lib/apps/fabro-server/src/server/run_records.rs +++ b/lib/apps/fabro-server/src/server/run_records.rs @@ -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) -> RunLifecycleRecord { diff --git a/lib/apps/fabro-server/src/worker_runtime.rs b/lib/apps/fabro-server/src/worker_runtime.rs index 9d2029340..1171e7153 100644 --- a/lib/apps/fabro-server/src/worker_runtime.rs +++ b/lib/apps/fabro-server/src/worker_runtime.rs @@ -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, - pub(crate) detail: String, + pub(crate) code: Option, + 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> = 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(), }) }); diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 32d707583..f9f2380ed 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -316,14 +316,14 @@ pub async fn run(request: RunRequest) -> Result { } 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 { ); start(&runtime, graphs, observers, wiring()).await } - resumed => resumed, + outcome => outcome, } } }; diff --git a/lib/components/fabro-petri/tests/resume.rs b/lib/components/fabro-petri/tests/resume.rs index 9f07cb148..bf607dfed 100644 --- a/lib/components/fabro-petri/tests/resume.rs +++ b/lib/components/fabro-petri/tests/resume.rs @@ -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 } diff --git a/lib/foundation/fabro-proc/src/lib.rs b/lib/foundation/fabro-proc/src/lib.rs index ff3103f5f..6ca06b1af 100644 --- a/lib/foundation/fabro-proc/src/lib.rs +++ b/lib/foundation/fabro-proc/src/lib.rs @@ -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::{ diff --git a/lib/foundation/fabro-proc/src/process_lock.rs b/lib/foundation/fabro-proc/src/process_lock.rs index 3b1873dc6..90c7f7df6 100644 --- a/lib/foundation/fabro-proc/src/process_lock.rs +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -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 { +/// `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> { 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!(