use std::io;
use std::path::Path;
use std::time::Duration;
const LEADER_RETRIES: u32 = 50;
const LEADER_BACKOFF: Duration = Duration::from_millis(10);
pub trait PgidFinder {
fn find_holder_pgid(&self, inbox_dir: &Path) -> io::Result<Option<i32>>;
}
#[derive(Debug, Clone)]
pub struct ProcFsFinder {
proc_root: std::path::PathBuf,
leader_retries: u32,
leader_backoff: Duration,
}
impl Default for ProcFsFinder {
fn default() -> Self {
Self {
proc_root: std::path::PathBuf::from("/proc"),
leader_retries: LEADER_RETRIES,
leader_backoff: LEADER_BACKOFF,
}
}
}
impl ProcFsFinder {
#[cfg(test)] pub fn with_root(proc_root: std::path::PathBuf) -> Self {
Self {
proc_root,
..Self::default()
}
}
#[cfg(test)] pub fn with_leader_retry(self, retries: u32, backoff: Duration) -> Self {
Self {
leader_retries: retries,
leader_backoff: backoff,
..self
}
}
fn leader_pgid(&self, pid: i32) -> io::Result<i32> {
let mut pgid = read_pgid(&self.proc_root, pid)?;
let mut retries = self.leader_retries;
while pgid != pid && retries > 0 {
std::thread::sleep(self.leader_backoff);
pgid = read_pgid(&self.proc_root, pid)?;
retries -= 1;
}
if pgid == pid {
return Ok(pgid);
}
Err(io::Error::other(format!(
"pid {pid} holds the agent's inbox lock but reports process group \
{pgid} instead of its own pid: it is not a group leader, so that \
group is one lernie stop does not own (the executor's \
setpgid/setsid has not landed, or failed — ARCH §2.9). Refusing to \
signal it; re-run `lernie stop` once the executor has settled."
)))
}
}
impl PgidFinder for ProcFsFinder {
fn find_holder_pgid(&self, inbox_dir: &Path) -> io::Result<Option<i32>> {
let target = match std::fs::canonicalize(inbox_dir) {
Ok(p) => p,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(e),
};
for entry in std::fs::read_dir(&self.proc_root)?.filter_map(Result::ok) {
let Some(pid) = parse_pid_dir_name(&entry.file_name()) else {
continue;
};
if pid_holds(&entry.path(), &target) {
return self.leader_pgid(pid).map(Some);
}
}
Ok(None)
}
}
fn parse_pid_dir_name(name: &std::ffi::OsStr) -> Option<i32> {
name.to_str().and_then(|s| s.parse::<i32>().ok())
}
fn pid_holds(proc_pid: &Path, target: &Path) -> bool {
let fd_dir = proc_pid.join("fd");
let entries = match std::fs::read_dir(&fd_dir) {
Ok(e) => e,
Err(_) => return false,
};
for entry in entries.filter_map(Result::ok) {
match std::fs::read_link(entry.path()) {
Ok(link) if link == *target => return true,
_ => continue,
}
}
false
}
fn read_pgid(proc_root: &Path, pid: i32) -> io::Result<i32> {
let stat_path = proc_root.join(pid.to_string()).join("stat");
let raw = std::fs::read_to_string(&stat_path)?;
let after_comm = raw.rsplit_once(')').map(|(_, rest)| rest).ok_or_else(|| {
io::Error::new(
io::ErrorKind::InvalidData,
format!("malformed /proc/{pid}/stat"),
)
})?;
let mut fields = after_comm.split_whitespace();
fields.next(); fields.next(); let pgid_str = fields.next().ok_or_else(|| {
io::Error::new(
io::ErrorKind::InvalidData,
format!("missing pgid field in /proc/{pid}/stat"),
)
})?;
pgid_str.parse::<i32>().map_err(|e| {
io::Error::new(
io::ErrorKind::InvalidData,
format!("pgid parse: {e} in /proc/{pid}/stat"),
)
})
}
#[cfg(test)]
mod tests;