#![cfg(all(unix, feature = "process"))]
use std::error::Error;
use std::io::{BufRead, BufReader, ErrorKind};
use std::os::unix::process::{CommandExt, ExitStatusExt};
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};
use lgwks_std::process::{
child_has_exited_without_reaping, kill_process_group, process_group_exists,
};
use crate::rng::Rng;
use crate::seeded_sweep::{SWEEP_SEEDS, fold, fold_usize, initial_trace, word_of};
type TestResult = Result<(), Box<dyn Error>>;
const ESRCH: i32 = 3;
const EPERM: i32 = 1;
const SIGKILL: i32 = 9;
const VACANT_FLOOR: i32 = 4_194_305;
const SETTLE: Duration = Duration::from_secs(10);
const STEPS: usize = 24;
const MAX_GROUPS: usize = 4;
const MAX_MEMBERS: usize = 3;
const DRAWN_IDS: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Phase {
Live,
Killed,
Reaped,
}
impl Phase {
const fn code(self) -> u64 {
match self {
Self::Live => 1,
Self::Killed => 2,
Self::Reaped => 3,
}
}
}
struct Group {
leader: Child,
id: i32,
phase: Phase,
}
impl Group {
fn spawn(members: usize) -> Result<Self, Box<dyn Error>> {
let mut script = "sleep 30 & ".repeat(members);
script.push_str("echo ready; exec sleep 30");
let leader = Command::new("sh")
.arg("-c")
.arg(&script)
.process_group(0)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.spawn()?;
let id = i32::try_from(leader.id())?;
let mut group = Self {
leader,
id,
phase: Phase::Live,
};
let stdout = group
.leader
.stdout
.take()
.ok_or("the leader's stdout was not piped")?;
let mut ready = String::new();
BufReader::new(stdout).read_line(&mut ready)?;
assert_eq!(
ready.trim_end(),
"ready",
"group {id}: the leader must report its members forked"
);
Ok(group)
}
fn kill(&mut self) -> TestResult {
kill_process_group(self.id)?;
self.phase = Phase::Killed;
Ok(())
}
fn reap(&mut self) -> Result<Option<i32>, Box<dyn Error>> {
let status = self.leader.wait()?;
self.phase = Phase::Reaped;
let settled = settles_absent(self.id)? || holder_group(self.id)? == Some(self.id);
let survivors = if settled {
String::new()
} else {
group_listing(self.id)?
};
assert!(
settled,
"group {} outlived {SETTLE:?} after its reap; still in it: {survivors}",
self.id
);
Ok(status.signal())
}
fn stop(&mut self) -> Result<Option<i32>, Box<dyn Error>> {
self.kill()?;
self.reap()
}
}
impl Drop for Group {
fn drop(&mut self) {
if self.phase != Phase::Reaped {
drop(kill_process_group(self.id));
drop(self.leader.wait());
}
}
}
fn settles_absent(id: i32) -> Result<bool, Box<dyn Error>> {
let deadline = Instant::now()
.checked_add(SETTLE)
.ok_or("the settle deadline overflows the clock")?;
while Instant::now() < deadline {
if !process_group_exists(id)? {
return Ok(true);
}
std::thread::yield_now();
}
Ok(false)
}
fn holder_group(pid: i32) -> Result<Option<i32>, Box<dyn Error>> {
let listed = Command::new("ps")
.args(["-o", "pgid=", "-p", &pid.to_string()])
.stdin(Stdio::null())
.stderr(Stdio::null())
.output()?;
if !listed.status.success() {
return Ok(None);
}
let text = String::from_utf8_lossy(&listed.stdout);
Ok(text.trim().parse().ok())
}
fn group_listing(id: i32) -> Result<String, Box<dyn Error>> {
let listed = Command::new("ps")
.args(["-A", "-o", "pid=,pgid=,ppid=,stat=,comm="])
.stdin(Stdio::null())
.stderr(Stdio::null())
.output()?;
let text = String::from_utf8_lossy(&listed.stdout);
let wanted = id.to_string();
let rows: Vec<&str> = text
.lines()
.filter(|row| row.split_whitespace().nth(1) == Some(wanted.as_str()))
.collect();
Ok(rows.join(" | "))
}
fn exit_observed(pid: i32) -> Result<bool, Box<dyn Error>> {
let deadline = Instant::now()
.checked_add(SETTLE)
.ok_or("the settle deadline overflows the clock")?;
while Instant::now() < deadline {
if child_has_exited_without_reaping(pid)? {
return Ok(true);
}
std::thread::yield_now();
}
Ok(false)
}
fn non_positive(rng: &mut Rng) -> Result<i32, Box<dyn Error>> {
let magnitude = i32::try_from(rng.below(usize::try_from(i32::MAX)?))?;
Ok(0_i32.saturating_sub(magnitude))
}
fn refused_ids(seed: u64) -> Result<Vec<i32>, Box<dyn Error>> {
let mut rng = Rng::new(seed);
let mut ids = vec![0, -1, i32::MIN, i32::MIN.saturating_add(1)];
for _ in 0..DRAWN_IDS {
ids.push(non_positive(&mut rng)?);
}
Ok(ids)
}
fn vacant_ids(seed: u64) -> Result<Vec<i32>, Box<dyn Error>> {
let mut rng = Rng::new(seed);
let span = usize::try_from(i32::MAX.saturating_sub(VACANT_FLOOR))?;
let mut ids = vec![VACANT_FLOOR, i32::MAX];
for _ in 0..DRAWN_IDS {
ids.push(VACANT_FLOOR.saturating_add(i32::try_from(rng.below(span))?));
}
Ok(ids)
}
fn is_esrch(error: &std::io::Error) -> bool {
error.raw_os_error() == Some(ESRCH)
}
fn zombie_group_refusal(error: &std::io::Error) -> bool {
cfg!(target_os = "macos") && error.raw_os_error() == Some(EPERM)
}
fn check_probe(group: &Group, at: &str) -> TestResult {
let present = process_group_exists(group.id)?;
match group.phase {
Phase::Live => assert!(present, "{at}: a live group must be present"),
Phase::Killed => {}
Phase::Reaped => assert!(
!present || holder_group(group.id)? == Some(group.id),
"{at}: a reaped group must be absent"
),
}
Ok(())
}
fn check_kill(group: &mut Group, at: &str) -> TestResult {
match group.phase {
Phase::Live => group.kill(),
Phase::Killed => {
if let Err(error) = kill_process_group(group.id) {
assert!(
is_esrch(&error) || zombie_group_refusal(&error),
"{at}: a second kill must succeed or say ESRCH, got {error}"
);
}
Ok(())
}
Phase::Reaped => Ok(()),
}
}
fn check_observe(group: &Group, at: &str) -> TestResult {
match group.phase {
Phase::Live => {
let exited = child_has_exited_without_reaping(group.id)?;
assert!(!exited, "{at}: a live leader must not read as exited");
}
Phase::Killed => {
let exited = exit_observed(group.id)?;
assert!(
exited,
"{at}: a killed leader must read as exited within {SETTLE:?}"
);
}
Phase::Reaped => assert!(
child_has_exited_without_reaping(group.id).is_err()
|| holder_group(group.id)?.is_some(),
"{at}: a reaped leader is no longer a child to observe"
),
}
Ok(())
}
fn check_reap(group: &mut Group, at: &str) -> TestResult {
let signal = match group.phase {
Phase::Live => group.stop()?,
Phase::Killed => group.reap()?,
Phase::Reaped => return Ok(()),
};
assert_eq!(signal, Some(SIGKILL), "{at}: the leader ends by SIGKILL");
Ok(())
}
fn lifecycle_step(rng: &mut Rng, trace: &mut u64, groups: &mut Vec<Group>, at: &str) -> TestResult {
let live = groups
.iter()
.filter(|group| group.phase != Phase::Reaped)
.count();
let op = if live == 0 || (groups.len() < MAX_GROUPS && rng.below(4) == 0) {
0
} else {
rng.below(4).saturating_add(1)
};
fold_usize(trace, op);
if op == 0 {
let members = rng.below(MAX_MEMBERS.saturating_add(1));
fold_usize(trace, members);
groups.push(Group::spawn(members)?);
return Ok(());
}
let slot = rng.below(groups.len());
let group = groups
.get_mut(slot)
.ok_or_else(|| format!("{at}: slot {slot} is out of range"))?;
fold_usize(trace, slot);
fold(trace, group.phase.code());
let checked = match op {
1 => check_probe(group, at),
2 => check_kill(group, at),
3 => check_observe(group, at),
_ => check_reap(group, at),
};
checked?;
fold(trace, group.phase.code());
Ok(())
}
fn lifecycle(seed: u64) -> Result<u64, Box<dyn Error>> {
let mut rng = Rng::new(seed);
let mut trace = initial_trace();
let mut groups: Vec<Group> = Vec::new();
for step in 0..STEPS {
let at = format!("seed {seed:#x} step {step}");
lifecycle_step(&mut rng, &mut trace, &mut groups, &at)?;
}
let drain = format!("seed {seed:#x} drain");
for group in &mut groups {
check_reap(group, &drain)?;
fold(&mut trace, group.phase.code());
}
Ok(trace)
}
#[test]
fn a_seeded_lifecycle_agrees_with_the_model_at_every_step() -> TestResult {
for seed in SWEEP_SEEDS {
lifecycle(seed)?;
}
Ok(())
}
#[test]
fn the_same_seed_replays_the_same_lifecycle_trace() -> TestResult {
for seed in SWEEP_SEEDS {
assert_eq!(
lifecycle(seed)?,
lifecycle(seed)?,
"seed {seed:#x}: the same seed must replay the same lifecycle"
);
}
Ok(())
}
#[test]
fn distinct_seeds_drive_distinct_lifecycles() -> TestResult {
let [first, second, ..] = SWEEP_SEEDS;
assert_ne!(
lifecycle(first)?,
lifecycle(second)?,
"seeds {first:#x} and {second:#x} must schedule different lifecycles"
);
Ok(())
}
#[test]
fn every_non_positive_id_is_refused_by_the_probe() -> TestResult {
for seed in SWEEP_SEEDS {
for id in refused_ids(seed)? {
assert_eq!(
process_group_exists(id).map_err(|error| error.kind()),
Err(ErrorKind::InvalidInput),
"seed {seed:#x}: id {id} names the caller or a process, never a group"
);
}
}
Ok(())
}
#[test]
fn every_non_positive_id_is_refused_by_the_kill() -> TestResult {
for seed in SWEEP_SEEDS {
for id in refused_ids(seed)? {
assert_eq!(
kill_process_group(id).map_err(|error| error.kind()),
Err(ErrorKind::InvalidInput),
"seed {seed:#x}: id {id} must be refused before any signal is sent"
);
}
}
Ok(())
}
#[test]
fn every_non_positive_pid_is_refused_by_the_exit_observation() -> TestResult {
for seed in SWEEP_SEEDS {
for pid in refused_ids(seed)? {
assert_eq!(
child_has_exited_without_reaping(pid).map_err(|error| error.kind()),
Err(ErrorKind::InvalidInput),
"seed {seed:#x}: pid {pid} names no child"
);
}
}
Ok(())
}
#[test]
fn ids_beyond_every_pid_ceiling_name_no_group() -> TestResult {
for seed in SWEEP_SEEDS {
for id in vacant_ids(seed)? {
assert!(
!process_group_exists(id)?,
"seed {seed:#x}: id {id} is above every pid ceiling, so no group can hold it"
);
}
}
Ok(())
}
#[test]
fn killing_a_vacant_group_reports_esrch_not_success() -> TestResult {
for seed in SWEEP_SEEDS {
for id in vacant_ids(seed)? {
let killed = kill_process_group(id);
assert!(
killed.as_ref().is_err_and(is_esrch),
"seed {seed:#x}: killing vacant group {id} must say ESRCH, got {killed:?}"
);
}
}
Ok(())
}
#[test]
fn observing_a_vacant_pid_is_an_error_not_an_exit() -> TestResult {
for seed in SWEEP_SEEDS {
for pid in vacant_ids(seed)? {
assert!(
child_has_exited_without_reaping(pid).is_err(),
"seed {seed:#x}: pid {pid} is not a child, so no exit can be observed"
);
}
}
Ok(())
}
#[test]
fn probing_delivers_no_signal() -> TestResult {
for seed in SWEEP_SEEDS {
let mut rng = Rng::new(seed);
let mut group = Group::spawn(rng.below(MAX_MEMBERS.saturating_add(1)))?;
for probe in 0..rng.below(512).saturating_add(64) {
assert!(
process_group_exists(group.id)?,
"seed {seed:#x} probe {probe}: the probe must find the live group"
);
}
assert!(
!child_has_exited_without_reaping(group.id)?,
"seed {seed:#x}: signal zero must not have ended the leader"
);
group.stop()?;
}
Ok(())
}
#[test]
fn a_group_kill_reaches_every_member() -> TestResult {
for seed in SWEEP_SEEDS {
let members = Rng::new(seed).below(MAX_MEMBERS).saturating_add(1);
let mut group = Group::spawn(members)?;
group.stop()?;
assert!(
!process_group_exists(group.id)? || holder_group(group.id)? == Some(group.id),
"seed {seed:#x}: all {members} members must be gone after the group kill"
);
}
Ok(())
}
#[test]
fn a_killed_leader_reports_sigkill_when_reaped() -> TestResult {
for seed in SWEEP_SEEDS {
let mut group = Group::spawn(Rng::new(seed).below(MAX_MEMBERS.saturating_add(1)))?;
assert_eq!(
group.stop()?,
Some(SIGKILL),
"seed {seed:#x}: the reaped status must name the signal the kill sent"
);
}
Ok(())
}
#[test]
fn an_exit_observation_never_reaps() -> TestResult {
for seed in SWEEP_SEEDS {
let mut rng = Rng::new(seed);
let mut group = Group::spawn(rng.below(MAX_MEMBERS.saturating_add(1)))?;
group.kill()?;
assert!(
exit_observed(group.id)?,
"seed {seed:#x}: the killed leader must read as exited"
);
for again in 0..rng.below(64).saturating_add(1) {
assert!(
child_has_exited_without_reaping(group.id)?,
"seed {seed:#x} observation {again}: an observation must leave the exit waitable"
);
}
assert_eq!(
group.reap()?,
Some(SIGKILL),
"seed {seed:#x}: the reap after many observations still returns the real status"
);
}
Ok(())
}
#[test]
fn a_reaped_child_is_no_longer_observable() -> TestResult {
for seed in SWEEP_SEEDS {
let mut group = Group::spawn(Rng::new(seed).below(MAX_MEMBERS.saturating_add(1)))?;
group.stop()?;
assert!(
child_has_exited_without_reaping(group.id).is_err()
|| holder_group(group.id)?.is_some(),
"seed {seed:#x}: a reaped pid must not be reported as an exit a second time"
);
}
Ok(())
}
#[test]
fn a_killed_but_unreaped_group_is_probed_without_error() -> TestResult {
for seed in SWEEP_SEEDS {
let mut rng = Rng::new(seed);
let mut group = Group::spawn(rng.below(MAX_MEMBERS.saturating_add(1)))?;
group.kill()?;
for probe in 0..rng.below(64).saturating_add(1) {
assert!(
process_group_exists(group.id).is_ok(),
"seed {seed:#x} probe {probe}: a group held by a zombie leader is observable"
);
}
group.reap()?;
}
Ok(())
}
#[test]
fn two_tenants_groups_are_isolated() -> TestResult {
for seed in SWEEP_SEEDS {
let mut rng = Rng::new(seed);
let mut first: Vec<Group> = Vec::new();
let mut second: Vec<Group> = Vec::new();
for _ in 0..rng.below(3).saturating_add(1) {
first.push(Group::spawn(rng.below(MAX_MEMBERS))?);
second.push(Group::spawn(rng.below(MAX_MEMBERS))?);
}
while !first.is_empty() {
let mut victim = first.swap_remove(rng.below(first.len()));
victim.stop()?;
for (index, survivor) in second.iter().enumerate() {
assert!(
process_group_exists(survivor.id)?,
"seed {seed:#x}: tenant two's group {index} must survive tenant one's kill"
);
assert!(
!child_has_exited_without_reaping(survivor.id)?,
"seed {seed:#x}: tenant two's leader {index} must still be running"
);
}
}
for mut survivor in second {
survivor.stop()?;
}
}
Ok(())
}
#[test]
fn concurrent_probes_agree_with_the_model() -> TestResult {
const PROBERS: usize = 16;
const PROBES: usize = 256;
for seed in SWEEP_SEEDS {
let mut live = Group::spawn(1)?;
let expected = [(live.id, true), (VACANT_FLOOR, false)];
let disagreements = std::thread::scope(|scope| {
let mut probers = Vec::with_capacity(PROBERS);
for prober in 0..PROBERS {
let mut rng = Rng::new(seed.wrapping_add(word_of(prober)));
probers.push(scope.spawn(move || {
let mut wrong = 0_usize;
for _ in 0..PROBES {
let Some(&(id, present)) = expected.get(rng.below(expected.len())) else {
continue;
};
if process_group_exists(id).ok() != Some(present) {
wrong = wrong.saturating_add(1);
}
}
wrong
}));
}
probers
.into_iter()
.map(|prober| match prober.join() {
Ok(wrong) => wrong,
Err(_) => usize::MAX,
})
.fold(0_usize, usize::saturating_add)
});
assert_eq!(
disagreements, 0,
"seed {seed:#x}: {PROBERS} concurrent probers must each see the model's answer"
);
live.stop()?;
}
Ok(())
}
#[test]
fn a_kill_retried_before_the_reap_is_idempotent() -> TestResult {
for seed in SWEEP_SEEDS {
let mut rng = Rng::new(seed);
let mut group = Group::spawn(rng.below(MAX_MEMBERS.saturating_add(1)))?;
group.kill()?;
for retry in 0..rng.below(16).saturating_add(1) {
if let Err(error) = kill_process_group(group.id) {
assert!(
is_esrch(&error) || zombie_group_refusal(&error),
"seed {seed:#x} retry {retry}: a retried kill must succeed or say ESRCH, got {error}"
);
}
}
assert_eq!(
group.reap()?,
Some(SIGKILL),
"seed {seed:#x}: the retries do not change how the leader ended"
);
}
Ok(())
}
#[test]
fn a_leader_that_exits_by_itself_leaves_its_group_absent_once_reaped() -> TestResult {
for seed in SWEEP_SEEDS {
let code = i32::try_from(Rng::new(seed).below(64))?;
let mut leader = Command::new("sh")
.arg("-c")
.arg(format!("exit {code}"))
.process_group(0)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()?;
let id = i32::try_from(leader.id())?;
assert!(
exit_observed(id)?,
"seed {seed:#x}: the leader's own exit must be observed"
);
let status = leader.wait()?;
assert_eq!(
status.code(),
Some(code),
"seed {seed:#x}: the reap returns the exit code"
);
assert!(
settles_absent(id)? || holder_group(id)? == Some(id),
"seed {seed:#x}: a group whose leader exited and was reaped must be absent"
);
}
Ok(())
}