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(()),