use chrono::Utc;
use kranz_engine::error::EngineError;
use kranz_engine::event_log::{EventLog, LockForce};
use kranz_engine::events::{Event, EventKind};
use kranz_engine::paths::MissionPaths;
use std::io::Write;
use std::time::Duration;
const MISSION: &str = "m-test";
const NEVER: Duration = Duration::from_secs(3600);
fn paths(dir: &std::path::Path) -> MissionPaths {
MissionPaths::new(dir, MISSION)
}
fn lifecycle(text: &str) -> EventKind {
EventKind::UserMessage {
text: text.to_string(),
interrupt: false,
}
}
fn delta(content: &str) -> EventKind {
EventKind::WorkerMessage {
run_id: "r-1".to_string(),
tag: "text".to_string(),
content: content.to_string(),
}
}
fn write_raw_log(path: &std::path::Path, lines: &[String]) {
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
let mut f = std::fs::File::create(path).unwrap();
for line in lines {
writeln!(f, "{line}").unwrap();
}
}
fn raw_event(seq: u64, kind: EventKind) -> String {
serde_json::to_string(&Event {
seq,
ts: Utc::now(),
mission_id: MISSION.to_string(),
kind,
})
.unwrap()
}
fn now_epoch_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs()
}
#[test]
fn append_redacting_scrubs_secret_before_persisting_event() {
let tmp = tempfile::tempdir().unwrap();
let paths = paths(tmp.path());
let mut log = EventLog::acquire(&paths, MISSION, NEVER, LockForce::No).unwrap();
let secret = "sk-ant-api03-AbCdEf_123-xyz";
let (event, findings) = log
.append_redacting(EventKind::UserMessage {
text: format!("use {secret}"),
interrupt: false,
})
.unwrap();
drop(log);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].rule_id, "anthropic-api-key");
match event.kind {
EventKind::UserMessage { text, .. } => assert_eq!(text, "use [REDACTED]"),
other => panic!("wrong event: {other:?}"),
}
let raw = std::fs::read_to_string(paths.events_file()).unwrap();
assert!(!raw.contains(secret), "event log leaked secret: {raw}");
assert!(raw.contains("[REDACTED]"));
}
#[test]
fn append_emits_secret_redacted_audit_event() {
let tmp = tempfile::tempdir().unwrap();
let paths = paths(tmp.path());
let mut log = EventLog::acquire(&paths, MISSION, NEVER, LockForce::No).unwrap();
let secret = "sk-ant-api03-AbCdEf_123-xyz";
log.append(EventKind::UserMessage {
text: format!("use {secret}"),
interrupt: false,
})
.unwrap();
drop(log);
let events = EventLog::read_events(&paths.events_file()).unwrap();
assert_eq!(events.len(), 2);
assert!(matches!(
&events[0].kind,
EventKind::UserMessage { text, .. } if text == "use [REDACTED]"
));
assert!(matches!(
&events[1].kind,
EventKind::SecretRedacted {
rule_id,
location,
..
} if rule_id == "anthropic-api-key" && location.contains("/payload/text")
));
}
#[cfg(target_os = "linux")]
fn identity_token_for(pid: u32) -> String {
let boot_id = std::fs::read_to_string("/proc/sys/kernel/random/boot_id").unwrap();
let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).unwrap();
let rest = &stat[stat.rfind(')').unwrap() + 1..];
let ticks = rest.split_whitespace().nth(19).unwrap();
format!("{}:{}", boot_id.trim(), ticks)
}
#[cfg(target_os = "macos")]
fn identity_token_for(pid: u32) -> String {
let out = std::process::Command::new("ps")
.env("LC_ALL", "C")
.env("TZ", "UTC")
.args(["-p", &pid.to_string(), "-o", "lstart="])
.output()
.unwrap();
assert!(
out.status.success(),
"ps -o lstart= must succeed for a live pid"
);
String::from_utf8_lossy(&out.stdout).trim().to_string()
}
#[cfg(target_os = "macos")]
fn ps_can_execute() -> bool {
std::process::Command::new("ps")
.args(["-p", &std::process::id().to_string(), "-o", "command="])
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
#[cfg(unix)]
struct LiveHolder(std::process::Child);
#[cfg(unix)]
impl LiveHolder {
fn spawn() -> Self {
LiveHolder(
std::process::Command::new("sleep")
.arg("300")
.spawn()
.expect("spawn sleep child"),
)
}
fn pid(&self) -> u32 {
self.0.id()
}
}
#[cfg(unix)]
impl Drop for LiveHolder {
fn drop(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
#[test]
fn acquire_creates_dirs_and_lock() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let before = now_epoch_secs();
let log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert!(p.mission_dir().is_dir());
assert!(p.runs_dir().is_dir());
assert!(p.control_dir().is_dir());
assert!(p.lock_file().is_file());
let contents = std::fs::read_to_string(p.lock_file()).unwrap();
let mut lines = contents.lines();
assert_eq!(lines.next().unwrap(), std::process::id().to_string());
let acquired: u64 = lines
.next()
.expect("second lock line: acquire epoch secs")
.parse()
.expect("acquire time must be an integer");
assert!(
acquired >= before && acquired <= now_epoch_secs(),
"acquire time {acquired} outside [{before}, now]"
);
#[cfg(target_os = "linux")]
assert_eq!(
lines
.next()
.expect("third lock line: identity token")
.trim(),
identity_token_for(std::process::id()),
"the recorded token must be OUR OWN process identity"
);
#[cfg(target_os = "macos")]
{
let recorded = lines
.next()
.expect("third lock line: identity token")
.trim()
.to_string();
assert!(
!recorded.is_empty(),
"a non-empty token is always recorded on macOS, wrapped or not"
);
if ps_can_execute() {
assert_eq!(
recorded,
identity_token_for(std::process::id()),
"the recorded token must be OUR OWN process identity"
);
} else {
eprintln!(
"SKIP-UNDER-WRAP (gate-sandbox-supervision-dogfood): \
acquire_creates_dirs_and_lock — /bin/ps cannot execute inside the gate \
sandbox wrap; the ps-computed token comparison is skipped (non-empty \
recording still asserted)"
);
}
}
assert_eq!(log.last_seq(), 0);
}
#[test]
fn second_acquire_fails_with_lock_held_naming_pid() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let _held = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
let err = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap_err();
match err {
EngineError::LockHeld(msg) => {
assert!(
msg.contains(&std::process::id().to_string()),
"message should name the holding pid: {msg}"
);
}
other => panic!("expected LockHeld, got {other:?}"),
}
}
#[test]
fn even_if_live_steals_lock_from_live_holder_if_not_live_refuses() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let first = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
let err = EventLog::acquire(&p, MISSION, NEVER, LockForce::IfNotLive).unwrap_err();
match err {
EngineError::LockHeld(msg) => assert!(
msg.contains("dangerously-steal-live-lock"),
"live-holder refusal must point at the stronger flag: {msg}"
),
other => panic!("expected LockHeld, got {other:?}"),
}
let mut stolen = EventLog::acquire(&p, MISSION, NEVER, LockForce::EvenIfLive).unwrap();
stolen.append(lifecycle("after steal")).unwrap();
drop(first);
drop(stolen);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 1);
}
#[test]
fn drop_removes_lock_file() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert!(p.lock_file().exists());
drop(log);
assert!(!p.lock_file().exists());
let _again = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
}
#[test]
fn stolen_from_drop_leaves_stealers_lock() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let first = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
let mut stealer = EventLog::acquire(&p, MISSION, NEVER, LockForce::EvenIfLive).unwrap();
assert!(p.lock_file().exists());
drop(first);
assert!(
p.lock_file().exists(),
"stolen-from Drop must not remove the stealer's lock"
);
stealer.append(lifecycle("still holding")).unwrap();
drop(stealer);
assert!(!p.lock_file().exists());
}
#[test]
fn failed_acquire_releases_lock() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("hello")).unwrap();
}
let err = EventLog::acquire(&p, "m-other", NEVER, LockForce::No).unwrap_err();
assert!(matches!(err, EngineError::InvalidState(_)), "got {err:?}");
assert!(!p.lock_file().exists());
let _ok = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
}
#[test]
fn append_assigns_contiguous_seq_and_round_trips() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
let e1 = log.append(lifecycle("one")).unwrap();
let e2 = log.append(lifecycle("two")).unwrap();
let e3 = log
.append(EventKind::MilestoneStarted {
milestone_id: "ms-1".to_string(),
start_sha: "abc123".to_string(),
})
.unwrap();
assert_eq!((e1.seq, e2.seq, e3.seq), (1, 2, 3));
assert_eq!(e1.mission_id, MISSION);
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(events[0].seq, 1);
match &events[2].kind {
EventKind::MilestoneStarted {
milestone_id,
start_sha,
} => {
assert_eq!(milestone_id, "ms-1");
assert_eq!(start_sha, "abc123");
}
other => panic!("wrong kind round-tripped: {other:?}"),
}
}
#[test]
fn lifecycle_events_are_durable_immediately() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("durable")).unwrap();
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].seq, 1);
}
#[test]
fn deltas_buffer_and_lifecycle_drains_in_order() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("L1")).unwrap();
log.append(delta("d1")).unwrap();
log.append(delta("d2")).unwrap();
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 1);
log.append(lifecycle("L2")).unwrap();
let events = EventLog::read_events(&p.events_file()).unwrap();
let kinds: Vec<&str> = events.iter().map(|e| e.kind.type_name()).collect();
assert_eq!(
kinds,
vec![
"user.message",
"worker.message",
"worker.message",
"user.message"
]
);
assert_eq!(
events.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![1, 2, 3, 4]
);
let contents: Vec<String> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::WorkerMessage { content, .. } => Some(content.clone()),
_ => None,
})
.collect();
assert_eq!(contents, vec!["d1", "d2"]);
}
#[test]
fn throttle_flushes_buffer_by_age() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, Duration::from_millis(30), LockForce::No).unwrap();
log.append(delta("d1")).unwrap();
assert_eq!(
EventLog::read_events(&p.events_file()).unwrap().len(),
0,
"young delta stays buffered"
);
std::thread::sleep(Duration::from_millis(60));
log.append(delta("d2")).unwrap();
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 2);
}
#[test]
fn flush_if_due_drains_idle_buffer_by_age() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, Duration::from_millis(30), LockForce::No).unwrap();
log.append(delta("d1")).unwrap();
assert_eq!(
EventLog::read_events(&p.events_file()).unwrap().len(),
0,
"young delta stays buffered"
);
std::thread::sleep(Duration::from_millis(60));
let due = log.flush_if_due().unwrap();
assert!(due, "flush_if_due must report that it drained the buffer");
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 1);
}
#[test]
fn buffer_age_reports_oldest_and_none_when_empty() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.buffer_age(), None, "fresh log has no buffered deltas");
log.append(delta("d1")).unwrap();
assert!(
log.buffer_age().is_some(),
"buffer_age must report the oldest delta's age"
);
let due = log.flush_if_due().unwrap();
assert!(
!due,
"flush_if_due must not drain a buffer younger than the throttle"
);
assert_eq!(
EventLog::read_events(&p.events_file()).unwrap().len(),
0,
"delta must remain buffered"
);
}
#[test]
fn explicit_flush_drains_buffer() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(delta("d1")).unwrap();
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 0);
log.flush().unwrap();
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 1);
}
#[test]
fn drop_flushes_buffered_deltas() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(delta("d1")).unwrap();
log.append(delta("d2")).unwrap();
}
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[1].seq, 2);
}
#[test]
fn reacquire_resumes_seq_from_existing_log() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("one")).unwrap();
log.append(lifecycle("two")).unwrap();
}
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 2);
let e = log.append(lifecycle("three")).unwrap();
assert_eq!(e.seq, 3);
drop(log);
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 3);
}
#[test]
fn read_events_refuses_seq_gap() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(
&p.events_file(),
&[raw_event(1, lifecycle("a")), raw_event(3, lifecycle("b"))],
);
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(matches!(err, EngineError::LogCorruption(_)), "got {err:?}");
}
#[test]
fn read_events_refuses_duplicate_seq() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(
&p.events_file(),
&[raw_event(1, lifecycle("a")), raw_event(1, lifecycle("b"))],
);
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(matches!(err, EngineError::LogCorruption(_)), "got {err:?}");
}
#[test]
fn read_events_requires_seq_starting_at_one() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(&p.events_file(), &[raw_event(2, lifecycle("a"))]);
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(matches!(err, EngineError::LogCorruption(_)), "got {err:?}");
}
#[test]
fn torn_final_line_is_dropped_silently() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(
&p.events_file(),
&[
raw_event(1, lifecycle("a")),
raw_event(2, lifecycle("b")),
r#"{"seq":3,"ts":"2026-01-01T00:0"#.to_string(), ],
);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events.last().unwrap().seq, 2);
}
#[test]
fn torn_middle_line_is_corruption() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(
&p.events_file(),
&[
raw_event(1, lifecycle("a")),
"{not json".to_string(),
raw_event(2, lifecycle("b")),
],
);
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(matches!(err, EngineError::LogCorruption(_)), "got {err:?}");
}
#[test]
fn empty_log_reads_as_no_events() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
write_raw_log(&p.events_file(), &[]);
assert!(EventLog::read_events(&p.events_file()).unwrap().is_empty());
}
fn append_raw_bytes(path: &std::path::Path, bytes: &[u8]) {
let mut f = std::fs::OpenOptions::new().append(true).open(path).unwrap();
f.write_all(bytes).unwrap();
}
#[test]
fn reacquire_truncates_torn_final_line_without_newline() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("one")).unwrap();
log.append(lifecycle("two")).unwrap();
}
append_raw_bytes(&p.events_file(), br#"{"seq":3,"ts":"2026-01-01T00:0"#);
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 2, "torn line must not count toward seq");
log.append(lifecycle("three")).unwrap();
log.append(lifecycle("four")).unwrap();
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(
events.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![1, 2, 3, 4]
);
}
#[test]
fn reacquire_truncates_torn_final_line_with_newline() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("one")).unwrap();
log.append(lifecycle("two")).unwrap();
}
append_raw_bytes(&p.events_file(), b"{\"seq\":3,\"ts\":\"2026-01-01T00:0\n");
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 2);
log.append(lifecycle("three")).unwrap();
log.append(lifecycle("four")).unwrap();
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(
events.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![1, 2, 3, 4]
);
}
#[test]
fn reacquire_repairs_valid_final_line_missing_its_newline() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
std::fs::create_dir_all(p.events_file().parent().unwrap()).unwrap();
let mut f = std::fs::File::create(p.events_file()).unwrap();
writeln!(f, "{}", raw_event(1, lifecycle("a"))).unwrap();
write!(f, "{}", raw_event(2, lifecycle("b"))).unwrap(); drop(f);
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 2, "unterminated valid line must survive");
log.append(lifecycle("c")).unwrap();
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(
events.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![1, 2, 3]
);
}
#[test]
fn reacquire_truncates_log_that_is_only_a_torn_line() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
std::fs::create_dir_all(p.events_file().parent().unwrap()).unwrap();
std::fs::write(p.events_file(), b"{\"seq\":1,\"ts").unwrap();
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 0);
log.append(lifecycle("first")).unwrap();
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1]);
}
#[test]
fn torn_final_line_splitting_multibyte_char_is_dropped_on_read() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("one")).unwrap();
}
append_raw_bytes(&p.events_file(), b"{\"seq\":2,\"ts\":\"2026\xE2");
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].seq, 1);
}
#[test]
fn reacquire_truncates_torn_line_splitting_multibyte_char() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(lifecycle("one")).unwrap();
}
append_raw_bytes(&p.events_file(), b"{\"seq\":2,\"ts\":\"2026\xE2");
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
assert_eq!(log.last_seq(), 1);
log.append(lifecycle("two")).unwrap();
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1, 2]);
}
#[test]
fn read_events_after_returns_suffix() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
{
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
for i in 1..=5 {
log.append(lifecycle(&format!("e{i}"))).unwrap();
}
}
let tail = EventLog::read_events_after(&p.events_file(), 2).unwrap();
assert_eq!(
tail.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![3, 4, 5]
);
let all = EventLog::read_events_after(&p.events_file(), 0).unwrap();
assert_eq!(all.len(), 5);
let none = EventLog::read_events_after(&p.events_file(), 5).unwrap();
assert!(none.is_empty());
}
#[test]
fn read_tail_events_window_on_line_boundary_keeps_the_full_line() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let lines: Vec<String> = (1..=3).map(|i| raw_event(i, lifecycle("e"))).collect();
write_raw_log(&p.events_file(), &lines);
let window = (lines[1].len() + 1 + lines[2].len() + 1) as u64;
let tail = EventLog::read_tail_events(&p.events_file(), window).unwrap();
assert_eq!(
tail.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![2, 3],
"a boundary-aligned window must not drop its first complete line"
);
}
#[test]
fn read_tail_events_window_mid_line_drops_only_the_torn_head() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let lines: Vec<String> = (1..=3).map(|i| raw_event(i, lifecycle("e"))).collect();
write_raw_log(&p.events_file(), &lines);
let window = (lines[2].len() + 1 + 3) as u64;
let tail = EventLog::read_tail_events(&p.events_file(), window).unwrap();
assert_eq!(tail.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![3]);
}
#[test]
fn read_tail_events_window_covering_whole_file_returns_all_events() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
let lines: Vec<String> = (1..=3).map(|i| raw_event(i, lifecycle("e"))).collect();
write_raw_log(&p.events_file(), &lines);
let total: u64 = lines.iter().map(|l| (l.len() + 1) as u64).sum();
for window in [total, total + 1024, u64::MAX] {
let tail = EventLog::read_tail_events(&p.events_file(), window).unwrap();
assert_eq!(
tail.iter().map(|e| e.seq).collect::<Vec<_>>(),
vec![1, 2, 3],
"window {window} covers the whole file"
);
}
}
#[test]
fn read_tail_events_empty_file_reads_as_no_events() {
let dir = tempfile::tempdir().unwrap();
let p = paths(dir.path());
std::fs::create_dir_all(p.events_file().parent().unwrap()).unwrap();
std::fs::write(p.events_file(), b"").unwrap();
assert!(EventLog::read_tail_events(&p.events_file(), 4096)
.unwrap()
.is_empty());
}
#[cfg(unix)]
#[test]
fn dead_holder_lock_is_stolen_at_every_tier() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
for force in [LockForce::No, LockForce::IfNotLive, LockForce::EvenIfLive] {
std::fs::write(paths.lock_file(), i32::MAX.to_string()).unwrap();
let log = EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), force)
.unwrap_or_else(|e| panic!("dead holder must be stolen at tier {force:?}: {e}"));
drop(log); }
}
#[cfg(unix)]
#[test]
fn alive_holder_lock_needs_the_dangerous_tier() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
let holder = LiveHolder::spawn();
std::fs::write(paths.lock_file(), holder.pid().to_string()).unwrap();
let err =
EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), LockForce::No).unwrap_err();
match err {
EngineError::LockHeld(msg) => assert!(
msg.contains("--force-lock") && msg.contains(&holder.pid().to_string()),
"no-force refusal keeps today's message shape: {msg}"
),
other => panic!("expected LockHeld, got {other:?}"),
}
let err = EventLog::acquire(
&paths,
"m-lock",
Duration::from_millis(50),
LockForce::IfNotLive,
)
.unwrap_err();
match err {
EngineError::LockHeld(msg) => {
assert!(msg.contains("ALIVE"), "must say the holder is alive: {msg}");
assert!(
msg.contains(&format!("ps -p {}", holder.pid())),
"must suggest identifying the holder: {msg}"
);
assert!(
msg.contains("--dangerously-steal-live-lock"),
"must name the stronger flag: {msg}"
);
}
other => panic!("expected LockHeld, got {other:?}"),
}
let log = EventLog::acquire(
&paths,
"m-lock",
Duration::from_millis(50),
LockForce::EvenIfLive,
)
.expect("EvenIfLive must steal even from a live holder");
drop(log);
drop(holder);
}
#[test]
fn unknown_holder_lock_yields_to_any_force_tier() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
std::fs::write(paths.lock_file(), "not-a-pid").unwrap();
let err =
EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), LockForce::No).unwrap_err();
match err {
EngineError::LockHeld(msg) => assert!(
msg.contains("--force-lock"),
"unknown-holder refusal mentions --force-lock: {msg}"
),
other => panic!("expected LockHeld, got {other:?}"),
}
for force in [LockForce::IfNotLive, LockForce::EvenIfLive] {
std::fs::write(paths.lock_file(), "not-a-pid").unwrap();
let log = EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), force)
.unwrap_or_else(|e| panic!("unknown holder must yield to {force:?}: {e}"));
drop(log);
}
std::fs::write(paths.lock_file(), "-7").unwrap();
let err =
EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), LockForce::No).unwrap_err();
assert!(matches!(err, EngineError::LockHeld(_)));
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn reused_pid_lock_is_stale_at_every_tier() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
let holder = LiveHolder::spawn();
for force in [LockForce::No, LockForce::IfNotLive, LockForce::EvenIfLive] {
std::fs::write(
paths.lock_file(),
format!(
"{}\n{}\nsome-other-boot-id:12345\n",
holder.pid(),
now_epoch_secs()
),
)
.unwrap();
let log = EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), force)
.unwrap_or_else(|e| panic!("reused pid means dead writer; {force:?} must steal: {e}"));
drop(log);
}
drop(holder);
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn live_holder_with_matching_token_survives_clock_steps() {
#[cfg(target_os = "macos")]
if !ps_can_execute() {
eprintln!(
"SKIP-UNDER-WRAP (gate-sandbox-supervision-dogfood): \
live_holder_with_matching_token_survives_clock_steps — /bin/ps cannot execute \
inside the gate sandbox wrap, so the reference token cannot be computed; \
skipping"
);
return;
}
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
let holder = LiveHolder::spawn();
let token = identity_token_for(holder.pid());
assert!(!token.is_empty(), "live child must have an identity token");
std::fs::write(
paths.lock_file(),
format!("{}\n{}\n{}\n", holder.pid(), now_epoch_secs() - 3600, token),
)
.unwrap();
let err =
EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), LockForce::No).unwrap_err();
assert!(matches!(err, EngineError::LockHeld(_)), "got {err:?}");
let err = EventLog::acquire(
&paths,
"m-lock",
Duration::from_millis(50),
LockForce::IfNotLive,
)
.unwrap_err();
match err {
EngineError::LockHeld(msg) => {
assert!(
msg.contains("ALIVE"),
"matching token ⇒ alive holder refusal: {msg}"
)
}
other => panic!("expected LockHeld, got {other:?}"),
}
drop(holder);
}
#[cfg(unix)]
#[test]
fn two_line_lock_without_token_degrades_to_plain_liveness() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
let holder = LiveHolder::spawn();
std::fs::write(
paths.lock_file(),
format!("{}\n{}\n", holder.pid(), now_epoch_secs() - 3600),
)
.unwrap();
let err = EventLog::acquire(
&paths,
"m-lock",
Duration::from_millis(50),
LockForce::IfNotLive,
)
.unwrap_err();
match err {
EngineError::LockHeld(msg) => {
assert!(msg.contains("ALIVE"), "tokenless lock ⇒ plain alive: {msg}")
}
other => panic!("expected LockHeld, got {other:?}"),
}
drop(holder);
}
#[cfg(unix)]
#[test]
fn garbage_acquire_time_degrades_to_plain_liveness() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
let holder = LiveHolder::spawn();
std::fs::write(paths.lock_file(), format!("{}\nnot-a-time\n", holder.pid())).unwrap();
let err = EventLog::acquire(
&paths,
"m-lock",
Duration::from_millis(50),
LockForce::IfNotLive,
)
.unwrap_err();
assert!(matches!(err, EngineError::LockHeld(_)), "got {err:?}");
drop(holder);
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn own_pid_with_foreign_token_is_provably_reused() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
std::fs::write(
paths.lock_file(),
format!(
"{}\n{}\nsome-other-boot-id:12345\n",
std::process::id(),
now_epoch_secs()
),
)
.unwrap();
let log = EventLog::acquire(&paths, "m-lock", Duration::from_millis(50), LockForce::No)
.expect("a foreign token on our own pid proves the recorder is dead");
drop(log);
}
#[test]
fn lock_holder_is_alive_understands_every_lock_format() {
use kranz_engine::event_log::lock_holder_is_alive;
let tmp = tempfile::tempdir().unwrap();
let lock = tmp.path().join("events.jsonl.lock");
assert!(
!lock_holder_is_alive(&lock),
"missing lock has no live holder"
);
std::fs::write(&lock, "garbage\n").unwrap();
assert!(
lock_holder_is_alive(&lock),
"unparseable lock is conservatively alive"
);
std::fs::write(
&lock,
format!("{}\n{}\n", std::process::id(), now_epoch_secs()),
)
.unwrap();
assert!(lock_holder_is_alive(&lock), "our own pid is alive");
#[cfg(unix)]
{
std::fs::write(
&lock,
format!("{}\n{}\nsome-token\n", i32::MAX, now_epoch_secs()),
)
.unwrap();
assert!(
!lock_holder_is_alive(&lock),
"a multi-line lock with a dead pid must read NOT alive"
);
}
}
#[test]
fn mission_lock_is_live_delegates_to_the_canonical_probe() {
use kranz_engine::mission_catalog::mission_lock_is_live;
let tmp = tempfile::tempdir().unwrap();
let p = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(p.mission_dir()).unwrap();
assert!(!mission_lock_is_live(&p), "missing lock is not live");
std::fs::write(p.lock_file(), "not-a-pid\n").unwrap();
assert!(
mission_lock_is_live(&p),
"unparseable lock is conservatively live"
);
std::fs::write(
p.lock_file(),
format!("{}\n{}\n", std::process::id(), now_epoch_secs()),
)
.unwrap();
assert!(mission_lock_is_live(&p), "a live holder (us) is live");
#[cfg(unix)]
{
std::fs::write(
p.lock_file(),
format!("{}\n{}\ntok\n", i32::MAX, now_epoch_secs()),
)
.unwrap();
assert!(
!mission_lock_is_live(&p),
"dead-holder multi-line lock is not live"
);
}
}
#[cfg(unix)]
#[test]
fn racing_acquires_on_a_dead_lock_admit_exactly_one_winner() {
let tmp = tempfile::tempdir().unwrap();
let paths = MissionPaths::new(tmp.path(), "m-lock");
std::fs::create_dir_all(paths.mission_dir()).unwrap();
std::fs::write(paths.lock_file(), i32::MAX.to_string()).unwrap();
const RACERS: usize = 16;
let barrier = std::sync::Arc::new(std::sync::Barrier::new(RACERS));
let root = tmp.path().to_path_buf();
let handles: Vec<_> = (0..RACERS)
.map(|_| {
let barrier = std::sync::Arc::clone(&barrier);
let paths = MissionPaths::new(&root, "m-lock");
std::thread::spawn(move || {
barrier.wait();
EventLog::acquire(&paths, "m-lock", NEVER, LockForce::No)
})
})
.collect();
let results: Vec<_> = handles.into_iter().map(|h| h.join().unwrap()).collect();
let winners = results.iter().filter(|r| r.is_ok()).count();
assert_eq!(winners, 1, "exactly one racer may steal a dead-holder lock");
for r in &results {
if let Err(e) = r {
assert!(
matches!(e, EngineError::LockHeld(_)),
"losers must see LockHeld (the winner is alive), got: {e:?}"
);
}
}
}
fn seeded_log(dir: &std::path::Path) -> MissionPaths {
let p = paths(dir);
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(EventKind::MissionCreated {
goal: "one".to_string(),
base_branch: "main".to_string(),
mission_branch: format!("kranz/mission-{MISSION}"),
config: kranz_engine::types::MissionConfig::default(),
})
.unwrap();
for text in ["two", "three"] {
log.append(lifecycle(text)).unwrap();
}
drop(log);
p
}
fn append_line(path: &std::path::Path, line: &str) {
let mut f = std::fs::OpenOptions::new().append(true).open(path).unwrap();
writeln!(f, "{line}").unwrap();
}
#[test]
fn the_writer_seals_every_line_and_readers_accept_it() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let raw = std::fs::read_to_string(p.events_file()).unwrap();
for line in raw.lines() {
let value: serde_json::Value = serde_json::from_str(line).unwrap();
assert!(
value.get("h").and_then(|v| v.as_str()).is_some(),
"every written line carries a chain hash: {line}"
);
}
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 3);
assert!(matches!(&events[0].kind, EventKind::MissionCreated { goal, .. } if goal == "one"));
}
#[test]
fn fractional_costs_keep_their_bits_and_integrity_across_reopen() {
let tmp = tempfile::tempdir().unwrap();
kranz_engine::paths::load_or_create_authority_key(tmp.path()).unwrap();
let p = paths(tmp.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
let mut costs = vec![-0.0, f64::MIN_POSITIVE, 1e-100, 1e100, f64::MAX];
for value in [0.3917785_f64, 0.095758_f64] {
for bits in value.to_bits() - 2..=value.to_bits() + 2 {
costs.push(f64::from_bits(bits));
}
}
for &cost in &costs {
log.append(EventKind::WorkerCompleted {
run_id: "r-cost".into(),
result: kranz_engine::types::RunResult::Pass,
tokens: kranz_engine::types::TokenUsage::default(),
cost_usd: Some(cost),
report: None,
})
.unwrap();
}
drop(log);
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), costs.len());
for (event, cost) in events.iter().zip(costs) {
let EventKind::WorkerCompleted {
cost_usd: Some(observed),
..
} = &event.kind
else {
panic!("unexpected event: {event:?}");
};
assert_eq!(observed.to_bits(), cost.to_bits());
}
let mut reopened = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
reopened.append(lifecycle("after reopen")).unwrap();
drop(reopened);
assert_eq!(
EventLog::read_events(&p.events_file()).unwrap().len(),
events.len() + 1
);
let raw = std::fs::read_to_string(p.events_file()).unwrap();
let changed = raw.replacen("\"costUsd\":-0.0", "\"costUsd\":0.5", 1);
assert_ne!(raw, changed);
std::fs::write(p.events_file(), changed).unwrap();
assert!(matches!(
EventLog::read_events(&p.events_file()),
Err(EngineError::LogCorruption(_))
));
}
#[test]
fn legacy_float_seals_still_read_and_accept_versioned_appends() {
let tmp = tempfile::tempdir().unwrap();
let key = kranz_engine::paths::load_or_create_authority_key(tmp.path()).unwrap();
let p = paths(tmp.path());
let original = 4.613131942124616e-9_f64;
let stored = 4.6131319421246164e-9_f64;
let mut previous = String::new();
let mut lines = Vec::new();
for (index, kind) in [
EventKind::WorkerCompleted {
run_id: "r-legacy".into(),
result: kranz_engine::types::RunResult::Pass,
tokens: kranz_engine::types::TokenUsage::default(),
cost_usd: Some(original),
report: None,
},
EventKind::ConfigChanged {
patch: serde_json::json!({"thresholds": [original], "count": 7, "label": "0.095758"}),
},
]
.into_iter()
.enumerate()
{
let event = Event {
seq: index as u64 + 1,
ts: Utc::now(),
mission_id: MISSION.into(),
kind,
};
let body = serde_json::to_string(&event).unwrap();
let hash =
kranz_engine::standards_waiver::sha256_hex(format!("{previous}{body}").as_bytes());
let mut value = serde_json::to_value(&event).unwrap();
if index == 0 {
value["payload"]["costUsd"] = stored.into();
} else {
value["payload"]["patch"]["thresholds"][0] = stored.into();
}
value["h"] = hash.clone().into();
value["m"] = kranz_engine::hooks::hmac_sha256_hex(&key, hash.as_bytes()).into();
lines.push(serde_json::to_string(&value).unwrap());
previous = hash;
}
write_raw_log(&p.events_file(), &lines);
let legacy_bytes = std::fs::read(p.events_file()).unwrap();
for events in [
EventLog::read_events(&p.events_file()).unwrap(),
EventLog::read_tail_events(&p.events_file(), u64::MAX).unwrap(),
] {
assert!(
matches!(events[0].kind, EventKind::WorkerCompleted { cost_usd: Some(cost), .. } if cost.to_bits() == original.to_bits())
);
let EventKind::ConfigChanged { patch } = &events[1].kind else {
panic!("wrong event")
};
assert_eq!(
patch["thresholds"][0].as_f64().unwrap().to_bits(),
original.to_bits()
);
assert_eq!(patch["count"].as_u64(), Some(7));
assert_eq!(patch["label"].as_str(), Some("0.095758"));
}
let mut reopened = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
reopened.append(lifecycle("new writer")).unwrap();
drop(reopened);
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 3);
let bytes = std::fs::read(p.events_file()).unwrap();
assert!(
bytes.starts_with(&legacy_bytes),
"upgrade must not rewrite legacy evidence"
);
let last: serde_json::Value = serde_json::from_slice(
bytes
.split(|b| *b == b'\n')
.rfind(|line| !line.is_empty())
.unwrap(),
)
.unwrap();
assert_eq!(last["v"], 2);
}
#[test]
fn versioned_seals_refuse_numeric_downgrades_and_unknown_versions() {
let tmp = tempfile::tempdir().unwrap();
let p = paths(tmp.path());
let mut log = EventLog::acquire(&p, MISSION, NEVER, LockForce::No).unwrap();
log.append(EventKind::WorkerCompleted {
run_id: "r-version".into(),
result: kranz_engine::types::RunResult::Pass,
tokens: kranz_engine::types::TokenUsage::default(),
cost_usd: Some(4.613131942124616e-9),
report: None,
})
.unwrap();
drop(log);
let raw = std::fs::read_to_string(p.events_file()).unwrap();
let value: serde_json::Value = serde_json::from_str(&raw).unwrap();
assert_eq!(value["v"], 2);
for version in [
None,
Some(serde_json::json!(1)),
Some(serde_json::json!(3)),
Some(serde_json::json!(2.0)),
Some(serde_json::json!("2")),
] {
let mut changed = value.clone();
match version {
None => {
changed.as_object_mut().unwrap().remove("v");
}
Some(version) => {
changed["v"] = version;
}
}
changed["payload"]["costUsd"] = serde_json::json!(4.6131319421246164e-9);
std::fs::write(
p.events_file(),
serde_json::to_string(&changed).unwrap() + "\n",
)
.unwrap();
assert!(matches!(
EventLog::read_events(&p.events_file()),
Err(EngineError::LogCorruption(_))
));
}
let mut unsealed = value;
unsealed.as_object_mut().unwrap().remove("h");
unsealed.as_object_mut().unwrap().remove("m");
std::fs::write(
p.events_file(),
serde_json::to_string(&unsealed).unwrap() + "\n",
)
.unwrap();
assert!(matches!(
EventLog::read_events(&p.events_file()),
Err(EngineError::LogCorruption(_))
));
}
#[test]
fn a_forged_well_formed_append_is_refused() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
append_line(&p.events_file(), &raw_event(4, lifecycle("forged")));
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("integrity chain")),
"a forged append must be corruption, not truth: {err:?}"
);
}
#[test]
fn an_in_place_payload_rewrite_is_refused() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let raw = std::fs::read_to_string(p.events_file()).unwrap();
let rewritten = raw.replace("\"two\"", "\"rewritten\"");
assert_ne!(rewritten, raw, "fixture: the rewrite must actually apply");
std::fs::write(p.events_file(), rewritten).unwrap();
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("integrity chain broken")),
"{err:?}"
);
}
#[test]
fn dropping_the_chain_mid_log_is_refused() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let raw = std::fs::read_to_string(p.events_file()).unwrap();
let mut lines: Vec<String> = raw.lines().map(str::to_string).collect();
let mut value: serde_json::Value = serde_json::from_str(lines.last().unwrap()).unwrap();
let object = value.as_object_mut().unwrap();
object.remove("h");
object.remove("m");
*lines.last_mut().unwrap() = serde_json::to_string(&value).unwrap();
std::fs::write(p.events_file(), lines.join("\n") + "\n").unwrap();
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("integrity chain dropped")),
"{err:?}"
);
}
#[test]
fn a_legacy_unchained_log_still_reads() {
let tmp = tempfile::tempdir().unwrap();
let p = paths(tmp.path());
write_raw_log(
&p.events_file(),
&[
raw_event(1, lifecycle("one")),
raw_event(2, lifecycle("two")),
],
);
assert_eq!(EventLog::read_events(&p.events_file()).unwrap().len(), 2);
}
#[test]
fn a_seq_gap_still_reports_as_a_seq_discontinuity() {
let tmp = tempfile::tempdir().unwrap();
let p = paths(tmp.path());
write_raw_log(
&p.events_file(),
&[
raw_event(1, lifecycle("one")),
raw_event(3, lifecycle("three")),
],
);
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("seq discontinuity")),
"{err:?}"
);
}
#[test]
fn a_recomputed_chain_without_the_mac_is_refused() {
let tmp = tempfile::tempdir().unwrap();
kranz_engine::paths::load_or_create_authority_key(tmp.path()).unwrap();
let p = seeded_log(tmp.path());
let raw = std::fs::read_to_string(p.events_file()).unwrap();
assert!(
raw.lines().all(|l| l.contains("\"m\":")),
"fixture: the writer must MAC every line when the key exists"
);
let mut events = EventLog::read_events(&p.events_file()).unwrap();
let mut forged = events[0].clone();
forged.seq = 4;
forged.kind = lifecycle("forged");
events.push(forged);
std::fs::write(
p.events_file(),
kranz_engine::event_log::seal_events(&events, None).unwrap(),
)
.unwrap();
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("mac")),
"a chain the attacker recomputed must fail the MAC: {err:?}"
);
}
#[test]
fn a_mac_under_the_wrong_key_is_refused() {
let tmp = tempfile::tempdir().unwrap();
kranz_engine::paths::load_or_create_authority_key(tmp.path()).unwrap();
let p = seeded_log(tmp.path());
let events = EventLog::read_events(&p.events_file()).unwrap();
std::fs::write(
p.events_file(),
kranz_engine::event_log::seal_events(&events, Some(b"not the authority key")).unwrap(),
)
.unwrap();
let err = EventLog::read_events(&p.events_file()).unwrap_err();
assert!(
matches!(&err, EngineError::LogCorruption(m) if m.contains("mac does not verify")),
"{err:?}"
);
}
#[test]
fn a_log_shorter_than_the_snapshot_refuses_to_resume() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let events = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(events.len(), 3);
let state = kranz_engine::reducer::fold(&events).unwrap();
assert_eq!(state.last_seq, 3);
kranz_engine::reducer::write_snapshot(&state, &p.state_file()).unwrap();
let raw = std::fs::read_to_string(p.events_file()).unwrap();
let kept: Vec<&str> = raw.lines().take(1).collect();
std::fs::write(p.events_file(), kept.join("\n") + "\n").unwrap();
let short = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(short.len(), 1, "the truncated log still parses cleanly");
let err = kranz_engine::event_log::check_no_rollback(&p, &short).unwrap_err();
let message = err.to_string();
assert!(message.contains("seq 1"), "names the log's end: {message}");
assert!(
message.contains("seq 3"),
"names the snapshot's mark: {message}"
);
}
#[test]
fn a_snapshot_behind_the_log_is_not_a_rollback() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let events = EventLog::read_events(&p.events_file()).unwrap();
let stale = kranz_engine::reducer::fold(&events[..1]).unwrap();
assert_eq!(stale.last_seq, 1);
kranz_engine::reducer::write_snapshot(&stale, &p.state_file()).unwrap();
kranz_engine::event_log::check_no_rollback(&p, &events).unwrap();
}
#[test]
fn a_missing_snapshot_is_not_a_rollback() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let events = EventLog::read_events(&p.events_file()).unwrap();
kranz_engine::event_log::check_no_rollback(&p, &events).unwrap();
}
#[test]
fn truncation_is_refused_even_when_the_snapshot_is_gone() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let recorded = kranz_engine::paths::read_high_water(&p.repo_root, MISSION)
.expect("a sealed mission records its high-water mark");
assert_eq!(recorded, 3);
let _ = std::fs::remove_file(p.state_file());
let full = std::fs::read_to_string(p.events_file()).unwrap();
let first_line = full.lines().next().unwrap();
std::fs::write(p.events_file(), format!("{first_line}\n")).unwrap();
let short = EventLog::read_events(&p.events_file()).unwrap();
assert_eq!(short.len(), 1, "a clean prefix still parses on its own");
let err = kranz_engine::event_log::check_no_rollback(&p, &short).unwrap_err();
let message = err.to_string();
assert!(message.contains("high-water mark"), "{message}");
assert!(
message.contains("seq 1") && message.contains("seq 3"),
"{message}"
);
}
#[test]
fn high_water_mark_never_lowers() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
kranz_engine::paths::record_high_water(&p.repo_root, MISSION, 1).unwrap();
assert_eq!(
kranz_engine::paths::read_high_water(&p.repo_root, MISSION),
Some(3)
);
}
#[test]
fn acquire_refuses_when_the_authority_key_is_unreadable() {
let tmp = tempfile::tempdir().unwrap();
let p = seeded_log(tmp.path());
let key_path = kranz_engine::paths::authority_key_path(&p.repo_root).unwrap();
std::fs::write(&key_path, b"").unwrap();
let err = EventLog::acquire(&p, MISSION, NEVER, LockForce::No)
.expect_err("acquire must refuse without a readable key");
let message = err.to_string();
assert!(message.contains("authority key"), "{message}");
}