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 live = store
.current("memory-a", None)
.await
.expect("current")
.expect("a revised memory is still current");
assert_eq!(
live.version, 2,
"current must be the revision, not the v1 it superseded"
);
assert!(
store
.current("memory-never-written", None)
.await
.expect("current of nothing")
.is_none()
);
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()
);
store
.remember(&make("chain-a", "team-chain", "support", json!({"n": 1})))
.await
.expect("chain root");
let chain_a = store
.version("chain-a", 1)
.await
.expect("read root")
.expect("root exists");
let mut chain_b = make("chain-b", "team-chain", "support", json!({"n": 2}));
chain_b.derived_from = vec![Selected {
id: chain_a.id.clone(),
version: chain_a.version,
digest: chain_a.selection_digest(),
}];
store.remember(&chain_b).await.expect("chain middle");
let chain_b = store
.version("chain-b", 1)
.await
.expect("read middle")
.expect("middle exists");
let mut chain_c = make("chain-c", "team-chain", "support", json!({"n": 3}));
chain_c.derived_from = vec![Selected {
id: chain_b.id.clone(),
version: chain_b.version,
digest: chain_b.selection_digest(),
}];
store.remember(&chain_c).await.expect("chain leaf");
store.forget("chain-b").await.expect("erase the middle");
assert_eq!(
store
.forget_cascading("chain-a")
.await
.expect("cascade through the tombstone"),
2,
"the cascade must erase the root and the leaf — and count only what it \
erased, not the tombstone it passed through"
);
assert!(
store
.version("chain-c", 1)
.await
.expect("leaf lookup")
.is_none(),
"a derivative reached only through an already-erased intermediate \
survived the cascade — the poisoned source's summary of a summary is \
still readable"
);
assert_eq!(
store
.forget_cascading("chain-a")
.await
.expect("cascading a second time is not an error"),
0,
"a cascade over an already-erased id reported an erasure it did not \
perform — the count is what an erasure request is answered with"
);
store
.remember(&make(
"exp-u",
"team-expiry-chain",
"support",
json!({"n": 1}),
))
.await
.expect("expiry-chain root");
let exp_u = store
.version("exp-u", 1)
.await
.expect("read root")
.expect("root exists");
let mut exp_e = make("exp-e", "team-expiry-chain", "support", json!({"n": 2}));
exp_e.expires_at = Some(at(1_760_000_150));
exp_e.derived_from = vec![Selected {
id: exp_u.id.clone(),
version: exp_u.version,
digest: exp_u.selection_digest(),
}];
store.remember(&exp_e).await.expect("expiring middle");
let exp_e = store
.version("exp-e", 1)
.await
.expect("read middle")
.expect("middle exists");
let mut exp_d = make("exp-d", "team-expiry-chain", "support", json!({"n": 3}));
exp_d.derived_from = vec![Selected {
id: exp_e.id.clone(),
version: exp_e.version,
digest: exp_e.selection_digest(),
}];
store.remember(&exp_d).await.expect("expiry-chain leaf");
assert_eq!(
store
.sweep_expired(at(1_760_000_150))
.await
.expect("expire the middle"),
1,
"exactly the middle link expires"
);
assert_eq!(
store
.forget_cascading("exp-u")
.await
.expect("cascade through the expired tombstone"),
2,
"the cascade must erase the root and the leaf, counting only what it \
erased — not the expired tombstone it passed through"
);
assert!(
store
.version("exp-d", 1)
.await
.expect("leaf lookup")
.is_none(),
"a derivative reached only through an *expired* intermediate survived \
the cascade — the expiry sweep severed the lineage that a later \
erasure of the poisoned source needed to route through"
);
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"
);
assert!(
store
.current("memory-expiring", Some(at(1_760_000_099)))
.await
.expect("current before expiry")
.is_some()
);
assert!(
store
.current("memory-expiring", Some(at(1_760_000_100)))
.await
.expect("current at expiry")
.is_none(),
"current must apply the inclusive expiry cutoff"
);
assert!(
store
.current("memory-expiring", None)
.await
.expect("current without cutoff")
.is_some(),
"no cutoff means no expiry check, exactly as on recall"
);
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()
);
assert!(
store
.current("memory-expiring", None)
.await
.expect("current after sweep")
.is_none(),
"a swept id has no current version under any cutoff"
);
let mut fading = make(
"memory-fading",
"team-fading",
"support",
json!({"fading": true}),
);
fading.access_retention_seconds = Some(60);
store.remember(&fading).await.expect("fading memory");
assert_eq!(
store
.recall(&Recall::about("team-fading").at(at(1_760_000_059)))
.await
.expect("inside the untouched window")
.len(),
1
);
assert!(
store
.recall(&Recall::about("team-fading").at(at(1_760_000_060)))
.await
.expect("after the untouched window")
.is_empty(),
"an untouched sliding-retention memory outlived its own window — the \
window must open at the write, or an item nobody ever touches is \
immortal"
);
assert_eq!(
store
.sweep_expired(at(1_760_000_060))
.await
.expect("sweep the untouched window"),
1,
"the sweep must collect an untouched sliding-retention memory"
);
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_040))
.await
.expect("touch sliding retention");
assert_eq!(
store
.recall(&Recall::about("team-retention").at(at(1_760_000_099)))
.await
.expect("inside the touched window")
.len(),
1,
"a journaled touch did not slide the window"
);
assert!(
store
.current("memory-sliding", Some(at(1_760_000_100)))
.await
.expect("current after the touched window")
.is_none(),
"current must honour the sliding window exactly as recall does"
);
store
.touch(&["memory-sliding".to_owned()], at(1_760_000_190))
.await
.expect("touch near the ceiling");
assert_eq!(
store
.recall(&Recall::about("team-retention").at(at(1_760_000_199)))
.await
.expect("just under the ceiling")
.len(),
1
);
assert!(
store
.recall(&Recall::about("team-retention").at(at(1_760_000_200)))
.await
.expect("at the ceiling")
.is_empty(),
"a touch extended a memory past its immutable `expires_at` — the \
ceiling is the disposal date an operator can state, and sliding \
retention may only shorten a life, never lengthen it"
);
assert_eq!(
store
.sweep_expired(at(1_760_000_200))
.await
.expect("sweep at the ceiling"),
1,
"the sweep must erase at the ceiling even though the touched window \
reaches past it"
);
assert!(
store
.version("memory-sliding", 1)
.await
.expect("swept sliding version")
.is_none()
);
let mut lazarus = make("memory-lazarus", "team-lazarus", "support", json!({"n": 1}));
lazarus.access_retention_seconds = Some(60);
store.remember(&lazarus).await.expect("lazarus memory");
store
.touch(&["memory-lazarus".to_owned()], at(1_760_000_100))
.await
.expect("touch a lapsed window");
assert_eq!(
store
.recall(&Recall::about("team-lazarus").at(at(1_760_000_150)))
.await
.expect("inside the reopened window")
.len(),
1,
"a touch on a lapsed-but-unswept id must resurrect it — expiry \
becomes fact at the sweep, and until then the id is alive to touch"
);
assert!(
store
.recall(&Recall::about("team-lazarus").at(at(1_760_000_160)))
.await
.expect("after the reopened window")
.is_empty(),
"the reopened window is a window, not immortality"
);
let sel = |item: &MemoryItem| Selected {
id: item.id.clone(),
version: item.version,
digest: item.selection_digest(),
};
store
.remember(&make("roll-a", "team-roll", "support", json!({"n": 1})))
.await
.expect("roll source a");
store
.remember(&make("roll-b", "team-roll", "support", json!({"n": 2})))
.await
.expect("roll source b");
let roll_a = store
.version("roll-a", 1)
.await
.expect("read a")
.expect("a exists");
let roll_b = store
.version("roll-b", 1)
.await
.expect("read b")
.expect("b exists");
let mut summary = make("roll-s", "team-roll", "support", json!({"sum": "of a"}));
summary.derived_from = vec![sel(&roll_a)];
store.remember(&summary).await.expect("summary v1");
let summary_v1 = store
.version("roll-s", 1)
.await
.expect("read v1")
.expect("v1 exists");
let mut over = make("roll-t", "team-roll", "support", json!({"sum": "of s v1"}));
over.derived_from = vec![sel(&summary_v1)];
store.remember(&over).await.expect("summary of the summary");
let mut summary2 = make("roll-s", "team-roll", "support", json!({"sum": "of b"}));
summary2.created_at = at(1_760_000_001);
summary2.derived_from = vec![sel(&roll_b)];
store.remember(&summary2).await.expect("summary v2");
let v1 = store
.version("roll-s", 1)
.await
.expect("v1 read")
.expect("v1 kept before the cascade");
assert_eq!(v1.derived_from[0].id, "roll-a");
assert!(
store
.derivatives("roll-a")
.await
.expect("derivatives of a")
.is_empty(),
"a re-derived summary is no longer a live derivative of its old source"
);
assert_eq!(
store
.forget_cascading("roll-a")
.await
.expect("cascade from the superseded source"),
3,
"the cascade must erase roll-a, roll-s's superseded v1 (counted once) \
and roll-t, which was derived from exactly that version"
);
assert!(
store
.version("roll-a", 1)
.await
.expect("a lookup")
.is_none()
);
assert!(
store
.version("roll-s", 1)
.await
.expect("superseded v1 lookup")
.is_none(),
"the superseded summary version that absorbed the erased source is \
still readable — supersession sheltered absorbed content from the \
cascade"
);
assert!(
store
.version("roll-t", 1)
.await
.expect("roll-t lookup")
.is_none(),
"a derivative of the superseded version survived — the cascade did \
not route through the version it erased"
);
let kept = store
.version("roll-s", 2)
.await
.expect("v2 read")
.expect("the clean re-derivation survives");
assert_eq!(kept.derived_from[0].id, "roll-b");
assert!(
store
.recall(&Recall::about("team-roll"))
.await
.expect("recall after the cascade")
.iter()
.any(|item| item.id == "roll-s"),
"the re-derived current summary must stay recallable"
);
assert!(
store
.version("roll-b", 1)
.await
.expect("b lookup")
.is_some(),
"the clean source is untouched"
);
let mut order_old = make("order-old", "team-order", "alpha", json!({"n": 1}));
order_old.created_at = at(1_760_000_300);
store
.remember(&order_old)
.await
.expect("older, first purpose");
let mut order_new = make("order-new", "team-order", "zeta", json!({"n": 2}));
order_new.created_at = at(1_760_000_400);
store
.remember(&order_new)
.await
.expect("newest, last purpose");
let mut order_trusted = make("order-trusted", "team-order", "middle", json!({"n": 3}));
order_trusted.trust = Trust::Trusted;
order_trusted.created_at = at(1_760_000_100);
store
.remember(&order_trusted)
.await
.expect("trusted, oldest");
let picked = store
.recall(&Recall::about("team-order").limit(2))
.await
.expect("ordered recall");
assert_eq!(
picked.iter().map(|i| i.id.as_str()).collect::<Vec<_>>(),
vec!["order-trusted", "order-new"],
"a purpose-less recall must rank trust first and then recency \
globally across purposes — not whichever purpose sorts first"
);
let all = store
.recall(&Recall::about("team-order").limit(3))
.await
.expect("full ordered recall");
assert_eq!(
all.iter().map(|i| i.id.as_str()).collect::<Vec<_>>(),
vec!["order-trusted", "order-new", "order-old"]
);
assert_eq!(
store.subject_ids("team-order").await.expect("subject ids"),
vec!["order-new", "order-old", "order-trusted"],
"subject_ids must name every current id of the subject"
);
assert!(
store
.subject_ids("team-nobody")
.await
.expect("empty subject")
.is_empty()
);
}
#[allow(clippy::too_many_lines)]
pub async fn event_erasure(store: Arc<dyn crate::case::EventStore>) {
use crate::core::{
CorrelationKey, InboundEvent, Phase, RunId, StepId, Subscription, Timestamp,
};
use serde_json::json;
let at = |seconds| Timestamp::from_unix_timestamp(seconds).expect("representable time");
let event = |id: &str, payload: serde_json::Value| InboundEvent {
source: "counterparty".to_owned(),
id: id.to_owned(),
kind: "reply".to_owned(),
correlation: vec![CorrelationKey::new("order", id.to_owned())],
payload,
};
let sub = |run: RunId, effect, id: &str| Subscription {
run,
case: None,
effect,
step: StepId(1),
phase: Phase::Forward,
kind: "reply".to_owned(),
correlation: vec![CorrelationKey::new("order", id.to_owned())],
};
let run = RunId::generate();
let effect = key(201);
let pii = json!({"pii": "the counterparty's content"});
let first = event("A-1", pii.clone());
assert!(store.buffer(&first, at(1_000)).await.expect("buffer"));
store
.subscribe(&sub(run, effect, "A-1"), at(1_001))
.await
.expect("subscribe");
let claimed = store
.claim_for(&sub(run, effect, "A-1"), at(1_002))
.await
.expect("claim")
.expect("a buffered event matches");
assert_eq!(
claimed.event.payload, pii,
"delivery hands over the payload"
);
let recovered = store
.claim_for(&sub(run, effect, "A-1"), at(1_003))
.await
.expect("re-claim")
.expect("the run's own claim is returned to it");
assert_eq!(
recovered.event.payload, pii,
"a payload stripped at the claim is lost to a run that crashed \
between claim and resume — stripping must wait for the unsubscribe"
);
store.unsubscribe(run, effect).await.expect("unsubscribe");
store
.subscribe(&sub(run, effect, "A-1"), at(1_004))
.await
.expect("resubscribe");
let after = store
.claim_for(&sub(run, effect, "A-1"), at(1_005))
.await
.expect("claim after the strip")
.expect("the row keeps its identity");
assert_eq!(
after.event.payload,
serde_json::Value::Null,
"a delivered payload outlived its delivery in the buffer — the \
journal holds the delivered copy, and the buffer needs only the \
(source, id) identity"
);
store.unsubscribe(run, effect).await.expect("clean up");
assert!(
!store.buffer(&first, at(1_006)).await.expect("replay"),
"shedding the payload deleted the dedup row — a replay of the \
delivered message was accepted as new, so erasure reopened the door \
it was supposed to close"
);
let unclaimed = event("B-1", json!({"pii": "unclaimed"}));
assert!(store.buffer(&unclaimed, at(2_000)).await.expect("buffer"));
assert!(
store
.erase_payload("counterparty", "B-1")
.await
.expect("erase the unclaimed payload"),
"the row exists, so the erasure must report having acted"
);
assert_eq!(
store
.sweep_unclaimed(at(2_100), "aged out")
.await
.expect("sweep"),
1,
"an erased event still dead-letters — erasure removes content, never \
accounting"
);
let letters = store.dead_letters(10).await.expect("dead letters");
let erased = letters
.iter()
.find(|letter| letter.event.id == "B-1")
.expect("the erased event is still listed");
assert_eq!(
erased.event.payload,
serde_json::Value::Null,
"an erased unclaimed event still yielded its payload"
);
assert_eq!(erased.reason, "aged out");
assert_eq!(
erased.event.correlation,
vec![CorrelationKey::new("order", "B-1")],
"correlation keys are the identity the operator routes on, not content"
);
let dead = event("C-1", json!({"pii": "dead"}));
assert!(store.buffer(&dead, at(3_000)).await.expect("buffer"));
assert_eq!(
store
.sweep_unclaimed(at(3_100), "nobody came")
.await
.expect("sweep"),
1
);
let letters = store.dead_letters(10).await.expect("dead letters");
assert_eq!(
letters
.iter()
.find(|letter| letter.event.id == "C-1")
.expect("listed")
.event
.payload,
json!({"pii": "dead"})
);
assert!(
store
.erase_payload("counterparty", "C-1")
.await
.expect("erase the dead letter")
);
let letters = store.dead_letters(10).await.expect("dead letters");
assert_eq!(
letters
.iter()
.find(|letter| letter.event.id == "C-1")
.expect("still listed")
.event
.payload,
serde_json::Value::Null,
"an erased dead letter still yielded its payload"
);
assert!(
!store
.erase_payload("counterparty", "never-arrived")
.await
.expect("absent row")
);
}
fn admitted_under(run: RunId, key: &str) -> 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,
canon: crate::core::canon::VERSION,
idempotency_key: Some(key.to_owned()),
},
)
}
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,
canon: crate::core::canon::VERSION,
idempotency_key: None,
},
)
}
fn concluded(run: RunId, outcome: &str) -> Append {
Append::new(
run,
RecordKind::RunConcluded {
outcome: outcome.to_owned(),
reason: None,
exhaustion: None,
live_spend: crate::core::Spend::default(),
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;
acquire_claims_and_never_renews(fresh, &mut r).await;
a_renewal_extends_only_what_is_still_held(fresh, &mut r).await;
a_batch_spans_one_run(fresh, &mut r).await;
a_stale_epoch_cannot_seal(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;
the_discovery_index_pages_in_a_total_order(fresh, &mut r).await;
abandonment_names_the_dead_not_the_finished(fresh, &mut r).await;
an_admission_key_admits_one_run(fresh, &mut r).await;
a_key_is_claimed_by_the_append_that_writes_the_run(fresh, &mut r).await;
a_refused_admission_leaves_its_key_free(fresh, &mut r).await;
an_unkeyed_admission_claims_nothing(fresh, &mut r).await;
retiring_a_key_frees_it_and_leaves_the_run(fresh, &mut r).await;
r
}
async fn retiring_a_key_frees_it_and_leaves_the_run(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let key = "conformance/retired";
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
if store
.append(lease.epoch, vec![admitted_under(run, key)])
.await
.is_err()
{
return;
}
match store
.forget_admissions(crate::core::Timestamp::from_unix_timestamp(0).expect("epoch"))
.await
{
Ok(0) => {}
Ok(n) => r.record(
"admission keys",
format!("a cutoff older than every claim retired {n} of them"),
),
Err(e) => {
r.record("admission keys", format!("forget_admissions failed: {e}"));
return;
}
}
if !matches!(store.admitted_as(key).await, Ok(Some(_))) {
r.record(
"admission keys",
"a key was released by a retirement that reported retiring nothing",
);
}
match store
.forget_admissions(
crate::core::Timestamp::from_unix_timestamp(4_102_444_800).expect("year 2100"),
)
.await
{
Ok(1) => {}
Ok(n) => r.record(
"admission keys",
format!("retiring one claim past its window reported {n}"),
),
Err(e) => {
r.record("admission keys", format!("forget_admissions failed: {e}"));
return;
}
}
match store.admitted_as(key).await {
Ok(None) => {}
Ok(Some(held)) => r.record(
"admission keys",
format!("a retired key is still held by {held}, so retirement reclaims nothing"),
),
Err(e) => r.record("admission keys", format!("admitted_as failed: {e}")),
}
match store.read(run, 1).await {
Ok(records) if !records.is_empty() => {}
Ok(_) => r.record(
"admission keys",
"retiring a key erased the run it named — a retention policy must not \
edit history, and the key's authority is the RunAdmitted record",
),
Err(e) => r.record("admission keys", format!("the run became unreadable: {e}")),
}
}
async fn an_admission_key_admits_one_run(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let key = "conformance/one-message";
let first = RunId::generate();
let Ok(lease) = store.acquire(first, "conformance", LEASE).await else {
return;
};
if let Err(e) = store
.append(lease.epoch, vec![admitted_under(first, key)])
.await
{
r.record(
"admission keys",
format!("the first admission under a fresh key was refused: {e}"),
);
return;
}
let second = RunId::generate();
let Ok(lease) = store.acquire(second, "conformance", LEASE).await else {
return;
};
match store
.append(lease.epoch, vec![admitted_under(second, key)])
.await
{
Err(crate::core::StoreError::DuplicateAdmission { run, .. }) => {
if run != first.to_string() {
r.record(
"admission keys",
format!(
"the refusal named '{run}' as the key's holder, but {first} holds it — a caller sent there reads a run that answers a different question"
),
);
}
}
Ok(_) => r.record(
"admission keys",
"a second run was admitted under a key already issued — an emitter's redelivery would start a second run, which is the whole failure this refuses"
.to_owned(),
),
Err(e) => r.record(
"admission keys",
format!("the duplicate was refused, but not as DuplicateAdmission: {e}"),
),
}
match store.admitted_as(key).await {
Ok(Some(held)) if held == first => {}
Ok(other) => r.record(
"admission keys",
format!("admitted_as must name {first}, said {other:?}"),
),
Err(e) => r.record("admission keys", format!("admitted_as failed: {e}")),
}
}
async fn a_key_is_claimed_by_the_append_that_writes_the_run(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let key = "conformance/not-yet";
match store.admitted_as(key).await {
Ok(None) => {}
Ok(Some(run)) => {
r.record(
"admission keys",
format!("a key nothing has admitted under is held by {run}"),
);
return;
}
Err(e) => {
r.record("admission keys", format!("admitted_as failed: {e}"));
return;
}
}
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
if let Ok(Some(held)) = store.admitted_as(key).await {
r.record(
"admission keys",
format!("taking a lease claimed the key for {held}; only the append may"),
);
}
let _ = store
.append(lease.epoch, vec![admitted_under(run, key)])
.await;
match store.admitted_as(key).await {
Ok(Some(held)) if held == run => {}
Ok(other) => r.record(
"admission keys",
format!("after the append the key must be held by {run}, said {other:?}"),
),
Err(e) => r.record("admission keys", format!("admitted_as failed: {e}")),
}
}
async fn a_refused_admission_leaves_its_key_free(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let key = "conformance/refused";
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
let other = RunId::generate();
let refused = store
.append(lease.epoch, vec![admitted_under(run, key), admitted(other)])
.await;
if refused.is_ok() {
r.record(
"admission keys",
"a batch spanning two runs was accepted; the rest of this check cannot mean anything"
.to_owned(),
);
return;
}
match store.admitted_as(key).await {
Ok(None) => {}
Ok(Some(held)) => r.record(
"admission keys",
format!(
"a refused append left the key held by {held} — the key committed without the run it names, so every future retry is answered with a run that has no journal"
),
),
Err(e) => r.record("admission keys", format!("admitted_as failed: {e}")),
}
}
async fn an_unkeyed_admission_claims_nothing(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
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![admitted(run)]).await {
r.record(
"admission keys",
format!("an admission carrying no key was refused: {e}"),
);
return;
}
}
if let Ok(Some(held)) = store.admitted_as("").await {
r.record(
"admission keys",
format!("an unkeyed admission was filed under the empty key, held by {held}"),
);
}
}
async fn abandonment_names_the_dead_not_the_finished(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let dead = RunId::generate();
let finished = RunId::generate();
let alive = RunId::generate();
let Ok(_) = store.acquire(dead, "instance-a", SHORT).await else {
return;
};
let Ok(f) = store.acquire(finished, "instance-b", SHORT).await else {
return;
};
if store.release_lease(finished, f.epoch).await.is_err() {
r.record("recovery", "release_lease of a held lease failed");
return;
}
let Ok(_) = store.acquire(alive, "instance-c", LEASE).await else {
return;
};
tokio::time::sleep(Duration::from_millis(1_500)).await;
match store.abandoned_runs(10).await {
Ok(listed) => {
if !listed.contains(&dead) {
r.record(
"recovery",
"a lease that expired while still naming an owner was not listed — \
the run its instance died holding is stranded forever",
);
}
if listed.contains(&finished) {
r.record(
"recovery",
"a released lease was listed as abandoned — every clean exit would \
be 'recovered' on every tick",
);
}
if listed.contains(&alive) {
r.record(
"recovery",
"a live lease was listed as abandoned — the sweep would take a run \
over from under its working owner's heartbeat",
);
}
}
Err(e) => r.record("recovery", format!("abandoned_runs failed: {e}")),
}
}
async fn the_discovery_index_pages_in_a_total_order(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let mut written = Vec::new();
for _ in 0..4 {
let run = RunId::generate();
let Ok(lease) = store.acquire(run, "conformance", LEASE).await else {
return;
};
if store
.append(lease.epoch, vec![admitted(run)])
.await
.is_err()
{
return;
}
written.push(run);
}
let whole = match store.recent_runs(None, 100).await {
Ok(rows) => rows,
Err(e) => {
r.record("chaining", format!("recent_runs failed: {e}"));
return;
}
};
if whole.len() != 4 {
r.record(
"chaining",
format!("recent_runs returned {} rows for 4 runs", whole.len()),
);
return;
}
let mut sorted = whole.clone();
sorted.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| b.0.cmp(&a.0)));
if whole != sorted {
r.record(
"chaining",
"recent_runs is not ordered by (updated_at, run) descending, so a \
cursor into it cannot resume a page"
.to_owned(),
);
}
match store.recent_runs(None, 2).await {
Ok(rows) if rows.len() > 2 => r.record(
"chaining",
format!("recent_runs(limit=2) returned {} rows", rows.len()),
),
Err(e) => r.record("chaining", format!("recent_runs(limit=2) failed: {e}")),
Ok(_) => {}
}
let mut paged = Vec::new();
let mut cursor: Option<(u64, RunId)> = None;
for _ in 0..4 {
let Ok(rows) = store.recent_runs(cursor, 2).await else {
r.record("chaining", "a cursored recent_runs page failed".to_owned());
return;
};
if rows.is_empty() {
break;
}
cursor = rows.last().map(|(run, updated)| (*updated, *run));
paged.extend(rows);
}
if paged != whole {
r.record(
"chaining",
format!(
"paging recent_runs two at a time did not reassemble the index: \
{} rows paged against {} whole — a page boundary either repeats \
a run or drops one, and both read as a healthy listing",
paged.len(),
whole.len()
),
);
}
}
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", SHORT).await else {
return;
};
tokio::time::sleep(Duration::from_millis(1_500)).await;
let Ok(again) = store.acquire(run, "instance-b", LEASE).await else {
r.record("fencing", "an expired lease could not be claimed");
return;
};
if again.epoch <= a.epoch {
r.record(
"fencing",
format!(
"claiming an expired lease must advance the epoch ({} did not move past \
{}) — without the bump the dead owner's writes are never fenced",
again.epoch, a.epoch
),
);
}
}
async fn acquire_claims_and_never_renews(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let run = RunId::generate();
let Ok(_) = store.acquire(run, "instance-a", LEASE).await else {
return;
};
match store.acquire(run, "instance-a", LEASE).await {
Err(crate::core::StoreError::LeaseHeld { .. }) => {}
Err(e) => r.record(
"fencing",
format!("a same-owner acquire of a held lease must fail as LeaseHeld, not {e}"),
),
Ok(lease) => r.record(
"fencing",
format!(
"a same-owner acquire of a held lease succeeded at epoch {} — acquire \
renewed instead of claiming, so two executors on one instance share \
one epoch and fencing cannot tell them apart",
lease.epoch
),
),
}
}
async fn a_renewal_extends_only_what_is_still_held(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 {
return;
};
match store.renew(run, "instance-a", lease.epoch, LEASE).await {
Ok(renewed) if renewed.epoch == lease.epoch => {}
Ok(renewed) => r.record(
"fencing",
format!(
"a renewal moved the epoch ({} became {}) — the owner is fenced \
against its own in-flight writes",
lease.epoch, renewed.epoch
),
),
Err(e) => r.record("fencing", format!("renewing a held lease failed: {e}")),
}
if store
.renew(run, "instance-b", lease.epoch, LEASE)
.await
.is_ok()
{
r.record(
"fencing",
"another owner renewed a lease it does not hold — renewal became theft",
);
}
if store
.renew(run, "instance-a", lease.epoch + 1, LEASE)
.await
.is_ok()
{
r.record(
"fencing",
"a renewal with a fabricated epoch succeeded — a caller invented \
authority it never acquired",
);
}
if store.release_lease(run, lease.epoch).await.is_err() {
r.record("fencing", "release_lease of a held lease failed");
return;
}
match store.renew(run, "instance-a", lease.epoch, LEASE).await {
Err(_) => {}
Ok(_) => r.record(
"fencing",
"a renewal re-took a lease the owner had already released — the \
heartbeat racing its own run's conclusion leaves a live lease over \
a concluded run, and the sweep recovers a clean exit forever",
),
}
match store.abandoned_runs(10).await {
Ok(listed) if listed.contains(&run) => r.record(
"recovery",
"a released run appears abandoned after a refused renewal — the \
renewal wrote something it had no claim to",
),
_ => {}
}
}
async fn a_batch_spans_one_run(fresh: Factory<'_>, r: &mut Report) {
r.checked += 1;
let store = fresh().await;
let a = RunId::generate();
let b = RunId::generate();
let Ok(lease) = store.acquire(a, "conformance", LEASE).await else {
return;
};
if store
.append(lease.epoch, vec![admitted(a), admitted(b)])
.await
.is_ok()
{
r.record(
"atomicity",
"a batch spanning two runs was accepted — the second run's record was \
sealed into the first run's chain under the first run's fence",
);
return;
}
for run in [a, b] {
if store.read(run, 1).await.map_or(0, |v| v.len()) != 0 {
r.record(
"atomicity",
"a refused cross-run batch still left records behind",
);
}
}
}
async fn a_stale_epoch_cannot_seal(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 store
.append(lease.epoch, vec![admitted(run)])
.await
.is_err()
{
return;
}
if store.seal(run, lease.epoch + 7, "succeeded").await.is_ok() {
r.record(
"fencing",
"a seal under a fabricated epoch was accepted — a fenced writer froze \
a chain it does not own",
);
}
if let Err(e) = store.seal(run, lease.epoch, "succeeded").await {
r.record(
"fencing",
format!("a seal under the current epoch was refused: {e}"),
);
}
}
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");
}