lib

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

mod.rs (42954B)


      1 use crate::backend::map_backend;
      2 use radroots_event::SignedEvent;
      3 use radroots_event_codec::Codec;
      4 use radroots_storage::{
      5     Error, EventStore,
      6     event::{
      7         AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventAdmission,
      8         EventCursor, EventId, EventPage, EventPosition, EventQuery, EventQueryBounds,
      9         EventSequence, SourceGeneration, StoredEventProvenance, StoredRawEvent,
     10         StoredVerifiedEvent, StoredVisibleEvent, VisibilityEvaluation, VisibilityInput,
     11         VisibilitySnapshot, evaluate_visibility,
     12     },
     13     status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
     14 };
     15 use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool};
     16 use std::{
     17     path::PathBuf,
     18     sync::{Arc, Mutex},
     19 };
     20 
     21 use crate::lock::WriterLock;
     22 use crate::status::StorageLifecycle;
     23 
     24 #[derive(Clone)]
     25 pub struct SqliteStorage {
     26     pub(crate) pool: SqlitePool,
     27     pub(crate) private_pool: SqlitePool,
     28     pub(crate) generation: SourceGeneration,
     29     pub(crate) mode: EventStoreMode,
     30     pub(crate) lifecycle: Arc<StorageLifecycle>,
     31     pub(crate) backup_root: Option<Arc<PathBuf>>,
     32     pub(crate) paths: Option<Arc<crate::Paths>>,
     33     pub(crate) reliability: Arc<Mutex<crate::backup::ReliabilityState>>,
     34 }
     35 
     36 struct StoredEventRow {
     37     position: EventPosition,
     38     raw_json: String,
     39     event: SignedEvent,
     40     stage: AdmissionStage,
     41 }
     42 
     43 impl SqliteStorage {
     44     #[allow(dead_code)] // Single-pool in-memory scaffold retained for focused backend tests.
     45     pub(crate) fn new(
     46         pool: SqlitePool,
     47         generation: SourceGeneration,
     48         mode: EventStoreMode,
     49     ) -> Self {
     50         Self {
     51             private_pool: pool.clone(),
     52             pool,
     53             generation,
     54             mode,
     55             lifecycle: Arc::new(StorageLifecycle::scaffold(mode)),
     56             backup_root: None,
     57             paths: None,
     58             reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())),
     59         }
     60     }
     61 
     62     #[allow(dead_code)] // Two-pool in-memory scaffold retained for focused backend tests.
     63     pub(crate) fn with_private_pool(
     64         pool: SqlitePool,
     65         private_pool: SqlitePool,
     66         generation: SourceGeneration,
     67         mode: EventStoreMode,
     68     ) -> Self {
     69         Self {
     70             pool,
     71             private_pool,
     72             generation,
     73             mode,
     74             lifecycle: Arc::new(StorageLifecycle::scaffold(mode)),
     75             backup_root: None,
     76             paths: None,
     77             reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())),
     78         }
     79     }
     80 
     81     pub(crate) fn from_opened(
     82         pool: SqlitePool,
     83         private_pool: SqlitePool,
     84         generation: SourceGeneration,
     85         options: &crate::OpenOptions,
     86         writer_lock: Option<WriterLock>,
     87     ) -> Self {
     88         let mode = if options.mode().is_writable() {
     89             EventStoreMode::ReadWrite
     90         } else {
     91             EventStoreMode::ReadOnly
     92         };
     93         Self {
     94             pool,
     95             private_pool,
     96             generation,
     97             mode,
     98             lifecycle: Arc::new(StorageLifecycle::new(
     99                 options.mode(),
    100                 options.busy_timeout(),
    101                 writer_lock,
    102             )),
    103             backup_root: options
    104                 .backup_root()
    105                 .map(|path| Arc::new(path.to_path_buf())),
    106             paths: Some(Arc::new(options.paths().clone())),
    107             reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())),
    108         }
    109     }
    110 
    111     pub(crate) const fn pool(&self) -> &SqlitePool {
    112         &self.pool
    113     }
    114 
    115     pub(crate) const fn private_pool(&self) -> &SqlitePool {
    116         &self.private_pool
    117     }
    118 
    119     pub(crate) const fn event_mode(&self) -> EventStoreMode {
    120         self.mode
    121     }
    122 
    123     async fn selected(
    124         &self,
    125         query: &EventQuery,
    126         minimum_stage: AdmissionStage,
    127     ) -> Result<(Vec<StoredEventRow>, Option<EventCursor>), Error> {
    128         self.validate_cursor(query)?;
    129         let after = query
    130             .bounds()
    131             .cursor()
    132             .map_or(0, |cursor| cursor.sequence().get());
    133         let fetch_limit = u64::from(query.bounds().limit()) + 1;
    134         let mut builder = QueryBuilder::<Sqlite>::new(
    135             "SELECT source_generation, source_sequence, signed_event, admission_stage, \
    136                     admitted_contract_id, admitted_registry_version \
    137              FROM radroots_runtime_events WHERE source_generation = ",
    138         );
    139         builder.push_bind(self.generation.as_bytes().as_slice());
    140         builder.push(" AND source_sequence > ");
    141         builder.push_bind(i64_from_u64(after)?);
    142         if minimum_stage == AdmissionStage::Verified {
    143             builder.push(" AND admission_stage IN ('verified', 'visible')");
    144         } else if minimum_stage == AdmissionStage::Visible {
    145             builder.push(" AND admission_stage = 'visible'");
    146         }
    147         if !query.event_ids().is_empty() {
    148             builder.push(" AND event_id IN (");
    149             let mut separated = builder.separated(", ");
    150             for event_id in query.event_ids() {
    151                 separated.push_bind(event_id.as_bytes().to_vec());
    152             }
    153             separated.push_unseparated(")");
    154         }
    155         builder.push(" ORDER BY source_sequence LIMIT ");
    156         builder.push_bind(i64_from_u64(fetch_limit)?);
    157 
    158         let rows = builder
    159             .build()
    160             .fetch_all(&self.pool)
    161             .await
    162             .map_err(map_backend)?;
    163         let mut decoded = rows
    164             .iter()
    165             .map(|row| self.decode_event_row(row))
    166             .collect::<Result<Vec<_>, _>>()?;
    167         let next = if decoded.len() > usize::from(query.bounds().limit()) {
    168             decoded.truncate(usize::from(query.bounds().limit()));
    169             decoded.last().map(|row| row.position)
    170         } else {
    171             None
    172         };
    173         Ok((decoded, next))
    174     }
    175 
    176     async fn current_event_rows(&self) -> Result<Vec<StoredEventRow>, Error> {
    177         let rows = sqlx::query(
    178             "SELECT source_generation, source_sequence, signed_event, admission_stage,
    179                     admitted_contract_id, admitted_registry_version
    180              FROM radroots_runtime_events
    181              WHERE source_generation = ?
    182              ORDER BY source_sequence",
    183         )
    184         .bind(self.generation.as_bytes().as_slice())
    185         .fetch_all(&self.pool)
    186         .await
    187         .map_err(map_backend)?;
    188         rows.iter().map(|row| self.decode_event_row(row)).collect()
    189     }
    190 
    191     fn visibility_for_rows(&self, rows: &[StoredEventRow]) -> Result<VisibilityEvaluation, Error> {
    192         evaluate_visibility(
    193             self.generation,
    194             rows.iter()
    195                 .map(|row| VisibilityInput::new(row.position, &row.event, row.stage)),
    196         )
    197     }
    198 
    199     fn validate_cursor(&self, query: &EventQuery) -> Result<(), Error> {
    200         if query
    201             .bounds()
    202             .cursor()
    203             .is_some_and(|cursor| cursor.generation() != self.generation)
    204         {
    205             return Err(Error::SourceGenerationChanged);
    206         }
    207         Ok(())
    208     }
    209 
    210     fn decode_event_row(&self, row: &sqlx::sqlite::SqliteRow) -> Result<StoredEventRow, Error> {
    211         let generation = source_generation(row.try_get("source_generation").map_err(map_corrupt)?)?;
    212         if generation != self.generation {
    213             return Err(Error::CorruptStoredEvent);
    214         }
    215         let sequence = event_sequence(row.try_get("source_sequence").map_err(map_corrupt)?)?;
    216         let raw_json = String::from_utf8(row.try_get("signed_event").map_err(map_corrupt)?)
    217             .map_err(|_| Error::CorruptStoredEvent)?;
    218         let event =
    219             Codec::decode_signed_event(raw_json.as_str()).map_err(|_| Error::CorruptStoredEvent)?;
    220         let stage = admission_stage(row.try_get("admission_stage").map_err(map_corrupt)?)?;
    221         let contract_id = row
    222             .try_get::<Option<String>, _>("admitted_contract_id")
    223             .map_err(map_corrupt)?;
    224         let registry_version = row
    225             .try_get::<Option<i64>, _>("admitted_registry_version")
    226             .map_err(map_corrupt)?
    227             .map(|value| u32::try_from(value).map_err(|_| Error::CorruptStoredEvent))
    228             .transpose()?;
    229         if contract_id.is_some() != registry_version.is_some()
    230             || registry_version.is_some_and(|version| {
    231                 version != radroots_event::contract::RegistryVersion::CURRENT.get()
    232             })
    233             || contract_id.as_deref().is_some_and(|stored| {
    234                 radroots_event::contract::registry_v7::validate_event_contract_for_admission(
    235                     event.envelope(),
    236                     stored,
    237                 )
    238                 .is_err()
    239             })
    240         {
    241             return Err(Error::CorruptStoredEvent);
    242         }
    243         Ok(StoredEventRow {
    244             position: EventPosition::new(generation, sequence),
    245             raw_json,
    246             event,
    247             stage,
    248         })
    249     }
    250 
    251     async fn store_provenance(
    252         transaction: &mut sqlx::Transaction<'_, Sqlite>,
    253         admission: &EventAdmission,
    254     ) -> Result<(), Error> {
    255         let provenance = admission.provenance();
    256         let cursor = provenance.cursor().map_or("", |cursor| cursor.as_str());
    257         sqlx::query(
    258             "INSERT OR IGNORE INTO radroots_runtime_event_provenance (
    259                event_id, transport_id, target_fingerprint, observed_at_unix_ms, cursor
    260              ) VALUES (?, ?, ?, ?, ?)",
    261         )
    262         .bind(admission.event_id().as_bytes().as_slice())
    263         .bind(provenance.transport_id().as_str())
    264         .bind(provenance.target().as_str())
    265         .bind(i64_from_u64(provenance.observed_at_unix_ms())?)
    266         .bind(cursor)
    267         .execute(&mut **transaction)
    268         .await
    269         .map_err(map_backend)?;
    270         Ok(())
    271     }
    272 
    273     pub(crate) async fn admit_transaction(
    274         &self,
    275         transaction: &mut sqlx::Transaction<'_, Sqlite>,
    276         admission: EventAdmission,
    277     ) -> Result<AdmissionReceipt, Error> {
    278         let existing = sqlx::query(
    279             "SELECT source_generation, source_sequence, signed_event, admission_stage,
    280                     admitted_contract_id, admitted_registry_version
    281              FROM radroots_runtime_events WHERE event_id = ?",
    282         )
    283         .bind(admission.event_id().as_bytes().as_slice())
    284         .fetch_optional(&mut **transaction)
    285         .await
    286         .map_err(map_backend)?;
    287 
    288         let (position, disposition) = if let Some(row) = existing {
    289             let stored = self.decode_event_row(&row)?;
    290             let metadata = admission_contract_metadata(&admission);
    291             if stored.raw_json.as_bytes() != admission.event().raw_json().as_bytes() {
    292                 return Err(Error::EventConflict);
    293             }
    294             if admission.stage() < stored.stage {
    295                 return Err(Error::AdmissionRegression);
    296             }
    297             let disposition = if admission.stage() == stored.stage {
    298                 AdmissionDisposition::Duplicate
    299             } else {
    300                 sqlx::query(
    301                     "UPDATE radroots_runtime_events
    302                      SET admission_stage = ?, updated_at_unix_ms = MAX(updated_at_unix_ms, ?),
    303                          admitted_contract_id = COALESCE(admitted_contract_id, ?),
    304                          admitted_registry_version = COALESCE(admitted_registry_version, ?)
    305                      WHERE event_id = ?",
    306                 )
    307                 .bind(stage_name(admission.stage()))
    308                 .bind(i64_from_u64(admission.provenance().observed_at_unix_ms())?)
    309                 .bind(metadata.map(|value| value.0))
    310                 .bind(metadata.map(|value| i64::from(value.1)))
    311                 .bind(admission.event_id().as_bytes().as_slice())
    312                 .execute(&mut **transaction)
    313                 .await
    314                 .map_err(map_backend)?;
    315                 AdmissionDisposition::Advanced
    316             };
    317             (stored.position, disposition)
    318         } else {
    319             let metadata = admission_contract_metadata(&admission);
    320             let next = sqlx::query_scalar::<_, i64>(
    321                 "UPDATE radroots_runtime_source_generations
    322                  SET sequence_head = sequence_head + 1
    323                  WHERE generation = ? AND state = 'active'
    324                  RETURNING sequence_head",
    325             )
    326             .bind(self.generation.as_bytes().as_slice())
    327             .fetch_optional(&mut **transaction)
    328             .await
    329             .map_err(map_backend)?
    330             .ok_or(Error::SourceGenerationChanged)?;
    331             let sequence = event_sequence(next)?;
    332             let observed_at = i64_from_u64(admission.provenance().observed_at_unix_ms())?;
    333             sqlx::query(
    334                 "INSERT INTO radroots_runtime_events (
    335                    source_generation, source_sequence, event_id, admission_stage,
    336                    signed_event, admitted_at_unix_ms, updated_at_unix_ms,
    337                    admitted_contract_id, admitted_registry_version
    338                  ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
    339             )
    340             .bind(self.generation.as_bytes().as_slice())
    341             .bind(next)
    342             .bind(admission.event_id().as_bytes().as_slice())
    343             .bind(stage_name(admission.stage()))
    344             .bind(admission.event().raw_json().as_bytes())
    345             .bind(observed_at)
    346             .bind(observed_at)
    347             .bind(metadata.map(|value| value.0))
    348             .bind(metadata.map(|value| i64::from(value.1)))
    349             .execute(&mut **transaction)
    350             .await
    351             .map_err(map_backend)?;
    352             (
    353                 EventPosition::new(self.generation, sequence),
    354                 AdmissionDisposition::Inserted,
    355             )
    356         };
    357         Self::store_provenance(transaction, &admission).await?;
    358         Ok(AdmissionReceipt::new(
    359             *admission.event_id(),
    360             position,
    361             admission.stage(),
    362             disposition,
    363         ))
    364     }
    365 }
    366 
    367 fn admission_contract_metadata(admission: &EventAdmission) -> Option<(&'static str, u32)> {
    368     admission.visible_event().map(|event| {
    369         (
    370             event.admitted_event().validated_event().contract_id(),
    371             radroots_event::contract::RegistryVersion::CURRENT.get(),
    372         )
    373     })
    374 }
    375 
    376 #[cfg_attr(coverage_nightly, coverage(off))]
    377 impl EventStore for SqliteStorage {
    378     fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> {
    379         Box::pin(async move {
    380             let row = sqlx::query(
    381                 "SELECT
    382                    COUNT(*) AS raw_events,
    383                    COALESCE(SUM(CASE WHEN admission_stage IN ('verified', 'visible') THEN 1 ELSE 0 END), 0)
    384                      AS verified_events
    385                  FROM radroots_runtime_events WHERE source_generation = ?",
    386             )
    387             .bind(self.generation.as_bytes().as_slice())
    388             .fetch_one(&self.pool)
    389             .await
    390             .map_err(map_backend)?;
    391             let visibility_rows = self.current_event_rows().await?;
    392             let visible_events = u64::try_from(
    393                 self.visibility_for_rows(&visibility_rows)?
    394                     .snapshot()
    395                     .visible_event_ids()
    396                     .len(),
    397             )
    398             .map_err(|_| Error::CorruptStoredEvent)?;
    399             EventStoreStatus::new(
    400                 self.generation,
    401                 self.mode,
    402                 EventStoreHealth::Available,
    403                 u64_from_i64(row.try_get("raw_events").map_err(map_corrupt)?)?,
    404                 u64_from_i64(row.try_get("verified_events").map_err(map_corrupt)?)?,
    405                 visible_events,
    406             )
    407         })
    408     }
    409 
    410     fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> {
    411         Box::pin(async move {
    412             if self.mode == EventStoreMode::ReadOnly {
    413                 return Err(Error::BackendUnavailable);
    414             }
    415             let mut transaction = self
    416                 .pool
    417                 .begin_with("BEGIN IMMEDIATE")
    418                 .await
    419                 .map_err(map_backend)?;
    420             let receipt = self.admit_transaction(&mut transaction, admission).await?;
    421             transaction.commit().await.map_err(map_backend)?;
    422             Ok(receipt)
    423         })
    424     }
    425 
    426     fn query_raw(
    427         &self,
    428         query: EventQuery,
    429     ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> {
    430         Box::pin(async move {
    431             let (rows, next) = self.selected(&query, AdmissionStage::Raw).await?;
    432             let items = rows
    433                 .into_iter()
    434                 .map(|row| StoredRawEvent::new(row.position, row.event, row.stage))
    435                 .collect();
    436             EventPage::new(self.generation, items, next, query.bounds())
    437         })
    438     }
    439 
    440     fn query_verified(
    441         &self,
    442         query: EventQuery,
    443     ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> {
    444         Box::pin(async move {
    445             let (rows, next) = self.selected(&query, AdmissionStage::Verified).await?;
    446             let items = rows
    447                 .into_iter()
    448                 .map(|row| StoredVerifiedEvent::new(row.position, row.event))
    449                 .collect();
    450             EventPage::new(self.generation, items, next, query.bounds())
    451         })
    452     }
    453 
    454     fn query_visible(
    455         &self,
    456         query: EventQuery,
    457     ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> {
    458         Box::pin(async move {
    459             self.validate_cursor(&query)?;
    460             let after = query
    461                 .bounds()
    462                 .cursor()
    463                 .map_or(0, |cursor| cursor.sequence().get());
    464             let rows = self.current_event_rows().await?;
    465             let visibility = self.visibility_for_rows(&rows)?;
    466             let mut selected = rows
    467                 .into_iter()
    468                 .filter(|row| {
    469                     row.position.sequence().get() > after
    470                         && query.selects(row.event.id())
    471                         && visibility.is_visible(row.event.id())
    472                 })
    473                 .take(usize::from(query.bounds().limit()) + 1)
    474                 .collect::<Vec<_>>();
    475             let next = if selected.len() > usize::from(query.bounds().limit()) {
    476                 selected.truncate(usize::from(query.bounds().limit()));
    477                 selected.last().map(|row| row.position)
    478             } else {
    479                 None
    480             };
    481             let items = selected
    482                 .into_iter()
    483                 .map(|row| StoredVisibleEvent::new(row.position, row.event))
    484                 .collect();
    485             EventPage::new(self.generation, items, next, query.bounds())
    486         })
    487     }
    488 
    489     fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> {
    490         Box::pin(async move {
    491             let rows = self.current_event_rows().await?;
    492             Ok(self.visibility_for_rows(&rows)?.into_snapshot())
    493         })
    494     }
    495 
    496     fn query_provenance(
    497         &self,
    498         event_id: EventId,
    499         bounds: EventQueryBounds,
    500     ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> {
    501         Box::pin(async move {
    502             if bounds
    503                 .cursor()
    504                 .is_some_and(|cursor| cursor.generation() != self.generation)
    505             {
    506                 return Err(Error::SourceGenerationChanged);
    507             }
    508             let event_row = sqlx::query(
    509                 "SELECT source_generation, source_sequence
    510                  FROM radroots_runtime_events WHERE event_id = ?",
    511             )
    512             .bind(event_id.as_bytes().as_slice())
    513             .fetch_optional(&self.pool)
    514             .await
    515             .map_err(map_backend)?
    516             .ok_or(Error::EventNotFound)?;
    517             let generation = source_generation(
    518                 event_row
    519                     .try_get("source_generation")
    520                     .map_err(map_corrupt)?,
    521             )?;
    522             if generation != self.generation {
    523                 return Err(Error::CorruptStoredEvent);
    524             }
    525             let position = EventPosition::new(
    526                 generation,
    527                 event_sequence(event_row.try_get("source_sequence").map_err(map_corrupt)?)?,
    528             );
    529             let after = bounds.cursor().map_or(0, |cursor| cursor.sequence().get());
    530             let items = if position.sequence().get() <= after {
    531                 Vec::new()
    532             } else {
    533                 let rows = sqlx::query(
    534                     "SELECT transport_id, target_fingerprint, observed_at_unix_ms, cursor
    535                      FROM radroots_runtime_event_provenance
    536                      WHERE event_id = ?
    537                      ORDER BY observed_at_unix_ms, transport_id, target_fingerprint, cursor
    538                      LIMIT ?",
    539                 )
    540                 .bind(event_id.as_bytes().as_slice())
    541                 .bind(i64::from(bounds.limit()))
    542                 .fetch_all(&self.pool)
    543                 .await
    544                 .map_err(map_backend)?;
    545                 rows.iter()
    546                     .map(|row| {
    547                         let cursor = row.try_get::<String, _>("cursor").map_err(map_corrupt)?;
    548                         StoredEventProvenance::from_stored_parts(
    549                             position,
    550                             row.try_get::<String, _>("transport_id")
    551                                 .map_err(map_corrupt)?
    552                                 .as_str(),
    553                             row.try_get::<String, _>("target_fingerprint")
    554                                 .map_err(map_corrupt)?
    555                                 .as_str(),
    556                             u64_from_i64(row.try_get("observed_at_unix_ms").map_err(map_corrupt)?)?,
    557                             (!cursor.is_empty()).then_some(cursor.as_str()),
    558                         )
    559                     })
    560                     .collect::<Result<Vec<_>, Error>>()?
    561             };
    562             EventPage::new(self.generation, items, None, bounds)
    563         })
    564     }
    565 }
    566 
    567 const fn stage_name(stage: AdmissionStage) -> &'static str {
    568     match stage {
    569         AdmissionStage::Raw => "raw",
    570         AdmissionStage::Verified => "verified",
    571         AdmissionStage::Visible => "visible",
    572     }
    573 }
    574 
    575 fn admission_stage(value: String) -> Result<AdmissionStage, Error> {
    576     match value.as_str() {
    577         "raw" => Ok(AdmissionStage::Raw),
    578         "verified" => Ok(AdmissionStage::Verified),
    579         "visible" => Ok(AdmissionStage::Visible),
    580         _ => Err(Error::CorruptStoredEvent),
    581     }
    582 }
    583 
    584 fn source_generation(value: Vec<u8>) -> Result<SourceGeneration, Error> {
    585     SourceGeneration::new(value.try_into().map_err(|_| Error::CorruptStoredEvent)?)
    586         .map_err(|_| Error::CorruptStoredEvent)
    587 }
    588 
    589 fn event_sequence(value: i64) -> Result<EventSequence, Error> {
    590     EventSequence::new(u64_from_i64(value)?).map_err(|_| Error::CorruptStoredEvent)
    591 }
    592 
    593 fn i64_from_u64(value: u64) -> Result<i64, Error> {
    594     i64::try_from(value).map_err(|_| Error::CorruptStoredEvent)
    595 }
    596 
    597 fn u64_from_i64(value: i64) -> Result<u64, Error> {
    598     u64::try_from(value).map_err(|_| Error::CorruptStoredEvent)
    599 }
    600 
    601 fn map_corrupt(_: sqlx::Error) -> Error {
    602     Error::CorruptStoredEvent
    603 }
    604 
    605 #[cfg(test)]
    606 #[cfg_attr(coverage_nightly, coverage(off))]
    607 mod tests {
    608     use super::*;
    609     use crate::migration::runtime::{MIGRATIONS, migration_sql};
    610     use radroots_event::{
    611         SignedEvent,
    612         admission::{AdmissionPolicy, RawEvent, VisibilityPolicy, VisibleEvent},
    613         wire::Nip01EventWire,
    614     };
    615     use radroots_storage::event::EventQueryBounds;
    616     use radroots_transport::{
    617         Target, TransportId,
    618         source::{EventProvenance, FetchCursor, ObservedEvent},
    619     };
    620     use sqlx::sqlite::SqlitePoolOptions;
    621 
    622     struct Allow;
    623 
    624     impl radroots_event::admission::SignatureVerifier for Allow {
    625         fn verify_signature(
    626             &self,
    627             _event: &radroots_event::Event,
    628         ) -> Result<(), radroots_event::Error> {
    629             Ok(())
    630         }
    631     }
    632 
    633     impl AdmissionPolicy for Allow {
    634         type Error = core::convert::Infallible;
    635 
    636         fn policy_id(&self) -> &'static str {
    637             "test.storage-sqlite.admission.v1"
    638         }
    639 
    640         fn admit(
    641             &self,
    642             _event: &radroots_event::admission::ContractValidatedEvent,
    643         ) -> Result<(), Self::Error> {
    644             Ok(())
    645         }
    646     }
    647 
    648     impl VisibilityPolicy for Allow {
    649         type Error = core::convert::Infallible;
    650 
    651         fn policy_id(&self) -> &'static str {
    652             "test.storage-sqlite.visibility.v1"
    653         }
    654 
    655         fn make_visible(
    656             &self,
    657             _event: &radroots_event::admission::AdmittedEvent,
    658         ) -> Result<(), Self::Error> {
    659             Ok(())
    660         }
    661     }
    662 
    663     async fn store(generation: SourceGeneration) -> SqliteStorage {
    664         let pool = SqlitePoolOptions::new()
    665             .max_connections(1)
    666             .connect("sqlite::memory:")
    667             .await
    668             .expect("memory SQLite");
    669         sqlx::query("PRAGMA foreign_keys = ON")
    670             .execute(&pool)
    671             .await
    672             .expect("foreign keys");
    673         for migration in MIGRATIONS {
    674             sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL"))
    675                 .execute(&pool)
    676                 .await
    677                 .expect("runtime migration");
    678         }
    679         sqlx::query(
    680             "INSERT INTO radroots_runtime_source_generations (
    681                generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms
    682              ) VALUES (?, 0, 'active', 1, NULL)",
    683         )
    684         .bind(generation.as_bytes().as_slice())
    685         .execute(&pool)
    686         .await
    687         .expect("source generation");
    688         SqliteStorage::new(pool, generation, EventStoreMode::ReadWrite)
    689     }
    690 
    691     fn signed_event(content: &str, pretty: bool) -> SignedEvent {
    692         signed_event_with(content, pretty, 1_800_000_100, 0, vec![])
    693     }
    694 
    695     fn signed_event_with(
    696         content: &str,
    697         pretty: bool,
    698         created_at: u64,
    699         kind: u32,
    700         tags: Vec<Vec<String>>,
    701     ) -> SignedEvent {
    702         let mut wire = Nip01EventWire {
    703             id: "0".repeat(64),
    704             pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
    705             created_at,
    706             kind,
    707             tags,
    708             content: content.to_owned(),
    709             sig: "42".repeat(64),
    710             extra: Default::default(),
    711         };
    712         wire.id = wire
    713             .computed_event_id()
    714             .expect("canonical event id")
    715             .to_hex();
    716         let value = serde_json::json!({
    717             "id": &wire.id,
    718             "pubkey": &wire.pubkey,
    719             "created_at": wire.created_at,
    720             "kind": wire.kind,
    721             "tags": &wire.tags,
    722             "content": &wire.content,
    723             "sig": &wire.sig,
    724         });
    725         let raw_json = if pretty {
    726             serde_json::to_string_pretty(&value).expect("pretty event JSON")
    727         } else {
    728             value.to_string()
    729         };
    730         SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event")
    731     }
    732 
    733     fn observed(event: SignedEvent, at: u64, cursor: Option<&str>) -> ObservedEvent {
    734         let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    735         let mut provenance =
    736             EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at)
    737                 .expect("provenance");
    738         if let Some(cursor) = cursor {
    739             provenance = provenance.with_cursor(FetchCursor::parse(cursor).expect("cursor"));
    740         }
    741         ObservedEvent::new(event, provenance)
    742     }
    743 
    744     fn verified(event: &SignedEvent) -> radroots_event::VerifiedEvent {
    745         RawEvent::new(event.envelope().clone())
    746             .verify_id()
    747             .expect("event id")
    748             .verify_signature(&Allow)
    749             .expect("signature")
    750     }
    751 
    752     fn visible(event: &SignedEvent) -> VisibleEvent {
    753         let verified = verified(event);
    754         let validated = if event.envelope().kind_u32() == 5 {
    755             verified
    756                 .validate_contract_for_admission("radroots.social.deletion_request.v1")
    757                 .expect("admission-selected contract")
    758         } else {
    759             verified.validate_contract().expect("contract")
    760         };
    761         validated
    762             .admit_with(&Allow)
    763             .expect("admission")
    764             .make_visible_with(&Allow)
    765             .expect("visibility")
    766     }
    767 
    768     #[tokio::test]
    769     async fn admission_is_idempotent_monotonic_and_conflict_safe() {
    770         let generation = SourceGeneration::new([7; 32]).expect("generation");
    771         let store = store(generation).await;
    772         let event = signed_event(
    773             "{\"display_name\":\"Moss Street Farm\",\"bot\":false}",
    774             false,
    775         );
    776 
    777         let inserted = store
    778             .admit(EventAdmission::raw(observed(event.clone(), 10, None)))
    779             .await
    780             .expect("insert raw");
    781         assert_eq!(inserted.disposition(), AdmissionDisposition::Inserted);
    782         assert_eq!(inserted.position().sequence().get(), 1);
    783 
    784         let advanced = store
    785             .admit(
    786                 EventAdmission::visible(
    787                     observed(event.clone(), 11, Some("relay-page-1")),
    788                     visible(&event),
    789                 )
    790                 .expect("visible admission"),
    791             )
    792             .await
    793             .expect("advance visible");
    794         assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced);
    795         assert_eq!(advanced.position(), inserted.position());
    796         let metadata = sqlx::query(
    797             "SELECT admitted_contract_id, admitted_registry_version
    798              FROM radroots_runtime_events WHERE event_id = ?",
    799         )
    800         .bind(event.id().as_bytes().as_slice())
    801         .fetch_one(&store.pool)
    802         .await
    803         .expect("event contract metadata");
    804         assert_eq!(
    805             metadata.get::<String, _>("admitted_contract_id"),
    806             "radroots.profile.metadata.v1"
    807         );
    808         assert_eq!(metadata.get::<i64, _>("admitted_registry_version"), 7);
    809 
    810         let duplicate = store
    811             .admit(
    812                 EventAdmission::visible(
    813                     observed(event.clone(), 11, Some("relay-page-1")),
    814                     visible(&event),
    815                 )
    816                 .expect("visible admission"),
    817             )
    818             .await
    819             .expect("duplicate visible");
    820         assert_eq!(duplicate.disposition(), AdmissionDisposition::Duplicate);
    821         assert_eq!(
    822             store
    823                 .admit(EventAdmission::raw(observed(event.clone(), 12, None)))
    824                 .await,
    825             Err(Error::AdmissionRegression)
    826         );
    827 
    828         let same_id_different_bytes = signed_event(
    829             "{\"display_name\":\"Moss Street Farm\",\"bot\":false}",
    830             true,
    831         );
    832         assert_eq!(same_id_different_bytes.id(), event.id());
    833         assert_eq!(
    834             store
    835                 .admit(EventAdmission::raw(observed(
    836                     same_id_different_bytes,
    837                     13,
    838                     None,
    839                 )))
    840                 .await,
    841             Err(Error::EventConflict)
    842         );
    843     }
    844 
    845     #[tokio::test]
    846     async fn queries_preserve_stage_bounds_cursors_and_exact_provenance() {
    847         let generation = SourceGeneration::new([8; 32]).expect("generation");
    848         let store = store(generation).await;
    849         let empty_status = store.status().await.expect("empty status");
    850         assert_eq!(empty_status.raw_events(), 0);
    851         assert_eq!(empty_status.verified_events(), 0);
    852         assert_eq!(empty_status.visible_events(), 0);
    853         let raw_event = signed_event("raw", false);
    854         let visible_event =
    855             signed_event("{\"display_name\":\"Visible Farm\",\"bot\":false}", false);
    856         store
    857             .admit(EventAdmission::raw(observed(raw_event.clone(), 20, None)))
    858             .await
    859             .expect("raw event");
    860         store
    861             .admit(
    862                 EventAdmission::visible(
    863                     observed(visible_event.clone(), 21, Some("relay-page-2")),
    864                     visible(&visible_event),
    865                 )
    866                 .expect("visible admission"),
    867             )
    868             .await
    869             .expect("visible event");
    870 
    871         let first = store
    872             .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds")))
    873             .await
    874             .expect("first page");
    875         assert_eq!(first.items().len(), 1);
    876         assert_eq!(first.items()[0].event(), &raw_event);
    877         let next = first.next_cursor().expect("continuation cursor");
    878         let second = store
    879             .query_raw(EventQuery::all(
    880                 EventQueryBounds::first(1).expect("bounds").after(next),
    881             ))
    882             .await
    883             .expect("second page");
    884         assert_eq!(second.items()[0].event(), &visible_event);
    885         assert!(second.next_cursor().is_none());
    886 
    887         let verified_page = store
    888             .query_verified(EventQuery::all(
    889                 EventQueryBounds::first(10).expect("bounds"),
    890             ))
    891             .await
    892             .expect("verified page");
    893         assert_eq!(verified_page.items().len(), 1);
    894         assert_eq!(verified_page.items()[0].event(), &visible_event);
    895         let visible_page = store
    896             .query_visible(
    897                 EventQuery::for_ids(
    898                     EventQueryBounds::first(10).expect("bounds"),
    899                     vec![*visible_event.id()],
    900                 )
    901                 .expect("id query"),
    902             )
    903             .await
    904             .expect("visible page");
    905         assert_eq!(visible_page.items().len(), 1);
    906 
    907         let provenance = store
    908             .query_provenance(
    909                 *visible_event.id(),
    910                 EventQueryBounds::first(10).expect("bounds"),
    911             )
    912             .await
    913             .expect("provenance");
    914         assert_eq!(provenance.items().len(), 1);
    915         assert_eq!(provenance.items()[0].provenance().observed_at_unix_ms(), 21);
    916         assert_eq!(
    917             provenance.items()[0]
    918                 .provenance()
    919                 .cursor()
    920                 .expect("cursor")
    921                 .as_str(),
    922             "relay-page-2"
    923         );
    924 
    925         let status = store.status().await.expect("status");
    926         assert_eq!(status.raw_events(), 2);
    927         assert_eq!(status.verified_events(), 1);
    928         assert_eq!(status.visible_events(), 1);
    929         let foreign_cursor = EventPosition::new(
    930             SourceGeneration::new([9; 32]).expect("foreign generation"),
    931             EventSequence::new(1).expect("sequence"),
    932         );
    933         assert_eq!(
    934             store
    935                 .query_raw(EventQuery::all(
    936                     EventQueryBounds::first(1)
    937                         .expect("bounds")
    938                         .after(foreign_cursor),
    939                 ))
    940                 .await,
    941             Err(Error::SourceGenerationChanged)
    942         );
    943     }
    944 
    945     #[tokio::test]
    946     async fn verified_replacement_and_same_count_advance_match_memory_after_reopen() {
    947         let generation = SourceGeneration::new([19; 32]).unwrap();
    948         let store = store(generation).await;
    949         let memory = radroots_storage::memory::MemoryStorage::new(generation);
    950         let old = signed_event_with(
    951             r#"{"display_name":"Old Farm","bot":false}"#,
    952             false,
    953             10,
    954             0,
    955             vec![],
    956         );
    957         let newer = signed_event_with("malformed profile", false, 20, 0, vec![]);
    958         let admissions = [
    959             EventAdmission::visible(observed(old.clone(), 10, None), visible(&old)).unwrap(),
    960             EventAdmission::raw(observed(newer.clone(), 20, None)),
    961             EventAdmission::verified(observed(newer.clone(), 20, None), verified(&newer)).unwrap(),
    962         ];
    963         let mut previous_digest = None;
    964         for (index, admission) in admissions.into_iter().enumerate() {
    965             assert_eq!(
    966                 store.admit(admission.clone()).await.unwrap(),
    967                 memory.admit(admission).await.unwrap()
    968             );
    969             let reopened =
    970                 SqliteStorage::new(store.pool.clone(), generation, EventStoreMode::ReadWrite);
    971             let snapshot = reopened.rebuild_visibility().await.unwrap();
    972             assert_eq!(snapshot, memory.rebuild_visibility().await.unwrap());
    973             let query = EventQuery::all(EventQueryBounds::first(10).unwrap());
    974             assert_eq!(
    975                 reopened.query_visible(query.clone()).await.unwrap(),
    976                 memory.query_visible(query).await.unwrap()
    977             );
    978             if index == 2 {
    979                 assert_ne!(previous_digest, Some(snapshot.digest()));
    980                 assert!(snapshot.visible_event_ids().is_empty());
    981                 assert_eq!(snapshot.current_heads()[0].event_id, *newer.id());
    982                 assert_eq!(reopened.status().await.unwrap().raw_events(), 2);
    983             } else {
    984                 assert_eq!(snapshot.visible_event_ids(), &[*old.id()]);
    985             }
    986             previous_digest = Some(snapshot.digest());
    987         }
    988     }
    989 
    990     #[tokio::test]
    991     async fn visibility_rebuild_survives_reopen_and_matches_current_head_deletion_queries() {
    992         let generation = SourceGeneration::new([18; 32]).expect("generation");
    993         let store = store(generation).await;
    994         let old = signed_event_with(
    995             r#"{"display_name":"Old Farm","bot":false}"#,
    996             false,
    997             1_800_000_100,
    998             0,
    999             vec![],
   1000         );
   1001         let current = signed_event_with(
   1002             r#"{"display_name":"Current Farm","bot":false}"#,
   1003             false,
   1004             1_800_000_200,
   1005             0,
   1006             vec![],
   1007         );
   1008         let deletion = signed_event_with(
   1009             "retired profile",
   1010             false,
   1011             1_800_000_300,
   1012             5,
   1013             vec![vec!["e".to_owned(), current.id().to_hex()]],
   1014         );
   1015         for (event, observed_at) in [
   1016             (old.clone(), 100),
   1017             (current.clone(), 200),
   1018             (deletion.clone(), 300),
   1019         ] {
   1020             store
   1021                 .admit(
   1022                     EventAdmission::visible(
   1023                         observed(event.clone(), observed_at, None),
   1024                         visible(&event),
   1025                     )
   1026                     .expect("visible admission"),
   1027                 )
   1028                 .await
   1029                 .expect("admit visible event");
   1030         }
   1031 
   1032         let before = store
   1033             .rebuild_visibility()
   1034             .await
   1035             .expect("visibility rebuild");
   1036         let reopened =
   1037             SqliteStorage::new(store.pool.clone(), generation, EventStoreMode::ReadWrite);
   1038         let after = reopened
   1039             .rebuild_visibility()
   1040             .await
   1041             .expect("reopened visibility rebuild");
   1042         assert_eq!(before, after);
   1043         assert_eq!(after.current_heads()[0].event_id, *current.id());
   1044         assert_eq!(after.visible_event_ids(), &[*deletion.id()]);
   1045         assert_eq!(after.suppressed_event_ids(), &[*current.id()]);
   1046         assert_eq!(after.superseded_event_ids(), &[*old.id()]);
   1047         let page = reopened
   1048             .query_visible(EventQuery::all(
   1049                 EventQueryBounds::first(10).expect("bounds"),
   1050             ))
   1051             .await
   1052             .expect("visible page");
   1053         assert_eq!(page.items().len(), 1);
   1054         assert_eq!(page.items()[0].event().id(), deletion.id());
   1055         assert_eq!(reopened.status().await.expect("status").visible_events(), 1);
   1056     }
   1057 
   1058     #[tokio::test]
   1059     async fn corrupt_rows_fail_closed_and_source_history_is_immutable() {
   1060         let generation = SourceGeneration::new([10; 32]).expect("generation");
   1061         let store = store(generation).await;
   1062         let event = signed_event(
   1063             "{\"display_name\":\"Corruption Probe\",\"bot\":false}",
   1064             false,
   1065         );
   1066         store
   1067             .admit(EventAdmission::raw(observed(event.clone(), 30, None)))
   1068             .await
   1069             .expect("event");
   1070 
   1071         let wrong_generation = SourceGeneration::new([11; 32]).expect("generation");
   1072         let mismatched_store = SqliteStorage::new(
   1073             store.pool.clone(),
   1074             wrong_generation,
   1075             EventStoreMode::ReadWrite,
   1076         );
   1077         assert_eq!(
   1078             mismatched_store
   1079                 .admit(EventAdmission::raw(observed(event, 31, None)))
   1080                 .await,
   1081             Err(Error::CorruptStoredEvent)
   1082         );
   1083 
   1084         assert!(
   1085             sqlx::query("DELETE FROM radroots_runtime_events")
   1086                 .execute(&store.pool)
   1087                 .await
   1088                 .is_err()
   1089         );
   1090         assert!(
   1091             sqlx::query("DELETE FROM radroots_runtime_source_generations")
   1092                 .execute(&store.pool)
   1093                 .await
   1094                 .is_err()
   1095         );
   1096         sqlx::query("DROP TRIGGER radroots_runtime_events_contract_metadata_guard")
   1097             .execute(&store.pool)
   1098             .await
   1099             .expect("drop metadata guard for corruption probe");
   1100         sqlx::query(
   1101             "UPDATE radroots_runtime_events
   1102              SET admitted_contract_id = 'radroots.social.geochat.v1'",
   1103         )
   1104         .execute(&store.pool)
   1105         .await
   1106         .expect("forge contract metadata");
   1107         assert_eq!(
   1108             store
   1109                 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
   1110                 .await,
   1111             Err(Error::CorruptStoredEvent)
   1112         );
   1113         sqlx::query(
   1114             "UPDATE radroots_runtime_events
   1115              SET admitted_contract_id = 'radroots.social.geochat.v1',
   1116                  admitted_registry_version = 7",
   1117         )
   1118         .execute(&store.pool)
   1119         .await
   1120         .expect("forge mismatched selected contract");
   1121         assert_eq!(
   1122             store
   1123                 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
   1124                 .await,
   1125             Err(Error::CorruptStoredEvent)
   1126         );
   1127         sqlx::query(
   1128             "UPDATE radroots_runtime_events
   1129              SET admitted_registry_version = 6",
   1130         )
   1131         .execute(&store.pool)
   1132         .await
   1133         .expect("forge registry version");
   1134         assert_eq!(
   1135             store
   1136                 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
   1137                 .await,
   1138             Err(Error::CorruptStoredEvent)
   1139         );
   1140         sqlx::query("PRAGMA ignore_check_constraints = ON")
   1141             .execute(&store.pool)
   1142             .await
   1143             .expect("disable checks for corruption probe");
   1144         sqlx::query("UPDATE radroots_runtime_events SET admission_stage = 'corrupt'")
   1145             .execute(&store.pool)
   1146             .await
   1147             .expect("forge corrupt stage");
   1148         assert_eq!(
   1149             store
   1150                 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
   1151                 .await,
   1152             Err(Error::CorruptStoredEvent)
   1153         );
   1154     }
   1155 }