use super::*;
#[cfg(unix)]
use std::os::unix::fs::{MetadataExt, PermissionsExt};
#[cfg(unix)]
fn current_uid() -> u32 {
unix_impl::current_uid()
}
fn heartbeat(pid: u32) -> WalpinHeartbeat {
let now = now_epoch_secs();
WalpinHeartbeat {
pid,
process_role: "session".to_string(),
started_at: process_start_time_secs(std::process::id()).unwrap_or(0),
oldest_tx_age_secs: 45.0,
oldest_tx_label: Some("test_span".to_string()),
oldest_tx_started_at: Some(now - 45),
updated_at: now,
sweep_interval_ms: 5_000,
attribution_basis: Some("origin".to_string()),
}
}
#[test]
fn sidecar_dir_is_db_scoped_sibling() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("khive.db");
assert_eq!(sidecar_dir_for(&db), dir.path().join("khive.db.walpin"));
}
include!("walpin/environment_tests.rs");
#[test]
fn windows_handle_kind_requires_expected_type_without_reparse_data() {
const DIRECTORY: u32 = 0x10;
const REPARSE_POINT: u32 = 0x400;
assert!(windows_attribute_tag_is_acceptable(DIRECTORY, 0, true));
assert!(windows_attribute_tag_is_acceptable(0, 0, false));
assert!(!windows_attribute_tag_is_acceptable(0, 0, true));
assert!(!windows_attribute_tag_is_acceptable(DIRECTORY, 0, false));
assert!(!windows_attribute_tag_is_acceptable(
DIRECTORY | REPARSE_POINT,
0,
true
));
assert!(!windows_attribute_tag_is_acceptable(DIRECTORY, 1, true));
}
#[test]
fn windows_final_path_comparison_requires_exact_handle_resolution() {
let expected: Vec<u16> = r"\\?\C:\data\khive.db.walpin".encode_utf16().collect();
let same = expected.clone();
let redirected: Vec<u16> = r"\\?\C:\other\khive.db.walpin".encode_utf16().collect();
assert!(windows_final_path_matches(&expected, &same));
assert!(!windows_final_path_matches(&expected, &redirected));
}
#[test]
fn windows_relative_child_names_are_single_components() {
assert!(windows_relative_child_name_is_safe("42.json"));
assert!(!windows_relative_child_name_is_safe(""));
assert!(!windows_relative_child_name_is_safe("."));
assert!(!windows_relative_child_name_is_safe(".."));
assert!(!windows_relative_child_name_is_safe("..\\42.json"));
assert!(!windows_relative_child_name_is_safe("nested/42.json"));
assert!(!windows_relative_child_name_is_safe("42\0.json"));
}
#[test]
fn windows_owner_dacl_accepts_token_user_owner_and_rejects_broader_shapes() {
const ACCESS_ALLOWED: u8 = 0;
const OBJECT_AND_CONTAINER_INHERIT: u8 = 0x03;
const FILE_ALL_ACCESS: u32 = 0x001f_01ff;
assert!(windows_owner_dacl_is_restricted(
1,
ACCESS_ALLOWED,
OBJECT_AND_CONTAINER_INHERIT,
FILE_ALL_ACCESS,
true,
true,
true,
));
assert!(!windows_owner_dacl_is_restricted(
2,
ACCESS_ALLOWED,
OBJECT_AND_CONTAINER_INHERIT,
FILE_ALL_ACCESS,
true,
true,
true,
));
assert!(!windows_owner_dacl_is_restricted(
1,
ACCESS_ALLOWED,
OBJECT_AND_CONTAINER_INHERIT,
FILE_ALL_ACCESS,
false,
true,
true,
));
assert!(!windows_owner_dacl_is_restricted(
1,
ACCESS_ALLOWED,
OBJECT_AND_CONTAINER_INHERIT,
FILE_ALL_ACCESS,
true,
true,
false,
));
}
#[test]
fn windows_owner_dacl_rejects_group_owner() {
const ACCESS_ALLOWED: u8 = 0;
const OBJECT_AND_CONTAINER_INHERIT: u8 = 0x03;
const FILE_ALL_ACCESS: u32 = 0x001f_01ff;
assert!(!windows_owner_dacl_is_restricted(
1,
ACCESS_ALLOWED,
OBJECT_AND_CONTAINER_INHERIT,
FILE_ALL_ACCESS,
true,
false,
true,
));
}
#[cfg(unix)]
#[test]
fn ensure_sidecar_dir_creates_0700_owned_dir() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).expect("should create");
let meta = fs::symlink_metadata(&dir).unwrap();
assert!(meta.is_dir());
assert_eq!(meta.permissions().mode() & 0o777, 0o700);
assert_eq!(meta.uid(), current_uid());
}
#[cfg(unix)]
#[test]
fn ensure_sidecar_dir_refuses_wrong_mode() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
fs::create_dir(&dir).unwrap();
fs::set_permissions(&dir, fs::Permissions::from_mode(0o755)).unwrap();
let err = ensure_sidecar_dir(&dir).expect_err("wrong mode must be refused");
assert!(err.to_string().contains("expected 0700"));
}
#[cfg(unix)]
#[test]
fn ensure_sidecar_dir_refuses_symlink() {
let root = tempfile::tempdir().unwrap();
let real = root.path().join("real_dir");
fs::create_dir(&real).unwrap();
let link = root.path().join("khive.db.walpin");
std::os::unix::fs::symlink(&real, &link).unwrap();
let err = ensure_sidecar_dir(&link).expect_err("symlink must be refused");
assert!(err.to_string().contains("symlink"));
}
#[test]
#[cfg(unix)]
fn ensure_sidecar_dir_accepts_current_user_link_at_end_of_parent_walk() {
let root = tempfile::tempdir().unwrap();
let real = root.path().join("real_ancestor");
fs::create_dir(&real).unwrap();
let link = root.path().join("linked_ancestor");
std::os::unix::fs::symlink(&real, &link).unwrap();
assert_eq!(fs::symlink_metadata(&link).unwrap().uid(), current_uid());
let dir = link.join("khive.db.walpin");
ensure_sidecar_dir(&dir).expect("the parent walk's final link is a trusted ancestor");
let meta = fs::symlink_metadata(real.join("khive.db.walpin")).unwrap();
assert!(meta.is_dir());
assert_eq!(meta.permissions().mode() & 0o777, 0o700);
assert_eq!(meta.uid(), current_uid());
}
#[test]
#[cfg(unix)]
fn ensure_sidecar_dir_refuses_ancestor_link_in_nonsticky_writable_parent() {
use khive_fs::directory_walk::{AncestorLinkCondition, AncestorLinkRefusal};
let root = tempfile::tempdir().unwrap();
let real = root.path().join("real_ancestor");
fs::create_dir(&real).unwrap();
let payload = real.join("payload");
fs::write(&payload, b"keep").unwrap();
let link = root.path().join("linked_ancestor");
std::os::unix::fs::symlink(&real, &link).unwrap();
fs::set_permissions(root.path(), fs::Permissions::from_mode(0o777)).unwrap();
let err = ensure_sidecar_dir(&link.join("khive.db.walpin"))
.expect_err("a nonsticky writable link parent must be refused");
let refusal = err
.get_ref()
.and_then(|cause| cause.downcast_ref::<AncestorLinkRefusal>())
.expect("the shared typed refusal must survive the caller");
assert_eq!(refusal.condition, AncestorLinkCondition::ParentPermissions);
assert_eq!(
refusal.component,
std::ffi::OsString::from("linked_ancestor")
);
assert_eq!(refusal.link_uid, current_uid());
assert_eq!(refusal.parent_uid, Some(current_uid()));
assert_eq!(refusal.parent_mode, Some(0o40777));
assert!(!real.join("khive.db.walpin").exists());
assert_eq!(fs::read(payload).unwrap(), b"keep");
assert!(fs::symlink_metadata(link).unwrap().file_type().is_symlink());
}
#[test]
fn write_then_read_heartbeat_roundtrips() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let hb = heartbeat(std::process::id());
write_heartbeat(&dir, &hb).expect("write should succeed");
let content = fs::read_to_string(dir.join(format!("{}.json", hb.pid))).unwrap();
let read_back: WalpinHeartbeat = serde_json::from_str(&content).unwrap();
assert_eq!(read_back, hb);
}
#[test]
fn heartbeat_deserializes_pre_rename_interval_ms_field() {
let json = r#"{
"pid": 4242,
"process_role": "session",
"started_at": 1000,
"oldest_tx_age_secs": 45.0,
"oldest_tx_label": "test_span",
"updated_at": 1045,
"interval_ms": 60000
}"#;
let hb: WalpinHeartbeat = serde_json::from_str(json).unwrap();
assert_eq!(hb.sweep_interval_ms, 60_000);
}
#[test]
fn beacon_deserializes_pre_rename_interval_ms_field() {
let json = r#"{
"pid": 4242,
"process_role": "session",
"started_at": 1000,
"interval_ms": 60000
}"#;
let b: WalpinBeacon = serde_json::from_str(json).unwrap();
assert_eq!(b.sweep_interval_ms, 60_000);
}
#[cfg(unix)]
#[test]
fn write_heartbeat_refuses_symlinked_target() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let real = root.path().join("elsewhere.txt");
fs::write(&real, b"nope").unwrap();
let hb = heartbeat(999_999);
let target = dir.join(format!("{}.json", hb.pid));
std::os::unix::fs::symlink(&real, &target).unwrap();
let err = write_heartbeat(&dir, &hb).expect_err("symlinked target must be refused");
assert!(err.to_string().contains("symlink"));
assert_eq!(fs::read_to_string(&real).unwrap(), "nope");
}
#[cfg(unix)]
#[test]
fn sidecar_dir_distinguishes_non_utf8_db_names_on_disk() {
use std::os::unix::ffi::OsStrExt;
let root = tempfile::tempdir().unwrap();
let name_a = std::ffi::OsStr::from_bytes(b"khive-\xffdb.sqlite");
let name_b = std::ffi::OsStr::from_bytes(b"khive-\xfedb.sqlite");
let dir_a = sidecar_dir_for(&root.path().join(name_a));
let dir_b = sidecar_dir_for(&root.path().join(name_b));
assert_ne!(
dir_a, dir_b,
"distinct db names must produce distinct sidecar paths"
);
let hb = heartbeat(std::process::id());
if let Err(e) = write_heartbeat(&dir_a, &hb) {
eprintln!(
"skipping sidecar_dir_distinguishes_non_utf8_db_names_on_disk: filesystem \
rejected a non-UTF-8 sidecar directory name ({e}); this platform's filesystem \
does not support the case under test"
);
return;
}
write_heartbeat(&dir_b, &hb).expect("write to second non-UTF-8 sidecar");
let mut entries: Vec<_> = fs::read_dir(root.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.file_name())
.collect();
entries.sort();
assert_eq!(
entries.len(),
2,
"distinct non-UTF-8 database names must produce two distinct sidecar directories \
on disk, not collide onto one: got {entries:?}"
);
}
#[test]
fn remove_heartbeat_is_idempotent_when_absent() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
remove_heartbeat(&dir, 123_456).expect("removing an absent entry is a no-op");
}
#[test]
fn is_process_alive_true_for_self_false_for_reserved_pid() {
assert!(is_process_alive(std::process::id()));
assert!(!is_process_alive(0));
}
#[test]
fn process_start_time_resolves_for_self() {
let start = process_start_time_secs(std::process::id());
assert!(
start.is_some(),
"must resolve this process's own start time"
);
let now = now_epoch_secs();
assert!(
start.unwrap() <= now,
"start time must not be in the future"
);
}
fn beacon(pid: u32) -> WalpinBeacon {
WalpinBeacon {
pid,
process_role: "session".to_string(),
started_at: process_start_time_secs(std::process::id()).unwrap_or(0),
sweep_interval_ms: 5_000,
}
}
#[cfg(unix)]
#[test]
fn enumerate_live_reports_and_retains_a_genuinely_live_entry() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let hb = heartbeat(std::process::id());
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
let reporting: Vec<_> = report.reporting().collect();
assert_eq!(reporting.len(), 1);
assert_eq!(reporting[0].pid, hb.pid);
assert!(report.fully_attributed());
assert!(dir.join(format!("{}.json", hb.pid)).exists());
}
#[cfg(unix)]
#[test]
fn epoch_abs_diff_saturates_instead_of_wrapping() {
assert_eq!(epoch_abs_diff(5, 3), 2);
assert_eq!(epoch_abs_diff(3, 5), 2);
assert_eq!(epoch_abs_diff(0, 0), 0);
assert_eq!(epoch_abs_diff(1, i64::MIN), u64::MAX);
assert_eq!(epoch_abs_diff(i64::MIN, i64::MAX), u64::MAX);
assert_eq!(epoch_abs_diff(-1, i64::MAX), 1u64 << 63);
}
#[cfg(unix)]
#[test]
fn enumerate_live_extreme_timestamp_classifies_unknown_not_fresh() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut hb = heartbeat(std::process::id());
hb.oldest_tx_started_at = None;
hb.updated_at = i64::MIN;
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(
report.reporting().next().is_none(),
"an extreme updated_at must never classify as fresh"
);
assert!(
report
.entries
.iter()
.any(|e| matches!(e, WalpinPidHealth::Unknown { pid, .. } if *pid == hb.pid)),
"the extreme-timestamp entry must stay Unknown, not vanish"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_bounded_caps_listing_with_sentinel_marker() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let live = heartbeat(std::process::id());
write_heartbeat(&dir, &live).unwrap();
for pid in [2_000_000_001u32, 2_000_000_002] {
let mut hb = heartbeat(std::process::id());
hb.pid = pid;
write_heartbeat(&dir, &hb).unwrap();
}
let report = enumerate_live_bounded(
&dir,
Duration::from_secs(5),
1,
EnumerationPurpose::Attribution,
)
.unwrap();
let markers = report
.entries
.iter()
.filter(|e| {
matches!(
e,
WalpinPidHealth::Unknown { pid: 0, reason }
if reason.contains("enumeration cap")
)
})
.count();
assert_eq!(
markers, 1,
"a truncated listing must surface exactly one sentinel Unknown marker"
);
assert!(
!report.fully_attributed(),
"a capped enumeration can never claim full attribution"
);
assert!(report.entries.len() <= 2, "got {:?}", report.entries);
}
#[cfg(unix)]
#[test]
fn enumerate_live_bounded_caps_hidden_entry_scan_with_sentinel_marker() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let live = heartbeat(std::process::id());
write_heartbeat(&dir, &live).unwrap();
for i in 0..64 {
std::fs::write(dir.join(format!(".junk{i}")), b"x").unwrap();
}
let report = enumerate_live_bounded(
&dir,
Duration::from_secs(5),
4,
EnumerationPurpose::Attribution,
)
.unwrap();
let markers = report
.entries
.iter()
.filter(|e| {
matches!(
e,
WalpinPidHealth::Unknown { pid: 0, reason }
if reason.contains("enumeration cap")
)
})
.count();
assert_eq!(
markers, 1,
"a hidden-entry flood must surface exactly one sentinel Unknown marker"
);
assert!(
!report.fully_attributed(),
"an enumeration cut short by hidden entries can never claim full attribution"
);
}
#[cfg(unix)]
#[test]
fn housekeeping_reaps_only_stale_producer_temps_with_dead_identity() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let dead_pid = 2_000_000_000;
let stale_dead = dir.join(format!(".{dead_pid}.beacon.tmp"));
fs::write(&stale_dead, serde_json::to_vec(&beacon(dead_pid)).unwrap()).unwrap();
fs::File::options()
.write(true)
.open(&stale_dead)
.unwrap()
.set_modified(SystemTime::now() - Duration::from_secs(3_600))
.unwrap();
let fresh_dead = dir.join(format!(".{}.json.tmp", dead_pid + 1));
fs::write(
&fresh_dead,
serde_json::to_vec(&heartbeat(dead_pid + 1)).unwrap(),
)
.unwrap();
let live_pid = std::process::id();
let stale_live = dir.join(format!(".{live_pid}.beacon.tmp"));
fs::write(&stale_live, serde_json::to_vec(&beacon(live_pid)).unwrap()).unwrap();
fs::File::options()
.write(true)
.open(&stale_live)
.unwrap()
.set_modified(SystemTime::now() - Duration::from_secs(3_600))
.unwrap();
let report = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.orphan_temps_reaped, 1);
assert!(
!stale_dead.exists(),
"a stale dead-producer temp must be reaped"
);
assert!(fresh_dead.exists(), "a fresh temp may still be in flight");
assert!(
stale_live.exists(),
"a live producer's temp must never be reaped"
);
}
#[cfg(unix)]
#[test]
fn housekeeping_reports_a_live_producer_temp_whose_start_time_is_unreadable_as_unknown() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let live_pid = std::process::id();
let stale_live = dir.join(format!(".{live_pid}.beacon.tmp"));
fs::write(&stale_live, serde_json::to_vec(&beacon(live_pid)).unwrap()).unwrap();
fs::File::options()
.write(true)
.open(&stale_live)
.unwrap()
.set_modified(SystemTime::now() - Duration::from_secs(3_600))
.unwrap();
set_stale_orphan_temp_start_time_override(None);
let report = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.orphan_temps_reaped, 0,
"identity that could not be verified is never trustworthy reap evidence"
);
assert!(
stale_live.exists(),
"an uninspectable producer temp must be retained, not silently dropped"
);
let unknown: Vec<u32> = report.unknown_pids().collect();
assert!(
unknown.contains(&live_pid),
"a live producer temp whose start time cannot be read must be reported as \
Unknown, not silently skipped: {unknown:?}"
);
}
#[cfg(unix)]
#[test]
fn housekeeping_retains_and_reports_malformed_or_mismatched_dead_producer_temps() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let malformed_pid = 2_000_000_010;
let malformed = dir.join(format!(".{malformed_pid}.beacon.tmp"));
fs::write(&malformed, b"not valid json").unwrap();
fs::File::options()
.write(true)
.open(&malformed)
.unwrap()
.set_modified(SystemTime::now() - Duration::from_secs(3_600))
.unwrap();
let mismatched_pid = 2_000_000_020;
let mismatched = dir.join(format!(".{mismatched_pid}.beacon.tmp"));
fs::write(
&mismatched,
serde_json::to_vec(&beacon(mismatched_pid + 1)).unwrap(),
)
.unwrap();
fs::File::options()
.write(true)
.open(&mismatched)
.unwrap()
.set_modified(SystemTime::now() - Duration::from_secs(3_600))
.unwrap();
let report = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.orphan_temps_reaped, 0,
"neither a malformed nor a mismatched dead-PID temp is trustworthy reap evidence"
);
assert!(
malformed.exists(),
"a malformed dead-PID temp must survive cleanup as unknown evidence"
);
assert!(
mismatched.exists(),
"a dead-PID temp whose recorded identity does not match its filename must survive \
cleanup as unknown evidence"
);
let unknown: Vec<u32> = report.unknown_pids().collect();
assert!(
unknown.contains(&malformed_pid),
"malformed evidence must be reported, not silently dropped: {unknown:?}"
);
assert!(
unknown.contains(&mismatched_pid),
"mismatched evidence must be reported, not silently dropped: {unknown:?}"
);
}
#[cfg(unix)]
#[test]
fn housekeeping_refuses_a_symlinked_producer_temp_without_touching_its_target() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let target = root.path().join("forensic-evidence.json");
fs::write(&target, b"keep me").unwrap();
let link = dir.join(".2000000000.beacon.tmp");
std::os::unix::fs::symlink(&target, &link).unwrap();
let report = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.orphan_temps_reaped, 0);
assert_eq!(fs::read(&target).unwrap(), b"keep me");
assert!(link.is_symlink(), "suspicious evidence must be retained");
}
#[cfg(unix)]
#[test]
fn read_only_inspection_classifies_dead_residue_without_deleting_it() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut dead = beacon(2_000_000_000);
dead.started_at = 1;
write_beacon(&dir, &dead).unwrap();
let path = beacon_path(&dir, dead.pid);
let report = inspect_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![dead.pid]);
assert!(
path.exists(),
"diagnostics must never delete sidecar evidence"
);
assert_eq!(report.orphan_temps_reaped, 0);
}
#[cfg(unix)]
#[test]
fn enumerate_live_uncapped_population_has_no_sentinel() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let live = heartbeat(std::process::id());
write_heartbeat(&dir, &live).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(
report
.entries
.iter()
.all(|e| !matches!(e, WalpinPidHealth::Unknown { pid: 0, .. })),
"an in-budget population must not carry the cap sentinel"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_deletes_dead_pid_entry() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut hb = heartbeat(std::process::id());
hb.pid = 2_000_000_000;
hb.started_at = 12345;
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(report.entries.is_empty());
assert!(!dir.join(format!("{}.json", hb.pid)).exists());
}
#[cfg(unix)]
#[test]
fn enumerate_live_deletes_mismatched_start_time_entry() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut hb = heartbeat(std::process::id());
hb.started_at = 1;
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(
report.entries.is_empty(),
"mismatched identity must fail the gate"
);
assert!(!dir.join(format!("{}.json", hb.pid)).exists());
}
#[cfg(unix)]
#[test]
fn enumerate_live_deletes_stale_updated_at_entry() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut hb = heartbeat(std::process::id());
hb.oldest_tx_started_at = None;
hb.updated_at = now_epoch_secs() - 3600; write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.reporting().count(), 0);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![hb.pid]);
assert!(!report.fully_attributed());
assert!(!dir.join(format!("{}.json", hb.pid)).exists());
}
#[cfg(unix)]
#[test]
fn enumerate_live_subsecond_sweep_interval_does_not_collapse_freshness_window() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let hb = heartbeat(std::process::id());
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_millis(200)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"must not be spuriously stale"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_refuses_symlinked_entry_as_unknown_without_touching_target() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let real = root.path().join("elsewhere.txt");
fs::write(&real, b"precious").unwrap();
let link = dir.join("42.json");
std::os::unix::fs::symlink(&real, &link).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(report.reporting().count() == 0);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![42]);
assert!(!report.fully_attributed());
assert_eq!(fs::read_to_string(&real).unwrap(), "precious");
assert!(
link.exists(),
"the symlink itself must not be deleted either"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_refuses_non_owned_entry_before_reading_contents() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let hb = heartbeat(std::process::id());
write_heartbeat(&dir, &hb).unwrap();
let meta = fs::symlink_metadata(dir.join(format!("{}.json", hb.pid))).unwrap();
assert_eq!(
meta.uid(),
current_uid(),
"self-written entries are owned by the current user, exercising the accept path"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_refuses_non_compliant_directory_wholesale() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
fs::create_dir(&dir).unwrap();
fs::set_permissions(&dir, fs::Permissions::from_mode(0o755)).unwrap();
let err = enumerate_live(&dir, Duration::from_secs(5))
.expect_err("non-compliant directory must be refused, not silently enumerated");
assert!(err.to_string().contains("expected 0700"));
}
#[cfg(unix)]
#[test]
fn enumerate_live_missing_directory_is_ok_empty_not_a_failure() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(report.entries.is_empty());
}
#[cfg(unix)]
#[test]
fn enumerate_live_classifies_registered_silent_beacon_with_no_heartbeat() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let b = beacon(std::process::id());
write_beacon(&dir, &b).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.reporting().count(), 0);
assert_eq!(
report.registered_silent_pids().collect::<Vec<_>>(),
vec![std::process::id()]
);
assert!(report.fully_attributed());
}
#[cfg(unix)]
#[test]
fn enumerate_live_reporting_wins_over_registered_silent_for_same_pid() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
write_heartbeat(&dir, &heartbeat(pid)).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.reporting().count(), 1);
assert_eq!(report.registered_silent_pids().count(), 0);
}
#[cfg(unix)]
#[test]
fn housekeeping_uses_the_five_second_legacy_cadence_fallback() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let mut hb = heartbeat(pid);
hb.oldest_tx_started_at = None;
hb.sweep_interval_ms = 0;
hb.updated_at = now_epoch_secs() - 4;
write_heartbeat(&dir, &hb).unwrap();
let report = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"a 4s-old legacy record is inside the ADR-091 15s fallback window; the daemon's \
500ms checkpoint cadence would incorrectly narrow that window to 3s: {report:?}"
);
assert!(
dir.join(format!("{pid}.json")).exists(),
"healthy housekeeping must retain a live legacy record"
);
}
#[cfg(unix)]
#[test]
fn housekeeping_preserves_malformed_unknown_for_no_progress_attribution() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
let heartbeat_path = dir.join(format!("{pid}.json"));
fs::write(&heartbeat_path, b"{not-json").unwrap();
let housekeeping = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(housekeeping.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert_eq!(housekeeping.registered_silent_pids().count(), 0);
assert!(
heartbeat_path.exists(),
"ordinary housekeeping must preserve malformed live-PID evidence"
);
let attribution = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(attribution.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert_eq!(
attribution.registered_silent_pids().count(),
0,
"a fresh beacon must never exonerate a PID whose heartbeat is malformed"
);
assert!(!attribution.fully_attributed());
}
#[cfg(unix)]
#[test]
fn housekeeping_preserves_stale_unknown_for_no_progress_attribution() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
let mut hb = heartbeat(pid);
hb.oldest_tx_started_at = None;
hb.sweep_interval_ms = 1_000;
hb.updated_at = now_epoch_secs() - 30;
write_heartbeat(&dir, &hb).unwrap();
let heartbeat_path = dir.join(format!("{pid}.json"));
let housekeeping = housekeep_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(housekeeping.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert_eq!(housekeeping.registered_silent_pids().count(), 0);
assert!(
heartbeat_path.exists(),
"ordinary housekeeping must preserve a live PID's stale heartbeat evidence"
);
let attribution = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(attribution.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert_eq!(
attribution.registered_silent_pids().count(),
0,
"a fresh beacon must never exonerate a PID whose heartbeat went stale"
);
assert!(!attribution.fully_attributed());
}
#[cfg(unix)]
#[test]
fn enumerate_live_deletes_dead_beacon() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut b = beacon(std::process::id());
b.pid = 2_000_000_001;
b.started_at = 12345;
write_beacon(&dir, &b).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(report.entries.is_empty());
assert!(!beacon_path(&dir, b.pid).exists());
}
#[cfg(unix)]
#[test]
fn write_beacon_refuses_symlinked_target() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let real = root.path().join("elsewhere.txt");
fs::write(&real, b"nope").unwrap();
let b = beacon(999_998);
let target = beacon_path(&dir, b.pid);
std::os::unix::fs::symlink(&real, &target).unwrap();
let err = write_beacon(&dir, &b).expect_err("symlinked target must be refused");
assert!(err.to_string().contains("symlink"));
assert_eq!(fs::read_to_string(&real).unwrap(), "nope");
}
#[cfg(unix)]
#[test]
fn enumerate_live_classifies_stale_beacon_as_unknown() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
let beacon_file = fs::OpenOptions::new()
.write(true)
.open(dir.join(format!("{pid}.beacon")))
.unwrap();
beacon_file
.set_modified(SystemTime::now() - Duration::from_secs(3600))
.unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.registered_silent_pids().count(), 0);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert!(!report.fully_attributed());
assert!(
!beacon_path(&dir, pid).exists(),
"a stale beacon must be deleted, not left to re-classify next sweep"
);
}
#[cfg(unix)]
#[test]
fn enumerate_live_stale_heartbeat_with_fresh_beacon_stays_unknown_not_registered_silent() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
let mut hb = heartbeat(pid);
hb.oldest_tx_started_at = None;
hb.updated_at = now_epoch_secs() - 3600;
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.reporting().count(), 0);
assert_eq!(
report.registered_silent_pids().collect::<Vec<_>>(),
Vec::<u32>::new(),
"a co-existing fresh beacon must not rescue a PID with a stale heartbeat"
);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![pid]);
}
#[test]
#[cfg(target_os = "macos")]
fn census_holders_macos_discovers_self_as_a_holder_of_an_open_db_file() {
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
let file = fs::File::create(&db_path).unwrap();
let census = census_holders(&db_path).expect("census must succeed for a live target");
assert!(
census.holders.contains(&std::process::id()),
"this process holds {db_path:?} open and must appear in its own OS-derived census"
);
assert!(
!census.truncated,
"self was found; the self-canary must not report truncation on its own"
);
drop(file);
}
#[test]
#[cfg(unix)]
fn heartbeat_freshness_uses_producer_cadence_not_enumerator_interval() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let mut slow = heartbeat(pid);
slow.oldest_tx_started_at = None;
slow.sweep_interval_ms = 60_000;
slow.updated_at = now_epoch_secs() - 30;
write_heartbeat(&dir, &slow).unwrap();
let report = enumerate_live(&dir, Duration::from_millis(500)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"a heartbeat 30s old under a recorded 60s cadence is fresh: {report:?}"
);
let mut fast = heartbeat(pid);
fast.oldest_tx_started_at = None;
fast.sweep_interval_ms = 1_000;
fast.updated_at = now_epoch_secs() - 30;
write_heartbeat(&dir, &fast).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(60)).unwrap();
assert_eq!(report.reporting().count(), 0);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![pid]);
}
#[test]
#[cfg(unix)]
fn enumerate_live_new_style_heartbeat_uses_mtime_not_stale_updated_at() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let mut hb = heartbeat(std::process::id());
hb.updated_at = now_epoch_secs() - 3600;
write_heartbeat(&dir, &hb).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"a new-style record with a fresh mtime must classify live regardless \
of a stale `updated_at` body field: {report:?}"
);
}
#[test]
#[cfg(unix)]
fn enumerate_live_new_style_heartbeat_stale_via_mtime_despite_fresh_updated_at() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let hb = heartbeat(pid); write_heartbeat(&dir, &hb).unwrap();
let heartbeat_path = dir.join(format!("{pid}.json"));
let file = fs::OpenOptions::new()
.write(true)
.open(&heartbeat_path)
.unwrap();
file.set_modified(SystemTime::now() - Duration::from_secs(3600))
.unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(report.reporting().count(), 0);
assert_eq!(report.unknown_pids().collect::<Vec<_>>(), vec![pid]);
assert!(
!heartbeat_path.exists(),
"a new-style entry stale by mtime must be deleted, not merely unreported"
);
}
#[test]
#[cfg(unix)]
fn stale_window_from_boundary_inclusive_and_floors_subsecond_cadence_at_three_seconds() {
assert_eq!(stale_window_from(Duration::from_secs(2)), 6);
assert_eq!(stale_window_from(Duration::from_millis(100)), 3);
assert_eq!(stale_window_from(Duration::from_millis(999)), 3);
assert_eq!(stale_window_from(Duration::from_secs(1)), 3);
}
#[test]
#[cfg(unix)]
fn enumerate_live_new_style_heartbeat_boundary_is_inclusive() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let mut hb = heartbeat(pid);
hb.sweep_interval_ms = 2_000;
write_heartbeat(&dir, &hb).unwrap();
let heartbeat_path = dir.join(format!("{pid}.json"));
let touch = |age_secs: u64| {
let file = fs::OpenOptions::new()
.write(true)
.open(&heartbeat_path)
.unwrap();
file.set_modified(SystemTime::now() - Duration::from_secs(age_secs))
.unwrap();
};
touch(6);
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"exactly 3x the declared interval old must still classify live: {report:?}"
);
touch(7);
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.reporting().count(),
0,
"one second past 3x the declared interval must classify stale: {report:?}"
);
}
#[test]
#[cfg(unix)]
fn enumerate_live_new_style_heartbeat_subsecond_cadence_floors_at_three_seconds() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let mut hb = heartbeat(pid);
hb.sweep_interval_ms = 100;
write_heartbeat(&dir, &hb).unwrap();
let heartbeat_path = dir.join(format!("{pid}.json"));
let file = fs::OpenOptions::new()
.write(true)
.open(&heartbeat_path)
.unwrap();
file.set_modified(SystemTime::now() - Duration::from_secs(2))
.unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert_eq!(
report.reporting().count(),
1,
"a 100ms declared cadence must floor its window at 3s, not 1s: 2s old must \
still classify live: {report:?}"
);
}
#[test]
fn attribution_is_evidence_backed_fails_closed_on_missing_or_unrecognized_value() {
let mut hb = heartbeat(std::process::id());
hb.attribution_basis = Some("origin".to_string());
assert!(hb.attribution_is_evidence_backed());
hb.attribution_basis = Some("fallback".to_string());
assert!(!hb.attribution_is_evidence_backed());
hb.attribution_basis = None;
assert!(
!hb.attribution_is_evidence_backed(),
"a missing attribution_basis must never be read as evidence-backed"
);
hb.attribution_basis = Some("some-future-value".to_string());
assert!(
!hb.attribution_is_evidence_backed(),
"an unrecognized value must degrade to fallback-confidence, never guess origin"
);
}
#[test]
#[cfg(unix)]
fn fifo_sidecar_entry_is_refused_without_blocking() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
use std::os::unix::ffi::OsStrExt;
let fifo = dir.join("999999941.json");
let c_path = std::ffi::CString::new(fifo.as_os_str().as_bytes()).unwrap();
let rc = unsafe { libc::mkfifo(c_path.as_ptr(), 0o600) };
assert_eq!(rc, 0, "mkfifo failed: {}", io::Error::last_os_error());
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(
report.unknown_pids().any(|p| p == 999_999_941),
"a FIFO entry must classify its PID as unknown: {report:?}"
);
assert_eq!(
report.registered_silent_pids().collect::<Vec<_>>(),
vec![pid]
);
}
#[test]
#[cfg(unix)]
fn oversized_sidecar_entry_is_refused() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
write_beacon(&dir, &beacon(pid)).unwrap();
fs::write(dir.join("999999942.json"), vec![b'x'; 128 * 1024]).unwrap();
let report = enumerate_live(&dir, Duration::from_secs(5)).unwrap();
assert!(
report.unknown_pids().any(|p| p == 999_999_942),
"an oversized entry must classify its PID as unknown: {report:?}"
);
}
#[test]
#[cfg(unix)]
fn readdir_null_is_error_distinguishes_eof_from_a_real_read_error() {
assert!(
!unix_impl::readdir_null_is_error(0),
"errno == 0 after a NULL readdir means ordinary end-of-directory"
);
assert!(
unix_impl::readdir_null_is_error(libc::EIO),
"a nonzero errno after a NULL readdir means the walk failed mid-stream"
);
}
#[test]
#[cfg(unix)]
fn readdir_failure_mid_walk_propagates_from_list_names_to_enumerate_live() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
unix_impl::set_list_names_readdir_fault(libc::EIO);
let err = enumerate_live(&dir, Duration::from_secs(5)).expect_err(
"a readdir failure mid-walk must propagate as an error, never be folded into a \
truncated-but-otherwise-complete listing",
);
assert_eq!(err.raw_os_error(), Some(libc::EIO));
}
#[test]
#[cfg(unix)]
fn remove_if_same_a_replacement_landing_in_the_recheck_to_unlink_window_is_still_removed() {
use std::os::unix::fs::MetadataExt;
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let handle = unix_impl::SidecarDirHandle::open_or_create(&dir).unwrap();
let name = ".999999999.beacon.tmp";
fs::write(dir.join(name), b"stale evidence").unwrap();
let expected = handle
.read_checked_entry(name)
.unwrap()
.expect("the stale file must be readable and checked");
let expected_inode = fs::metadata(dir.join(name)).unwrap().ino();
let replacement_path = dir.join(name);
let replacement_tmp_path = dir.join(".999999999.beacon.tmp.replacement");
let hook_ran = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let hook_ran_writer = std::sync::Arc::clone(&hook_ran);
let replacement_inode = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
let replacement_inode_writer = std::sync::Arc::clone(&replacement_inode);
unix_impl::set_remove_if_same_race_hook(move || {
fs::write(
&replacement_tmp_path,
b"a producer's brand-new in-flight write",
)
.unwrap();
let new_inode = fs::metadata(&replacement_tmp_path).unwrap().ino();
fs::rename(&replacement_tmp_path, &replacement_path).unwrap();
replacement_inode_writer.store(new_inode, std::sync::atomic::Ordering::SeqCst);
hook_ran_writer.store(true, std::sync::atomic::Ordering::SeqCst);
});
let removed = handle
.remove_if_same(name, &expected)
.expect("remove_if_same must not error when a replacement lands mid-call");
assert!(
hook_ran.load(std::sync::atomic::Ordering::SeqCst),
"the race hook must actually run inside remove_if_same for this test to prove \
anything about the recheck-to-unlink window"
);
assert_ne!(
replacement_inode.load(std::sync::atomic::Ordering::SeqCst),
expected_inode,
"the hook must swap in a genuinely different inode, not rewrite the checked \
one in place, or this test cannot distinguish the race from an ordinary \
same-inode removal"
);
assert!(
removed,
"remove_if_same reports success because the name-based unlink always succeeds, \
even though what it removed is no longer the inode it verified"
);
assert!(
!dir.join(name).exists(),
"a replacement landing in the recheck-to-unlink window is removed too — the \
documented residual race, not one this function actually closes"
);
}
#[test]
#[cfg(any(target_os = "macos", target_os = "linux"))]
fn census_holders_discovers_holder_through_hard_link_path() {
let Some(self_pid) = census_visible_self_pid() else {
eprintln!("skipping census test: process identity source is unavailable");
return;
};
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
fs::File::create(&db_path).unwrap();
let link_path = root.path().join("test-link.db");
fs::hard_link(&db_path, &link_path).unwrap();
let file = fs::File::open(&link_path).unwrap();
let census = census_holders(&db_path).expect("census must succeed for a live target");
assert!(
census.holders.contains(&self_pid),
"this process holds the db open via hard link {link_path:?} and must appear in the census for {db_path:?}"
);
drop(file);
}
#[test]
#[cfg(target_os = "macos")]
fn negotiate_buffer_converges_when_the_set_stops_growing() {
let probe = std::cell::Cell::new(0usize);
let sizes = [4usize, 8, 8]; let (items, truncated) = negotiate_buffer::<i32>(
|| {
let i = probe.get().min(sizes.len() - 1);
sizes[i] as std::os::raw::c_int
},
|_buf_ptr, buf_bytes| {
let i = probe.get();
probe.set(i + 1);
if i < 2 {
buf_bytes
} else {
(buf_bytes as usize - 4) as std::os::raw::c_int
}
},
&|| false,
)
.expect("negotiation must succeed once the set stabilizes");
assert!(
!truncated,
"a snapshot that ends up strictly under capacity must not be marked truncated"
);
assert!(!items.is_empty());
}
#[test]
#[cfg(target_os = "macos")]
fn negotiate_buffer_reports_truncated_after_exhausting_retries() {
let (items, truncated) = negotiate_buffer::<i32>(
|| 4 as std::os::raw::c_int,
|_buf_ptr, buf_bytes| buf_bytes,
&|| false,
)
.expect("negotiation must still return a (possibly truncated) result, not error");
assert!(
truncated,
"a buffer that stays exactly full across every retry must be reported truncated"
);
assert!(!items.is_empty());
}
#[test]
#[cfg(target_os = "macos")]
fn negotiate_buffer_propagates_a_failed_size_call() {
let result = negotiate_buffer::<i32>(
|| -1 as std::os::raw::c_int,
|_buf_ptr, buf_bytes| buf_bytes,
&|| false,
);
assert!(result.is_err(), "a non-positive size probe must error out");
}
#[test]
#[cfg(target_os = "macos")]
fn macos_pid_genuinely_gone_only_true_for_esrch() {
assert!(macos_pid_genuinely_gone(Some(libc::ESRCH)));
assert!(!macos_pid_genuinely_gone(Some(libc::EPERM)));
assert!(!macos_pid_genuinely_gone(Some(libc::EACCES)));
assert!(!macos_pid_genuinely_gone(None));
}
#[test]
#[cfg(target_os = "macos")]
fn proc_pidfdinfo_returned_expected_size_boundary() {
let expected = std::mem::size_of::<u64>(); assert!(
proc_pidfdinfo_returned_expected_size(expected as i32, expected),
"an exact match on the expected struct size must be ok"
);
assert!(
!proc_pidfdinfo_returned_expected_size(expected as i32 - 1, expected),
"a positive but short byte count must be an inspection failure"
);
assert!(
!proc_pidfdinfo_returned_expected_size(0, expected),
"a zero return must be an inspection failure"
);
assert!(
!proc_pidfdinfo_returned_expected_size(-1, expected),
"a negative return must be an inspection failure"
);
}
#[test]
#[cfg(target_os = "linux")]
fn census_holders_linux_discovers_self_as_a_holder_of_an_open_db_file() {
let procfs_usable = fs::read_dir("/proc").is_ok()
&& fs::read_dir("/proc/self/fd").is_ok()
&& census_visible_self_pid().is_some();
if !procfs_usable {
eprintln!("skipping census test: procfs holder inspection is unavailable");
return;
}
let self_pid = census_visible_self_pid().expect("checked above");
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
let file = fs::File::create(&db_path).unwrap();
let census = census_holders(&db_path).expect("census must succeed for a live target");
assert!(
census.holders.contains(&self_pid),
"this process holds {db_path:?} open and must appear in its own OS-derived census"
);
let global_census_supported = fs::metadata("/proc/self/ns/pid")
.is_ok_and(|meta| pid_ns_is_init(meta.ino()))
&& proc_mount_is_visibility_restricted() == Some(false);
if global_census_supported {
assert!(
!census.truncated,
"self was found in an unrestricted init-namespace procfs census"
);
} else {
assert!(
census.truncated,
"a namespace- or mount-restricted procfs census must stay incomplete"
);
}
drop(file);
}
#[test]
#[cfg(target_os = "linux")]
fn linux_proc_gone_only_true_for_not_found() {
assert!(linux_proc_gone(&io::Error::from(io::ErrorKind::NotFound)));
assert!(!linux_proc_gone(&io::Error::from(
io::ErrorKind::PermissionDenied
)));
}
#[test]
#[cfg(target_os = "linux")]
fn pid_ns_is_init_only_true_for_the_fixed_kernel_inode() {
assert!(pid_ns_is_init(PROC_PID_INIT_INO));
assert!(!pid_ns_is_init(PROC_PID_INIT_INO + 1));
assert!(!pid_ns_is_init(0));
assert!(!pid_ns_is_init(12345));
}
#[test]
#[cfg(target_os = "linux")]
fn proc_mount_restricts_visibility_classifies_hidepid_and_subset() {
assert!(!proc_mount_restricts_visibility(
"rw,nosuid,nodev,noexec,relatime"
));
assert!(proc_mount_restricts_visibility("rw,hidepid=2"));
assert!(proc_mount_restricts_visibility("rw,hidepid=invisible"));
assert!(proc_mount_restricts_visibility("rw,hidepid=ptraceable"));
assert!(proc_mount_restricts_visibility("rw,subset=pid"));
assert!(!proc_mount_restricts_visibility("hidepid=0"));
assert!(!proc_mount_restricts_visibility("rw,hidepid=off"));
assert!(proc_mount_restricts_visibility("rw,hidepid"));
}
#[test]
#[cfg(target_os = "linux")]
fn proc_mounts_restricted_in_is_any_restrictive_across_stacked_mounts() {
let clean = "36 25 0:16 / /proc rw,nosuid,nodev,noexec,relatime - proc proc rw";
let restricted = "99 25 0:34 / /proc rw,relatime - proc proc rw,hidepid=2";
let clean_then_restricted = format!("{clean}\n{restricted}");
let restricted_then_clean = format!("{restricted}\n{clean}");
assert_eq!(proc_mounts_restricted_in(clean), Some(false));
assert_eq!(proc_mounts_restricted_in(restricted), Some(true));
assert_eq!(
proc_mounts_restricted_in(&clean_then_restricted),
Some(true)
);
assert_eq!(
proc_mounts_restricted_in(&restricted_then_clean),
Some(true)
);
assert_eq!(
proc_mounts_restricted_in("36 25 0:16 / /sys rw - sysfs sysfs rw"),
None
);
}
#[test]
#[cfg(target_os = "linux")]
fn proc_mount_is_visibility_restricted_reads_this_hosts_own_proc_mount() {
assert!(
proc_mount_is_visibility_restricted().is_some(),
"expected to find and parse this process's own /proc mount entry in \
/proc/self/mountinfo"
);
}
#[test]
fn census_result_is_complete_reflects_uninspectable_pids() {
let complete = CensusResult {
holders: std::collections::HashSet::from([1, 2]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
assert!(complete.is_complete());
let incomplete = CensusResult {
holders: std::collections::HashSet::from([1]),
uninspectable_pids: vec![7],
truncated: false,
budget_exhausted: false,
};
assert!(!incomplete.is_complete());
}
#[test]
#[cfg(any(target_os = "macos", target_os = "linux"))]
fn a_spent_budget_truncates_the_census_instead_of_failing_it() {
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
let handle = fs::File::create(&db_path).unwrap();
let unbounded = census_holders(&db_path).expect("unbounded census of a live target");
assert!(
!unbounded.budget_exhausted,
"an unbounded census has no budget to spend"
);
let bounded = census_holders_until_within(&db_path, || false, Duration::ZERO)
.expect("a spent budget returns the partial census, never an error");
assert!(
bounded.budget_exhausted,
"a walk stopped by its budget must say so; a caller cannot otherwise \
tell a bounded answer from a complete one"
);
assert!(
bounded.truncated,
"budget exhaustion is an incompleteness signal, folded into the same \
field every other incompleteness uses"
);
assert!(
!bounded.is_complete(),
"the consumer branches on is_complete(); a bounded census that reads \
complete is the defect this bound would otherwise introduce"
);
drop(handle);
}
#[test]
#[cfg(any(target_os = "macos", target_os = "linux"))]
fn a_cancel_and_a_spent_budget_leave_by_opposite_exits() {
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
let handle = fs::File::create(&db_path).unwrap();
let cancelled =
census_holders_until(&db_path, || true).expect_err("a cancellation is an error by design");
assert_eq!(
cancelled.kind(),
io::ErrorKind::Interrupted,
"a cancelled census reports Interrupted"
);
let bounded = census_holders_until_within(&db_path, || false, Duration::ZERO)
.expect("a spent budget is not a cancellation");
assert!(bounded.budget_exhausted);
drop(handle);
}
#[test]
#[cfg(any(target_os = "macos", target_os = "linux"))]
fn an_unspent_budget_is_invisible_in_the_result() {
let root = tempfile::tempdir().unwrap();
let db_path = root.path().join("test.db");
let handle = fs::File::create(&db_path).unwrap();
let bounded = census_holders_until_within(&db_path, || false, Duration::from_secs(300))
.expect("census of a live target");
assert!(
!bounded.budget_exhausted,
"a five-minute budget cannot be spent by one process walk on a test host"
);
drop(handle);
}
#[test]
fn census_result_is_complete_reflects_truncated() {
let truncated = CensusResult {
holders: std::collections::HashSet::from([1]),
uninspectable_pids: Vec::new(),
truncated: true,
budget_exhausted: false,
};
assert!(!truncated.is_complete());
}
#[test]
fn self_canary_marks_truncated_when_own_pid_missing() {
let scanner_pid = 41;
let mut census = CensusResult {
holders: std::collections::HashSet::from([42]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
census.apply_self_canary_for(Some(scanner_pid));
assert!(
census.truncated,
"a census that discovered other holders but not the calling process itself is \
positive proof of a missed enumeration and must be marked incomplete"
);
}
#[test]
fn self_canary_leaves_a_correct_census_untouched() {
let scanner_pid = 41;
let mut census = CensusResult {
holders: std::collections::HashSet::from([scanner_pid]),
uninspectable_pids: Vec::new(),
truncated: false,
budget_exhausted: false,
};
census.apply_self_canary_for(Some(scanner_pid));
assert!(
!census.truncated,
"self was found; the canary must not fire"
);
}
#[test]
fn self_canary_does_not_clear_an_existing_truncated_flag() {
let scanner_pid = 41;
let mut census = CensusResult {
holders: std::collections::HashSet::from([scanner_pid]),
uninspectable_pids: Vec::new(),
truncated: true,
budget_exhausted: false,
};
census.apply_self_canary_for(Some(scanner_pid));
assert!(
census.truncated,
"the self-canary only ever sets `truncated`; it must never clear a flag another \
step already raised"
);
}
#[test]
fn self_canary_fails_closed_when_process_identity_is_unavailable() {
let mut census = CensusResult::default();
census.apply_self_canary_for(None);
assert!(census.truncated);
}
#[cfg(all(test, windows))]
mod windows_tests {
use super::*;
use std::fs;
#[test]
fn ensure_sidecar_dir_creates_directory() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).expect("should create");
let meta = fs::symlink_metadata(&dir).unwrap();
assert!(meta.is_dir());
}
#[test]
fn ensure_sidecar_dir_is_idempotent_and_revalidates_dacl() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).expect("first create should succeed");
ensure_sidecar_dir(&dir)
.expect("second call must re-open and re-validate the existing dir, not fail");
}
#[test]
fn ensure_sidecar_dir_refuses_preexisting_dir_with_default_acl() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
fs::create_dir(&dir).unwrap();
let err = ensure_sidecar_dir(&dir)
.expect_err("a pre-existing dir without the exact owner-only DACL must be refused");
assert!(err.to_string().contains("owner"), "unexpected error: {err}");
}
#[test]
fn ensure_sidecar_dir_refuses_symlinked_target() {
let root = tempfile::tempdir().unwrap();
let real = root.path().join("real_dir");
fs::create_dir(&real).unwrap();
let link = root.path().join("khive.db.walpin");
std::os::windows::fs::symlink_dir(&real, &link).expect(
"creating a directory symlink requires Developer Mode or an elevated \
process on the Windows CI runner",
);
let err = ensure_sidecar_dir(&link)
.expect_err("a reparse-point sidecar path must be refused, never followed");
assert!(
err.to_string().contains("reparse"),
"unexpected error: {err}"
);
}
#[test]
fn write_heartbeat_creates_then_replaces_then_removes() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let path = dir.join(format!("{pid}.json"));
let first = heartbeat(pid);
write_heartbeat(&dir, &first).expect("initial create must succeed");
let read_back: WalpinHeartbeat =
serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(read_back, first);
let mut second = heartbeat(pid);
second.oldest_tx_label = Some("replaced".to_string());
write_heartbeat(&dir, &second).expect("replacing an already-existing target must succeed");
let read_back: WalpinHeartbeat =
serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(read_back, second);
assert_ne!(read_back, first);
remove_heartbeat(&dir, pid).expect("remove must succeed");
assert!(!path.exists());
}
#[test]
fn replacing_heartbeat_keeps_old_target_until_rename() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let path = dir.join(format!("{pid}.json"));
let first = heartbeat(pid);
write_heartbeat(&dir, &first).unwrap();
let old_body = fs::read(&path).unwrap();
let mut second = heartbeat(pid);
second.oldest_tx_label = Some("replacement".to_string());
let hook_ran = std::rc::Rc::new(std::cell::Cell::new(false));
let hook_ran_inside = std::rc::Rc::clone(&hook_ran);
super::super::windows_impl::set_before_target_rename_hook(move || {
assert_eq!(
fs::read(&path).expect("old target must still exist before rename"),
old_body,
"the old heartbeat must remain at the target until replacement"
);
hook_ran_inside.set(true);
});
write_heartbeat(&dir, &second).expect("replacement write must succeed");
assert!(hook_ran.get(), "the inspection-to-rename hook must run");
assert_eq!(
fs::read(dir.join(format!("{pid}.json"))).unwrap(),
serde_json::to_vec(&second).unwrap()
);
}
#[test]
fn write_heartbeat_refuses_directory_target_without_removing_it() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let hb = heartbeat(std::process::id());
let target = dir.join(format!("{}.json", hb.pid));
fs::create_dir(&target).unwrap();
write_heartbeat(&dir, &hb).expect_err("directory target must be refused");
assert!(fs::symlink_metadata(&target).unwrap().is_dir());
}
#[test]
fn write_heartbeat_refuses_reparse_target_without_removing_it() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let outside = root.path().join("outside.txt");
fs::write(&outside, b"untouched").unwrap();
let hb = heartbeat(std::process::id());
let target = dir.join(format!("{}.json", hb.pid));
std::os::windows::fs::symlink_file(&outside, &target).expect(
"creating a file symlink requires Developer Mode or an elevated \
process on the Windows CI runner",
);
write_heartbeat(&dir, &hb).expect_err("reparse target must be refused");
assert!(fs::symlink_metadata(&target)
.unwrap()
.file_type()
.is_symlink());
assert_eq!(fs::read(&outside).unwrap(), b"untouched");
}
#[test]
fn repeated_heartbeat_write_validates_sidecar_root_once() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let first = heartbeat(pid);
write_heartbeat(&dir, &first).expect("initial write must create the sidecar");
let before = super::super::windows_impl::open_dir_handle_call_count();
let mut replacement = first;
replacement.oldest_tx_label = Some("replacement".to_string());
write_heartbeat(&dir, &replacement).expect("replacement write must succeed");
let validations = super::super::windows_impl::open_dir_handle_call_count() - before;
assert_eq!(
validations, 1,
"an existing sidecar root must be fully validated exactly once per record write"
);
}
#[test]
fn remove_heartbeat_on_missing_sidecar_dir_is_a_noop() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
remove_heartbeat(&dir, 4242).expect("missing sidecar dir must be a no-op");
assert!(
!dir.exists(),
"removal must never create the sidecar dir as a side effect"
);
}
#[test]
fn touch_heartbeat_refreshes_mtime_without_changing_content() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let hb = heartbeat(pid);
write_heartbeat(&dir, &hb).unwrap();
let path = dir.join(format!("{pid}.json"));
let before = fs::metadata(&path).unwrap().modified().unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
touch_heartbeat(&dir, pid).expect("touch of an existing heartbeat must succeed");
let after = fs::metadata(&path).unwrap().modified().unwrap();
assert!(
after > before,
"touch must advance the mtime; an unchanged timestamp means the \
refresh was a no-op (before {before:?}, after {after:?})"
);
let content_after: WalpinHeartbeat =
serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(
content_after, hb,
"touch is metadata-only; the body must be unchanged"
);
}
#[test]
fn touch_heartbeat_fails_when_entry_is_absent() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
ensure_sidecar_dir(&dir).unwrap();
let err =
touch_heartbeat(&dir, 99999).expect_err("touching a nonexistent heartbeat must fail");
assert!(
err.to_string().contains("does not exist"),
"unexpected error: {err}"
);
}
#[test]
fn beacon_write_touch_remove_cycle() {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("khive.db.walpin");
let pid = std::process::id();
let b = beacon(pid);
let path = beacon_path(&dir, pid);
write_beacon(&dir, &b).expect("beacon create must succeed");
assert!(path.exists());
let before = fs::metadata(&path).unwrap().modified().unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
touch_beacon(&dir, pid).expect("beacon touch must succeed");
let after = fs::metadata(&path).unwrap().modified().unwrap();
assert!(
after > before,
"beacon touch must advance the mtime; an unchanged timestamp means \
the refresh was a no-op (before {before:?}, after {after:?})"
);
remove_beacon(&dir, pid).expect("beacon remove must succeed");
assert!(!path.exists());
}
}