use std::time::{Duration, Instant};
#[cfg(windows)]
use std::env;
#[cfg(windows)]
use std::fs;
#[cfg(windows)]
use std::io::{BufRead, BufReader, Read, Write};
#[cfg(windows)]
use std::path::PathBuf;
#[cfg(any(windows, unix))]
use std::process::Command;
#[cfg(any(windows, target_os = "linux"))]
use std::process::Stdio;
use std::thread;
#[cfg(target_os = "linux")]
use std::{
ffi::{CString, OsString},
os::unix::ffi::{OsStrExt, OsStringExt},
};
const CHILD_EXIT_WAIT: Duration = Duration::from_secs(30);
use running_process::{
run_command, run_command_bounded, CommandSpec, NativeProcess, ProcessConfig, ProcessError,
ReadStatus, StderrMode, StdinMode, StreamKind,
};
#[cfg(unix)]
use running_process::{run_std_command_bounded_with_options, BoundedRunOptions};
fn stdio_scripted() -> String {
let exe = std::env::current_exe().expect("current test executable");
let profile = exe
.parent()
.and_then(std::path::Path::parent)
.expect("test executable should live in <profile>/deps");
profile
.join(format!(
"testbin-stdio-scripted{}",
std::env::consts::EXE_SUFFIX
))
.to_string_lossy()
.into_owned()
}
fn config(
command: CommandSpec,
capture: bool,
stdin_mode: StdinMode,
nice: Option<i32>,
) -> ProcessConfig {
ProcessConfig {
command,
cwd: None,
env: None,
capture,
stderr_mode: StderrMode::Stdout,
creationflags: None,
create_process_group: false,
stdin_mode,
nice,
address_space_limit_bytes: None,
}
}
#[test]
fn captures_stderr_in_stdout_by_default() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('out'); print('err', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert!(process.captured_stdout().iter().any(|line| line == b"out"));
assert!(process.captured_stdout().iter().any(|line| line == b"err"));
assert!(process.captured_stderr().is_empty());
}
#[test]
fn run_command_returns_raw_output_and_exit_code() {
let output = run_command(
ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; sys.stdout.buffer.write(b'out\\n'); sys.stderr.buffer.write(b'err\\n'); sys.exit(7)"
.into(),
]),
false,
StdinMode::Null,
None,
)
},
Some(CHILD_EXIT_WAIT),
)
.unwrap();
assert_eq!(output.exit_code, 7);
assert_eq!(output.stdout, b"out\n");
assert_eq!(output.stderr, b"err\n");
}
#[test]
fn run_command_drains_stdout_and_stderr_concurrently() {
let output = run_command(
ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; sys.stderr.buffer.write(b'e' * 262144); sys.stderr.flush(); sys.stdout.buffer.write(b'ok\\n'); sys.stdout.flush()"
.into(),
]),
false,
StdinMode::Null,
None,
)
},
Some(CHILD_EXIT_WAIT),
)
.unwrap();
assert_eq!(output.exit_code, 0);
assert_eq!(output.stdout, b"ok\n");
assert_eq!(output.stderr.len(), 262144);
assert!(output.stderr.iter().all(|byte| *byte == b'e'));
}
#[test]
fn run_command_timeout_kills_child_and_returns_timeout() {
let started = Instant::now();
let result = run_command(
config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(10)".into(),
]),
false,
StdinMode::Null,
None,
),
Some(Duration::from_millis(100)),
);
assert!(matches!(result, Err(ProcessError::Timeout)));
assert!(
started.elapsed() < Duration::from_secs(5),
"timeout path did not kill promptly"
);
}
#[test]
fn bounded_run_stops_allocating_after_output_limit() {
let result = run_command_bounded(
ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; sys.stdout.buffer.write(b'x' * 4194304); sys.stdout.flush()"
.into(),
]),
false,
StdinMode::Null,
None,
)
},
Some(CHILD_EXIT_WAIT),
1024,
);
assert!(matches!(
result,
Err(ProcessError::OutputLimitExceeded { limit: 1024 })
));
}
#[cfg(target_os = "linux")]
#[test]
fn bounded_std_command_preserves_non_utf8_process_inputs() {
let temp = tempfile::tempdir().unwrap();
let program = temp.path().join(OsString::from_vec(b"shell-\xff".to_vec()));
std::os::unix::fs::symlink("/bin/sh", &program).unwrap();
let mut command = Command::new(program);
command
.args(["-c", "printf %s \"$RP_BYTES\"; printf %s \"$1\"", "sh"])
.arg(OsString::from_vec(b"arg-\xfe".to_vec()))
.env("RP_BYTES", OsString::from_vec(b"environment-\xff".to_vec()));
let output =
running_process::run_std_command_bounded(command, Some(CHILD_EXIT_WAIT), 4096).unwrap();
assert_eq!(output.stdout, b"environment-\xffarg-\xfe");
}
#[cfg(target_os = "linux")]
#[test]
fn bounded_std_command_options_preserve_nonzero_timeout_and_output_limit_results() {
let mut nonzero = Command::new("sh");
nonzero.args(["-c", "printf stdout; printf stderr >&2; exit 7"]);
let output = run_std_command_bounded_with_options(
nonzero,
Some(CHILD_EXIT_WAIT),
4096,
BoundedRunOptions::default(),
)
.unwrap();
assert_eq!(output.exit_code, 7);
assert_eq!(output.stdout, b"stdout");
assert_eq!(output.stderr, b"stderr");
let started = Instant::now();
let mut slow = Command::new("sh");
slow.args(["-c", "sleep 30"]);
assert!(matches!(
run_std_command_bounded_with_options(
slow,
Some(Duration::from_millis(100)),
4096,
BoundedRunOptions::default(),
),
Err(ProcessError::Timeout)
));
assert!(
started.elapsed() < Duration::from_secs(2),
"bounded timeout did not request cleanup promptly"
);
let mut overflowing = Command::new("sh");
overflowing.args(["-c", "yes bounded-output"]);
assert!(matches!(
run_std_command_bounded_with_options(
overflowing,
Some(CHILD_EXIT_WAIT),
64,
BoundedRunOptions::default(),
),
Err(ProcessError::OutputLimitExceeded { limit: 64 })
));
}
#[cfg(target_os = "linux")]
#[test]
fn bounded_std_command_owner_death_composes_with_user_pre_exec_and_group_policy() {
use std::os::unix::process::CommandExt;
let working_directory = tempfile::tempdir().unwrap();
let working_directory_bytes = CString::new(working_directory.path().as_os_str().as_bytes())
.expect("temporary paths do not contain NUL");
let expected_directory = working_directory.path().display().to_string();
let mut command = Command::new("sh");
command.args([
"-c",
"pgid=$(ps -o pgid= -p $$ | tr -d '[:space:]'); test \"$pgid\" = \"$$\" && pwd",
]);
unsafe {
command.pre_exec(move || {
if libc::chdir(working_directory_bytes.as_ptr()) == -1 {
return Err(std::io::Error::last_os_error());
}
Ok(())
});
}
let output = run_std_command_bounded_with_options(
command,
Some(CHILD_EXIT_WAIT),
4096,
BoundedRunOptions::default()
.kill_when_owner_dies(true)
.nice(None),
)
.unwrap();
assert_eq!(output.exit_code, 0);
assert_eq!(
std::str::from_utf8(&output.stdout).unwrap().trim_end(),
expected_directory
);
}
#[cfg(target_os = "linux")]
#[test]
fn helper_force_killed_bounded_owner_reaps_child() {
if std::env::var("RUNNING_PROCESS_BOUNDED_OWNER_DEATH_HELPER")
.ok()
.as_deref()
!= Some("1")
{
return;
}
let mut command = Command::new("sh");
command.args([
"-c",
"echo $$ > \"$RUNNING_PROCESS_BOUNDED_OWNER_DEATH_PID_FILE\"; exec sleep 30",
]);
let result = run_std_command_bounded_with_options(
command,
Some(Duration::from_secs(120)),
4096,
BoundedRunOptions::default()
.kill_when_owner_dies(true)
.nice(None),
);
assert!(
matches!(result, Err(ProcessError::Timeout)),
"helper must run until its bounded timeout unless its owner dies: {result:?}"
);
}
#[cfg(target_os = "linux")]
#[test]
fn bounded_std_command_owner_death_kills_child_after_owner_is_forced_down() {
let temporary = tempfile::tempdir().unwrap();
let pid_file = temporary.path().join("bounded-owner-death.pid");
let mut owner = Command::new(std::env::current_exe().unwrap())
.arg("--exact")
.arg(helper_test_name(
"helper_force_killed_bounded_owner_reaps_child",
))
.arg("--nocapture")
.env("RUNNING_PROCESS_BOUNDED_OWNER_DEATH_HELPER", "1")
.env("RUNNING_PROCESS_BOUNDED_OWNER_DEATH_PID_FILE", &pid_file)
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap();
let startup_deadline = Instant::now() + Duration::from_secs(5);
let child_pid = loop {
if let Ok(pid) = std::fs::read_to_string(&pid_file) {
break pid
.trim()
.parse::<u32>()
.expect("helper wrote a numeric child pid");
}
if let Some(status) = owner.try_wait().unwrap() {
panic!("bounded owner exited before reporting child pid: {status}");
}
assert!(
Instant::now() < startup_deadline,
"bounded owner did not report a child pid"
);
thread::sleep(Duration::from_millis(20));
};
owner.kill().unwrap();
owner.wait().unwrap();
let death_deadline = Instant::now() + Duration::from_secs(5);
while linux_pid_exists(child_pid) {
assert!(
Instant::now() < death_deadline,
"bounded child {child_pid} survived owner death"
);
thread::sleep(Duration::from_millis(20));
}
}
#[cfg(target_os = "linux")]
fn linux_pid_exists(pid: u32) -> bool {
match std::fs::read_to_string(format!("/proc/{pid}/stat")) {
Ok(stat) => {
let state = stat
.rsplit_once(") ")
.and_then(|(_, tail)| tail.as_bytes().first().copied());
state != Some(b'Z')
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => false,
Err(_) => {
(unsafe { libc::kill(pid as libc::pid_t, 0) == 0 })
|| std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
}
}
}
#[cfg(target_os = "linux")]
#[test]
fn bounded_run_cancels_readers_held_by_escaped_descendant() {
let started = Instant::now();
let result = run_command_bounded(
config(
CommandSpec::Argv(vec![
"sh".into(),
"-c".into(),
"setsid sh -c 'sleep 3' & sleep 30".into(),
]),
false,
StdinMode::Null,
None,
),
Some(Duration::from_millis(100)),
4096,
);
assert!(matches!(result, Err(ProcessError::Timeout)));
assert!(
started.elapsed() < Duration::from_secs(1),
"escaped descendant kept capture readers alive for {:?}",
started.elapsed()
);
}
#[test]
fn captures_stdout_and_stderr_separately_when_requested() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('out'); print('err', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert!(process.captured_stdout().iter().any(|line| line == b"out"));
assert!(process.captured_stderr().iter().any(|line| line == b"err"));
}
#[test]
fn stream_reads_report_timeout_then_eof() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.2); print('ready')".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
assert_eq!(
process.read_stream(StreamKind::Stdout, Some(Duration::from_millis(10))),
ReadStatus::Timeout
);
assert!(matches!(
process.read_stream(StreamKind::Stdout, Some(Duration::from_secs(30))),
ReadStatus::Line(line) if line == b"ready"
));
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(
process.read_stream(StreamKind::Stdout, Some(Duration::from_millis(10))),
ReadStatus::Eof
);
}
#[test]
fn normalizes_crlf_and_preserves_invalid_bytes() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; sys.stdout.buffer.write(b'bad:\\xff\\r\\nnext\\rthird\\n'); sys.stdout.flush()"
.into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(
process.captured_stdout(),
vec![b"bad:\xff".to_vec(), b"next\rthird".to_vec()]
);
}
#[test]
fn raw_stream_drain_preserves_delimiters_non_utf8_and_eof() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
stdio_scripted(),
"outhex:63726c660d0a6c660a7461696cff".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
assert_eq!(process.wait(Some(CHILD_EXIT_WAIT)).unwrap(), 0);
assert_eq!(
process.drain_stream_raw(StreamKind::Stdout),
b"crlf\r\nlf\ntail\xff"
);
assert!(process.drain_stream_raw(StreamKind::Stdout).is_empty());
assert_eq!(
process.captured_stdout(),
vec![b"crlf".to_vec(), b"lf".to_vec(), b"tail\xff".to_vec()]
);
}
#[test]
fn raw_stream_drain_is_incremental_before_and_after_eof() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
stdio_scripted(),
"outhex:66697273740d0a".into(),
"sleep-ms:250".into(),
"outhex:7365636f6e64ff".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
while !process.has_pending_stream(StreamKind::Stdout) {
assert!(
Instant::now() < deadline,
"first raw chunk was not observed"
);
std::thread::sleep(Duration::from_millis(10));
}
assert_eq!(process.drain_stream_raw(StreamKind::Stdout), b"first\r\n");
assert_eq!(process.wait(Some(CHILD_EXIT_WAIT)).unwrap(), 0);
assert_eq!(process.drain_stream_raw(StreamKind::Stdout), b"second\xff");
assert!(process.drain_stream_raw(StreamKind::Stdout).is_empty());
}
#[test]
fn raw_stream_drain_keeps_pipes_independent_under_high_volume() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
stdio_scripted(),
"repeat-outhex:262144:6f".into(),
"repeat-errhex:262144:65".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
assert_eq!(process.wait(Some(CHILD_EXIT_WAIT)).unwrap(), 0);
let stdout = process.drain_stream_raw(StreamKind::Stdout);
let stderr = process.drain_stream_raw(StreamKind::Stderr);
assert_eq!(stdout.len(), 262144);
assert!(stdout.iter().all(|byte| *byte == b'o'));
assert_eq!(stderr.len(), 262144);
assert!(stderr.iter().all(|byte| *byte == b'e'));
}
#[test]
fn supports_piped_stdin_filter_execution() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; data = sys.stdin.buffer.read(); sys.stdout.buffer.write(data[::-1])"
.into(),
]),
true,
StdinMode::Piped,
None,
)
});
process.start().unwrap();
process.write_stdin(b"abc").unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(process.captured_stdout(), vec![b"cba".to_vec()]);
}
#[test]
fn captured_output_can_be_cleared_to_release_memory() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('alpha'); print('beta', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(process.captured_stream_bytes(StreamKind::Stdout), 5);
assert_eq!(process.captured_stream_bytes(StreamKind::Stderr), 4);
assert_eq!(process.captured_combined_bytes(), 9);
assert_eq!(process.clear_captured_stream(StreamKind::Stdout), 5);
assert!(process.captured_stdout().is_empty());
assert_eq!(process.captured_stream_bytes(StreamKind::Stdout), 0);
assert_eq!(process.clear_captured_combined(), 9);
assert!(process.captured_combined().is_empty());
assert_eq!(process.captured_combined_bytes(), 0);
}
#[test]
#[cfg(not(windows))]
fn applies_positive_nice_before_exec() {
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import os; print(os.nice(0))".into(),
]),
true,
StdinMode::Inherit,
Some(5),
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
let observed = String::from_utf8(process.captured_stdout()[0].clone())
.unwrap()
.parse::<i32>()
.unwrap();
assert!(observed >= 5);
}
#[cfg(unix)]
#[test]
fn bounded_std_command_applies_positive_nice_before_exec() {
let mut command = Command::new("python");
command.args(["-c", "import os; print(os.nice(0))"]);
let output = run_std_command_bounded_with_options(
command,
Some(CHILD_EXIT_WAIT),
4096,
BoundedRunOptions::default()
.kill_when_owner_dies(false)
.nice(Some(5)),
)
.unwrap();
assert_eq!(output.exit_code, 0);
let observed = String::from_utf8(output.stdout)
.unwrap()
.trim()
.parse::<i32>()
.unwrap();
assert!(observed >= 5);
}
#[test]
fn start_twice_returns_already_started() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
assert!(matches!(process.start(), Err(ProcessError::AlreadyStarted)));
let _ = process.kill();
}
#[test]
fn write_stdin_before_start_returns_not_running() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "pass".into()]),
false,
StdinMode::Piped,
None,
));
assert!(matches!(
process.write_stdin(b"hello"),
Err(ProcessError::NotRunning)
));
}
#[test]
fn write_stdin_without_piped_returns_stdin_unavailable() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
assert!(matches!(
process.write_stdin(b"hello"),
Err(ProcessError::StdinUnavailable)
));
let _ = process.kill();
}
#[test]
fn kill_before_start_returns_not_running() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "pass".into()]),
false,
StdinMode::Inherit,
None,
));
assert!(matches!(process.kill(), Err(ProcessError::NotRunning)));
}
#[test]
fn wait_before_start_returns_not_running() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "pass".into()]),
false,
StdinMode::Inherit,
None,
));
assert!(matches!(
process.wait(Some(Duration::from_secs(1))),
Err(ProcessError::NotRunning)
));
}
#[test]
fn wait_timeout_returns_timeout_error() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(10)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
assert!(matches!(
process.wait(Some(Duration::from_millis(100))),
Err(ProcessError::Timeout)
));
let _ = process.kill();
}
#[test]
fn read_combined_returns_events_from_both_streams() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('out'); sys.stdout.flush(); print('err', file=sys.stderr); sys.stderr.flush()".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
let mut events = Vec::new();
loop {
match process.read_combined(Some(Duration::from_millis(100))) {
ReadStatus::Line(event) => events.push(event),
ReadStatus::Eof => break,
ReadStatus::Timeout => break,
}
}
assert!(events
.iter()
.any(|e| e.stream == StreamKind::Stdout && e.line == b"out"));
assert!(events
.iter()
.any(|e| e.stream == StreamKind::Stderr && e.line == b"err"));
}
#[test]
fn drain_combined_returns_all_pending() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('a'); print('b', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
std::thread::sleep(Duration::from_millis(50));
let events = process.drain_combined();
assert!(events.len() >= 2);
}
#[test]
fn has_pending_combined_reports_correctly() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "print('hello')".into()]),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
let deadline = Instant::now() + CHILD_EXIT_WAIT;
while !process.has_pending_combined() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(10));
}
assert!(
process.has_pending_combined(),
"combined output never became pending"
);
process.drain_combined();
assert!(!process.has_pending_combined());
}
#[test]
fn captured_combined_includes_both_streams() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('out'); print('err', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
let combined = process.captured_combined();
assert!(combined
.iter()
.any(|e| e.stream == StreamKind::Stdout && e.line == b"out"));
assert!(combined
.iter()
.any(|e| e.stream == StreamKind::Stderr && e.line == b"err"));
}
#[test]
fn captured_combined_bytes_and_clear() {
let process = NativeProcess::new(ProcessConfig {
stderr_mode: StderrMode::Pipe,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; print('ab'); print('cd', file=sys.stderr)".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(process.captured_combined_bytes(), 4);
assert_eq!(process.clear_captured_combined(), 4);
assert_eq!(process.captured_combined_bytes(), 0);
assert!(process.captured_combined().is_empty());
}
#[test]
fn shell_command_captures_output() {
let process = NativeProcess::new(config(
CommandSpec::Shell("echo shell-works".into()),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
let stdout = process.captured_stdout();
assert!(
stdout.iter().any(|line| {
let text = String::from_utf8_lossy(line);
text.contains("shell-works")
}),
"expected 'shell-works' in output, got: {:?}",
stdout,
);
}
#[test]
fn custom_cwd_is_respected() {
let tmp = std::env::temp_dir();
let process = NativeProcess::new(ProcessConfig {
cwd: Some(tmp.clone()),
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import os; print(os.getcwd())".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
let output = String::from_utf8(process.captured_stdout()[0].clone()).unwrap();
let expected = std::fs::canonicalize(&tmp).unwrap_or(tmp);
let actual = std::fs::canonicalize(output.trim()).unwrap_or_else(|_| output.trim().into());
assert_eq!(actual, expected);
}
#[test]
fn custom_env_is_applied() {
let mut env_vars = vec![("RP_TEST_VAR".into(), "hello_coverage".into())];
if let Ok(path) = std::env::var("PATH") {
env_vars.push(("PATH".into(), path));
}
#[cfg(windows)]
if let Ok(root) = std::env::var("SystemRoot") {
env_vars.push(("SystemRoot".into(), root));
}
let process = NativeProcess::new(ProcessConfig {
env: Some(env_vars),
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import os; print(os.environ.get('RP_TEST_VAR', 'MISSING'))".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(process.captured_stdout(), vec![b"hello_coverage".to_vec()]);
}
#[test]
fn stdin_null_produces_empty_input() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import sys; data=sys.stdin.buffer.read(); print(len(data))".into(),
]),
true,
StdinMode::Null,
None,
));
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(process.captured_stdout(), vec![b"0".to_vec()]);
}
#[test]
fn poll_returns_none_while_running_then_exit_code() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.3)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let status = process.poll().unwrap();
assert!(status.is_none(), "expected None, got {:?}", status);
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
let status = process.poll().unwrap();
assert_eq!(status, Some(0));
}
#[test]
fn close_kills_running_process() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
process.close().unwrap();
}
#[test]
fn close_on_already_finished_is_noop() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "pass".into()]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
process.close().unwrap();
}
#[test]
fn terminate_kills_running_process() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
));
process.start().unwrap();
process.terminate().unwrap();
}
#[cfg(windows)]
#[test]
fn kill_cancels_capture_io_when_grandchild_orphans_pipe() {
let script = "\
import os, subprocess, sys, time;\
print('PARENT_PID=' + str(os.getpid()), flush=True);\
gc = subprocess.Popen([sys.executable, '-c', 'import time; time.sleep(60)']);\
print('GRANDCHILD_PID=' + str(gc.pid), flush=True);\
time.sleep(60)";
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), script.into()]),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let mut grandchild_pid: Option<u32> = None;
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
match process.read_combined(Some(Duration::from_millis(200))) {
ReadStatus::Line(event) => {
let line = String::from_utf8_lossy(&event.line).into_owned();
if let Some(rest) = line.strip_prefix("GRANDCHILD_PID=") {
grandchild_pid = rest.trim().parse::<u32>().ok();
break;
}
}
ReadStatus::Timeout => continue,
ReadStatus::Eof => panic!("parent exited before announcing grandchild"),
}
}
let grandchild_pid = grandchild_pid.expect("did not observe GRANDCHILD_PID line");
let prior = env::var_os("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS");
env::set_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS", "5000");
let kill_start = Instant::now();
let kill_result = process.kill();
let kill_elapsed = kill_start.elapsed();
match prior {
Some(v) => env::set_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS", v),
None => env::remove_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS"),
}
kill_result.expect("kill() returned an error");
assert!(
kill_elapsed < Duration::from_secs(1),
"kill() took {kill_elapsed:?} with 5 s safety-net deadline; \
CancelIoEx fast path is not interrupting the reader thread",
);
let _ = Command::new("taskkill")
.args(["/F", "/T", "/PID", &grandchild_pid.to_string()])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status();
}
#[test]
fn kill_returns_when_grandchild_inherits_stdout_pipe() {
let marker = format!(
"running-process-619-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("system clock before epoch")
.as_nanos()
);
let grandchild_code = format!("import time; time.sleep(60) # {marker}");
let script = format!(
"\
import os, subprocess, sys, time;\
print('PARENT_PID=' + str(os.getpid()), flush=True);\
gc = subprocess.Popen([sys.executable, '-c', {grandchild_code:?}]);\
print('GRANDCHILD_PID=' + str(gc.pid), flush=True);\
time.sleep(60)"
);
let process_config = config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), script]),
true,
StdinMode::Inherit,
None,
);
#[cfg(unix)]
let process_config = ProcessConfig {
create_process_group: true,
..process_config
};
let process = NativeProcess::new(process_config);
process.start().unwrap();
let mut grandchild_pid: Option<u32> = None;
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
match process.read_combined(Some(Duration::from_millis(200))) {
ReadStatus::Line(event) => {
let line = String::from_utf8_lossy(&event.line).into_owned();
if let Some(rest) = line.strip_prefix("GRANDCHILD_PID=") {
grandchild_pid = rest.trim().parse::<u32>().ok();
break;
}
}
ReadStatus::Timeout => continue,
ReadStatus::Eof => panic!("parent exited before announcing grandchild"),
}
}
let grandchild_pid = grandchild_pid.expect("did not observe GRANDCHILD_PID line");
#[cfg(unix)]
let is_original_grandchild_running = || {
Command::new("ps")
.args([
"-ww",
"-o",
"stat=",
"-o",
"command=",
"-p",
&grandchild_pid.to_string(),
])
.output()
.ok()
.filter(|output| output.status.success())
.map(|output| {
let state_and_command = String::from_utf8_lossy(&output.stdout);
state_and_command.contains(&marker)
&& !state_and_command.trim_start().starts_with('Z')
})
.unwrap_or(false)
};
#[cfg(unix)]
assert!(
is_original_grandchild_running(),
"grandchild identity marker was not observable before kill"
);
let kill_start = Instant::now();
process.kill().expect("kill() returned an error");
let kill_elapsed = kill_start.elapsed();
assert!(
kill_elapsed < Duration::from_secs(5),
"kill() blocked for {kill_elapsed:?}; expected bounded return after grandchild orphan",
);
#[cfg(windows)]
{
let _ = Command::new("taskkill")
.args(["/F", "/T", "/PID", &grandchild_pid.to_string()])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status();
}
#[cfg(unix)]
{
let gone_deadline = Instant::now() + Duration::from_secs(2);
while is_original_grandchild_running() && Instant::now() < gone_deadline {
thread::sleep(Duration::from_millis(10));
}
let survived = is_original_grandchild_running();
if survived {
unsafe {
libc::kill(grandchild_pid as i32, libc::SIGKILL);
}
}
assert!(
!survived,
"grandchild {grandchild_pid} survived kill() on its owned process group"
);
}
}
#[test]
fn wait_returns_when_grandchild_orphans_pipe_on_natural_exit() {
let script = "\
import os, subprocess, sys;\
print('PARENT_PID=' + str(os.getpid()), flush=True);\
gc = subprocess.Popen([sys.executable, '-c', 'import time; time.sleep(60)']);\
print('GRANDCHILD_PID=' + str(gc.pid), flush=True);\
sys.stdout.flush()";
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), script.into()]),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let mut grandchild_pid: Option<u32> = None;
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
match process.read_combined(Some(Duration::from_millis(200))) {
ReadStatus::Line(event) => {
let line = String::from_utf8_lossy(&event.line).into_owned();
if let Some(rest) = line.strip_prefix("GRANDCHILD_PID=") {
grandchild_pid = rest.trim().parse::<u32>().ok();
break;
}
}
ReadStatus::Timeout => continue,
ReadStatus::Eof => break,
}
}
let grandchild_pid = grandchild_pid.expect("did not observe GRANDCHILD_PID line");
let wait_start = Instant::now();
let result = process.wait(Some(Duration::from_secs(30)));
let wait_elapsed = wait_start.elapsed();
#[cfg(windows)]
{
let _ = Command::new("taskkill")
.args(["/F", "/T", "/PID", &grandchild_pid.to_string()])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status();
}
#[cfg(not(windows))]
unsafe {
libc::kill(grandchild_pid as i32, libc::SIGKILL);
}
assert!(
result.is_ok(),
"wait() returned {result:?} instead of the child's exit code",
);
assert!(
wait_elapsed < Duration::from_secs(10),
"wait() blocked for {wait_elapsed:?}; the natural-exit capture drain \
is not bounded (grandchild-orphaned pipe wedge, issue #590)",
);
}
#[test]
fn natural_exit_still_captures_late_grandchild_output() {
let script = "\
import subprocess, sys;\
subprocess.Popen([sys.executable, '-c', \
\"import time,sys; time.sleep(0.3); print('LATE_GRANDCHILD_OUTPUT', flush=True); sys.stdout.flush()\"]);\
sys.exit(0)";
let prior = std::env::var_os("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS");
std::env::set_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS", "15000");
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), script.into()]),
true,
StdinMode::Inherit,
None,
));
process.start().unwrap();
let mut saw_late_output = false;
let deadline = Instant::now() + Duration::from_secs(12);
while Instant::now() < deadline {
match process.read_combined(Some(Duration::from_millis(200))) {
ReadStatus::Line(event) => {
if String::from_utf8_lossy(&event.line).contains("LATE_GRANDCHILD_OUTPUT") {
saw_late_output = true;
break;
}
}
ReadStatus::Timeout => continue,
ReadStatus::Eof => break,
}
}
let _ = process.wait(Some(Duration::from_secs(5)));
match prior {
Some(v) => std::env::set_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS", v),
None => std::env::remove_var("RUNNING_PROCESS_KILL_DRAIN_TIMEOUT_MS"),
}
assert!(
saw_late_output,
"late grandchild output written after the direct child exited was \
dropped; the natural-exit drain cancelled the reader too eagerly",
);
}
#[test]
fn pid_returns_some_after_start() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
));
assert!(process.pid().is_none());
process.start().unwrap();
assert!(process.pid().is_some());
let _ = process.kill();
}
#[test]
#[cfg(not(windows))]
fn create_process_group_sets_new_pgid() {
let process = NativeProcess::new(ProcessConfig {
create_process_group: true,
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import os; print(os.getpgid(0) == os.getpid())".into(),
]),
true,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
let code = process.wait(Some(CHILD_EXIT_WAIT)).unwrap();
assert_eq!(code, 0);
assert_eq!(process.captured_stdout(), vec![b"True".to_vec()]);
}
fn helper_test_name(name: &str) -> String {
match module_path!().split_once("::") {
Some((_crate_root, module)) => format!("{module}::{name}"),
None => name.to_owned(),
}
}
#[test]
fn helper_test_name_is_module_qualified_for_the_category_target() {
assert_eq!(
helper_test_name("helper_x"),
"process_core_test::helper_x",
"libtest matches `--exact` against the module-qualified name"
);
}
#[test]
#[cfg(windows)]
fn helper_force_killed_parent_reaps_native_child() {
if env::var("RUNNING_PROCESS_CORE_HELPER").ok().as_deref() != Some("1") {
return;
}
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
println!("CHILD_PID={}", process.pid().unwrap());
std::io::stdout().flush().unwrap();
thread::sleep(Duration::from_secs(30));
}
#[test]
#[cfg(windows)]
fn force_killed_parent_reaps_native_child_on_windows() {
let current_exe = env::current_exe().unwrap();
let mut owner = Command::new(current_exe)
.arg("--exact")
.arg(helper_test_name(
"helper_force_killed_parent_reaps_native_child",
))
.arg("--nocapture")
.env("RUNNING_PROCESS_CORE_HELPER", "1")
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let child_pid = {
let stdout = owner.stdout.take().unwrap();
let mut reader = BufReader::new(stdout);
let mut line = String::new();
loop {
line.clear();
let read = reader.read_line(&mut line).unwrap();
assert!(read != 0, "helper exited before reporting child pid");
if line.starts_with("CHILD_PID=") {
break line
.trim()
.trim_start_matches("CHILD_PID=")
.parse::<u32>()
.unwrap();
}
}
};
owner.kill().unwrap();
owner.wait().unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
if !pid_exists(child_pid) {
break;
}
assert!(
std::time::Instant::now() < deadline,
"child {child_pid} survived owner death"
);
thread::sleep(Duration::from_millis(50));
}
}
#[test]
#[cfg(windows)]
fn helper_force_killed_parent_logs_native_child() {
if env::var("RUNNING_PROCESS_CORE_HELPER_LOGGED")
.ok()
.as_deref()
!= Some("1")
{
return;
}
let process = NativeProcess::new(ProcessConfig {
..config(
CommandSpec::Argv(vec![
"python".into(),
"-c".into(),
"import time; time.sleep(0.1)".into(),
]),
false,
StdinMode::Inherit,
None,
)
});
process.start().unwrap();
println!("OWNER_READY");
std::io::stdout().flush().unwrap();
thread::sleep(Duration::from_secs(30));
}
#[test]
#[cfg(windows)]
fn repeated_force_killed_parents_leave_no_logged_native_children_on_windows() {
let current_exe = env::current_exe().unwrap();
let log_path = unique_pid_log_path();
let owner_count = if std::env::consts::ARCH == "aarch64" {
4
} else {
6
};
let mut owners = Vec::new();
for _ in 0..owner_count {
let mut owner = Command::new(¤t_exe)
.arg("--exact")
.arg(helper_test_name(
"helper_force_killed_parent_logs_native_child",
))
.arg("--nocapture")
.env("RUNNING_PROCESS_CORE_HELPER_LOGGED", "1")
.env("RUNNING_PROCESS_CHILD_PID_LOG_PATH", &log_path)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
{
let stdout = owner.stdout.take().unwrap();
let mut reader = BufReader::new(stdout);
let mut line = String::new();
loop {
line.clear();
let read = reader.read_line(&mut line).unwrap();
if read == 0 {
let mut stderr = String::new();
if let Some(stderr_pipe) = owner.stderr.as_mut() {
let _ = stderr_pipe.read_to_string(&mut stderr);
}
panic!(
"helper exited before reporting readiness; stderr:\n{}",
stderr.trim_end()
);
}
if line.trim() == "OWNER_READY" {
break;
}
}
}
owners.push(owner);
}
for owner in &mut owners {
owner.kill().unwrap();
owner.wait().unwrap();
}
let child_pids = read_logged_pids(&log_path);
assert_eq!(child_pids.len(), owner_count);
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
let all_dead = child_pids.iter().all(|pid| !pid_exists(*pid));
if all_dead {
break;
}
assert!(
std::time::Instant::now() < deadline,
"some logged child pids survived owner death: {child_pids:?}"
);
thread::sleep(Duration::from_millis(50));
}
let _ = fs::remove_file(&log_path);
}
#[cfg(windows)]
fn unique_pid_log_path() -> PathBuf {
let suffix = format!(
"{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
);
env::temp_dir().join(format!("running-process-native-child-pids-{suffix}.log"))
}
#[cfg(windows)]
fn read_logged_pids(path: &PathBuf) -> Vec<u32> {
let content = fs::read_to_string(path).unwrap();
content
.lines()
.map(str::trim)
.filter(|line| !line.is_empty())
.map(|line| line.parse::<u32>().unwrap())
.collect()
}
#[cfg(windows)]
fn pid_exists(pid: u32) -> bool {
use winapi::um::handleapi::CloseHandle;
use winapi::um::processthreadsapi::{GetExitCodeProcess, OpenProcess};
use winapi::um::winnt::PROCESS_QUERY_LIMITED_INFORMATION;
const STILL_ACTIVE: u32 = 259;
let handle = unsafe { OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid) };
if handle.is_null() {
return false;
}
let mut exit_code = 0u32;
let ok = unsafe { GetExitCodeProcess(handle, &mut exit_code) } != 0;
unsafe {
CloseHandle(handle);
}
ok && exit_code == STILL_ACTIVE
}
#[test]
fn returncode_auto_updates_without_poll() {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["python".into(), "-c".into(), "print('hello')".into()]),
true,
StdinMode::Null,
None,
));
process.start().unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
while Instant::now() < deadline {
if process.returncode().is_some() {
break;
}
std::thread::sleep(Duration::from_millis(10));
}
assert!(
process.returncode().is_some(),
"returncode should auto-update via background waiter thread without calling poll()"
);
assert_eq!(process.returncode(), Some(0));
}
#[cfg(target_os = "linux")]
#[test]
fn a_one_millisecond_wait_finds_a_child_that_already_exited() {
fn is_zombie(pid: u32) -> bool {
std::fs::read_to_string(format!("/proc/{pid}/stat"))
.ok()
.and_then(|stat| {
let rest = stat.get(stat.rfind(')')? + 2..)?;
rest.chars().next()
})
== Some('Z')
}
for attempt in 0..30 {
let process = NativeProcess::new(config(
CommandSpec::Argv(vec!["/bin/sh".into(), "-c".into(), "exit 0".into()]),
false,
StdinMode::Inherit,
None,
));
process
.start()
.expect("start a child that exits immediately");
let pid = process.pid().expect("child pid");
let deadline = Instant::now() + Duration::from_secs(5);
while !is_zombie(pid) && process.returncode().is_none() {
assert!(Instant::now() < deadline, "child never exited");
std::thread::sleep(Duration::from_micros(200));
}
let code = process
.wait(Some(Duration::from_millis(1)))
.unwrap_or_else(|error| {
panic!("attempt {attempt}: a child that already exited was not found: {error:?}")
});
assert_eq!(code, 0);
}
}