use crate::fd::Fd;
use crate::fs::{mmap_madvise, path_exists, path_lstat_exists, readahead};
use crate::inotify::{InotifyEvent, decode_events};
use crate::io::{DrainState, SinkResult};
use crate::proc::{
ProcDir, clock_ticks_per_second, parse_proc_cmdline_bytes, parse_proc_status, path_lstat,
path_stat, path_uid, read_proc_cmdline_at, stat_at, uid, uid_at,
};
use crate::reactor::Reactor;
use crate::signal::{install_shutdown_flag, install_shutdown_flag_guard, shutdown_requested};
use crate::socket::{
StaleSocketPolicy, UnixConnectResult, UnixSocketAddr, UnixSocketBindOptions, bind, chmod,
connect,
};
use crate::spawn::{
CancelPolicy, ExitStatus, Process, ProcessGroup, SessionExitPolicy, SpawnBackend,
SpawnFdPolicy, SpawnOptions, SpawnOptionsBuilder, make_pty, orphan_session,
pty_session_enumerated, pty_window, reset_pty_enumeration, spawn_managed, spawn_start,
};
use std::fs::{File, remove_file};
use std::io::Write;
use std::os::unix::ffi::OsStringExt;
use std::os::unix::fs::{FileTypeExt, PermissionsExt};
use std::os::unix::io::{AsRawFd, RawFd};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
fn with_temp_readahead_file<T>(f: impl FnOnce(File, &std::path::Path) -> T) -> T {
let path = std::env::temp_dir().join(format!(
"coreshift_test_readahead_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let mut file = File::options()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&path)
.unwrap();
file.write_all(b"readahead test data").unwrap();
file.sync_all().unwrap();
let result = f(file, &path);
let _ = remove_file(&path);
result
}
fn temp_socket_path(name: &str) -> PathBuf {
std::env::temp_dir().join(format!(
"coreshift_test_{name}_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
))
}
fn bind_test_unix_listener(path: &Path, unlink_stale: bool) -> crate::socket::UnixListener {
bind(
UnixSocketAddr::Path(path),
UnixSocketBindOptions {
stale_socket_policy: if unlink_stale {
StaleSocketPolicy::UnlinkSocketOnly
} else {
StaleSocketPolicy::Preserve
},
mode: None,
},
)
.unwrap()
}
fn connect_test_unix_stream(addr: UnixSocketAddr<'_>) -> crate::socket::UnixStream {
match connect(addr).unwrap() {
UnixConnectResult::Connected(stream) => stream,
UnixConnectResult::InProgress(stream) => stream.finish_connect().unwrap(),
}
}
fn assert_readahead_result(result: Result<(), crate::CoreError>) {
match result {
Ok(()) => {}
Err(err) if err.raw_os_error() == Some(libc::ENOSYS) => {
eprintln!("skipping readahead test: unsupported on this target");
}
Err(err) => panic!("readahead failed unexpectedly: {err}"),
}
}
struct RawFdRef(RawFd);
static TEST_SHUTDOWN_FLAG_A: AtomicBool = AtomicBool::new(false);
static TEST_SHUTDOWN_FLAG_B: AtomicBool = AtomicBool::new(false);
static SIGNAL_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
impl AsRawFd for RawFdRef {
fn as_raw_fd(&self) -> RawFd {
self.0
}
}
#[test]
fn test_log_backend_selection() {
use crate::log::{LogBackend, LogLevel, Logger};
let logger_null = Logger::new(LogBackend::Null);
logger_null.log(LogLevel::Info, "Test", "This should be discarded");
let logger_stderr = Logger::new(LogBackend::Stderr);
logger_stderr.log(LogLevel::Info, "Test", "This should go to stderr");
}
#[test]
#[ignore]
fn test_signalfd_new() {
use crate::reactor::Reactor;
use crate::signal::{SIGUSR1, SignalRuntime};
let signals = SignalRuntime::set_with(&[SIGUSR1]).unwrap();
let old_mask = SignalRuntime::block_current_thread(&signals).unwrap();
let sfd = SignalRuntime::signalfd_new(&signals).unwrap();
let mut reactor = Reactor::new().unwrap();
let token = reactor.add(&sfd, true, false).unwrap();
unsafe {
libc::kill(libc::getpid(), SIGUSR1);
}
let mut events = Vec::new();
let n = reactor.wait(&mut events, 1, 1000).unwrap();
assert_eq!(n, 1);
assert_eq!(events[0].token, token);
assert!(events[0].readable);
SignalRuntime::restore_current_thread(&old_mask).unwrap();
}
#[test]
fn test_decode_inotify_events() {
let mut buf = Vec::new();
buf.extend_from_slice(&1i32.to_ne_bytes());
buf.extend_from_slice(&2u32.to_ne_bytes());
buf.extend_from_slice(&0u32.to_ne_bytes()); buf.extend_from_slice(&0u32.to_ne_bytes());
buf.extend_from_slice(&2i32.to_ne_bytes());
buf.extend_from_slice(&4u32.to_ne_bytes());
buf.extend_from_slice(&0u32.to_ne_bytes()); buf.extend_from_slice(&8u32.to_ne_bytes()); buf.extend_from_slice(b"file.txt");
let events = decode_events(&buf).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(
events[0],
InotifyEvent {
wd: 1,
mask: 2,
name: None,
}
);
assert_eq!(
events[1],
InotifyEvent {
wd: 2,
mask: 4,
name: Some(b"file.txt".to_vec()),
}
);
buf.extend_from_slice(&3i32.to_ne_bytes());
buf.extend_from_slice(&8u32.to_ne_bytes());
let err = decode_events(&buf).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_parse_proc_status() {
let content = "Name:\tcore_daemon\nState:\tR (running)\nUid:\t1000\t1000\t1000\t1000\nGid:\t1000\t1000\t1000\t1000\n";
let status = parse_proc_status(content).unwrap();
assert_eq!(status.name, "core_daemon");
assert_eq!(status.uid, 1000);
}
#[test]
fn test_parse_proc_status_missing_uid_is_invalid() {
let content = "Name:\tcore_daemon\nState:\tR (running)\n";
let err = parse_proc_status(content).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_parse_proc_cmdline_bytes_replaces_nul_separators() {
let content = b"/system/bin/sh\0-c\0echo hello\0";
assert_eq!(
parse_proc_cmdline_bytes(content),
"/system/bin/sh -c echo hello"
);
}
#[test]
fn test_read_proc_cmdline_at_uses_explicit_root() {
let proc_root = std::env::temp_dir().join(format!(
"coreshift_test_proc_root_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let pid_dir = proc_root.join("12345");
std::fs::create_dir_all(&pid_dir).unwrap();
std::fs::write(pid_dir.join("cmdline"), b"com.example.app\0:service\0").unwrap();
let cmdline = read_proc_cmdline_at(&proc_root, 12345).unwrap();
assert_eq!(cmdline, "com.example.app :service");
let _ = std::fs::remove_dir_all(&proc_root);
}
#[test]
fn test_proc_dir_pins_process_and_reads_via_dirfd() {
let proc_root = std::env::temp_dir().join(format!(
"coreshift_test_proc_dir_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let pid_dir = proc_root.join("4242");
std::fs::create_dir_all(&pid_dir).unwrap();
std::fs::write(pid_dir.join("cmdline"), b"com.example.manager\0").unwrap();
std::fs::write(pid_dir.join("status"), b"Name:\tmanager\nUid:\t12345\n").unwrap();
let dir = ProcDir::open_at(&proc_root, 4242).unwrap();
assert_eq!(dir.cmdline().unwrap(), "com.example.manager");
let status = dir.status().unwrap();
assert_eq!(status.name, "manager");
assert_eq!(status.uid, 12345);
assert_eq!(
dir.stat().unwrap().uid,
stat_at(&proc_root, 4242).unwrap().uid
);
let _ = std::fs::remove_dir_all(&proc_root);
}
#[test]
fn test_exec_context_validation() {
let res = SpawnOptions::builder(vec![], SpawnBackend::Fork).build();
assert!(res.is_err());
let res = SpawnOptions::builder(
vec!["valid".to_string(), "inv\0alid".to_string()],
SpawnBackend::Fork,
)
.build();
assert!(res.is_err());
let res = SpawnOptions::builder(
vec!["/bin/ls".to_string(), "-l".to_string()],
SpawnBackend::Fork,
)
.cwd("/tmp".to_string())
.build();
assert!(res.is_ok());
}
#[test]
fn test_posix_spawn_rejects_unsupported_cwd() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.cwd("/tmp".to_string())
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected unsupported posix_spawn cwd to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_posix_spawn_rejects_unsupported_setsid() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.pgroup(ProcessGroup::new(None, true))
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected unsupported posix_spawn setsid to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_posix_spawn_rejects_close_from3_fd_policy() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.fd_policy(SpawnFdPolicy::CloseFrom3)
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected unsupported fd policy to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_fork_reports_child_chdir_error() {
let err = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::Fork)
.cwd("/definitely/missing/coreshift-core-test".to_string())
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::ENOENT));
assert!(err.to_string().contains("spawn child chdir"));
}
#[test]
fn test_spawn_fd_policy_rejects_dirty_allowlist() {
let file = File::open("/dev/null").unwrap();
let fd = file.as_raw_fd();
let err = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.fd_policy(SpawnFdPolicy::Allowlist(vec![fd, fd]))
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_fork_close_from3_closes_inherited_fd() {
let file = File::open("/dev/null").unwrap();
let fd = file.as_raw_fd();
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
assert!(flags >= 0);
let ret = unsafe { libc::fcntl(fd, libc::F_SETFD, flags & !libc::FD_CLOEXEC) };
assert_eq!(ret, 0);
let script = format!("if [ -e /proc/$$/fd/{fd} ]; then echo open; else echo closed; fi");
let output = SpawnOptions::builder(
vec!["/bin/sh".to_string(), "-c".to_string(), script],
SpawnBackend::Fork,
)
.fd_policy(SpawnFdPolicy::CloseFrom3)
.capture_stdout()
.build()
.unwrap()
.run()
.unwrap();
assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "closed");
}
#[test]
fn test_exec_backends_support_isolated_pgroup_and_close_from3() {
for backend in [
SpawnBackend::Fork,
SpawnBackend::Vfork,
SpawnBackend::Clone3,
SpawnBackend::Clone3Pidfd,
] {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], backend)
.pgroup(ProcessGroup::new(None, true))
.fd_policy(SpawnFdPolicy::CloseFrom3)
.build()
.unwrap();
let running = spawn_start(opts).unwrap_or_else(|e| {
panic!("{backend:?} should accept isolated pgroup + CloseFrom3: {e}")
});
let _ = running.process.wait_blocking();
}
}
#[test]
fn test_exec_backends_reject_isolated_custom_leader() {
for backend in [
SpawnBackend::Fork,
SpawnBackend::Vfork,
SpawnBackend::Clone3,
SpawnBackend::Clone3Pidfd,
] {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], backend)
.pgroup(ProcessGroup::new(Some(5), true))
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("{backend:?} must reject isolated + custom leader"),
Err(e) => e,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
}
#[test]
fn test_posix_spawn_rejects_session_containment() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected posix_spawn session containment to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_session_containment_requires_isolated_pgroup() {
for backend in [
SpawnBackend::Fork,
SpawnBackend::Vfork,
SpawnBackend::Clone3,
SpawnBackend::Clone3Pidfd,
] {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], backend)
.session_containment()
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("{backend:?} must require isolated pgroup for containment"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
}
#[test]
fn test_session_containment_allows_normal_exec() {
for backend in [
SpawnBackend::Fork,
SpawnBackend::Vfork,
SpawnBackend::Clone3,
SpawnBackend::Clone3Pidfd,
] {
let out = SpawnOptions::builder(
vec!["/bin/echo".to_string(), "contained".to_string()],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stdout()
.build()
.unwrap()
.run()
.unwrap();
assert!(matches!(out.status, Some(ExitStatus::Exited(0))));
assert_eq!(String::from_utf8_lossy(&out.stdout).trim(), "contained");
}
}
#[test]
fn test_session_containment_blocks_setsid_escape() {
let setsid = ["/usr/bin/setsid", "/bin/setsid", "/system/bin/setsid"]
.into_iter()
.find(|p| crate::fs::path_exists(p));
let Some(setsid) = setsid else {
return;
};
let started = std::time::Instant::now();
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
format!("{setsid} sleep 60"),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stderr()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(
!output.timed_out,
"contained setsid must fail fast, not hang for the escaped sleep"
);
assert!(
started.elapsed().as_secs() < 8,
"contained setsid must exit promptly"
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
stderr.contains("setsid") && !stderr.is_empty(),
"expected setsid utility to report the denied syscall, got {stderr:?}"
);
}
fn session_pids(sid: libc::pid_t) -> Vec<libc::pid_t> {
let mut pids = Vec::new();
let Ok(entries) = std::fs::read_dir("/proc") else {
return pids;
};
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
let Ok(pid) = name.parse::<libc::pid_t>() else {
continue;
};
let Ok(stat) = std::fs::read_to_string(format!("/proc/{name}/stat")) else {
continue;
};
let Some(rest) = stat.rsplit_once(')') else {
continue;
};
let Some(rest) = rest.1.strip_prefix(' ') else {
continue;
};
let mut f = rest.split(' ');
let _ = f.next(); let _ = f.next(); let _ = f.next(); if f.next().and_then(|s| s.parse::<libc::pid_t>().ok()) == Some(sid) {
pids.push(pid);
}
}
pids
}
fn session_live_pids(sid: libc::pid_t) -> Vec<libc::pid_t> {
let mut pids = Vec::new();
let Ok(entries) = std::fs::read_dir("/proc") else {
return pids;
};
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
let Ok(pid) = name.parse::<libc::pid_t>() else {
continue;
};
let Ok(stat) = std::fs::read_to_string(format!("/proc/{name}/stat")) else {
continue;
};
let Some(rest) = stat.rsplit_once(')') else {
continue;
};
let Some(rest) = rest.1.strip_prefix(' ') else {
continue;
};
let mut f = rest.split(' ');
if f.next() == Some("Z") {
continue;
}
let _ = f.next(); let _ = f.next(); if f.next().and_then(|s| s.parse::<libc::pid_t>().ok()) == Some(sid) {
pids.push(pid);
}
}
pids
}
#[test]
fn parse_stat_starttime_extracts_field_22() {
let content = "123 (my proc) S 456 789 123 0 1 2 3 4 5 6 7 8 9 10 11 12 13 14 424242 15 16";
assert_eq!(crate::proc::parse_stat_starttime(content), Some(424242));
}
#[test]
fn parse_stat_starttime_tolerates_comm_with_spaces_and_parens() {
let content = "7 (bash (with ) parens) R 1 7 7 0 -1 4194560 0 1 2 3 4 5 6 7 8 9 10 11 999 12";
assert_eq!(crate::proc::parse_stat_starttime(content), Some(999));
}
#[test]
fn parse_stat_starttime_rejects_garbage() {
assert_eq!(crate::proc::parse_stat_starttime("not a stat line"), None);
assert_eq!(crate::proc::parse_stat_starttime("1 (x) R"), None);
}
fn drive_managed_with_input(
mut managed: crate::spawn::ManagedProcess,
buffer: Arc<Mutex<Vec<u8>>>,
mut on_output: impl FnMut(&str) -> Vec<u8>,
) -> Result<crate::spawn::Output, crate::CoreError> {
let mut reactor = Reactor::new().unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
let mut events = Vec::new();
loop {
if let Some(output) = managed.poll_completion(&mut reactor)? {
return Ok(output);
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
let bytes = on_output(&text);
if !bytes.is_empty() {
managed.write_input(&bytes)?;
}
}
}
fn accumulate_sink(buffer: Arc<Mutex<Vec<u8>>>) -> impl Fn(bool, &[u8]) -> SinkResult {
move |_is_stdout, bytes| {
buffer.lock().unwrap().extend_from_slice(bytes);
SinkResult::Accept
}
}
fn pty_contained_builder(
argv: Vec<String>,
backend: SpawnBackend,
grace_ms: u32,
) -> SpawnOptionsBuilder {
SpawnOptions::builder(argv, backend)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.pty()
.capture_stdout()
.timeout_ms(15_000)
.kill_grace_ms(grace_ms)
.cancel(CancelPolicy::Graceful)
}
fn pipe_contained_builder(
argv: Vec<String>,
backend: SpawnBackend,
grace_ms: u32,
) -> SpawnOptionsBuilder {
SpawnOptions::builder(argv, backend)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stdout()
.timeout_ms(15_000)
.kill_grace_ms(grace_ms)
.cancel(CancelPolicy::Graceful)
}
#[test]
fn test_pty_containment_allows_job_control_and_ctrl_c() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; echo shellpgid=$(ps -o pgid= -p $$ | tr -d ' '); \
sh -c 'echo fgpgid=$(ps -o pgid= -p $$ | tr -d \" \"); sleep 30'"
.to_string(),
],
backend,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let ctrl_c_sent = Arc::new(AtomicBool::new(false));
let sent = Arc::clone(&ctrl_c_sent);
let output = drive_managed_with_input(managed, Arc::clone(&buffer), move |text| {
if !sent.load(Ordering::SeqCst) && text.contains("fgpgid=") {
sent.store(true, Ordering::SeqCst);
vec![0x03] } else {
Vec::new()
}
})
.unwrap();
let merged = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
let shell_pgid = merged
.lines()
.find_map(|l| l.trim().strip_prefix("shellpgid=").map(str::trim))
.expect("shellpgid marker");
let fg_pgid = merged
.lines()
.find_map(|l| l.trim().strip_prefix("fgpgid=").map(str::trim))
.expect("fgpgid marker");
assert!(
shell_pgid != fg_pgid && !shell_pgid.is_empty() && !fg_pgid.is_empty(),
"{backend:?}: job control must put the foreground job in its own pgrp, got shell={shell_pgid} fg={fg_pgid}"
);
assert!(
ctrl_c_sent.load(Ordering::SeqCst),
"{backend:?}: Ctrl-C must have been sent (fg marker seen), got {merged:?}"
);
assert!(
!output.timed_out,
"{backend:?}: Ctrl-C must preempt the sleep before the timeout"
);
assert!(
session_pids(sid).is_empty(),
"{backend:?}: session must be empty after Ctrl-C, got {:?}",
session_pids(sid)
);
}
}
#[test]
fn test_pty_session_kill_is_total_across_groups() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; sleep 60 & echo bg1=$!; sleep 60 & echo bg2=$!; wait".to_string(),
],
backend,
200,
)
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
std::thread::sleep(std::time::Duration::from_millis(300));
let started = std::time::Instant::now();
let output = drive_managed(managed, true).unwrap();
assert!(
!output.timed_out,
"{backend:?}: session kill must terminate the whole session, got {output:?}"
);
assert!(
started.elapsed().as_secs() < 10,
"{backend:?}: session kill must be prompt"
);
let merged = String::from_utf8_lossy(&output.stdout);
let bg1: libc::pid_t = merged
.lines()
.find_map(|l| l.trim().strip_prefix("bg1="))
.and_then(|s| s.trim().parse().ok())
.expect("bg1 pid");
let bg2: libc::pid_t = merged
.lines()
.find_map(|l| l.trim().strip_prefix("bg2="))
.and_then(|s| s.trim().parse().ok())
.expect("bg2 pid");
assert!(
session_pids(sid).is_empty(),
"{backend:?}: session must be empty after cancellation, got {:?}",
session_pids(sid)
);
assert_ne!(bg1, bg2);
}
}
#[test]
fn test_pty_session_kill_continues_stopped_groups() {
let grace_ms = 8_000;
let managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; sleep 60 & echo bg=$!; kill -STOP $!; wait".to_string(),
],
SpawnBackend::Fork,
grace_ms,
)
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let started = std::time::Instant::now();
let output = drive_managed(managed, true).unwrap();
let elapsed = started.elapsed().as_millis();
assert!(
!output.timed_out,
"stopped group must be continued and die, got {output:?}"
);
assert!(
elapsed < grace_ms as u128,
"stopped group must die from TERM after SIGCONT, not the KILL escalation (elapsed {elapsed}ms vs grace {grace_ms}ms)"
);
assert!(session_pids(sid).is_empty(), "session must be empty");
}
#[test]
fn test_pty_session_kill_bounded_reenumeration_catches_raced_groups() {
let managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; i=0; while [ $i -lt 300 ]; do sleep 60 & i=$((i+1)); sleep 0.01; done; echo spawned=$i; wait".to_string(),
],
SpawnBackend::Fork,
200,
)
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let started = std::time::Instant::now();
let output = drive_managed(managed, true).unwrap();
assert!(
started.elapsed().as_secs() < 12,
"termination must stay bounded under a spawner, took {:?}",
started.elapsed()
);
let mut wait = 0u32;
while !session_pids(sid).is_empty() && wait < 50 {
std::thread::sleep(std::time::Duration::from_millis(50));
wait += 1;
}
assert!(
session_pids(sid).is_empty(),
"session must empty after bounded re-enumeration, still have {:?}",
session_pids(sid)
);
let _ = output;
}
#[test]
fn test_pty_session_kill_ends_before_escalation_deadline() {
let grace_ms = 8_000;
let managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; sleep 30 & wait".to_string(),
],
SpawnBackend::Fork,
grace_ms,
)
.build()
.unwrap(),
)
.unwrap();
let started = std::time::Instant::now();
let output = drive_managed(managed, true).unwrap();
let elapsed = started.elapsed().as_millis();
assert!(
!output.timed_out,
"clean TERM must EOF the master before the grace elapses, got {output:?}"
);
assert!(
elapsed < grace_ms as u128,
"must finish before the SIGKILL escalation (elapsed {elapsed}ms vs grace {grace_ms}ms)"
);
}
#[test]
fn test_pty_session_kill_leader_only_leaves_background_alive() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; sleep 60 & echo bg=$!; wait".to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.session_exit(SessionExitPolicy::LetMembersSurvive)
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let leader_pgid = sid;
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let mut leader_killed = false;
let started = std::time::Instant::now();
loop {
if !leader_killed {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg=")) {
let bg = rest.trim().parse::<libc::pid_t>().ok();
if bg.is_some_and(|bg| session_pids(sid).contains(&bg)) {
bg_pid = bg;
unsafe {
libc::kill(-leader_pgid, libc::SIGTERM);
}
leader_killed = true;
}
}
} else if let Some(bg) = bg_pid
&& started.elapsed() > std::time::Duration::from_millis(1500)
{
assert!(
session_pids(sid).contains(&bg),
"leader-only kill must leave the background job alive (bg={bg}, session={:?})",
session_pids(sid)
);
managed.request_cancel();
let output = drive_managed(managed, false)?;
assert!(
!output.timed_out,
"session-total kill must terminate the background job, got {output:?}"
);
assert!(
session_pids(sid).is_empty(),
"session must be empty after session-total kill, got {:?}",
session_pids(sid)
);
return Ok(());
}
reactor.wait(&mut events, 16, 50)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
}
}
#[test]
fn test_pty_session_enumeration_only_on_termination() -> Result<(), crate::CoreError> {
reset_pty_enumeration();
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pty_contained_builder(vec!["/bin/cat".to_string()], SpawnBackend::Fork, 200)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let started = std::time::Instant::now();
let mut wrote = false;
let mut resized = false;
let mut cancelled = false;
loop {
if let Some(_output) = managed.poll_completion(&mut reactor)? {
panic!("cat must stay alive until cancelled");
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
let elapsed = started.elapsed();
if !wrote {
managed.write_input(b"hello enumeration\n")?;
wrote = true;
}
if !resized && elapsed > std::time::Duration::from_millis(120) {
managed.resize_pty(50, 200)?;
resized = true;
}
assert!(
!pty_session_enumerated(sid),
"streaming/output/write/resize must never enumerate /proc"
);
if !cancelled && resized && elapsed > std::time::Duration::from_millis(250) {
managed.request_cancel();
cancelled = true;
}
if cancelled && elapsed > std::time::Duration::from_millis(250) {
break;
}
}
assert!(
!pty_session_enumerated(sid),
"no /proc enumeration before termination"
);
let output = drive_managed(managed, false)?;
let _ = output;
assert!(
pty_session_enumerated(sid),
"termination must enumerate the session's PGIDs"
);
assert!(session_pids(sid).is_empty(), "session must be empty");
Ok(())
}
#[test]
fn test_pty_natural_completion_waits_for_detached_background() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; sleep 1000 </dev/null >/dev/null 2>&1 & echo bg=$!; exit 0"
.to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let mut bg_seen_live = false;
let mut output = None;
let started = std::time::Instant::now();
loop {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if bg_pid.is_none()
&& let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg="))
{
bg_pid = rest.trim().parse::<libc::pid_t>().ok();
}
if let Some(bg) = bg_pid
&& session_pids(sid).contains(&bg)
&& unsafe { libc::kill(bg, 0) } == 0
{
bg_seen_live = true;
}
if let Some(out) = managed.poll_completion(&mut reactor)? {
output = Some(out);
break;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
break;
}
}
let output = output.expect("natural completion must eventually be reported");
assert!(
bg_seen_live,
"test setup broken: the detached background member was never observed live in the session"
);
assert!(
!output.timed_out,
"natural completion must not time out, got {output:?}"
);
assert!(
matches!(output.status, Some(crate::spawn::ExitStatus::Exited(0))),
"shell must exit 0, got {:?}",
output.status
);
assert!(
session_live_pids(sid).is_empty(),
"completion must be gated on an empty session, got {:?}",
session_live_pids(sid)
);
if let Some(bg) = bg_pid {
assert!(
unsafe { libc::kill(bg, 0) } != 0,
"detached background member must be swept before completion (bg={bg})"
);
}
Ok(())
}
#[test]
fn test_pty_natural_completion_sweeps_slave_holding_background() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; sleep 1000 & echo bg=$!; exit 0".to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let mut bg_seen_live = false;
let mut output = None;
let started = std::time::Instant::now();
loop {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if bg_pid.is_none()
&& let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg="))
{
bg_pid = rest.trim().parse::<libc::pid_t>().ok();
}
if let Some(bg) = bg_pid
&& session_pids(sid).contains(&bg)
&& unsafe { libc::kill(bg, 0) } == 0
{
bg_seen_live = true;
}
if let Some(out) = managed.poll_completion(&mut reactor)? {
output = Some(out);
break;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
break;
}
}
let output = output.expect("natural completion must not hang on a slave-holding member");
assert!(
bg_seen_live,
"test setup broken: the slave-holding background member was never observed live in the session"
);
assert!(
matches!(output.status, Some(crate::spawn::ExitStatus::Exited(0))),
"shell must exit 0, got {:?}",
output.status
);
assert!(
session_live_pids(sid).is_empty(),
"completion must be gated on an empty session, got {:?}",
session_live_pids(sid)
);
assert!(
output.swept_members,
"F14: the natural-completion sweep that emptied the session must be \
surfaced on Output (slave-holding member was SIGKILLed), got {:?}",
output.swept_members
);
if let Some(bg) = bg_pid {
assert!(
unsafe { libc::kill(bg, 0) } != 0,
"slave-holding background member must be swept before completion (bg={bg})"
);
}
Ok(())
}
#[test]
fn test_let_members_survive_reports_no_swept_members() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pty_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; sleep 1000 & echo bg=$!; exit 0".to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.session_exit(SessionExitPolicy::LetMembersSurvive)
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut output = None;
let started = std::time::Instant::now();
loop {
if let Some(out) = managed.poll_completion(&mut reactor)? {
output = Some(out);
break;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
break;
}
}
let output = output.expect("LetMembersSurvive must complete on leader-reap");
assert!(
!output.swept_members,
"F14: LetMembersSurvive must not report swept members, got {:?}",
output.swept_members
);
assert!(
matches!(output.status, Some(crate::spawn::ExitStatus::Exited(0))),
"shell must exit 0, got {:?}",
output.status
);
assert!(
!session_live_pids(sid).is_empty(),
"LetMembersSurvive must leave contained members running"
);
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
let bg = text
.lines()
.find_map(|l| l.trim().strip_prefix("bg="))
.and_then(|rest| rest.trim().parse::<libc::pid_t>().ok())
.expect("test setup broken: bg= never printed");
assert_eq!(unsafe { libc::kill(bg, libc::SIGKILL) }, 0, "cleanup kill");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while !session_live_pids(sid).is_empty() {
assert!(
std::time::Instant::now() < deadline,
"direct SIGKILL must empty the session, sid={sid}"
);
std::thread::sleep(std::time::Duration::from_millis(100));
}
Ok(())
}
#[test]
fn test_pipe_natural_completion_waits_for_detached_background() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pipe_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; sleep 1000 </dev/null >/dev/null 2>&1 & echo bg=$!; exit 0"
.to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let mut bg_seen_live = false;
let mut output = None;
let started = std::time::Instant::now();
loop {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if bg_pid.is_none()
&& let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg="))
{
bg_pid = rest.trim().parse::<libc::pid_t>().ok();
}
if let Some(bg) = bg_pid
&& session_pids(sid).contains(&bg)
&& unsafe { libc::kill(bg, 0) } == 0
{
bg_seen_live = true;
}
if let Some(out) = managed.poll_completion(&mut reactor)? {
output = Some(out);
break;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
break;
}
}
let output = output.expect("natural completion must eventually be reported");
assert!(
bg_seen_live,
"test setup broken: the detached background member was never observed live in the session"
);
assert!(
!output.timed_out,
"natural completion must not time out, got {output:?}"
);
assert!(
matches!(output.status, Some(crate::spawn::ExitStatus::Exited(0))),
"shell must exit 0, got {:?}",
output.status
);
assert!(
session_live_pids(sid).is_empty(),
"completion must be gated on an empty session, got {:?}",
session_live_pids(sid)
);
if let Some(bg) = bg_pid {
assert!(
unsafe { libc::kill(bg, 0) } != 0,
"detached background member must be swept before completion (bg={bg})"
);
}
Ok(())
}
#[test]
fn test_spawn_blocking_natural_completion_sweeps_detached_background() {
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; (sleep 1000 </dev/null >/dev/null 2>&1 &); exit 0".to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Graceful)
.build()
.unwrap()
.run()
.unwrap();
let sid = out.pid;
assert!(
matches!(out.status, Some(crate::spawn::ExitStatus::Exited(0))),
"shell must exit 0, got {:?}",
out.status
);
assert!(
!out.timed_out,
"blocking spawn must not time out, got {out:?}"
);
assert!(
session_live_pids(sid).is_empty(),
"completion must be gated on an empty session, got {:?}",
session_live_pids(sid)
);
assert!(
out.swept_members,
"F14: the natural-completion sweep that emptied the session must be \
surfaced on Output (detached member was SIGKILLed), got {:?}",
out.swept_members
);
}
#[test]
fn test_pipe_cancel_sweeps_leader_reap_survivor() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pipe_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' TERM; sleep 1000 & echo bg=$!; trap - TERM; sleep 30".to_string(),
],
SpawnBackend::Fork,
3_000,
)
.timeout_ms(1_000)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let mut bg_seen_live = false;
let mut output = None;
let started = std::time::Instant::now();
loop {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if bg_pid.is_none()
&& let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg="))
{
bg_pid = rest.trim().parse::<libc::pid_t>().ok();
}
if let Some(bg) = bg_pid
&& session_pids(sid).contains(&bg)
&& unsafe { libc::kill(bg, 0) } == 0
{
bg_seen_live = true;
}
if let Some(out) = managed.poll_completion(&mut reactor)? {
output = Some(out);
break;
}
reactor.wait(&mut events, 16, 50)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
break;
}
}
let output = output.expect("cancel must eventually be reported");
assert!(
bg_seen_live,
"test setup broken: the TERM-immune background member was never observed live in the session"
);
assert!(
output.timed_out,
"cancel must report timed_out, got {output:?}"
);
assert!(
output.status.is_some(),
"the leader must be reaped on cancel, got {:?}",
output.status
);
assert!(
session_live_pids(sid).is_empty(),
"cancel completion must be gated on an empty session, got {:?}",
session_live_pids(sid)
);
if let Some(bg) = bg_pid {
assert!(
unsafe { libc::kill(bg, 0) } != 0,
"TERM-immune background member must be swept before cancel completion (bg={bg})"
);
}
Ok(())
}
#[test]
fn test_spawn_blocking_cancel_sweeps_leader_reap_survivor() {
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' TERM; sleep 1000 & echo bg=$!; trap - TERM; sleep 30".to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stdout()
.timeout_ms(1_000)
.kill_grace_ms(3_000)
.cancel(CancelPolicy::Graceful)
.build()
.unwrap()
.run()
.unwrap();
let sid = out.pid;
let merged = String::from_utf8_lossy(&out.stdout);
let bg = merged
.lines()
.find_map(|l| l.trim().strip_prefix("bg=").map(str::trim))
.and_then(|s| s.parse::<libc::pid_t>().ok())
.expect("bg= marker must be present in the captured output");
assert!(
out.timed_out,
"blocking cancel must report timed_out, got {out:?}"
);
assert!(
out.status.is_some(),
"the leader must be reaped on cancel, got {:?}",
out.status
);
assert!(
session_live_pids(sid).is_empty(),
"blocking cancel must sweep the contained session before returning, got {:?}",
session_live_pids(sid)
);
assert!(
unsafe { libc::kill(bg, 0) } != 0,
"TERM-immune background member must be swept before the blocking call returns (bg={bg})"
);
}
#[test]
fn test_orphan_session_reaper_sweeps_detached_members() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pipe_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' TERM; sleep 1000 & echo bg=$!; sleep 30".to_string(),
],
SpawnBackend::Fork,
200,
)
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let started = std::time::Instant::now();
while bg_pid.is_none() {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg=")) {
let bg = rest.trim().parse::<libc::pid_t>().ok();
if bg.is_some_and(|bg| session_pids(sid).contains(&bg)) {
bg_pid = bg;
break;
}
}
let _ = managed.poll_completion(&mut reactor)?;
reactor.wait(&mut events, 16, 20)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
assert!(
started.elapsed() < std::time::Duration::from_secs(5),
"the TERM-immune background member was never observed live in the session"
);
}
let bg = bg_pid.unwrap();
orphan_session(sid);
let mut emptied = false;
for _ in 0..40 {
if session_live_pids(sid).is_empty() && unsafe { libc::kill(bg, 0) } != 0 {
emptied = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
assert!(
emptied,
"detached reaper must empty the orphaned session, live={:?} bg={bg}",
session_live_pids(sid)
);
Ok(())
}
#[test]
fn reaper_prune_removes_only_confirmed_reaped_pids() {
let mut set = std::collections::HashSet::new();
set.insert(101);
set.insert(202);
let snapshot: Vec<libc::pid_t> = vec![101];
let reaped: Vec<libc::pid_t> = snapshot
.iter()
.copied()
.filter(|pid| *pid == 101) .collect();
crate::spawn::prune_reaped_by_reaped_only(&mut set, &reaped);
assert!(!set.contains(&101), "reaped pid must be pruned");
assert!(
set.contains(&202),
"pid registered concurrently must survive the prune"
);
}
#[test]
fn orphan_session_refuses_unreadable_or_recycled_sid() {
let nonexistent = 1_000_000_000;
assert_eq!(crate::proc::starttime(nonexistent), None);
crate::spawn::orphan_session(nonexistent);
assert_eq!(
crate::spawn::orphaned_sessions_len(),
0,
"a sid whose leader starttime is unreadable must not be registered"
);
}
#[test]
fn session_scan_failures_are_fail_closed() {
let bad_root = std::path::Path::new("/nonexistent");
assert!(
crate::spawn::session_pgids_at(bad_root, 1234).is_err(),
"unreadable /proc must surface as Err, not an empty set"
);
assert!(
crate::spawn::session_sweep_at(bad_root, 1234).is_err(),
"unreadable /proc must surface as Err, not a false 'session empty'"
);
}
#[test]
fn session_scan_seam_reads_a_real_root() {
let pid = unsafe { libc::fork() };
assert!(pid >= 0);
if pid == 0 {
unsafe {
libc::setsid();
libc::pause();
}
unsafe { libc::_exit(0) };
}
let sid = pid; let pgids = crate::spawn::session_pgids_at(std::path::Path::new("/proc"), sid).unwrap();
assert!(
pgids.contains(&sid),
"the child's own pgid must be enumerable in its own session, got {pgids:?}"
);
assert!(
!crate::spawn::session_sweep_at(std::path::Path::new("/proc"), sid).unwrap(),
"a session with a live member must sweep as non-empty"
);
let mut status = 0;
loop {
let r = unsafe { libc::waitpid(sid, &mut status, 0) };
if r == sid {
break;
}
assert_eq!(r, -1, "waitpid of swept child failed");
assert_eq!(
std::io::Error::last_os_error().raw_os_error(),
Some(libc::EINTR)
);
}
assert!(
crate::spawn::session_sweep_at(std::path::Path::new("/proc"), sid).unwrap(),
"a session with no live member must sweep as empty"
);
}
#[test]
fn drop_pty_managed_process_kills_session_members() -> Result<(), crate::CoreError> {
let buffer = Arc::new(Mutex::new(Vec::new()));
let sink_buf = Arc::clone(&buffer);
let mut managed = spawn_managed(
pipe_contained_builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"trap '' HUP; sleep 1000 & echo bg=$!; exit 0".to_string(),
],
SpawnBackend::Fork,
200,
)
.pty()
.chunk_sink(accumulate_sink(sink_buf))
.build()
.unwrap(),
)
.unwrap();
let sid = managed.pid();
let mut reactor = Reactor::new()?;
managed.register_with_reactor(&mut reactor)?;
let mut events = Vec::new();
let mut bg_pid: Option<libc::pid_t> = None;
let started = std::time::Instant::now();
loop {
let text = String::from_utf8_lossy(&buffer.lock().unwrap()).to_string();
if bg_pid.is_none()
&& let Some(rest) = text.lines().find_map(|l| l.trim().strip_prefix("bg="))
{
let bg = rest.trim().parse::<libc::pid_t>().ok();
if bg.is_some_and(|bg| session_pids(sid).contains(&bg)) {
bg_pid = bg;
}
}
if let Some(output) = managed.poll_completion(&mut reactor)? {
let live = session_live_pids(sid);
assert!(
live.is_empty(),
"natural completion under Sweep must empty the session: \
output={output:?} live={live:?}"
);
return Ok(());
}
if bg_pid.is_some()
&& !std::path::Path::new(&format!("/proc/{sid}")).exists()
&& !session_pids(sid).is_empty()
{
break;
}
reactor.wait(&mut events, 16, 20)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
assert!(
started.elapsed() < std::time::Duration::from_secs(5),
"leader never reaped with a live member (bg={bg_pid:?} sid={sid})"
);
}
let bg = bg_pid.unwrap();
drop(managed);
let mut emptied = false;
for _ in 0..40 {
if session_live_pids(sid).is_empty() && unsafe { libc::kill(bg, 0) } != 0 {
emptied = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
assert!(
emptied,
"Drop must empty the session even with a reaped leader, live={:?} bg={bg}",
session_live_pids(sid)
);
Ok(())
}
#[test]
fn test_pipe_containment_still_denies_setpgid() {
let out = SpawnOptions::builder(
vec![
"/bin/bash".to_string(),
"-c".to_string(),
"set -m; sleep 30 & echo bg=$!; sleep 0.3; \
echo shell=$(ps -o pgid= -p $$ | tr -d ' '); \
echo bgpgid=$(ps -o pgid= -p $! | tr -d ' '); \
kill $! 2>/dev/null; wait 2>/dev/null; exit 0"
.to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.capture_stdout()
.capture_stderr()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap()
.run()
.unwrap();
let merged = String::from_utf8_lossy(&out.stdout);
let shell_pgid = merged
.lines()
.find_map(|l| l.trim().strip_prefix("shell=").map(str::trim))
.expect("shell pgid");
let bg_pgid = merged
.lines()
.find_map(|l| l.trim().strip_prefix("bgpgid=").map(str::trim))
.expect("bg pgid");
assert_eq!(
shell_pgid, bg_pgid,
"pipe containment must still deny setpgid (job stays in shell's pgrp), got shell={shell_pgid} bg={bg_pgid}"
);
assert!(!out.timed_out, "pipe child must exit cleanly");
}
#[test]
fn test_pty_containment_still_blocks_session_escape_syscalls() {
if !crate::fs::path_exists("/usr/bin/python3") {
return;
}
let setsid_probe = concat!(
"import os, time\n",
"pid = os.fork()\n",
"if pid == 0:\n",
" try:\n",
" os.setsid()\n",
" print('gc_setsid_ok', flush=True)\n",
" except OSError as e:\n",
" print('gc_setsid_errno=%d' % e.errno, flush=True)\n",
" time.sleep(1)\n",
"else:\n",
" time.sleep(1)\n",
);
let out = SpawnOptions::builder(
vec![
"/usr/bin/python3".to_string(),
"-c".to_string(),
setsid_probe.to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap()
.run()
.unwrap();
let merged = String::from_utf8_lossy(&out.stdout);
assert!(
merged.contains("gc_setsid_errno=1"),
"pty containment must deny a grandchild's setsid, got {merged:?}"
);
assert!(
!out.timed_out,
"denied setsid must exit promptly, not hold the master"
);
let syscall_probes = r#"
import os
try:
os.setns(99, 0); print("setns_ok")
except OSError as e:
print("setns_errno=%d" % e.errno)
try:
os.unshare(os.CLONE_NEWUSER); print("unshare_ok")
except OSError as e:
print("unshare_errno=%d" % e.errno)
"#;
let out = SpawnOptions::builder(
vec![
"/usr/bin/python3".to_string(),
"-c".to_string(),
syscall_probes.to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.session_containment()
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap()
.run()
.unwrap();
let merged = String::from_utf8_lossy(&out.stdout);
assert!(
merged.contains("setns_errno=1"),
"pty containment must deny setns (EPERM, not the natural EBADF), got {merged:?}"
);
assert!(
merged.contains("unshare_errno=1"),
"pty containment must deny unshare (EPERM, uncontained it succeeds), got {merged:?}"
);
}
#[test]
fn test_pty_rejects_non_isolated_pgroup() {
for backend in [
SpawnBackend::Fork,
SpawnBackend::Vfork,
SpawnBackend::Clone3,
SpawnBackend::Clone3Pidfd,
] {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], backend)
.pty()
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("{backend:?} pty must require an isolated pgroup"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
}
#[test]
fn test_posix_spawn_rejects_pty() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.pgroup(ProcessGroup::new(None, true))
.pty()
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected posix_spawn pty to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_pty_rejects_stdin_buffer() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::Fork)
.pgroup(ProcessGroup::new(None, true))
.pty()
.stdin(b"hello".to_vec())
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("pty + stdin buffer must be rejected"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_pty_merged_io_and_controlling_terminal() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let script = "test -t 0 && test -t 1 && test -t 2 && echo stdio-is-tty; \
readlink /proc/self/fd/0; readlink /proc/self/fd/1; \
echo out-line; echo err-line >&2; \
echo sid=$(ps -o sid= -p $$); echo pgid=$(ps -o pgid= -p $$)";
let out = SpawnOptions::builder(
vec!["/bin/sh".to_string(), "-c".to_string(), script.to_string()],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.build()
.unwrap()
.run()
.unwrap();
assert!(
matches!(out.status, Some(ExitStatus::Exited(0))),
"{backend:?}: pty child must exit 0, got {out:?}"
);
let merged = String::from_utf8_lossy(&out.stdout);
assert!(
merged.contains("stdio-is-tty"),
"{backend:?}: all of stdin/stdout/stderr must be a tty, got {merged:?}"
);
let pts_lines: Vec<_> = merged
.lines()
.filter(|l| l.starts_with("/dev/pts/"))
.collect();
assert!(
pts_lines.len() >= 2,
"{backend:?}: fd 0/1 must be a pty slave, got {merged:?}"
);
let out_pos = merged.find("out-line").expect("merged stdout");
let err_pos = merged.find("err-line").expect("merged stderr");
assert!(
out_pos < err_pos,
"{backend:?}: stdout/stderr must be merged in order, got {merged:?}"
);
let sid = merged
.lines()
.find_map(|l| l.trim_start().strip_prefix("sid="))
.expect("sid line")
.trim();
let pgid = merged
.lines()
.find_map(|l| l.trim_start().strip_prefix("pgid="))
.expect("pgid line")
.trim();
assert_eq!(
sid, pgid,
"{backend:?}: child must be a session+group leader, got {merged:?}"
);
assert_eq!(
sid,
out.pid.to_string(),
"{backend:?}: child must be its own session leader (sid==pid), got {merged:?}"
);
}
}
#[test]
fn test_pty_resize_applies_winsize() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"sleep 1; stty size".to_string(),
],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
managed.resize_pty(24, 100).unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(!output.timed_out, "{backend:?}: pty resize child must exit");
let merged = String::from_utf8_lossy(&output.stdout);
assert!(
merged.trim() == "24 100",
"{backend:?}: TIOCSWINSZ must propagate to the child, got {merged:?}"
);
}
}
#[test]
fn test_pty_with_caller_pair_window_and_env() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let (master, slave) = make_pty().expect("make_pty");
let (rows, cols) = pty_window(&master, 30, 120).expect("pty_window");
assert_eq!((rows, cols), (30, 120), "{backend:?}: winsize read-back");
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"echo \"$LINES:$COLUMNS\"; stty size".to_string(),
],
backend,
)
.env(vec![format!("LINES={rows}"), format!("COLUMNS={cols}")])
.pgroup(ProcessGroup::new(None, true))
.pty_with(master, slave)
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(!output.timed_out, "{backend:?}: pty_with child must exit");
let merged = String::from_utf8_lossy(&output.stdout);
let env_line = merged.lines().next().unwrap_or("");
let stty_line = merged.lines().nth(1).unwrap_or("");
assert_eq!(
env_line, "30:120",
"{backend:?}: env LINES/COLUMNS must match the applied window, got {merged:?}"
);
assert_eq!(
stty_line, "30 120",
"{backend:?}: child must observe the applied winsize, got {merged:?}"
);
}
}
#[test]
fn test_pty_without_env_inherits_winsize() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let (master, slave) = make_pty().expect("make_pty");
let (rows, cols) = pty_window(&master, 40, 80).expect("pty_window");
assert_eq!((rows, cols), (40, 80), "{backend:?}: winsize read-back");
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"stty size; sh -c 'echo ${LINES:-unset}:${COLUMNS:-unset}'".to_string(),
],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty_with(master, slave)
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(!output.timed_out, "{backend:?}: pty_with child must exit");
let merged = String::from_utf8_lossy(&output.stdout);
let stty_line = merged.lines().next().unwrap_or("");
assert_eq!(
stty_line, "40 80",
"{backend:?}: child must observe the applied winsize, got {merged:?}"
);
}
}
#[test]
fn test_pty_write_input_round_trips_to_the_child() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"read line; echo got:$line".to_string(),
],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let n = managed.write_input(b"hello\r").unwrap();
assert_eq!(n, 6, "{backend:?}: write_input must accept every byte");
let output = drive_managed(managed, false).unwrap();
assert!(!output.timed_out, "{backend:?}: echo child must exit");
let merged = String::from_utf8_lossy(&output.stdout);
assert!(
merged.contains("got:hello"),
"{backend:?}: child must read the written line, got {merged:?}"
);
}
}
#[test]
fn test_pty_write_input_ctrl_c_signals_the_foreground_group() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"echo start; sleep 30".to_string(),
],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
managed.write_input(b"\x03").unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(
!output.timed_out,
"{backend:?}: Ctrl-C must preempt the sleep before the timeout"
);
let status = output.status;
assert!(
matches!(
status,
Some(ExitStatus::Signaled(libc::SIGINT)) | Some(ExitStatus::Exited(130))
),
"{backend:?}: Ctrl-C must terminate the child by SIGINT, got {status:?}"
);
}
}
#[test]
fn test_pty_write_input_rejects_non_pty_spawn() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/sleep".to_string(), "0.2".to_string()],
SpawnBackend::Fork,
)
.capture_stdout()
.build()
.unwrap(),
)
.unwrap();
match managed.write_input(b"x") {
Err(e) => assert_eq!(e.raw_os_error(), Some(libc::EINVAL)),
Ok(n) => panic!("write_input on a pipe-capture spawn must be refused, wrote {n}"),
}
}
#[test]
fn test_pty_write_input_nonblock_round_trips_to_the_child() {
for backend in [SpawnBackend::Fork, SpawnBackend::Clone3Pidfd] {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"read line; echo got:$line".to_string(),
],
backend,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let n = managed.write_input_nonblock(b"hello\r").unwrap();
assert_eq!(
n,
Some(6),
"{backend:?}: write_input_nonblock must accept every byte on a fresh pty"
);
let output = drive_managed(managed, false).unwrap();
assert!(!output.timed_out, "{backend:?}: echo child must exit");
let merged = String::from_utf8_lossy(&output.stdout);
assert!(
merged.contains("got:hello"),
"{backend:?}: child must read the written line, got {merged:?}"
);
}
}
#[test]
fn test_drain_maps_pty_master_eio_to_eof() {
let master = unsafe {
libc::open(
c"/dev/ptmx".as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(master >= 0);
assert_eq!(unsafe { libc::grantpt(master) }, 0);
assert_eq!(unsafe { libc::unlockpt(master) }, 0);
let mut name = [0 as libc::c_char; 4096];
assert_eq!(
unsafe { libc::ptsname_r(master, name.as_mut_ptr(), name.len()) },
0
);
let slave = unsafe {
libc::open(
name.as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(slave >= 0);
let master_fd = Fd::new(master, "test pty master").unwrap();
let slave_fd = Fd::new(slave, "test pty slave").unwrap();
let mut drain: DrainState<fn(&[u8]) -> bool> =
DrainState::new(None, None, Some(master_fd), None, 1024, None, None, true).unwrap();
assert!(!drain.read_fd(true).unwrap());
drop(slave_fd);
let mut reads = 0;
while !drain.read_fd(true).unwrap() {
reads += 1;
if reads > 100 {
panic!("pty master EIO was not mapped to EOF");
}
}
assert!(drain.is_done());
let master2 = unsafe {
libc::open(
c"/dev/ptmx".as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(master2 >= 0);
assert_eq!(unsafe { libc::grantpt(master2) }, 0);
assert_eq!(unsafe { libc::unlockpt(master2) }, 0);
let mut name2 = [0 as libc::c_char; 4096];
assert_eq!(
unsafe { libc::ptsname_r(master2, name2.as_mut_ptr(), name2.len()) },
0
);
let slave2 = unsafe {
libc::open(
name2.as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(slave2 >= 0);
let master2_fd = Fd::new(master2, "test pty master 2").unwrap();
let slave2_fd = Fd::new(slave2, "test pty slave 2").unwrap();
let mut drain2: DrainState<fn(&[u8]) -> bool> =
DrainState::new(None, None, Some(master2_fd), None, 1024, None, None, false).unwrap();
drop(slave2_fd);
let err = loop {
match drain2.read_fd(true) {
Ok(true) => panic!("non-pty drain must not treat EIO as EOF"),
Ok(false) => continue,
Err(e) => break e,
}
};
assert_eq!(err.raw_os_error(), Some(libc::EIO));
}
#[test]
fn test_vfork_reports_child_chdir_error() {
let err = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::Vfork)
.cwd("/definitely/missing/coreshift-core-test".to_string())
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::ENOENT));
assert!(err.to_string().contains("spawn child chdir"));
}
#[test]
fn test_clone3_reports_child_chdir_error() {
let err = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::Clone3)
.cwd("/definitely/missing/coreshift-core-test".to_string())
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::ENOENT));
assert!(err.to_string().contains("spawn child chdir"));
}
#[test]
fn test_fork_close_from3_closes_many_inherited_fds() {
let mut files = Vec::new();
let mut fds = Vec::new();
for _ in 0..200 {
let f = File::open("/dev/null").unwrap();
let fd = f.as_raw_fd();
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
assert!(flags >= 0);
let ret = unsafe { libc::fcntl(fd, libc::F_SETFD, flags & !libc::FD_CLOEXEC) };
assert_eq!(ret, 0);
fds.push(fd);
files.push(f);
}
let probe = *fds.last().unwrap();
let script = format!("if [ -e /proc/$$/fd/{probe} ]; then echo open; else echo closed; fi");
let output = SpawnOptions::builder(
vec!["/bin/sh".to_string(), "-c".to_string(), script],
SpawnBackend::Fork,
)
.fd_policy(SpawnFdPolicy::CloseFrom3)
.capture_stdout()
.build()
.unwrap()
.run()
.unwrap();
assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "closed");
drop(files);
}
#[test]
fn test_fork_allowlist_preserves_allowed_fd() {
let file = File::open("/dev/null").unwrap();
let fd = file.as_raw_fd();
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
assert!(flags >= 0);
let ret = unsafe { libc::fcntl(fd, libc::F_SETFD, flags & !libc::FD_CLOEXEC) };
assert_eq!(ret, 0);
let script = format!("if [ -e /proc/$$/fd/{fd} ]; then echo open; else echo closed; fi");
let output = SpawnOptions::builder(
vec!["/bin/sh".to_string(), "-c".to_string(), script],
SpawnBackend::Fork,
)
.fd_policy(SpawnFdPolicy::Allowlist(vec![fd]))
.capture_stdout()
.build()
.unwrap()
.run()
.unwrap();
assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "open");
}
#[test]
fn test_process_echild() {
let p = Process::new(999999);
let res = p.wait_step();
assert!(res.is_err());
}
#[test]
fn test_spawn_start_wait_false_validation() {
assert!(spawn_start(base_wait_false_builder().build().unwrap()).is_ok());
assert!(spawn_start(base_wait_false_builder().capture_stdout().build().unwrap()).is_err());
assert!(
spawn_start(
base_wait_false_builder()
.stdin(vec![1, 2, 3])
.build()
.unwrap()
)
.is_err()
);
}
fn base_wait_false_builder() -> SpawnOptionsBuilder {
SpawnOptions::builder(
vec!["/bin/sh".to_string(), "-c".to_string(), "true".to_string()],
SpawnBackend::PosixSpawn,
)
.wait(false)
.max_output(1024)
.kill_grace_ms(1000)
}
#[test]
fn test_spawn_wait_false_child_is_reaped_by_reaper() {
let output = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.wait(false)
.build()
.unwrap()
.run()
.unwrap();
let pid = output.pid;
assert!(pid > 0);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let mut status: libc::c_int = 0;
let r = unsafe { libc::waitpid(pid, &mut status, libc::WNOHANG) };
if r == pid {
break;
}
if r < 0 && std::io::Error::last_os_error().raw_os_error() == Some(libc::ECHILD) {
return;
}
if std::time::Instant::now() > deadline {
panic!("wait=false child {pid} was never reaped (zombie leak)");
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
}
#[test]
fn test_pdeath_signal_kills_child_when_parent_exits() {
use std::io::Read;
let mut fds = [0i32; 2];
assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
for f in fds {
let fl = unsafe { libc::fcntl(f, libc::F_GETFD) };
assert!(fl >= 0);
assert_eq!(
unsafe { libc::fcntl(f, libc::F_SETFD, fl | libc::FD_CLOEXEC) },
0
);
}
let (r, w) = (fds[0], fds[1]);
let daemon = unsafe { libc::fork() };
assert!(daemon >= 0);
if daemon == 0 {
let opts = SpawnOptions::builder(
vec!["/bin/sleep".to_string(), "1000".to_string()],
SpawnBackend::Fork,
)
.wait(false)
.pdeath_signal(libc::SIGKILL)
.build()
.unwrap();
let running = spawn_start(opts).unwrap();
let gp = running.process.pid();
unsafe {
let pid_str = gp.to_string();
let bytes = pid_str.as_bytes();
let mut written = 0;
while written < bytes.len() {
let n = libc::write(w, bytes[written..].as_ptr().cast(), bytes.len() - written);
if n < 0 {
if std::io::Error::last_os_error().raw_os_error() == Some(libc::EINTR) {
continue;
}
break;
}
written += n as usize;
}
libc::_exit(0);
}
}
unsafe { libc::close(w) };
let mut pid_buf = String::new();
{
use std::os::unix::io::FromRawFd;
let rfile = unsafe { std::fs::File::from_raw_fd(r) };
let mut br = std::io::BufReader::new(rfile);
br.read_to_string(&mut pid_buf).unwrap();
}
let gp: i32 = pid_buf.trim().parse().unwrap();
assert!(gp > 0);
let mut status: libc::c_int = 0;
assert_eq!(unsafe { libc::waitpid(daemon, &mut status, 0) }, daemon);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let alive = std::path::Path::new(&format!("/proc/{gp}")).exists();
if !alive {
break;
}
if std::time::Instant::now() > deadline {
unsafe { libc::kill(gp, libc::SIGKILL) };
panic!("grandchild {gp} survived the daemon's exit — pdeath_signal not delivered");
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
}
#[test]
fn test_pdeath_signal_off_means_child_survives_parent_exit() {
use std::io::Read;
let mut fds = [0i32; 2];
assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
for f in fds {
let fl = unsafe { libc::fcntl(f, libc::F_GETFD) };
assert!(fl >= 0);
assert_eq!(
unsafe { libc::fcntl(f, libc::F_SETFD, fl | libc::FD_CLOEXEC) },
0
);
}
let (r, w) = (fds[0], fds[1]);
let daemon = unsafe { libc::fork() };
assert!(daemon >= 0);
if daemon == 0 {
let opts = SpawnOptions::builder(
vec!["/bin/sleep".to_string(), "1000".to_string()],
SpawnBackend::Fork,
)
.wait(false)
.build()
.unwrap();
let running = spawn_start(opts).unwrap();
let gp = running.process.pid();
unsafe {
let pid_str = gp.to_string();
let bytes = pid_str.as_bytes();
let mut written = 0;
while written < bytes.len() {
let n = libc::write(w, bytes[written..].as_ptr().cast(), bytes.len() - written);
if n < 0 {
if std::io::Error::last_os_error().raw_os_error() == Some(libc::EINTR) {
continue;
}
break;
}
written += n as usize;
}
libc::_exit(0);
}
}
unsafe { libc::close(w) };
let mut pid_buf = String::new();
{
use std::os::unix::io::FromRawFd;
let rfile = unsafe { std::fs::File::from_raw_fd(r) };
let mut br = std::io::BufReader::new(rfile);
br.read_to_string(&mut pid_buf).unwrap();
}
let gp: i32 = pid_buf.trim().parse().unwrap();
assert!(gp > 0);
let mut status: libc::c_int = 0;
assert_eq!(unsafe { libc::waitpid(daemon, &mut status, 0) }, daemon);
std::thread::sleep(std::time::Duration::from_millis(300));
assert!(
std::path::Path::new(&format!("/proc/{gp}")).exists(),
"child {gp} died without pdeath_signal"
);
unsafe {
libc::kill(gp, libc::SIGKILL);
libc::waitpid(gp, &mut status, 0);
}
}
#[test]
fn test_pdeath_signal_rejected_on_posix_spawn() {
let err = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.pdeath_signal(libc::SIGKILL)
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_reactor_wait_zero_events() {
let mut reactor = Reactor::new().unwrap();
let mut events = Vec::new();
let res = reactor.wait(&mut events, 0, 0);
assert!(res.is_ok());
assert_eq!(res.unwrap(), 0);
}
#[test]
fn test_reactor_add_priority_registers_fd() {
let mut reactor = Reactor::new().unwrap();
let fd = Fd::eventfd(0).unwrap();
let token = reactor.add_priority(&fd).unwrap();
let mut events = Vec::new();
assert_eq!(reactor.wait(&mut events, 4, 0).unwrap(), 0);
assert!(events.is_empty());
assert!(format!("{token:?}").starts_with("Token("));
}
#[test]
fn test_reactor_mod_toggles_read_interest() {
let mut reactor = Reactor::new().unwrap();
let fd = Fd::eventfd(0).unwrap();
let token = reactor.add(&fd, true, false).unwrap();
let mut events = Vec::new();
assert_eq!(reactor.wait(&mut events, 4, 0).unwrap(), 0);
fd.write_u64(1).unwrap();
events.clear();
assert_eq!(reactor.wait(&mut events, 4, 100).unwrap(), 1);
assert_eq!(events[0].token, token);
assert!(events[0].readable);
assert_eq!(fd.read_u64().unwrap(), Some(1));
reactor.mod_(&fd, token, false, false).unwrap();
fd.write_u64(1).unwrap();
events.clear();
assert_eq!(
reactor.wait(&mut events, 4, 50).unwrap(),
0,
"cleared readable interest must not deliver events"
);
reactor.mod_(&fd, token, true, false).unwrap();
events.clear();
assert_eq!(reactor.wait(&mut events, 4, 100).unwrap(), 1);
assert_eq!(events[0].token, token);
assert!(events[0].readable);
}
#[test]
fn test_fd_slice_io_helpers() {
let mut fds = [0; 2];
unsafe { libc::pipe(fds.as_mut_ptr()) };
let r = Fd::new(fds[0], "pipe").unwrap();
let w = Fd::new(fds[1], "pipe").unwrap();
let input = b"hello";
let written = w.write_slice(input).unwrap().unwrap();
assert_eq!(written, input.len());
let mut buf = [0u8; 8];
let read = r.read_slice(&mut buf).unwrap().unwrap();
assert_eq!(read, input.len());
assert_eq!(&buf[..read], input);
}
#[test]
fn test_unix_socket_bind_connect_roundtrip() {
let path = temp_socket_path("unix_roundtrip");
let listener = bind_test_unix_listener(&path, false);
let client = connect_test_unix_stream(UnixSocketAddr::Path(&path));
let server = loop {
if let Some(server) = listener.accept().unwrap() {
break server;
}
std::thread::yield_now();
};
let written = client.fd().write_slice(b"ping").unwrap().unwrap();
assert_eq!(written, 4);
let mut buf = [0u8; 8];
let read = loop {
if let Some(read) = server.fd().read_slice(&mut buf).unwrap() {
break read;
}
std::thread::yield_now();
};
assert_eq!(&buf[..read], b"ping");
chmod(UnixSocketAddr::Path(&path), 0o600).unwrap();
let mode = std::fs::metadata(&path).unwrap().permissions().mode();
assert_eq!(mode & 0o777, 0o600);
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_abstract_bind_connect_roundtrip() {
let name = format!(
"coreshift_abstract_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
);
let listener = bind(
UnixSocketAddr::Abstract(name.as_bytes()),
UnixSocketBindOptions::default(),
)
.unwrap();
let client = connect_test_unix_stream(UnixSocketAddr::Abstract(name.as_bytes()));
let server = loop {
if let Some(server) = listener.accept().unwrap() {
break server;
}
std::thread::yield_now();
};
let written = client.fd().write_slice(b"pong").unwrap().unwrap();
assert_eq!(written, 4);
let mut buf = [0u8; 8];
let read = loop {
if let Some(read) = server.fd().read_slice(&mut buf).unwrap() {
break read;
}
std::thread::yield_now();
};
assert_eq!(&buf[..read], b"pong");
}
#[test]
fn test_unix_socket_connect_result_finish_is_safe() {
let path = temp_socket_path("unix_connect_result");
let _listener = bind_test_unix_listener(&path, false);
let client = match connect(UnixSocketAddr::Path(&path)).unwrap() {
UnixConnectResult::Connected(stream) => {
assert_eq!(stream.check_connect_error().unwrap(), None);
stream
}
UnixConnectResult::InProgress(stream) => stream.finish_connect().unwrap(),
};
drop(client);
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_unlinks_stale_path_when_requested() {
let path = temp_socket_path("unix_unlink_stale");
{
let _listener = bind_test_unix_listener(&path, false);
}
let _listener = bind_test_unix_listener(&path, true);
assert!(std::fs::metadata(&path).unwrap().file_type().is_socket());
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_does_not_unlink_stale_path_by_default() {
let path = temp_socket_path("unix_no_unlink");
File::create(&path).unwrap();
let result = bind(
UnixSocketAddr::Path(&path),
UnixSocketBindOptions {
stale_socket_policy: StaleSocketPolicy::Preserve,
mode: None,
},
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EADDRINUSE));
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_unlink_socket_only_preserves_regular_file() {
let path = temp_socket_path("unix_regular_preserved");
File::create(&path).unwrap();
let result = bind(
UnixSocketAddr::Path(&path),
UnixSocketBindOptions {
stale_socket_policy: StaleSocketPolicy::UnlinkSocketOnly,
mode: None,
},
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EEXIST));
assert!(std::fs::metadata(&path).unwrap().is_file());
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_nul_path_rejected() {
let path = PathBuf::from(std::ffi::OsString::from_vec(b"/tmp/core\0socket".to_vec()));
let result = bind(
UnixSocketAddr::Path(&path),
UnixSocketBindOptions {
stale_socket_policy: StaleSocketPolicy::Preserve,
mode: None,
},
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
let path = PathBuf::from(std::ffi::OsString::from_vec(b"/tmp/core\0socket".to_vec()));
let result = connect(UnixSocketAddr::Path(&path));
let err = match result {
Ok(_) => panic!("connect unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_unix_socket_overlong_path_and_abstract_rejected() {
let path = PathBuf::from("/tmp".to_string() + &"/x".repeat(120));
let result = bind(
UnixSocketAddr::Path(&path),
UnixSocketBindOptions::default(),
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::ENAMETOOLONG));
let name = vec![b'x'; 108];
let result = bind(
UnixSocketAddr::Abstract(&name),
UnixSocketBindOptions::default(),
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::ENAMETOOLONG));
}
#[test]
fn test_unix_socket_accept_returns_none_on_eagain() {
let path = temp_socket_path("unix_accept_eagain");
let listener = bind_test_unix_listener(&path, false);
assert!(listener.accept().unwrap().is_none());
let _ = remove_file(&path);
}
#[test]
fn test_unix_socket_chmod_path_only() {
let regular = temp_socket_path("unix_chmod_regular");
File::create(®ular).unwrap();
let err = chmod(UnixSocketAddr::Path(®ular), 0o600).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
let _ = remove_file(®ular);
let name = b"coreshift_abstract_chmod_rejected";
let err = chmod(UnixSocketAddr::Abstract(name), 0o600).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
let result = bind(
UnixSocketAddr::Abstract(name),
UnixSocketBindOptions {
stale_socket_policy: StaleSocketPolicy::Preserve,
mode: Some(0o600),
},
);
let err = match result {
Ok(_) => panic!("bind unexpectedly succeeded"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_unix_socket_peer_credentials_when_supported() {
let path = temp_socket_path("unix_peer_cred");
let listener = bind_test_unix_listener(&path, false);
let client = connect_test_unix_stream(UnixSocketAddr::Path(&path));
let server = loop {
if let Some(server) = listener.accept().unwrap() {
break server;
}
std::thread::yield_now();
};
if let Some(cred) = server.peer_cred().unwrap() {
assert_eq!(cred.pid, Some(std::process::id() as i32));
assert_eq!(cred.uid, unsafe { libc::geteuid() });
assert_eq!(cred.gid, unsafe { libc::getegid() });
}
drop(client);
let _ = remove_file(&path);
}
fn early_exit_on_stop(chunk: &[u8]) -> bool {
chunk.windows(4).any(|window| window == b"stop")
}
#[test]
fn test_drain_early_exit_is_explicit_state_not_eof() {
let mut fds = [0; 2];
unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC | libc::O_NONBLOCK) };
let r = Fd::new(fds[0], "pipe").unwrap();
let w = Fd::new(fds[1], "pipe").unwrap();
w.write_slice(b"stop").unwrap();
let mut drain = DrainState::new(
None,
None,
Some(r),
None,
128,
Some(early_exit_on_stop),
None,
false,
)
.unwrap();
assert!(drain.read_fd(true).unwrap());
assert!(drain.stdout_early_exited());
assert!(!drain.output_limit_exceeded());
}
#[test]
fn test_drain_pause_preserves_registration_and_resume_rearms() {
let mut fds = [0; 2];
unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC | libc::O_NONBLOCK) };
let r = Fd::new(fds[0], "pipe").unwrap();
let w = Fd::new(fds[1], "pipe").unwrap();
let mut drain: DrainState<fn(&[u8]) -> bool> =
DrainState::new(None, None, Some(r), None, 1024, None, None, false).unwrap();
let mut reactor = Reactor::new().unwrap();
drain.register_with_reactor(&mut reactor).unwrap();
let token = drain.stdout_slot.as_ref().unwrap().token.unwrap();
drain.pause_stdout(&mut reactor).unwrap();
assert_eq!(
drain.stdout_slot.as_ref().unwrap().token,
Some(token),
"pause must keep the fd registered (MOD), not del it"
);
w.write_slice(b"x").unwrap();
let mut events = Vec::new();
assert_eq!(
reactor.wait(&mut events, 4, 50).unwrap(),
0,
"paused readable interest must not deliver events"
);
assert!(drain.resume_stdout(&mut reactor).unwrap());
let mut events = Vec::new();
assert_eq!(reactor.wait(&mut events, 4, 100).unwrap(), 1);
assert_eq!(events[0].token, token);
assert!(events[0].readable);
assert!(!drain.read_fd(true).unwrap());
assert!(!drain.stdout_paused());
}
#[test]
fn test_pty_set_writable_arms_and_disarms_writable_interest() {
let master = unsafe {
libc::open(
c"/dev/ptmx".as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(master >= 0);
assert_eq!(unsafe { libc::grantpt(master) }, 0);
assert_eq!(unsafe { libc::unlockpt(master) }, 0);
let mut name = [0 as libc::c_char; 4096];
assert_eq!(
unsafe { libc::ptsname_r(master, name.as_mut_ptr(), name.len()) },
0
);
let slave = unsafe {
libc::open(
name.as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_CLOEXEC,
)
};
assert!(slave >= 0);
let master_fd = Fd::new(master, "test pty master").unwrap();
let _slave_fd = Fd::new(slave, "test pty slave").unwrap();
let mut drain: DrainState<fn(&[u8]) -> bool> =
DrainState::new(None, None, Some(master_fd), None, 1024, None, None, true).unwrap();
let mut reactor = Reactor::new().unwrap();
drain.register_with_reactor(&mut reactor).unwrap();
let token = drain.stdout_slot.as_ref().unwrap().token.unwrap();
assert_eq!(drain.pty_input_token(), Some(token));
drain.set_pty_writable(&mut reactor, true).unwrap();
assert_eq!(drain.stdout_slot.as_ref().unwrap().token, Some(token));
assert!(drain.stdout_slot.as_ref().unwrap().writable);
assert!(drain.stdout_slot.as_ref().unwrap().readable);
let mut events = Vec::new();
assert_eq!(reactor.wait(&mut events, 4, 100).unwrap(), 1);
assert_eq!(events[0].token, token);
assert!(events[0].writable);
drain.set_pty_writable(&mut reactor, false).unwrap();
assert_eq!(drain.stdout_slot.as_ref().unwrap().token, Some(token));
assert!(!drain.stdout_slot.as_ref().unwrap().writable);
assert!(drain.stdout_slot.as_ref().unwrap().readable);
let mut events = Vec::new();
assert_eq!(
reactor.wait(&mut events, 4, 50).unwrap(),
0,
"disarmed writable interest must not deliver events"
);
}
#[test]
fn test_pty_set_writable_rejects_non_pty() {
let mut fds = [0; 2];
unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC | libc::O_NONBLOCK) };
let r = Fd::new(fds[0], "pipe").unwrap();
let _w = Fd::new(fds[1], "pipe").unwrap();
let mut drain: DrainState<fn(&[u8]) -> bool> =
DrainState::new(None, None, Some(r), None, 1024, None, None, false).unwrap();
let mut reactor = Reactor::new().unwrap();
drain.register_with_reactor(&mut reactor).unwrap();
assert_eq!(drain.pty_input_token(), None);
let err = drain.set_pty_writable(&mut reactor, true).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_managed_pty_input_token_and_writable_forward() {
let mut managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"read line; echo got:$line".to_string(),
],
SpawnBackend::Fork,
)
.pgroup(ProcessGroup::new(None, true))
.pty()
.capture_stdout()
.timeout_ms(10_000)
.kill_grace_ms(200)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let mut reactor = Reactor::new().unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
assert!(managed.pty_input_token().is_some());
assert!(managed.set_pty_writable(&mut reactor, true).is_ok());
assert!(managed.set_pty_writable(&mut reactor, false).is_ok());
assert_eq!(managed.write_input_nonblock(b"hello\r").unwrap(), Some(6));
let mut events = Vec::new();
let output = loop {
if let Some(output) = managed.poll_completion(&mut reactor).unwrap() {
break output;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout).unwrap();
for event in &events {
managed.handle_reactor_event(&mut reactor, event).unwrap();
}
};
assert!(!output.timed_out);
let merged = String::from_utf8_lossy(&output.stdout);
assert!(
merged.contains("got:hello"),
"child must read the written line, got {merged:?}"
);
}
#[test]
fn test_spawn_output_limit_is_combined_and_errors() {
let err = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"printf 12345; printf 67890 >&2".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.capture_stderr()
.max_output(8)
.build()
.unwrap()
.run()
.unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EOVERFLOW));
}
#[test]
fn test_spawn_output_exact_limit_eof_is_not_overflow() {
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"printf 12345".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.max_output(5)
.build()
.unwrap()
.run()
.unwrap();
assert_eq!(out.stdout, b"12345");
}
#[test]
fn test_spawn_zero_output_limit_without_output_is_not_overflow() {
let out = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::PosixSpawn)
.capture_stdout()
.capture_stderr()
.max_output(0)
.build()
.unwrap()
.run()
.unwrap();
assert!(out.stdout.is_empty());
assert!(out.stderr.is_empty());
}
fn drive_managed(
mut managed: crate::spawn::ManagedProcess,
cancel: bool,
) -> Result<crate::spawn::Output, crate::CoreError> {
let mut reactor = Reactor::new().unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
if cancel {
managed.request_cancel();
}
let mut events = Vec::new();
loop {
if let Some(output) = managed.poll_completion(&mut reactor)? {
return Ok(output);
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout)?;
for event in &events {
managed.handle_reactor_event(&mut reactor, event)?;
}
}
}
#[test]
fn test_spawn_managed_returns_full_output() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/echo".to_string(), "managed".to_string()],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.capture_stderr()
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert_eq!(output.status, Some(ExitStatus::Exited(0)));
assert_eq!(output.stdout, b"managed\n");
assert!(output.stderr.is_empty());
assert!(!output.timed_out);
}
#[test]
fn test_spawn_managed_closed_stdio_fds_relocated_and_full_output() {
let mut pipe = [0i32; 2];
assert_eq!(unsafe { libc::pipe(pipe.as_mut_ptr()) }, 0);
let pid = unsafe { libc::fork() };
assert!(pid >= 0, "fork failed");
if pid == 0 {
unsafe {
libc::close(pipe[0]);
}
let _ = unsafe { libc::close(0) };
let _ = unsafe { libc::close(1) };
let _ = unsafe { libc::close(2) };
let result = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"i=0; while [ $i -lt 4000 ]; do echo chunk-$i; i=$((i+1)); done".to_string(),
],
SpawnBackend::Fork,
)
.capture_stdout()
.build()
.unwrap(),
)
.and_then(|managed| drive_managed(managed, false));
let ok = matches!(result, Ok(output)
if output.status == Some(ExitStatus::Exited(0))
&& output.stdout.starts_with(b"chunk-0\n")
&& output.stdout.ends_with(b"chunk-3999\n"));
let buf = [ok as i8, 0];
unsafe {
libc::write(pipe[1], buf.as_ptr().cast(), buf.len());
libc::close(pipe[1]);
libc::_exit(if ok { 0 } else { 1 });
}
}
unsafe {
libc::close(pipe[1]);
}
let mut got = [0i8; 2];
let mut off = 0usize;
while off < got.len() {
let n = unsafe { libc::read(pipe[0], got[off..].as_mut_ptr().cast(), got.len() - off) };
if n <= 0 {
break;
}
off += n as usize;
}
unsafe {
libc::close(pipe[0]);
}
let mut wait_status = 0;
unsafe {
libc::waitpid(pid, &mut wait_status, 0);
}
assert_eq!(got[0], 1, "child spawn failed or stdout truncated");
}
#[test]
fn test_spawn_managed_exact_limit_is_not_overflow() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/echo".to_string(), "1234".to_string()],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.max_output(5)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert_eq!(output.stdout, b"1234\n");
}
#[test]
fn test_spawn_managed_reports_output_overflow() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/echo".to_string(), "12345".to_string()],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.max_output(5)
.build()
.unwrap(),
)
.unwrap();
let err = drive_managed(managed, false).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EOVERFLOW));
}
#[test]
fn test_spawn_managed_timeout_kills_own_process_group() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/sleep".to_string(), "10".to_string()],
SpawnBackend::PosixSpawn,
)
.pgroup(ProcessGroup::new(Some(0), false))
.timeout_ms(20)
.kill_grace_ms(20)
.cancel(CancelPolicy::Graceful)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(output.timed_out);
assert!(matches!(output.status, Some(ExitStatus::Signaled(_))));
}
#[test]
fn test_spawn_managed_explicit_cancel_reaps_child() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/sleep".to_string(), "10".to_string()],
SpawnBackend::PosixSpawn,
)
.pgroup(ProcessGroup::new(Some(0), false))
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, true).unwrap();
assert!(!output.timed_out);
assert!(matches!(output.status, Some(ExitStatus::Signaled(_))));
}
#[test]
fn test_managed_pid_available_after_completion() {
let managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/echo".to_string(), "managed".to_string()],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.build()
.unwrap(),
)
.unwrap();
let pid = managed.pid();
assert!(pid > 0);
let output = drive_managed(managed, false).unwrap();
assert_eq!(output.pid, pid);
}
#[test]
fn test_managed_register_twice_is_idempotent() {
let mut managed = spawn_managed(
SpawnOptions::builder(
vec!["/bin/echo".to_string(), "managed".to_string()],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.build()
.unwrap(),
)
.unwrap();
let mut reactor = Reactor::new().unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
let mut events = Vec::new();
loop {
if let Some(output) = managed.poll_completion(&mut reactor).unwrap() {
assert_eq!(output.stdout, b"managed\n");
return;
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout).unwrap();
for event in &events {
managed.handle_reactor_event(&mut reactor, event).unwrap();
}
}
}
#[test]
fn test_spawn_chunk_sink_forwards_chunks_streaming() {
let received_stdout = Arc::new(Mutex::new(Vec::<u8>::new()));
let received_stderr = Arc::new(Mutex::new(Vec::<u8>::new()));
let r_out = Arc::clone(&received_stdout);
let r_err = Arc::clone(&received_stderr);
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"printf hello; printf world; printf err >&2".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.capture_stderr()
.chunk_sink(move |is_stdout, bytes| {
if is_stdout {
r_out.lock().unwrap().extend_from_slice(bytes);
} else {
r_err.lock().unwrap().extend_from_slice(bytes);
}
SinkResult::Accept
})
.build()
.unwrap()
.run()
.unwrap();
assert_eq!(&*received_stdout.lock().unwrap(), b"helloworld");
assert_eq!(&*received_stderr.lock().unwrap(), b"err");
assert!(out.stdout.is_empty());
assert!(out.stderr.is_empty());
assert!(out.stdout_pending.is_none());
assert!(out.stderr_pending.is_none());
assert_eq!(out.status, Some(ExitStatus::Exited(0)));
}
#[test]
fn test_spawn_chunk_sink_pause_resume_lossless_with_consumer() {
const CAP: usize = 1; let outstanding = Arc::new(AtomicUsize::new(0));
let received = Arc::new(Mutex::new(Vec::<u8>::new()));
let stop = Arc::new(AtomicBool::new(false));
let sink_out = Arc::clone(&outstanding);
let sink_recv = Arc::clone(&received);
let sink = move |_is_stdout, bytes: &[u8]| {
if sink_out.load(Ordering::SeqCst) >= CAP {
return SinkResult::Pause;
}
sink_out.fetch_add(1, Ordering::SeqCst);
sink_recv.lock().unwrap().extend_from_slice(bytes);
SinkResult::Accept
};
let consumer_out = Arc::clone(&outstanding);
let consumer_stop = Arc::clone(&stop);
let consumer = std::thread::spawn(move || {
while !consumer_stop.load(Ordering::SeqCst) {
consumer_out.store(0, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(1));
}
});
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"dd if=/dev/zero bs=1 count=300000 2>/dev/null | tr '\\000' x".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.chunk_sink(sink)
.build()
.unwrap()
.run()
.unwrap();
if let Some(p) = out.stdout_pending {
let pending = p;
loop {
if outstanding.load(Ordering::SeqCst) < CAP {
outstanding.fetch_add(1, Ordering::SeqCst);
received.lock().unwrap().extend_from_slice(&pending);
break;
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
}
stop.store(true, Ordering::SeqCst);
consumer.join().unwrap();
let got = received.lock().unwrap();
assert_eq!(got.len(), 300000);
assert!(got.iter().all(|&b| b == b'x'));
assert!(out.stdout.is_empty());
}
#[test]
fn test_spawn_chunk_sink_pause_at_end_pending_flushed() {
let received = Arc::new(Mutex::new(Vec::<u8>::new()));
let out = SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"printf helddata".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.chunk_sink(move |_is_stdout, _bytes| SinkResult::Pause)
.build()
.unwrap()
.run()
.unwrap();
assert!(received.lock().unwrap().is_empty());
assert_eq!(out.stdout_pending.as_deref(), Some(b"helddata".as_slice()));
assert!(out.stdout.is_empty());
assert_eq!(out.status, Some(ExitStatus::Exited(0)));
}
#[test]
fn test_managed_chunk_sink_pause_resume_lossless() {
const CAP: usize = 1;
let outstanding = Arc::new(AtomicUsize::new(0));
let received = Arc::new(Mutex::new(Vec::<u8>::new()));
let stop = Arc::new(AtomicBool::new(false));
let sink_out = Arc::clone(&outstanding);
let sink_recv = Arc::clone(&received);
let sink = move |_is_stdout, bytes: &[u8]| {
if sink_out.load(Ordering::SeqCst) >= CAP {
return SinkResult::Pause;
}
sink_out.fetch_add(1, Ordering::SeqCst);
sink_recv.lock().unwrap().extend_from_slice(bytes);
SinkResult::Accept
};
let consumer_out = Arc::clone(&outstanding);
let consumer_stop = Arc::clone(&stop);
let consumer = std::thread::spawn(move || {
while !consumer_stop.load(Ordering::SeqCst) {
consumer_out.store(0, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(1));
}
});
let mut managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"dd if=/dev/zero bs=1 count=300000 2>/dev/null | tr '\\000' x".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.chunk_sink(sink)
.build()
.unwrap(),
)
.unwrap();
let mut reactor = Reactor::new().unwrap();
managed.register_with_reactor(&mut reactor).unwrap();
let mut events = Vec::new();
let mut output = None;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(20);
while std::time::Instant::now() < deadline {
if let Some(out) = managed.poll_completion(&mut reactor).unwrap() {
output = Some(out);
break;
}
if managed.stdout_paused() {
let _ = managed.resume_stdout(&mut reactor);
}
if managed.stderr_paused() {
let _ = managed.resume_stderr(&mut reactor);
}
let timeout = managed
.next_deadline()
.map(|at| {
at.saturating_duration_since(std::time::Instant::now())
.as_millis()
.min(i32::MAX as u128) as i32
})
.unwrap_or(-1);
reactor.wait(&mut events, 16, timeout).unwrap();
for event in &events {
managed.handle_reactor_event(&mut reactor, event).unwrap();
}
}
stop.store(true, Ordering::SeqCst);
consumer.join().unwrap();
let out = output.expect("managed streaming child completed");
if let Some(p) = out.stdout_pending {
let pending = p;
loop {
if outstanding.load(Ordering::SeqCst) < CAP {
outstanding.fetch_add(1, Ordering::SeqCst);
received.lock().unwrap().extend_from_slice(&pending);
break;
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
}
let got = received.lock().unwrap();
assert_eq!(got.len(), 300000);
assert!(got.iter().all(|&b| b == b'x'));
}
#[test]
fn test_kill_rejects_zero_pid() {
let p = Process::new(0);
assert!(p.kill(libc::SIGTERM).is_err());
}
#[test]
fn test_kill_group_rejects_nonpositive_pgid() {
let p = Process::new(1);
assert!(p.kill_group(0, libc::SIGTERM).is_err());
assert!(p.kill_group(-1, libc::SIGTERM).is_err());
assert!(p.kill_group(99999999, libc::SIGTERM).is_ok());
}
#[test]
fn test_fork_rejects_isolated_with_custom_leader() {
let opts = SpawnOptions::builder(vec!["/bin/true".to_string()], SpawnBackend::Fork)
.pgroup(ProcessGroup::new(Some(12345), true))
.build()
.unwrap();
let err = match spawn_start(opts) {
Ok(_) => panic!("expected isolated+custom-leader to fail"),
Err(err) => err,
};
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_managed_timeout_with_overflow_returns_partial() {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"printf 1234567890; sleep 30 & exec sleep 30".to_string(),
],
SpawnBackend::PosixSpawn,
)
.capture_stdout()
.max_output(5)
.timeout_ms(100)
.kill_grace_ms(50)
.cancel(CancelPolicy::Kill)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
assert!(output.timed_out);
assert!(matches!(output.status, Some(ExitStatus::Signaled(_))));
}
#[test]
fn test_writer_state_epipe() {
use crate::io::writer::WriterState;
let mut fds = [0; 2];
unsafe { libc::pipe(fds.as_mut_ptr()) };
let r = Fd::new(fds[0], "pipe").unwrap();
let w = Fd::new(fds[1], "pipe").unwrap();
let mut writer = WriterState::new(Some(vec![0u8; 1024 * 1024].into_boxed_slice()));
drop(r);
let mut last_res = Ok(false);
for _ in 0..100 {
last_res = writer.write_to_fd(&w);
if last_res.is_err() || (last_res.is_ok() && writer.buf.is_none()) {
break;
}
}
assert!(last_res.is_ok());
assert!(last_res.unwrap());
assert!(writer.buf.is_none());
}
#[test]
fn test_path_existence() {
let temp_file = std::env::temp_dir().join("coreshift_test_path");
let path_str = temp_file.to_str().unwrap();
std::fs::write(&temp_file, "test").unwrap();
assert!(path_exists(path_str));
assert!(path_lstat_exists(path_str));
std::fs::remove_file(&temp_file).unwrap();
assert!(!path_exists(path_str));
assert!(!path_lstat_exists(path_str));
}
#[test]
fn test_path_uid_temp_file() {
let path = std::env::temp_dir().join(format!(
"coreshift_test_uid_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::write(&path, b"uid").unwrap();
let uid = path_uid(&path).unwrap();
assert_eq!(uid, unsafe { libc::geteuid() });
remove_file(&path).unwrap();
}
#[test]
fn test_path_stat_reports_identity_fields() {
let path = std::env::temp_dir().join(format!(
"coreshift_test_path_stat_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::write(&path, b"identity").unwrap();
let stat = path_stat(&path).unwrap();
assert_eq!(stat.uid, unsafe { libc::geteuid() });
assert!(stat.inode > 0);
remove_file(&path).unwrap();
}
#[test]
fn test_path_stat_follows_symlink_and_lstat_reports_link() {
let dir = std::env::temp_dir().join(format!(
"coreshift_test_symlink_stat_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir(&dir).unwrap();
let target = dir.join("target");
let link = dir.join("link");
std::fs::write(&target, b"identity").unwrap();
std::os::unix::fs::symlink(&target, &link).unwrap();
let target_stat = path_stat(&target).unwrap();
let followed = path_stat(&link).unwrap();
let link_stat = path_lstat(&link).unwrap();
assert_eq!(followed.inode, target_stat.inode);
assert_ne!(link_stat.inode, target_stat.inode);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_clock_ticks_per_second_checked_result() {
let ticks = clock_ticks_per_second().unwrap();
assert!(ticks > 0);
}
#[test]
fn test_proc_uid_current_process() {
let uid = uid(std::process::id() as i32).unwrap();
assert_eq!(uid, unsafe { libc::geteuid() });
}
#[test]
fn test_proc_uid_at_uses_explicit_root() {
let proc_root = std::env::temp_dir().join(format!(
"coreshift_test_proc_uid_root_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let pid_dir = proc_root.join("54321");
std::fs::create_dir_all(&pid_dir).unwrap();
let uid = uid_at(&proc_root, 54321).unwrap();
assert_eq!(uid, unsafe { libc::geteuid() });
let _ = std::fs::remove_dir_all(&proc_root);
}
#[test]
fn test_proc_stat_at_uses_explicit_root() {
let proc_root = std::env::temp_dir().join(format!(
"coreshift_test_proc_stat_root_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let pid_dir = proc_root.join("54321");
std::fs::create_dir_all(&pid_dir).unwrap();
let stat = stat_at(&proc_root, 54321).unwrap();
assert_eq!(stat.uid, unsafe { libc::geteuid() });
assert!(stat.inode > 0);
let _ = std::fs::remove_dir_all(&proc_root);
}
#[test]
fn test_proc_uid_invalid_pid_returns_error() {
let err = uid(999_999).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::ENOENT));
}
#[test]
fn test_path_uid_missing_path_returns_error() {
let path = std::env::temp_dir().join(format!(
"coreshift_missing_uid_{}_{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let err = path_uid(&path).unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::ENOENT));
}
#[test]
fn test_install_shutdown_flag_compiles_and_sets_helper_state() {
let _lock = SIGNAL_TEST_LOCK.lock().unwrap();
TEST_SHUTDOWN_FLAG_A.store(false, Ordering::Release);
install_shutdown_flag(&TEST_SHUTDOWN_FLAG_A).unwrap();
assert!(!shutdown_requested(&TEST_SHUTDOWN_FLAG_A));
TEST_SHUTDOWN_FLAG_A.store(true, Ordering::Release);
assert!(shutdown_requested(&TEST_SHUTDOWN_FLAG_A));
TEST_SHUTDOWN_FLAG_A.store(false, Ordering::Release);
}
#[test]
fn test_install_shutdown_flag_repeated_install_does_not_panic() {
let _lock = SIGNAL_TEST_LOCK.lock().unwrap();
TEST_SHUTDOWN_FLAG_A.store(false, Ordering::Release);
TEST_SHUTDOWN_FLAG_B.store(false, Ordering::Release);
install_shutdown_flag(&TEST_SHUTDOWN_FLAG_A).unwrap();
install_shutdown_flag(&TEST_SHUTDOWN_FLAG_B).unwrap();
assert!(!shutdown_requested(&TEST_SHUTDOWN_FLAG_A));
assert!(!shutdown_requested(&TEST_SHUTDOWN_FLAG_B));
}
#[test]
fn test_spawn_resets_inherited_ignored_signals() {
let _lock = SIGNAL_TEST_LOCK.lock().unwrap();
let saved = unsafe {
[
libc::signal(libc::SIGHUP, libc::SIG_IGN),
libc::signal(libc::SIGINT, libc::SIG_IGN),
libc::signal(libc::SIGQUIT, libc::SIG_IGN),
]
};
assert_ne!(saved[0], libc::SIG_ERR);
assert_ne!(saved[1], libc::SIG_ERR);
assert_ne!(saved[2], libc::SIG_ERR);
for backend in [SpawnBackend::Clone3Pidfd, SpawnBackend::Fork] {
let managed = spawn_managed(
SpawnOptions::builder(
vec![
"/bin/sh".to_string(),
"-c".to_string(),
"kill -INT $$; echo SURVIVED".to_string(),
],
backend,
)
.capture_stdout()
.timeout_ms(5000)
.cancel(CancelPolicy::Graceful)
.build()
.unwrap(),
)
.unwrap();
let output = drive_managed(managed, false).unwrap();
let stdout = String::from_utf8_lossy(&output.stdout).to_string();
assert!(
!stdout.contains("SURVIVED"),
"{backend:?}: child inherited SIG_IGN for SIGINT (SURVIVED printed)"
);
assert!(
matches!(output.status, Some(ExitStatus::Signaled(libc::SIGINT))),
"{backend:?}: expected child killed by SIGINT, got {:?}",
output.status
);
}
unsafe {
libc::signal(libc::SIGHUP, saved[0]);
libc::signal(libc::SIGINT, saved[1]);
libc::signal(libc::SIGQUIT, saved[2]);
}
}
#[test]
fn test_shutdown_flag_guard_restores_previous_handler() {
let _lock = SIGNAL_TEST_LOCK.lock().unwrap();
let mut ignore_action: libc::sigaction = unsafe { std::mem::zeroed() };
let mut old_action: libc::sigaction = unsafe { std::mem::zeroed() };
ignore_action.sa_sigaction = libc::SIG_IGN;
unsafe { libc::sigemptyset(&mut ignore_action.sa_mask) };
let ret = unsafe { libc::sigaction(libc::SIGTERM, &ignore_action, &mut old_action) };
assert_eq!(ret, 0);
{
let _guard = install_shutdown_flag_guard(&TEST_SHUTDOWN_FLAG_A).unwrap();
}
let mut current_action: libc::sigaction = unsafe { std::mem::zeroed() };
let ret = unsafe { libc::sigaction(libc::SIGTERM, std::ptr::null(), &mut current_action) };
assert_eq!(ret, 0);
assert_eq!(current_action.sa_sigaction, libc::SIG_IGN);
let ret = unsafe { libc::sigaction(libc::SIGTERM, &old_action, std::ptr::null_mut()) };
assert_eq!(ret, 0);
}
#[test]
fn test_reactor_setup_signalfd_restores_previous_mask_on_drop() {
let mut before: libc::sigset_t = unsafe { std::mem::zeroed() };
let ret = unsafe { libc::pthread_sigmask(libc::SIG_SETMASK, std::ptr::null(), &mut before) };
assert_eq!(ret, 0);
let was_blocked = unsafe { libc::sigismember(&before, libc::SIGCHLD) };
{
let mut reactor = Reactor::new().unwrap();
reactor.setup_signalfd().unwrap();
let mut during: libc::sigset_t = unsafe { std::mem::zeroed() };
let ret =
unsafe { libc::pthread_sigmask(libc::SIG_SETMASK, std::ptr::null(), &mut during) };
assert_eq!(ret, 0);
assert_eq!(unsafe { libc::sigismember(&during, libc::SIGCHLD) }, 1);
}
let mut after: libc::sigset_t = unsafe { std::mem::zeroed() };
let ret = unsafe { libc::pthread_sigmask(libc::SIG_SETMASK, std::ptr::null(), &mut after) };
assert_eq!(ret, 0);
assert_eq!(
unsafe { libc::sigismember(&after, libc::SIGCHLD) },
was_blocked
);
}
#[test]
fn test_reactor_setup_signalfd_rejects_repeated_setup() {
let mut reactor = Reactor::new().unwrap();
reactor.setup_signalfd().unwrap();
let err = reactor.setup_signalfd().unwrap_err();
assert_eq!(err.raw_os_error(), Some(libc::EINVAL));
}
#[test]
fn test_reactor_signalfd_drop_on_other_thread_does_not_clobber_mask() {
let mut reactor = Reactor::new().unwrap();
reactor.setup_signalfd().unwrap();
let result = std::sync::Arc::new(std::sync::Mutex::new(None));
let result_clone = std::sync::Arc::clone(&result);
let handle = std::thread::spawn(move || {
let mut sentinel: libc::sigset_t = unsafe { std::mem::zeroed() };
unsafe { libc::sigemptyset(&mut sentinel) };
unsafe { libc::sigaddset(&mut sentinel, libc::SIGUSR2) };
let ret =
unsafe { libc::pthread_sigmask(libc::SIG_SETMASK, &sentinel, std::ptr::null_mut()) };
assert_eq!(ret, 0);
drop(reactor);
let mut current: libc::sigset_t = unsafe { std::mem::zeroed() };
let ret =
unsafe { libc::pthread_sigmask(libc::SIG_SETMASK, std::ptr::null(), &mut current) };
assert_eq!(ret, 0);
let sentinel_kept = unsafe { libc::sigismember(¤t, libc::SIGUSR2) };
*result_clone.lock().unwrap() = Some(sentinel_kept);
});
handle.join().unwrap();
assert_eq!(result.lock().unwrap().take(), Some(1));
}
#[test]
fn test_readahead_small_temp_file() {
with_temp_readahead_file(|file, _| {
assert_readahead_result(readahead(file, 0, 16));
});
}
#[test]
fn test_readahead_zero_length() {
with_temp_readahead_file(|file, _| {
assert_readahead_result(readahead(file, 0, 0));
});
}
#[test]
fn test_readahead_offset_beyond_eof() {
with_temp_readahead_file(|file, _| {
assert_readahead_result(readahead(file, 1 << 20, 16));
});
}
#[test]
fn test_readahead_invalid_fd() {
let result = readahead(RawFdRef(-1), 0, 16);
match result {
Err(err) if err.raw_os_error() == Some(libc::EBADF) => {}
Err(err) if err.raw_os_error() == Some(libc::ENOSYS) => {
eprintln!("skipping readahead test: unsupported on this target");
}
Err(err) => panic!("expected EBADF from invalid fd, got: {err}"),
Ok(()) => panic!("expected invalid fd to fail"),
}
}
#[test]
fn test_mmap_madvise_offset_zero() {
with_temp_readahead_file(|file, _| {
let result = mmap_madvise(file, 0, 16, false);
if let Err(err) = &result
&& err.raw_os_error() == Some(libc::ENOSYS)
{
eprintln!("skipping mmap_madvise test: unsupported on this target");
return;
}
result.unwrap();
});
}
#[test]
fn test_mmap_madvise_rejects_unaligned_offset() {
with_temp_readahead_file(|file, _| {
let result = mmap_madvise(file, 1, 16, false);
assert_eq!(result.unwrap_err().raw_os_error(), Some(libc::EINVAL));
});
}