#![cfg(all(unix, feature = "rt"))]
#[path = "support/journal.rs"]
mod shared;
use shared::{ProbeGuard, TempGuard, key_for as key, pause, scratch_dir};
use std::path::{Path, PathBuf};
use std::time::Duration;
use lgwks_bot::effect::EffectKey;
use lgwks_bot::journal::{EffectEvent, EffectJournal, FileJournal};
use lgwks_bot::rt::clock::Clock;
const DIGEST: &str = "f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef";
const BUDGET: Duration = Duration::from_secs(600);
const LEFT: Duration = Duration::from_secs(600 - 250);
const HALF: Duration = Duration::from_secs(300);
const HALF_LEFT: Duration = Duration::from_secs(300 - 250);
const SPENT: Duration = Duration::from_secs(250);
type TestResult = Result<(), Box<dyn std::error::Error>>;
const PROBE_ENV: &str = "LGWKS_CLOCK_KILL_PROBE";
const PROBE_JOURNAL: &str = "LGWKS_CLOCK_KILL_JOURNAL";
const PROBE_MARKER: &str = "LGWKS_CLOCK_KILL_MARKER";
const PROBE_ROW: &str = "LGWKS_CLOCK_KILL_ROW";
fn scratch_and_guard(name: &str) -> Result<(PathBuf, TempGuard), Box<dyn std::error::Error>> {
let dir = scratch_dir(name)?;
let guard = TempGuard(dir.clone());
Ok((dir, guard))
}
fn step_key() -> Result<EffectKey, Box<dyn std::error::Error>> {
key("1", DIGEST)
}
fn admit(journal: &mut FileJournal, key: EffectKey) -> Result<(), Box<dyn std::error::Error>> {
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
Ok(())
}
fn child_body() -> TestResult {
let journal_path = journal_path()?;
let marker = marker_path()?;
let snapshot_path = journal_path.with_extension("elapsed");
let mut journal = FileJournal::open(&journal_path)?;
let key = step_key()?;
admit(&mut journal, key)?;
let clock = Clock::virtual_at(Duration::ZERO);
clock.advance(SPENT)?;
let snapshot = clock.snapshot();
std::fs::write(&snapshot_path, nanos_of(snapshot.elapsed())?.to_le_bytes())?;
let file = std::fs::File::open(&snapshot_path)?;
file.sync_all()?;
std::fs::write(&marker, b"acked")?;
for _ in 0..600 {
pause(100);
}
Err("the probe child parked for its whole bound and was never killed".into())
}
fn spawn_probe(
test_name: &str,
row: usize,
journal_path: &Path,
marker_path: &Path,
) -> Result<ProbeGuard, Box<dyn std::error::Error>> {
Ok(ProbeGuard(Some(
std::process::Command::new(std::env::current_exe()?)
.args([test_name, "--exact", "--nocapture"])
.env(PROBE_ENV, "1")
.env(PROBE_ROW, row.to_string())
.env(PROBE_JOURNAL, journal_path)
.env(PROBE_MARKER, marker_path)
.spawn()?,
)))
}
fn journal_path() -> Result<PathBuf, Box<dyn std::error::Error>> {
let path = std::env::var_os(PROBE_JOURNAL)
.ok_or("the probe child was started without a journal path")?;
Ok(PathBuf::from(path))
}
fn marker_path() -> Result<PathBuf, Box<dyn std::error::Error>> {
let path = std::env::var_os(PROBE_MARKER)
.ok_or("the probe child was started without a marker path")?;
Ok(PathBuf::from(path))
}
fn kill_after_marker(
mut guard: ProbeGuard,
marker: &Path,
test_name: &str,
) -> Result<(), Box<dyn std::error::Error>> {
guard.kill_after_marker(marker, test_name)
}
fn read_child_elapsed(journal_path: &Path) -> Result<Duration, Box<dyn std::error::Error>> {
let path = journal_path.with_extension("elapsed");
let bytes = std::fs::read(&path)?;
if bytes.len() != std::mem::size_of::<u64>() {
{
let refusal =
Err(format!("the child's elapsed record is {} bytes, not 8", bytes.len()).into());
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read_child_elapsed: returning an error to the caller");
return refusal;
};
}
let mut raw = [0_u8; 8];
raw.copy_from_slice(&bytes);
Ok(Duration::from_nanos(u64::from_le_bytes(raw)))
}
fn nanos_of(elapsed: Duration) -> Result<u64, Box<dyn std::error::Error>> {
u64::try_from(elapsed.as_nanos())
.map_err(|error| format!("{elapsed:?} is longer than the record can hold: {error}").into())
}
struct Fixture {
journal: PathBuf,
marker: PathBuf,
_dir: TempGuard,
}
impl Fixture {
fn kill_after(&self, test_name: &str, row: usize) -> Result<(), Box<dyn std::error::Error>> {
let child = spawn_probe(test_name, row, &self.journal, &self.marker)?;
kill_after_marker(child, &self.marker, test_name)
}
fn reopened(&self) -> Result<FileJournal, Box<dyn std::error::Error>> {
Ok(FileJournal::open(&self.journal)?)
}
fn restarted_journal(&self) -> Result<FileJournal, Box<dyn std::error::Error>> {
let mut journal = self.reopened()?;
admit(&mut journal, step_key()?)?;
Ok(journal)
}
fn child_elapsed(&self) -> Result<Duration, Box<dyn std::error::Error>> {
read_child_elapsed(&self.journal)
}
fn recovered_budget(&self) -> Result<Duration, Box<dyn std::error::Error>> {
Ok(clock_from_elapsed(self.child_elapsed()?).remaining_from(BUDGET))
}
fn new(name: &str) -> Result<Self, Box<dyn std::error::Error>> {
let (dir, guard) = scratch_and_guard(name)?;
Ok(Self {
journal: dir.join("journal.log"),
marker: dir.join("marker"),
_dir: guard,
})
}
}
fn agree<T: std::fmt::Debug + PartialEq>(observed: T, expected: T, claim: &str) {
assert_eq!(observed, expected, "{claim}");
}
fn refuse<T, E>(result: Result<T, E>, why: &str) -> Result<E, Box<dyn std::error::Error>> {
match result {
Ok(_) => Err(why.into()),
Err(refusal) => Ok(refusal),
}
}
mod claim {
pub const BUDGET_LEFT: &str = "the restart must find the budget the child had left";
pub const ONE_ATTEMPT: &str = "the refused replay added no second attempt";
pub const RECORD_LENGTH: &str = "the refusal names the length it refused";
}
#[derive(Clone, Copy)]
enum Row {
Budget,
Duration,
Idempotent,
NoRerun,
}
impl Row {
const ALL: [Row; 4] = [Row::Budget, Row::Duration, Row::Idempotent, Row::NoRerun];
const fn scratch(self) -> &'static str {
match self {
Self::Budget => "budget",
Self::Duration => "duration",
Self::Idempotent => "idempotent",
Self::NoRerun => "rerun",
}
}
fn check(self, fixture: &Fixture) -> Result<(), Box<dyn std::error::Error>> {
match self {
Self::Budget => budget_survives(fixture),
Self::Duration => budget_is_a_duration(fixture),
Self::Idempotent => replay_is_idempotent(fixture),
Self::NoRerun => step_did_not_rerun(fixture),
}
}
}
fn budget_survives(row: &Fixture) -> Result<(), Box<dyn std::error::Error>> {
let empty = FileJournal::open(&row.journal)?.recover().is_empty();
assert!(
!empty,
"the acknowledged intent must be on the disk after a SIGKILL"
);
agree(row.recovered_budget()?, LEFT, claim::BUDGET_LEFT);
assert!(
!clock_from_elapsed(row.child_elapsed()?).is_exhausted(BUDGET),
"a budget with time left must not read as exhausted after the kill",
);
Ok(())
}
fn budget_is_a_duration(row: &Fixture) -> Result<(), Box<dyn std::error::Error>> {
let spent = row.child_elapsed()?;
let recovered = clock_from_elapsed(spent);
let observed = (
spent,
recovered.remaining_from(BUDGET),
recovered.remaining_from(HALF),
);
let expected = (SPENT, LEFT, HALF_LEFT);
agree(
observed,
expected,
"the restart reads back the duration the child spent, and the same duration \
recovers against a different budget too: the value is carried by the record, \
not by the budget it was spent under",
);
Ok(())
}
fn replay_is_idempotent(row: &Fixture) -> Result<(), Box<dyn std::error::Error>> {
let refusal = refuse(
row.restarted_journal(),
"a journal already holding attempt 1 must refuse to open for a replay",
)?;
let reason = refusal.to_string();
for expected in ["cannot follow what is committed", "attempt 1"] {
assert!(
reason.contains(expected),
"the refusal said {reason:?}, not {expected:?}"
);
}
let attempts = FileJournal::open(&row.journal)?.recover().len();
agree(attempts, 1, claim::ONE_ATTEMPT);
Ok(())
}
fn step_did_not_rerun(row: &Fixture) -> Result<(), Box<dyn std::error::Error>> {
let journal = row.reopened()?;
let attempts = journal.recover().len();
agree(
attempts,
1,
"the restart must find exactly the attempt the killed child admitted: a second \
copy is a rerun, and a rerun performs the effect twice",
);
let resumed = Clock::virtual_at(row.child_elapsed()?);
let remaining = resumed.snapshot().remaining_from(BUDGET);
agree(
remaining,
LEFT,
"resuming from the recovered reading must spend nothing: a budget that came back \
shorter was charged twice, and one that came back longer was never charged at all",
);
Ok(())
}
fn row_from_env() -> Option<usize> {
std::env::var(PROBE_ROW).ok()?.parse().ok()
}
fn clock_from_elapsed(elapsed: Duration) -> lgwks_bot::rt::clock::ClockSnapshot {
let clock = Clock::virtual_at(elapsed);
clock.snapshot()
}
#[test]
fn the_remaining_budget_survives_a_real_sigkill() -> TestResult {
match row_from_env() {
Some(_) => child_body(),
None => run_every_row(),
}
}
#[test]
fn a_torn_snapshot_after_the_kill_is_refused_not_read_as_a_shorter_step() -> TestResult {
let fixture = Fixture::new("torn")?;
let mut journal = fixture.reopened()?;
admit(&mut journal, step_key()?)?;
drop(journal);
std::fs::write(fixture.journal.with_extension("elapsed"), [1_u8, 2, 3, 4])?;
std::fs::write(&fixture.marker, b"acked")?;
let torn = "a half-written snapshot must be refused; reading it as a shorter step is how \
a restart silently loses budget";
let refusal = refuse(fixture.child_elapsed(), torn)?;
agree(
refusal.to_string(),
String::from("the child's elapsed record is 4 bytes, not 8"),
claim::RECORD_LENGTH,
);
Ok(())
}
fn run_every_row() -> TestResult {
let this = std::thread::current();
let name = this.name().ok_or("the test has no name")?;
for (index, row) in Row::ALL.iter().enumerate() {
let fixture = Fixture::new(row.scratch())?;
fixture.kill_after(name, index)?;
row.check(&fixture)?;
}
Ok(())
}