#[cfg(unix)]
use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::io;
#[cfg(all(unix, feature = "process"))]
pub fn kill_process_group(pgid: i32) -> io::Result<()> {
if pgid <= 0 {
let refusal = Err(invalid_pgid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "kill_process_group: returning an error to the caller");
return refusal;
}
let pid = pid_from_raw(pgid)?;
rustix::process::kill_process_group(pid, rustix::process::Signal::KILL).map_err(errno_to_io)
}
#[cfg(all(not(unix), feature = "process"))]
pub fn kill_process_group(pgid: i32) -> io::Result<()> {
let _ = pgid;
Err(unsupported("process kill_process_group is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
pub fn process_group_exists(pgid: i32) -> io::Result<bool> {
if pgid <= 0 {
let refusal = Err(invalid_pgid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "process_group_exists: returning an error to the caller");
return refusal;
}
let pid = pid_from_raw(pgid)?;
match rustix::process::test_kill_process_group(pid) {
Ok(()) | Err(rustix::io::Errno::PERM) => Ok(true),
Err(rustix::io::Errno::SRCH) => Ok(false),
Err(errno) => Err(errno_to_io(errno)),
}
}
#[cfg(all(not(unix), feature = "process"))]
pub fn process_group_exists(pgid: i32) -> io::Result<bool> {
let _ = pgid;
Err(unsupported("process process_group_exists is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
pub fn child_has_exited_without_reaping(pid: i32) -> io::Result<bool> {
if pid <= 0 {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "child_has_exited_without_reaping: returning an error to the caller");
return refusal;
}
let Some(target) = rustix::process::Pid::from_raw(pid) else {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "child_has_exited_without_reaping: returning an error to the caller");
return refusal;
};
let observed = rustix::process::waitid(
rustix::process::WaitId::Pid(target),
rustix::process::WaitIdOptions::EXITED
| rustix::process::WaitIdOptions::NOHANG
| rustix::process::WaitIdOptions::NOWAIT,
)
.map_err(errno_to_io)?;
Ok(observed.is_some())
}
#[cfg(all(not(unix), feature = "process"))]
pub fn child_has_exited_without_reaping(pid: i32) -> io::Result<bool> {
let _ = pid;
Err(unsupported(
"process child_has_exited_without_reaping is Unix-only",
))
}
#[cfg(all(unix, feature = "process"))]
fn errno_to_io(errno: rustix::io::Errno) -> io::Error {
io::Error::from_raw_os_error(errno.raw_os_error())
}
#[cfg(all(unix, feature = "process"))]
fn pid_from_raw(pgid: i32) -> io::Result<rustix::process::Pid> {
rustix::process::Pid::from_raw(pgid).ok_or_else(invalid_pgid)
}
#[cfg(all(unix, feature = "process"))]
fn invalid_pgid() -> io::Error {
io::Error::new(
io::ErrorKind::InvalidInput,
"process group id must be a positive pid",
)
}
#[cfg(all(unix, feature = "process"))]
fn invalid_pid() -> io::Error {
io::Error::new(
io::ErrorKind::InvalidInput,
"process id must be a positive pid",
)
}
#[cfg(all(not(unix), feature = "process"))]
fn unsupported(message: &'static str) -> io::Error {
io::Error::new(io::ErrorKind::Unsupported, message)
}
pub const MAX_CAPTURED_DESCENDANTS: usize = 4096;
#[cfg(unix)]
const MAX_CAPTURED_DEPTH: usize = 64;
#[cfg(all(target_os = "linux", feature = "process"))]
const MAX_CAPTURED_THREADS: usize = 64;
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
#[non_exhaustive]
pub enum ContainmentMechanism {
#[default]
ProcessGroupOnly,
ProcChildrenTree,
ProcessTableSnapshot,
}
impl ContainmentMechanism {
#[must_use]
pub const fn strongest(self, other: Self) -> Self {
match (self, other) {
(Self::ProcChildrenTree, _) | (_, Self::ProcChildrenTree) => Self::ProcChildrenTree,
(Self::ProcessTableSnapshot, _) | (_, Self::ProcessTableSnapshot) => {
Self::ProcessTableSnapshot
}
(Self::ProcessGroupOnly, Self::ProcessGroupOnly) => Self::ProcessGroupOnly,
}
}
#[must_use]
pub const fn read_a_table(self) -> bool {
!matches!(self, Self::ProcessGroupOnly)
}
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
#[non_exhaustive]
pub struct DescendantSet {
mechanism: ContainmentMechanism,
pids: Vec<i32>,
truncated: bool,
}
impl DescendantSet {
#[must_use]
pub const fn mechanism(&self) -> ContainmentMechanism {
self.mechanism
}
#[must_use]
pub fn pids(&self) -> &[i32] {
&self.pids
}
#[must_use]
pub fn len(&self) -> usize {
self.pids.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.pids.is_empty()
}
#[must_use]
pub fn contains(&self, pid: i32) -> bool {
self.pids.binary_search(&pid).is_ok()
}
#[must_use]
pub const fn is_truncated(&self) -> bool {
self.truncated
}
pub fn absorb(&mut self, other: &Self) {
self.mechanism = self.mechanism.strongest(other.mechanism);
self.truncated = self.truncated || other.truncated;
for pid in &other.pids {
match self.pids.binary_search(pid) {
Ok(_) => {}
Err(at) if self.pids.len() < MAX_CAPTURED_DESCENDANTS => {
self.pids.insert(at, *pid);
}
Err(_) => self.truncated = true,
}
}
self.pids.sort_unstable();
self.pids.dedup();
}
}
#[cfg(all(unix, feature = "process"))]
pub fn capture_descendants(root: i32) -> io::Result<DescendantSet> {
if root <= 0 {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "capture_descendants: returning an error to the caller");
return refusal;
}
#[cfg(target_os = "linux")]
if let Some(set) = capture_proc_children(root) {
return Ok(set);
}
capture_table_snapshot(root)
}
#[cfg(all(not(unix), feature = "process"))]
pub fn capture_descendants(root: i32) -> io::Result<DescendantSet> {
let _ = root;
Err(unsupported("process capture_descendants is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
pub fn kill_process(pid: i32) -> io::Result<()> {
if pid <= 0 {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "kill_process: returning an error to the caller");
return refusal;
}
let target = pid_from_raw(pid)?;
rustix::process::kill_process(target, rustix::process::Signal::KILL).map_err(errno_to_io)
}
#[cfg(all(not(unix), feature = "process"))]
pub fn kill_process(pid: i32) -> io::Result<()> {
let _ = pid;
Err(unsupported("process kill_process is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
pub fn process_exists(pid: i32) -> io::Result<bool> {
if pid <= 0 {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "process_exists: returning an error to the caller");
return refusal;
}
let target = pid_from_raw(pid)?;
match rustix::process::test_kill_process(target) {
Ok(()) | Err(rustix::io::Errno::PERM) => Ok(true),
Err(rustix::io::Errno::SRCH) => Ok(false),
Err(errno) => Err(errno_to_io(errno)),
}
}
#[cfg(all(not(unix), feature = "process"))]
pub fn process_exists(pid: i32) -> io::Result<bool> {
let _ = pid;
Err(unsupported("process process_exists is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
pub fn running_processes(pids: &[i32]) -> io::Result<BTreeSet<i32>> {
if pids.is_empty() {
return Ok(BTreeSet::new());
}
let mut running = BTreeSet::new();
#[cfg(target_os = "linux")]
for pid in pids {
match read_proc_state(*pid) {
Some(state) => {
if state != PROC_ZOMBIE {
running.insert(*pid);
}
}
None => {
if process_exists(*pid)? {
running.insert(*pid);
}
}
}
}
#[cfg(not(target_os = "linux"))]
{
let snapshot = read_table_states()?;
for pid in pids {
match snapshot.get(pid) {
Some(state) => {
if !state.starts_with('Z') {
running.insert(*pid);
}
}
None => {
if process_exists(*pid)? {
running.insert(*pid);
}
}
}
}
}
Ok(running)
}
#[cfg(all(not(unix), feature = "process"))]
pub fn running_processes(pids: &[i32]) -> io::Result<BTreeSet<i32>> {
let _ = pids;
Err(unsupported("process running_processes is Unix-only"))
}
#[cfg(all(target_os = "linux", feature = "process"))]
const PROC_ZOMBIE: char = 'Z';
#[cfg(all(target_os = "linux", feature = "process"))]
fn read_proc_state(pid: i32) -> Option<char> {
let path = std::path::Path::new("/proc")
.join(pid.to_string())
.join("stat");
let stat = std::fs::read_to_string(path).ok()?;
let after_name = stat.get(stat.rfind(')')?.saturating_add(1)..)?;
after_name.split_ascii_whitespace().next()?.chars().next()
}
#[cfg(all(not(target_os = "linux"), unix, feature = "process"))]
fn read_table_states() -> io::Result<BTreeMap<i32, String>> {
let mut states = BTreeMap::new();
for row in read_ps_table(&["pid=", "stat="], "read_table_states")? {
let mut columns = row.split_ascii_whitespace();
let (Some(pid), Some(stat)) = (columns.next(), columns.next()) else {
continue;
};
let Ok(pid) = pid.parse::<i32>() else {
continue;
};
states.insert(pid, stat.to_owned());
}
Ok(states)
}
#[cfg(all(unix, feature = "process"))]
fn read_ps_table(columns: &[&str], caller: &'static str) -> io::Result<Vec<String>> {
let mut arguments = vec!["-A"];
for column in columns {
arguments.push("-o");
arguments.push(column);
}
let output = run_ps(&arguments, caller)?;
if !output.status.success() {
let refusal = Err(io::Error::other(format!(
"lgwks_std::process ({caller}): the process table could not be read with \
`ps`, so no process outside the supervised group could be named",
)));
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "read_ps_table: returning an error to the caller");
return refusal;
}
Ok(String::from_utf8_lossy(&output.stdout)
.lines()
.map(str::to_owned)
.collect())
}
#[cfg(all(unix, feature = "process"))]
fn run_ps(arguments: &[&str], caller: &'static str) -> io::Result<std::process::Output> {
let spawned = std::process::Command::new("ps")
.args(arguments)
.env("TZ", "UTC0")
.env("LC_ALL", "C")
.stdin(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.output();
match spawned {
Ok(output) => Ok(output),
Err(error) => {
let refusal = Err(io::Error::new(
error.kind(),
format!("lgwks_std::process ({caller}): `ps` could not be run: {error}"),
));
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "run_ps: returning an error to the caller");
refusal
}
}
}
#[cfg(all(target_os = "linux", feature = "process"))]
fn capture_proc_children(root: i32) -> Option<DescendantSet> {
read_proc_children(root).ok()?;
Some(walk_descendants(
root,
ContainmentMechanism::ProcChildrenTree,
|pid| match read_proc_children(pid) {
Ok(children) => Some(children),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Some(Vec::new()),
Err(error) => {
#[cfg(feature = "trace")]
crate::trace::debug!(
pid,
error = %error,
"descendant capture: a child list could not be read"
);
#[cfg(not(feature = "trace"))]
let _unreported = error;
None
}
},
))
}
#[cfg(all(target_os = "linux", feature = "process"))]
fn read_proc_children(pid: i32) -> std::io::Result<Vec<i32>> {
let task_dir = std::path::Path::new("/proc")
.join(pid.to_string())
.join("task");
let mut children = Vec::new();
let mut threads = 0_usize;
for entry in std::fs::read_dir(task_dir)? {
let entry = entry?;
threads = threads.saturating_add(1);
if threads > MAX_CAPTURED_THREADS {
break;
}
let Ok(text) = std::fs::read_to_string(entry.path().join("children")) else {
continue;
};
children.extend(
text.split_ascii_whitespace()
.filter_map(|pid| pid.parse::<i32>().ok()),
);
}
Ok(children)
}
#[cfg(all(unix, feature = "process"))]
fn walk_descendants(
root: i32,
mechanism: ContainmentMechanism,
mut children_of: impl FnMut(i32) -> Option<Vec<i32>>,
) -> DescendantSet {
let mut set = DescendantSet {
mechanism,
pids: Vec::new(),
truncated: false,
};
let mut frontier = vec![root];
let mut visited: BTreeSet<i32> = BTreeSet::from([root]);
for _ in 0..MAX_CAPTURED_DEPTH {
let mut next = Vec::new();
for parent in frontier {
let Some(listed) = children_of(parent) else {
set.truncated = true;
continue;
};
for child in listed {
if child <= 0 || !visited.insert(child) {
continue;
}
if set.pids.len() >= MAX_CAPTURED_DESCENDANTS {
set.truncated = true;
return finish(set);
}
set.pids.push(child);
next.push(child);
}
}
if next.is_empty() {
return finish(set);
}
frontier = next;
}
set.truncated = true;
finish(set)
}
#[cfg(all(unix, feature = "process"))]
fn finish(mut set: DescendantSet) -> DescendantSet {
set.pids.sort_unstable();
set.pids.dedup();
set
}
#[cfg(all(unix, feature = "process"))]
fn capture_table_snapshot(root: i32) -> io::Result<DescendantSet> {
let rows = read_ps_table(&["pid=", "ppid="], "capture_table_snapshot")?;
let mut children_of: BTreeMap<i32, Vec<i32>> = BTreeMap::new();
for row in &rows {
let mut columns = row.split_ascii_whitespace();
let (Some(pid), Some(parent)) = (columns.next(), columns.next()) else {
continue;
};
let (Ok(pid), Ok(parent)) = (pid.parse::<i32>(), parent.parse::<i32>()) else {
continue;
};
if pid <= 0 {
continue;
}
children_of.entry(parent).or_default().push(pid);
}
Ok(walk_descendants(
root,
ContainmentMechanism::ProcessTableSnapshot,
|pid| Some(children_of.remove(&pid).into_iter().flatten().collect()),
))
}
pub const MAX_START_TOKEN_BYTES: usize = 128;
const PROC_SCHEME: &str = "proc:";
const PS_SCHEME: &str = "ps:";
const IDENTITY_SEPARATOR: char = '/';
#[cfg(all(unix, feature = "process"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)]
enum StartScheme {
Proc,
Ps,
}
#[derive(Clone, Debug, Eq, PartialEq, Hash)]
pub struct ProcessIdentity {
pid: i32,
started: String,
}
impl ProcessIdentity {
pub fn new(pid: i32, started: &str) -> Result<Self, ProcessIdentityError> {
if pid <= 0 {
let refusal = Err(ProcessIdentityError::NonPositivePid { pid });
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "ProcessIdentity::new: returning an error to the caller");
return refusal;
}
validate_start(started)?;
Ok(Self {
pid,
started: started.to_owned(),
})
}
#[must_use]
pub const fn pid(&self) -> i32 {
self.pid
}
#[must_use]
pub fn started(&self) -> &str {
&self.started
}
#[cfg(all(unix, feature = "process"))]
fn scheme(&self) -> StartScheme {
if self.started.starts_with(PROC_SCHEME) {
StartScheme::Proc
} else {
StartScheme::Ps
}
}
}
impl std::fmt::Display for ProcessIdentity {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}{IDENTITY_SEPARATOR}{}", self.pid, self.started)
}
}
impl std::str::FromStr for ProcessIdentity {
type Err = ProcessIdentityError;
fn from_str(text: &str) -> Result<Self, Self::Err> {
let Some((pid, started)) = text.split_once(IDENTITY_SEPARATOR) else {
let refusal = Err(ProcessIdentityError::MissingSeparator);
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "ProcessIdentity::from_str: returning an error to the caller");
return refusal;
};
let canonical = !pid.is_empty()
&& pid.bytes().all(|byte| byte.is_ascii_digit())
&& (pid == "0" || !pid.starts_with('0'));
let Some(pid) = canonical.then(|| pid.parse::<i32>().ok()).flatten() else {
let refusal = Err(ProcessIdentityError::MalformedPid);
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "ProcessIdentity::from_str: returning an error to the caller");
return refusal;
};
Self::new(pid, started)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ProcessIdentityError {
MissingSeparator,
MalformedPid,
NonPositivePid {
pid: i32,
},
EmptyStart,
StartTooLong {
len: usize,
},
StartNotPrintable {
at: usize,
},
UnknownScheme,
}
impl std::fmt::Display for ProcessIdentityError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match *self {
Self::MissingSeparator => f.write_str("a process identity needs `<pid>/<start>`"),
Self::MalformedPid => f.write_str("a process identity's pid is not a decimal i32"),
Self::NonPositivePid { pid } => {
write!(f, "a process identity's pid must be positive, got {pid}")
}
Self::EmptyStart => f.write_str("a process identity's start token is empty"),
Self::StartTooLong { len } => write!(
f,
"a process identity's start token is {len} bytes, past the {MAX_START_TOKEN_BYTES}-byte bound"
),
Self::StartNotPrintable { at } => write!(
f,
"a process identity's start token holds a non-printable byte at offset {at}"
),
Self::UnknownScheme => {
f.write_str("a process identity's start token names no reading this crate produces")
}
}
}
}
impl std::error::Error for ProcessIdentityError {}
fn validate_start(started: &str) -> Result<(), ProcessIdentityError> {
let refusal = if started.is_empty() {
Some(ProcessIdentityError::EmptyStart)
} else if started.len() > MAX_START_TOKEN_BYTES {
Some(ProcessIdentityError::StartTooLong { len: started.len() })
} else {
started
.bytes()
.position(|byte| !byte.is_ascii_graphic())
.map(|at| ProcessIdentityError::StartNotPrintable { at })
};
if let Some(error) = refusal {
let refusal = Err(error);
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "validate_start: returning an error to the caller");
return refusal;
}
if started.starts_with(PROC_SCHEME) || started.starts_with(PS_SCHEME) {
Ok(())
} else {
let refusal = Err(ProcessIdentityError::UnknownScheme);
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "validate_start: returning an error to the caller");
refusal
}
}
#[cfg(all(unix, feature = "process"))]
pub fn identify_process(pid: i32) -> io::Result<Option<ProcessIdentity>> {
if pid <= 0 {
let refusal = Err(invalid_pid());
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "identify_process: returning an error to the caller");
return refusal;
}
#[cfg(target_os = "linux")]
let started = match proc_start(pid) {
Ok(started) => started,
Err(error) => {
#[cfg(feature = "trace")]
crate::trace::debug!(pid, error = %error, "identify_process: /proc was unreadable; reading ps");
#[cfg(not(feature = "trace"))]
let _unreported = error;
ps_start(pid)?
}
};
#[cfg(not(target_os = "linux"))]
let started = ps_start(pid)?;
Ok(started.map(|started| ProcessIdentity { pid, started }))
}
#[cfg(all(not(unix), feature = "process"))]
pub fn identify_process(pid: i32) -> io::Result<Option<ProcessIdentity>> {
let _ = pid;
Err(unsupported("process identify_process is Unix-only"))
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum OrphanReap {
Signalled {
descendants: DescendantSet,
},
LeaderGone,
LeaderReused {
holder: ProcessIdentity,
},
}
#[cfg(all(unix, feature = "process"))]
pub fn reap_orphaned_group(leader: &ProcessIdentity) -> io::Result<OrphanReap> {
let holder = match leader.scheme() {
StartScheme::Ps => ps_start(leader.pid)?,
#[cfg(target_os = "linux")]
StartScheme::Proc => proc_start(leader.pid)?,
#[cfg(not(target_os = "linux"))]
StartScheme::Proc => {
let refusal = Err(io::Error::new(
io::ErrorKind::Unsupported,
"a /proc process identity can only be read on Linux",
));
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "reap_orphaned_group: returning an error to the caller");
return refusal;
}
};
let Some(started) = holder else {
return Ok(OrphanReap::LeaderGone);
};
if started != leader.started {
return Ok(OrphanReap::LeaderReused {
holder: ProcessIdentity {
pid: leader.pid,
started,
},
});
}
let descendants = match capture_descendants(leader.pid) {
Ok(captured) => captured,
Err(error) => {
#[cfg(feature = "trace")]
crate::trace::debug!(pid = leader.pid, error = %error, "reap_orphaned_group: the descendants could not be captured");
#[cfg(not(feature = "trace"))]
let _unreported = error;
DescendantSet::default()
}
};
stopped_unless_refused(kill_process_group(leader.pid))?;
stopped_unless_refused(kill_process(leader.pid))?;
for pid in descendants.pids() {
stopped_unless_refused(kill_process(*pid))?;
}
Ok(OrphanReap::Signalled { descendants })
}
#[cfg(all(not(unix), feature = "process"))]
pub fn reap_orphaned_group(leader: &ProcessIdentity) -> io::Result<OrphanReap> {
let _ = leader;
Err(unsupported("process reap_orphaned_group is Unix-only"))
}
#[cfg(all(unix, feature = "process"))]
fn stopped_unless_refused(signalled: io::Result<()>) -> io::Result<()> {
match signalled {
Err(error) if error.raw_os_error() == Some(rustix::io::Errno::SRCH.raw_os_error()) => {
Ok(())
}
other => other,
}
}
#[cfg(all(target_os = "linux", feature = "process"))]
fn proc_start(pid: i32) -> io::Result<Option<String>> {
const START_FIELD_AFTER_NAME: usize = 19;
let boot = std::fs::read_to_string("/proc/sys/kernel/random/boot_id")?;
let stat_path = std::path::Path::new("/proc")
.join(pid.to_string())
.join("stat");
let stat = match std::fs::read_to_string(stat_path) {
Ok(stat) => stat,
Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
Err(error) => {
let refusal = Err(error);
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "proc_start: returning an error to the caller");
return refusal;
}
};
let ticks = stat
.rfind(')')
.and_then(|end| stat.get(end.saturating_add(1)..))
.and_then(|fields| fields.split_ascii_whitespace().nth(START_FIELD_AFTER_NAME));
let Some(ticks) = ticks else {
let refusal = Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("/proc/{pid}/stat carried no start tick"),
));
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "proc_start: returning an error to the caller");
return refusal;
};
Ok(Some(format!("{PROC_SCHEME}{}:{ticks}", boot.trim())))
}
#[cfg(all(unix, feature = "process"))]
fn ps_start(pid: i32) -> io::Result<Option<String>> {
let pid_text = pid.to_string();
let output = run_ps(&["-o", "lstart=", "-p", &pid_text], "ps_start")?;
let text = String::from_utf8_lossy(&output.stdout);
let words: Vec<&str> = text.split_ascii_whitespace().collect();
if words.is_empty() {
if process_exists(pid)? {
let refusal = Err(io::Error::other(format!(
"lgwks_std::process (ps_start): `ps` listed no start for live pid {pid}"
)));
#[cfg(feature = "trace")]
crate::trace::debug!(error = ?refusal.as_ref().err(), "ps_start: returning an error to the caller");
return refusal;
}
return Ok(None);
}
Ok(Some(format!("{PS_SCHEME}{}", words.join("-"))))
}
#[cfg(test)]
mod tests {
use std::fmt::Write as _;
#[cfg(unix)]
use std::os::unix::process::CommandExt as _;
use super::*;
#[test]
fn invalid_pgid_is_rejected() {
assert!(kill_process_group(0).is_err());
assert!(kill_process_group(-1).is_err());
}
#[test]
#[cfg(unix)]
fn nonexistent_group_surfs_os_error() {
assert!(kill_process_group(i32::MAX).is_err());
}
#[test]
#[cfg(unix)]
fn a_live_group_is_present_and_a_reaped_one_is_absent() -> Result<(), Box<dyn std::error::Error>>
{
use std::os::unix::process::CommandExt;
let mut child = std::process::Command::new("sleep")
.arg("30")
.process_group(0)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let group = i32::try_from(child.id())?;
assert!(
process_group_exists(group)?,
"the live child owns its group"
);
child.kill()?;
child.wait()?;
assert!(!process_group_exists(group)?, "the reaped group is absent");
Ok(())
}
#[test]
#[cfg(unix)]
fn the_callers_group_and_nonpositive_ids_are_not_probed() {
for refused in [0, -1] {
assert_eq!(
process_group_exists(refused).map_err(|error| error.kind()),
Err(std::io::ErrorKind::InvalidInput),
"group id {refused} names the caller or a process, not a group"
);
}
}
#[test]
fn invalid_pid_is_rejected_for_the_exit_observation() {
assert!(child_has_exited_without_reaping(0).is_err());
assert!(child_has_exited_without_reaping(-1).is_err());
}
#[test]
#[cfg(unix)]
fn a_non_child_pid_is_an_error_not_an_observation() {
assert!(child_has_exited_without_reaping(i32::MAX).is_err());
}
#[test]
#[cfg(unix)]
fn a_live_child_is_reported_not_yet_exited() -> Result<(), Box<dyn std::error::Error>> {
let mut child = std::process::Command::new("sleep")
.arg("30")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let pid = i32::try_from(child.id())?;
assert!(
!child_has_exited_without_reaping(pid)?,
"live child reported exited"
);
child.kill()?;
let status = child.wait()?;
assert!(
!status.success(),
"the killed probe child must report its signal, not success"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn an_exited_child_is_reported_exited_and_stays_reapable()
-> Result<(), Box<dyn std::error::Error>> {
let mut child = std::process::Command::new("true")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let pid = i32::try_from(child.id())?;
let mut observed_exit = false;
for _ in 0..100_000 {
if child_has_exited_without_reaping(pid)? {
observed_exit = true;
break;
}
std::thread::yield_now();
}
assert!(
observed_exit,
"an exited child must be observed within the bound"
);
let status = child.wait()?;
assert!(status.success(), "`true` must report success after reaping");
Ok(())
}
const VACANT: i32 = i32::MAX;
const ESRCH: i32 = 3;
#[cfg(unix)]
struct Tree {
leader: std::process::Child,
root: i32,
descendants: Vec<i32>,
dir: std::path::PathBuf,
}
#[cfg(unix)]
impl Tree {
fn grow(depth: usize) -> Result<Self, Box<dyn std::error::Error>> {
let dir = std::env::temp_dir().join(format!(
"lgwks-std-tree-{}-{}",
std::process::id(),
depth
));
let _ignored = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir)?;
let mut script = String::new();
for level in 0..depth {
let file = dir.join(format!("{level}.pid"));
let _written = write!(
script,
"sh -c 'echo $$ > {}; exec sleep 30' & ",
file.display()
);
}
script.push_str("exec sleep 30");
let leader = std::process::Command::new("sh")
.arg("-c")
.arg(&script)
.process_group(0)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let root = i32::try_from(leader.id())?;
let mut descendants = Vec::with_capacity(depth);
for level in 0..depth {
let file = dir.join(format!("{level}.pid"));
let mut pid = None;
for _ in 0..2000 {
if let Ok(text) = std::fs::read_to_string(&file) {
pid = text.trim().parse::<i32>().ok();
if pid.is_some() {
break;
}
}
std::thread::park_timeout(std::time::Duration::from_millis(5));
}
let pid = pid.ok_or_else(|| format!("level {level} never recorded its pid"))?;
descendants.push(pid);
}
Ok(Self {
leader,
root,
descendants,
dir,
})
}
}
#[cfg(unix)]
impl Drop for Tree {
fn drop(&mut self) {
let _ignored = kill_process_group(self.root);
let _reaped = self.leader.wait();
let _removed = std::fs::remove_dir_all(&self.dir);
}
}
#[test]
#[cfg(unix)]
fn a_non_positive_root_is_refused_before_any_table_is_read()
-> Result<(), Box<dyn std::error::Error>> {
for refused in [0, -1, i32::MIN] {
assert_eq!(
capture_descendants(refused).map_err(|error| error.kind()),
Err(std::io::ErrorKind::InvalidInput),
"root {refused} names the caller or no process, so no tree is walked"
);
}
Ok(())
}
#[test]
#[cfg(unix)]
fn a_non_positive_pid_is_refused_by_both_per_process_primitives()
-> Result<(), Box<dyn std::error::Error>> {
for refused in [0, -1, i32::MIN] {
assert_eq!(
process_exists(refused).map_err(|error| error.kind()),
Err(std::io::ErrorKind::InvalidInput),
"pid {refused} names the caller or no process, so nothing is probed"
);
assert_eq!(
kill_process(refused).map_err(|error| error.kind()),
Err(std::io::ErrorKind::InvalidInput),
"pid {refused} must be refused before any signal leaves the process"
);
}
Ok(())
}
#[test]
#[cfg(unix)]
fn a_vacant_pid_is_absent_and_signalling_it_says_esrch()
-> Result<(), Box<dyn std::error::Error>> {
assert!(
!process_exists(VACANT)?,
"an id above every pid ceiling names no process"
);
assert_eq!(
kill_process(VACANT).map_err(|error| error.raw_os_error()),
Err(Some(ESRCH)),
"signalling a vacant pid must say it does not exist, never succeed"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn the_capture_finds_every_descendant_of_a_live_root() -> Result<(), Box<dyn std::error::Error>>
{
let tree = Tree::grow(3)?;
let captured = capture_descendants(tree.root)?;
assert!(
captured.mechanism().read_a_table(),
"a live capture must name a mechanism that read a table, got {:?}",
captured.mechanism()
);
assert!(
!captured.is_truncated(),
"a depth-3 tree is inside every bound"
);
let mut expected = tree.descendants.clone();
expected.sort_unstable();
assert_eq!(
captured.pids(),
expected.as_slice(),
"the capture must be exactly the tree the pids describe"
);
assert!(
!captured.contains(tree.root),
"the root leads the group; a capture that listed it would double-signal the leader"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn a_descendant_that_left_the_group_is_still_captured() -> Result<(), Box<dyn std::error::Error>>
{
let dir = std::env::temp_dir().join(format!(
"lgwks-std-escape-{}-{}",
std::process::id(),
match std::thread::current().name() {
Some(name) => name.to_owned(),
None => format!("{:?}", std::thread::current().id()),
}
));
let _ignored = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir)?;
let pid_file = dir.join("escaped.pid");
let mut leader = std::process::Command::new("sh")
.arg("-c")
.arg(format!(
"python3 -c 'import os,time; os.setsid(); open(\"{}\",\"w\").write(str(os.getpid())); time.sleep(30)' & exec sleep 30",
pid_file.display()
))
.process_group(0)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let root = i32::try_from(leader.id())?;
let mut escaped = None;
for _ in 0..4000 {
if let Ok(text) = std::fs::read_to_string(&pid_file) {
escaped = text.trim().parse::<i32>().ok();
if escaped.is_some() {
break;
}
}
std::thread::park_timeout(std::time::Duration::from_millis(5));
}
let Some(escaped) = escaped else {
let _killed = kill_process_group(root);
let _reaped = leader.wait();
let _removed = std::fs::remove_dir_all(&dir);
return Err("this host cannot leave a process group, so nothing was exercised".into());
};
let captured = capture_descendants(root)?;
assert!(
captured.contains(escaped),
"a descendant that called setsid is still a child of the root, so the capture must \
hold pid {escaped}: captured {:?} of root {root}",
captured.pids()
);
assert!(
process_exists(escaped)?,
"the premise: the escapee is running and needs a pid signal to stop"
);
kill_process(escaped)?;
let _killed = kill_process_group(root);
let _reaped = leader.wait();
let _removed = std::fs::remove_dir_all(&dir);
Ok(())
}
#[test]
#[cfg(unix)]
fn a_signal_to_one_captured_pid_stops_that_pid_and_leaves_its_siblings()
-> Result<(), Box<dyn std::error::Error>> {
let tree = Tree::grow(2)?;
let captured = capture_descendants(tree.root)?;
let Some(first) = captured.pids().first().copied() else {
return Err("the capture held no descendant to signal".into());
};
kill_process(first)?;
assert!(
all_stop_running(&[first])?,
"the signalled pid {first} must have stopped running; it may remain in the \
table as an unreaped zombie, which is not a survivor"
);
let sibling = captured
.pids()
.iter()
.copied()
.find(|pid| *pid != first)
.ok_or("the capture held one descendant, so no sibling could survive")?;
assert_eq!(
running_processes(&[sibling])?
.into_iter()
.collect::<Vec<i32>>(),
vec![sibling],
"a per-process signal must reach exactly one process: pid {sibling} is untouched"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn a_capture_of_a_live_root_never_reports_a_truncated_tree_of_its_own_shape()
-> Result<(), Box<dyn std::error::Error>> {
let tree = Tree::grow(1)?;
let captured = capture_descendants(tree.root)?;
assert!(
!captured.is_truncated(),
"a one-descendant tree is inside every declared bound, so the flag must stay clear"
);
assert!(
captured.pids().len() <= MAX_CAPTURED_DESCENDANTS,
"a capture is bounded by MAX_CAPTURED_DESCENDANTS, got {}",
captured.pids().len()
);
Ok(())
}
#[test]
fn absorbing_two_captures_unions_them_and_keeps_the_stronger_mechanism() {
let mut first = DescendantSet {
mechanism: ContainmentMechanism::ProcessTableSnapshot,
pids: vec![30, 10],
truncated: true,
};
let second = DescendantSet {
mechanism: ContainmentMechanism::ProcChildrenTree,
pids: vec![10, 20],
truncated: false,
};
first.absorb(&second);
assert_eq!(
first.pids(),
[10, 20, 30].as_slice(),
"a merged capture is the sorted union of both"
);
assert_eq!(
first.mechanism(),
ContainmentMechanism::ProcChildrenTree,
"the merged mechanism must not understate what was read"
);
assert!(
first.is_truncated(),
"a truncation in either half survives the merge"
);
}
#[test]
fn a_merge_past_the_declared_ceiling_is_refused_and_reported()
-> Result<(), std::num::TryFromIntError> {
let mut first = DescendantSet {
mechanism: ContainmentMechanism::ProcChildrenTree,
pids: (0..MAX_CAPTURED_DESCENDANTS)
.map(i32::try_from)
.collect::<Result<_, _>>()?,
truncated: false,
};
let before = first.pids().len();
first.absorb(&DescendantSet {
mechanism: ContainmentMechanism::ProcChildrenTree,
pids: vec![i32::MAX],
truncated: false,
});
assert_eq!(
first.pids().len(),
before,
"a capture at the ceiling must not grow past it"
);
assert!(
first.is_truncated(),
"a refused id must leave the set reported as a prefix, not as the whole tree"
);
Ok(())
}
#[cfg(unix)]
fn all_stop_running(pids: &[i32]) -> Result<bool, Box<dyn std::error::Error>> {
for _ in 0..4000 {
if running_processes(pids)?.is_empty() {
return Ok(true);
}
std::thread::park_timeout(std::time::Duration::from_millis(5));
}
Ok(false)
}
#[test]
#[cfg(unix)]
fn a_live_process_reads_one_identity_that_round_trips_through_its_text()
-> Result<(), Box<dyn std::error::Error>> {
let pid = i32::try_from(std::process::id())?;
let first = identify_process(pid)?.ok_or("this process has no identity")?;
let second = identify_process(pid)?.ok_or("this process lost its identity")?;
assert_eq!(first, second, "one process read twice is one identity");
assert_eq!(first.pid(), pid);
let parsed: ProcessIdentity = first.to_string().parse()?;
assert_eq!(parsed, first, "the stored text names the same process");
Ok(())
}
#[test]
#[cfg(unix)]
fn a_vacant_pid_has_no_identity_and_a_non_positive_one_is_refused()
-> Result<(), Box<dyn std::error::Error>> {
assert_eq!(identify_process(VACANT)?, None, "no process holds {VACANT}");
for refused in [0, -1, i32::MIN] {
assert_eq!(
identify_process(refused).map_err(|error| error.kind()),
Err(std::io::ErrorKind::InvalidInput),
"pid {refused} names no single process"
);
}
Ok(())
}
#[test]
fn a_damaged_record_is_refused_by_the_arm_that_names_its_damage() {
let long = format!("ps:{}", "x".repeat(MAX_START_TOKEN_BYTES));
let cases: [(&str, ProcessIdentityError); 9] = [
("42", ProcessIdentityError::MissingSeparator),
("4x2/ps:Mon", ProcessIdentityError::MalformedPid),
("+42/ps:Mon", ProcessIdentityError::MalformedPid),
("042/ps:Mon", ProcessIdentityError::MalformedPid),
("0/ps:Mon", ProcessIdentityError::NonPositivePid { pid: 0 }),
("42/", ProcessIdentityError::EmptyStart),
(
"42/ps:Mon Oct",
ProcessIdentityError::StartNotPrintable { at: 6 },
),
("42/when:Mon", ProcessIdentityError::UnknownScheme),
("", ProcessIdentityError::MissingSeparator),
];
for (text, expected) in cases {
assert_eq!(
text.parse::<ProcessIdentity>(),
Err(expected),
"{text:?} is refused by its own arm"
);
}
assert_eq!(
ProcessIdentity::new(42, &long),
Err(ProcessIdentityError::StartTooLong {
len: MAX_START_TOKEN_BYTES.saturating_add(3)
}),
"a token one prefix past the bound is refused before it is kept"
);
}
#[test]
#[cfg(unix)]
fn a_matching_leader_has_its_group_and_every_descendant_stopped()
-> Result<(), Box<dyn std::error::Error>> {
let mut tree = Tree::grow(3)?;
let leader = identify_process(tree.root)?.ok_or("the live leader has no identity")?;
let OrphanReap::Signalled { descendants } = reap_orphaned_group(&leader)? else {
return Err("a leader that still holds its pid must be signalled".into());
};
let mut expected = tree.descendants.clone();
expected.sort_unstable();
assert_eq!(
descendants.pids(),
expected.as_slice(),
"the reap captures exactly the tree below the leader"
);
let status = tree.leader.wait()?;
assert!(!status.success(), "the leader ends by the reap's signal");
assert!(
all_stop_running(&expected)?,
"every descendant of the reaped group stops running: {expected:?}"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn a_forged_start_on_a_live_pid_is_never_signalled() -> Result<(), Box<dyn std::error::Error>> {
let mut tree = Tree::grow(1)?;
let real = identify_process(tree.root)?.ok_or("the live leader has no identity")?;
let forged = ProcessIdentity::new(real.pid(), &format!("{}0", real.started()))?;
assert_eq!(
reap_orphaned_group(&forged)?,
OrphanReap::LeaderReused { holder: real },
"a pid whose start differs from the record is another process"
);
assert!(
tree.leader.try_wait()?.is_none(),
"the process holding the number was not signalled"
);
assert_eq!(
running_processes(&tree.descendants)?
.into_iter()
.collect::<Vec<i32>>(),
tree.descendants,
"nor was anything below it"
);
Ok(())
}
#[test]
#[cfg(unix)]
fn a_leader_that_is_gone_is_left_alone() -> Result<(), Box<dyn std::error::Error>> {
let mut child = std::process::Command::new("true")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let pid = i32::try_from(child.id())?;
let recorded = identify_process(pid)?.ok_or("the child has no identity")?;
child.wait()?;
let reaped = reap_orphaned_group(&recorded)?;
assert!(
!matches!(reaped, OrphanReap::Signalled { .. }),
"a reaped leader's record signals nothing, got {reaped:?}"
);
Ok(())
}
}