Stop a surviving worker before ending its lease from outside

A worker leads a process group of its own, so it outlives a server
crash. The restarted server released the run's lease from outside and
launched a resume while that worker could still be running: a worker
whose lease is released keeps acting until its next write, beside its
successor. A delete likewise dropped the lease without knowing the
worker was gone.

Now each worker holds a lock on `worker.lock` in its run's scratch
directory for its whole life, taken before anything else. It is a POSIX
record lock: the kernel frees it only when the worker exits, and names
the process that holds it. Before the server ends a lease from outside
(the relaunch after a restart or a store interruption, and a delete), it
kills whatever process still holds the lock, with its process group, and
waits until the lock is free. A worker that finds the lock held does not
start.

- fabro-proc: ProcessLock::try_hold and stop_lock_holder, tested with
  this test binary as the holding process.
- fabro-config: RunScratch::worker_lock_path.
- New scenario: a worker that outlives the server is gone before the
  resume's worker launches, and the run succeeds once. It fails without
  the server-side stop.

A host stage process runs in a process group of its own and still
outlives its killed worker, as it did at a worker crash.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-28 10:33:26 -04:00 • committed by Scott Werner
parent 37116758fa
commit cb54a0df90
8 changed files with 436 additions and 10 deletions

View file

@ -5,6 +5,8 @@ use std::time::Duration;
use anyhow::{Context, Result, anyhow};
use fabro_client::ServerTarget;
#[cfg(unix)]
use fabro_config::RunScratch;
use fabro_config::Storage;
use fabro_interview::{
AnswerSubmission, ControlInterviewer, WORKER_CONTROL_INVALID_CURSOR_REASON,
@ -22,6 +24,8 @@ use futures::{SinkExt, StreamExt};
use jsonwebtoken::dangerous::insecure_decode;
#[cfg(unix)]
use nix::unistd;
#[cfg(unix)]
use tokio::fs;
#[cfg(test)]
use tokio::io::DuplexStream;
use tokio::net::TcpStream;
@ -65,6 +69,10 @@ pub(crate) async fn execute(
) -> Result<()> {
let _ = fabro_proc::title_init();
set_worker_title(&run_id, initial_worker_title_phase(mode));
// Held until this process exits: while it is, the server ends no lease
// of this run from outside.
#[cfg(unix)]
let _running = hold_worker_lock(&run_dir, &run_id).await?;
let target = server.parse::<ServerTarget>()?;
let client = server_client::connect_server_target_with_bearer(&target, worker_token).await?;
@ -93,6 +101,29 @@ pub(crate) async fn execute(
.await
}
/// Take the run's worker lock for this process's whole life, before
/// anything else. The server ends a worker's lease from outside only once
/// the lock is free, which the kernel makes it only when the process that
/// held it is gone: a worker that outlived a server crash is stopped
/// before its run resumes. Another process holding the lock is a worker of
/// the same run still running, and this one does not start beside it.
#[cfg(unix)]
async fn hold_worker_lock(run_dir: &Path, run_id: &RunId) -> Result<fabro_proc::ProcessLock> {
fs::create_dir_all(run_dir)
.await
.with_context(|| format!("creating the run directory {}", run_dir.display()))?;
let path = RunScratch::new(run_dir).worker_lock_path();
fabro_proc::ProcessLock::try_hold(&path)
.await
.with_context(|| format!("taking the worker lock {}", path.display()))?
.ok_or_else(|| {
anyhow!(
"another worker of run {run_id} is still running: it holds {}",
path.display()
)
})
}
const WORKER_TOKEN_SCOPE: &str = "run:worker";
const WORKER_RUN_TOOLS_SCOPE: &str = "agent:run_tools";

View file

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

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

View file

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

View file

@ -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<Relaunch> {
stop_previous_worker(state, run_id).await?;
let held = match state.petri_runs.release_for_restart(run_id).await {
Ok(()) => true,
Err(StoreError::NotFound { .. }) => false,
@ -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())

View file

@ -143,6 +143,13 @@ impl RunScratch {
self.root.join("runtime")
}
/// The lock the run's worker holds for its whole life: while it is
/// held, the worker is running.
#[must_use]
pub fn worker_lock_path(&self) -> PathBuf {
self.root.join("worker.lock")
}
pub fn create(&self) -> std::io::Result<()> {
std::fs::create_dir_all(self.worktree_dir())?;
std::fs::create_dir_all(self.runtime_dir())?;

View file

@ -11,6 +11,8 @@
mod flock;
#[cfg(unix)]
mod pre_exec;
#[cfg(unix)]
mod process_lock;
mod signal;
mod title;
@ -22,6 +24,8 @@ pub use pre_exec::pre_exec_pdeathsig;
pub use pre_exec::pre_exec_setpgid;
#[cfg(unix)]
pub use pre_exec::pre_exec_setsid;
#[cfg(unix)]
pub use process_lock::{LockHolder, ProcessLock, stop_lock_holder};
pub use signal::{process_exists, process_group_alive, process_running, process_running_strict};
#[cfg(unix)]
pub use signal::{

View file

@ -0,0 +1,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<Option<Self>> {
let file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(path)
.await?
.into_std()
.await;
let mut lock = whole_file(consts::WRITE_LOCK);
// SAFETY: fcntl(F_SETLK) on a valid descriptor with a valid flock
// struct; F_SETLK does not wait.
if unsafe { libc::fcntl(file.as_raw_fd(), libc::F_SETLK, &raw mut lock) } == 0 {
return Ok(Some(Self { _file: file }));
}
let error = io::Error::last_os_error();
match error.raw_os_error() {
Some(libc::EACCES | libc::EAGAIN) => Ok(None),
_ => Err(error),
}
}
}
/// 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> {
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<Option<u32>> {
let mut lock = whole_file(consts::WRITE_LOCK);
// SAFETY: fcntl(F_GETLK) on a valid descriptor with a valid flock struct
// only reads the lock table.
if unsafe { libc::fcntl(file.as_raw_fd(), libc::F_GETLK, &raw mut lock) } != 0 {
return Err(io::Error::last_os_error());
}
if lock.l_type == consts::UNLOCKED {
return Ok(None);
}
u32::try_from(lock.l_pid).map(Some).map_err(|_| {
io::Error::other(format!(
"the lock on the file is held, by no process id ({})",
lock.l_pid
))
})
}
/// A lock request covering the whole file, however long it grows.
fn whole_file(kind: libc::c_short) -> libc::flock {
// SAFETY: flock is a plain C struct, for which all zeros is valid.
let mut lock: libc::flock = unsafe { std::mem::zeroed() };
lock.l_type = kind;
lock.l_whence = consts::FROM_START;
lock.l_start = 0;
lock.l_len = 0;
lock
}
#[cfg(test)]
#[expect(
clippy::disallowed_methods,
clippy::print_stdout,
reason = "the tests start this test binary as the process that holds a lock, and it says so \
on its stdout"
)]
mod tests {
use std::process::Stdio;
use tokio::io::{AsyncBufReadExt as _, BufReader};
use tokio::process::{Child, Command};
use super::*;
/// Set for the test binary started as a lock's holder: the lock's path.
const HOLD: &str = "FABRO_PROC_TEST_HOLD_LOCK";
/// Not a test: the process the other tests start to hold a lock. It
/// holds the lock at `$FABRO_PROC_TEST_HOLD_LOCK`, says so, and waits
/// to be killed. Without the variable it does nothing.
#[tokio::test]
async fn hold_the_lock_until_killed() {
let Some(path) = std::env::var_os(HOLD) else {
return;
};
let lock = ProcessLock::try_hold(Path::new(&path))
.await
.expect("the lock file opens")
.expect("the lock is free");
println!("held");
time::sleep(Duration::from_mins(1)).await;
drop(lock);
}
/// This test binary, holding the lock at `path` in a process group of
/// its own, once it says it holds it.
async fn holder_process(path: &Path) -> Child {
let mut command = Command::new(std::env::current_exe().expect("the test binary"));
command
.args([
"--exact",
"process_lock::tests::hold_the_lock_until_killed",
"--nocapture",
])
.env(HOLD, path)
.stdout(Stdio::piped())
.stderr(Stdio::null())
.kill_on_drop(true);
crate::pre_exec_setpgid(command.as_std_mut());
let mut child = command.spawn().expect("the holder starts");
let stdout = child.stdout.take().expect("the holder's stdout");
let mut lines = BufReader::new(stdout).lines();
while let Some(line) = lines.next_line().await.expect("the holder's stdout reads") {
if line.trim() == "held" {
return child;
}
}
panic!("the holder exited before it held the lock");
}
#[tokio::test]
async fn a_missing_or_free_lock_has_no_holder() {
let dir = tempfile::tempdir().expect("a temp dir");
let path = dir.path().join("worker.lock");
assert_eq!(
stop_lock_holder(&path, Duration::from_secs(1))
.await
.expect("the check runs"),
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"
);
}
}