diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index e8120b462..3aceaa24a 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -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::()?; 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 { + 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"; diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index 7beb2476c..8a3482286 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -779,6 +779,77 @@ 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_other_than(&run_id, 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(); +} + +/// 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 c008f69d5..b977b1d7f 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -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 diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 5c3933cab..a21ce3692 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -2788,8 +2788,12 @@ 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| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, format!("{err:#}")))?; state.petri_runs.worker_exited(id); let delete_outcome = delete_run_sandbox_resource(state, id, force).await?; diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 0be8d2675..2d3201c35 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -38,7 +38,7 @@ use std::collections::HashMap; use std::sync::Arc; -use std::time::Instant; +use std::time::{Duration, Instant}; use fabro_config::{ EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, SettingsLayer, Storage, @@ -700,8 +700,9 @@ enum Relaunch { /// 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 +/// 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 @@ -717,6 +718,7 @@ async fn relaunch( run_id: RunId, run_state: &fabro_store::RunProjection, ) -> anyhow::Result { + 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, @@ -760,6 +762,44 @@ async fn relaunch( 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 = Storage::new(state.server_storage_dir()) + .run_scratch(&run_id) + .worker_lock_path(); + let holder = 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 { + 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 root: its worker's run directory. fn scratch_root(state: &AppState, run_id: RunId) -> std::path::PathBuf { Storage::new(state.server_storage_dir()) diff --git a/lib/foundation/fabro-config/src/storage.rs b/lib/foundation/fabro-config/src/storage.rs index e7a249798..8a76c9d73 100644 --- a/lib/foundation/fabro-config/src/storage.rs +++ b/lib/foundation/fabro-config/src/storage.rs @@ -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())?; diff --git a/lib/foundation/fabro-proc/src/lib.rs b/lib/foundation/fabro-proc/src/lib.rs index 7dcc9fe97..ff3103f5f 100644 --- a/lib/foundation/fabro-proc/src/lib.rs +++ b/lib/foundation/fabro-proc/src/lib.rs @@ -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::{LockHolder, 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 new file mode 100644 index 000000000..ceaadc8d0 --- /dev/null +++ b/lib/foundation/fabro-proc/src/process_lock.rs @@ -0,0 +1,264 @@ +//! 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::{sigkill, sigkill_process_group}; + +/// 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> { + 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), + } + } +} + +/// 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 { + 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) => return Err(error), + }; + let Some(pid) = holder(&file)? else { + return Ok(LockHolder::None); + }; + sigkill_process_group(pid); + sigkill(pid); + let deadline = Instant::now() + patience; + loop { + match holder(&file)? { + None => return Ok(LockHolder::Stopped { 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> { + 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"), + LockHolder::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"), + LockHolder::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, LockHolder::Stopped { 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" + ); + } +}