use super::{
epoch_abs_diff, io, is_process_alive, now_epoch_secs, process_start_time_secs,
producer_temp_identity, stale_window_from, stale_window_secs, unix_impl, Duration, Path,
ProducerTempKind, WalpinBeacon, WalpinHeartbeat, WalpinPidHealth, WalpinReport,
START_TIME_EPSILON_SECS, UNIX_EPOCH,
};
#[cfg(unix)]
pub(super) const MAX_SIDECAR_ENTRIES: usize = 512;
#[cfg(unix)]
const CAP_SENTINEL_PID: u32 = 0;
#[cfg(unix)]
#[derive(Clone, Copy, PartialEq, Eq)]
pub(super) enum EnumerationPurpose {
Attribution,
Housekeeping,
Diagnostics,
}
#[cfg(unix)]
impl EnumerationPurpose {
fn removes_uncertain_evidence(self) -> bool {
self == Self::Attribution
}
fn removes_dead_or_reused_evidence(self) -> bool {
self != Self::Diagnostics
}
fn removes_orphan_temps(self) -> bool {
self != Self::Diagnostics
}
}
#[cfg(unix)]
enum OrphanTempVerdict {
Skip,
Reap(unix_impl::CheckedEntry),
Untrusted(&'static str),
}
#[cfg(all(unix, test))]
thread_local! {
static STALE_ORPHAN_TEMP_START_TIME_OVERRIDE: std::cell::Cell<Option<Option<i64>>> =
const { std::cell::Cell::new(None) };
}
#[cfg(all(unix, test))]
pub(super) fn set_stale_orphan_temp_start_time_override(value: Option<i64>) {
STALE_ORPHAN_TEMP_START_TIME_OVERRIDE.with(|cell| cell.set(Some(value)));
}
#[cfg(unix)]
fn stale_orphan_temp_actual_start(pid: u32) -> Option<i64> {
#[cfg(test)]
if let Some(overridden) = STALE_ORPHAN_TEMP_START_TIME_OVERRIDE.with(|cell| cell.take()) {
return overridden;
}
process_start_time_secs(pid)
}
#[cfg(unix)]
fn stale_orphan_temp(
handle: &unix_impl::SidecarDirHandle,
name: &str,
pid: u32,
kind: ProducerTempKind,
now: i64,
stale_after_secs: i64,
) -> io::Result<OrphanTempVerdict> {
let entry = match handle.read_checked_entry(name) {
Ok(Some(entry)) => entry,
Ok(None) => return Ok(OrphanTempVerdict::Skip),
Err(e) if e.kind() == io::ErrorKind::PermissionDenied => {
return Ok(OrphanTempVerdict::Untrusted(
"refused: producer temp not owned by current user",
));
}
Err(e) => return Err(e),
};
let modified_at = entry
.mtime
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs() as i64)
.unwrap_or(0);
if now.saturating_sub(modified_at) <= stale_after_secs {
return Ok(OrphanTempVerdict::Skip);
}
let recorded_identity = match kind {
ProducerTempKind::Heartbeat => serde_json::from_slice::<WalpinHeartbeat>(&entry.body)
.ok()
.map(|record| (record.pid, record.started_at)),
ProducerTempKind::Beacon => serde_json::from_slice::<WalpinBeacon>(&entry.body)
.ok()
.map(|record| (record.pid, record.started_at)),
};
let Some((recorded_pid, recorded_start)) = recorded_identity else {
return Ok(OrphanTempVerdict::Untrusted(
"refused: producer temp body does not parse as its recorded kind",
));
};
if recorded_pid != pid {
return Ok(OrphanTempVerdict::Untrusted(
"refused: producer temp identity does not match its filename",
));
}
if !is_process_alive(pid) {
return Ok(OrphanTempVerdict::Reap(entry));
}
let Some(actual_start) = stale_orphan_temp_actual_start(pid) else {
return Ok(OrphanTempVerdict::Untrusted(
"refused: producer temp process start time unavailable",
));
};
if epoch_abs_diff(actual_start, recorded_start) > START_TIME_EPSILON_SECS {
return Ok(OrphanTempVerdict::Reap(entry));
}
Ok(OrphanTempVerdict::Skip)
}
#[cfg(unix)]
pub(super) fn enumerate_live_bounded(
dir: &Path,
sweep_interval: Duration,
max_entries: usize,
purpose: EnumerationPurpose,
) -> io::Result<WalpinReport> {
let handle = match unix_impl::SidecarDirHandle::open_if_exists(dir) {
Ok(Some(h)) => h,
Ok(None) => return Ok(WalpinReport::default()),
Err(e) => return Err(e),
};
let now = now_epoch_secs();
let fallback_window_secs = stale_window_from(sweep_interval);
let mut heartbeats: std::collections::HashMap<u32, WalpinHeartbeat> = Default::default();
let mut beacon_pids: std::collections::HashSet<u32> = Default::default();
let mut unknown: Vec<(u32, &'static str)> = Vec::new();
let mut wedged: std::collections::HashSet<u32> = Default::default();
let (names, producer_temps, truncated) = handle.list_names(max_entries)?;
if truncated {
unknown.push((
CAP_SENTINEL_PID,
"refused: sidecar entry count exceeds enumeration cap",
));
}
let mut cleanup_would_reap = 0usize;
let mut orphan_temps_reaped = 0usize;
for name in producer_temps {
let Some((pid, kind)) = producer_temp_identity(&name) else {
continue;
};
match stale_orphan_temp(&handle, &name, pid, kind, now, fallback_window_secs) {
Ok(OrphanTempVerdict::Reap(entry)) => {
cleanup_would_reap = cleanup_would_reap.saturating_add(1);
if purpose.removes_orphan_temps() {
match handle.remove_if_same(&name, &entry) {
Ok(true) => orphan_temps_reaped = orphan_temps_reaped.saturating_add(1),
Ok(false) => unknown.push((
pid,
"producer temp changed while orphan cleanup was in progress",
)),
Err(_) => unknown
.push((pid, "refused: producer temp changed to an untrusted entry")),
}
}
}
Ok(OrphanTempVerdict::Skip) => {}
Ok(OrphanTempVerdict::Untrusted(reason)) => unknown.push((pid, reason)),
Err(_) => unknown.push((
pid,
"refused: untrusted producer temp (symlink, non-regular, or oversized)",
)),
}
}
for name in names {
let is_heartbeat = name.ends_with(".json");
let is_beacon = name.ends_with(".beacon");
if !is_heartbeat && !is_beacon {
continue;
}
let Some(pid) = name
.rsplit_once('.')
.and_then(|(stem, _)| stem.parse::<u32>().ok())
else {
continue;
};
let (body, mtime) = match handle.read_checked(&name) {
Ok(Some(v)) => v,
Ok(None) => continue, Err(e) if e.kind() == io::ErrorKind::PermissionDenied => {
unknown.push((pid, "refused: sidecar entry not owned by current user"));
continue;
}
Err(_) => {
unknown.push((
pid,
"refused: untrusted sidecar entry (symlink, non-regular, or oversized)",
));
continue;
}
};
let mtime_secs = mtime
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0);
if is_heartbeat {
let heartbeat: WalpinHeartbeat = match serde_json::from_slice(&body) {
Ok(hb) => hb,
Err(_) => {
if purpose.removes_uncertain_evidence()
|| (purpose.removes_dead_or_reused_evidence() && !is_process_alive(pid))
{
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "malformed walpin heartbeat entry"));
continue;
}
};
if heartbeat.pid != pid {
if purpose.removes_uncertain_evidence() {
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "walpin heartbeat PID does not match its entry name"));
continue;
}
let alive = is_process_alive(heartbeat.pid);
let actual_start = if alive {
process_start_time_secs(heartbeat.pid)
} else {
None
};
let identity_ok = actual_start
.map(|actual| {
epoch_abs_diff(actual, heartbeat.started_at) <= START_TIME_EPSILON_SECS
})
.unwrap_or(false);
if !identity_ok {
let positively_dead_or_reused = !alive
|| actual_start.is_some_and(|actual| {
epoch_abs_diff(actual, heartbeat.started_at) > START_TIME_EPSILON_SECS
});
if purpose.removes_uncertain_evidence()
|| (purpose.removes_dead_or_reused_evidence() && positively_dead_or_reused)
{
let _ = handle.unlink_tolerant(&name);
} else {
wedged.insert(pid);
unknown.push((pid, "walpin heartbeat identity could not be verified"));
}
continue;
}
let window = stale_window_secs(heartbeat.sweep_interval_ms, fallback_window_secs);
let hb_fresh = if heartbeat.oldest_tx_started_at.is_some() {
epoch_abs_diff(now, mtime_secs) <= window as u64
} else {
epoch_abs_diff(now, heartbeat.updated_at) <= window as u64
};
if !hb_fresh {
if purpose.removes_uncertain_evidence() {
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "stale walpin heartbeat"));
continue;
}
heartbeats.insert(heartbeat.pid, heartbeat);
} else {
let beacon: WalpinBeacon = match serde_json::from_slice(&body) {
Ok(b) => b,
Err(_) => {
if purpose.removes_uncertain_evidence()
|| (purpose.removes_dead_or_reused_evidence() && !is_process_alive(pid))
{
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "malformed walpin beacon entry"));
continue;
}
};
if beacon.pid != pid {
if purpose.removes_uncertain_evidence() {
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "walpin beacon PID does not match its entry name"));
continue;
}
let alive = is_process_alive(beacon.pid);
let actual_start = if alive {
process_start_time_secs(beacon.pid)
} else {
None
};
let identity_ok = actual_start
.map(|actual| epoch_abs_diff(actual, beacon.started_at) <= START_TIME_EPSILON_SECS)
.unwrap_or(false);
if !identity_ok {
let positively_dead_or_reused = !alive
|| actual_start.is_some_and(|actual| {
epoch_abs_diff(actual, beacon.started_at) > START_TIME_EPSILON_SECS
});
if purpose.removes_uncertain_evidence()
|| (purpose.removes_dead_or_reused_evidence() && positively_dead_or_reused)
{
let _ = handle.unlink_tolerant(&name);
} else {
wedged.insert(pid);
unknown.push((pid, "walpin beacon identity could not be verified"));
}
continue;
}
let window = stale_window_secs(beacon.sweep_interval_ms, fallback_window_secs);
let fresh = epoch_abs_diff(now, mtime_secs) <= window as u64;
if !fresh {
if purpose.removes_uncertain_evidence() {
let _ = handle.unlink_tolerant(&name);
}
wedged.insert(pid);
unknown.push((pid, "stale walpin beacon"));
continue;
}
beacon_pids.insert(beacon.pid);
}
}
for (pid, _) in &unknown {
wedged.insert(*pid);
}
let mut entries: Vec<WalpinPidHealth> = Vec::new();
for (pid, hb) in heartbeats {
if !wedged.contains(&pid) {
entries.push(WalpinPidHealth::Reporting(hb));
}
beacon_pids.remove(&pid);
}
for pid in beacon_pids {
if wedged.contains(&pid) {
continue; }
entries.push(WalpinPidHealth::RegisteredSilent { pid });
}
for (pid, reason) in unknown {
entries.push(WalpinPidHealth::Unknown { pid, reason });
}
Ok(WalpinReport {
entries,
sidecar_listing_truncated: truncated,
cleanup_would_reap,
orphan_temps_reaped,
})
}