lib

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

mod.rs (25709B)


      1 use crate::SqliteStorage;
      2 use crate::backend::map_backend;
      3 use radroots_storage::{
      4     Error, Journal,
      5     journal::{
      6         BoxFuture, CancellationState, EventId, IdempotencyDigest, IdempotencyKey, JournalRevision,
      7         JournalState, JournalTransition, OperationId, OperationInstanceId, OperationRecord,
      8         PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX,
      9         RecoveryPoint, RecoveryReason, RecoveryRecord,
     10     },
     11 };
     12 use sqlx::{Row, Sqlite};
     13 
     14 #[cfg_attr(coverage_nightly, coverage(off))]
     15 impl Journal for SqliteStorage {
     16     fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> {
     17         Box::pin(async move {
     18             self.require_journal_writer()?;
     19             let mut transaction = self
     20                 .pool()
     21                 .begin_with("BEGIN IMMEDIATE")
     22                 .await
     23                 .map_err(map_backend)?;
     24             let receipt = prepare_transaction(&mut transaction, operation).await?;
     25             transaction.commit().await.map_err(map_backend)?;
     26             Ok(receipt)
     27         })
     28     }
     29 
     30     fn operation(
     31         &self,
     32         instance_id: OperationInstanceId,
     33     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
     34         Box::pin(async move {
     35             sqlx::query("SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?")
     36                 .bind(instance_id.as_bytes().as_slice())
     37                 .fetch_optional(self.pool())
     38                 .await
     39                 .map_err(map_backend)?
     40                 .as_ref()
     41                 .map(decode_record)
     42                 .transpose()
     43         })
     44     }
     45 
     46     fn by_idempotency_key(
     47         &self,
     48         operation_id: OperationId,
     49         idempotency_key: IdempotencyKey,
     50     ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
     51         Box::pin(async move {
     52             sqlx::query(
     53                 "SELECT * FROM radroots_runtime_journal_operations
     54                  WHERE operation_id = ? AND idempotency_key = ?",
     55             )
     56             .bind(operation_id.as_str().as_bytes())
     57             .bind(idempotency_key.as_str())
     58             .fetch_optional(self.pool())
     59             .await
     60             .map_err(map_backend)?
     61             .as_ref()
     62             .map(decode_record)
     63             .transpose()
     64         })
     65     }
     66 
     67     fn transition(
     68         &self,
     69         transition: JournalTransition,
     70     ) -> BoxFuture<'_, Result<OperationRecord, Error>> {
     71         Box::pin(async move {
     72             self.require_journal_writer()?;
     73             let mut transaction = self
     74                 .pool()
     75                 .begin_with("BEGIN IMMEDIATE")
     76                 .await
     77                 .map_err(map_backend)?;
     78             let next = transition_transaction(&mut transaction, transition).await?;
     79             transaction.commit().await.map_err(map_backend)?;
     80             Ok(next)
     81         })
     82     }
     83 
     84     fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> {
     85         Box::pin(async move {
     86             if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX {
     87                 return Err(Error::InvalidJournalQueryLimit);
     88             }
     89             sqlx::query(
     90                 "SELECT * FROM radroots_runtime_journal_operations
     91                  WHERE stage = 'recoverable'
     92                  ORDER BY updated_at_unix_ms, instance_id LIMIT ?",
     93             )
     94             .bind(i64::from(limit))
     95             .fetch_all(self.pool())
     96             .await
     97             .map_err(map_backend)?
     98             .iter()
     99             .map(decode_record)
    100             .collect()
    101         })
    102     }
    103 }
    104 
    105 #[cfg_attr(coverage_nightly, coverage(off))]
    106 pub(crate) async fn prepare_transaction(
    107     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    108     operation: PrepareOperation,
    109 ) -> Result<PrepareReceipt, Error> {
    110     if let Some(row) = sqlx::query(
    111         "SELECT * FROM radroots_runtime_journal_operations
    112          WHERE idempotency_key = ?",
    113     )
    114     .bind(operation.idempotency_key().as_str())
    115     .fetch_optional(&mut **transaction)
    116     .await
    117     .map_err(map_backend)?
    118     {
    119         let record = decode_record(&row)?;
    120         if record.operation_id() != operation.operation_id()
    121             || record.input_digest() != operation.input_digest()
    122             || record.instance_id() != operation.instance_id()
    123         {
    124             return Err(Error::IdempotencyConflict);
    125         }
    126         return Ok(PrepareReceipt::new(PrepareDisposition::Replay, record));
    127     }
    128     if sqlx::query_scalar::<_, i64>(
    129         "SELECT 1 FROM radroots_runtime_journal_operations WHERE instance_id = ?",
    130     )
    131     .bind(operation.instance_id().as_bytes().as_slice())
    132     .fetch_optional(&mut **transaction)
    133     .await
    134     .map_err(map_backend)?
    135     .is_some()
    136     {
    137         return Err(Error::OperationIdentityMismatch);
    138     }
    139 
    140     let record = operation.into_record()?;
    141     insert_record(transaction, &record).await?;
    142     Ok(PrepareReceipt::new(PrepareDisposition::Created, record))
    143 }
    144 
    145 #[cfg_attr(coverage_nightly, coverage(off))]
    146 pub(crate) async fn transition_transaction(
    147     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    148     transition: JournalTransition,
    149 ) -> Result<OperationRecord, Error> {
    150     let row =
    151         sqlx::query("SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?")
    152             .bind(transition.instance_id().as_bytes().as_slice())
    153             .fetch_optional(&mut **transaction)
    154             .await
    155             .map_err(map_backend)?
    156             .ok_or(Error::OperationNotFound)?;
    157     let current = decode_record(&row)?;
    158     let next = current.transition(&transition)?;
    159     let (stage, event_id, recovery, committed_at) = encode_state(next.state());
    160     let result = sqlx::query(
    161         "UPDATE radroots_runtime_journal_operations SET
    162            revision = ?, stage = ?, event_id = ?, recovery_record = ?,
    163            cancellation_state = ?, committed_at_unix_ms = ?, updated_at_unix_ms = ?
    164          WHERE instance_id = ? AND revision = ?",
    165     )
    166     .bind(i64_from_u64(next.revision().get())?)
    167     .bind(stage)
    168     .bind(event_id)
    169     .bind(recovery)
    170     .bind(cancellation_name(next.cancellation()))
    171     .bind(committed_at.map(i64_from_u64).transpose()?)
    172     .bind(i64_from_u64(updated_at(&next))?)
    173     .bind(next.instance_id().as_bytes().as_slice())
    174     .bind(i64_from_u64(current.revision().get())?)
    175     .execute(&mut **transaction)
    176     .await
    177     .map_err(map_backend)?;
    178     if result.rows_affected() != 1 {
    179         return Err(Error::JournalRevisionConflict);
    180     }
    181     Ok(next)
    182 }
    183 
    184 impl SqliteStorage {
    185     fn require_journal_writer(&self) -> Result<(), Error> {
    186         if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
    187             return Err(Error::BackendUnavailable);
    188         }
    189         Ok(())
    190     }
    191 }
    192 
    193 #[cfg_attr(coverage_nightly, coverage(off))]
    194 async fn insert_record(
    195     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    196     record: &OperationRecord,
    197 ) -> Result<(), Error> {
    198     let (stage, event_id, recovery, committed_at) = encode_state(record.state());
    199     sqlx::query(
    200         "INSERT INTO radroots_runtime_journal_operations (
    201            instance_id, operation_id, idempotency_key, input_digest,
    202            prepared_at_unix_ms, revision, stage, event_id, recovery_record,
    203            cancellation_state, committed_at_unix_ms, updated_at_unix_ms
    204          ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
    205     )
    206     .bind(record.instance_id().as_bytes().as_slice())
    207     .bind(record.operation_id().as_str().as_bytes())
    208     .bind(record.idempotency_key().as_str())
    209     .bind(record.input_digest().as_bytes().as_slice())
    210     .bind(i64_from_u64(record.prepared_at_unix_ms())?)
    211     .bind(i64_from_u64(record.revision().get())?)
    212     .bind(stage)
    213     .bind(event_id)
    214     .bind(recovery)
    215     .bind(cancellation_name(record.cancellation()))
    216     .bind(committed_at.map(i64_from_u64).transpose()?)
    217     .bind(i64_from_u64(updated_at(record))?)
    218     .execute(&mut **transaction)
    219     .await
    220     .map_err(map_backend)?;
    221     Ok(())
    222 }
    223 
    224 pub(crate) fn decode_record(row: &sqlx::sqlite::SqliteRow) -> Result<OperationRecord, Error> {
    225     let instance_id = OperationInstanceId::new(array(
    226         row.try_get::<Vec<u8>, _>("instance_id")
    227             .map_err(map_corrupt)?,
    228     )?)
    229     .map_err(|_| Error::CorruptJournalRecord)?;
    230     let operation_text = String::from_utf8(
    231         row.try_get::<Vec<u8>, _>("operation_id")
    232             .map_err(map_corrupt)?,
    233     )
    234     .map_err(|_| Error::CorruptJournalRecord)?;
    235     let operation_id =
    236         OperationId::parse(operation_text.as_str()).map_err(|_| Error::CorruptJournalRecord)?;
    237     let idempotency_key = IdempotencyKey::parse(
    238         row.try_get::<String, _>("idempotency_key")
    239             .map_err(map_corrupt)?,
    240     )
    241     .map_err(|_| Error::CorruptJournalRecord)?;
    242     let input_digest = IdempotencyDigest::new(array(
    243         row.try_get::<Vec<u8>, _>("input_digest")
    244             .map_err(map_corrupt)?,
    245     )?);
    246     let prepared_at = u64_from_i64(row.try_get("prepared_at_unix_ms").map_err(map_corrupt)?)?;
    247     let revision =
    248         JournalRevision::new(u64_from_i64(row.try_get("revision").map_err(map_corrupt)?)?)
    249             .map_err(|_| Error::CorruptJournalRecord)?;
    250     let event_id = row
    251         .try_get::<Option<Vec<u8>>, _>("event_id")
    252         .map_err(map_corrupt)?
    253         .map(|bytes| array(bytes).map(EventId::from_bytes))
    254         .transpose()?;
    255     let recovery = row
    256         .try_get::<Option<Vec<u8>>, _>("recovery_record")
    257         .map_err(map_corrupt)?
    258         .map(|bytes| decode_recovery(bytes.as_slice()))
    259         .transpose()?;
    260     let committed_at = row
    261         .try_get::<Option<i64>, _>("committed_at_unix_ms")
    262         .map_err(map_corrupt)?
    263         .map(u64_from_i64)
    264         .transpose()?;
    265     let state = decode_state(
    266         row.try_get::<String, _>("stage")
    267             .map_err(map_corrupt)?
    268             .as_str(),
    269         event_id,
    270         recovery,
    271         committed_at,
    272     )?;
    273     let cancellation = cancellation(
    274         row.try_get::<String, _>("cancellation_state")
    275             .map_err(map_corrupt)?
    276             .as_str(),
    277     )?;
    278     OperationRecord::from_parts(
    279         instance_id,
    280         operation_id,
    281         idempotency_key,
    282         input_digest,
    283         prepared_at,
    284         revision,
    285         state,
    286         cancellation,
    287     )
    288     .map_err(|_| Error::CorruptJournalRecord)
    289 }
    290 
    291 type EncodedState = (&'static str, Option<Vec<u8>>, Option<Vec<u8>>, Option<u64>);
    292 
    293 fn encode_state(state: &JournalState) -> EncodedState {
    294     match state {
    295         JournalState::Prepared => ("prepared", None, None, None),
    296         JournalState::Signed { event_id } => {
    297             ("signed", Some(event_id.as_bytes().to_vec()), None, None)
    298         }
    299         JournalState::Recoverable(record) => {
    300             let event_id = match record.point() {
    301                 RecoveryPoint::Prepared => None,
    302                 RecoveryPoint::Signed { event_id } => Some(event_id.as_bytes().to_vec()),
    303             };
    304             ("recoverable", event_id, Some(encode_recovery(record)), None)
    305         }
    306         JournalState::Committed {
    307             event_id,
    308             committed_at_unix_ms,
    309         } => (
    310             "committed",
    311             Some(event_id.as_bytes().to_vec()),
    312             None,
    313             Some(*committed_at_unix_ms),
    314         ),
    315     }
    316 }
    317 
    318 fn decode_state(
    319     stage: &str,
    320     event_id: Option<EventId>,
    321     recovery: Option<RecoveryRecord>,
    322     committed_at: Option<u64>,
    323 ) -> Result<JournalState, Error> {
    324     match (stage, event_id, recovery, committed_at) {
    325         ("prepared", None, None, None) => Ok(JournalState::Prepared),
    326         ("signed", Some(event_id), None, None) => Ok(JournalState::Signed { event_id }),
    327         ("recoverable", event_id, Some(recovery), None)
    328             if recovery_event_id(&recovery) == event_id =>
    329         {
    330             Ok(JournalState::Recoverable(recovery))
    331         }
    332         ("committed", Some(event_id), None, Some(committed_at_unix_ms)) => {
    333             Ok(JournalState::Committed {
    334                 event_id,
    335                 committed_at_unix_ms,
    336             })
    337         }
    338         _ => Err(Error::CorruptJournalRecord),
    339     }
    340 }
    341 
    342 fn recovery_event_id(record: &RecoveryRecord) -> Option<EventId> {
    343     match record.point() {
    344         RecoveryPoint::Prepared => None,
    345         RecoveryPoint::Signed { event_id } => Some(*event_id),
    346     }
    347 }
    348 
    349 fn encode_recovery(record: &RecoveryRecord) -> Vec<u8> {
    350     let mut bytes = Vec::with_capacity(48);
    351     bytes.push(1);
    352     match record.point() {
    353         RecoveryPoint::Prepared => bytes.push(0),
    354         RecoveryPoint::Signed { event_id } => {
    355             bytes.push(1);
    356             bytes.extend_from_slice(event_id.as_bytes());
    357         }
    358     }
    359     bytes.push(recovery_reason_byte(record.reason()));
    360     bytes.extend_from_slice(&record.attempt().to_be_bytes());
    361     match record.retry_not_before_unix_ms() {
    362         Some(value) => {
    363             bytes.push(1);
    364             bytes.extend_from_slice(&value.to_be_bytes());
    365         }
    366         None => bytes.push(0),
    367     }
    368     bytes
    369 }
    370 
    371 fn decode_recovery(bytes: &[u8]) -> Result<RecoveryRecord, Error> {
    372     let mut offset = 0;
    373     if take_byte(bytes, &mut offset)? != 1 {
    374         return Err(Error::CorruptJournalRecord);
    375     }
    376     let point = match take_byte(bytes, &mut offset)? {
    377         0 => RecoveryPoint::Prepared,
    378         1 => RecoveryPoint::Signed {
    379             event_id: EventId::from_bytes(take_array(bytes, &mut offset)?),
    380         },
    381         _ => return Err(Error::CorruptJournalRecord),
    382     };
    383     let reason = recovery_reason(take_byte(bytes, &mut offset)?)?;
    384     let attempt = u32::from_be_bytes(take_array(bytes, &mut offset)?);
    385     let retry = match take_byte(bytes, &mut offset)? {
    386         0 => None,
    387         1 => Some(u64::from_be_bytes(take_array(bytes, &mut offset)?)),
    388         _ => return Err(Error::CorruptJournalRecord),
    389     };
    390     if offset != bytes.len() {
    391         return Err(Error::CorruptJournalRecord);
    392     }
    393     RecoveryRecord::new(point, reason, attempt, retry).map_err(|_| Error::CorruptJournalRecord)
    394 }
    395 
    396 const fn recovery_reason_byte(reason: RecoveryReason) -> u8 {
    397     match reason {
    398         RecoveryReason::CancelledBeforeCommit => 0,
    399         RecoveryReason::SignerUnavailable => 1,
    400         RecoveryReason::TransportUnavailable => 2,
    401         RecoveryReason::StorageUnavailable => 3,
    402         RecoveryReason::DeadlineExceeded => 4,
    403         RecoveryReason::Interrupted => 5,
    404     }
    405 }
    406 
    407 const fn recovery_reason(value: u8) -> Result<RecoveryReason, Error> {
    408     match value {
    409         0 => Ok(RecoveryReason::CancelledBeforeCommit),
    410         1 => Ok(RecoveryReason::SignerUnavailable),
    411         2 => Ok(RecoveryReason::TransportUnavailable),
    412         3 => Ok(RecoveryReason::StorageUnavailable),
    413         4 => Ok(RecoveryReason::DeadlineExceeded),
    414         5 => Ok(RecoveryReason::Interrupted),
    415         _ => Err(Error::CorruptJournalRecord),
    416     }
    417 }
    418 
    419 const fn cancellation_name(value: CancellationState) -> &'static str {
    420     match value {
    421         CancellationState::NotRequested => "not_requested",
    422         CancellationState::CancelledBeforeCommit => "cancelled_before_commit",
    423         CancellationState::ObservedAfterCommit => "observed_after_commit",
    424     }
    425 }
    426 
    427 const fn cancellation(value: &str) -> Result<CancellationState, Error> {
    428     match value.as_bytes() {
    429         b"not_requested" => Ok(CancellationState::NotRequested),
    430         b"cancelled_before_commit" => Ok(CancellationState::CancelledBeforeCommit),
    431         b"observed_after_commit" => Ok(CancellationState::ObservedAfterCommit),
    432         _ => Err(Error::CorruptJournalRecord),
    433     }
    434 }
    435 
    436 fn updated_at(record: &OperationRecord) -> u64 {
    437     match record.state() {
    438         JournalState::Committed {
    439             committed_at_unix_ms,
    440             ..
    441         } => *committed_at_unix_ms,
    442         JournalState::Recoverable(recovery) => recovery
    443             .retry_not_before_unix_ms()
    444             .unwrap_or(record.prepared_at_unix_ms()),
    445         JournalState::Prepared | JournalState::Signed { .. } => record.prepared_at_unix_ms(),
    446     }
    447 }
    448 
    449 fn take_byte(bytes: &[u8], offset: &mut usize) -> Result<u8, Error> {
    450     let value = bytes
    451         .get(*offset)
    452         .copied()
    453         .ok_or(Error::CorruptJournalRecord)?;
    454     *offset += 1;
    455     Ok(value)
    456 }
    457 
    458 fn take_array<const N: usize>(bytes: &[u8], offset: &mut usize) -> Result<[u8; N], Error> {
    459     let end = offset.checked_add(N).ok_or(Error::CorruptJournalRecord)?;
    460     let value = bytes
    461         .get(*offset..end)
    462         .ok_or(Error::CorruptJournalRecord)?
    463         .try_into()
    464         .map_err(|_| Error::CorruptJournalRecord)?;
    465     *offset = end;
    466     Ok(value)
    467 }
    468 
    469 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
    470     bytes.try_into().map_err(|_| Error::CorruptJournalRecord)
    471 }
    472 
    473 fn i64_from_u64(value: u64) -> Result<i64, Error> {
    474     i64::try_from(value).map_err(|_| Error::CorruptJournalRecord)
    475 }
    476 
    477 fn u64_from_i64(value: i64) -> Result<u64, Error> {
    478     u64::try_from(value).map_err(|_| Error::CorruptJournalRecord)
    479 }
    480 
    481 fn map_corrupt(_: sqlx::Error) -> Error {
    482     Error::CorruptJournalRecord
    483 }
    484 
    485 #[cfg(test)]
    486 #[cfg_attr(coverage_nightly, coverage(off))]
    487 mod tests {
    488     use super::*;
    489     use crate::migration::runtime::{MIGRATIONS, migration_sql};
    490     use radroots_storage::{
    491         Journal, event::SourceGeneration, journal::JournalStage, status::EventStoreMode,
    492     };
    493     use sqlx::sqlite::SqlitePoolOptions;
    494 
    495     async fn store(mode: EventStoreMode) -> SqliteStorage {
    496         let pool = SqlitePoolOptions::new()
    497             .max_connections(1)
    498             .connect("sqlite::memory:")
    499             .await
    500             .expect("memory SQLite");
    501         sqlx::query("PRAGMA foreign_keys = ON")
    502             .execute(&pool)
    503             .await
    504             .expect("foreign keys");
    505         for migration in MIGRATIONS {
    506             sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL"))
    507                 .execute(&pool)
    508                 .await
    509                 .expect("runtime migration");
    510         }
    511         SqliteStorage::new(
    512             pool,
    513             SourceGeneration::new([21; 32]).expect("generation"),
    514             mode,
    515         )
    516     }
    517 
    518     fn instance(byte: u8) -> OperationInstanceId {
    519         OperationInstanceId::new([byte; 16]).expect("instance")
    520     }
    521 
    522     fn key(byte: u8) -> IdempotencyKey {
    523         IdempotencyKey::parse(format!("journal-{byte:02x}")).expect("key")
    524     }
    525 
    526     fn prepare(
    527         instance_id: OperationInstanceId,
    528         key_byte: u8,
    529         digest: u8,
    530         at: u64,
    531     ) -> PrepareOperation {
    532         PrepareOperation::new(
    533             instance_id,
    534             OperationId::SyncPush,
    535             key(key_byte),
    536             IdempotencyDigest::new([digest; 32]),
    537             at,
    538         )
    539         .expect("prepare")
    540     }
    541 
    542     #[test]
    543     fn recovery_decoder_and_state_guard_reject_noncanonical_encodings() {
    544         let event_id = EventId::from_bytes([42; 32]);
    545         let recovery = RecoveryRecord::new(
    546             RecoveryPoint::Signed { event_id },
    547             RecoveryReason::Interrupted,
    548             1,
    549             None,
    550         )
    551         .expect("recovery");
    552         assert_eq!(
    553             decode_state("recoverable", None, Some(recovery.clone()), None),
    554             Err(Error::CorruptJournalRecord)
    555         );
    556         assert_eq!(decode_recovery(&[2]), Err(Error::CorruptJournalRecord));
    557         let mut encoded = encode_recovery(&recovery);
    558         encoded.push(0);
    559         assert_eq!(
    560             decode_recovery(encoded.as_slice()),
    561             Err(Error::CorruptJournalRecord)
    562         );
    563     }
    564 
    565     #[tokio::test]
    566     async fn prepare_replays_exact_identity_and_rejects_conflicts() {
    567         let store = store(EventStoreMode::ReadWrite).await;
    568         let request = prepare(instance(1), 1, 2, 100);
    569         let created = store.prepare(request.clone()).await.expect("created");
    570         assert_eq!(created.disposition(), PrepareDisposition::Created);
    571         assert_eq!(created.record().revision(), JournalRevision::INITIAL);
    572         let replay = store.prepare(request).await.expect("replay");
    573         assert_eq!(replay.disposition(), PrepareDisposition::Replay);
    574         assert_eq!(replay.record(), created.record());
    575         assert_eq!(
    576             store.prepare(prepare(instance(1), 1, 3, 100)).await,
    577             Err(Error::IdempotencyConflict)
    578         );
    579         assert_eq!(
    580             store.prepare(prepare(instance(1), 2, 2, 100)).await,
    581             Err(Error::OperationIdentityMismatch)
    582         );
    583         assert_eq!(
    584             store
    585                 .operation(instance(1))
    586                 .await
    587                 .expect("lookup")
    588                 .expect("record"),
    589             *created.record()
    590         );
    591         assert_eq!(
    592             store
    593                 .by_idempotency_key(OperationId::SyncPush, key(1))
    594                 .await
    595                 .expect("key lookup")
    596                 .expect("record"),
    597             *created.record()
    598         );
    599         assert!(
    600             store
    601                 .by_idempotency_key(OperationId::FarmPublish, key(1))
    602                 .await
    603                 .expect("wrong operation lookup")
    604                 .is_none()
    605         );
    606     }
    607 
    608     #[tokio::test]
    609     async fn lifecycle_recovery_commit_and_cancellation_round_trip() {
    610         let store = store(EventStoreMode::ReadWrite).await;
    611         let instance_id = instance(3);
    612         let event_id = EventId::from_bytes([4; 32]);
    613         let prepared = store
    614             .prepare(prepare(instance_id, 3, 3, 100))
    615             .await
    616             .expect("prepare")
    617             .record()
    618             .clone();
    619         let signed = store
    620             .transition(JournalTransition::signed(
    621                 instance_id,
    622                 prepared.revision(),
    623                 event_id,
    624             ))
    625             .await
    626             .expect("signed");
    627         assert_eq!(signed.state().stage(), JournalStage::Signed);
    628         assert_eq!(
    629             store
    630                 .transition(JournalTransition::signed(
    631                     instance_id,
    632                     prepared.revision(),
    633                     event_id,
    634                 ))
    635                 .await,
    636             Err(Error::JournalRevisionConflict)
    637         );
    638 
    639         let recovery = RecoveryRecord::new(
    640             RecoveryPoint::Signed { event_id },
    641             RecoveryReason::TransportUnavailable,
    642             2,
    643             Some(200),
    644         )
    645         .expect("recovery");
    646         let recoverable = store
    647             .transition(JournalTransition::recoverable(
    648                 instance_id,
    649                 signed.revision(),
    650                 recovery.clone(),
    651             ))
    652             .await
    653             .expect("recoverable");
    654         assert_eq!(
    655             store.recoverable(10).await.expect("recovery query"),
    656             vec![recoverable.clone()]
    657         );
    658         assert_eq!(recoverable.state(), &JournalState::Recoverable(recovery));
    659 
    660         let resumed = store
    661             .transition(JournalTransition::resume(
    662                 instance_id,
    663                 recoverable.revision(),
    664             ))
    665             .await
    666             .expect("resume");
    667         assert_eq!(resumed.state().stage(), JournalStage::Signed);
    668         let committed = store
    669             .transition(JournalTransition::committed(
    670                 instance_id,
    671                 resumed.revision(),
    672                 event_id,
    673                 250,
    674             ))
    675             .await
    676             .expect("commit");
    677         assert_eq!(committed.state().stage(), JournalStage::Committed);
    678         let cancelled = store
    679             .transition(JournalTransition::cancelled(
    680                 instance_id,
    681                 committed.revision(),
    682                 260,
    683             ))
    684             .await
    685             .expect("post-commit cancellation");
    686         assert_eq!(cancelled.state(), committed.state());
    687         assert_eq!(
    688             cancelled.cancellation(),
    689             CancellationState::ObservedAfterCommit
    690         );
    691         assert!(
    692             store
    693                 .recoverable(10)
    694                 .await
    695                 .expect("empty recovery")
    696                 .is_empty()
    697         );
    698     }
    699 
    700     #[tokio::test]
    701     async fn cancellation_corruption_bounds_and_read_only_mode_fail_closed() {
    702         let store = store(EventStoreMode::ReadWrite).await;
    703         let instance_id = instance(5);
    704         let prepared = store
    705             .prepare(prepare(instance_id, 5, 5, 500))
    706             .await
    707             .expect("prepare")
    708             .record()
    709             .clone();
    710         let cancelled = store
    711             .transition(JournalTransition::cancelled(
    712                 instance_id,
    713                 prepared.revision(),
    714                 501,
    715             ))
    716             .await
    717             .expect("cancel");
    718         assert_eq!(cancelled.state().stage(), JournalStage::Recoverable);
    719         assert_eq!(
    720             store.recoverable(0).await,
    721             Err(Error::InvalidJournalQueryLimit)
    722         );
    723         assert_eq!(
    724             store.recoverable(RECOVERABLE_QUERY_LIMIT_MAX + 1).await,
    725             Err(Error::InvalidJournalQueryLimit)
    726         );
    727 
    728         sqlx::query(
    729             "UPDATE radroots_runtime_journal_operations
    730              SET recovery_record = X'0100FF' WHERE instance_id = ?",
    731         )
    732         .bind(instance_id.as_bytes().as_slice())
    733         .execute(store.pool())
    734         .await
    735         .expect("forge corrupt recovery");
    736         assert_eq!(
    737             store.operation(instance_id).await,
    738             Err(Error::CorruptJournalRecord)
    739         );
    740 
    741         let read_only = SqliteStorage::new(
    742             store.pool().clone(),
    743             SourceGeneration::new([21; 32]).expect("generation"),
    744             EventStoreMode::ReadOnly,
    745         );
    746         assert_eq!(
    747             read_only.prepare(prepare(instance(6), 6, 6, 600)).await,
    748             Err(Error::BackendUnavailable)
    749         );
    750     }
    751 }