use std::sync::Arc;
use std::time::Duration;
use crate::core::{Digest, EffectKey, RunId, StepId};
use crate::journal::{Append, JournalStore, Record, RecordKind};
#[derive(Debug, Clone)]
pub struct Violation {
pub invariant: &'static str,
pub detail: String,
}
#[derive(Debug, Default)]
pub struct Report {
pub violations: Vec<Violation>,
pub checked: usize,
}
impl Report {
pub(crate) fn record(&mut self, invariant: &'static str, detail: impl Into<String>) {
self.violations.push(Violation {
invariant,
detail: detail.into(),
});
}
pub fn assert_conforms(&self, store: &str) {
assert!(
self.checked > 0,
"the conformance battery ran no checks against {store}"
);
assert!(
self.violations.is_empty(),
"{store} does not satisfy its store contract:\n{}",
self.violations
.iter()
.map(|v| format!(" • {}: {}", v.invariant, v.detail))
.collect::<Vec<_>>()
.join("\n")
);
}
}
pub type Factory<'a> =
&'a (dyn Fn() -> std::pin::Pin<Box<dyn Future<Output = Arc<dyn JournalStore>> + Send>> + Sync);
const LEASE: Duration = Duration::from_mins(5);
const SHORT: Duration = Duration::from_secs(1);
fn admitted(run: RunId) -> Append {
Append::new(
run,
RecordKind::RunAdmitted {
agent: "conformance".into(),
input: serde_json::Value::Null,
policy: None,
},
)
}
fn started(run: RunId, key: EffectKey) -> Append {
Append::new(
run,
RecordKind::EffectStarted {
descriptor: crate::core::EffectDescriptor::nullary("probe"),
recovery: crate::core::Recovery::Retry,
mutates: true,
attempt: 1,
backoff_ms: 0,
},
)
.step(StepId(0))
.effect(key)
}
fn key(n: u8) -> EffectKey {
EffectKey::derive(
StepId(0),
crate::core::Phase::Forward,
u32::from(n),
1,
"probe",
&[n],
)
}
pub async fn check(fresh: Factory<'_>) -> Report {
let mut r = Report::default();
head_of_an_unwritten_run_is_genesis(fresh, &mut r).await;
append_assigns_contiguous_seq(fresh, &mut r).await;
append_links_each_record_to_the_previous(fresh, &mut r).await;
a_sound_chain_verifies(fresh, &mut r).await;
a_second_start_for_one_effect_is_rejected(fresh, &mut r).await;
the_same_effect_key_in_another_run_is_allowed(fresh, &mut r).await;
a_stale_epoch_is_fenced(fresh, &mut r).await;
a_live_lease_is_not_stolen(fresh, &mut r).await;
a_takeover_advances_the_epoch(fresh, &mut r).await;
a_rejected_batch_writes_nothing(fresh, &mut r).await;
read_starts_where_it_is_told(fresh, &mut r).await;
a_stop_request_needs_no_lease(fresh, &mut r).await;
the_first_asker_stays_on_the_record(fresh, &mut r).await;
an_attestation_survives_the_round_trip(fresh, &mut r).await;
the_log_only_grows(fresh, &mut r).await;
r
}
async fn the_log_only_grows(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let mut sealed = Vec::new();
let mut history: Vec<(u64, Digest)> = Vec::new();
for _ in 0..4u8 {
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
if store
.append(lease.epoch, vec![admitted(run)])
.await
.is_err()
{
r.record("merkle-log", "append failed under a fresh lease");
return;
}
if store.seal(run, lease.epoch, "succeeded").await.is_err() {
r.record("merkle-log", "seal failed");
return;
}
sealed.push(run);
let Ok(cp) = store.checkpoint().await else {
return;
};
if cp.size != sealed.len() as u64 {
r.record(
"merkle-log",
format!(
"{} runs are sealed and the checkpoint commits to {} — a run \
that sealed without entering the log is a run no checkpoint \
covers",
sealed.len(),
cp.size
),
);
return;
}
history.push((cp.size, cp.root));
}
let Ok(latest) = store.checkpoint().await else {
return;
};
for (i, run) in sealed.iter().enumerate() {
match store.inclusion_proof(*run).await {
Ok(Some(inc)) => {
let leaf = crate::core::merkle::leaf_hash(&inc.seal);
if !crate::core::merkle::verify_inclusion(
&leaf,
usize::try_from(inc.index).unwrap_or(usize::MAX),
usize::try_from(inc.size).unwrap_or(0),
&inc.proof,
&latest.root,
) {
r.record(
"merkle-log",
format!("run {i} is in the log and cannot prove it"),
);
}
}
Ok(None) => r.record(
"merkle-log",
format!("run {i} sealed and the store has no position for it"),
),
Err(e) => r.record("merkle-log", format!("inclusion_proof failed: {e}")),
}
}
for (size, root) in &history {
let proof = match store.consistency_proof(*size).await {
Ok(p) => p,
Err(e) => {
r.record(
"merkle-log",
format!("consistency_proof({size}) failed: {e}"),
);
continue;
}
};
if !crate::core::merkle::verify_consistency(
usize::try_from(*size).unwrap_or(0),
root,
usize::try_from(latest.size).unwrap_or(0),
&latest.root,
&proof,
) {
r.record(
"merkle-log",
format!(
"the log could not prove it only appended between size {size} \
and {} — so a checkpoint published then proves nothing now",
latest.size
),
);
}
}
}
async fn a_stop_request_needs_no_lease(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "the-owner", LEASE).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
if store
.append(lease.epoch, vec![admitted(run)])
.await
.is_err()
{
r.record("chaining", "append failed under a fresh lease");
return;
}
match store.request_cancel(run, "operator", "stop it").await {
Ok(true) => {}
Ok(false) => r.record(
"cancellation",
"the first stop request on a run reported that one already existed",
),
Err(e) => r.record(
"cancellation",
format!(
"a stop request from somebody who does not hold the lease was refused ({e}) — \
then the only party who can stop a running agent is the process running it"
),
),
}
match store.cancellation(run).await {
Ok(Some(c)) if c.actor == "operator" && c.reason == "stop it" => {}
Ok(Some(c)) => r.record(
"cancellation",
format!("the stop request read back as {c:?}, which is not what was written"),
),
Ok(None) => r.record(
"cancellation",
"a recorded stop request was not readable — an operator's intervention \
that the owner cannot see is not an intervention",
),
Err(e) => r.record("cancellation", format!("cancellation() failed: {e}")),
}
match store.cancellation(RunId::generate()).await {
Ok(None) => {}
Ok(Some(_)) => r.record(
"cancellation",
"an untouched run reported a stop request — every run would unwind",
),
Err(e) => r.record("cancellation", format!("cancellation() failed: {e}")),
}
}
async fn the_first_asker_stays_on_the_record(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
if !matches!(store.request_cancel(run, "alice", "first").await, Ok(true)) {
r.record("cancellation", "the first stop request was not recorded");
return;
}
match store.request_cancel(run, "bob", "second").await {
Ok(false) => {}
Ok(true) => r.record(
"cancellation",
"a second stop request reported itself as the first — a retry would \
read as a fresh intervention",
),
Err(e) => r.record(
"cancellation",
format!("a repeated stop request failed: {e}"),
),
}
match store.cancellation(run).await {
Ok(Some(c)) if c.actor == "alice" && c.reason == "first" => {}
Ok(other) => r.record(
"cancellation",
format!(
"the stop request now reads {other:?} — the second asker overwrote the \
first, so the permanent record names the wrong person"
),
),
Err(e) => r.record("cancellation", format!("cancellation() failed: {e}")),
}
}
async fn an_attestation_survives_the_round_trip(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
let Ok(written) = store.append(lease.epoch, vec![admitted(run)]).await else {
r.record("attestation", "append failed under a fresh lease");
return;
};
let Some(expected) = written.first().and_then(|w| w.attestation.clone()) else {
return;
};
match store.read(run, 1).await {
Ok(back) => match back.first().and_then(|b| b.attestation.clone()) {
Some(got) if got == expected => {}
Some(got) => r.record(
"attestation",
format!(
"the signature read back is not the one written: wrote {:?}, read {:?}",
expected.key_id, got.key_id
),
),
None => r.record(
"attestation",
"a signed record read back unsigned — the chain still verifies, and \
authorship is gone without anything reporting it",
),
},
Err(e) => r.record("attestation", format!("read failed: {e}")),
}
}
async fn head_of_an_unwritten_run_is_genesis(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
match store.head(RunId::generate()).await {
Ok(head) if head.seq == 0 && head.hash == Digest::ZERO => {}
Ok(head) => r.record(
"chaining",
format!(
"a run with no records must start at genesis (seq 0, zero hash), got seq {} — \
a chain that starts anywhere else cannot be verified from the beginning",
head.seq
),
),
Err(e) => r.record(
"chaining",
format!("head() of an unwritten run must succeed at genesis, got {e}"),
),
}
}
async fn append_assigns_contiguous_seq(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
for i in 0..3u8 {
if let Err(e) = store.append(lease.epoch, vec![admitted(run)]).await {
r.record("chaining", format!("append {i} failed: {e}"));
return;
}
}
match store.read(run, 1).await {
Ok(records) => {
let seqs: Vec<u64> = records.iter().map(Record::seq).collect();
if seqs != vec![1, 2, 3] {
r.record(
"chaining",
format!(
"seq must be contiguous from 1, got {seqs:?} — a gap means a record was \
lost and verification cannot tell that from tampering"
),
);
}
}
Err(e) => r.record("chaining", format!("read failed: {e}")),
}
}
async fn append_links_each_record_to_the_previous(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
for _ in 0..3 {
let _ = store.append(lease.epoch, vec![admitted(run)]).await;
}
let Ok(records) = store.read(run, 1).await else {
return;
};
let mut prev = Digest::ZERO;
for rec in &records {
if rec.prev_hash != prev {
r.record(
"chaining",
format!(
"record {} does not link to its predecessor — the chain is what makes \
deletion and reordering detectable, and an unlinked record is neither",
rec.seq()
),
);
return;
}
prev = rec.hash;
}
}
async fn a_sound_chain_verifies(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
for _ in 0..4 {
let _ = store.append(lease.epoch, vec![admitted(run)]).await;
}
if let Err(e) = store.verify(run).await {
r.record(
"chaining",
format!("a chain this store wrote itself must verify, got {e}"),
);
}
}
async fn a_second_start_for_one_effect_is_rejected(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
let k = key(1);
if let Err(e) = store.append(lease.epoch, vec![started(run, k)]).await {
r.record("exactly-once", format!("the first start must succeed: {e}"));
return;
}
match store.append(lease.epoch, vec![started(run, k)]).await {
Err(crate::core::StoreError::DuplicateEffect(_)) => {}
Err(e) => r.record(
"exactly-once",
format!("a second start must be rejected as DuplicateEffect, got {e}"),
),
Ok(_) => r.record(
"exactly-once",
"a second EffectStarted for one effect key was accepted. Exactly-once is a \
storage constraint here, not a code path — replay reads a completed effect \
back, and this is what stops anything that gets past it from performing the \
call twice",
),
}
}
async fn the_same_effect_key_in_another_run_is_allowed(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let k = key(2);
for _ in 0..2 {
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
if let Err(e) = store
.append(lease.epoch, vec![started(run, k)].clone())
.await
{
r.record(
"exactly-once",
format!(
"the constraint is per-run: two runs performing the same effect are two \
performances, and rejecting the second would make a batch of identical \
items unrunnable. Got {e}"
),
);
return;
}
}
}
async fn a_stale_epoch_is_fenced(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(first) = store.acquire(run, "instance-a", SHORT).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
tokio::time::sleep(Duration::from_millis(1_200)).await;
let Ok(second) = store.acquire(run, "instance-b", LEASE).await else {
r.record(
"fencing",
"an expired lease must be claimable by another instance, or a crashed \
owner strands its run forever",
);
return;
};
if second.epoch <= first.epoch {
r.record(
"fencing",
"a takeover must advance the epoch, or the previous owner cannot be fenced",
);
return;
}
match store.append(first.epoch, vec![admitted(run)]).await {
Err(crate::core::StoreError::Fenced { .. }) => {}
Err(e) => r.record("fencing", format!("expected Fenced, got {e}")),
Ok(_) => r.record(
"fencing",
"a superseded owner's write was accepted. The check must happen inside the same \
transaction that writes: a read-then-write leaves a window a paused instance \
wakes up into, which is split-brain",
),
}
}
async fn a_live_lease_is_not_stolen(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(_held) = store.acquire(run, "instance-a", LEASE).await else {
return;
};
match store.acquire(run, "instance-b", LEASE).await {
Err(crate::core::StoreError::LeaseHeld { .. }) => {}
Err(e) => r.record(
"fencing",
format!(
"taking a live lease must fail as LeaseHeld, not {e} — a writer that is \
merely early must not be told it has been superseded"
),
),
Ok(_) => r.record(
"fencing",
"a live lease was handed to a second instance. Two owners writing one chain is \
the split-brain the epoch exists to prevent",
),
}
}
async fn a_takeover_advances_the_epoch(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(a) = store.acquire(run, "instance-a", LEASE).await else {
return;
};
let Ok(again) = store.acquire(run, "instance-a", LEASE).await else {
return;
};
if again.epoch != a.epoch {
r.record(
"fencing",
format!(
"renewing an owned lease must keep the epoch ({} became {}) — bumping it \
fences the owner against itself, and its in-flight writes start failing",
a.epoch, again.epoch
),
);
}
}
async fn a_rejected_batch_writes_nothing(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
let k = key(3);
let _ = store.append(lease.epoch, vec![started(run, k)]).await;
let before = store.read(run, 1).await.map_or(0, |v| v.len());
let batch = vec![admitted(run), started(run, k)];
if store.append(lease.epoch, batch).await.is_ok() {
r.record(
"atomicity",
"a batch containing a duplicate effect start was accepted",
);
return;
}
let after = store.read(run, 1).await.map_or(0, |v| v.len());
if after != before {
r.record(
"atomicity",
format!(
"a rejected batch left {} record(s) behind. The whole batch commits or none \
of it does — a partially written step describes something that never \
happened",
after - before
),
);
}
}
async fn read_starts_where_it_is_told(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
for _ in 0..4 {
let _ = store.append(lease.epoch, vec![admitted(run)]).await;
}
match store.read(run, 3).await {
Ok(records) => {
let seqs: Vec<u64> = records.iter().map(Record::seq).collect();
if seqs != vec![3, 4] {
r.record(
"chaining",
format!("read(from=3) must return seq 3 onward, got {seqs:?}"),
);
}
}
Err(e) => r.record("chaining", format!("read(from=3) failed: {e}")),
}
}