mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
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.
This commit is contained in:
parent
6f5da61497
commit
5899850541
2 changed files with 96 additions and 1 deletions
|
|
@ -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!();
|
||||
|
|
|
|||
|
|
@ -1710,6 +1710,7 @@ async fn delete_run_internal(state: &Arc<AppState>, 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<AppState>, id: RunId) -> Result<(), Res
|
|||
Ok(())
|
||||
}
|
||||
|
||||
async fn terminate_worker_for_deletion(worker_pid: Option<u32>, worker_pgid: Option<u32>) {
|
||||
#[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(()),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue