lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 0965a817a211042ebee91b0744a387f90042d091
parent 1d205a28895c148969ae55d4050d143a4cb49ef3
Author: triesap <tyson@radroots.org>
Date:   Sat,  1 Aug 2026 21:50:46 +0000

storage: define atomic workflow commit operations

- define prepared signed enqueued delivered and ingested workflows
- bind replay to opaque commit identities and canonical digests
- document all-or-nothing cancellation and durable commit behavior
- prove failure isolation exact replay conflict and dyn safety

Diffstat:
Mcrates/storage/src/atomic.rs | 334++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/storage/src/error.rs | 10++++++++++
Mcrates/storage/src/lib.rs | 1+
Acrates/storage/tests/atomic.rs | 160+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
4 files changed, 504 insertions(+), 1 deletion(-)

diff --git a/crates/storage/src/atomic.rs b/crates/storage/src/atomic.rs @@ -1 +1,333 @@ -//! Atomic storage workflow contracts. +//! High-level all-or-nothing storage workflow contracts. +//! +//! [`AtomicStorage::commit`] is the local durable commit boundary. Dropping its +//! future before that boundary must leave no partial mutation. Cancellation +//! observed after a successful commit cannot claim rollback; replaying the same +//! commit identity and digest returns the original receipt. + +use radroots_event::{EventId, SignedEvent}; +use radroots_transport::BoxFuture; + +use crate::{ + Error, + event::{AdmissionReceipt, EventAdmission}, + journal::{JournalRevision, OperationInstanceId, OperationRecord, PrepareOperation}, + outbox::{DeliveryAttemptEvidence, EnqueueOutboxItem, OutboxRecord}, + projection::{ProjectionCheckpoint, ProjectionStatus}, +}; + +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct AtomicCommitId([u8; 16]); + +impl AtomicCommitId { + pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { + if bytes_are_zero(&bytes) { + return Err(Error::InvalidAtomicCommitId); + } + Ok(Self(bytes)) + } + pub const fn as_bytes(&self) -> &[u8; 16] { + &self.0 + } +} + +/// Domain-owner digest of the complete canonical workflow input. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct AtomicCommitDigest([u8; 32]); + +impl AtomicCommitDigest { + pub const fn new(bytes: [u8; 32]) -> Self { + Self(bytes) + } + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AtomicWorkflowKind { + Prepared, + Signed, + Enqueued, + Delivered, + Ingested, +} + +/// Journal transition inputs for a newly signed operation. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CommitSigned { + instance_id: OperationInstanceId, + expected_revision: JournalRevision, + event: SignedEvent, +} + +impl CommitSigned { + pub fn new( + instance_id: OperationInstanceId, + expected_revision: JournalRevision, + event: SignedEvent, + ) -> Self { + Self { + instance_id, + expected_revision, + event, + } + } + pub const fn instance_id(&self) -> OperationInstanceId { + self.instance_id + } + pub const fn expected_revision(&self) -> JournalRevision { + self.expected_revision + } + pub const fn event(&self) -> &SignedEvent { + &self.event + } +} + +/// Atomic canonical admission, journal commit, and delivery-plan enqueue. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CommitEnqueued { + instance_id: OperationInstanceId, + expected_revision: JournalRevision, + admission: EventAdmission, + outbox: EnqueueOutboxItem, + committed_at_unix_ms: u64, +} + +impl CommitEnqueued { + pub fn new( + instance_id: OperationInstanceId, + expected_revision: JournalRevision, + admission: EventAdmission, + outbox: EnqueueOutboxItem, + committed_at_unix_ms: u64, + ) -> Result<Self, Error> { + if committed_at_unix_ms == 0 + || admission.event_id() != outbox.request().payload().event().id() + || outbox.operation_instance_id() != instance_id + { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(Self { + instance_id, + expected_revision, + admission, + outbox, + committed_at_unix_ms, + }) + } + pub const fn instance_id(&self) -> OperationInstanceId { + self.instance_id + } + pub const fn expected_revision(&self) -> JournalRevision { + self.expected_revision + } + pub const fn admission(&self) -> &EventAdmission { + &self.admission + } + pub const fn outbox(&self) -> &EnqueueOutboxItem { + &self.outbox + } + pub const fn committed_at_unix_ms(&self) -> u64 { + self.committed_at_unix_ms + } +} + +/// Atomic inbound admission with an optional projection checkpoint advance. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CommitIngested { + admission: EventAdmission, + projection: Option<ProjectionCheckpoint>, +} + +impl CommitIngested { + pub const fn new(admission: EventAdmission, projection: Option<ProjectionCheckpoint>) -> Self { + Self { + admission, + projection, + } + } + pub const fn admission(&self) -> &EventAdmission { + &self.admission + } + pub const fn projection(&self) -> Option<&ProjectionCheckpoint> { + self.projection.as_ref() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum AtomicWorkflow { + Prepared(PrepareOperation), + Signed(Box<CommitSigned>), + Enqueued(Box<CommitEnqueued>), + Delivered(Box<DeliveryAttemptEvidence>), + Ingested(Box<CommitIngested>), +} + +impl AtomicWorkflow { + pub const fn kind(&self) -> AtomicWorkflowKind { + match self { + Self::Prepared(_) => AtomicWorkflowKind::Prepared, + Self::Signed(_) => AtomicWorkflowKind::Signed, + Self::Enqueued(_) => AtomicWorkflowKind::Enqueued, + Self::Delivered(_) => AtomicWorkflowKind::Delivered, + Self::Ingested(_) => AtomicWorkflowKind::Ingested, + } + } +} + +/// Idempotent request for one high-level durable workflow. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AtomicCommit { + commit_id: AtomicCommitId, + digest: AtomicCommitDigest, + requested_at_unix_ms: u64, + workflow: AtomicWorkflow, +} + +impl AtomicCommit { + pub fn new( + commit_id: AtomicCommitId, + digest: AtomicCommitDigest, + requested_at_unix_ms: u64, + workflow: AtomicWorkflow, + ) -> Result<Self, Error> { + if requested_at_unix_ms == 0 { + return Err(Error::InvalidAtomicCommitTimestamp); + } + Ok(Self { + commit_id, + digest, + requested_at_unix_ms, + workflow, + }) + } + pub const fn commit_id(&self) -> AtomicCommitId { + self.commit_id + } + pub const fn digest(&self) -> AtomicCommitDigest { + self.digest + } + pub const fn requested_at_unix_ms(&self) -> u64 { + self.requested_at_unix_ms + } + pub const fn workflow(&self) -> &AtomicWorkflow { + &self.workflow + } +} + +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AtomicCommitDisposition { + Committed, + Replay, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum AtomicCommitOutcome { + Prepared { + journal: OperationRecord, + }, + Signed { + journal: OperationRecord, + event_id: EventId, + }, + Enqueued { + journal: OperationRecord, + admission: AdmissionReceipt, + outbox: Box<OutboxRecord>, + }, + Delivered { + outbox: Box<OutboxRecord>, + }, + Ingested { + admission: AdmissionReceipt, + projection: Option<Box<ProjectionStatus>>, + }, +} + +impl AtomicCommitOutcome { + pub const fn kind(&self) -> AtomicWorkflowKind { + match self { + Self::Prepared { .. } => AtomicWorkflowKind::Prepared, + Self::Signed { .. } => AtomicWorkflowKind::Signed, + Self::Enqueued { .. } => AtomicWorkflowKind::Enqueued, + Self::Delivered { .. } => AtomicWorkflowKind::Delivered, + Self::Ingested { .. } => AtomicWorkflowKind::Ingested, + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AtomicCommitReceipt { + commit_id: AtomicCommitId, + digest: AtomicCommitDigest, + disposition: AtomicCommitDisposition, + committed_at_unix_ms: u64, + outcome: AtomicCommitOutcome, +} + +impl AtomicCommitReceipt { + pub fn new( + request: &AtomicCommit, + disposition: AtomicCommitDisposition, + committed_at_unix_ms: u64, + outcome: AtomicCommitOutcome, + ) -> Result<Self, Error> { + if committed_at_unix_ms < request.requested_at_unix_ms() + || outcome.kind() != request.workflow().kind() + { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(Self { + commit_id: request.commit_id(), + digest: request.digest(), + disposition, + committed_at_unix_ms, + outcome, + }) + } + pub const fn commit_id(&self) -> AtomicCommitId { + self.commit_id + } + pub const fn digest(&self) -> AtomicCommitDigest { + self.digest + } + pub const fn disposition(&self) -> AtomicCommitDisposition { + self.disposition + } + pub const fn committed_at_unix_ms(&self) -> u64 { + self.committed_at_unix_ms + } + pub const fn outcome(&self) -> &AtomicCommitOutcome { + &self.outcome + } +} + +/// Backend-neutral all-or-nothing workflow commit SPI. +pub trait AtomicStorage: Send + Sync { + /// Commits every mutation in the workflow or none of them. Exact replay + /// returns the original outcome; identity reuse with another digest or kind + /// is [`Error::AtomicCommitConflict`]. + fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>>; + fn receipt( + &self, + commit_id: AtomicCommitId, + ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>>; +} + +const fn bytes_are_zero(bytes: &[u8; 16]) -> bool { + let mut index = 0; + while index < bytes.len() { + if bytes[index] != 0 { + return false; + } + index += 1; + } + true +} diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs @@ -107,6 +107,11 @@ pub enum Error { CorruptReliabilityOperation, InvalidIntegrityStatus, InvalidStorageStatus, + InvalidAtomicCommitId, + InvalidAtomicCommitTimestamp, + AtomicCommitConflict, + AtomicWorkflowMismatch, + AtomicCommitFailed, } impl fmt::Display for Error { @@ -233,6 +238,11 @@ impl fmt::Display for Error { Self::CorruptReliabilityOperation => "storage reliability operation is corrupt", Self::InvalidIntegrityStatus => "storage integrity status is invalid", Self::InvalidStorageStatus => "storage status is invalid", + Self::InvalidAtomicCommitId => "storage atomic commit id is invalid", + Self::InvalidAtomicCommitTimestamp => "storage atomic commit timestamp is invalid", + Self::AtomicCommitConflict => "storage atomic commit conflicts with durable state", + Self::AtomicWorkflowMismatch => "storage atomic workflow components do not match", + Self::AtomicCommitFailed => "storage atomic workflow did not commit", }) } } diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs @@ -14,6 +14,7 @@ pub mod private_artifact; pub mod projection; pub mod status; +pub use atomic::AtomicStorage; pub use backup::StorageReliability; pub use error::Error; pub use event::EventStore; diff --git a/crates/storage/tests/atomic.rs b/crates/storage/tests/atomic.rs @@ -0,0 +1,160 @@ +use futures_executor::block_on; +use radroots_protocol::runtime::v1::OperationId; +use radroots_storage::{ + AtomicStorage, Error, + atomic::{ + AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, + AtomicCommitOutcome, AtomicCommitReceipt, AtomicWorkflow, + }, + journal::{IdempotencyDigest, IdempotencyKey, OperationInstanceId, PrepareOperation}, +}; +use radroots_transport::BoxFuture; +use std::{ + collections::BTreeMap, + sync::{ + Mutex, + atomic::{AtomicBool, Ordering}, + }, +}; + +struct ReferenceAtomic { + receipts: Mutex<BTreeMap<AtomicCommitId, AtomicCommitReceipt>>, + fail_before_commit: AtomicBool, +} + +impl ReferenceAtomic { + fn new() -> Self { + Self { + receipts: Mutex::new(BTreeMap::new()), + fail_before_commit: AtomicBool::new(false), + } + } +} + +impl AtomicStorage for ReferenceAtomic { + fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> { + Box::pin(async move { + let mut receipts = self.receipts.lock().expect("atomic test lock"); + if let Some(existing) = receipts.get(&request.commit_id()) { + if existing.digest() != request.digest() + || existing.outcome().kind() != request.workflow().kind() + { + return Err(Error::AtomicCommitConflict); + } + return AtomicCommitReceipt::new( + &request, + AtomicCommitDisposition::Replay, + existing.committed_at_unix_ms(), + existing.outcome().clone(), + ); + } + if self.fail_before_commit.swap(false, Ordering::SeqCst) { + return Err(Error::AtomicCommitFailed); + } + let outcome = match request.workflow().clone() { + AtomicWorkflow::Prepared(operation) => AtomicCommitOutcome::Prepared { + journal: operation.into_record()?, + }, + _ => return Err(Error::AtomicCommitFailed), + }; + let receipt = AtomicCommitReceipt::new( + &request, + AtomicCommitDisposition::Committed, + request.requested_at_unix_ms(), + outcome, + )?; + receipts.insert(request.commit_id(), receipt.clone()); + Ok(receipt) + }) + } + + fn receipt( + &self, + commit_id: AtomicCommitId, + ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>> { + Box::pin(async move { + Ok(self + .receipts + .lock() + .expect("atomic test lock") + .get(&commit_id) + .cloned()) + }) + } +} + +fn request(commit_byte: u8, digest_byte: u8) -> AtomicCommit { + let prepare = PrepareOperation::new( + OperationInstanceId::new([9; 16]).expect("operation instance"), + OperationId::TradePrivateArtifactSeal, + IdempotencyKey::parse("atomic-test-key").expect("idempotency key"), + IdempotencyDigest::new([3; 32]), + 100, + ) + .expect("prepare operation"); + AtomicCommit::new( + AtomicCommitId::new([commit_byte; 16]).expect("commit id"), + AtomicCommitDigest::new([digest_byte; 32]), + 100, + AtomicWorkflow::Prepared(prepare), + ) + .expect("atomic commit") +} + +#[test] +fn failure_before_commit_leaves_no_partial_receipt() { + let store = ReferenceAtomic::new(); + let request = request(1, 2); + store.fail_before_commit.store(true, Ordering::SeqCst); + assert_eq!( + block_on(store.commit(request.clone())), + Err(Error::AtomicCommitFailed) + ); + assert_eq!( + block_on(store.receipt(request.commit_id())).expect("receipt query"), + None + ); + + let committed = block_on(store.commit(request.clone())).expect("commit"); + assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed); + assert!( + block_on(store.receipt(request.commit_id())) + .expect("receipt query") + .is_some() + ); +} + +#[test] +fn exact_commit_replays_and_digest_reuse_conflicts() { + let store = ReferenceAtomic::new(); + let original_request = request(1, 2); + let original = block_on(store.commit(original_request.clone())).expect("commit"); + let replay = block_on(store.commit(original_request)).expect("replay"); + assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); + assert_eq!(replay.outcome(), original.outcome()); + assert_eq!( + block_on(store.commit(request(1, 4))), + Err(Error::AtomicCommitConflict) + ); +} + +#[test] +fn atomic_contract_is_dyn_compatible_and_rejects_invalid_identity_and_time() { + fn accepts_dyn(_: &dyn AtomicStorage) {} + let store = ReferenceAtomic::new(); + accepts_dyn(&store); + assert_eq!( + AtomicCommitId::new([0; 16]), + Err(Error::InvalidAtomicCommitId) + ); + let valid = request(1, 2); + assert_eq!( + AtomicCommit::new( + valid.commit_id(), + valid.digest(), + 0, + valid.workflow().clone() + ), + Err(Error::InvalidAtomicCommitTimestamp) + ); +}