lib

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

atomic.rs (7117B)


      1 //! High-level all-or-nothing storage workflow contracts.
      2 //!
      3 //! [`AtomicStorage::commit`] is the local durable commit boundary. Dropping its
      4 //! future before that boundary must leave no partial mutation. Cancellation
      5 //! observed after a successful commit cannot claim rollback; replaying the same
      6 //! commit identity and digest returns the original receipt.
      7 
      8 use radroots_transport::BoxFuture;
      9 
     10 use crate::{
     11     Error,
     12     event::{AdmissionReceipt, EventAdmission},
     13     projection::{ProjectionCheckpoint, ProjectionStatus},
     14 };
     15 
     16 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     17 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     18 pub struct AtomicCommitId([u8; 16]);
     19 
     20 impl AtomicCommitId {
     21     pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
     22         if bytes_are_zero(&bytes) {
     23             return Err(Error::InvalidAtomicCommitId);
     24         }
     25         Ok(Self(bytes))
     26     }
     27     pub const fn as_bytes(&self) -> &[u8; 16] {
     28         &self.0
     29     }
     30 }
     31 
     32 /// Domain-owner digest of the complete canonical workflow input.
     33 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     34 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     35 pub struct AtomicCommitDigest([u8; 32]);
     36 
     37 impl AtomicCommitDigest {
     38     pub const fn new(bytes: [u8; 32]) -> Self {
     39         Self(bytes)
     40     }
     41     pub const fn as_bytes(&self) -> &[u8; 32] {
     42         &self.0
     43     }
     44 }
     45 
     46 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     47 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
     48 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     49 pub enum AtomicWorkflowKind {
     50     Ingested,
     51 }
     52 
     53 /// Atomic inbound admission with an optional projection checkpoint advance.
     54 #[derive(Clone, Debug, Eq, PartialEq)]
     55 pub struct CommitIngested {
     56     admission: EventAdmission,
     57     projection: Option<ProjectionCheckpoint>,
     58 }
     59 
     60 impl CommitIngested {
     61     pub const fn new(admission: EventAdmission, projection: Option<ProjectionCheckpoint>) -> Self {
     62         Self {
     63             admission,
     64             projection,
     65         }
     66     }
     67     pub const fn admission(&self) -> &EventAdmission {
     68         &self.admission
     69     }
     70     pub const fn projection(&self) -> Option<&ProjectionCheckpoint> {
     71         self.projection.as_ref()
     72     }
     73 }
     74 
     75 #[derive(Clone, Debug, Eq, PartialEq)]
     76 pub enum AtomicWorkflow {
     77     Ingested(Box<CommitIngested>),
     78 }
     79 
     80 impl AtomicWorkflow {
     81     pub const fn kind(&self) -> AtomicWorkflowKind {
     82         match self {
     83             Self::Ingested(_) => AtomicWorkflowKind::Ingested,
     84         }
     85     }
     86 }
     87 
     88 /// Idempotent request for one high-level durable workflow.
     89 #[derive(Clone, Debug, Eq, PartialEq)]
     90 pub struct AtomicCommit {
     91     commit_id: AtomicCommitId,
     92     digest: AtomicCommitDigest,
     93     requested_at_unix_ms: u64,
     94     workflow: AtomicWorkflow,
     95 }
     96 
     97 impl AtomicCommit {
     98     pub fn new(
     99         commit_id: AtomicCommitId,
    100         digest: AtomicCommitDigest,
    101         requested_at_unix_ms: u64,
    102         workflow: AtomicWorkflow,
    103     ) -> Result<Self, Error> {
    104         if requested_at_unix_ms == 0 {
    105             return Err(Error::InvalidAtomicCommitTimestamp);
    106         }
    107         Ok(Self {
    108             commit_id,
    109             digest,
    110             requested_at_unix_ms,
    111             workflow,
    112         })
    113     }
    114     pub const fn commit_id(&self) -> AtomicCommitId {
    115         self.commit_id
    116     }
    117     pub const fn digest(&self) -> AtomicCommitDigest {
    118         self.digest
    119     }
    120     pub const fn requested_at_unix_ms(&self) -> u64 {
    121         self.requested_at_unix_ms
    122     }
    123     pub const fn workflow(&self) -> &AtomicWorkflow {
    124         &self.workflow
    125     }
    126 }
    127 
    128 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    129 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    130 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    131 pub enum AtomicCommitDisposition {
    132     Committed,
    133     Replay,
    134 }
    135 
    136 #[derive(Clone, Debug, Eq, PartialEq)]
    137 pub enum AtomicCommitOutcome {
    138     Ingested {
    139         admission: AdmissionReceipt,
    140         projection: Option<Box<ProjectionStatus>>,
    141     },
    142 }
    143 
    144 impl AtomicCommitOutcome {
    145     pub const fn kind(&self) -> AtomicWorkflowKind {
    146         match self {
    147             Self::Ingested { .. } => AtomicWorkflowKind::Ingested,
    148         }
    149     }
    150 }
    151 
    152 #[derive(Clone, Debug, Eq, PartialEq)]
    153 pub struct AtomicCommitReceipt {
    154     commit_id: AtomicCommitId,
    155     digest: AtomicCommitDigest,
    156     disposition: AtomicCommitDisposition,
    157     committed_at_unix_ms: u64,
    158     outcome: AtomicCommitOutcome,
    159 }
    160 
    161 impl AtomicCommitReceipt {
    162     pub fn new(
    163         request: &AtomicCommit,
    164         disposition: AtomicCommitDisposition,
    165         committed_at_unix_ms: u64,
    166         outcome: AtomicCommitOutcome,
    167     ) -> Result<Self, Error> {
    168         if [
    169             committed_at_unix_ms < request.requested_at_unix_ms(),
    170             outcome.kind() != request.workflow().kind(),
    171         ]
    172         .contains(&true)
    173         {
    174             return Err(Error::AtomicWorkflowMismatch);
    175         }
    176         Ok(Self {
    177             commit_id: request.commit_id(),
    178             digest: request.digest(),
    179             disposition,
    180             committed_at_unix_ms,
    181             outcome,
    182         })
    183     }
    184 
    185     /// Reconstructs and validates a receipt at a durable backend boundary.
    186     pub fn from_durable_parts(
    187         commit_id: AtomicCommitId,
    188         digest: AtomicCommitDigest,
    189         disposition: AtomicCommitDisposition,
    190         requested_at_unix_ms: u64,
    191         committed_at_unix_ms: u64,
    192         workflow_kind: AtomicWorkflowKind,
    193         outcome: AtomicCommitOutcome,
    194     ) -> Result<Self, Error> {
    195         if requested_at_unix_ms == 0
    196             || committed_at_unix_ms < requested_at_unix_ms
    197             || outcome.kind() != workflow_kind
    198         {
    199             return Err(Error::AtomicWorkflowMismatch);
    200         }
    201         Ok(Self {
    202             commit_id,
    203             digest,
    204             disposition,
    205             committed_at_unix_ms,
    206             outcome,
    207         })
    208     }
    209     pub const fn commit_id(&self) -> AtomicCommitId {
    210         self.commit_id
    211     }
    212     pub const fn digest(&self) -> AtomicCommitDigest {
    213         self.digest
    214     }
    215     pub const fn disposition(&self) -> AtomicCommitDisposition {
    216         self.disposition
    217     }
    218     pub const fn committed_at_unix_ms(&self) -> u64 {
    219         self.committed_at_unix_ms
    220     }
    221     pub const fn outcome(&self) -> &AtomicCommitOutcome {
    222         &self.outcome
    223     }
    224 }
    225 
    226 /// Backend-neutral all-or-nothing workflow commit SPI.
    227 pub trait AtomicStorage: Send + Sync {
    228     /// Commits every mutation in the workflow or none of them. Exact replay
    229     /// returns the original outcome; identity reuse with another digest or kind
    230     /// is [`Error::AtomicCommitConflict`].
    231     fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>>;
    232     fn receipt(
    233         &self,
    234         commit_id: AtomicCommitId,
    235     ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>>;
    236 }
    237 
    238 const fn bytes_are_zero(bytes: &[u8; 16]) -> bool {
    239     let mut index = 0;
    240     while index < bytes.len() {
    241         if bytes[index] != 0 {
    242             return false;
    243         }
    244         index += 1;
    245     }
    246     true
    247 }