policy.rs (10659B)
1 //! Explicit clocks, identifiers, deadlines, and retry decisions. 2 3 use std::sync::Arc; 4 5 use radroots_signing::Signer; 6 use radroots_storage::{ 7 EventStore, Journal, Outbox, ProjectionStore, atomic::AtomicStorage, 8 authored_atomic::AuthoredAtomicStorage, status::StorageStatusProvider, 9 }; 10 use radroots_transport::{EventSink, EventSource}; 11 12 use crate::Engine; 13 14 const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000; 15 16 /// Exact backend-neutral storage capability required by sync orchestration. 17 pub trait SyncStorage: 18 EventStore 19 + Journal 20 + Outbox 21 + ProjectionStore 22 + AtomicStorage 23 + AuthoredAtomicStorage 24 + StorageStatusProvider 25 { 26 } 27 28 impl<T> SyncStorage for T where 29 T: EventStore 30 + Journal 31 + Outbox 32 + ProjectionStore 33 + AtomicStorage 34 + AuthoredAtomicStorage 35 + StorageStatusProvider 36 { 37 } 38 39 /// Sync operation class used for identity and deadline policy. 40 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 41 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 42 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 43 #[non_exhaustive] 44 pub enum OperationKind { 45 Ingest, 46 Projection, 47 Pull, 48 Sign, 49 Deliver, 50 } 51 52 /// Opaque host-generated identity for one synchronization operation. 53 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 54 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 55 pub struct SyncId([u8; 16]); 56 57 impl SyncId { 58 pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { 59 let mut index = 0; 60 while index < bytes.len() { 61 if bytes[index] != 0 { 62 return Ok(Self(bytes)); 63 } 64 index += 1; 65 } 66 Err(Error::InvalidSyncId) 67 } 68 69 pub const fn as_bytes(&self) -> &[u8; 16] { 70 &self.0 71 } 72 } 73 74 /// Host clock used instead of reading ambient time inside orchestration. 75 pub trait Clock: Send + Sync { 76 fn now_unix_ms(&self) -> Result<u64, Error>; 77 } 78 79 /// Host identity source used instead of ambient randomness or global counters. 80 pub trait IdSource: Send + Sync { 81 fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error>; 82 } 83 84 /// Bounded time budgets applied to individual orchestration calls. 85 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 86 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 87 pub struct DeadlinePolicy { 88 pull_timeout_ms: u64, 89 sign_timeout_ms: u64, 90 delivery_timeout_ms: u64, 91 } 92 93 impl DeadlinePolicy { 94 pub const fn new( 95 pull_timeout_ms: u64, 96 sign_timeout_ms: u64, 97 delivery_timeout_ms: u64, 98 ) -> Result<Self, Error> { 99 if !valid_timeout(pull_timeout_ms) 100 || !valid_timeout(sign_timeout_ms) 101 || !valid_timeout(delivery_timeout_ms) 102 { 103 return Err(Error::InvalidDeadlinePolicy); 104 } 105 Ok(Self { 106 pull_timeout_ms, 107 sign_timeout_ms, 108 delivery_timeout_ms, 109 }) 110 } 111 112 pub const fn timeout_ms(self, operation: OperationKind) -> u64 { 113 match operation { 114 // Ingest performs local verification and one atomic commit. It 115 // shares the inbound operation budget with pull orchestration. 116 OperationKind::Ingest => self.pull_timeout_ms, 117 OperationKind::Projection => self.pull_timeout_ms, 118 OperationKind::Pull => self.pull_timeout_ms, 119 OperationKind::Sign => self.sign_timeout_ms, 120 OperationKind::Deliver => self.delivery_timeout_ms, 121 } 122 } 123 124 pub fn deadline_unix_ms( 125 self, 126 operation: OperationKind, 127 now_unix_ms: u64, 128 ) -> Result<u64, Error> { 129 if now_unix_ms == 0 { 130 return Err(Error::ClockUnavailable); 131 } 132 now_unix_ms 133 .checked_add(self.timeout_ms(operation)) 134 .ok_or(Error::DeadlineOverflow) 135 } 136 } 137 138 const fn valid_timeout(value: u64) -> bool { 139 value != 0 && value <= MAX_OPERATION_TIMEOUT_MS 140 } 141 142 /// Builder for an [`Engine`] with explicit optional transport capabilities. 143 pub struct EngineBuilder { 144 storage: Arc<dyn SyncStorage>, 145 source: Option<Arc<dyn EventSource>>, 146 sink: Option<Arc<dyn EventSink>>, 147 signer: Option<Arc<dyn Signer>>, 148 clock: Arc<dyn Clock>, 149 ids: Arc<dyn IdSource>, 150 deadlines: DeadlinePolicy, 151 } 152 153 impl EngineBuilder { 154 pub(crate) fn new( 155 storage: Arc<dyn SyncStorage>, 156 clock: Arc<dyn Clock>, 157 ids: Arc<dyn IdSource>, 158 deadlines: DeadlinePolicy, 159 ) -> Self { 160 Self { 161 storage, 162 source: None, 163 sink: None, 164 signer: None, 165 clock, 166 ids, 167 deadlines, 168 } 169 } 170 171 #[must_use] 172 pub fn source(mut self, source: Arc<dyn EventSource>) -> Self { 173 self.source = Some(source); 174 self 175 } 176 177 #[must_use] 178 pub fn sink(mut self, sink: Arc<dyn EventSink>) -> Self { 179 self.sink = Some(sink); 180 self 181 } 182 183 #[must_use] 184 pub fn signer(mut self, signer: Arc<dyn Signer>) -> Self { 185 self.signer = Some(signer); 186 self 187 } 188 189 pub fn build(self) -> Result<Engine, Error> { 190 if self.signer.is_some() && self.sink.is_none() { 191 return Err(Error::SignerWithoutSink); 192 } 193 if self.source.is_none() && self.sink.is_none() { 194 return Err(Error::MissingTransportCapability); 195 } 196 Ok(Engine { 197 storage: self.storage, 198 source: self.source, 199 sink: self.sink, 200 signer: self.signer, 201 clock: self.clock, 202 ids: self.ids, 203 deadlines: self.deadlines, 204 }) 205 } 206 } 207 208 /// Sync composition and host-policy error. 209 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 210 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 211 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 212 #[non_exhaustive] 213 pub enum Error { 214 InvalidSyncId, 215 InvalidDeadlinePolicy, 216 ClockUnavailable, 217 DeadlineOverflow, 218 MissingTransportCapability, 219 SignerWithoutSink, 220 VerificationFailed, 221 PolicyRejected, 222 StorageConflict, 223 /// Capacity failure; original effects may already be durable and require reconciliation. 224 StorageSpaceInsufficient, 225 StorageFailed, 226 InvalidIngestReceipt, 227 InvalidPullRequest, 228 MissingSource, 229 InvalidSourcePage, 230 InvalidProjectionRequest, 231 ReducerFailed, 232 InvalidReducerOutput, 233 InvalidPushRequest, 234 MissingSigner, 235 SignerCapabilityUnavailable, 236 SignerFailed, 237 SignerDeadlineExceeded, 238 SigningCancelled, 239 SigningIndeterminate, 240 WorkClaimConflict, 241 InvalidSignerOutput, 242 AdmissionFailed, 243 InvalidDeliveryRequest, 244 DeliveryDeferred, 245 MissingSink, 246 InvalidStatusRequest, 247 } 248 249 impl core::fmt::Display for Error { 250 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { 251 formatter.write_str(match self { 252 Self::InvalidSyncId => "sync identity must not be all zero", 253 Self::InvalidDeadlinePolicy => "sync deadline policy is outside its bounds", 254 Self::ClockUnavailable => "sync clock did not provide a valid timestamp", 255 Self::DeadlineOverflow => "sync deadline overflowed", 256 Self::MissingTransportCapability => "sync engine requires a source or sink", 257 Self::SignerWithoutSink => "sync signer requires a sink", 258 Self::VerificationFailed => "sync event verification failed", 259 Self::PolicyRejected => "sync admission policy rejected the event", 260 Self::StorageConflict => "sync input conflicts with durable storage state", 261 Self::StorageSpaceInsufficient => "sync storage space is insufficient", 262 Self::StorageFailed => "sync storage operation failed", 263 Self::InvalidIngestReceipt => "sync storage returned an invalid ingest receipt", 264 Self::InvalidPullRequest => "sync pull request is outside its bounds", 265 Self::MissingSource => "sync engine has no event source", 266 Self::InvalidSourcePage => "sync source returned an invalid page", 267 Self::InvalidProjectionRequest => "sync projection request is invalid", 268 Self::ReducerFailed => "sync projection reducer failed", 269 Self::InvalidReducerOutput => "sync projection reducer returned invalid progress", 270 Self::InvalidPushRequest => "sync push request is invalid", 271 Self::MissingSigner => "sync engine has no signer", 272 Self::SignerCapabilityUnavailable => { 273 "sync signer did not declare one usable replay capability" 274 } 275 Self::SignerFailed => "sync signer did not produce an event", 276 Self::SignerDeadlineExceeded => "sync signer exceeded its deadline", 277 Self::SigningCancelled => "sync signing was durably cancelled", 278 Self::SigningIndeterminate => { 279 "sync signing may have produced a non-replayable remote effect" 280 } 281 Self::WorkClaimConflict => "sync authored work is claimed by another execution", 282 Self::InvalidSignerOutput => "sync signer output failed canonical verification", 283 Self::AdmissionFailed => "sync local admission did not complete", 284 Self::InvalidDeliveryRequest => "sync delivery request is invalid", 285 Self::DeliveryDeferred => "sync delivery retry is not yet eligible", 286 Self::MissingSink => "sync engine has no event sink", 287 Self::InvalidStatusRequest => "sync status request is invalid", 288 }) 289 } 290 } 291 292 impl std::error::Error for Error {} 293 294 #[cfg(feature = "serde")] 295 impl<'de> serde::Deserialize<'de> for SyncId { 296 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 297 where 298 D: serde::Deserializer<'de>, 299 { 300 let bytes = <[u8; 16] as serde::Deserialize>::deserialize(deserializer)?; 301 Self::new(bytes).map_err(serde::de::Error::custom) 302 } 303 } 304 305 #[cfg(feature = "serde")] 306 impl<'de> serde::Deserialize<'de> for DeadlinePolicy { 307 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 308 where 309 D: serde::Deserializer<'de>, 310 { 311 #[derive(serde::Deserialize)] 312 #[serde(deny_unknown_fields)] 313 struct Wire { 314 pull_timeout_ms: u64, 315 sign_timeout_ms: u64, 316 delivery_timeout_ms: u64, 317 } 318 319 let wire = Wire::deserialize(deserializer)?; 320 Self::new( 321 wire.pull_timeout_ms, 322 wire.sign_timeout_ms, 323 wire.delivery_timeout_ms, 324 ) 325 .map_err(serde::de::Error::custom) 326 } 327 }