lib

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

journal.rs (19017B)


      1 use futures_executor::block_on;
      2 use radroots_event::EventId;
      3 use radroots_protocol::runtime::v1::OperationId;
      4 use radroots_storage::{
      5     Error,
      6     journal::{
      7         CancellationState, IdempotencyDigest, IdempotencyKey, Journal, JournalRevision,
      8         JournalStage, JournalState, JournalTransition, OperationInstanceId, OperationRecord,
      9         PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX,
     10         RecoveryPoint, RecoveryReason, RecoveryRecord,
     11     },
     12 };
     13 use radroots_transport::BoxFuture;
     14 use std::sync::Mutex;
     15 
     16 struct MemoryJournal(Mutex<Vec<OperationRecord>>);
     17 
     18 impl MemoryJournal {
     19     fn new() -> Self {
     20         Self(Mutex::new(Vec::new()))
     21     }
     22 }
     23 
     24 impl Journal for MemoryJournal {
     25     fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> {
     26         Box::pin(async move {
     27             let mut records = self.0.lock().expect("test journal lock");
     28             if let Some(record) = records
     29                 .iter()
     30                 .find(|record| record.idempotency_key() == operation.idempotency_key())
     31             {
     32                 if record.operation_id() != operation.operation_id()
     33                     || record.input_digest() != operation.input_digest()
     34                     || record.instance_id() != operation.instance_id()
     35                 {
     36                     return Err(Error::IdempotencyConflict);
     37                 }
     38                 return Ok(PrepareReceipt::new(
     39                     PrepareDisposition::Replay,
     40                     record.clone(),
     41                 ));
     42             }
     43             if records
     44                 .iter()
     45                 .any(|record| record.instance_id() == operation.instance_id())
     46             {
     47                 return Err(Error::OperationIdentityMismatch);
     48             }
     49             let record = operation.into_record()?;
     50             records.push(record.clone());
     51             Ok(PrepareReceipt::new(PrepareDisposition::Created, record))
     52         })
     53     }
     54 
     55     fn operation(
     56         &self,
     57         instance_id: OperationInstanceId,
     58     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
     59         Box::pin(async move {
     60             Ok(self
     61                 .0
     62                 .lock()
     63                 .expect("test journal lock")
     64                 .iter()
     65                 .find(|record| record.instance_id() == instance_id)
     66                 .cloned())
     67         })
     68     }
     69 
     70     fn by_idempotency_key(
     71         &self,
     72         operation_id: OperationId,
     73         idempotency_key: IdempotencyKey,
     74     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
     75         Box::pin(async move {
     76             Ok(self
     77                 .0
     78                 .lock()
     79                 .expect("test journal lock")
     80                 .iter()
     81                 .find(|record| {
     82                     record.operation_id() == operation_id
     83                         && record.idempotency_key() == &idempotency_key
     84                 })
     85                 .cloned())
     86         })
     87     }
     88 
     89     fn transition(
     90         &self,
     91         transition: JournalTransition,
     92     ) -> BoxFuture<'_, Result<OperationRecord, Error>> {
     93         Box::pin(async move {
     94             let mut records = self.0.lock().expect("test journal lock");
     95             let record = records
     96                 .iter_mut()
     97                 .find(|record| record.instance_id() == transition.instance_id())
     98                 .ok_or(Error::OperationNotFound)?;
     99             let next = record.transition(&transition)?;
    100             *record = next.clone();
    101             Ok(next)
    102         })
    103     }
    104 
    105     fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> {
    106         Box::pin(async move {
    107             if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX {
    108                 return Err(Error::InvalidJournalQueryLimit);
    109             }
    110             Ok(self
    111                 .0
    112                 .lock()
    113                 .expect("test journal lock")
    114                 .iter()
    115                 .filter(|record| record.state().stage() == JournalStage::Recoverable)
    116                 .take(usize::from(limit))
    117                 .cloned()
    118                 .collect())
    119         })
    120     }
    121 }
    122 
    123 fn instance(byte: u8) -> OperationInstanceId {
    124     OperationInstanceId::new([byte; 16]).expect("operation instance")
    125 }
    126 
    127 fn key(byte: u8) -> IdempotencyKey {
    128     IdempotencyKey::parse(format!("sync-push-{byte:02x}")).expect("idempotency key")
    129 }
    130 
    131 fn event_id(byte: &str) -> EventId {
    132     EventId::parse(byte.repeat(64)).expect("event id")
    133 }
    134 
    135 fn prepare(instance_id: OperationInstanceId, digest: u8, at: u64) -> PrepareOperation {
    136     PrepareOperation::new(
    137         instance_id,
    138         OperationId::SyncPush,
    139         key(instance_id.as_bytes()[0]),
    140         IdempotencyDigest::new([digest; 32]),
    141         at,
    142     )
    143     .expect("prepare operation")
    144 }
    145 
    146 #[test]
    147 fn prepare_replays_exact_input_and_rejects_conflicts() {
    148     let journal = MemoryJournal::new();
    149     let dynamic: &dyn Journal = &journal;
    150     let operation = prepare(instance(1), 2, 100);
    151     let created = block_on(dynamic.prepare(operation.clone())).expect("created");
    152     assert_eq!(created.disposition(), PrepareDisposition::Created);
    153     assert_eq!(created.record().revision(), JournalRevision::INITIAL);
    154 
    155     let replay = block_on(journal.prepare(operation)).expect("replay");
    156     assert_eq!(replay.disposition(), PrepareDisposition::Replay);
    157     assert_eq!(replay.record(), created.record());
    158     assert_eq!(
    159         block_on(journal.prepare(prepare(instance(1), 3, 100))),
    160         Err(Error::IdempotencyConflict)
    161     );
    162     let conflicting_kind = PrepareOperation::new(
    163         instance(1),
    164         OperationId::FarmPublish,
    165         key(1),
    166         IdempotencyDigest::new([2; 32]),
    167         100,
    168     )
    169     .expect("conflicting operation kind");
    170     assert_eq!(
    171         block_on(journal.prepare(conflicting_kind)),
    172         Err(Error::IdempotencyConflict)
    173     );
    174 
    175     let debug = format!("{:?}", key(1));
    176     assert!(debug.contains("[REDACTED]"));
    177     assert!(!debug.contains("sync-push-01"));
    178 }
    179 
    180 #[test]
    181 fn lifecycle_is_optimistic_monotonic_and_commit_bound() {
    182     let journal = MemoryJournal::new();
    183     let instance_id = instance(4);
    184     let event_id = event_id("a");
    185     let prepared = block_on(journal.prepare(prepare(instance_id, 5, 100)))
    186         .expect("prepare")
    187         .record()
    188         .clone();
    189     let signed = block_on(journal.transition(JournalTransition::signed(
    190         instance_id,
    191         prepared.revision(),
    192         event_id,
    193     )))
    194     .expect("signed");
    195     assert_eq!(signed.state().stage(), JournalStage::Signed);
    196     assert_eq!(
    197         block_on(journal.transition(JournalTransition::signed(
    198             instance_id,
    199             prepared.revision(),
    200             event_id,
    201         ))),
    202         Err(Error::JournalRevisionConflict)
    203     );
    204 
    205     let committed = block_on(journal.transition(JournalTransition::committed(
    206         instance_id,
    207         signed.revision(),
    208         event_id,
    209         150,
    210     )))
    211     .expect("committed");
    212     assert_eq!(committed.state().stage(), JournalStage::Committed);
    213     assert_eq!(
    214         block_on(
    215             journal.transition(JournalTransition::recoverable(
    216                 instance_id,
    217                 committed.revision(),
    218                 RecoveryRecord::new(
    219                     RecoveryPoint::Signed { event_id },
    220                     RecoveryReason::TransportUnavailable,
    221                     1,
    222                     None,
    223                 )
    224                 .expect("recovery"),
    225             ))
    226         ),
    227         Err(Error::JournalOperationCommitted)
    228     );
    229 }
    230 
    231 #[test]
    232 fn cancellation_before_commit_recovers_and_after_commit_preserves_commit() {
    233     let journal = MemoryJournal::new();
    234     let before_id = instance(6);
    235     let prepared = block_on(journal.prepare(prepare(before_id, 7, 100)))
    236         .expect("prepare before")
    237         .record()
    238         .clone();
    239     let cancelled = block_on(journal.transition(JournalTransition::cancelled(
    240         before_id,
    241         prepared.revision(),
    242         110,
    243     )))
    244     .expect("cancel before commit");
    245     assert_eq!(cancelled.state().stage(), JournalStage::Recoverable);
    246     assert_eq!(
    247         cancelled.cancellation(),
    248         CancellationState::CancelledBeforeCommit
    249     );
    250     assert_eq!(
    251         block_on(journal.recoverable(10)).expect("recovery").len(),
    252         1
    253     );
    254 
    255     let resumed =
    256         block_on(journal.transition(JournalTransition::resume(before_id, cancelled.revision())))
    257             .expect("resume");
    258     assert_eq!(resumed.state(), &JournalState::Prepared);
    259     assert_eq!(resumed.cancellation(), CancellationState::NotRequested);
    260 
    261     let after_id = instance(8);
    262     let event_id = event_id("b");
    263     let prepared = block_on(journal.prepare(prepare(after_id, 9, 200)))
    264         .expect("prepare after")
    265         .record()
    266         .clone();
    267     let signed = block_on(journal.transition(JournalTransition::signed(
    268         after_id,
    269         prepared.revision(),
    270         event_id,
    271     )))
    272     .expect("signed after");
    273     let committed = block_on(journal.transition(JournalTransition::committed(
    274         after_id,
    275         signed.revision(),
    276         event_id,
    277         220,
    278     )))
    279     .expect("committed after");
    280     let observed = block_on(journal.transition(JournalTransition::cancelled(
    281         after_id,
    282         committed.revision(),
    283         230,
    284     )))
    285     .expect("cancel observed after commit");
    286     assert_eq!(observed.state().stage(), JournalStage::Committed);
    287     assert_eq!(
    288         observed.cancellation(),
    289         CancellationState::ObservedAfterCommit
    290     );
    291 }
    292 
    293 #[test]
    294 fn invalid_records_and_inputs_fail_closed() {
    295     assert_eq!(
    296         OperationInstanceId::new([0; 16]),
    297         Err(Error::InvalidOperationInstanceId)
    298     );
    299     assert_eq!(
    300         IdempotencyKey::parse(" bad-key"),
    301         Err(Error::InvalidIdempotencyKey)
    302     );
    303     assert_eq!(
    304         RecoveryRecord::new(
    305             RecoveryPoint::Prepared,
    306             RecoveryReason::Interrupted,
    307             0,
    308             None,
    309         ),
    310         Err(Error::InvalidRecoveryAttempt)
    311     );
    312 }
    313 
    314 #[test]
    315 fn journal_value_and_state_validation_matrix_is_complete() {
    316     let instance_id = instance(1);
    317     assert_eq!(instance_id.as_bytes(), &[1; 16]);
    318     for invalid in ["", " leading", "trailing ", "bad\nkey"] {
    319         assert_eq!(
    320             IdempotencyKey::parse(invalid),
    321             Err(Error::InvalidIdempotencyKey)
    322         );
    323     }
    324     assert_eq!(
    325         IdempotencyKey::parse("x".repeat(radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES + 1)),
    326         Err(Error::InvalidIdempotencyKey)
    327     );
    328     assert_eq!(
    329         IdempotencyKey::parse("x".repeat(4 * 1024 * 1024)),
    330         Err(Error::InvalidIdempotencyKey)
    331     );
    332     let maximum =
    333         IdempotencyKey::parse("x".repeat(radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES))
    334             .expect("exact maximum key");
    335     assert_eq!(
    336         maximum.as_str().len(),
    337         radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES
    338     );
    339     let idempotency_key = key(1);
    340     assert_eq!(idempotency_key.as_str(), "sync-push-01");
    341     let digest = IdempotencyDigest::new([2; 32]);
    342     assert_eq!(digest.as_bytes(), &[2; 32]);
    343     assert_eq!(JournalRevision::new(0), Err(Error::InvalidJournalRevision));
    344     assert_eq!(JournalRevision::new(2).unwrap().get(), 2);
    345     assert_eq!(
    346         RecoveryRecord::new(
    347             RecoveryPoint::Prepared,
    348             RecoveryReason::Interrupted,
    349             1,
    350             Some(0),
    351         ),
    352         Err(Error::InvalidRecoveryDeadline)
    353     );
    354     let recovery = RecoveryRecord::new(
    355         RecoveryPoint::Signed {
    356             event_id: event_id("a"),
    357         },
    358         RecoveryReason::TransportUnavailable,
    359         2,
    360         Some(150),
    361     )
    362     .unwrap();
    363     assert!(matches!(recovery.point(), RecoveryPoint::Signed { .. }));
    364     assert_eq!(recovery.reason(), RecoveryReason::TransportUnavailable);
    365     assert_eq!(recovery.attempt(), 2);
    366     assert_eq!(recovery.retry_not_before_unix_ms(), Some(150));
    367     for state in [
    368         JournalState::Prepared,
    369         JournalState::Signed {
    370             event_id: event_id("a"),
    371         },
    372         JournalState::Recoverable(recovery.clone()),
    373         JournalState::Committed {
    374             event_id: event_id("a"),
    375             committed_at_unix_ms: 150,
    376         },
    377     ] {
    378         let expected = match state {
    379             JournalState::Prepared => JournalStage::Prepared,
    380             JournalState::Signed { .. } => JournalStage::Signed,
    381             JournalState::Recoverable(_) => JournalStage::Recoverable,
    382             JournalState::Committed { .. } => JournalStage::Committed,
    383         };
    384         assert_eq!(state.stage(), expected);
    385     }
    386 
    387     let build = |state, cancellation| {
    388         OperationRecord::from_parts(
    389             instance_id,
    390             OperationId::SyncPush,
    391             idempotency_key.clone(),
    392             digest,
    393             100,
    394             JournalRevision::INITIAL,
    395             state,
    396             cancellation,
    397         )
    398     };
    399     assert_eq!(
    400         OperationRecord::from_parts(
    401             instance_id,
    402             OperationId::SyncPush,
    403             idempotency_key.clone(),
    404             digest,
    405             0,
    406             JournalRevision::INITIAL,
    407             JournalState::Prepared,
    408             CancellationState::NotRequested,
    409         ),
    410         Err(Error::InvalidOperationTimestamp)
    411     );
    412     for result in [
    413         build(
    414             JournalState::Committed {
    415                 event_id: event_id("a"),
    416                 committed_at_unix_ms: 0,
    417             },
    418             CancellationState::NotRequested,
    419         ),
    420         build(
    421             JournalState::Committed {
    422                 event_id: event_id("a"),
    423                 committed_at_unix_ms: 99,
    424             },
    425             CancellationState::NotRequested,
    426         ),
    427         build(
    428             JournalState::Recoverable(
    429                 RecoveryRecord::new(
    430                     RecoveryPoint::Prepared,
    431                     RecoveryReason::Interrupted,
    432                     1,
    433                     Some(99),
    434                 )
    435                 .unwrap(),
    436             ),
    437             CancellationState::NotRequested,
    438         ),
    439         build(
    440             JournalState::Committed {
    441                 event_id: event_id("a"),
    442                 committed_at_unix_ms: 100,
    443             },
    444             CancellationState::CancelledBeforeCommit,
    445         ),
    446         build(
    447             JournalState::Recoverable(recovery.clone()),
    448             CancellationState::ObservedAfterCommit,
    449         ),
    450         build(
    451             JournalState::Recoverable(recovery.clone()),
    452             CancellationState::CancelledBeforeCommit,
    453         ),
    454         build(
    455             JournalState::Prepared,
    456             CancellationState::ObservedAfterCommit,
    457         ),
    458         build(
    459             JournalState::Signed {
    460                 event_id: event_id("a"),
    461             },
    462             CancellationState::CancelledBeforeCommit,
    463         ),
    464     ] {
    465         assert_eq!(result, Err(Error::CorruptJournalRecord));
    466     }
    467 
    468     let operation = prepare(instance_id, 2, 100);
    469     assert_eq!(operation.instance_id(), instance_id);
    470     assert_eq!(operation.operation_id(), OperationId::SyncPush);
    471     assert_eq!(operation.idempotency_key(), &idempotency_key);
    472     assert_eq!(operation.input_digest(), digest);
    473     let record = operation.into_record().unwrap();
    474     assert_eq!(record.instance_id(), instance_id);
    475     assert_eq!(record.operation_id(), OperationId::SyncPush);
    476     assert_eq!(record.idempotency_key(), &idempotency_key);
    477     assert_eq!(record.input_digest(), digest);
    478     assert_eq!(record.prepared_at_unix_ms(), 100);
    479     assert_eq!(record.revision(), JournalRevision::INITIAL);
    480     assert_eq!(record.cancellation(), CancellationState::NotRequested);
    481     let receipt = PrepareReceipt::new(PrepareDisposition::Created, record.clone());
    482     assert_eq!(receipt.disposition(), PrepareDisposition::Created);
    483     assert_eq!(receipt.record(), &record);
    484 
    485     assert_eq!(
    486         record.transition(&JournalTransition::signed(
    487             instance(2),
    488             record.revision(),
    489             event_id("a"),
    490         )),
    491         Err(Error::OperationIdentityMismatch)
    492     );
    493     assert_eq!(
    494         record.transition(&JournalTransition::committed(
    495             instance_id,
    496             record.revision(),
    497             event_id("a"),
    498             100,
    499         )),
    500         Err(Error::InvalidJournalTransition)
    501     );
    502     assert_eq!(
    503         record.transition(&JournalTransition::cancelled(
    504             instance_id,
    505             record.revision(),
    506             99,
    507         )),
    508         Err(Error::InvalidJournalTransition)
    509     );
    510     let signed = record
    511         .transition(&JournalTransition::signed(
    512             instance_id,
    513             record.revision(),
    514             event_id("a"),
    515         ))
    516         .unwrap();
    517     assert_eq!(
    518         signed.transition(&JournalTransition::committed(
    519             instance_id,
    520             signed.revision(),
    521             event_id("b"),
    522             101,
    523         )),
    524         Err(Error::InvalidJournalTransition)
    525     );
    526     let recoverable = signed
    527         .transition(&JournalTransition::recoverable(
    528             instance_id,
    529             signed.revision(),
    530             recovery,
    531         ))
    532         .unwrap();
    533     let resumed = recoverable
    534         .transition(&JournalTransition::resume(
    535             instance_id,
    536             recoverable.revision(),
    537         ))
    538         .unwrap();
    539     assert!(matches!(resumed.state(), JournalState::Signed { .. }));
    540     assert_eq!(
    541         JournalTransition::resume(instance_id, resumed.revision()).instance_id(),
    542         instance_id
    543     );
    544 
    545     assert_eq!(
    546         PrepareOperation::new(
    547             instance(9),
    548             OperationId::SyncPush,
    549             key(9),
    550             IdempotencyDigest::new([9; 32]),
    551             0,
    552         ),
    553         Err(Error::InvalidOperationTimestamp)
    554     );
    555 
    556     let prepared = prepare(instance(10), 10, 100).into_record().unwrap();
    557     let signed = prepared
    558         .transition(&JournalTransition::signed(
    559             instance(10),
    560             prepared.revision(),
    561             event_id("c"),
    562         ))
    563         .unwrap();
    564     assert_eq!(
    565         signed.transition(&JournalTransition::committed(
    566             instance(10),
    567             signed.revision(),
    568             event_id("c"),
    569             99,
    570         )),
    571         Err(Error::InvalidJournalTransition)
    572     );
    573     let cancelled_signed = signed
    574         .transition(&JournalTransition::cancelled(
    575             instance(10),
    576             signed.revision(),
    577             100,
    578         ))
    579         .unwrap();
    580     assert_eq!(
    581         cancelled_signed.cancellation(),
    582         CancellationState::CancelledBeforeCommit
    583     );
    584 
    585     let cancellation_recovery = RecoveryRecord::new(
    586         RecoveryPoint::Prepared,
    587         RecoveryReason::CancelledBeforeCommit,
    588         1,
    589         None,
    590     )
    591     .unwrap();
    592     let recovered = prepared
    593         .transition(&JournalTransition::recoverable(
    594             instance(10),
    595             prepared.revision(),
    596             cancellation_recovery,
    597         ))
    598         .unwrap();
    599     assert_eq!(
    600         recovered.cancellation(),
    601         CancellationState::CancelledBeforeCommit
    602     );
    603 }