use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::path::{Path, PathBuf};
use std::time::Duration;
const PID_FILE: &str = "pid";
pub(crate) fn pid_path(state_dir: &Path, ident: &str) -> PathBuf {
state_dir.join(ident).join(PID_FILE)
}
pub(crate) fn start_time(pid: u32) -> Option<String> {
imp::start_time(pid)
}
pub(crate) fn record(state_dir: &Path, ident: &str, pid: u32) {
let Some(start) = start_time(pid) else {
clear(state_dir, ident);
return;
};
let path = pid_path(state_dir, ident);
if let Err(e) = std::fs::write(&path, format!("{pid} {start}\n")) {
tracing::warn!(path = %path.display(), "could not record the child's pid: {e}");
}
}
pub(crate) fn clear(state_dir: &Path, ident: &str) {
let _ = std::fs::remove_file(pid_path(state_dir, ident));
}
fn read(state_dir: &Path, ident: &str) -> Option<(u32, String)> {
let text = std::fs::read_to_string(pid_path(state_dir, ident)).ok()?;
let (pid, start) = text.trim().split_once(' ')?;
Some((pid.parse().ok()?, start.to_string()))
}
fn alive(pid: u32) -> bool {
unsafe {
libc::kill(pid as i32, 0) == 0
|| std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}
}
async fn wait_gone(pid: u32, within: Duration) -> bool {
let deadline = tokio::time::Instant::now() + within;
while alive(pid) {
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
true
}
pub(crate) async fn reap(state_dir: &Path, ident: &str, grace: Duration) -> Option<u32> {
let recorded = read(state_dir, ident);
clear(state_dir, ident);
let (pid, start) = recorded?;
if pid == 0 || !alive(pid) || start_time(pid).as_deref() != Some(start.as_str()) {
return None;
}
tracing::warn!(
workload = ident,
pid,
"reaping the orphaned child of a previous native runtime on this state dir"
);
unsafe { libc::kill(pid as i32, libc::SIGTERM) };
if !wait_gone(pid, grace).await {
unsafe { libc::kill(pid as i32, libc::SIGKILL) };
wait_gone(pid, Duration::from_secs(5)).await;
}
Some(pid)
}
pub(crate) fn is_held(bind_ip: Ipv4Addr, port: u16) -> bool {
let addr = SocketAddr::new(IpAddr::V4(bind_ip), port);
std::net::TcpStream::connect_timeout(&addr, Duration::from_millis(250))
.is_ok_and(|c| c.local_addr().ok() != Some(addr))
}
pub(crate) async fn is_held_settled(bind_ip: Ipv4Addr, port: u16) -> bool {
if !is_held(bind_ip, port) {
return false;
}
tokio::time::sleep(Duration::from_millis(300)).await;
is_held(bind_ip, port)
}
pub(crate) fn describe_holder(port: u16) -> String {
let out = std::process::Command::new("lsof")
.args(["-nP", &format!("-iTCP:{port}"), "-sTCP:LISTEN", "-Fpc"])
.output();
let Ok(out) = out else {
return "pid unknown".to_string();
};
let text = String::from_utf8_lossy(&out.stdout);
let pid = text.lines().find_map(|l| l.strip_prefix('p'));
let cmd = text.lines().find_map(|l| l.strip_prefix('c'));
match (pid, cmd) {
(Some(p), Some(c)) => format!("pid {p} ({c})"),
(Some(p), None) => format!("pid {p}"),
_ => "pid unknown".to_string(),
}
}
#[cfg(target_os = "macos")]
mod imp {
pub fn start_time(pid: u32) -> Option<String> {
let pid = i32::try_from(pid).ok().filter(|p| *p > 0)?;
let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
let size = std::mem::size_of::<libc::proc_bsdinfo>() as libc::c_int;
let rc = unsafe {
libc::proc_pidinfo(
pid,
libc::PROC_PIDTBSDINFO,
0,
(&mut info as *mut libc::proc_bsdinfo).cast(),
size,
)
};
if rc < size {
return None;
}
Some(format!("{}.{}", info.pbi_start_tvsec, info.pbi_start_tvusec))
}
}
#[cfg(target_os = "linux")]
mod imp {
pub fn start_time(pid: u32) -> Option<String> {
let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
let rest = &stat[stat.rfind(')')? + 1..];
rest.split_whitespace().nth(19).map(str::to_string)
}
}
#[cfg(not(any(target_os = "macos", target_os = "linux")))]
mod imp {
pub fn start_time(_pid: u32) -> Option<String> {
None
}
}
#[cfg(all(test, unix))]
mod tests {
use super::*;
#[test]
fn recording_a_dead_pid_removes_the_file_instead_of_leaving_a_stale_one() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(tmp.path().join("w")).unwrap();
record(tmp.path(), "w", std::process::id());
assert!(pid_path(tmp.path(), "w").exists());
let mut child = std::process::Command::new("true").spawn().unwrap();
let dead = child.id();
child.wait().unwrap();
record(tmp.path(), "w", dead);
assert!(!pid_path(tmp.path(), "w").exists());
}
}