#![cfg(feature = "heavy-tests")]
use std::time::Duration;
use tokio::process::Command;
#[cfg(unix)]
#[tokio::test]
async fn test_unix_process_group_cleanup() {
use std::process::Stdio;
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("sleep 30 & sleep 30 & wait")
.stdout(Stdio::null())
.stderr(Stdio::null())
.stdin(Stdio::null());
unsafe {
cmd.pre_exec(|| {
use nix::unistd::{setpgid, Pid};
setpgid(Pid::from_raw(0), Pid::from_raw(0)).map_err(std::io::Error::other)?;
Ok(())
});
}
let mut child = cmd.spawn().expect("Failed to spawn test process");
let pid = child.id().expect("Failed to get PID");
let pid_string = pid.to_string();
tokio::time::sleep(Duration::from_millis(500)).await;
let check_output = std::process::Command::new("ps")
.args(["-o", "pid,pgid", "-p", &pid_string])
.output()
.expect("Failed to check process group");
assert!(
check_output.status.success(),
"Process should be running before termination"
);
use nix::sys::signal::{killpg, Signal};
use nix::unistd::Pid;
killpg(Pid::from_raw(pid as i32), Signal::SIGTERM).expect("Failed to kill process group");
let _ = child.wait().await;
let mut terminated = false;
for _ in 0..10 {
tokio::time::sleep(Duration::from_millis(200)).await;
let check_output = std::process::Command::new("ps")
.args(["-p", &pid_string])
.output()
.expect("Failed to check process");
if !check_output.status.success() {
terminated = true;
break;
}
}
assert!(terminated, "Process should be terminated after killpg");
}
#[cfg(windows)]
#[tokio::test]
async fn test_windows_job_object_cleanup() {
println!("Windows job object test not yet implemented");
}
#[cfg(unix)]
mod apply_completion_cleanup {
use conflux::process_manager::{
cleanup_process_group_verified, configure_process_group, ManagedChild,
ProcessGroupQuiescence,
};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::time::Duration;
use tempfile::TempDir;
use tokio::process::Command;
fn git(workspace: &Path, args: &[&str]) -> std::process::Output {
std::process::Command::new("git")
.args(args)
.current_dir(workspace)
.output()
.expect("git should run")
}
fn init_worktree(workspace: &Path) {
git(workspace, &["init"]);
git(workspace, &["config", "user.email", "test@example.com"]);
git(workspace, &["config", "user.name", "Test User"]);
std::fs::write(workspace.join("README.md"), "initial\n").unwrap();
git(workspace, &["add", "README.md"]);
let out = git(workspace, &["commit", "-m", "initial"]);
assert!(out.status.success(), "initial commit must succeed");
}
fn write_script(dir: &Path, name: &str, body: &str) -> PathBuf {
let path = dir.join(name);
std::fs::write(&path, body).expect("script should be written");
path
}
fn group_has_members(pgid: u32) -> bool {
unsafe { libc::killpg(pgid as i32, 0) == 0 }
}
async fn spawn_leader(dir: &Path, descendant_script: &Path) -> (ManagedChild, u32) {
let leader_script = write_script(
dir,
"leader.sh",
&format!(
"#!/bin/sh\n\
sh {descendant} >/dev/null 2>&1 </dev/null &\n\
sleep 120\n",
descendant = descendant_script.display()
),
);
let mut cmd = Command::new("sh");
cmd.arg(leader_script.as_os_str())
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("leader should spawn");
let managed = ManagedChild::new(child).expect("managed leader");
let pgid = managed.id().expect("leader pid");
(managed, pgid)
}
async fn wait_for(path: &Path) {
for _ in 0..100 {
if path.exists() {
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
panic!("expected {} to appear", path.display());
}
#[tokio::test]
async fn apply_completion_cleanup_waits_for_descendant_holding_index_lock() {
let temp = TempDir::new().unwrap();
let workspace = temp.path().join("worktree");
std::fs::create_dir_all(&workspace).unwrap();
init_worktree(&workspace);
let index_lock = workspace.join(".git").join("index.lock");
let released = temp.path().join("released");
let holder = write_script(
temp.path(),
"holder.sh",
&format!(
"#!/bin/sh\n\
release() {{ rm -f {lock}; echo released > {released}; exit 0; }}\n\
trap release TERM\n\
: > {lock}\n\
while :; do sleep 0.2; done\n",
lock = index_lock.display(),
released = released.display()
),
);
let (mut leader, pgid) = spawn_leader(temp.path(), &holder).await;
wait_for(&index_lock).await;
let leader_pid = leader.id().expect("leader pid") as i32;
unsafe { libc::kill(leader_pid, libc::SIGTERM) };
let _ = tokio::time::timeout(Duration::from_secs(5), leader.wait())
.await
.expect("leader should exit");
assert!(
index_lock.exists(),
"the descendant must still hold index.lock after leader exit"
);
assert!(
group_has_members(pgid),
"leader exit must not be mistaken for process-group quiescence"
);
let blocked = git(&workspace, &["commit", "--allow-empty", "-m", "would race"]);
assert!(
!blocked.status.success(),
"Git finalization must be impossible while the descendant holds index.lock"
);
let report = cleanup_process_group_verified(pgid, 2_000, 10_000, Some("apply"), None).await;
assert_eq!(
report.quiescence(),
ProcessGroupQuiescence::Confirmed,
"cleanup must confirm quiescence: {}",
report.diagnostics()
);
assert!(
released.exists() && !index_lock.exists(),
"the descendant must have released the lock itself before quiescence was confirmed"
);
assert!(
!group_has_members(pgid),
"confirmed quiescence must mean no owned members remain"
);
let finalized = git(
&workspace,
&["commit", "--allow-empty", "-m", "Apply: change"],
);
assert!(
finalized.status.success(),
"Git finalization must succeed after confirmed cleanup: {}",
String::from_utf8_lossy(&finalized.stderr)
);
}
#[tokio::test]
async fn apply_completion_cleanup_reports_unconfirmed_when_budget_exhausted() {
let temp = TempDir::new().unwrap();
let workspace = temp.path().join("worktree");
std::fs::create_dir_all(&workspace).unwrap();
init_worktree(&workspace);
let index_lock = workspace.join(".git").join("index.lock");
let holder = write_script(
temp.path(),
"holder.sh",
&format!(
"#!/bin/sh\n\
trap '' TERM\n\
: > {lock}\n\
while :; do sleep 0.2; done\n",
lock = index_lock.display()
),
);
let (mut leader, pgid) = spawn_leader(temp.path(), &holder).await;
wait_for(&index_lock).await;
let leader_pid = leader.id().expect("leader pid") as i32;
unsafe { libc::kill(leader_pid, libc::SIGTERM) };
let _ = tokio::time::timeout(Duration::from_secs(5), leader.wait())
.await
.expect("leader should exit");
let report = cleanup_process_group_verified(pgid, 2_000, 0, Some("apply"), None).await;
let members_after = group_has_members(pgid);
let lock_after = index_lock.exists();
unsafe { libc::killpg(pgid as i32, libc::SIGKILL) };
assert!(
!report.is_confirmed(),
"an unswept group must never be reported as quiescent: {}",
report.diagnostics()
);
assert_eq!(report.quiescence(), ProcessGroupQuiescence::MembersRemain);
assert!(
report.diagnostics().contains("cleanup budget expired"),
"diagnostics must be actionable: {}",
report.diagnostics()
);
assert!(members_after, "the survivor must still be observable");
assert!(
lock_after,
"an unconfirmed cleanup must leave the worktree untouched, including its lock file"
);
}
#[tokio::test]
async fn apply_completion_cleanup_confirms_after_forced_termination() {
let temp = TempDir::new().unwrap();
let started = temp.path().join("started");
let holder = write_script(
temp.path(),
"holder.sh",
&format!(
"#!/bin/sh\n\
trap '' TERM\n\
: > {started}\n\
while :; do sleep 0.2; done\n",
started = started.display()
),
);
let (mut leader, pgid) = spawn_leader(temp.path(), &holder).await;
wait_for(&started).await;
let leader_pid = leader.id().expect("leader pid") as i32;
unsafe { libc::kill(leader_pid, libc::SIGTERM) };
let _ = tokio::time::timeout(Duration::from_secs(5), leader.wait())
.await
.expect("leader should exit");
let report = cleanup_process_group_verified(pgid, 200, 10_000, Some("apply"), None).await;
assert_eq!(
report.quiescence(),
ProcessGroupQuiescence::Confirmed,
"force-kill escalation must reach quiescence: {}",
report.diagnostics()
);
assert!(
report.force_killed(),
"a SIGTERM-immune descendant must be recorded as force killed"
);
assert!(!group_has_members(pgid));
}
}
#[cfg(unix)]
mod absolute_runtime_limit {
use conflux::ai_command_runner::{AiCommandRunner, OutputLine, RunCommandScope};
use conflux::config::OrchestratorConfig;
use conflux::process_manager::CommandTermination;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Mutex;
const CHATTY_FOREVER: &str = "while :; do echo tick; sleep 0.05; done";
const CHATTY_FOREVER_SIGTERM_IMMUNE: &str =
"sh -c 'trap \"\" TERM; while :; do sleep 0.2; done' >/dev/null 2>&1 </dev/null & \
while :; do echo tick; sleep 0.05; done";
const SAFETY: Duration = Duration::from_secs(60);
fn group_has_members(pgid: i32) -> bool {
unsafe { libc::killpg(pgid, 0) == 0 }
}
fn reap_and_report_survival(pgid: i32) -> bool {
let survived = group_has_members(pgid);
if survived {
unsafe { libc::killpg(pgid, libc::SIGKILL) };
}
survived
}
fn runner_with_limit(scope: &RunCommandScope, max_runtime_secs: u64) -> AiCommandRunner {
runner_with_limits(scope, max_runtime_secs, None)
}
fn runner_with_limits(
scope: &RunCommandScope,
max_runtime_secs: u64,
acceptance_max_runtime_secs: Option<u64>,
) -> AiCommandRunner {
let config = OrchestratorConfig {
command_queue_stagger_delay_ms: Some(0),
command_queue_max_retries: Some(1),
command_queue_retry_delay_ms: Some(0),
command_inactivity_timeout_secs: Some(0),
command_inactivity_timeout_max_retries: Some(0),
command_max_runtime_secs: Some(max_runtime_secs),
acceptance_max_runtime_secs,
..OrchestratorConfig::default()
};
AiCommandRunner::for_run(&config, Arc::new(Mutex::new(None)), scope.clone())
}
async fn wait_for_pgid(handle: &conflux::process_manager::StreamingChildHandle) -> i32 {
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
if let Some(pid) = handle.id() {
return pid as i32;
}
assert!(
std::time::Instant::now() < deadline,
"the command never reported a real pid"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
#[tokio::test]
async fn continuous_output_does_not_extend_the_absolute_deadline() {
let scope = RunCommandScope::new();
let runner = runner_with_limit(&scope, 1);
let (mut handle, mut rx) = runner
.execute_streaming_with_retry(CHATTY_FOREVER, None, Some("apply"), Some("change-a"))
.await
.expect("an open scope admits the command");
let pgid = wait_for_pgid(&handle).await;
let mut lines = 0usize;
let drain = tokio::spawn(async move {
while let Some(line) = rx.recv().await {
if matches!(line, OutputLine::Stdout(_)) {
lines += 1;
}
}
lines
});
let status = tokio::time::timeout(SAFETY, handle.wait())
.await
.expect("the absolute limit must end a command that never stops emitting output")
.expect("the runner publishes a final status");
let termination = handle.termination().await;
let cleanup = handle.process_group_cleanup().await;
let emitted = drain.await.expect("the drain task joins");
let survived = reap_and_report_survival(pgid);
assert!(
emitted > 0,
"arrangement failed: the command must have been emitting output"
);
assert!(
!status.success(),
"a command stopped by its runtime limit is not a success"
);
assert_eq!(
termination,
CommandTermination::RuntimeLimit,
"the reason must be the runtime limit, not an ordinary exit"
);
assert!(
!termination.permits_retry(),
"runtime-limit termination closes retry admission for the invocation"
);
assert!(
cleanup.is_confirmed(),
"the owned group must be proven quiescent: {}",
cleanup.diagnostics()
);
assert!(!survived, "process group {pgid} survived its runtime limit");
}
#[tokio::test]
async fn acceptance_stays_bounded_when_the_common_limit_is_disabled() {
const ACCEPTANCE_OPERATION_TYPE: &str = "acceptance";
let scope = RunCommandScope::new();
let runner = runner_with_limits(&scope, 0, Some(1));
assert_eq!(runner.queue_config().max_runtime_secs, 0);
assert_eq!(
runner.queue_config().effective_max_runtime_secs(None),
0,
"every other class is genuinely unbounded here"
);
assert_eq!(
runner
.queue_config()
.effective_max_runtime_secs(Some(ACCEPTANCE_OPERATION_TYPE)),
1,
"a disabled common limit leaves the dedicated Acceptance limit binding"
);
let (mut handle, mut rx) = runner
.execute_streaming_with_retry(
CHATTY_FOREVER,
None,
Some(ACCEPTANCE_OPERATION_TYPE),
Some("change-a"),
)
.await
.expect("an open scope admits the command");
let pgid = wait_for_pgid(&handle).await;
let mut lines = 0usize;
let drain = tokio::spawn(async move {
while let Some(line) = rx.recv().await {
if matches!(line, OutputLine::Stdout(_)) {
lines += 1;
}
}
lines
});
let status = tokio::time::timeout(SAFETY, handle.wait())
.await
.expect("the dedicated Acceptance limit must end a command that never stops emitting")
.expect("the runner publishes a final status");
let termination = handle.termination().await;
let cleanup = handle.process_group_cleanup().await;
let emitted = drain.await.expect("the drain task joins");
let survived = reap_and_report_survival(pgid);
assert!(
emitted > 0,
"arrangement failed: the command must have been emitting output"
);
assert!(!status.success());
assert_eq!(
termination,
CommandTermination::RuntimeLimit,
"Acceptance ends on its own absolute limit, not on an ordinary exit"
);
assert!(
!termination.permits_retry(),
"retry admission closes for a timed-out Acceptance invocation"
);
assert!(
cleanup.is_confirmed(),
"the owned group must be proven quiescent: {}",
cleanup.diagnostics()
);
assert!(
!survived,
"process group {pgid} survived the Acceptance runtime limit"
);
assert_eq!(
runner.queue_config().max_runtime_secs,
0,
"the dedicated Acceptance limit is applied at the Acceptance call site only"
);
}
#[tokio::test]
async fn zero_disables_the_absolute_deadline() {
let scope = RunCommandScope::new();
let runner = runner_with_limit(&scope, 0);
let (mut handle, mut rx) = runner
.execute_streaming_with_retry(CHATTY_FOREVER, None, Some("apply"), Some("change-a"))
.await
.expect("an open scope admits the command");
let pgid = wait_for_pgid(&handle).await;
let drain = tokio::spawn(async move { while rx.recv().await.is_some() {} });
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
group_has_members(pgid),
"a disabled deadline must not terminate the command"
);
assert!(
tokio::time::timeout(Duration::from_millis(100), handle.wait())
.await
.is_err(),
"a disabled deadline must leave the invocation running"
);
handle.terminate().expect("cancellation is admissible");
let _ = tokio::time::timeout(SAFETY, handle.wait())
.await
.expect("cancellation must end the command");
let termination = handle.termination().await;
drain.await.expect("the drain task joins");
let survived = reap_and_report_survival(pgid);
assert_eq!(
termination,
CommandTermination::Cancelled,
"with the deadline disabled the command ends by cancellation, not by runtime limit"
);
assert!(!survived, "process group {pgid} survived cancellation");
}
#[tokio::test]
async fn runtime_limit_without_provable_cleanup_reports_diagnostics() {
let scope = RunCommandScope::new();
let mut runner = runner_with_limit(&scope, 1);
runner.set_process_group_cleanup_timeout_ms(0);
let (mut handle, mut rx) = runner
.execute_streaming_with_retry(
CHATTY_FOREVER_SIGTERM_IMMUNE,
None,
Some("apply"),
Some("change-a"),
)
.await
.expect("an open scope admits the command");
let pgid = wait_for_pgid(&handle).await;
let drain = tokio::spawn(async move { while rx.recv().await.is_some() {} });
let _ = tokio::time::timeout(SAFETY, handle.wait())
.await
.expect("the absolute limit must still end the invocation");
let termination = handle.termination().await;
let cleanup = handle.process_group_cleanup().await;
drain.await.expect("the drain task joins");
reap_and_report_survival(pgid);
assert_eq!(
termination,
CommandTermination::RuntimeLimit,
"the reason stays the runtime limit even when cleanup is unprovable"
);
assert!(
!termination.permits_retry(),
"no retry may be admitted for an invocation stopped by its runtime limit"
);
assert!(
!cleanup.is_confirmed(),
"an unswept group must never be acknowledged as terminated: {}",
cleanup.diagnostics()
);
assert!(
!cleanup.diagnostics().is_empty(),
"unconfirmed cleanup must carry actionable diagnostics"
);
}
}
#[cfg(unix)]
#[tokio::test]
async fn test_process_group_isolation() {
use std::process::Stdio;
use tokio::process::Command;
let parent_pgid = unsafe { libc::getpgid(0) };
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("echo $PPID")
.stdout(Stdio::piped())
.stderr(Stdio::null());
unsafe {
cmd.pre_exec(|| {
use nix::unistd::{setpgid, Pid};
setpgid(Pid::from_raw(0), Pid::from_raw(0)).map_err(std::io::Error::other)?;
Ok(())
});
}
let child = cmd.spawn().expect("Failed to spawn test process");
let child_pid = child.id().expect("Failed to get child PID");
let child_pgid = unsafe { libc::getpgid(child_pid as i32) };
assert_ne!(
parent_pgid, child_pgid,
"Child should be in a different process group"
);
let _ = child.wait_with_output().await;
}
#[cfg(unix)]
mod run_scope_process_cleanup {
use conflux::ai_command_runner::{AiCommandRunner, OutputLine, RunCommandScope};
use conflux::config::OrchestratorConfig;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
const SIGTERM_IMMUNE_GROUP: &str =
"sh -c 'trap \"\" TERM; while :; do sleep 0.2; done' >/dev/null 2>&1 </dev/null & sleep 300";
fn scoped_runner(scope: &RunCommandScope) -> AiCommandRunner {
let config = OrchestratorConfig {
command_queue_stagger_delay_ms: Some(0),
command_queue_max_retries: Some(1),
command_queue_retry_delay_ms: Some(0),
command_inactivity_timeout_secs: Some(0),
command_inactivity_timeout_max_retries: Some(0),
..OrchestratorConfig::default()
};
AiCommandRunner::for_run(&config, Arc::new(Mutex::new(None)), scope.clone())
}
fn group_has_members(pgid: i32) -> bool {
unsafe { libc::killpg(pgid, 0) == 0 }
}
fn reap_and_report_survival(pgid: i32) -> bool {
let survived = group_has_members(pgid);
if survived {
unsafe { libc::killpg(pgid, libc::SIGKILL) };
}
survived
}
async fn launch_owned_group(
runner: &AiCommandRunner,
change_id: &str,
) -> (
conflux::process_manager::StreamingChildHandle,
tokio::sync::mpsc::Receiver<OutputLine>,
i32,
) {
let (handle, rx) = runner
.execute_streaming_with_retry(
SIGTERM_IMMUNE_GROUP,
None,
Some("apply"),
Some(change_id),
)
.await
.expect("an open scope admits the command");
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
if let Some(pid) = handle.id() {
return (handle, rx, pid as i32);
}
assert!(
std::time::Instant::now() < deadline,
"the command never reported a real pid"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
#[tokio::test]
async fn run_scope_global_cancellation_cleans_process_group() {
let scope = RunCommandScope::new();
let cancel = CancellationToken::new();
scope.link_cancellation(cancel.clone());
let runner = scoped_runner(&scope);
let (handle, _rx, pgid) = launch_owned_group(&runner, "change-a").await;
assert!(
group_has_members(pgid),
"arrangement failed: the owned group must exist before shutdown"
);
cancel.cancel();
drop(handle);
let cleanup = scope.wait_quiescent(Duration::from_secs(30)).await;
let survived = reap_and_report_survival(pgid);
assert!(
cleanup.is_quiescent(),
"the barrier must prove quiescence before terminal Stopped: {}",
cleanup.diagnostics()
);
assert!(
!survived,
"process group {pgid} survived global cancellation"
);
}
#[tokio::test]
async fn run_scope_run_fatal_cleans_process_group() {
let scope = RunCommandScope::new();
let runner = scoped_runner(&scope);
let (handle, _rx, pgid) = launch_owned_group(&runner, "change-a").await;
assert!(
group_has_members(pgid),
"arrangement failed: the owned group must exist before shutdown"
);
let (error_tx, mut errors) = tokio::sync::mpsc::channel::<String>(8);
error_tx
.send("Background merge failed for 'change-a'".to_string())
.await
.expect("the prompt Error is published before any waiting");
drop(error_tx);
scope.close();
assert!(
scope.is_closed(),
"run-fatal shutdown closes command admission"
);
drop(handle);
let cleanup = scope.wait_quiescent(Duration::from_secs(30)).await;
let survived = reap_and_report_survival(pgid);
assert!(
cleanup.is_quiescent(),
"failure return must follow proven cleanup: {}",
cleanup.diagnostics()
);
assert!(
!survived,
"process group {pgid} survived run-fatal shutdown"
);
let mut emitted = Vec::new();
while let Some(message) = errors.recv().await {
emitted.push(message);
}
assert_eq!(
emitted.len(),
1,
"exactly one prompt global Error for the run-fatal outcome, got {emitted:?}"
);
let (mut refused, mut refused_rx) = runner
.execute_streaming_with_retry("echo never", None, Some("archive"), Some("change-b"))
.await
.expect("the call returns a refusal rather than launching");
while refused_rx.recv().await.is_some() {}
assert!(!refused.wait().await.expect("status").success());
}
#[tokio::test]
async fn run_scope_tui_quit_cleans_process_group_after_timeout() {
let scope = RunCommandScope::new();
let mut runner = scoped_runner(&scope);
runner.set_process_group_cleanup_timeout_ms(0);
let (handle, mut rx, pgid) = launch_owned_group(&runner, "change-a").await;
assert!(
group_has_members(pgid),
"arrangement failed: the owned group must exist before shutdown"
);
scope.close();
drop(handle);
while rx.recv().await.is_some() {}
let deadline = std::time::Instant::now() + Duration::from_secs(10);
while scope.retained_process_ids().is_empty() {
if std::time::Instant::now() >= deadline {
reap_and_report_survival(pgid);
panic!("the unprovable cleanup must retain its owned identity");
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
let orchestrator = tokio::spawn(async move {
std::future::pending::<()>().await;
Ok(())
});
let outcome = conflux::tui::shutdown_local_orchestrator_task(
Some(orchestrator),
Some(CancellationToken::new()),
Some(scope.clone()),
Duration::from_millis(50),
)
.await;
let survived = reap_and_report_survival(pgid);
assert_eq!(
outcome,
conflux::tui::LocalOrchestratorShutdownOutcome::AbortedAfterTimeout,
"the timeout outcome stays distinguishable from graceful completion"
);
assert!(
!survived,
"process group {pgid} survived the local TUI timeout abort"
);
assert!(
scope.retained_process_ids().is_empty(),
"forced cleanup must verify, not merely signal"
);
}
}