use std::sync::Arc;
use std::time::Duration;
use processkit::{Command, Error, ErrorReason, OutputBufferPolicy, ProcessGroup, RunningProcess};
use crate::common::*;
const SEED_COUNT: u64 = 48;
const BASE_SEED: u64 = 0x00C0_FFEE_1234_5678;
const FLOOD_LINES: u32 = 3_000;
const REAP_GRACE: Duration = Duration::from_secs(20);
const MAX_REPORTED: usize = 12;
struct Rng {
state: u64,
}
impl Rng {
fn new(seed: u64) -> Self {
Rng { state: seed }
}
fn next_u64(&mut self) -> u64 {
self.state = self.state.wrapping_add(0x9E37_79B9_7F4A_7C15);
let mut z = self.state;
z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
z ^ (z >> 31)
}
fn below(&mut self, n: u32) -> u32 {
(self.next_u64() % u64::from(n)) as u32
}
fn one_in(&mut self, n: u32) -> bool {
self.below(n) == 0
}
fn choose<T: Copy>(&mut self, items: &[T]) -> T {
items[self.below(items.len() as u32) as usize]
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ChildClass {
ShortLived,
LongLived,
StdoutFlooder,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Terminal {
Wait,
OutputString,
Finish,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ChildPlan {
class: ChildClass,
keep_stdin: bool,
take_stdin: bool,
start_kill: bool,
inspect: bool,
wait_for_line: bool,
terminal: Terminal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum GroupOp {
Mechanism,
ShutdownRef,
#[cfg(feature = "process-control")]
Members,
#[cfg(feature = "process-control")]
MembersInfo,
#[cfg(feature = "process-control")]
SuspendResume,
#[cfg(feature = "process-control")]
Signal(processkit::Signal),
#[cfg(feature = "stats")]
Stats,
}
#[derive(Debug, Clone)]
struct ComboPlan {
children: Vec<ChildPlan>,
group_ops: Vec<GroupOp>,
interleave_ms: u64,
}
fn group_op_pool() -> Vec<GroupOp> {
#[allow(unused_mut)]
let mut pool = vec![GroupOp::Mechanism, GroupOp::ShutdownRef];
#[cfg(feature = "process-control")]
{
pool.push(GroupOp::Members);
pool.push(GroupOp::MembersInfo);
pool.push(GroupOp::SuspendResume);
pool.push(GroupOp::Signal(processkit::Signal::Term));
pool.push(GroupOp::Signal(processkit::Signal::Int));
pool.push(GroupOp::Signal(processkit::Signal::Hup));
}
#[cfg(feature = "stats")]
{
pool.push(GroupOp::Stats);
}
pool
}
fn gen_combo_plan(seed: u64) -> ComboPlan {
let mut rng = Rng::new(seed);
let n_children = 2 + rng.below(3); let mut children = Vec::with_capacity(n_children as usize);
for _ in 0..n_children {
let class = rng.choose(&[
ChildClass::ShortLived,
ChildClass::LongLived,
ChildClass::StdoutFlooder,
]);
let keep_stdin = rng.one_in(2);
let take_stdin = rng.one_in(2);
let start_kill = rng.one_in(3);
let inspect = rng.one_in(2);
let terminal = rng.choose(&[Terminal::Wait, Terminal::OutputString, Terminal::Finish]);
let wait_for_line = matches!(class, ChildClass::StdoutFlooder)
&& matches!(terminal, Terminal::Wait | Terminal::Finish)
&& rng.one_in(2);
children.push(ChildPlan {
class,
keep_stdin,
take_stdin,
start_kill,
inspect,
wait_for_line,
terminal,
});
}
let pool = group_op_pool();
let n_ops = 2 + rng.below(4); let mut group_ops = Vec::with_capacity(n_ops as usize);
for _ in 0..n_ops {
group_ops.push(rng.choose(pool.as_slice()));
}
let interleave_ms = u64::from(5 + rng.below(35));
ComboPlan {
children,
group_ops,
interleave_ms,
}
}
fn child_command(plan: &ChildPlan) -> Command {
let mut cmd = match plan.class {
ChildClass::ShortLived => quick_exit(),
ChildClass::LongLived => long_sleeper(),
ChildClass::StdoutFlooder => {
line_emitter(FLOOD_LINES).output_buffer(OutputBufferPolicy::bounded(1000))
}
};
if plan.keep_stdin {
cmd = cmd.keep_stdin_open();
}
cmd
}
fn is_impossible_error(e: &Error) -> bool {
if matches!(
e.reason(),
ErrorReason::CassetteMiss { .. }
| ErrorReason::Parse { .. }
| ErrorReason::NotFound { .. }
| ErrorReason::Spawn { .. }
| ErrorReason::OutputTooLarge { .. }
) {
return true;
}
#[cfg(feature = "limits")]
if matches!(e.reason(), ErrorReason::ResourceLimit { .. }) {
return true;
}
false
}
fn check(result: Result<(), Error>, what: &str) -> Result<(), String> {
match result {
Ok(()) => Ok(()),
Err(e) if is_impossible_error(&e) => {
Err(format!("{what} returned impossible error: {e:?}"))
}
Err(_) => Ok(()),
}
}
async fn run_child(mut child: RunningProcess, plan: ChildPlan) -> Result<(), String> {
if plan.inspect {
let _ = child.pid();
let _ = child.elapsed();
let _ = child.stdout_line_count();
}
if plan.take_stdin {
let _ = child.take_stdin();
}
if plan.start_kill {
check(child.start_kill(), "start_kill")?;
}
if plan.wait_for_line {
check(
child
.wait_for_line(|l| !l.is_empty(), Duration::from_secs(5))
.await
.map(|_| ()),
"wait_for_line",
)?;
}
let terminal = match plan.terminal {
Terminal::Wait => child.wait().await.map(|_| ()),
Terminal::OutputString => child.output_string().await.map(|_| ()),
Terminal::Finish => child.finish().await.map(|_| ()),
};
check(terminal, "terminal verb")
}
async fn run_group_op(group: &ProcessGroup, op: GroupOp) -> Result<(), String> {
match op {
GroupOp::Mechanism => {
let _ = group.mechanism();
Ok(())
}
GroupOp::ShutdownRef => check(group.shutdown_ref().await, "group.shutdown_ref"),
#[cfg(feature = "process-control")]
GroupOp::Members => check(group.members().map(|_| ()), "group.members"),
#[cfg(feature = "process-control")]
GroupOp::MembersInfo => check(group.members_info().map(|_| ()), "group.members_info"),
#[cfg(feature = "process-control")]
GroupOp::SuspendResume => {
check(group.suspend(), "group.suspend")?;
check(group.resume(), "group.resume")
}
#[cfg(feature = "process-control")]
GroupOp::Signal(sig) => check(group.signal(sig), "group.signal"),
#[cfg(feature = "stats")]
GroupOp::Stats => check(group.stats().map(|_| ()), "group.stats"),
}
}
async fn run_combo(seed: u64) -> Result<(), String> {
let plan = gen_combo_plan(seed);
let group = match ProcessGroup::new() {
Ok(g) => Arc::new(g),
Err(e) => return Err(format!("ProcessGroup::new failed: {e:?}")),
};
let mut child_handles = Vec::with_capacity(plan.children.len());
for child_plan in &plan.children {
let child = match group.start(&child_command(child_plan)).await {
Ok(c) => c,
Err(e) => return Err(format!("group.start failed: {e:?}")),
};
let child_plan = child_plan.clone();
child_handles.push(tokio::spawn(run_child(child, child_plan)));
}
let mut group_handles = Vec::with_capacity(plan.group_ops.len());
for op in plan.group_ops {
let g = Arc::clone(&group);
group_handles.push(tokio::spawn(async move { run_group_op(&g, op).await }));
}
tokio::time::sleep(Duration::from_millis(plan.interleave_ms)).await;
let _ = group.kill_all();
let mut failures = Vec::new();
for handle in child_handles {
match tokio::time::timeout(REAP_GRACE, handle).await {
Err(_) => failures.push(
"a child handle was not reaped within the grace (survivor/zombie)".to_string(),
),
Ok(Err(join)) if join.is_panic() => {
failures.push("a child task panicked (faulted background/op)".to_string());
}
Ok(Err(_)) => failures.push("a child task was cancelled".to_string()),
Ok(Ok(Err(msg))) => failures.push(msg),
Ok(Ok(Ok(()))) => {}
}
}
for handle in group_handles {
match tokio::time::timeout(REAP_GRACE, handle).await {
Err(_) => {
failures.push("a group-op task did not settle within the grace".to_string());
}
Ok(Err(join)) if join.is_panic() => {
failures.push("a group-op task panicked".to_string());
}
Ok(Err(_)) => failures.push("a group-op task was cancelled".to_string()),
Ok(Ok(Err(msg))) => failures.push(msg),
Ok(Ok(Ok(()))) => {}
}
}
let group = match Arc::into_inner(group) {
Some(g) => g,
None => {
failures.push("internal: group Arc still shared after all tasks joined".to_string());
return Err(failures.join("; "));
}
};
#[cfg(feature = "process-control")]
{
let deadline = std::time::Instant::now() + REAP_GRACE;
loop {
match group.members() {
Ok(m) if m.is_empty() => break,
Ok(_) if std::time::Instant::now() >= deadline => {
failures.push("group still reports live members after teardown".to_string());
break;
}
Ok(_) => tokio::time::sleep(Duration::from_millis(20)).await,
Err(e) if is_impossible_error(&e) => {
failures.push(format!("group.members returned impossible error: {e:?}"));
break;
}
Err(_) => break,
}
}
}
if let Err(e) = group.shutdown().await
&& is_impossible_error(&e)
{
failures.push(format!("group.shutdown returned impossible error: {e:?}"));
}
if failures.is_empty() {
Ok(())
} else {
Err(failures.join("; "))
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn randomized_lifecycle_interleavings() {
if skip_unless_enabled("randomized_lifecycle_interleavings") {
return;
}
let warmup = quick_exit().output_string().await.expect("warm up a child");
assert!(warmup.is_success());
let before_fds = open_handle_count();
#[cfg(target_os = "linux")]
let before_cgroups = own_cgroup_v2_parent().map(|parent| own_processkit_cgroup_dirs(&parent));
let seeds: Vec<u64> = match std::env::var("PROCESSKIT_STRESS_SEED") {
Ok(raw) => {
let seed: u64 = raw
.trim()
.parse()
.unwrap_or_else(|_| panic!("PROCESSKIT_STRESS_SEED must be a u64, got {raw:?}"));
eprintln!("[stress] interleave: replaying single seed {seed}");
vec![seed]
}
Err(_) => (0..SEED_COUNT).map(|i| BASE_SEED.wrapping_add(i)).collect(),
};
let mut failures = Vec::new();
for &seed in &seeds {
if let Err(msg) = run_combo(seed).await {
failures.push(format!("seed {seed}: {msg}"));
if failures.len() >= MAX_REPORTED {
break;
}
}
}
assert!(
failures.is_empty(),
"randomized interleaving harness found invariant violations \
(re-run one with PROCESSKIT_STRESS_SEED=<seed>):\n{}",
failures.join("\n")
);
if let (Some(before), Some(after)) = (before_fds, open_handle_count()) {
assert!(
after <= before + 32,
"fd/handle count grew across the interleaving sweep: {before} -> {after}"
);
}
#[cfg(target_os = "linux")]
if let (Some(parent), Some(before)) = (own_cgroup_v2_parent(), before_cgroups) {
let leaked: Vec<_> = own_processkit_cgroup_dirs(&parent)
.into_iter()
.filter(|p| !before.contains(p))
.collect();
assert!(
leaked.is_empty(),
"cgroup directories leaked after the interleaving sweep: {leaked:?}"
);
}
}
#[test]
fn combo_plans_are_seed_deterministic() {
for seed in [1u64, 42, 12_345, BASE_SEED] {
assert_eq!(
format!("{:?}", gen_combo_plan(seed)),
format!("{:?}", gen_combo_plan(seed)),
"the same seed must reproduce an identical plan"
);
}
assert_ne!(
format!("{:?}", gen_combo_plan(1)),
format!("{:?}", gen_combo_plan(2)),
"different seeds should generally produce different plans"
);
}