use std::{
collections::HashMap,
os::unix::process::ExitStatusExt,
process::ExitStatus,
sync::{Mutex, OnceLock},
time::{Duration, Instant},
};
use tracing::{debug, info};
use crate::runtime;
const RETENTION: Duration = Duration::from_secs(30);
struct Filed {
status: ExitStatus,
filed: Instant,
claimant: Option<String>,
}
fn mailbox() -> &'static Mutex<HashMap<i32, Filed>> {
static MAILBOX: OnceLock<Mutex<HashMap<i32, Filed>>> = OnceLock::new();
MAILBOX.get_or_init(|| Mutex::new(HashMap::new()))
}
pub fn reap_pending() {
if !runtime::init_mode() {
return;
}
loop {
let mut status: libc::c_int = 0;
let pid = unsafe { libc::waitpid(-1, &mut status, libc::WNOHANG) };
if pid <= 0 {
break;
}
debug!("init broker reaped pid {pid}");
mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(
pid,
Filed {
status: ExitStatus::from_raw(status),
filed: Instant::now(),
claimant: None,
},
);
}
}
pub fn take(pid: i32) -> Option<ExitStatus> {
if !runtime::init_mode() {
return None;
}
reap_pending();
let mut box_ = mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if box_.get(&pid).is_some_and(|filed| filed.claimant.is_some()) {
return None;
}
box_.remove(&pid).map(|filed| filed.status)
}
pub fn publish(pid: i32, claimant: &str, status: ExitStatus) {
debug!("routing exit status of pid {pid} to '{claimant}'");
mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(
pid,
Filed {
status,
filed: Instant::now(),
claimant: Some(claimant.to_string()),
},
);
}
pub fn take_for(pid: i32, claimant: &str) -> Option<ExitStatus> {
reap_pending();
let mut box_ = mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let claimable = box_.get(&pid).is_some_and(|filed| {
filed
.claimant
.as_deref()
.is_none_or(|addressee| addressee == claimant)
});
if !claimable {
return None;
}
box_.remove(&pid).map(|filed| filed.status)
}
pub fn drop_claims(claimant: &str) {
mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.retain(|_, filed| filed.claimant.as_deref() != Some(claimant));
}
pub fn sweep_orphans() {
reap_pending();
let mut box_ = mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
box_.retain(|pid, filed| {
if filed.claimant.is_some() || filed.filed.elapsed() < RETENTION {
return true;
}
info!("init reaped adopted orphan pid {pid} ({:?})", filed.status);
false
});
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn broker_is_inert_outside_init_mode() {
assert!(take(1).is_none());
reap_pending();
sweep_orphans();
assert!(
mailbox()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty()
);
}
#[test]
fn routed_status_reaches_only_its_claimant() {
let pid = -4242;
publish(pid, "nightly", ExitStatus::from_raw(1024));
assert!(take_for(pid, "someone-else").is_none());
assert!(take(pid).is_none());
let claimed = take_for(pid, "nightly").expect("claimant reads its own status");
assert_eq!(claimed.code(), Some(4));
assert!(take_for(pid, "nightly").is_none());
}
#[test]
fn dropping_claims_clears_stale_runs() {
let pid = -4243;
publish(pid, "nightly", ExitStatus::from_raw(0));
drop_claims("nightly");
assert!(take_for(pid, "nightly").is_none());
}
}