From 5899850541156ba9cec80684aa19c90329bf0374 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Tue, 7 Apr 2026 16:04:26 -0400 Subject: [PATCH] fix(run): terminate active workers on force removal Active runs deleted through rm --force were removed from server state without signalling the worker process, which could leave detached workers orphaned after test cleanup. Terminate the tracked worker process group before deleting run state and cover it with an integration regression. --- lib/crates/fabro-cli/tests/it/cmd/rm.rs | 52 ++++++++++++++++++++++++- lib/crates/fabro-server/src/server.rs | 45 +++++++++++++++++++++ 2 files changed, 96 insertions(+), 1 deletion(-) diff --git a/lib/crates/fabro-cli/tests/it/cmd/rm.rs b/lib/crates/fabro-cli/tests/it/cmd/rm.rs index f1ce74938..1ce9f1153 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/rm.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/rm.rs @@ -5,7 +5,8 @@ use serde_json::Value; use crate::support::unique_run_id; use super::support::{ - setup_completed_fast_dry_run, setup_created_fast_dry_run, setup_local_sandbox_run, + output_stdout, resolve_run, setup_completed_fast_dry_run, setup_created_fast_dry_run, + setup_local_sandbox_run, wait_for_no_process_match, wait_for_status, write_gated_workflow, }; #[test] @@ -149,6 +150,55 @@ fn rm_force_deletes_run_without_sandbox_json_when_store_has_sandbox() { ); } +#[test] +fn rm_force_terminates_active_run_worker() { + let context = test_context!(); + let _gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly"); + + let output = context + .run_cmd() + .env("OPENAI_API_KEY", "test") + .args([ + "--detach", + "--provider", + "openai", + "--sandbox", + "local", + "--no-retro", + "slow.fabro", + ]) + .output() + .expect("run --detach should execute"); + assert!( + output.status.success(), + "run --detach failed:\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + + let run_id = output_stdout(&output).trim().to_string(); + let run = resolve_run(&context, &run_id); + wait_for_status(&run.run_dir, &["running"]); + + let mut filters = context.filters(); + filters.push(( + r"\b[0-9A-HJKMNP-TV-Z]{12}\b".to_string(), + "[ULID]".to_string(), + )); + let mut cmd = context.command(); + cmd.args(["rm", "--force", &run_id]); + fabro_snapshot!(filters, cmd, @" + success: true + exit_code: 0 + ----- stdout ----- + ----- stderr ----- + [ULID] + "); + + assert!(!run.run_dir.exists(), "run directory should be deleted"); + wait_for_no_process_match(&format!("fabro {} ", &run_id[..12])); +} + #[test] fn rm_partial_failure_reports_which_identifiers_failed() { let context = test_context!(); diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 39660b5df..8f69e84b9 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -1710,6 +1710,7 @@ async fn delete_run_internal(state: &Arc, id: RunId) -> Result<(), Res if let Some(cancel_tx) = managed_run.cancel_tx.take() { let _ = cancel_tx.send(()); } + terminate_worker_for_deletion(managed_run.worker_pid, managed_run.worker_pgid).await; if let Some(run_dir) = managed_run.run_dir.take() { remove_run_dir(&run_dir).map_err(|err| { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() @@ -1736,6 +1737,50 @@ async fn delete_run_internal(state: &Arc, id: RunId) -> Result<(), Res Ok(()) } +async fn terminate_worker_for_deletion(worker_pid: Option, worker_pgid: Option) { + #[cfg(unix)] + if let Some(process_group_id) = worker_pgid.or(worker_pid) { + fabro_proc::sigterm_process_group(process_group_id); + + let deadline = Instant::now() + WORKER_CANCEL_GRACE; + while Instant::now() < deadline && fabro_proc::process_group_alive(process_group_id) { + sleep(Duration::from_millis(50)).await; + } + + if fabro_proc::process_group_alive(process_group_id) { + fabro_proc::sigkill_process_group(process_group_id); + + let kill_deadline = Instant::now() + Duration::from_secs(1); + while Instant::now() < kill_deadline + && fabro_proc::process_group_alive(process_group_id) + { + sleep(Duration::from_millis(50)).await; + } + } + + return; + } + + #[cfg(not(unix))] + if let Some(worker_pid) = worker_pid { + fabro_proc::sigterm(worker_pid); + + let deadline = Instant::now() + WORKER_CANCEL_GRACE; + while Instant::now() < deadline && fabro_proc::process_alive(worker_pid) { + sleep(Duration::from_millis(50)).await; + } + + if fabro_proc::process_alive(worker_pid) { + fabro_proc::sigkill(worker_pid); + + let kill_deadline = Instant::now() + Duration::from_secs(1); + while Instant::now() < kill_deadline && fabro_proc::process_alive(worker_pid) { + sleep(Duration::from_millis(50)).await; + } + } + } +} + fn remove_run_dir(run_dir: &std::path::Path) -> std::io::Result<()> { match std::fs::remove_dir_all(run_dir) { Ok(()) => Ok(()),