use std::collections::BTreeMap;
use std::ffi::{CString, OsStr, OsString};
use std::io;
use std::os::fd::OwnedFd;
use std::os::raw::c_char;
use std::os::unix::prelude::OsStrExt;
use std::os::unix::process::ExitStatusExt;
use std::process::ExitStatus;
use libc::execvp;
use nix::errno::Errno;
use nix::fcntl::OFlag;
use nix::sys::wait::{waitpid, WaitPidFlag, WaitStatus};
use nix::unistd::{fork, pipe2, read, write, ForkResult, Pid as NixPid};
use crate::Pid;
unsafe extern "C" {
static mut environ: *mut *mut c_char;
}
#[derive(Debug)]
pub struct SuspendedLaunchedProcess {
pid: NixPid,
public_pid: Pid,
pipes: Option<SuspendPipes>,
}
#[derive(Debug)]
struct SuspendPipes {
resume_tx: OwnedFd,
exec_error_rx: OwnedFd,
}
impl SuspendedLaunchedProcess {
pub fn launch_in_suspended_state(
command_name: &OsStr,
command_args: &[OsString],
env_vars: &[(OsString, OsString)],
) -> crate::Result<Self> {
let argv_strings: Vec<CString> = std::iter::once(command_name)
.chain(command_args.iter().map(OsString::as_os_str))
.map(cstring_from_os_str)
.collect::<io::Result<_>>()?;
let argv: Vec<*const c_char> = null_terminated_ptrs(&argv_strings);
let envp_strings = (!env_vars.is_empty())
.then(|| build_env(env_vars))
.transpose()?;
let envp: Option<Vec<*const c_char>> = envp_strings.as_deref().map(null_terminated_ptrs);
let (resume_rp, resume_sp) = pipe2(OFlag::O_CLOEXEC).map_err(nix_error)?;
let (execerr_rp, execerr_sp) = pipe2(OFlag::O_CLOEXEC).map_err(nix_error)?;
match unsafe { fork() }.map_err(nix_error)? {
ForkResult::Child => {
drop((resume_sp, execerr_rp));
Self::run_child(resume_rp, execerr_sp, &argv, envp.as_deref())
}
ForkResult::Parent { child } => {
drop((resume_rp, execerr_sp));
let Some(public_pid) = Pid::new(child.as_raw()) else {
drop((resume_sp, execerr_rp));
reap(child);
return Err(io::Error::other("fork returned a non-positive child pid").into());
};
Ok(Self {
pid: child,
public_pid,
pipes: Some(SuspendPipes {
resume_tx: resume_sp,
exec_error_rx: execerr_rp,
}),
})
}
}
}
pub fn pid(&self) -> Pid {
self.public_pid
}
const EXECERR_MSG_FOOTER: [u8; 4] = *b"NOEX";
pub fn unsuspend_and_run(mut self) -> crate::Result<RunningProcess> {
let result = self.unsuspend_inner().map_err(crate::Error::from);
if result.is_err() {
reap(self.pid);
}
result
}
fn unsuspend_inner(&mut self) -> io::Result<RunningProcess> {
let SuspendPipes {
resume_tx,
exec_error_rx,
} = self
.pipes
.take()
.ok_or_else(|| io::Error::other("process was already resumed"))?;
write(&resume_tx, &[0x42])?;
drop(resume_tx);
loop {
let mut bytes = [0; 8];
match read(&exec_error_rx, &mut bytes) {
Ok(0) => break, Ok(8) => {
let [a, b, c, d, e, f, g, h] = bytes;
let errno_bytes = [a, b, c, d];
let footer = [e, f, g, h];
if footer != Self::EXECERR_MSG_FOOTER {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("invalid exec error pipe footer: {bytes:?}"),
));
}
return Err(io::Error::from_raw_os_error(i32::from_be_bytes(
errno_bytes,
)));
}
Ok(_) => {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"short read on exec error pipe",
));
}
Err(Errno::EINTR) => {}
Err(err) => return Err(err.into()),
}
}
Ok(RunningProcess {
state: ChildState::Running(self.pid),
})
}
fn run_child(
recv_end_of_resume_pipe: OwnedFd,
send_end_of_execerr_pipe: OwnedFd,
argv: &[*const c_char],
envp: Option<&[*const c_char]>,
) -> ! {
loop {
let mut buf = [0];
match read(&recv_end_of_resume_pipe, &mut buf) {
Ok(0) => Self::exit_child(0),
Ok(_) => {
let _ = unsafe {
match envp {
Some(envp) => {
environ = envp.as_ptr().cast_mut().cast();
execvp(argv[0], argv.as_ptr())
}
None => execvp(argv[0], argv.as_ptr()),
}
};
let [a, b, c, d] = Errno::last_raw().to_be_bytes();
let [e, f, g, h] = Self::EXECERR_MSG_FOOTER;
let _ = write(send_end_of_execerr_pipe, &[a, b, c, d, e, f, g, h]);
Self::exit_child(1)
}
Err(Errno::EINTR) => {}
Err(_) => Self::exit_child(1),
}
}
}
fn exit_child(status: i32) -> ! {
unsafe { libc::_exit(status) }
}
}
fn waitpid_retry(pid: NixPid, flags: Option<WaitPidFlag>) -> nix::Result<WaitStatus> {
loop {
match waitpid(pid, flags) {
Err(Errno::EINTR) => {}
result => return result,
}
}
}
fn nix_error(error: Errno) -> crate::Error {
io::Error::from(error).into()
}
fn reap(pid: NixPid) {
let _ = waitpid_retry(pid, None);
}
fn process_exit_status(status: WaitStatus) -> io::Result<ExitStatus> {
let raw = match status {
WaitStatus::Exited(_, code) => code << 8,
WaitStatus::Signaled(_, signal, dumped_core) => {
signal as i32 | if dumped_core { 0x80 } else { 0 }
}
_ => {
return Err(io::Error::other(format!(
"unexpected child status: {status:?}"
)))
}
};
Ok(ExitStatus::from_raw(raw))
}
impl Drop for SuspendedLaunchedProcess {
fn drop(&mut self) {
if self.pipes.take().is_none() {
return;
}
reap(self.pid);
}
}
fn cstring_from_os_str(os_str: &OsStr) -> io::Result<CString> {
CString::new(os_str.as_bytes()).map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidInput,
"nul byte found in command arguments",
)
})
}
#[must_use = "dropping without wait may leave the child running"]
pub struct RunningProcess {
state: ChildState,
}
#[derive(Clone, Copy, Debug)]
enum ChildState {
Running(NixPid),
Exited(ExitStatus),
Waited,
}
impl std::fmt::Debug for RunningProcess {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RunningProcess")
.field("state", &self.state)
.finish()
}
}
impl RunningProcess {
pub fn try_wait(&mut self) -> crate::Result<Option<ExitStatus>> {
let pid = match self.state {
ChildState::Running(pid) => pid,
ChildState::Exited(status) => return Ok(Some(status)),
ChildState::Waited => return Ok(None),
};
match waitpid_retry(pid, Some(WaitPidFlag::WNOHANG)) {
Ok(WaitStatus::StillAlive) => Ok(None),
Ok(status) => {
let status = process_exit_status(status)?;
self.state = ChildState::Exited(status);
Ok(Some(status))
}
Err(err) => Err(nix_error(err)),
}
}
pub fn wait(mut self) -> crate::Result<ExitStatus> {
match std::mem::replace(&mut self.state, ChildState::Waited) {
ChildState::Running(pid) => {
let status = waitpid_retry(pid, None).map_err(nix_error)?;
Ok(process_exit_status(status)?)
}
ChildState::Exited(status) => Ok(status),
ChildState::Waited => Err(io::Error::other("process was already waited").into()),
}
}
}
impl Drop for RunningProcess {
fn drop(&mut self) {
if let ChildState::Running(pid) = self.state {
let _ = waitpid_retry(pid, Some(WaitPidFlag::WNOHANG));
}
}
}
fn null_terminated_ptrs(strings: &[CString]) -> Vec<*const c_char> {
strings
.iter()
.map(|c| c.as_ptr())
.chain(std::iter::once(std::ptr::null()))
.collect()
}
fn build_env(env_vars: &[(OsString, OsString)]) -> io::Result<Vec<CString>> {
use std::os::unix::ffi::OsStringExt;
let mut vars: BTreeMap<OsString, OsString> = std::env::vars_os().collect();
for (name, val) in env_vars {
vars.insert(name.clone(), val.clone());
}
vars.into_iter()
.map(|(mut k, v)| {
k.reserve_exact(v.len() + 2);
k.push("=");
k.push(&v);
CString::new(k.into_vec()).map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidInput,
"nul byte found in environment variables",
)
})
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::TempDir;
use nix::sys::signal::Signal;
use std::os::unix::ffi::{OsStrExt, OsStringExt};
use std::os::unix::fs::symlink;
use std::thread;
use std::time::{Duration, Instant};
const ENV_HELPER: &str = "linux::process::tests::stackpulse_process_helper_env_probe";
const EXIT_HELPER: &str = "linux::process::tests::stackpulse_process_helper_exit_7";
const PATH_HELPER: &str = "linux::process::tests::stackpulse_process_helper_path_override";
const CHILD_PATH_ENV: &str = "STACKPULSE_CHILD_PATH";
const PATH_EXECUTABLE: &str = "stackpulse-child-path-executable";
fn current_test_binary() -> OsString {
std::env::current_exe()
.expect("current test binary")
.into_os_string()
}
fn ignored_test_args(test_name: &str) -> [OsString; 3] {
[
OsString::from("--ignored"),
OsString::from("--exact"),
OsString::from(test_name),
]
}
#[test]
fn dropping_suspended_launch_reaps_child() {
let launched =
SuspendedLaunchedProcess::launch_in_suspended_state(OsStr::new("unused"), &[], &[])
.expect("launch suspended child");
let pid = NixPid::from_raw(launched.pid().get());
drop(launched);
assert!(matches!(
waitpid(pid, Some(WaitPidFlag::WNOHANG)),
Err(Errno::ECHILD)
));
}
#[test]
fn failed_unsuspend_reaps_child() {
let launched = SuspendedLaunchedProcess::launch_in_suspended_state(
OsStr::new("/path/that/does/not/exist/stackpulse-bogus"),
&[],
&[],
)
.expect("launch suspended child");
let pid = NixPid::from_raw(launched.pid().get());
let result = launched.unsuspend_and_run();
assert!(result.is_err());
assert!(matches!(
waitpid(pid, Some(WaitPidFlag::WNOHANG)),
Err(Errno::ECHILD)
));
}
#[test]
fn suspended_launch_runs_command_with_environment_overrides() {
let command = current_test_binary();
let args = ignored_test_args(ENV_HELPER);
let launched = SuspendedLaunchedProcess::launch_in_suspended_state(
command.as_os_str(),
&args,
&[(OsString::from("STACKPULSE_TEST_ENV"), OsString::from("ok"))],
)
.expect("launch suspended child");
let running = launched.unsuspend_and_run().expect("resume child");
let status = running.wait().expect("wait child");
assert!(status.success());
}
#[test]
fn suspended_launch_resolves_commands_with_the_child_path() {
let caller_path = TempDir::new("process-caller-path");
let executable_dir = TempDir::new("process-child-path");
symlink(
current_test_binary(),
executable_dir.path().join(PATH_EXECUTABLE),
)
.expect("create child PATH executable");
let args = ignored_test_args(PATH_HELPER);
let launched = SuspendedLaunchedProcess::launch_in_suspended_state(
current_test_binary().as_os_str(),
&args,
&[
(
OsString::from("PATH"),
caller_path.path().as_os_str().to_owned(),
),
(
OsString::from(CHILD_PATH_ENV),
executable_dir.path().as_os_str().to_owned(),
),
],
)
.expect("launch PATH test helper");
let status = launched
.unsuspend_and_run()
.expect("resume PATH test helper")
.wait()
.expect("wait for PATH test helper");
assert!(status.success());
}
#[test]
fn running_process_reports_none_after_it_has_been_waited() {
let mut process = RunningProcess {
state: ChildState::Waited,
};
assert!(process.try_wait().expect("try wait without pid").is_none());
assert_eq!(
process.wait().unwrap_err().to_string(),
"process was already waited"
);
}
#[test]
fn exit_status_preserves_signal_and_core_dump() {
let pid = NixPid::from_raw(42);
let status = process_exit_status(WaitStatus::Signaled(pid, Signal::SIGTERM, true))
.expect("convert wait status");
assert_eq!(status.signal(), Some(libc::SIGTERM));
assert!(status.core_dumped());
}
#[test]
fn exit_status_rejects_nonterminal_wait_status() {
assert_eq!(
process_exit_status(WaitStatus::StillAlive)
.unwrap_err()
.to_string(),
"unexpected child status: StillAlive"
);
}
#[test]
fn try_wait_reports_missing_child() {
let mut process = RunningProcess {
state: ChildState::Running(NixPid::from_raw(i32::MAX)),
};
let error = process.try_wait().expect_err("missing child should fail");
process.state = ChildState::Waited;
assert_eq!(error.raw_os_error(), Some(libc::ECHILD));
}
#[test]
fn try_wait_caches_exited_process_status() {
let command = current_test_binary();
let args = ignored_test_args(EXIT_HELPER);
let launched =
SuspendedLaunchedProcess::launch_in_suspended_state(command.as_os_str(), &args, &[])
.expect("launch suspended child");
let mut running = launched.unsuspend_and_run().expect("resume child");
let deadline = Instant::now() + Duration::from_secs(5);
loop {
if let Some(status) = running.try_wait().expect("try wait child") {
assert_eq!(status.code(), Some(7));
assert_eq!(
running
.try_wait()
.expect("try wait reaped child")
.and_then(|status| status.code()),
Some(7)
);
assert_eq!(running.wait().expect("wait reaped child").code(), Some(7));
return;
}
if Instant::now() >= deadline {
if let ChildState::Running(pid) = running.state {
unsafe {
libc::kill(pid.as_raw(), libc::SIGKILL);
}
}
let _ = running.wait();
panic!("child did not exit");
}
thread::sleep(Duration::from_millis(10));
}
}
#[test]
fn cstring_conversions_reject_nul_bytes() {
assert!(cstring_from_os_str(OsStr::from_bytes(b"abc\0def")).is_err());
assert!(build_env(&[(
OsString::from_vec(b"BAD\0NAME".to_vec()),
OsString::from("x")
)])
.is_err());
assert!(build_env(&[(
OsString::from("BAD_VALUE"),
OsString::from_vec(b"x\0y".to_vec())
)])
.is_err());
}
#[test]
#[ignore]
fn stackpulse_process_helper_env_probe() {
assert_eq!(std::env::var("STACKPULSE_TEST_ENV").as_deref(), Ok("ok"));
}
#[test]
#[ignore]
fn stackpulse_process_helper_exit_7() {
std::process::exit(7);
}
#[test]
#[ignore]
fn stackpulse_process_helper_path_override() {
let child_path = std::env::var_os(CHILD_PATH_ENV).expect("child PATH");
let args = ignored_test_args(EXIT_HELPER);
let launched = SuspendedLaunchedProcess::launch_in_suspended_state(
OsStr::new(PATH_EXECUTABLE),
&args,
&[(OsString::from("PATH"), child_path)],
)
.expect("launch executable from child PATH");
let status = launched
.unsuspend_and_run()
.expect("resume child PATH executable")
.wait()
.expect("wait for child PATH executable");
assert_eq!(status.code(), Some(7));
}
}