lib

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

journal.rs (20512B)


      1 //! Durable operation journal contracts.
      2 //!
      3 //! [`JournalState::Committed`] is the local durable commit point. Cancellation
      4 //! before that point records recoverable work and may be resumed with the same
      5 //! idempotency key. Cancellation observed after it never claims rollback:
      6 //! callers must receive committed state and may continue pending delivery.
      7 
      8 use core::fmt;
      9 pub use radroots_event::EventId;
     10 pub use radroots_protocol::runtime::v1::OperationId;
     11 pub use radroots_transport::BoxFuture;
     12 
     13 use crate::Error;
     14 
     15 /// Maximum UTF-8 bytes in a journal idempotency key.
     16 pub const IDEMPOTENCY_KEY_MAX_BYTES: usize = 256;
     17 /// Maximum records returned by one recoverable-work query.
     18 pub const RECOVERABLE_QUERY_LIMIT_MAX: u16 = 256;
     19 
     20 /// Host-generated identity for one durable operation execution.
     21 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     22 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     23 pub struct OperationInstanceId([u8; 16]);
     24 
     25 impl OperationInstanceId {
     26     pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
     27         if all_zero(&bytes) {
     28             return Err(Error::InvalidOperationInstanceId);
     29         }
     30         Ok(Self(bytes))
     31     }
     32 
     33     pub const fn as_bytes(&self) -> &[u8; 16] {
     34         &self.0
     35     }
     36 }
     37 
     38 const fn all_zero(bytes: &[u8; 16]) -> bool {
     39     let mut index = 0;
     40     while index < bytes.len() {
     41         if bytes[index] != 0 {
     42             return false;
     43         }
     44         index += 1;
     45     }
     46     true
     47 }
     48 
     49 /// Validated caller-owned idempotency key.
     50 #[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)]
     51 pub struct IdempotencyKey(String);
     52 
     53 impl IdempotencyKey {
     54     pub fn parse(value: impl AsRef<str>) -> Result<Self, Error> {
     55         let value = value.as_ref();
     56         if value.is_empty()
     57             || value.len() > IDEMPOTENCY_KEY_MAX_BYTES
     58             || value != value.trim()
     59             || value.chars().any(char::is_control)
     60         {
     61             return Err(Error::InvalidIdempotencyKey);
     62         }
     63         Ok(Self(value.to_owned()))
     64     }
     65 
     66     pub fn as_str(&self) -> &str {
     67         self.0.as_str()
     68     }
     69 }
     70 
     71 impl fmt::Debug for IdempotencyKey {
     72     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     73         formatter
     74             .debug_struct("IdempotencyKey")
     75             .field("value", &"[REDACTED]")
     76             .field("bytes", &self.0.len())
     77             .finish()
     78     }
     79 }
     80 
     81 #[cfg(feature = "serde")]
     82 impl serde::Serialize for IdempotencyKey {
     83     fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
     84     where
     85         S: serde::Serializer,
     86     {
     87         serializer.serialize_str(self.as_str())
     88     }
     89 }
     90 
     91 #[cfg(feature = "serde")]
     92 impl<'de> serde::Deserialize<'de> for IdempotencyKey {
     93     fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
     94     where
     95         D: serde::Deserializer<'de>,
     96     {
     97         let value = <String as serde::Deserialize>::deserialize(deserializer)?;
     98         Self::parse(value).map_err(serde::de::Error::custom)
     99     }
    100 }
    101 
    102 /// SHA-256 digest of canonical operation input, computed by its domain owner.
    103 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    104 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    105 pub struct IdempotencyDigest([u8; 32]);
    106 
    107 impl IdempotencyDigest {
    108     pub const fn new(bytes: [u8; 32]) -> Self {
    109         Self(bytes)
    110     }
    111 
    112     pub const fn as_bytes(&self) -> &[u8; 32] {
    113         &self.0
    114     }
    115 }
    116 
    117 /// Non-zero optimistic journal revision.
    118 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    119 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
    120 pub struct JournalRevision(u64);
    121 
    122 impl JournalRevision {
    123     pub const INITIAL: Self = Self(1);
    124 
    125     pub const fn new(value: u64) -> Result<Self, Error> {
    126         if value == 0 {
    127             return Err(Error::InvalidJournalRevision);
    128         }
    129         Ok(Self(value))
    130     }
    131 
    132     pub const fn get(self) -> u64 {
    133         self.0
    134     }
    135 
    136     fn next(self) -> Result<Self, Error> {
    137         self.0
    138             .checked_add(1)
    139             .map(Self)
    140             .ok_or(Error::CorruptJournalRecord)
    141     }
    142 }
    143 
    144 /// Durable lifecycle stage.
    145 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    146 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    147 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    148 pub enum JournalStage {
    149     Prepared,
    150     Signed,
    151     Recoverable,
    152     Committed,
    153 }
    154 
    155 /// Point from which recoverable work resumes.
    156 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    157 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    158 #[derive(Clone, Debug, Eq, PartialEq)]
    159 pub enum RecoveryPoint {
    160     Prepared,
    161     Signed { event_id: EventId },
    162 }
    163 
    164 /// Stable class of recoverable interruption.
    165 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    166 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    167 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    168 pub enum RecoveryReason {
    169     CancelledBeforeCommit,
    170     SignerUnavailable,
    171     TransportUnavailable,
    172     StorageUnavailable,
    173     DeadlineExceeded,
    174     Interrupted,
    175 }
    176 
    177 /// Durable recovery evidence without backend or secret detail.
    178 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    179 #[derive(Clone, Debug, Eq, PartialEq)]
    180 pub struct RecoveryRecord {
    181     point: RecoveryPoint,
    182     reason: RecoveryReason,
    183     attempt: u32,
    184     retry_not_before_unix_ms: Option<u64>,
    185 }
    186 
    187 impl RecoveryRecord {
    188     pub const fn new(
    189         point: RecoveryPoint,
    190         reason: RecoveryReason,
    191         attempt: u32,
    192         retry_not_before_unix_ms: Option<u64>,
    193     ) -> Result<Self, Error> {
    194         if attempt == 0 {
    195             return Err(Error::InvalidRecoveryAttempt);
    196         }
    197         if matches!(retry_not_before_unix_ms, Some(0)) {
    198             return Err(Error::InvalidRecoveryDeadline);
    199         }
    200         Ok(Self {
    201             point,
    202             reason,
    203             attempt,
    204             retry_not_before_unix_ms,
    205         })
    206     }
    207 
    208     pub const fn point(&self) -> &RecoveryPoint {
    209         &self.point
    210     }
    211 
    212     pub const fn reason(&self) -> RecoveryReason {
    213         self.reason
    214     }
    215 
    216     pub const fn attempt(&self) -> u32 {
    217         self.attempt
    218     }
    219 
    220     pub const fn retry_not_before_unix_ms(&self) -> Option<u64> {
    221         self.retry_not_before_unix_ms
    222     }
    223 }
    224 
    225 /// Cancellation observation relative to the local durable commit point.
    226 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    227 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    228 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    229 pub enum CancellationState {
    230     NotRequested,
    231     CancelledBeforeCommit,
    232     ObservedAfterCommit,
    233 }
    234 
    235 /// State-specific durable journal data.
    236 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    237 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    238 #[derive(Clone, Debug, Eq, PartialEq)]
    239 pub enum JournalState {
    240     Prepared,
    241     Signed {
    242         event_id: EventId,
    243     },
    244     Recoverable(RecoveryRecord),
    245     /// Local state is durably committed; remote delivery may remain pending.
    246     Committed {
    247         event_id: EventId,
    248         committed_at_unix_ms: u64,
    249     },
    250 }
    251 
    252 impl JournalState {
    253     pub const fn stage(&self) -> JournalStage {
    254         match self {
    255             Self::Prepared => JournalStage::Prepared,
    256             Self::Signed { .. } => JournalStage::Signed,
    257             Self::Recoverable(_) => JournalStage::Recoverable,
    258             Self::Committed { .. } => JournalStage::Committed,
    259         }
    260     }
    261 }
    262 
    263 /// Durable operation record.
    264 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    265 #[derive(Clone, Debug, Eq, PartialEq)]
    266 pub struct OperationRecord {
    267     instance_id: OperationInstanceId,
    268     operation_id: OperationId,
    269     idempotency_key: IdempotencyKey,
    270     input_digest: IdempotencyDigest,
    271     prepared_at_unix_ms: u64,
    272     revision: JournalRevision,
    273     state: JournalState,
    274     cancellation: CancellationState,
    275 }
    276 
    277 impl OperationRecord {
    278     #[allow(clippy::too_many_arguments)]
    279     pub fn from_parts(
    280         instance_id: OperationInstanceId,
    281         operation_id: OperationId,
    282         idempotency_key: IdempotencyKey,
    283         input_digest: IdempotencyDigest,
    284         prepared_at_unix_ms: u64,
    285         revision: JournalRevision,
    286         state: JournalState,
    287         cancellation: CancellationState,
    288     ) -> Result<Self, Error> {
    289         if prepared_at_unix_ms == 0 {
    290             return Err(Error::InvalidOperationTimestamp);
    291         }
    292         validate_state(prepared_at_unix_ms, &state, cancellation)?;
    293         Ok(Self {
    294             instance_id,
    295             operation_id,
    296             idempotency_key,
    297             input_digest,
    298             prepared_at_unix_ms,
    299             revision,
    300             state,
    301             cancellation,
    302         })
    303     }
    304 
    305     pub const fn instance_id(&self) -> OperationInstanceId {
    306         self.instance_id
    307     }
    308 
    309     pub const fn operation_id(&self) -> OperationId {
    310         self.operation_id
    311     }
    312 
    313     pub const fn idempotency_key(&self) -> &IdempotencyKey {
    314         &self.idempotency_key
    315     }
    316 
    317     pub const fn input_digest(&self) -> IdempotencyDigest {
    318         self.input_digest
    319     }
    320 
    321     pub const fn prepared_at_unix_ms(&self) -> u64 {
    322         self.prepared_at_unix_ms
    323     }
    324 
    325     pub const fn revision(&self) -> JournalRevision {
    326         self.revision
    327     }
    328 
    329     pub const fn state(&self) -> &JournalState {
    330         &self.state
    331     }
    332 
    333     pub const fn cancellation(&self) -> CancellationState {
    334         self.cancellation
    335     }
    336 
    337     /// Applies one optimistic transition without allowing lifecycle regressions.
    338     pub fn transition(&self, transition: &JournalTransition) -> Result<Self, Error> {
    339         if transition.instance_id != self.instance_id {
    340             return Err(Error::OperationIdentityMismatch);
    341         }
    342         if transition.expected_revision != self.revision {
    343             return Err(Error::JournalRevisionConflict);
    344         }
    345 
    346         let (state, cancellation) = apply_transition(self, &transition.kind)?;
    347         Self::from_parts(
    348             self.instance_id,
    349             self.operation_id,
    350             self.idempotency_key.clone(),
    351             self.input_digest,
    352             self.prepared_at_unix_ms,
    353             self.revision.next()?,
    354             state,
    355             cancellation,
    356         )
    357     }
    358 }
    359 
    360 fn validate_state(
    361     prepared_at: u64,
    362     state: &JournalState,
    363     cancellation: CancellationState,
    364 ) -> Result<(), Error> {
    365     if let JournalState::Committed {
    366         committed_at_unix_ms,
    367         ..
    368     } = state
    369         && (*committed_at_unix_ms == 0 || *committed_at_unix_ms < prepared_at)
    370     {
    371         return Err(Error::CorruptJournalRecord);
    372     }
    373     if let JournalState::Recoverable(recovery) = state
    374         && recovery
    375             .retry_not_before_unix_ms()
    376             .is_some_and(|deadline| deadline < prepared_at)
    377     {
    378         return Err(Error::CorruptJournalRecord);
    379     }
    380     match state {
    381         JournalState::Committed { .. }
    382             if cancellation == CancellationState::CancelledBeforeCommit =>
    383         {
    384             Err(Error::CorruptJournalRecord)
    385         }
    386         JournalState::Recoverable(_) if cancellation == CancellationState::ObservedAfterCommit => {
    387             Err(Error::CorruptJournalRecord)
    388         }
    389         JournalState::Recoverable(recovery)
    390             if (recovery.reason() == RecoveryReason::CancelledBeforeCommit)
    391                 != (cancellation == CancellationState::CancelledBeforeCommit) =>
    392         {
    393             Err(Error::CorruptJournalRecord)
    394         }
    395         JournalState::Prepared | JournalState::Signed { .. }
    396             if cancellation != CancellationState::NotRequested =>
    397         {
    398             Err(Error::CorruptJournalRecord)
    399         }
    400         _ => Ok(()),
    401     }
    402 }
    403 
    404 fn apply_transition(
    405     record: &OperationRecord,
    406     transition: &JournalTransitionKind,
    407 ) -> Result<(JournalState, CancellationState), Error> {
    408     match (record.state(), transition) {
    409         (JournalState::Prepared, JournalTransitionKind::Signed { event_id })
    410             if record.cancellation() == CancellationState::NotRequested =>
    411         {
    412             Ok((
    413                 JournalState::Signed {
    414                     event_id: *event_id,
    415                 },
    416                 CancellationState::NotRequested,
    417             ))
    418         }
    419         (JournalState::Committed { .. }, JournalTransitionKind::Cancelled { observed_at })
    420             if *observed_at >= record.prepared_at_unix_ms() =>
    421         {
    422             Ok((
    423                 record.state().clone(),
    424                 CancellationState::ObservedAfterCommit,
    425             ))
    426         }
    427         (JournalState::Committed { .. }, _) => Err(Error::JournalOperationCommitted),
    428         (_, JournalTransitionKind::Recoverable { record: recovery }) => Ok((
    429             JournalState::Recoverable(recovery.clone()),
    430             if recovery.reason() == RecoveryReason::CancelledBeforeCommit {
    431                 CancellationState::CancelledBeforeCommit
    432             } else {
    433                 CancellationState::NotRequested
    434             },
    435         )),
    436         (JournalState::Recoverable(recovery), JournalTransitionKind::Resume) => Ok((
    437             match recovery.point() {
    438                 RecoveryPoint::Prepared => JournalState::Prepared,
    439                 RecoveryPoint::Signed { event_id } => JournalState::Signed {
    440                     event_id: *event_id,
    441                 },
    442             },
    443             CancellationState::NotRequested,
    444         )),
    445         (
    446             JournalState::Signed {
    447                 event_id: signed_id,
    448             },
    449             JournalTransitionKind::Committed {
    450                 event_id,
    451                 committed_at,
    452             },
    453         ) if signed_id == event_id && *committed_at >= record.prepared_at_unix_ms() => Ok((
    454             JournalState::Committed {
    455                 event_id: *event_id,
    456                 committed_at_unix_ms: *committed_at,
    457             },
    458             CancellationState::NotRequested,
    459         )),
    460         (JournalState::Prepared, JournalTransitionKind::Cancelled { observed_at })
    461             if *observed_at >= record.prepared_at_unix_ms() =>
    462         {
    463             cancelled(RecoveryPoint::Prepared)
    464         }
    465         (JournalState::Signed { event_id }, JournalTransitionKind::Cancelled { observed_at })
    466             if *observed_at >= record.prepared_at_unix_ms() =>
    467         {
    468             cancelled(RecoveryPoint::Signed {
    469                 event_id: *event_id,
    470             })
    471         }
    472         _ => Err(Error::InvalidJournalTransition),
    473     }
    474 }
    475 
    476 fn cancelled(point: RecoveryPoint) -> Result<(JournalState, CancellationState), Error> {
    477     Ok((
    478         JournalState::Recoverable(RecoveryRecord::new(
    479             point,
    480             RecoveryReason::CancelledBeforeCommit,
    481             1,
    482             None,
    483         )?),
    484         CancellationState::CancelledBeforeCommit,
    485     ))
    486 }
    487 
    488 /// Input for an idempotent prepare operation.
    489 #[derive(Clone, Debug, Eq, PartialEq)]
    490 pub struct PrepareOperation {
    491     instance_id: OperationInstanceId,
    492     operation_id: OperationId,
    493     idempotency_key: IdempotencyKey,
    494     input_digest: IdempotencyDigest,
    495     prepared_at_unix_ms: u64,
    496 }
    497 
    498 impl PrepareOperation {
    499     pub fn new(
    500         instance_id: OperationInstanceId,
    501         operation_id: OperationId,
    502         idempotency_key: IdempotencyKey,
    503         input_digest: IdempotencyDigest,
    504         prepared_at_unix_ms: u64,
    505     ) -> Result<Self, Error> {
    506         if prepared_at_unix_ms == 0 {
    507             return Err(Error::InvalidOperationTimestamp);
    508         }
    509         Ok(Self {
    510             instance_id,
    511             operation_id,
    512             idempotency_key,
    513             input_digest,
    514             prepared_at_unix_ms,
    515         })
    516     }
    517 
    518     pub const fn instance_id(&self) -> OperationInstanceId {
    519         self.instance_id
    520     }
    521 
    522     pub const fn operation_id(&self) -> OperationId {
    523         self.operation_id
    524     }
    525 
    526     pub const fn idempotency_key(&self) -> &IdempotencyKey {
    527         &self.idempotency_key
    528     }
    529 
    530     pub const fn input_digest(&self) -> IdempotencyDigest {
    531         self.input_digest
    532     }
    533 
    534     pub fn into_record(self) -> Result<OperationRecord, Error> {
    535         OperationRecord::from_parts(
    536             self.instance_id,
    537             self.operation_id,
    538             self.idempotency_key,
    539             self.input_digest,
    540             self.prepared_at_unix_ms,
    541             JournalRevision::INITIAL,
    542             JournalState::Prepared,
    543             CancellationState::NotRequested,
    544         )
    545     }
    546 }
    547 
    548 /// Result of preparing an idempotent operation.
    549 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    550 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    551 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    552 pub enum PrepareDisposition {
    553     Created,
    554     Replay,
    555 }
    556 
    557 #[derive(Clone, Debug, Eq, PartialEq)]
    558 pub struct PrepareReceipt {
    559     disposition: PrepareDisposition,
    560     record: OperationRecord,
    561 }
    562 
    563 impl PrepareReceipt {
    564     pub const fn new(disposition: PrepareDisposition, record: OperationRecord) -> Self {
    565         Self {
    566             disposition,
    567             record,
    568         }
    569     }
    570 
    571     pub const fn disposition(&self) -> PrepareDisposition {
    572         self.disposition
    573     }
    574 
    575     pub const fn record(&self) -> &OperationRecord {
    576         &self.record
    577     }
    578 }
    579 
    580 /// Validated optimistic transition request.
    581 #[derive(Clone, Debug, Eq, PartialEq)]
    582 pub struct JournalTransition {
    583     instance_id: OperationInstanceId,
    584     expected_revision: JournalRevision,
    585     kind: JournalTransitionKind,
    586 }
    587 
    588 #[derive(Clone, Debug, Eq, PartialEq)]
    589 enum JournalTransitionKind {
    590     Signed {
    591         event_id: EventId,
    592     },
    593     Recoverable {
    594         record: RecoveryRecord,
    595     },
    596     Resume,
    597     Committed {
    598         event_id: EventId,
    599         committed_at: u64,
    600     },
    601     Cancelled {
    602         observed_at: u64,
    603     },
    604 }
    605 
    606 impl JournalTransition {
    607     pub const fn signed(
    608         instance_id: OperationInstanceId,
    609         expected_revision: JournalRevision,
    610         event_id: EventId,
    611     ) -> Self {
    612         Self {
    613             instance_id,
    614             expected_revision,
    615             kind: JournalTransitionKind::Signed { event_id },
    616         }
    617     }
    618 
    619     pub const fn recoverable(
    620         instance_id: OperationInstanceId,
    621         expected_revision: JournalRevision,
    622         record: RecoveryRecord,
    623     ) -> Self {
    624         Self {
    625             instance_id,
    626             expected_revision,
    627             kind: JournalTransitionKind::Recoverable { record },
    628         }
    629     }
    630 
    631     pub const fn resume(
    632         instance_id: OperationInstanceId,
    633         expected_revision: JournalRevision,
    634     ) -> Self {
    635         Self {
    636             instance_id,
    637             expected_revision,
    638             kind: JournalTransitionKind::Resume,
    639         }
    640     }
    641 
    642     pub const fn committed(
    643         instance_id: OperationInstanceId,
    644         expected_revision: JournalRevision,
    645         event_id: EventId,
    646         committed_at: u64,
    647     ) -> Self {
    648         Self {
    649             instance_id,
    650             expected_revision,
    651             kind: JournalTransitionKind::Committed {
    652                 event_id,
    653                 committed_at,
    654             },
    655         }
    656     }
    657 
    658     pub const fn cancelled(
    659         instance_id: OperationInstanceId,
    660         expected_revision: JournalRevision,
    661         observed_at: u64,
    662     ) -> Self {
    663         Self {
    664             instance_id,
    665             expected_revision,
    666             kind: JournalTransitionKind::Cancelled { observed_at },
    667         }
    668     }
    669 
    670     pub const fn instance_id(&self) -> OperationInstanceId {
    671         self.instance_id
    672     }
    673 }
    674 
    675 /// Backend-neutral durable operation journal SPI.
    676 pub trait Journal: Send + Sync {
    677     /// Creates a record or replays the exact existing operation. A reused key
    678     /// with a different operation kind, instance, or digest is a conflict.
    679     fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>>;
    680 
    681     fn operation(
    682         &self,
    683         instance_id: OperationInstanceId,
    684     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>;
    685 
    686     fn by_idempotency_key(
    687         &self,
    688         operation_id: OperationId,
    689         idempotency_key: IdempotencyKey,
    690     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>;
    691 
    692     /// Applies one lifecycle transition atomically at its expected revision.
    693     fn transition(
    694         &self,
    695         transition: JournalTransition,
    696     ) -> BoxFuture<'_, Result<OperationRecord, Error>>;
    697 
    698     /// Returns recoverable records in backend-stable order.
    699     fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>>;
    700 }