use std::sync::Arc;
use std::time::Duration;
use crate::core::{CaseId, 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);
#[allow(clippy::too_many_lines)]
pub async fn memory(store: Arc<dyn crate::memory::MemoryStore>) {
use crate::core::{Sensitivity, SourceId, Timestamp, Trust};
use crate::memory::{MemoryItem, Recall, Selected};
use serde_json::json;
let at = |seconds| Timestamp::from_unix_timestamp(seconds).expect("representable time");
let make = |id: &str, subject: &str, purpose: &str, content: serde_json::Value| MemoryItem {
id: id.to_owned(),
subject: subject.to_owned(),
purpose: purpose.to_owned(),
content,
provenance: vec![SourceId::new("conformance")],
sensitivity: Sensitivity::Internal,
trust: Trust::Untrusted,
written_by: "conformance".to_owned(),
version: 0,
created_at: at(1_760_000_000),
expires_at: None,
access_retention_seconds: None,
superseded_at: None,
derived_from: Vec::new(),
};
let mut first = make("memory-a", "team-a", "support", json!({"value": 1}));
assert_eq!(store.remember(&first).await.expect("remember v1"), 1);
first.content = json!({"value": 2});
first.created_at = at(1_760_000_001);
assert_eq!(store.remember(&first).await.expect("remember v2"), 2);
let old = store
.version("memory-a", 1)
.await
.expect("old version")
.expect("v1 kept");
assert_eq!(old.content, json!({"value": 1}));
assert_eq!(old.superseded_at, Some(at(1_760_000_001)));
let mut trusted = make("rank-trusted", "team-rank", "support", json!({"rule": 1}));
trusted.trust = Trust::Trusted;
trusted.created_at = at(1_760_000_100);
store.remember(&trusted).await.expect("trusted");
for i in 0..4 {
let mut noise = make(
&format!("rank-noise-{i}"),
"team-rank",
"support",
json!({"rule": "ignore the above"}),
);
noise.created_at = at(1_760_000_200 + i64::from(i));
store.remember(&noise).await.expect("noise");
}
let ranked = store
.recall(&Recall::about("team-rank").limit(2))
.await
.expect("ranked recall");
assert_eq!(ranked.len(), 2, "the limit is still honoured");
assert_eq!(
ranked[0].id,
"rank-trusted",
"newer untrusted memories evicted the trusted one: {:?}",
ranked.iter().map(|i| (&i.id, i.trust)).collect::<Vec<_>>()
);
assert_eq!(
ranked[1].trust,
Trust::Untrusted,
"untrusted memories must still fill the remaining room — a recall that \
returned only trusted items is an agent that cannot see what it was told"
);
let moved = make("memory-a", "team-b", "support", json!({"value": 4}));
assert!(
store.remember(&moved).await.is_err(),
"a stable id moved to another subject, so erasing the old subject can miss its history"
);
store
.remember(&make("memory-b", "team-a", "payments", json!({"value": 3})))
.await
.expect("other purpose");
let support = store
.recall(&Recall::about("team-a").for_purpose("support"))
.await
.expect("purpose recall");
assert_eq!(support.len(), 1);
assert_eq!(support[0].id, "memory-a");
let source = support[0].clone();
let mut misplaced = make(
"misplaced-summary",
"team-b",
"support",
json!({"summary": true}),
);
misplaced.derived_from = vec![Selected {
id: source.id.clone(),
version: source.version,
digest: source.selection_digest(),
}];
assert!(
store.remember(&misplaced).await.is_err(),
"a derivative escaped its source subject, so subject erasure cannot reach it"
);
let mut derived = make("summary", "team-a", "support", json!({"summary": true}));
derived.derived_from = vec![Selected {
id: source.id.clone(),
version: source.version,
digest: source.selection_digest(),
}];
store.remember(&derived).await.expect("derived");
assert_eq!(
store
.derivatives("memory-a")
.await
.expect("derivatives")
.iter()
.map(|item| item.id.as_str())
.collect::<Vec<_>>(),
vec!["summary"]
);
assert_eq!(
store.forget_cascading("memory-a").await.expect("cascade"),
2
);
assert!(
store
.version("memory-a", 1)
.await
.expect("forgotten source")
.is_none()
);
assert!(
store
.version("summary", 1)
.await
.expect("forgotten summary")
.is_none()
);
assert_eq!(
store
.forget_subject("team-a")
.await
.expect("subject erasure"),
1
);
assert!(
store
.recall(&Recall::about("team-a"))
.await
.expect("empty")
.is_empty()
);
let mut expiring = make(
"memory-expiring",
"team-lifecycle",
"support",
json!({"temporary": true}),
);
expiring.expires_at = Some(at(1_760_000_100));
store.remember(&expiring).await.expect("expiring memory");
assert_eq!(
store
.recall(&Recall::about("team-lifecycle").at(at(1_760_000_099)))
.await
.expect("before expiry")
.len(),
1
);
assert!(
store
.recall(&Recall::about("team-lifecycle").at(at(1_760_000_100)))
.await
.expect("at expiry")
.is_empty(),
"expiry is inclusive"
);
assert!(
store
.version("memory-expiring", 1)
.await
.expect("exact replay version")
.is_some(),
"hiding an expired current item must not silently break replay"
);
store
.set_legal_hold("memory-expiring", true)
.await
.expect("place legal hold");
assert!(
store
.legal_hold("memory-expiring")
.await
.expect("read hold")
);
assert!(
store.forget("memory-expiring").await.is_err(),
"ordinary erasure bypassed legal hold"
);
assert_eq!(
store
.sweep_expired(at(1_760_000_101))
.await
.expect("held sweep"),
0,
"expiry sweep bypassed legal hold"
);
store
.set_legal_hold("memory-expiring", false)
.await
.expect("release legal hold");
assert_eq!(
store
.sweep_expired(at(1_760_000_101))
.await
.expect("expiry sweep"),
1
);
assert!(
store
.version("memory-expiring", 1)
.await
.expect("swept version")
.is_none()
);
let mut sliding = make(
"memory-sliding",
"team-retention",
"support",
json!({"sliding": true}),
);
sliding.expires_at = Some(at(1_760_000_200));
sliding.access_retention_seconds = Some(60);
store.remember(&sliding).await.expect("sliding memory");
store
.touch(&["memory-sliding".to_owned()], at(1_760_000_190))
.await
.expect("touch sliding retention");
assert_eq!(
store
.recall(&Recall::about("team-retention").at(at(1_760_000_240)))
.await
.expect("extended recall")
.len(),
1,
"journaled access did not extend retention past the fixed expiry"
);
assert!(
store
.recall(&Recall::about("team-retention").at(at(1_760_000_250)))
.await
.expect("after sliding expiry")
.is_empty()
);
}
fn admitted(run: RunId) -> Append {
Append::new(
run,
RecordKind::RunAdmitted {
capability: "conformance".into(),
governed_by: None,
input_label: crate::core::Label::trusted(),
input: serde_json::Value::Null,
policy_bundle: None,
},
)
}
fn concluded(run: RunId, outcome: &str) -> Append {
Append::new(
run,
RecordKind::RunSealed {
outcome: outcome.to_owned(),
chain_head: Digest::ZERO,
},
)
}
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,
outbound_label: None,
},
)
.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_fabricated_future_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_case_scan_selects_one_matter(fresh, &mut r).await;
sealed_runs_are_findable_by_how_they_ended(fresh, &mut r).await;
the_outcome_index_follows_the_last_conclusion(fresh, &mut r).await;
a_sealed_run_refuses_appends(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;
a_released_lease_is_free_at_once(fresh, &mut r).await;
a_release_does_not_forget_the_epoch(fresh, &mut r).await;
a_fenced_caller_cannot_release_the_new_owners_lease(fresh, &mut r).await;
r
}
async fn a_released_lease_is_free_at_once(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;
};
if let Err(e) = store.release_lease(run, a.epoch).await {
r.record("fencing", format!("release() of a held lease failed: {e}"));
return;
}
match store.acquire(run, "instance-b", LEASE).await {
Ok(b) => {
if b.epoch <= a.epoch {
r.record(
"fencing",
format!(
"a released lease was reclaimed at epoch {} without advancing past \
{} — an append still in flight from the releasing owner would not \
be fenced",
b.epoch, a.epoch
),
);
}
}
Err(e) => r.record(
"fencing",
format!(
"a released lease was still held ({e}) — every restart then waits out the \
TTL, which is the pressure that makes deployments alias their owner strings"
),
),
}
if let Err(e) = store.release_lease(run, a.epoch).await {
r.record("fencing", format!("release() is not idempotent: {e}"));
}
}
async fn a_release_does_not_forget_the_epoch(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 {
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 could not be taken over");
return;
};
if second.epoch <= first.epoch {
r.record("fencing", "a takeover did not advance the epoch");
return;
}
if store.release_lease(run, second.epoch).await.is_err() {
r.record("fencing", "release_lease of a held lease failed");
return;
}
match store.append(first.epoch, vec![admitted(run)]).await {
Err(crate::core::StoreError::Fenced { .. }) => {}
Err(e) => r.record(
"fencing",
format!("unexpected error from a stale append: {e}"),
),
Ok(_) => r.record(
"fencing",
"a fenced writer appended after the lease was released — releasing threw \
away the epoch, so there was nothing left to fence against",
),
}
match store.acquire(run, "instance-c", LEASE).await {
Ok(third) => {
if third.epoch <= first.epoch {
r.record(
"fencing",
format!(
"after a release the epoch restarted at {} — a writer fenced at {} \
now outranks the legitimate owner",
third.epoch, first.epoch
),
);
}
}
Err(e) => r.record(
"fencing",
format!("a released lease was not claimable: {e}"),
),
}
}
async fn a_fenced_caller_cannot_release_the_new_owners_lease(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(old) = store.acquire(run, "instance-a", SHORT).await else {
return;
};
tokio::time::sleep(Duration::from_millis(1_200)).await;
let Ok(new) = store.acquire(run, "instance-b", LEASE).await else {
r.record("fencing", "an expired lease could not be taken over");
return;
};
if new.epoch == old.epoch {
r.record("fencing", "a takeover did not advance the epoch");
return;
}
if let Err(e) = store.release_lease(run, old.epoch).await {
r.record(
"fencing",
format!("a fenced caller's release must be a no-op, not an error: {e}"),
);
}
match store.acquire(run, "instance-c", LEASE).await {
Err(crate::core::StoreError::LeaseHeld { .. }) => {}
Err(e) => r.record(
"fencing",
format!("unexpected error after a stale release: {e}"),
),
Ok(_) => r.record(
"fencing",
"a fenced caller released the lease of the instance that replaced it, handing \
the run to a third party while its rightful owner is still writing",
),
}
}
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_fabricated_future_epoch_is_fenced(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "instance-a", LEASE).await else {
r.record("fencing", "acquire() failed on a fresh run");
return;
};
let invented = lease.epoch.saturating_add(1);
match store.append(invented, vec![admitted(run)]).await {
Err(crate::core::StoreError::Fenced { .. }) => {}
Err(e) => r.record("fencing", format!("expected Fenced, got {e}")),
Ok(_) => r.record(
"fencing",
"an epoch the store never issued was accepted. A fencing token is proof of a \
lease, not a number a caller may outrank by adding one",
),
}
}
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 sealed_runs_are_findable_by_how_they_ended(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let mut quarantined = Vec::new();
for (outcome, count) in [("quarantined", 2usize), ("succeeded", 1)] {
for _ in 0..count {
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
let _ = store
.append(lease.epoch, vec![admitted(run), concluded(run, outcome)])
.await;
if store.seal(run, lease.epoch, outcome).await.is_err() {
return;
}
if outcome == "quarantined" {
quarantined.push(run);
}
}
}
match store.runs_by_outcome("quarantined", 100).await {
Ok(found) => {
if found.len() != 2 {
r.record(
"outcome index",
format!(
"runs_by_outcome must find both quarantined runs, got {}",
found.len()
),
);
}
if found.iter().any(|run| !quarantined.contains(run)) {
r.record(
"outcome index",
"runs_by_outcome returned a run that ended some other way".to_owned(),
);
}
}
Err(e) => r.record("outcome index", format!("runs_by_outcome failed: {e}")),
}
let newest = quarantined.last().copied();
match store.runs_by_outcome("quarantined", 1).await {
Ok(found) if found.len() == 1 && Some(found[0]) == newest => {}
Ok(found) if found.len() == 1 => r.record(
"outcome index",
"runs_by_outcome(limit=1) kept the oldest run — a bounded query in \
ascending order never surfaces the quarantine that just happened"
.to_owned(),
),
Ok(found) => r.record(
"outcome index",
format!("runs_by_outcome(limit=1) returned {} runs", found.len()),
),
Err(e) => r.record(
"outcome index",
format!("runs_by_outcome(limit=1) failed: {e}"),
),
}
}
async fn the_outcome_index_follows_the_last_conclusion(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;
};
if let Err(e) = store
.append(lease.epoch, vec![admitted(run), concluded(run, "failed")])
.await
{
r.record(
"outcome index",
format!("appending a conclusion failed: {e}"),
);
return;
}
match store.runs_by_outcome("failed", 10).await {
Ok(found) if found.contains(&run) => {}
Ok(_) => r.record(
"outcome index",
"a failed conclusion was not indexed — an unsealed run's failure is \
a finding nobody can find"
.to_owned(),
),
Err(e) => r.record("outcome index", format!("runs_by_outcome failed: {e}")),
}
if let Err(e) = store
.append(lease.epoch, vec![concluded(run, "succeeded")])
.await
{
r.record("outcome index", format!("re-concluding failed: {e}"));
return;
}
let _ = store.seal(run, lease.epoch, "succeeded").await;
match store.runs_by_outcome("failed", 10).await {
Ok(found) if found.contains(&run) => r.record(
"outcome index",
"a run that later succeeded is still listed as failed — the backlog \
page never drains, and a wrong answer reads exactly like a right one"
.to_owned(),
),
Ok(_) => {}
Err(e) => r.record("outcome index", format!("runs_by_outcome failed: {e}")),
}
match store.runs_by_outcome("succeeded", 10).await {
Ok(found) if found.contains(&run) => {}
Ok(_) => r.record(
"outcome index",
"the re-conclusion did not move the run into the succeeded listing".to_owned(),
),
Err(e) => r.record("outcome index", format!("runs_by_outcome failed: {e}")),
}
}
async fn a_sealed_run_refuses_appends(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;
};
if let Err(e) = store.append(lease.epoch, vec![admitted(run)]).await {
r.record("seal", format!("append before seal failed: {e}"));
return;
}
if let Err(e) = store.seal(run, lease.epoch, "succeeded").await {
r.record("seal", format!("seal failed: {e}"));
return;
}
let Ok(frozen) = store.head(run).await else {
r.record("seal", "head() failed after seal".to_owned());
return;
};
match store
.append(lease.epoch, vec![started(run, key(201))])
.await
{
Ok(_) => r.record(
"seal",
"a sealed run accepted an append under the current epoch — the true \
head has moved past the leaf every checkpoint attests"
.to_owned(),
),
Err(crate::core::StoreError::RunSealed { .. }) => {}
Err(e) => r.record(
"seal",
format!(
"an append after seal was refused, but as '{e}' rather than \
RunSealed — a caller cannot tell a frozen run from a store fault"
),
),
}
match store.head(run).await {
Ok(after) if after.hash == frozen.hash && after.seq == frozen.seq => {}
Ok(_) => r.record(
"seal",
"the refused append still moved the chain head".to_owned(),
),
Err(e) => r.record("seal", format!("head() failed after refusal: {e}")),
}
}
async fn a_case_scan_selects_one_matter(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let mine = CaseId::generate();
let theirs = CaseId::generate();
for (case, count) in [(mine, 2), (theirs, 1)] {
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
for _ in 0..count {
let _ = store
.append(lease.epoch, vec![admitted(run).case(case)])
.await;
}
}
match store.case_history(mine, 100).await {
Ok(records) => {
if records.len() != 2 {
r.record(
"case history",
format!(
"case_history must return this matter's records, got {} of 2",
records.len()
),
);
}
if records.iter().any(|rec| rec.body.case != Some(mine)) {
r.record(
"case history",
"case_history returned a record belonging to another matter".to_owned(),
);
}
}
Err(e) => r.record("case history", format!("case_history failed: {e}")),
}
match store.case_history(mine, 1).await {
Ok(records) if records.len() == 1 => {}
Ok(records) => r.record(
"case history",
format!("case_history(limit=1) returned {} records", records.len()),
),
Err(e) => r.record("case history", format!("case_history(limit=1) failed: {e}")),
}
}
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}")),
}
}
#[allow(clippy::too_many_lines)]
pub async fn authority(store: Arc<dyn crate::authority::AuthorityStore>) {
use crate::authority::{AuthorityError, AuthorityId, StandingAuthority};
use crate::core::{EffectKey, Spend, Timestamp};
let key = |n: u8| EffectKey::from_hex(&format!("{n:064x}")).expect("32 bytes of hex");
let at = Timestamp::UNIX_EPOCH;
let id = AuthorityId::new("conformance-accumulate");
store
.issue(&StandingAuthority::new(
"conformance-accumulate",
"conformance",
Spend::money(1_000),
))
.await
.expect("issue");
let first = store
.draw(&id, key(1), Spend::money(600), at)
.await
.expect("a draw within the ceiling");
assert_eq!(
first.remaining,
Spend::money(400),
"the receipt must report what is left after this draw"
);
let refused = store
.draw(&id, key(2), Spend::money(500), at)
.await
.expect_err("the ceiling is cumulative across draws, not per draw");
assert!(
matches!(refused, AuthorityError::Exhausted { .. }),
"an over-draw must be Exhausted, not {refused:?}"
);
store
.draw(&id, key(3), Spend::money(400), at)
.await
.expect("exactly the remainder is still available after a refusal");
let id = AuthorityId::new("conformance-retry");
store
.issue(&StandingAuthority::new(
"conformance-retry",
"conformance",
Spend::money(1_000),
))
.await
.expect("issue");
let original = store
.draw(&id, key(10), Spend::money(300), at)
.await
.expect("draw");
let repeat = store
.draw(&id, key(10), Spend::money(300), at)
.await
.expect("a retry is not a second draw");
assert_eq!(
original, repeat,
"a retried draw must return its original receipt"
);
let state = store.state(&id).await.expect("state").expect("issued");
assert_eq!(
state.drawn,
Spend::money(300),
"a retry spent the authority twice"
);
assert_eq!(state.draws, 1, "a retry counted as a second draw");
store
.revoke(&id, "conformance withdrew it", at)
.await
.expect("revoke");
let after = store
.draw(&id, key(10), Spend::money(300), at)
.await
.expect("a draw that already landed must stay landed on retry");
assert_eq!(after, original);
let refused = store
.draw(&id, key(11), Spend::money(1), at)
.await
.expect_err("a new draw against a revoked authority");
assert!(
matches!(refused, AuthorityError::Revoked { .. }),
"revoked must be distinguishable from exhausted, got {refused:?}"
);
store
.revoke(&id, "a second, different reason", at)
.await
.expect("revoking twice is a retry");
let state = store.state(&id).await.expect("state").expect("issued");
let revocation = state.revoked.expect("revoked");
assert_eq!(
revocation.reason, "conformance withdrew it",
"the first reason must stand; overwriting loses why it was withdrawn"
);
let terms = StandingAuthority::new("conformance-terms", "conformance", Spend::money(500));
store.issue(&terms).await.expect("issue");
store
.issue(&terms)
.await
.expect("an identical re-issue is a retried deploy");
let conflict = store
.issue(&StandingAuthority::new(
"conformance-terms",
"conformance",
Spend::money(999_999),
))
.await
.expect_err("a ceiling must not be editable under whoever agreed to it");
assert!(
matches!(conflict, AuthorityError::AlreadyIssued(_)),
"got {conflict:?}"
);
store
.issue(
&StandingAuthority::new("conformance-expiry", "conformance", Spend::money(500))
.expires_at(Timestamp::from_unix_timestamp(1_000).expect("representable")),
)
.await
.expect("issue");
let id = AuthorityId::new("conformance-expiry");
store
.draw(
&id,
key(20),
Spend::money(1),
Timestamp::from_unix_timestamp(999).expect("representable"),
)
.await
.expect("before expiry the authority stands");
let expired = store
.draw(
&id,
key(21),
Spend::money(1),
Timestamp::from_unix_timestamp(2_000).expect("representable"),
)
.await
.expect_err("past its expiry");
assert!(
matches!(expired, AuthorityError::Expired { .. }),
"expiry must read the caller's instant, not a store clock; got {expired:?}"
);
store
.issue(
&StandingAuthority::new("conformance-draws", "conformance", Spend::money(10_000))
.max_draws(1),
)
.await
.expect("issue");
let id = AuthorityId::new("conformance-draws");
store
.draw(&id, key(30), Spend::money(1), at)
.await
.expect("first");
let spent = store
.draw(&id, key(31), Spend::money(1), at)
.await
.expect_err("one draw was all it permitted");
assert!(
matches!(spent, AuthorityError::DrawsSpent { .. }),
"a draw ceiling must refuse with money still left; got {spent:?}"
);
let unknown = store
.draw(
&AuthorityId::new("conformance-never-issued"),
key(40),
Spend::money(1),
at,
)
.await
.expect_err("never issued");
assert!(
matches!(unknown, AuthorityError::Unknown(_)),
"got {unknown:?}"
);
assert!(
store
.state(&AuthorityId::new("conformance-never-issued"))
.await
.expect("state")
.is_none()
);
let id = AuthorityId::new("conformance-racing-retry");
store
.issue(&StandingAuthority::new(
"conformance-racing-retry",
"conformance",
Spend::money(500),
))
.await
.expect("issue");
let carriers: Vec<_> = (0..8)
.map(|_| {
let store = Arc::clone(&store);
let id = id.clone();
tokio::spawn(async move { store.draw(&id, key(50), Spend::money(500), at).await })
})
.collect();
let mut receipts = Vec::new();
for carrier in carriers {
receipts.push(
carrier
.await
.expect("a carrier panicked")
.expect("every carrier of the one key must get the receipt, not a refusal"),
);
}
for r in &receipts {
assert_eq!(
*r, receipts[0],
"two carriers of one dispatch key were given different receipts"
);
}
let state = store.state(&id).await.expect("state").expect("issued");
assert_eq!(
state.drawn,
Spend::money(500),
"racing retries spent the authority more than once"
);
assert_eq!(state.draws, 1, "racing retries counted as several draws");
}