use std::io;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::process::Child;
use tracing::{debug, warn};
#[cfg(unix)]
pub struct ProcessHandle {
pid: Option<u32>,
}
#[cfg(windows)]
pub struct ProcessHandle {
job: Option<JobObjectGuard>,
}
#[cfg(windows)]
struct JobObjectGuard {
handle: windows::Win32::Foundation::HANDLE,
}
#[cfg(windows)]
unsafe impl Send for JobObjectGuard {}
#[cfg(windows)]
unsafe impl Sync for JobObjectGuard {}
#[cfg(windows)]
impl Drop for JobObjectGuard {
fn drop(&mut self) {
use windows::Win32::Foundation::CloseHandle;
unsafe {
let _ = CloseHandle(self.handle);
}
debug!("Job object handle closed");
}
}
pub struct ManagedChild {
pub child: Child,
pub handle: ProcessHandle,
}
#[allow(dead_code)]
#[derive(Debug)]
pub enum TerminationOutcome {
Exited(std::process::ExitStatus),
ForceKilled(std::process::ExitStatus),
TimedOut,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommandTermination {
Exited,
InactivityTimeout,
RuntimeLimit,
Cancelled,
NotStarted,
}
impl CommandTermination {
#[allow(dead_code)] pub fn permits_retry(self) -> bool {
matches!(self, Self::Exited | Self::InactivityTimeout)
}
pub fn is_runtime_limit(self) -> bool {
matches!(self, Self::RuntimeLimit)
}
#[allow(dead_code)] pub fn as_str(self) -> &'static str {
match self {
Self::Exited => "exited",
Self::InactivityTimeout => "inactivity_timeout",
Self::RuntimeLimit => "runtime_limit",
Self::Cancelled => "cancelled",
Self::NotStarted => "not_started",
}
}
}
impl ManagedChild {
pub fn new(mut child: Child) -> io::Result<Self> {
let handle = Self::create_handle(&mut child)?;
Ok(Self { child, handle })
}
#[cfg(unix)]
fn create_handle(child: &mut Child) -> io::Result<ProcessHandle> {
Ok(ProcessHandle { pid: child.id() })
}
#[cfg(windows)]
fn create_handle(child: &mut Child) -> io::Result<ProcessHandle> {
let job = assign_to_job(child)?;
Ok(ProcessHandle { job: Some(job) })
}
pub fn terminate(&mut self) -> io::Result<()> {
self.handle.terminate(&self.child)
}
pub async fn force_kill(&mut self) -> io::Result<()> {
#[cfg(unix)]
{
self.handle.force_kill()
}
#[cfg(windows)]
{
self.child.kill().await
}
}
pub async fn terminate_with_timeout(
&mut self,
timeout: Duration,
) -> io::Result<TerminationOutcome> {
self.terminate()?;
match tokio::time::timeout(timeout, self.wait()).await {
Ok(status) => Ok(TerminationOutcome::Exited(status?)),
Err(_) => {
self.force_kill().await?;
match tokio::time::timeout(timeout, self.wait()).await {
Ok(status) => Ok(TerminationOutcome::ForceKilled(status?)),
Err(_) => Ok(TerminationOutcome::TimedOut),
}
}
}
}
#[allow(dead_code)]
pub fn id(&self) -> Option<u32> {
self.child.id()
}
pub async fn wait(&mut self) -> io::Result<std::process::ExitStatus> {
self.child.wait().await
}
#[allow(dead_code)]
pub async fn kill(&mut self) -> io::Result<()> {
self.child.kill().await
}
}
pub struct StreamingChildHandle {
cancel_tx: Option<tokio::sync::oneshot::Sender<()>>,
current_pid: Arc<AtomicU32>,
final_status_rx: tokio::sync::oneshot::Receiver<std::process::ExitStatus>,
cleanup_rx: tokio::sync::oneshot::Receiver<ProcessGroupCleanupReport>,
cleanup_report: Option<ProcessGroupCleanupReport>,
termination_rx: tokio::sync::oneshot::Receiver<CommandTermination>,
termination: Option<CommandTermination>,
}
#[allow(dead_code)] impl StreamingChildHandle {
pub fn new(
cancel_tx: tokio::sync::oneshot::Sender<()>,
current_pid: Arc<AtomicU32>,
final_status_rx: tokio::sync::oneshot::Receiver<std::process::ExitStatus>,
cleanup_rx: tokio::sync::oneshot::Receiver<ProcessGroupCleanupReport>,
termination_rx: tokio::sync::oneshot::Receiver<CommandTermination>,
) -> Self {
Self {
cancel_tx: Some(cancel_tx),
current_pid,
final_status_rx,
cleanup_rx,
cleanup_report: None,
termination_rx,
termination: None,
}
}
#[cfg(test)]
pub(crate) fn for_test(
cleanup: ProcessGroupCleanupReport,
termination: CommandTermination,
status: std::process::ExitStatus,
) -> Self {
let (cancel_tx, _cancel_rx) = tokio::sync::oneshot::channel::<()>();
let (status_tx, status_rx) = tokio::sync::oneshot::channel::<std::process::ExitStatus>();
let (cleanup_tx, cleanup_rx) = tokio::sync::oneshot::channel::<ProcessGroupCleanupReport>();
let (termination_tx, termination_rx) =
tokio::sync::oneshot::channel::<CommandTermination>();
let _ = status_tx.send(status);
let _ = cleanup_tx.send(cleanup);
let _ = termination_tx.send(termination);
Self::new(
cancel_tx,
Arc::new(AtomicU32::new(0)),
status_rx,
cleanup_rx,
termination_rx,
)
}
pub async fn termination(&mut self) -> CommandTermination {
if let Some(termination) = self.termination {
return termination;
}
let termination = (&mut self.termination_rx)
.await
.unwrap_or(CommandTermination::NotStarted);
self.termination = Some(termination);
termination
}
pub async fn process_group_cleanup(&mut self) -> ProcessGroupCleanupReport {
if let Some(report) = &self.cleanup_report {
return report.clone();
}
let report = match (&mut self.cleanup_rx).await {
Ok(report) => report,
Err(_) => ProcessGroupCleanupReport::missing(
"the command runner ended without publishing process-group cleanup evidence",
),
};
self.cleanup_report = Some(report.clone());
report
}
pub fn terminate(&mut self) -> io::Result<()> {
if let Some(tx) = self.cancel_tx.take() {
let _ = tx.send(());
}
Ok(())
}
pub async fn terminate_with_timeout(
&mut self,
timeout: Duration,
) -> io::Result<TerminationOutcome> {
self.terminate()?;
match tokio::time::timeout(timeout, &mut self.final_status_rx).await {
Ok(Ok(status)) => Ok(TerminationOutcome::Exited(status)),
Ok(Err(_)) => {
Ok(TerminationOutcome::ForceKilled({
#[cfg(unix)]
{
use std::os::unix::process::ExitStatusExt;
std::process::ExitStatus::from_raw(0)
}
#[cfg(not(unix))]
{
use std::os::windows::process::ExitStatusExt;
std::process::ExitStatus::from_raw(0)
}
}))
}
Err(_elapsed) => Ok(TerminationOutcome::TimedOut),
}
}
pub async fn kill(&mut self) -> io::Result<()> {
self.terminate()
}
pub async fn wait(&mut self) -> io::Result<std::process::ExitStatus> {
(&mut self.final_status_rx).await.map_err(|_| {
io::Error::new(io::ErrorKind::BrokenPipe, "streaming child handle dropped")
})
}
pub fn id(&self) -> Option<u32> {
let pid = self.current_pid.load(Ordering::SeqCst);
if pid == 0 {
None
} else {
Some(pid)
}
}
}
impl ProcessHandle {
#[cfg(unix)]
pub fn terminate(&self, _child: &Child) -> io::Result<()> {
use nix::sys::signal::{killpg, Signal};
use nix::unistd::Pid;
if let Some(pid) = self.pid {
debug!("Sending SIGTERM to process group {}", pid);
match killpg(Pid::from_raw(pid as i32), Signal::SIGTERM) {
Ok(_) => {
debug!("Successfully sent SIGTERM to process group {}", pid);
Ok(())
}
Err(e) => {
warn!("Failed to send SIGTERM to process group {}: {}", pid, e);
Err(io::Error::other(e))
}
}
} else {
warn!("No PID available for process group termination");
Ok(())
}
}
#[cfg(unix)]
pub fn force_kill(&self) -> io::Result<()> {
use nix::sys::signal::{killpg, Signal};
use nix::unistd::Pid;
if let Some(pid) = self.pid {
debug!("Sending SIGKILL to process group {}", pid);
match killpg(Pid::from_raw(pid as i32), Signal::SIGKILL) {
Ok(_) => {
debug!("Successfully sent SIGKILL to process group {}", pid);
Ok(())
}
Err(e) => {
warn!("Failed to send SIGKILL to process group {}: {}", pid, e);
Err(io::Error::other(e))
}
}
} else {
warn!("No PID available for process group force kill");
Ok(())
}
}
#[cfg(windows)]
pub fn terminate(&self, child: &Child) -> io::Result<()> {
if let Some(pid) = child.id() {
debug!("Terminating Windows process {}", pid);
Ok(())
} else {
warn!("No PID available for Windows process termination");
Ok(())
}
}
}
#[cfg(unix)]
#[allow(dead_code)]
pub fn configure_process_group(cmd: &mut tokio::process::Command) {
use nix::unistd::{setpgid, setsid, Pid};
unsafe {
cmd.pre_exec(|| {
match setsid() {
Ok(_) => {
debug!("Created new session (setsid) for child process");
Ok(())
}
Err(e) => {
warn!("Failed to create new session (setsid): {}", e);
match setpgid(Pid::from_raw(0), Pid::from_raw(0)) {
Ok(_) => {
debug!("Process group created successfully (fallback)");
Ok(())
}
Err(e) => {
warn!("Failed to create process group: {}", e);
Err(io::Error::other(e))
}
}
}
}
});
}
}
#[cfg(windows)]
fn assign_to_job(child: &Child) -> io::Result<JobObjectGuard> {
use std::mem::size_of;
use windows::Win32::Foundation::{CloseHandle, HANDLE};
use windows::Win32::System::JobObjects::*;
use windows::Win32::System::Threading::{OpenProcess, PROCESS_ALL_ACCESS};
unsafe {
let job = CreateJobObjectW(None, windows::core::PCWSTR::null())
.map_err(|e| io::Error::new(io::ErrorKind::Other, e))?;
let mut info = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default();
info.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE;
SetInformationJobObject(
job,
JobObjectExtendedLimitInformation,
&info as *const _ as *const std::ffi::c_void,
size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32,
)
.map_err(|e| io::Error::new(io::ErrorKind::Other, e))?;
let process_handle = OpenProcess(PROCESS_ALL_ACCESS, false, child.id().unwrap())
.map_err(|e| io::Error::new(io::ErrorKind::Other, e))?;
AssignProcessToJobObject(job, process_handle).map_err(|e| {
CloseHandle(process_handle);
CloseHandle(job);
io::Error::new(io::ErrorKind::Other, e)
})?;
CloseHandle(process_handle);
debug!("Process assigned to job object successfully");
Ok(JobObjectGuard { handle: job })
}
}
#[cfg(windows)]
pub fn configure_process_group(_cmd: &mut tokio::process::Command) {
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PostCleanupOutcome {
#[allow(dead_code)]
NoPgid,
Terminated,
Killed,
AlreadyGone,
}
pub const DEFAULT_PROCESS_GROUP_CLEANUP_TIMEOUT_MS: u64 = 10_000;
pub const DEFAULT_PROCESS_GROUP_SIGTERM_GRACE_MS: u64 = 2_000;
const PROCESS_GROUP_PROBE_INTERVAL_MS: u64 = 20;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProcessGroupQuiescence {
Confirmed,
NotApplicable,
MembersRemain,
Unverifiable,
}
impl ProcessGroupQuiescence {
pub fn is_confirmed(self) -> bool {
matches!(self, Self::Confirmed | Self::NotApplicable)
}
pub fn as_str(self) -> &'static str {
match self {
Self::Confirmed => "confirmed",
Self::NotApplicable => "not_applicable",
Self::MembersRemain => "members_remain",
Self::Unverifiable => "unverifiable",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProcessGroupCleanupReport {
quiescence: ProcessGroupQuiescence,
pgid: Option<u32>,
force_killed: bool,
already_gone: bool,
detail: String,
}
impl ProcessGroupCleanupReport {
fn new(
quiescence: ProcessGroupQuiescence,
pgid: Option<u32>,
force_killed: bool,
already_gone: bool,
detail: impl Into<String>,
) -> Self {
Self {
quiescence,
pgid,
force_killed,
already_gone,
detail: detail.into(),
}
}
pub fn not_applicable(reason: impl Into<String>) -> Self {
Self::new(
ProcessGroupQuiescence::NotApplicable,
None,
false,
false,
reason,
)
}
pub fn missing(reason: impl Into<String>) -> Self {
Self::new(
ProcessGroupQuiescence::Unverifiable,
None,
false,
false,
reason,
)
}
pub fn quiescence(&self) -> ProcessGroupQuiescence {
self.quiescence
}
#[allow(dead_code)] pub fn pgid(&self) -> Option<u32> {
self.pgid
}
pub fn force_killed(&self) -> bool {
self.force_killed
}
pub fn already_gone(&self) -> bool {
self.already_gone
}
pub fn is_confirmed(&self) -> bool {
self.quiescence.is_confirmed()
}
#[cfg(test)]
pub(crate) fn for_test(
quiescence: ProcessGroupQuiescence,
pgid: Option<u32>,
detail: &str,
) -> Self {
Self::new(quiescence, pgid, false, false, detail)
}
pub fn diagnostics(&self) -> String {
match self.pgid {
Some(pgid) => format!(
"process-group cleanup {} (pgid={}, force_killed={}): {}",
self.quiescence.as_str(),
pgid,
self.force_killed,
self.detail
),
None => format!(
"process-group cleanup {}: {}",
self.quiescence.as_str(),
self.detail
),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum SignalResult {
Delivered,
AlreadyGone,
Failed(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ProbeResult {
MembersRemain,
Empty,
Unknown(String),
}
pub(crate) trait ProcessGroupControl {
fn signal_term(&self) -> SignalResult;
fn signal_kill(&self) -> SignalResult;
fn probe(&self) -> ProbeResult;
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct CleanupBudget {
pub sigterm_grace: Duration,
pub total: Duration,
pub probe_interval: Duration,
}
impl CleanupBudget {
pub fn from_millis(sigterm_grace_ms: u64, total_ms: u64) -> Self {
Self {
sigterm_grace: Duration::from_millis(sigterm_grace_ms),
total: Duration::from_millis(total_ms),
probe_interval: Duration::from_millis(PROCESS_GROUP_PROBE_INTERVAL_MS),
}
}
}
async fn poll_until_empty<C: ProcessGroupControl>(
control: &C,
deadline: tokio::time::Instant,
probe_interval: Duration,
) -> ProbeResult {
loop {
let observation = control.probe();
if matches!(observation, ProbeResult::Empty) {
return ProbeResult::Empty;
}
let now = tokio::time::Instant::now();
if now >= deadline {
return observation;
}
let step = probe_interval.min(deadline - now);
tokio::time::sleep(step).await;
}
}
fn unconfirmed_from(observation: ProbeResult, phase: &str) -> (ProcessGroupQuiescence, String) {
match observation {
ProbeResult::Unknown(reason) => (
ProcessGroupQuiescence::Unverifiable,
format!("group membership could not be checked after {phase} ({reason})"),
),
_ => (
ProcessGroupQuiescence::MembersRemain,
format!(
"members were still alive after {phase} and the cleanup budget expired; \
inspect surviving processes with `ps -o pid,pgid,command -g <pgid>`"
),
),
}
}
pub(crate) async fn drive_process_group_cleanup<C: ProcessGroupControl>(
control: &C,
pgid: Option<u32>,
budget: CleanupBudget,
) -> ProcessGroupCleanupReport {
let start = tokio::time::Instant::now();
let deadline = start + budget.total;
let mut notes: Vec<String> = Vec::new();
match control.signal_term() {
SignalResult::AlreadyGone => {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
false,
true,
"process group was already gone when SIGTERM was attempted",
);
}
SignalResult::Delivered => {}
SignalResult::Failed(reason) => notes.push(format!("SIGTERM failed: {reason}")),
}
let graceful_deadline = (start + budget.sigterm_grace).min(deadline);
let graceful = poll_until_empty(control, graceful_deadline, budget.probe_interval).await;
if matches!(graceful, ProbeResult::Empty) {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
false,
false,
join_detail("no members remained after graceful termination", ¬es),
);
}
if tokio::time::Instant::now() >= deadline {
let (quiescence, detail) = unconfirmed_from(graceful, "SIGTERM");
return ProcessGroupCleanupReport::new(
quiescence,
pgid,
false,
false,
join_detail(
&format!("cleanup budget expired before SIGKILL escalation: {detail}"),
¬es,
),
);
}
let mut force_killed = false;
match control.signal_kill() {
SignalResult::AlreadyGone => {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
false,
false,
join_detail("process group exited before SIGKILL was delivered", ¬es),
);
}
SignalResult::Delivered => force_killed = true,
SignalResult::Failed(reason) => notes.push(format!("SIGKILL failed: {reason}")),
}
let forced = poll_until_empty(control, deadline, budget.probe_interval).await;
if matches!(forced, ProbeResult::Empty) {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
force_killed,
false,
join_detail("no members remained after forced termination", ¬es),
);
}
let (quiescence, detail) = unconfirmed_from(forced, "SIGKILL");
ProcessGroupCleanupReport::new(
quiescence,
pgid,
force_killed,
false,
join_detail(&detail, ¬es),
)
}
pub(crate) async fn drive_immediate_process_group_kill<C: ProcessGroupControl>(
control: &C,
pgid: Option<u32>,
budget: CleanupBudget,
) -> ProcessGroupCleanupReport {
let deadline = tokio::time::Instant::now() + budget.total;
let mut notes: Vec<String> = Vec::new();
let mut force_killed = false;
match control.signal_kill() {
SignalResult::AlreadyGone => {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
false,
true,
"process group was already gone when SIGKILL was attempted",
);
}
SignalResult::Delivered => force_killed = true,
SignalResult::Failed(reason) => notes.push(format!("SIGKILL failed: {reason}")),
}
let forced = poll_until_empty(control, deadline, budget.probe_interval).await;
if matches!(forced, ProbeResult::Empty) {
return ProcessGroupCleanupReport::new(
ProcessGroupQuiescence::Confirmed,
pgid,
force_killed,
false,
join_detail(
"no members remained after immediate forced termination",
¬es,
),
);
}
let (quiescence, detail) = unconfirmed_from(forced, "SIGKILL");
ProcessGroupCleanupReport::new(
quiescence,
pgid,
force_killed,
false,
join_detail(&detail, ¬es),
)
}
#[cfg(unix)]
pub async fn kill_process_group_immediately(
pgid: u32,
cleanup_timeout_ms: u64,
op: Option<&str>,
change_id: Option<&str>,
) -> ProcessGroupCleanupReport {
if pgid == 0 {
return ProcessGroupCleanupReport::not_applicable(
"no process group id was available for immediate termination",
);
}
let control = UnixProcessGroup { pgid };
let report = drive_immediate_process_group_kill(
&control,
Some(pgid),
CleanupBudget::from_millis(0, cleanup_timeout_ms),
)
.await;
if report.is_confirmed() {
debug!(
pgid,
op,
change_id,
quiescence = report.quiescence().as_str(),
"force-stop: {}",
report.diagnostics()
);
} else {
warn!(
pgid,
op,
change_id,
quiescence = report.quiescence().as_str(),
"force-stop: {}",
report.diagnostics()
);
}
report
}
#[cfg(not(unix))]
pub async fn kill_process_group_immediately(
pgid: u32,
_cleanup_timeout_ms: u64,
_op: Option<&str>,
_change_id: Option<&str>,
) -> ProcessGroupCleanupReport {
debug!("force-stop: no-op on non-Unix platform (pgid={})", pgid);
ProcessGroupCleanupReport::not_applicable(
"descendant lifetime is owned by the Windows job object",
)
}
fn join_detail(base: &str, notes: &[String]) -> String {
if notes.is_empty() {
base.to_string()
} else {
format!("{base} [{}]", notes.join("; "))
}
}
#[cfg(unix)]
struct UnixProcessGroup {
pgid: u32,
}
#[cfg(unix)]
impl UnixProcessGroup {
fn send(&self, signal: nix::sys::signal::Signal) -> SignalResult {
use nix::errno::Errno;
use nix::sys::signal::killpg;
use nix::unistd::Pid;
match killpg(Pid::from_raw(self.pgid as i32), signal) {
Ok(()) => SignalResult::Delivered,
Err(Errno::ESRCH) => SignalResult::AlreadyGone,
Err(e) => SignalResult::Failed(e.to_string()),
}
}
}
#[cfg(unix)]
impl ProcessGroupControl for UnixProcessGroup {
fn signal_term(&self) -> SignalResult {
self.send(nix::sys::signal::Signal::SIGTERM)
}
fn signal_kill(&self) -> SignalResult {
self.send(nix::sys::signal::Signal::SIGKILL)
}
fn probe(&self) -> ProbeResult {
use nix::errno::Errno;
use nix::sys::signal::killpg;
use nix::unistd::Pid;
match killpg(Pid::from_raw(self.pgid as i32), None) {
Ok(()) => ProbeResult::MembersRemain,
Err(Errno::ESRCH) => ProbeResult::Empty,
Err(e) => ProbeResult::Unknown(e.to_string()),
}
}
}
#[cfg(unix)]
pub async fn cleanup_process_group_verified(
pgid: u32,
sigterm_grace_ms: u64,
cleanup_timeout_ms: u64,
op: Option<&str>,
change_id: Option<&str>,
) -> ProcessGroupCleanupReport {
if pgid == 0 {
return ProcessGroupCleanupReport::not_applicable(
"no process group id was available for cleanup",
);
}
let control = UnixProcessGroup { pgid };
let report = drive_process_group_cleanup(
&control,
Some(pgid),
CleanupBudget::from_millis(sigterm_grace_ms, cleanup_timeout_ms),
)
.await;
if report.is_confirmed() {
debug!(
pgid,
op,
change_id,
quiescence = report.quiescence().as_str(),
"post-cleanup: {}",
report.diagnostics()
);
} else {
warn!(
pgid,
op,
change_id,
quiescence = report.quiescence().as_str(),
"post-cleanup: {}",
report.diagnostics()
);
}
report
}
#[cfg(not(unix))]
pub async fn cleanup_process_group_verified(
pgid: u32,
_sigterm_grace_ms: u64,
_cleanup_timeout_ms: u64,
_op: Option<&str>,
_change_id: Option<&str>,
) -> ProcessGroupCleanupReport {
debug!("post-cleanup: no-op on non-Unix platform (pgid={})", pgid);
ProcessGroupCleanupReport::not_applicable(
"descendant lifetime is owned by the Windows job object",
)
}
pub async fn cleanup_process_group(
pgid: u32,
sigterm_grace_ms: u64,
op: Option<&str>,
change_id: Option<&str>,
) -> PostCleanupOutcome {
if pgid == 0 {
return PostCleanupOutcome::NoPgid;
}
let report = cleanup_process_group_verified(
pgid,
sigterm_grace_ms,
DEFAULT_PROCESS_GROUP_CLEANUP_TIMEOUT_MS,
op,
change_id,
)
.await;
match report.quiescence() {
ProcessGroupQuiescence::NotApplicable => PostCleanupOutcome::NoPgid,
_ if report.already_gone() => PostCleanupOutcome::AlreadyGone,
_ if report.force_killed() => PostCleanupOutcome::Killed,
_ => PostCleanupOutcome::Terminated,
}
}
#[cfg(test)]
mod cleanup_driver_tests {
use super::*;
use std::cell::Cell;
struct FakeProcessGroup {
term: SignalResult,
kill: SignalResult,
probes_until_empty_after_term: Cell<Option<u32>>,
empties_on_kill: bool,
probe_error: Option<String>,
killed: Cell<bool>,
term_calls: Cell<u32>,
kill_calls: Cell<u32>,
probe_calls: Cell<u32>,
}
impl FakeProcessGroup {
fn new() -> Self {
Self {
term: SignalResult::Delivered,
kill: SignalResult::Delivered,
probes_until_empty_after_term: Cell::new(None),
empties_on_kill: true,
probe_error: None,
killed: Cell::new(false),
term_calls: Cell::new(0),
kill_calls: Cell::new(0),
probe_calls: Cell::new(0),
}
}
fn with_term(mut self, result: SignalResult) -> Self {
self.term = result;
self
}
fn empties_after_term_probes(self, probes: u32) -> Self {
self.probes_until_empty_after_term.set(Some(probes));
self
}
fn survives_sigkill(mut self) -> Self {
self.empties_on_kill = false;
self
}
fn with_probe_error(mut self, reason: &str) -> Self {
self.probe_error = Some(reason.to_string());
self
}
}
impl ProcessGroupControl for FakeProcessGroup {
fn signal_term(&self) -> SignalResult {
self.term_calls.set(self.term_calls.get() + 1);
self.term.clone()
}
fn signal_kill(&self) -> SignalResult {
self.kill_calls.set(self.kill_calls.get() + 1);
if matches!(self.kill, SignalResult::Delivered) {
self.killed.set(true);
}
self.kill.clone()
}
fn probe(&self) -> ProbeResult {
self.probe_calls.set(self.probe_calls.get() + 1);
if let Some(reason) = &self.probe_error {
return ProbeResult::Unknown(reason.clone());
}
if self.killed.get() && self.empties_on_kill {
return ProbeResult::Empty;
}
match self.probes_until_empty_after_term.get() {
Some(0) => ProbeResult::Empty,
Some(remaining) => {
self.probes_until_empty_after_term.set(Some(remaining - 1));
ProbeResult::MembersRemain
}
None => ProbeResult::MembersRemain,
}
}
}
fn budget() -> CleanupBudget {
CleanupBudget::from_millis(100, 1_000)
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_confirms_when_group_already_gone() {
let group = FakeProcessGroup::new().with_term(SignalResult::AlreadyGone);
let report = drive_process_group_cleanup(&group, Some(4242), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(report.already_gone());
assert!(!report.force_killed());
assert_eq!(group.kill_calls.get(), 0, "SIGKILL must not be escalated");
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_confirms_after_graceful_descendant_exit() {
let group = FakeProcessGroup::new().empties_after_term_probes(3);
let report = drive_process_group_cleanup(&group, Some(11), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(!report.force_killed());
assert_eq!(group.kill_calls.get(), 0);
assert!(group.probe_calls.get() >= 4);
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_confirms_after_forced_termination() {
let group = FakeProcessGroup::new();
let report = drive_process_group_cleanup(&group, Some(12), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(report.force_killed());
assert_eq!(group.kill_calls.get(), 1);
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_never_confirms_from_leader_exit_alone() {
let group = FakeProcessGroup::new().survives_sigkill();
let report = drive_process_group_cleanup(&group, Some(13), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::MembersRemain);
assert!(!report.is_confirmed());
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_reports_members_remain_when_budget_exhausted() {
let group = FakeProcessGroup::new();
let report =
drive_process_group_cleanup(&group, Some(14), CleanupBudget::from_millis(100, 0)).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::MembersRemain);
assert!(!report.is_confirmed());
assert_eq!(
group.kill_calls.get(),
0,
"an exhausted budget must not claim an escalation it never ran"
);
assert!(
report.diagnostics().contains("cleanup budget expired"),
"diagnostics must name the exhausted budget: {}",
report.diagnostics()
);
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_reports_unverifiable_when_membership_unknown() {
let group = FakeProcessGroup::new().with_probe_error("EPERM");
let report = drive_process_group_cleanup(&group, Some(15), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Unverifiable);
assert!(!report.is_confirmed());
assert!(report.diagnostics().contains("EPERM"));
}
#[tokio::test(start_paused = true)]
async fn process_group_cleanup_records_signal_failures_in_diagnostics() {
let group = FakeProcessGroup::new().with_term(SignalResult::Failed("EPERM".to_string()));
let report = drive_process_group_cleanup(&group, Some(16), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(report.force_killed());
assert!(
report.diagnostics().contains("SIGTERM failed: EPERM"),
"diagnostics must retain the signal failure: {}",
report.diagnostics()
);
}
#[tokio::test(start_paused = true)]
async fn force_stop_change_kill_never_sends_sigterm() {
let group = FakeProcessGroup::new();
let report = drive_immediate_process_group_kill(&group, Some(21), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(report.force_killed());
assert_eq!(group.kill_calls.get(), 1);
assert_eq!(
group.term_calls.get(),
0,
"an immediate force-stop must not open a SIGTERM grace window"
);
}
#[tokio::test(start_paused = true)]
async fn force_stop_change_kill_confirms_when_group_already_gone() {
let mut group = FakeProcessGroup::new();
group.kill = SignalResult::AlreadyGone;
let report = drive_immediate_process_group_kill(&group, Some(22), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Confirmed);
assert!(report.already_gone());
assert!(!report.force_killed());
assert_eq!(group.term_calls.get(), 0);
}
#[tokio::test(start_paused = true)]
async fn force_stop_change_kill_never_confirms_a_surviving_group() {
let group = FakeProcessGroup::new().survives_sigkill();
let report = drive_immediate_process_group_kill(&group, Some(23), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::MembersRemain);
assert!(!report.is_confirmed());
assert_eq!(group.term_calls.get(), 0);
}
#[tokio::test(start_paused = true)]
async fn force_stop_change_kill_reports_unverifiable_membership() {
let group = FakeProcessGroup::new().with_probe_error("EPERM");
let report = drive_immediate_process_group_kill(&group, Some(24), budget()).await;
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Unverifiable);
assert!(!report.is_confirmed());
assert!(report.diagnostics().contains("EPERM"));
}
#[test]
fn missing_cleanup_evidence_is_never_confirmed() {
let report = ProcessGroupCleanupReport::missing("runner published no cleanup evidence");
assert!(!report.is_confirmed());
assert_eq!(report.quiescence(), ProcessGroupQuiescence::Unverifiable);
let skipped = ProcessGroupCleanupReport::not_applicable("cleanup disabled");
assert!(skipped.is_confirmed());
}
}
#[cfg(all(test, unix))]
mod tests {
use super::*;
use tokio::process::Command;
#[tokio::test]
async fn terminate_with_timeout_exits_cleanly() {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("sleep 5")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("spawn sleep");
let mut child = ManagedChild::new(child).expect("managed child");
let outcome = child
.terminate_with_timeout(Duration::from_secs(1))
.await
.expect("terminate");
assert!(matches!(
outcome,
TerminationOutcome::Exited(_) | TerminationOutcome::ForceKilled(_)
));
}
#[tokio::test]
async fn terminate_with_timeout_force_kills() {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("trap '' TERM; while true; do sleep 1; done")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("spawn trap");
let mut child = ManagedChild::new(child).expect("managed child");
let outcome = child
.terminate_with_timeout(Duration::from_millis(200))
.await
.expect("terminate");
assert!(matches!(
outcome,
TerminationOutcome::Exited(_)
| TerminationOutcome::ForceKilled(_)
| TerminationOutcome::TimedOut
));
}
fn pgid_is_gone(pgid: u32) -> bool {
use nix::errno::Errno;
use nix::sys::signal::{killpg, Signal};
use nix::unistd::Pid;
match killpg(Pid::from_raw(pgid as i32), Signal::SIGKILL) {
Ok(()) => false, Err(Errno::ESRCH) => true, Err(_) => false,
}
}
#[tokio::test]
async fn successful_command_backgrounded_child_is_cleaned_up() {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("sleep 60 & exit 0")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("spawn");
let mut child = ManagedChild::new(child).expect("managed child");
let pgid = child.id().unwrap_or(0);
assert!(pgid > 0, "process must have a PID");
child.wait().await.expect("wait");
cleanup_process_group(pgid, 50, Some("test"), Some("regression-1.6")).await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
pgid_is_gone(pgid),
"process group {} should be gone after cleanup, but members remain",
pgid
);
}
#[tokio::test]
async fn failed_command_backgrounded_child_is_cleaned_up() {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("sleep 60 & exit 1")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("spawn");
let mut child = ManagedChild::new(child).expect("managed child");
let pgid = child.id().unwrap_or(0);
assert!(pgid > 0);
child.wait().await.expect("wait");
cleanup_process_group(pgid, 50, Some("test"), Some("regression-1.7")).await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
pgid_is_gone(pgid),
"process group {} should be gone after cleanup, but members remain",
pgid
);
}
#[tokio::test]
async fn cancellation_triggers_full_process_group_cleanup() {
let mut cmd = Command::new("sh");
cmd.arg("-c")
.arg("sleep 60 & sleep 60")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
configure_process_group(&mut cmd);
let child = cmd.spawn().expect("spawn");
let mut child = ManagedChild::new(child).expect("managed child");
let pgid = child.id().unwrap_or(0);
assert!(pgid > 0);
let _ = child
.terminate_with_timeout(Duration::from_millis(500))
.await;
cleanup_process_group(pgid, 50, Some("test"), Some("regression-1.8")).await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
pgid_is_gone(pgid),
"process group {} should be gone after cancellation + cleanup",
pgid
);
}
}