use pset::{ProcId, ProcessSet, PsetEvent};
use std::{
collections::HashMap,
fmt::Debug,
hash::Hash,
io::{self, Write},
os::unix::process::ExitStatusExt,
process::{Command, ExitStatus, Stdio},
thread,
time::Duration,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
enum Tag {
Cmd1,
Cmd2,
}
#[derive(Debug, Default)]
struct Seen {
stdout: Vec<u8>,
stderr: Vec<u8>,
status: Option<ExitStatus>,
output_after_exit: usize,
}
impl Seen {
fn out(&self) -> String {
String::from_utf8_lossy(&self.stdout).into_owned()
}
fn err(&self) -> String {
String::from_utf8_lossy(&self.stderr).into_owned()
}
fn code(&self) -> Option<i32> {
self.status.expect("process reported an exit").code()
}
}
fn drain<T: Clone + Eq + Hash>(set: &mut ProcessSet<T>) -> HashMap<T, Seen> {
let mut seen: HashMap<T, Seen> = HashMap::new();
while let Some((tag, event)) = set.wait_next().expect("waiting failed") {
let entry = seen.entry(tag).or_default();
match event {
PsetEvent::Stdout(data) => {
assert!(!data.is_empty(), "an empty chunk is not an event");
if entry.status.is_some() {
entry.output_after_exit += 1;
}
entry.stdout.extend_from_slice(&data);
}
PsetEvent::Stderr(data) => {
assert!(!data.is_empty(), "an empty chunk is not an event");
if entry.status.is_some() {
entry.output_after_exit += 1;
}
entry.stderr.extend_from_slice(&data);
}
PsetEvent::ProcessExited(status) => {
assert!(entry.status.is_none(), "one exit per process");
entry.status = Some(status);
}
}
}
assert!(set.is_empty(), "a drained set has nothing left");
for (_, entry) in seen.iter() {
assert_eq!(entry.output_after_exit, 0, "output arrived after the exit");
}
seen
}
fn sh_piped(script: &str) -> Command {
let mut cmd = sh(script);
cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
cmd
}
fn sh(script: &str) -> Command {
let mut cmd = Command::new("sh");
cmd.arg("-c").arg(script).stdin(Stdio::null());
cmd
}
fn sh_stdin(script: &str) -> Command {
let mut cmd = sh_piped(script);
cmd.stdin(Stdio::piped());
cmd
}
#[test]
fn two_processes_report_their_output_and_their_exits() {
let mut set = pset::create::<Tag>().unwrap();
set.spawn(Tag::Cmd1, sh_piped("echo one")).unwrap();
set.spawn(Tag::Cmd2, sh_piped("echo two")).unwrap();
let seen = drain(&mut set);
assert_eq!(seen.len(), 2);
assert_eq!(seen[&Tag::Cmd1].out(), "one\n");
assert_eq!(seen[&Tag::Cmd2].out(), "two\n");
assert_eq!(seen[&Tag::Cmd1].code(), Some(0));
assert_eq!(seen[&Tag::Cmd2].code(), Some(0));
}
#[test]
fn stdout_and_stderr_stay_apart() {
let mut set = pset::create::<&str>().unwrap();
set.spawn("both", sh_piped("echo out; echo err >&2"))
.unwrap();
let seen = drain(&mut set);
assert_eq!(seen["both"].out(), "out\n");
assert_eq!(seen["both"].err(), "err\n");
}
#[test]
fn exit_codes_are_reported_per_process() {
let mut set = pset::create::<i32>().unwrap();
for code in [0, 1, 7, 42] {
set.spawn(code, sh_piped(&format!("exit {code}"))).unwrap();
}
let seen = drain(&mut set);
assert_eq!(seen.len(), 4);
for (tag, entry) in seen {
assert_eq!(entry.code(), Some(tag));
}
}
#[test]
fn a_process_killed_by_a_signal_reports_the_signal() {
let mut set = pset::create::<()>().unwrap();
let (id, _) = set.spawn((), sh_piped("kill -TERM $$; sleep 30")).unwrap();
assert_eq!(set.len(), 1);
let _ = id;
let seen = drain(&mut set);
let status = seen[&()].status.unwrap();
assert_eq!(status.signal(), Some(libc::SIGTERM));
assert_eq!(status.code(), None);
}
#[test]
fn kill_takes_a_process_down_through_its_pidfd() {
let mut set = pset::create::<&str>().unwrap();
let (id, _): (ProcId, _) = set.spawn("sleeper", sh_piped("exec sleep 30")).unwrap();
assert!(
set.wait_next_timeout(Duration::from_millis(100))
.unwrap()
.is_none()
);
assert!(!set.is_empty(), "the sleeper is still in the set");
set.kill(id, libc::SIGKILL).unwrap();
let seen = drain(&mut set);
assert_eq!(
seen["sleeper"].status.unwrap().signal(),
Some(libc::SIGKILL)
);
}
#[test]
fn kill_all_takes_down_everything_still_running() {
let mut set = pset::create::<usize>().unwrap();
for i in 0..5 {
set.spawn(i, sh_piped("exec sleep 30")).unwrap();
}
set.kill_all(libc::SIGKILL).unwrap();
let seen = drain(&mut set);
assert_eq!(seen.len(), 5);
for (_, entry) in seen {
assert_eq!(entry.status.unwrap().signal(), Some(libc::SIGKILL));
}
}
#[test]
fn output_larger_than_a_pipe_buffer_arrives_whole() {
let mut set = pset::create::<&str>().unwrap();
set.spawn(
"loud",
sh_piped("yes 0123456789012345678901234567890123456789 | head -c 1048576"),
)
.unwrap();
let seen = drain(&mut set);
assert_eq!(seen["loud"].stdout.len(), 1024 * 1024);
assert_eq!(seen["loud"].code(), Some(0));
}
#[test]
fn output_written_just_before_exiting_is_not_lost() {
let mut set = pset::create::<usize>().unwrap();
for i in 0..8 {
set.spawn(i, sh_piped("head -c 200000 /dev/zero")).unwrap();
}
let seen = drain(&mut set);
assert_eq!(seen.len(), 8);
for (_, entry) in seen {
assert_eq!(entry.stdout.len(), 200_000);
assert_eq!(entry.code(), Some(0));
}
}
#[test]
fn binary_output_survives_unchanged() {
let mut set = pset::create::<&str>().unwrap();
set.spawn("bytes", sh_piped(r"printf '\000\001\377\012\200'"))
.unwrap();
let seen = drain(&mut set);
assert_eq!(seen["bytes"].stdout, vec![0o0, 0o1, 0o377, 0o12, 0o200]);
}
#[test]
fn many_processes_are_all_accounted_for() {
const COUNT: usize = 32;
let mut set = pset::create::<usize>().unwrap();
for i in 0..COUNT {
set.spawn(i, sh_piped(&format!("echo {i}; echo {i} >&2")))
.unwrap();
}
assert_eq!(set.len(), COUNT);
let seen = drain(&mut set);
assert_eq!(seen.len(), COUNT);
for (tag, entry) in seen {
assert_eq!(entry.out(), format!("{tag}\n"));
assert_eq!(entry.err(), format!("{tag}\n"));
assert_eq!(entry.code(), Some(0));
}
}
#[test]
fn several_processes_under_one_tag_share_it() {
let mut set = pset::create::<&str>().unwrap();
set.spawn("worker", sh_piped("echo a")).unwrap();
set.spawn("worker", sh_piped("echo b")).unwrap();
let mut exits = 0;
let mut output = Vec::new();
while let Some((tag, event)) = set.wait_next().unwrap() {
assert_eq!(tag, "worker");
match event {
PsetEvent::Stdout(data) => output.extend_from_slice(&data),
PsetEvent::Stderr(data) => panic!("unexpected stderr: {data:?}"),
PsetEvent::ProcessExited(_) => exits += 1,
}
}
assert_eq!(exits, 2);
output.sort_unstable();
assert_eq!(output, b"\n\nab");
}
#[test]
fn a_command_that_cannot_be_spawned_leaves_the_set_untouched() {
let mut set = pset::create::<&str>().unwrap();
let err = set
.spawn("missing", Command::new("/nonexistent/definitely-not-here"))
.unwrap_err();
assert_eq!(err.kind(), std::io::ErrorKind::NotFound);
assert!(set.is_empty());
assert!(set.wait_next().unwrap().is_none());
}
#[test]
fn the_set_stays_usable_after_everything_has_exited() {
let mut set = pset::create::<&str>().unwrap();
set.spawn("first", sh_piped("echo first")).unwrap();
assert_eq!(drain(&mut set)["first"].out(), "first\n");
set.spawn("second", sh_piped("echo second")).unwrap();
assert_eq!(drain(&mut set)["second"].out(), "second\n");
}
#[test]
fn a_piped_stdin_comes_back_from_spawn() {
let mut set = pset::create::<&str>().unwrap();
let (_id, stdin) = set.spawn("cat", sh_stdin("cat")).unwrap();
let mut stdin = stdin.expect("the command asked for a stdin pipe");
stdin.write_all(b"fed in\n").unwrap();
drop(stdin);
let seen = drain(&mut set);
assert_eq!(seen["cat"].out(), "fed in\n");
assert_eq!(seen["cat"].code(), Some(0));
}
#[test]
fn stdin_left_alone_by_the_command_is_left_alone_by_the_set() {
let mut set = pset::create::<&str>().unwrap();
let (_id, null_stdin) = set.spawn("null", sh_piped("exit 0")).unwrap();
let (_id, inherited_stdin) = set.spawn("inherited", Command::new("true")).unwrap();
assert!(null_stdin.is_none(), "a null stdin is not a pipe");
assert!(
inherited_stdin.is_none(),
"an inherited stdin is not a pipe"
);
let seen = drain(&mut set);
assert_eq!(seen["null"].code(), Some(0));
assert_eq!(seen["inherited"].code(), Some(0));
}
#[test]
fn the_child_reads_what_it_is_fed_before_the_pipe_is_closed() {
let mut set = pset::create::<&str>().unwrap();
let (_id, stdin) = set.spawn("cat", sh_stdin("cat")).unwrap();
let mut stdin = stdin.unwrap();
stdin.write_all(b"first\n").unwrap();
assert_eq!(collect_stdout(&mut set, b"first\n".len()), b"first\n");
assert!(
set.wait_next_timeout(Duration::from_millis(100))
.unwrap()
.is_none()
);
assert_eq!(set.len(), 1);
stdin.write_all(b"second\n").unwrap();
assert_eq!(collect_stdout(&mut set, b"second\n".len()), b"second\n");
drop(stdin);
let seen = drain(&mut set);
assert!(
seen["cat"].stdout.is_empty(),
"both lines already collected"
);
assert_eq!(seen["cat"].code(), Some(0));
}
fn collect_stdout<T: Clone + Debug>(set: &mut ProcessSet<T>, want: usize) -> Vec<u8> {
let mut out = Vec::new();
while out.len() < want {
match set.wait_next().unwrap() {
Some((_, PsetEvent::Stdout(data))) => out.extend_from_slice(&data),
other => panic!("expected stdout, got {other:?}"),
}
}
out
}
#[test]
fn input_larger_than_a_pipe_buffer_is_fed_in_whole() {
const SIZE: usize = 1024 * 1024;
let input: Vec<u8> = (0..SIZE).map(|i| (i % 251) as u8).collect();
let mut set = pset::create::<&str>().unwrap();
let (_id, stdin) = set.spawn("cat", sh_stdin("cat")).unwrap();
let mut stdin = stdin.unwrap();
let fed = input.clone();
let writer = thread::spawn(move || {
stdin.write_all(&fed).unwrap();
});
let seen = drain(&mut set);
writer.join().expect("the writer finished");
assert_eq!(seen["cat"].stdout.len(), SIZE);
assert_eq!(seen["cat"].stdout, input, "the bytes came back unchanged");
assert_eq!(seen["cat"].code(), Some(0));
}
#[test]
fn each_process_is_fed_through_its_own_stdin() {
let mut set = pset::create::<Tag>().unwrap();
let (_id1, stdin1) = set.spawn(Tag::Cmd1, sh_stdin("cat")).unwrap();
let (_id2, stdin2) = set.spawn(Tag::Cmd2, sh_stdin("cat")).unwrap();
let mut stdin1 = stdin1.unwrap();
let mut stdin2 = stdin2.unwrap();
stdin2.write_all(b"for two\n").unwrap();
stdin1.write_all(b"for one\n").unwrap();
drop(stdin2);
drop(stdin1);
let seen = drain(&mut set);
assert_eq!(seen[&Tag::Cmd1].out(), "for one\n");
assert_eq!(seen[&Tag::Cmd2].out(), "for two\n");
}
#[test]
fn feeding_a_process_that_has_already_exited_fails() {
let mut set = pset::create::<&str>().unwrap();
let (_id, stdin) = set.spawn("quitter", sh_stdin("exit 0")).unwrap();
let mut stdin = stdin.unwrap();
let seen = drain(&mut set);
assert_eq!(seen["quitter"].code(), Some(0));
let err = stdin.write_all(b"too late\n").unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::BrokenPipe);
}
#[test]
fn dropping_the_set_leaves_no_child_behind() {
let mut set = pset::create::<&str>().unwrap();
set.spawn("sleeper", sh_piped("echo $$; exec sleep 300"))
.unwrap();
let pid: i32 = match set.wait_next().unwrap() {
Some(("sleeper", PsetEvent::Stdout(data))) => String::from_utf8_lossy(&data)
.trim()
.parse()
.expect("a pid"),
other => panic!("expected the pid on stdout, got {other:?}"),
};
assert!(pid_exists(pid), "the sleeper runs until the set goes away");
drop(set);
assert!(!pid_exists(pid), "pid {pid} outlived the set");
}
fn pid_exists(pid: i32) -> bool {
Command::new("kill")
.arg("-0")
.arg(pid.to_string())
.stderr(Stdio::null())
.status()
.expect("kill")
.success()
}
#[test]
fn a_process_with_nothing_piped_still_reports_its_exit() {
let mut set = pset::create::<&str>().unwrap();
let mut cmd = sh("echo out; echo err >&2; exit 3");
cmd.stdout(Stdio::null()).stderr(Stdio::null());
set.spawn("quiet", cmd).unwrap();
assert_eq!(set.len(), 1);
let event = set.wait_next().unwrap();
let Some(("quiet", PsetEvent::ProcessExited(status))) = event else {
panic!("expected the exit and only the exit, got {event:?}");
};
assert_eq!(status.code(), Some(3));
assert!(set.wait_next().unwrap().is_none());
assert!(set.is_empty());
}
#[test]
fn an_unpiped_process_is_not_held_back_by_a_grandchild() {
let script = "sleep 2 & exit 0";
let mut piped = pset::create::<&str>().unwrap();
piped.spawn("piped", sh_piped(script)).unwrap();
assert!(
piped
.wait_next_timeout(Duration::from_millis(100))
.unwrap()
.is_none(),
"the grandchild still holds the piped stdout"
);
assert_eq!(piped.len(), 1, "still waiting for the pipe to close");
let mut unpiped = pset::create::<&str>().unwrap();
let mut cmd = sh(script);
cmd.stdout(Stdio::null()).stderr(Stdio::null());
unpiped.spawn("unpiped", cmd).unwrap();
let event = unpiped
.wait_next_timeout(Duration::from_secs(1))
.unwrap()
.expect("the exit arrives without waiting for the grandchild");
let ("unpiped", PsetEvent::ProcessExited(status)) = event else {
panic!("expected the exit, got {event:?}");
};
assert_eq!(status.code(), Some(0));
assert!(unpiped.is_empty());
}
#[test]
fn piped_and_unpiped_processes_share_a_set() {
let mut set = pset::create::<Tag>().unwrap();
set.spawn(Tag::Cmd1, sh_piped("echo talkative")).unwrap();
let mut quiet = sh("exit 5");
quiet.stdout(Stdio::null()).stderr(Stdio::null());
set.spawn(Tag::Cmd2, quiet).unwrap();
assert_eq!(set.len(), 2);
let seen = drain(&mut set);
assert_eq!(seen.len(), 2, "both processes reported something");
assert_eq!(seen[&Tag::Cmd1].out(), "talkative\n");
assert_eq!(seen[&Tag::Cmd1].code(), Some(0));
assert!(seen[&Tag::Cmd2].stdout.is_empty());
assert!(seen[&Tag::Cmd2].stderr.is_empty());
assert_eq!(seen[&Tag::Cmd2].code(), Some(5));
}