#[cfg(all(unix, not(any(target_os = "macos", target_os = "linux"))))]
use super::io_other;
#[cfg(any(target_os = "macos", target_os = "linux"))]
use super::{fs, Instant};
#[cfg(unix)]
use super::{io, Duration, Path};
#[derive(Debug, Clone, Default, PartialEq)]
pub struct CensusResult {
pub holders: std::collections::HashSet<u32>,
pub uninspectable_pids: Vec<u32>,
pub truncated: bool,
pub budget_exhausted: bool,
}
#[cfg(target_os = "linux")]
pub(super) fn census_visible_self_pid() -> Option<u32> {
fs::read_link("/proc/self")
.ok()?
.file_name()?
.to_str()?
.parse()
.ok()
}
#[cfg(not(target_os = "linux"))]
pub(super) fn census_visible_self_pid() -> Option<u32> {
Some(std::process::id())
}
impl CensusResult {
pub fn is_complete(&self) -> bool {
self.uninspectable_pids.is_empty() && !self.truncated
}
#[cfg(unix)]
fn apply_self_canary(&mut self) {
self.apply_self_canary_for(census_visible_self_pid());
}
pub(super) fn apply_self_canary_for(&mut self, expected_self: Option<u32>) {
if expected_self.is_none_or(|pid| !self.holders.contains(&pid)) {
self.truncated = true;
}
}
}
#[cfg(target_os = "macos")]
pub(super) fn macos_pid_genuinely_gone(errno: Option<i32>) -> bool {
errno == Some(libc::ESRCH)
}
#[cfg(target_os = "macos")]
pub(super) fn proc_pidfdinfo_returned_expected_size(
returned_bytes: i32,
expected_size: usize,
) -> bool {
returned_bytes > 0 && returned_bytes as usize == expected_size
}
#[cfg(target_os = "macos")]
const CENSUS_BUFFER_NEGOTIATION_ATTEMPTS: usize = 4;
#[cfg(target_os = "macos")]
pub(super) fn negotiate_buffer<T: Default + Clone>(
size_call: impl Fn() -> std::os::raw::c_int,
data_call: impl Fn(*mut std::os::raw::c_void, std::os::raw::c_int) -> std::os::raw::c_int,
should_stop: &impl Fn() -> bool,
) -> io::Result<(Vec<T>, bool)> {
let item_size = std::mem::size_of::<T>();
for attempt in 0..CENSUS_BUFFER_NEGOTIATION_ATTEMPTS {
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
let needed = size_call();
if needed <= 0 {
return Err(io::Error::last_os_error());
}
let needed_items = needed as usize / item_size + 1;
let item_count = needed_items + needed_items / 4 + 8;
let mut buf: Vec<T> = vec![T::default(); item_count];
let cap_bytes = (buf.len() * item_size) as std::os::raw::c_int;
let bytes = data_call(buf.as_mut_ptr() as *mut std::os::raw::c_void, cap_bytes);
if bytes <= 0 {
return Err(io::Error::last_os_error());
}
let filled_capacity = bytes as usize >= cap_bytes as usize;
let is_last_attempt = attempt + 1 == CENSUS_BUFFER_NEGOTIATION_ATTEMPTS;
if filled_capacity && !is_last_attempt {
continue;
}
let count = (bytes as usize / item_size).min(buf.len());
buf.truncate(count);
return Ok((buf, filled_capacity));
}
unreachable!("loop always returns or errors within CENSUS_BUFFER_NEGOTIATION_ATTEMPTS")
}
#[cfg(target_os = "macos")]
pub fn census_holders(db_path: &Path) -> io::Result<CensusResult> {
census_holders_until(db_path, || false)
}
#[cfg(target_os = "macos")]
pub(crate) fn census_holders_until<C>(db_path: &Path, should_stop: C) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
census_holders_inner(db_path, should_stop, None)
}
#[cfg(target_os = "macos")]
pub fn census_holders_until_within<C>(
db_path: &Path,
should_stop: C,
budget: Duration,
) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
census_holders_inner(db_path, should_stop, Some(Instant::now() + budget))
}
#[cfg(target_os = "macos")]
fn census_holders_inner<C>(
db_path: &Path,
should_stop: C,
deadline: Option<Instant>,
) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
let budget_spent = || deadline.is_some_and(|d| Instant::now() >= d);
use std::os::raw::{c_int, c_void};
use std::os::unix::fs::MetadataExt;
const PROC_ALL_PIDS: u32 = 1;
const PROC_PIDLISTFDS: c_int = 1;
const PROC_PIDFDVNODEPATHINFO: c_int = 2;
const PROX_FDTYPE_VNODE: u32 = 1;
const MAXPATHLEN: usize = 1024;
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
#[repr(C)]
#[derive(Clone, Default)]
struct ProcFdInfo {
proc_fd: i32,
proc_fdtype: u32,
}
#[repr(C)]
struct ProcFileInfo {
fi_openflags: u32,
fi_status: u32,
fi_offset: i64,
fi_type: i32,
fi_guardflags: u32,
}
#[repr(C)]
struct FsId {
val: [i32; 2],
}
#[repr(C)]
struct VinfoStat {
vst_dev: u32,
vst_mode: u16,
vst_nlink: u16,
vst_ino: u64,
vst_uid: u32,
vst_gid: u32,
vst_atime: i64,
vst_atimensec: i64,
vst_mtime: i64,
vst_mtimensec: i64,
vst_ctime: i64,
vst_ctimensec: i64,
vst_birthtime: i64,
vst_birthtimensec: i64,
vst_size: i64,
vst_blocks: i64,
vst_blksize: i32,
vst_flags: u32,
vst_gen: u32,
vst_rdev: u32,
vst_qspare: [i64; 2],
}
#[repr(C)]
struct VnodeInfo {
vi_stat: VinfoStat,
vi_type: i32,
vi_pad: i32,
vi_fsid: FsId,
}
#[repr(C)]
struct VnodeInfoPath {
vip_vi: VnodeInfo,
vip_path: [u8; MAXPATHLEN],
}
#[repr(C)]
struct VnodeFdInfoWithPath {
pfi: ProcFileInfo,
pvip: VnodeInfoPath,
}
#[link(name = "proc")]
extern "C" {
fn proc_listpids(kind: u32, typeinfo: u32, buffer: *mut c_void, buffersize: c_int)
-> c_int;
fn proc_pidinfo(
pid: c_int,
flavor: c_int,
arg: u64,
buffer: *mut c_void,
buffersize: c_int,
) -> c_int;
fn proc_pidfdinfo(
pid: c_int,
fd: c_int,
flavor: c_int,
buffer: *mut c_void,
buffersize: c_int,
) -> c_int;
}
let target_meta = fs::metadata(db_path)?;
let target_ident = (target_meta.dev() as u32, target_meta.ino());
let (pid_buf, pid_list_truncated): (Vec<i32>, bool) = negotiate_buffer(
|| unsafe { proc_listpids(PROC_ALL_PIDS, 0, std::ptr::null_mut(), 0) },
|buf_ptr, buf_bytes| unsafe { proc_listpids(PROC_ALL_PIDS, 0, buf_ptr, buf_bytes) },
&should_stop,
)?;
let mut holders = std::collections::HashSet::new();
let mut uninspectable: Vec<u32> = Vec::new();
let mut budget_exhausted = false;
'pids: for &pid in &pid_buf {
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
if budget_spent() {
budget_exhausted = true;
break 'pids;
}
if pid <= 0 {
continue;
}
let (fd_buf, fd_list_truncated): (Vec<ProcFdInfo>, bool) = match negotiate_buffer(
|| unsafe { proc_pidinfo(pid, PROC_PIDLISTFDS, 0, std::ptr::null_mut(), 0) },
|buf_ptr, buf_bytes| unsafe {
proc_pidinfo(pid, PROC_PIDLISTFDS, 0, buf_ptr, buf_bytes)
},
&should_stop,
) {
Ok(v) => v,
Err(e) => {
if e.kind() == io::ErrorKind::Interrupted {
return Err(e);
}
if !macos_pid_genuinely_gone(e.raw_os_error()) {
uninspectable.push(pid as u32);
}
continue;
}
};
if fd_list_truncated {
uninspectable.push(pid as u32);
}
for fdinfo in &fd_buf {
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
if budget_spent() {
budget_exhausted = true;
break 'pids;
}
if fdinfo.proc_fdtype != PROX_FDTYPE_VNODE {
continue;
}
let mut vinfo: VnodeFdInfoWithPath = unsafe { std::mem::zeroed() };
let vsize = unsafe {
proc_pidfdinfo(
pid,
fdinfo.proc_fd,
PROC_PIDFDVNODEPATHINFO,
&mut vinfo as *mut _ as *mut c_void,
std::mem::size_of::<VnodeFdInfoWithPath>() as c_int,
)
};
if !proc_pidfdinfo_returned_expected_size(
vsize,
std::mem::size_of::<VnodeFdInfoWithPath>(),
) {
if vsize <= 0 {
let errno = io::Error::last_os_error().raw_os_error();
if !macos_pid_genuinely_gone(errno) {
uninspectable.push(pid as u32);
}
} else {
uninspectable.push(pid as u32);
}
continue;
}
let vstat = &vinfo.pvip.vip_vi.vi_stat;
if (vstat.vst_dev, vstat.vst_ino) == target_ident {
holders.insert(pid as u32);
break;
}
}
}
uninspectable.sort_unstable();
uninspectable.dedup();
let mut census = CensusResult {
holders,
uninspectable_pids: uninspectable,
truncated: pid_list_truncated || budget_exhausted,
budget_exhausted,
};
census.apply_self_canary();
Ok(census)
}
#[cfg(target_os = "linux")]
pub(super) fn linux_proc_gone(err: &io::Error) -> bool {
err.kind() == io::ErrorKind::NotFound
}
#[cfg(target_os = "linux")]
pub(super) const PROC_PID_INIT_INO: u64 = 0xEFFFFFFC;
#[cfg(target_os = "linux")]
pub(super) fn pid_ns_is_init(ino: u64) -> bool {
ino == PROC_PID_INIT_INO
}
#[cfg(target_os = "linux")]
pub(super) fn proc_mount_restricts_visibility(options: &str) -> bool {
options.split(',').map(str::trim).any(|opt| {
if let Some(value) = opt.strip_prefix("hidepid=") {
!matches!(value, "0" | "off")
} else {
opt == "hidepid" || opt == "subset" || opt.starts_with("subset=")
}
})
}
#[cfg(target_os = "linux")]
pub(super) fn proc_mount_is_visibility_restricted() -> Option<bool> {
let mountinfo = fs::read_to_string("/proc/self/mountinfo").ok()?;
proc_mounts_restricted_in(&mountinfo)
}
#[cfg(target_os = "linux")]
pub(super) fn proc_mounts_restricted_in(mountinfo: &str) -> Option<bool> {
let mut found_any = false;
for line in mountinfo.lines() {
let Some((fields_part, super_part)) = line.split_once(" - ") else {
continue;
};
let fields: Vec<&str> = fields_part.split(' ').collect();
if fields.len() < 6 || fields[4] != "/proc" {
continue;
}
let mount_options = fields[5];
let super_fields: Vec<&str> = super_part.split(' ').collect();
if super_fields.first().copied() != Some("proc") {
continue;
}
let super_options = super_fields.get(2).copied().unwrap_or("");
found_any = true;
if proc_mount_restricts_visibility(mount_options)
|| proc_mount_restricts_visibility(super_options)
{
return Some(true);
}
}
if found_any {
Some(false)
} else {
None
}
}
#[cfg(target_os = "linux")]
pub fn census_holders(db_path: &Path) -> io::Result<CensusResult> {
census_holders_until(db_path, || false)
}
#[cfg(target_os = "linux")]
pub(crate) fn census_holders_until<C>(db_path: &Path, should_stop: C) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
census_holders_inner(db_path, should_stop, None)
}
#[cfg(target_os = "linux")]
pub fn census_holders_until_within<C>(
db_path: &Path,
should_stop: C,
budget: Duration,
) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
census_holders_inner(db_path, should_stop, Some(Instant::now() + budget))
}
#[cfg(target_os = "linux")]
fn census_holders_inner<C>(
db_path: &Path,
should_stop: C,
deadline: Option<Instant>,
) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
let budget_spent = || deadline.is_some_and(|d| Instant::now() >= d);
use std::os::unix::fs::MetadataExt;
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
let target_meta = fs::metadata(db_path)?;
let target_ident = (target_meta.dev(), target_meta.ino());
let mut holders = std::collections::HashSet::new();
let mut uninspectable: Vec<u32> = Vec::new();
let mut truncated = false;
match fs::metadata("/proc/self/ns/pid") {
Ok(meta) if pid_ns_is_init(meta.ino()) => {}
_ => truncated = true,
}
match proc_mount_is_visibility_restricted() {
Some(false) => {}
Some(true) | None => truncated = true,
}
let mut budget_exhausted = false;
let proc_dir = fs::read_dir("/proc")?;
'pids: for entry_result in proc_dir {
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
if budget_spent() {
truncated = true;
budget_exhausted = true;
break 'pids;
}
let proc_entry = match entry_result {
Ok(e) => e,
Err(_) => {
truncated = true;
continue;
}
};
let Some(pid) = proc_entry
.file_name()
.to_str()
.and_then(|s| s.parse::<u32>().ok())
else {
continue;
};
let fd_dir = proc_entry.path().join("fd");
let fds = match fs::read_dir(&fd_dir) {
Ok(fds) => fds,
Err(e) if linux_proc_gone(&e) => continue,
Err(_) => {
uninspectable.push(pid);
continue;
}
};
for fd_result in fds {
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
if budget_spent() {
truncated = true;
budget_exhausted = true;
break 'pids;
}
let fd_entry = match fd_result {
Ok(e) => e,
Err(_) => {
uninspectable.push(pid);
continue;
}
};
match fs::metadata(fd_entry.path()) {
Ok(meta) => {
if (meta.dev(), meta.ino()) == target_ident {
holders.insert(pid);
break;
}
}
Err(e) if e.kind() == io::ErrorKind::NotFound => {}
Err(_) => uninspectable.push(pid),
}
}
}
uninspectable.sort_unstable();
uninspectable.dedup();
let mut census = CensusResult {
holders,
uninspectable_pids: uninspectable,
truncated,
budget_exhausted,
};
census.apply_self_canary();
Ok(census)
}
#[cfg(all(unix, not(any(target_os = "macos", target_os = "linux"))))]
pub fn census_holders(_db_path: &Path) -> io::Result<CensusResult> {
Err(io_other(
"OS-derived holder census has no implementation on this Unix target",
))
}
#[cfg(all(unix, not(any(target_os = "macos", target_os = "linux"))))]
pub(crate) fn census_holders_until<C>(_db_path: &Path, should_stop: C) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
if should_stop() {
return Err(io::Error::new(
io::ErrorKind::Interrupted,
"WAL holder census cancelled",
));
}
Err(io_other(
"OS-derived holder census has no implementation on this Unix target",
))
}
#[cfg(all(unix, not(any(target_os = "macos", target_os = "linux"))))]
pub fn census_holders_until_within<C>(
db_path: &Path,
should_stop: C,
_budget: Duration,
) -> io::Result<CensusResult>
where
C: Fn() -> bool,
{
census_holders_until(db_path, should_stop)
}