lib

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

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 }