lib

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

atomic.rs (19148B)


      1 use crate::backend::map_backend;
      2 use crate::{SqliteStorage, projection};
      3 use radroots_storage::{
      4     Error,
      5     atomic::{
      6         AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId,
      7         AtomicCommitOutcome, AtomicCommitReceipt, AtomicStorage, AtomicWorkflow,
      8         AtomicWorkflowKind,
      9     },
     10     event::{
     11         AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventPosition,
     12         EventSequence, SourceGeneration,
     13     },
     14 };
     15 use sqlx::{Row, Sqlite};
     16 
     17 const RECEIPT_FORMAT_VERSION: u8 = 1;
     18 const RECEIPT_MAX_BYTES: usize = 4 * 1024 * 1024;
     19 
     20 #[cfg_attr(coverage_nightly, coverage(off))]
     21 impl AtomicStorage for SqliteStorage {
     22     fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> {
     23         Box::pin(async move {
     24             self.require_atomic_writer()?;
     25             let mut transaction = self
     26                 .pool()
     27                 .begin_with("BEGIN IMMEDIATE")
     28                 .await
     29                 .map_err(map_backend)?;
     30 
     31             let result = commit_transaction(self, &mut transaction, &request).await;
     32             match result {
     33                 Ok(receipt) => {
     34                     transaction.commit().await.map_err(map_backend)?;
     35                     Ok(receipt)
     36                 }
     37                 Err(primary) => {
     38                     let rollback = transaction.rollback().await;
     39                     Err(preserve_primary(primary, rollback))
     40                 }
     41             }
     42         })
     43     }
     44 
     45     fn receipt(
     46         &self,
     47         commit_id: AtomicCommitId,
     48     ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>> {
     49         Box::pin(async move {
     50             sqlx::query(
     51                 "SELECT commit_id, commit_digest, workflow_kind, requested_at_unix_ms,
     52                         committed_at_unix_ms, receipt
     53                  FROM radroots_runtime_atomic_commits WHERE commit_id = ?",
     54             )
     55             .bind(commit_id.as_bytes().as_slice())
     56             .fetch_optional(self.pool())
     57             .await
     58             .map_err(map_backend)?
     59             .as_ref()
     60             .map(decode_receipt_row)
     61             .transpose()
     62         })
     63     }
     64 }
     65 
     66 impl SqliteStorage {
     67     fn require_atomic_writer(&self) -> Result<(), Error> {
     68         if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
     69             return Err(Error::BackendUnavailable);
     70         }
     71         Ok(())
     72     }
     73 }
     74 
     75 async fn commit_transaction(
     76     storage: &SqliteStorage,
     77     transaction: &mut sqlx::Transaction<'_, Sqlite>,
     78     request: &AtomicCommit,
     79 ) -> Result<AtomicCommitReceipt, Error> {
     80     if let Some(row) = sqlx::query(
     81         "SELECT commit_id, commit_digest, workflow_kind, requested_at_unix_ms,
     82                 committed_at_unix_ms, receipt
     83          FROM radroots_runtime_atomic_commits WHERE commit_id = ?",
     84     )
     85     .bind(request.commit_id().as_bytes().as_slice())
     86     .fetch_optional(&mut **transaction)
     87     .await
     88     .map_err(map_backend)?
     89     {
     90         let committed = decode_receipt_row(&row)?;
     91         if [
     92             committed.digest() != request.digest(),
     93             committed.outcome().kind() != request.workflow().kind(),
     94         ]
     95         .contains(&true)
     96         {
     97             return Err(Error::AtomicCommitConflict);
     98         }
     99         return AtomicCommitReceipt::new(
    100             request,
    101             AtomicCommitDisposition::Replay,
    102             committed.committed_at_unix_ms(),
    103             committed.outcome().clone(),
    104         );
    105     }
    106 
    107     let outcome = execute_workflow(storage, transaction, request.workflow().clone()).await?;
    108     let committed_at_unix_ms = request.requested_at_unix_ms();
    109     let receipt = AtomicCommitReceipt::new(
    110         request,
    111         AtomicCommitDisposition::Committed,
    112         committed_at_unix_ms,
    113         outcome,
    114     )?;
    115     let snapshot = encode_outcome(receipt.outcome())?;
    116     sqlx::query(
    117         "INSERT INTO radroots_runtime_atomic_commits (
    118            commit_id, commit_digest, workflow_kind, requested_at_unix_ms,
    119            committed_at_unix_ms, receipt
    120          ) VALUES (?, ?, ?, ?, ?, ?)",
    121     )
    122     .bind(request.commit_id().as_bytes().as_slice())
    123     .bind(request.digest().as_bytes().as_slice())
    124     .bind(workflow_name(request.workflow().kind()))
    125     .bind(i64_from_u64(request.requested_at_unix_ms())?)
    126     .bind(i64_from_u64(committed_at_unix_ms)?)
    127     .bind(snapshot)
    128     .execute(&mut **transaction)
    129     .await
    130     .map_err(map_backend)?;
    131     Ok(receipt)
    132 }
    133 
    134 async fn execute_workflow(
    135     storage: &SqliteStorage,
    136     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    137     workflow: AtomicWorkflow,
    138 ) -> Result<AtomicCommitOutcome, Error> {
    139     match workflow {
    140         AtomicWorkflow::Ingested(ingested) => {
    141             let admission = storage
    142                 .admit_transaction(transaction, ingested.admission().clone())
    143                 .await?;
    144             let projection = match ingested.projection().cloned() {
    145                 Some(checkpoint) => Some(Box::new(
    146                     projection::checkpoint_transaction(transaction, checkpoint).await?,
    147                 )),
    148                 None => None,
    149             };
    150             Ok(AtomicCommitOutcome::Ingested {
    151                 admission,
    152                 projection,
    153             })
    154         }
    155     }
    156 }
    157 
    158 fn decode_receipt_row(row: &sqlx::sqlite::SqliteRow) -> Result<AtomicCommitReceipt, Error> {
    159     let commit_id = AtomicCommitId::new(array(
    160         row.try_get::<Vec<u8>, _>("commit_id")
    161             .map_err(map_corrupt)?,
    162     )?)
    163     .map_err(|_| Error::AtomicCommitFailed)?;
    164     let digest = AtomicCommitDigest::new(array(
    165         row.try_get::<Vec<u8>, _>("commit_digest")
    166             .map_err(map_corrupt)?,
    167     )?);
    168     let workflow_kind = workflow_kind(
    169         row.try_get::<String, _>("workflow_kind")
    170             .map_err(map_corrupt)?
    171             .as_str(),
    172     )?;
    173     let requested_at_unix_ms =
    174         u64_from_i64(row.try_get("requested_at_unix_ms").map_err(map_corrupt)?)?;
    175     let committed_at_unix_ms =
    176         u64_from_i64(row.try_get("committed_at_unix_ms").map_err(map_corrupt)?)?;
    177     let bytes = row.try_get::<Vec<u8>, _>("receipt").map_err(map_corrupt)?;
    178     let outcome = decode_outcome(bytes.as_slice())?;
    179     AtomicCommitReceipt::from_durable_parts(
    180         commit_id,
    181         digest,
    182         AtomicCommitDisposition::Committed,
    183         requested_at_unix_ms,
    184         committed_at_unix_ms,
    185         workflow_kind,
    186         outcome,
    187     )
    188     .map_err(|_| Error::AtomicCommitFailed)
    189 }
    190 
    191 fn encode_outcome(outcome: &AtomicCommitOutcome) -> Result<Vec<u8>, Error> {
    192     let mut bytes = Vec::with_capacity(512);
    193     bytes.push(RECEIPT_FORMAT_VERSION);
    194     match outcome {
    195         AtomicCommitOutcome::Ingested {
    196             admission,
    197             projection,
    198         } => {
    199             bytes.push(4);
    200             encode_admission(&mut bytes, admission);
    201             match projection {
    202                 Some(status) => {
    203                     bytes.push(1);
    204                     put_blob(
    205                         &mut bytes,
    206                         projection::encode_status_snapshot(status)?.as_slice(),
    207                     )?;
    208                 }
    209                 None => bytes.push(0),
    210             }
    211         }
    212     }
    213     if bytes.len() > RECEIPT_MAX_BYTES {
    214         return Err(Error::AtomicCommitFailed);
    215     }
    216     Ok(bytes)
    217 }
    218 
    219 fn decode_outcome(bytes: &[u8]) -> Result<AtomicCommitOutcome, Error> {
    220     if bytes.len() > RECEIPT_MAX_BYTES {
    221         return Err(Error::AtomicCommitFailed);
    222     }
    223     let mut cursor = Cursor::new(bytes);
    224     if cursor.byte()? != RECEIPT_FORMAT_VERSION {
    225         return Err(Error::AtomicCommitFailed);
    226     }
    227     let outcome = match cursor.byte()? {
    228         4 => {
    229             let admission = decode_admission(&mut cursor)?;
    230             let projection = match cursor.byte()? {
    231                 0 => None,
    232                 1 => Some(Box::new(
    233                     projection::decode_status_snapshot(cursor.blob()?)
    234                         .map_err(|_| Error::AtomicCommitFailed)?,
    235                 )),
    236                 _ => return Err(Error::AtomicCommitFailed),
    237             };
    238             AtomicCommitOutcome::Ingested {
    239                 admission,
    240                 projection,
    241             }
    242         }
    243         _ => return Err(Error::AtomicCommitFailed),
    244     };
    245     cursor.finish()?;
    246     Ok(outcome)
    247 }
    248 
    249 fn encode_admission(bytes: &mut Vec<u8>, receipt: &AdmissionReceipt) {
    250     bytes.extend_from_slice(receipt.event_id().as_bytes());
    251     bytes.extend_from_slice(receipt.position().generation().as_bytes());
    252     bytes.extend_from_slice(&receipt.position().sequence().get().to_be_bytes());
    253     bytes.push(match receipt.stage() {
    254         AdmissionStage::Raw => 0,
    255         AdmissionStage::Verified => 1,
    256         AdmissionStage::Visible => 2,
    257     });
    258     bytes.push(match receipt.disposition() {
    259         AdmissionDisposition::Inserted => 0,
    260         AdmissionDisposition::Advanced => 1,
    261         AdmissionDisposition::Duplicate => 2,
    262     });
    263 }
    264 
    265 fn decode_admission(cursor: &mut Cursor<'_>) -> Result<AdmissionReceipt, Error> {
    266     let event_id = radroots_storage::event::EventId::from_bytes(cursor.array()?);
    267     let generation =
    268         SourceGeneration::new(cursor.array()?).map_err(|_| Error::AtomicCommitFailed)?;
    269     let sequence = EventSequence::new(cursor.u64()?).map_err(|_| Error::AtomicCommitFailed)?;
    270     let stage = match cursor.byte()? {
    271         0 => AdmissionStage::Raw,
    272         1 => AdmissionStage::Verified,
    273         2 => AdmissionStage::Visible,
    274         _ => return Err(Error::AtomicCommitFailed),
    275     };
    276     let disposition = match cursor.byte()? {
    277         0 => AdmissionDisposition::Inserted,
    278         1 => AdmissionDisposition::Advanced,
    279         2 => AdmissionDisposition::Duplicate,
    280         _ => return Err(Error::AtomicCommitFailed),
    281     };
    282     Ok(AdmissionReceipt::new(
    283         event_id,
    284         EventPosition::new(generation, sequence),
    285         stage,
    286         disposition,
    287     ))
    288 }
    289 
    290 fn put_blob(bytes: &mut Vec<u8>, value: &[u8]) -> Result<(), Error> {
    291     let length = u32::try_from(value.len()).map_err(|_| Error::AtomicCommitFailed)?;
    292     bytes.extend_from_slice(&length.to_be_bytes());
    293     bytes.extend_from_slice(value);
    294     Ok(())
    295 }
    296 
    297 struct Cursor<'a> {
    298     bytes: &'a [u8],
    299     offset: usize,
    300 }
    301 
    302 impl<'a> Cursor<'a> {
    303     const fn new(bytes: &'a [u8]) -> Self {
    304         Self { bytes, offset: 0 }
    305     }
    306 
    307     fn byte(&mut self) -> Result<u8, Error> {
    308         let value = self
    309             .bytes
    310             .get(self.offset)
    311             .copied()
    312             .ok_or(Error::AtomicCommitFailed)?;
    313         self.offset += 1;
    314         Ok(value)
    315     }
    316 
    317     fn u32(&mut self) -> Result<u32, Error> {
    318         Ok(u32::from_be_bytes(self.array()?))
    319     }
    320 
    321     fn u64(&mut self) -> Result<u64, Error> {
    322         Ok(u64::from_be_bytes(self.array()?))
    323     }
    324 
    325     fn array<const N: usize>(&mut self) -> Result<[u8; N], Error> {
    326         self.take(N)?
    327             .try_into()
    328             .map_err(|_| Error::AtomicCommitFailed)
    329     }
    330 
    331     fn blob(&mut self) -> Result<&'a [u8], Error> {
    332         let length = usize::try_from(self.u32()?).map_err(|_| Error::AtomicCommitFailed)?;
    333         self.take(length)
    334     }
    335 
    336     fn take(&mut self, length: usize) -> Result<&'a [u8], Error> {
    337         let end = self
    338             .offset
    339             .checked_add(length)
    340             .ok_or(Error::AtomicCommitFailed)?;
    341         let value = self
    342             .bytes
    343             .get(self.offset..end)
    344             .ok_or(Error::AtomicCommitFailed)?;
    345         self.offset = end;
    346         Ok(value)
    347     }
    348 
    349     fn finish(self) -> Result<(), Error> {
    350         if self.offset == self.bytes.len() {
    351             Ok(())
    352         } else {
    353             Err(Error::AtomicCommitFailed)
    354         }
    355     }
    356 }
    357 
    358 fn preserve_primary<T>(primary: Error, _rollback: Result<(), T>) -> Error {
    359     primary
    360 }
    361 
    362 const fn workflow_name(kind: AtomicWorkflowKind) -> &'static str {
    363     match kind {
    364         AtomicWorkflowKind::Ingested => "ingested",
    365     }
    366 }
    367 
    368 fn workflow_kind(value: &str) -> Result<AtomicWorkflowKind, Error> {
    369     match value.as_bytes() {
    370         b"ingested" => Ok(AtomicWorkflowKind::Ingested),
    371         _ => Err(Error::AtomicCommitFailed),
    372     }
    373 }
    374 
    375 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
    376     bytes.try_into().map_err(|_| Error::AtomicCommitFailed)
    377 }
    378 
    379 fn i64_from_u64(value: u64) -> Result<i64, Error> {
    380     i64::try_from(value).map_err(|_| Error::AtomicCommitFailed)
    381 }
    382 
    383 fn u64_from_i64(value: i64) -> Result<u64, Error> {
    384     u64::try_from(value).map_err(|_| Error::AtomicCommitFailed)
    385 }
    386 
    387 fn map_corrupt(_: sqlx::Error) -> Error {
    388     Error::AtomicCommitFailed
    389 }
    390 
    391 #[cfg(test)]
    392 #[cfg_attr(coverage_nightly, coverage(off))]
    393 mod tests {
    394     use super::*;
    395     use crate::migration::runtime::{MIGRATIONS, migration_sql};
    396     use radroots_event::{SignedEvent, wire::Nip01EventWire};
    397     use radroots_storage::{
    398         atomic::{AtomicWorkflow, CommitIngested},
    399         event::{EventAdmission, SourceGeneration},
    400         projection::{ProjectionCheckpoint, ProjectionGeneration, ProjectionId},
    401         status::EventStoreMode,
    402     };
    403     use radroots_transport::{
    404         Target, TransportId,
    405         source::{EventProvenance, ObservedEvent},
    406     };
    407     use sqlx::sqlite::SqlitePoolOptions;
    408 
    409     async fn store(mode: EventStoreMode) -> SqliteStorage {
    410         let generation = SourceGeneration::new([71; 32]).expect("generation");
    411         let pool = SqlitePoolOptions::new()
    412             .max_connections(1)
    413             .connect("sqlite::memory:")
    414             .await
    415             .expect("memory SQLite");
    416         sqlx::query("PRAGMA foreign_keys = ON")
    417             .execute(&pool)
    418             .await
    419             .expect("foreign keys");
    420         for migration in MIGRATIONS {
    421             sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL"))
    422                 .execute(&pool)
    423                 .await
    424                 .expect("runtime migration");
    425         }
    426         sqlx::query(
    427             "INSERT INTO radroots_runtime_source_generations (
    428                generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms
    429              ) VALUES (?, 0, 'active', 1, NULL)",
    430         )
    431         .bind(generation.as_bytes().as_slice())
    432         .execute(&pool)
    433         .await
    434         .expect("source generation");
    435         SqliteStorage::new(pool, generation, mode)
    436     }
    437 
    438     fn signed_event() -> SignedEvent {
    439         let mut wire = Nip01EventWire {
    440             id: "0".repeat(64),
    441             pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
    442             created_at: 1_800_001_001,
    443             kind: 1,
    444             tags: vec![],
    445             content: "atomic inbound".to_owned(),
    446             sig: "42".repeat(64),
    447             extra: Default::default(),
    448         };
    449         wire.id = wire.computed_event_id().expect("event id").to_hex();
    450         let raw_json = serde_json::json!({
    451             "id": &wire.id,
    452             "pubkey": &wire.pubkey,
    453             "created_at": wire.created_at,
    454             "kind": wire.kind,
    455             "tags": &wire.tags,
    456             "content": &wire.content,
    457             "sig": &wire.sig,
    458         })
    459         .to_string();
    460         SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event")
    461     }
    462 
    463     fn admission(observed_at: u64) -> EventAdmission {
    464         let target = Target::new(TransportId::NOSTR, "wss://atomic.example").expect("target");
    465         let provenance = EventProvenance::new(
    466             TransportId::NOSTR,
    467             target.fingerprint().clone(),
    468             observed_at,
    469         )
    470         .expect("provenance");
    471         EventAdmission::raw(ObservedEvent::new(signed_event(), provenance))
    472     }
    473 
    474     fn commit(byte: u8, digest: u8, workflow: AtomicWorkflow) -> AtomicCommit {
    475         AtomicCommit::new(
    476             AtomicCommitId::new([byte; 16]).expect("commit id"),
    477             AtomicCommitDigest::new([digest; 32]),
    478             100,
    479             workflow,
    480         )
    481         .expect("atomic commit")
    482     }
    483 
    484     fn checkpoint() -> ProjectionCheckpoint {
    485         ProjectionCheckpoint::new(
    486             ProjectionId::parse("atomic_projection").expect("projection id"),
    487             ProjectionGeneration::new([81; 32]).expect("projection generation"),
    488             None,
    489             1,
    490             100,
    491         )
    492         .expect("checkpoint")
    493     }
    494 
    495     #[tokio::test]
    496     async fn inbound_commit_is_atomic_replayable_and_reconstructable() {
    497         let store = store(EventStoreMode::ReadWrite).await;
    498         let request = commit(
    499             1,
    500             1,
    501             AtomicWorkflow::Ingested(Box::new(CommitIngested::new(
    502                 admission(90),
    503                 Some(checkpoint()),
    504             ))),
    505         );
    506         let committed = store.commit(request.clone()).await.expect("commit");
    507         assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed);
    508         assert!(matches!(
    509             committed.outcome(),
    510             AtomicCommitOutcome::Ingested {
    511                 projection: Some(_),
    512                 ..
    513             }
    514         ));
    515 
    516         let replay = store.commit(request.clone()).await.expect("replay");
    517         assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
    518         assert_eq!(replay.outcome(), committed.outcome());
    519         let reconstructed = store
    520             .receipt(request.commit_id())
    521             .await
    522             .expect("receipt lookup")
    523             .expect("receipt");
    524         assert_eq!(reconstructed.outcome(), committed.outcome());
    525 
    526         let conflict = commit(
    527             1,
    528             2,
    529             AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(90), None))),
    530         );
    531         assert_eq!(
    532             store.commit(conflict).await,
    533             Err(Error::AtomicCommitConflict)
    534         );
    535     }
    536 
    537     #[tokio::test]
    538     async fn read_only_mode_rejects_atomic_mutation() {
    539         let store = store(EventStoreMode::ReadOnly).await;
    540         let request = commit(
    541             2,
    542             2,
    543             AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(90), None))),
    544         );
    545         assert_eq!(store.commit(request).await, Err(Error::BackendUnavailable));
    546     }
    547 
    548     #[test]
    549     fn receipt_decoder_rejects_retired_and_trailing_payloads() {
    550         assert_eq!(
    551             decode_outcome(&vec![0; RECEIPT_MAX_BYTES + 1]),
    552             Err(Error::AtomicCommitFailed)
    553         );
    554         assert_eq!(decode_outcome(&[0]), Err(Error::AtomicCommitFailed));
    555         assert_eq!(
    556             decode_outcome(&[RECEIPT_FORMAT_VERSION, 0]),
    557             Err(Error::AtomicCommitFailed)
    558         );
    559         let outcome = AtomicCommitOutcome::Ingested {
    560             admission: AdmissionReceipt::new(
    561                 radroots_storage::event::EventId::from_bytes([1; 32]),
    562                 EventPosition::new(
    563                     SourceGeneration::new([2; 32]).expect("generation"),
    564                     EventSequence::new(1).expect("sequence"),
    565                 ),
    566                 AdmissionStage::Raw,
    567                 AdmissionDisposition::Inserted,
    568             ),
    569             projection: None,
    570         };
    571         let mut bytes = encode_outcome(&outcome).expect("encode");
    572         bytes.push(0);
    573         assert_eq!(decode_outcome(&bytes), Err(Error::AtomicCommitFailed));
    574     }
    575 }