diff --git a/crates/nexum-sdk/src/keeper.rs b/crates/nexum-sdk/src/keeper.rs index 115b9725..7457acf3 100644 --- a/crates/nexum-sdk/src/keeper.rs +++ b/crates/nexum-sdk/src/keeper.rs @@ -287,9 +287,51 @@ fn read_u64(host: &H, key: &str) -> Result, Fault } } -/// Receipt-keyed idempotency journal: presence markers under a fixed -/// prefix. The marker value is empty - presence of the key is the -/// receipt - so re-recording is idempotent by construction. +/// RESERVED value tag: the marker owes a durable effect, its body and +/// `next_eligible` follow. +const RESERVED_TAG: u8 = 0x01; +/// COMMITTED value tag: the durable effect landed; nothing is owed. +const COMMITTED_TAG: u8 = 0x02; + +/// Presence class of a durable-effect marker under the tag-first value +/// scheme. RESERVED and COMMITTED can never collide by presence: the +/// leading tag byte tells them apart, and a corrupt tag falls to +/// [`Reserved`](Mark::Reserved) so a reconcile resubmits, never skips. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum Mark { + /// A submit is in flight or owed; the effect is not yet durable. + Reserved, + /// The upstream effect is durable; no resubmit is owed. + Committed, +} + +/// A live reservation enumerated from the journal: the receipt key +/// (prefix-stripped), the earliest epoch-seconds it may retry at, and +/// the stored body. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct Reservation { + /// Receipt key as passed to [`Journal::reserve`], prefix stripped. + pub key: String, + /// Earliest Unix-seconds a retry is eligible; `0` when unset. + pub next_eligible: u64, + /// Opaque reservation body stored alongside the marker. + pub body: Vec, +} + +/// Encode a RESERVED value: `[0x01] ++ be64(next_eligible) ++ body`. +fn reserved_value(next_eligible: u64, body: &[u8]) -> Vec { + let mut v = Vec::with_capacity(9 + body.len()); + v.push(RESERVED_TAG); + v.extend_from_slice(&next_eligible.to_be_bytes()); + v.extend_from_slice(body); + v +} + +/// Receipt-keyed idempotency journal under a fixed prefix. The observed +/// path ([`record`](Self::record) / [`contains`](Self::contains)) stores +/// an empty presence marker, so re-recording is idempotent. The +/// durable-effect path ([`reserve`](Self::reserve) / [`commit`](Self::commit)) +/// tags the value so [`mark`](Self::mark) tells RESERVED from COMMITTED. pub struct Journal<'h, H> { host: &'h H, prefix: &'static str, @@ -320,12 +362,89 @@ impl<'h, H: LocalStoreHost> Journal<'h, H> { } /// Whether the receipt is already journalled. + /// + /// Presence-only: reports the key exists, not whether it is RESERVED + /// or COMMITTED. MUST NOT guard a reserve/commit journal; use + /// [`mark`](Self::mark) there, else a corrupt marker is guard-skipped + /// yet never reconciled. pub fn contains(&self, receipt: &str) -> Result { Ok(self .host .get(&format!("{}{receipt}", self.prefix))? .is_some()) } + + /// Reserve `key`: write a RESERVED marker with `next_eligible` 0. + pub fn reserve(&self, key: &str, body: &[u8]) -> Result<(), Fault> { + self.host + .set(&format!("{}{key}", self.prefix), &reserved_value(0, body)) + } + + /// Re-park a reservation: rewrite its RESERVED marker with + /// `next_eligible = until`, body unchanged. + pub fn park(&self, key: &str, body: &[u8], until: u64) -> Result<(), Fault> { + self.host.set( + &format!("{}{key}", self.prefix), + &reserved_value(until, body), + ) + } + + /// Commit `key`: overwrite with the COMMITTED marker. Idempotent. + pub fn commit(&self, key: &str) -> Result<(), Fault> { + self.host + .set(&format!("{}{key}", self.prefix), &[COMMITTED_TAG]) + } + + /// Release `key`: delete the marker. No-op when absent. + pub fn release(&self, key: &str) -> Result<(), Fault> { + self.host.delete(&format!("{}{key}", self.prefix)) + } + + /// Classify `key`'s marker. Absent: `None`. Legacy empty or a + /// leading `0x02`: [`Committed`](Mark::Committed). Any other first + /// byte (`0x01`, unknown, corrupt): [`Reserved`](Mark::Reserved), so + /// a corrupt tag reconciles rather than skips. + pub fn mark(&self, key: &str) -> Result, Fault> { + let Some(v) = self.host.get(&format!("{}{key}", self.prefix))? else { + return Ok(None); + }; + Ok(Some(match v.first() { + None | Some(&COMMITTED_TAG) => Mark::Committed, + Some(_) => Mark::Reserved, + })) + } + + /// Enumerate every live reservation: each key whose value is + /// non-empty and not COMMITTED. Parse is fail-safe and agrees with + /// [`mark`](Self::mark) - a well-formed `0x01` marker yields its + /// `next_eligible` and body, any other non-committed value yields + /// `next_eligible` 0 and the bytes after the tag. Feeds the sweep + /// reconcile and operator enumerate. + pub fn pending(&self) -> Result, Fault> { + let mut out = Vec::new(); + for full in self.host.list_keys(self.prefix)? { + let Some(v) = self.host.get(&full)? else { + continue; + }; + if v.first().is_none_or(|&b| b == COMMITTED_TAG) { + continue; + } + let key = full.strip_prefix(self.prefix).unwrap_or(&full).to_owned(); + let (next_eligible, body) = if v[0] == RESERVED_TAG && v.len() >= 9 { + let mut be = [0u8; 8]; + be.copy_from_slice(&v[1..9]); + (u64::from_be_bytes(be), v[9..].to_vec()) + } else { + (0, v.get(1..).unwrap_or(&[]).to_vec()) + }; + out.push(Reservation { + key, + next_eligible, + body, + }); + } + Ok(out) + } } /// One poll dispatch's world view: chain, block height, and the block diff --git a/crates/nexum-sdk/tests/keeper.rs b/crates/nexum-sdk/tests/keeper.rs index c2e2ea9f..1d7f33ee 100644 --- a/crates/nexum-sdk/tests/keeper.rs +++ b/crates/nexum-sdk/tests/keeper.rs @@ -8,8 +8,8 @@ use alloy_primitives::{Address, B256, address, b256}; use nexum_sdk::host::{Fault, LocalStoreHost as _}; use nexum_sdk::keeper::{ - Gates, Journal, NEXT_BLOCK_PREFIX, NEXT_EPOCH_PREFIX, Poller, REFUSED_PREFIX, Retrier, - RetryAction, Tick, WATCH_PREFIX, WatchRef, WatchSet, watch_key, + Gates, Journal, Mark, NEXT_BLOCK_PREFIX, NEXT_EPOCH_PREFIX, Poller, REFUSED_PREFIX, + Reservation, Retrier, RetryAction, Tick, WATCH_PREFIX, WatchRef, WatchSet, watch_key, }; use nexum_sdk_test::MockHost; @@ -354,6 +354,162 @@ fn submitted_and_observed_keyspaces_are_disjoint() { assert!(!snapshot.contains_key("observed:0xuid")); } +// ---- durable-effect journal ---- + +#[test] +fn mark_distinguishes_reserved_committed_absent_and_legacy() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + + // Absent. + assert_eq!(journal.mark("absent").unwrap(), None); + + // Reserved. + journal.reserve("res", b"body").unwrap(); + assert_eq!(journal.mark("res").unwrap(), Some(Mark::Reserved)); + + // Committed. + journal.commit("com").unwrap(); + assert_eq!(journal.mark("com").unwrap(), Some(Mark::Committed)); + + // Legacy empty presence marker reads as Committed. + host.store.set("submitted:legacy", b"").unwrap(); + assert_eq!(journal.mark("legacy").unwrap(), Some(Mark::Committed)); + + // A corrupt single-byte tag reads as Reserved. + host.store.set("submitted:corrupt", b"\x7f").unwrap(); + assert_eq!(journal.mark("corrupt").unwrap(), Some(Mark::Reserved)); +} + +#[test] +fn reserve_then_commit_marks_committed() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + journal.reserve("uid", b"body").unwrap(); + journal.commit("uid").unwrap(); + assert_eq!(journal.mark("uid").unwrap(), Some(Mark::Committed)); +} + +#[test] +fn reserve_then_release_marks_absent() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + journal.reserve("uid", b"body").unwrap(); + journal.release("uid").unwrap(); + assert_eq!(journal.mark("uid").unwrap(), None); + assert!(host.store.is_empty()); +} + +#[test] +fn mark_and_pending_agree_on_every_non_committed_value() { + // The invariant: every value mark() calls Reserved, pending() must + // enumerate - including a corrupt single-byte tag with empty body. + // A guard-skipped-yet-never-reconciled marker is a dropped effect. + let host = MockHost::new(); + let journal = Journal::submitted(&host); + + journal.reserve("well_formed", b"body").unwrap(); + host.store.set("submitted:corrupt_tag", b"\x7f").unwrap(); + host.store.set("submitted:short_01", b"\x01").unwrap(); + + let listed: std::collections::BTreeSet = journal + .pending() + .unwrap() + .into_iter() + .map(|r| r.key) + .collect(); + + for key in ["well_formed", "corrupt_tag", "short_01"] { + assert_eq!( + journal.mark(key).unwrap(), + Some(Mark::Reserved), + "{key} must mark Reserved", + ); + assert!(listed.contains(key), "{key} must be enumerated"); + } + + // The corrupt tag enumerates with a zero next_eligible and empty body. + let corrupt = journal + .pending() + .unwrap() + .into_iter() + .find(|r| r.key == "corrupt_tag") + .unwrap(); + assert_eq!(corrupt.next_eligible, 0); + assert!(corrupt.body.is_empty()); +} + +#[test] +fn pending_ignores_committed_and_legacy_and_returns_bodies() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + + journal.park("a", b"aa", 1_700_000_000).unwrap(); + journal.reserve("b", b"bb").unwrap(); + journal.commit("c").unwrap(); + host.store.set("submitted:d", b"").unwrap(); + + let mut listed = journal.pending().unwrap(); + listed.sort_by(|x, y| x.key.cmp(&y.key)); + + assert_eq!( + listed, + vec![ + Reservation { + key: "a".into(), + next_eligible: 1_700_000_000, + body: b"aa".to_vec(), + }, + Reservation { + key: "b".into(), + next_eligible: 0, + body: b"bb".to_vec(), + }, + ], + ); +} + +#[test] +fn park_updates_next_eligible_without_touching_body() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + journal.reserve("uid", b"body").unwrap(); + journal.park("uid", b"body", 1_700_000_000).unwrap(); + + let listed = journal.pending().unwrap(); + assert_eq!( + listed, + vec![Reservation { + key: "uid".into(), + next_eligible: 1_700_000_000, + body: b"body".to_vec(), + }], + ); +} + +#[test] +fn double_commit_is_idempotent() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + journal.reserve("uid", b"body").unwrap(); + journal.commit("uid").unwrap(); + journal.commit("uid").unwrap(); + assert_eq!(journal.mark("uid").unwrap(), Some(Mark::Committed)); + assert_eq!(host.store.len(), 1); + assert_eq!( + host.store.snapshot().get("submitted:uid").unwrap(), + &vec![0x02], + ); +} + +#[test] +fn release_of_absent_is_a_no_op() { + let host = MockHost::new(); + let journal = Journal::submitted(&host); + journal.release("nope").unwrap(); + assert!(host.store.is_empty()); +} + // ---- retry ledger ---- fn seeded_watch(host: &MockHost) -> String {