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 }