lib

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

mod.rs (70605B)


      1 use crate::SqliteStorage;
      2 use crate::backend::map_backend;
      3 use radroots_storage::{
      4     Error, ProjectionStore,
      5     event::{EventPosition, SourceGeneration},
      6     projection::{
      7         ArtifactDigest, BoxFuture, EVENT_INDEX_SHARDS_MAX, EventId, EventIdRange,
      8         EventIndexCheckpoint, EventIndexManifest, EventIndexShard, EventIndexShardCheckpoint,
      9         EventIndexShardId, InvalidationReason, ProjectionCheckpoint, ProjectionDocument,
     10         ProjectionGeneration, ProjectionHealth, ProjectionId, ProjectionInvalidation,
     11         ProjectionRevision, ProjectionSnapshot, ProjectionStatus, RawSourceDigest, RebuildFailure,
     12         RebuildStage, RebuildTicket, RebuildTicketId, RebuildTransition,
     13     },
     14 };
     15 use sqlx::{Row, Sqlite, SqliteConnection};
     16 
     17 mod document_query;
     18 
     19 #[cfg_attr(coverage_nightly, coverage(off))]
     20 impl ProjectionStore for SqliteStorage {
     21     fn status(
     22         &self,
     23         projection_id: ProjectionId,
     24     ) -> BoxFuture<'_, Result<Option<ProjectionStatus>, Error>> {
     25         Box::pin(async move {
     26             sqlx::query(
     27                 "SELECT * FROM radroots_runtime_projection_checkpoints WHERE projection_id = ?",
     28             )
     29             .bind(projection_id.as_str())
     30             .fetch_optional(self.pool())
     31             .await
     32             .map_err(map_backend)?
     33             .as_ref()
     34             .map(decode_status)
     35             .transpose()
     36         })
     37     }
     38 
     39     fn checkpoint(
     40         &self,
     41         checkpoint: ProjectionCheckpoint,
     42     ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> {
     43         Box::pin(async move {
     44             self.require_projection_writer()?;
     45             let mut transaction = self
     46                 .pool()
     47                 .begin_with("BEGIN IMMEDIATE")
     48                 .await
     49                 .map_err(map_backend)?;
     50             let status = checkpoint_transaction(&mut transaction, checkpoint).await?;
     51             transaction.commit().await.map_err(map_backend)?;
     52             Ok(status)
     53         })
     54     }
     55 
     56     fn invalidate(
     57         &self,
     58         invalidation: ProjectionInvalidation,
     59     ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> {
     60         Box::pin(async move {
     61             self.require_projection_writer()?;
     62             let mut transaction = self
     63                 .pool()
     64                 .begin_with("BEGIN IMMEDIATE")
     65                 .await
     66                 .map_err(map_backend)?;
     67             let row = sqlx::query(
     68                 "SELECT * FROM radroots_runtime_projection_checkpoints WHERE projection_id = ?",
     69             )
     70             .bind(invalidation.projection_id().as_str())
     71             .fetch_optional(&mut *transaction)
     72             .await
     73             .map_err(map_backend)?
     74             .ok_or(Error::ProjectionCheckpointMismatch)?;
     75             let current = decode_status(&row)?;
     76             if current.generation() != invalidation.invalid_generation() {
     77                 return Err(Error::ProjectionCheckpointMismatch);
     78             }
     79             if let Some(existing) = load_invalidation(
     80                 &mut transaction,
     81                 invalidation.projection_id(),
     82                 invalidation.invalid_generation(),
     83             )
     84             .await?
     85             {
     86                 if existing != invalidation {
     87                     return Err(Error::ProjectionRevisionConflict);
     88                 }
     89             } else {
     90                 sqlx::query(
     91                     "INSERT INTO radroots_runtime_projection_invalidations (
     92                        projection_id, invalid_generation, replacement_generation, reason,
     93                        invalidated_at_unix_ms
     94                      ) VALUES (?, ?, ?, ?, ?)",
     95                 )
     96                 .bind(invalidation.projection_id().as_str())
     97                 .bind(invalidation.invalid_generation().as_bytes().as_slice())
     98                 .bind(invalidation.replacement_generation().as_bytes().as_slice())
     99                 .bind(reason_name(invalidation.reason()))
    100                 .bind(i64_from_u64(invalidation.invalidated_at_unix_ms())?)
    101                 .execute(&mut *transaction)
    102                 .await
    103                 .map_err(map_backend)?;
    104             }
    105             let next = ProjectionStatus::new(
    106                 invalidation.projection_id().clone(),
    107                 current.generation(),
    108                 ProjectionHealth::Invalidated,
    109                 current.checkpoint().cloned(),
    110                 None,
    111             )?;
    112             put_status_transaction(&mut transaction, &next).await?;
    113             transaction.commit().await.map_err(map_backend)?;
    114             Ok(next)
    115         })
    116     }
    117 
    118     fn request_rebuild(
    119         &self,
    120         ticket: RebuildTicket,
    121     ) -> BoxFuture<'_, Result<RebuildTicket, Error>> {
    122         Box::pin(async move {
    123             self.require_projection_writer()?;
    124             let mut transaction = self
    125                 .pool()
    126                 .begin_with("BEGIN IMMEDIATE")
    127                 .await
    128                 .map_err(map_backend)?;
    129             if let Some(row) = sqlx::query(
    130                 "SELECT * FROM radroots_runtime_projection_rebuilds WHERE ticket_id = ?",
    131             )
    132             .bind(ticket.ticket_id().as_bytes().as_slice())
    133             .fetch_optional(&mut *transaction)
    134             .await
    135             .map_err(map_backend)?
    136             {
    137                 let existing = decode_ticket(&mut transaction, &row).await?;
    138                 return if existing == ticket {
    139                     transaction.commit().await.map_err(map_backend)?;
    140                     Ok(existing)
    141                 } else {
    142                     Err(Error::ProjectionRevisionConflict)
    143                 };
    144             }
    145             let status_row = sqlx::query(
    146                 "SELECT * FROM radroots_runtime_projection_checkpoints WHERE projection_id = ?",
    147             )
    148             .bind(ticket.invalidation().projection_id().as_str())
    149             .fetch_optional(&mut *transaction)
    150             .await
    151             .map_err(map_backend)?
    152             .ok_or(Error::ProjectionCheckpointMismatch)?;
    153             let status = decode_status(&status_row)?;
    154             if status.generation() != ticket.invalidation().invalid_generation()
    155                 || status.health() != ProjectionHealth::Invalidated
    156                 || load_invalidation(
    157                     &mut transaction,
    158                     ticket.invalidation().projection_id(),
    159                     ticket.invalidation().invalid_generation(),
    160                 )
    161                 .await?
    162                 .as_ref()
    163                     != Some(ticket.invalidation())
    164             {
    165                 return Err(Error::ProjectionCheckpointMismatch);
    166             }
    167             insert_ticket(&mut transaction, &ticket).await?;
    168             let next = ProjectionStatus::new(
    169                 status.projection_id().clone(),
    170                 status.generation(),
    171                 ProjectionHealth::Rebuilding,
    172                 status.checkpoint().cloned(),
    173                 Some(ticket.ticket_id()),
    174             )?;
    175             put_status_transaction(&mut transaction, &next).await?;
    176             transaction.commit().await.map_err(map_backend)?;
    177             Ok(ticket)
    178         })
    179     }
    180 
    181     fn invalidation(
    182         &self,
    183         projection_id: ProjectionId,
    184         replacement_generation: ProjectionGeneration,
    185     ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>> {
    186         Box::pin(async move {
    187             sqlx::query(
    188                 "SELECT * FROM radroots_runtime_projection_invalidations
    189                  WHERE projection_id = ? AND replacement_generation = ?
    190                  ORDER BY invalidated_at_unix_ms DESC LIMIT 1",
    191             )
    192             .bind(projection_id.as_str())
    193             .bind(replacement_generation.as_bytes().as_slice())
    194             .fetch_optional(self.pool())
    195             .await
    196             .map_err(map_backend)?
    197             .as_ref()
    198             .map(decode_invalidation)
    199             .transpose()
    200         })
    201     }
    202 
    203     fn rebuild(
    204         &self,
    205         ticket_id: RebuildTicketId,
    206     ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>> {
    207         Box::pin(async move {
    208             let mut connection = self.pool().acquire().await.map_err(map_backend)?;
    209             let row = sqlx::query(
    210                 "SELECT * FROM radroots_runtime_projection_rebuilds WHERE ticket_id = ?",
    211             )
    212             .bind(ticket_id.as_bytes().as_slice())
    213             .fetch_optional(&mut *connection)
    214             .await
    215             .map_err(map_backend)?;
    216             match row {
    217                 Some(row) => decode_ticket(&mut connection, &row).await.map(Some),
    218                 None => Ok(None),
    219             }
    220         })
    221     }
    222 
    223     fn transition_rebuild(
    224         &self,
    225         transition: RebuildTransition,
    226     ) -> BoxFuture<'_, Result<RebuildTicket, Error>> {
    227         Box::pin(async move {
    228             self.require_projection_writer()?;
    229             let mut transaction = self
    230                 .pool()
    231                 .begin_with("BEGIN IMMEDIATE")
    232                 .await
    233                 .map_err(map_backend)?;
    234             let row = sqlx::query(
    235                 "SELECT * FROM radroots_runtime_projection_rebuilds WHERE ticket_id = ?",
    236             )
    237             .bind(transition.ticket_id().as_bytes().as_slice())
    238             .fetch_optional(&mut *transaction)
    239             .await
    240             .map_err(map_backend)?
    241             .ok_or(Error::ProjectionRevisionConflict)?;
    242             let current = decode_ticket(&mut transaction, &row).await?;
    243             let status_row = sqlx::query(
    244                 "SELECT * FROM radroots_runtime_projection_checkpoints WHERE projection_id = ?",
    245             )
    246             .bind(current.invalidation().projection_id().as_str())
    247             .fetch_optional(&mut *transaction)
    248             .await
    249             .map_err(map_backend)?
    250             .ok_or(Error::CorruptProjectionRecord)?;
    251             let current_status = decode_status(&status_row)?;
    252             if current_status.generation() != current.invalidation().invalid_generation()
    253                 || current_status.health() != ProjectionHealth::Rebuilding
    254                 || current_status.active_rebuild() != Some(current.ticket_id())
    255             {
    256                 return Err(Error::CorruptProjectionRecord);
    257             }
    258             let next = current.transition(transition)?;
    259             if next.stage() == RebuildStage::Completed {
    260                 let source = sqlx::query(
    261                     "SELECT generation, sequence_head
    262                      FROM radroots_runtime_source_generations WHERE state = 'active'",
    263                 )
    264                 .fetch_one(&mut *transaction)
    265                 .await
    266                 .map_err(map_backend)?;
    267                 let generation = SourceGeneration::new(array(
    268                     source
    269                         .try_get::<Vec<u8>, _>("generation")
    270                         .map_err(map_corrupt)?,
    271                 )?)
    272                 .map_err(|_| Error::CorruptProjectionRecord)?;
    273                 let sequence = u64_from_i64(
    274                     source
    275                         .try_get::<i64, _>("sequence_head")
    276                         .map_err(map_corrupt)?,
    277                 )?;
    278                 if generation != next.source_generation()
    279                     || sequence
    280                         != next
    281                             .source_high_water()
    282                             .map_or(0, |position| position.sequence().get())
    283                 {
    284                     return Err(Error::SourceGenerationChanged);
    285                 }
    286             }
    287             update_ticket(&mut transaction, &next, current.revision()).await?;
    288             let (generation, health, checkpoint, active_rebuild) = match next.stage() {
    289                 RebuildStage::Requested | RebuildStage::Running => (
    290                     current_status.generation(),
    291                     ProjectionHealth::Rebuilding,
    292                     current_status.checkpoint().cloned(),
    293                     Some(next.ticket_id()),
    294                 ),
    295                 RebuildStage::Completed => (
    296                     next.invalidation().replacement_generation(),
    297                     ProjectionHealth::Ready,
    298                     next.checkpoint().cloned(),
    299                     None,
    300                 ),
    301                 RebuildStage::Failed => (
    302                     current_status.generation(),
    303                     ProjectionHealth::Ready,
    304                     current_status.checkpoint().cloned(),
    305                     None,
    306                 ),
    307             };
    308             let status = ProjectionStatus::new(
    309                 next.invalidation().projection_id().clone(),
    310                 generation,
    311                 health,
    312                 checkpoint,
    313                 active_rebuild,
    314             )?;
    315             put_status_transaction(&mut transaction, &status).await?;
    316             transaction.commit().await.map_err(map_backend)?;
    317             Ok(next)
    318         })
    319     }
    320 
    321     fn event_index_manifest(
    322         &self,
    323         generation: ProjectionGeneration,
    324     ) -> BoxFuture<'_, Result<Option<EventIndexManifest>, Error>> {
    325         Box::pin(async move {
    326             let mut connection = self.pool().acquire().await.map_err(map_backend)?;
    327             load_manifest(&mut connection, generation).await
    328         })
    329     }
    330 
    331     fn put_event_index_manifest(
    332         &self,
    333         manifest: EventIndexManifest,
    334     ) -> BoxFuture<'_, Result<(), Error>> {
    335         Box::pin(async move {
    336             self.require_projection_writer()?;
    337             let mut transaction = self
    338                 .pool()
    339                 .begin_with("BEGIN IMMEDIATE")
    340                 .await
    341                 .map_err(map_backend)?;
    342             if let Some(existing) = load_manifest(&mut transaction, manifest.generation()).await? {
    343                 return if existing == manifest {
    344                     transaction.commit().await.map_err(map_backend)?;
    345                     Ok(())
    346                 } else {
    347                     Err(Error::CorruptProjectionRecord)
    348                 };
    349             }
    350             sqlx::query(
    351                 "INSERT INTO radroots_runtime_event_index_manifests (
    352                    projection_generation, total_events, target_shard_size,
    353                    first_published_at_unix_s, last_published_at_unix_s
    354                  ) VALUES (?, ?, ?, ?, ?)",
    355             )
    356             .bind(manifest.generation().as_bytes().as_slice())
    357             .bind(i64_from_u64(manifest.total_events())?)
    358             .bind(i64::from(manifest.target_shard_size()))
    359             .bind(i64_from_u64(manifest.first_published_at_unix_s())?)
    360             .bind(i64_from_u64(manifest.last_published_at_unix_s())?)
    361             .execute(&mut *transaction)
    362             .await
    363             .map_err(map_backend)?;
    364             for (ordinal, shard) in manifest.shards().iter().enumerate() {
    365                 sqlx::query(
    366                     "INSERT INTO radroots_runtime_event_index_shards (
    367                        projection_generation, shard_id, ordinal, artifact_path, event_count,
    368                        first_event_id, last_event_id, first_published_at_unix_s,
    369                        last_published_at_unix_s, artifact_digest
    370                      ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
    371                 )
    372                 .bind(manifest.generation().as_bytes().as_slice())
    373                 .bind(shard.shard_id().as_str())
    374                 .bind(i64::try_from(ordinal).map_err(|_| Error::CorruptProjectionRecord)?)
    375                 .bind(shard.artifact_path())
    376                 .bind(i64::from(shard.event_count()))
    377                 .bind(shard.event_ids().first().as_bytes().as_slice())
    378                 .bind(shard.event_ids().last().as_bytes().as_slice())
    379                 .bind(i64_from_u64(shard.first_published_at_unix_s())?)
    380                 .bind(i64_from_u64(shard.last_published_at_unix_s())?)
    381                 .bind(shard.sha256().as_bytes().as_slice())
    382                 .execute(&mut *transaction)
    383                 .await
    384                 .map_err(map_backend)?;
    385             }
    386             transaction.commit().await.map_err(map_backend)
    387         })
    388     }
    389 
    390     fn event_index_checkpoint(
    391         &self,
    392         generation: ProjectionGeneration,
    393     ) -> BoxFuture<'_, Result<Option<EventIndexCheckpoint>, Error>> {
    394         Box::pin(async move {
    395             sqlx::query(
    396                 "SELECT generated_at_unix_ms, checkpoint
    397                  FROM radroots_runtime_event_index_checkpoints
    398                  WHERE projection_generation = ?",
    399             )
    400             .bind(generation.as_bytes().as_slice())
    401             .fetch_optional(self.pool())
    402             .await
    403             .map_err(map_backend)?
    404             .as_ref()
    405             .map(|row| decode_index_checkpoint(row, generation))
    406             .transpose()
    407         })
    408     }
    409 
    410     fn put_event_index_checkpoint(
    411         &self,
    412         checkpoint: EventIndexCheckpoint,
    413     ) -> BoxFuture<'_, Result<(), Error>> {
    414         Box::pin(async move {
    415             self.require_projection_writer()?;
    416             let mut transaction = self
    417                 .pool()
    418                 .begin_with("BEGIN IMMEDIATE")
    419                 .await
    420                 .map_err(map_backend)?;
    421             if let Some(row) = sqlx::query(
    422                 "SELECT generated_at_unix_ms, checkpoint
    423                  FROM radroots_runtime_event_index_checkpoints
    424                  WHERE projection_generation = ?",
    425             )
    426             .bind(checkpoint.generation().as_bytes().as_slice())
    427             .fetch_optional(&mut *transaction)
    428             .await
    429             .map_err(map_backend)?
    430                 && checkpoint.generated_at_unix_ms()
    431                     < decode_index_checkpoint(&row, checkpoint.generation())?.generated_at_unix_ms()
    432             {
    433                 return Err(Error::InvalidEventIndexCheckpoint);
    434             }
    435             sqlx::query(
    436                 "INSERT INTO radroots_runtime_event_index_checkpoints (
    437                    projection_generation, generated_at_unix_ms, checkpoint
    438                  ) VALUES (?, ?, ?)
    439                  ON CONFLICT(projection_generation) DO UPDATE SET
    440                    generated_at_unix_ms = excluded.generated_at_unix_ms,
    441                    checkpoint = excluded.checkpoint",
    442             )
    443             .bind(checkpoint.generation().as_bytes().as_slice())
    444             .bind(i64_from_u64(checkpoint.generated_at_unix_ms())?)
    445             .bind(encode_index_checkpoint(&checkpoint)?)
    446             .execute(&mut *transaction)
    447             .await
    448             .map_err(map_backend)?;
    449             transaction.commit().await.map_err(map_backend)?;
    450             Ok(())
    451         })
    452     }
    453 
    454     fn put_projection_document(
    455         &self,
    456         projection_id: ProjectionId,
    457         generation: ProjectionGeneration,
    458         document: ProjectionDocument,
    459     ) -> BoxFuture<'_, Result<(), Error>> {
    460         Box::pin(async move {
    461             self.require_projection_writer()?;
    462             sqlx::query(
    463                 "INSERT INTO radroots_runtime_projection_documents (
    464                    projection_id, generation, document_key, value, value_sha256
    465                  ) VALUES (?, ?, ?, ?, ?)
    466                  ON CONFLICT(projection_id, generation, document_key) DO UPDATE SET
    467                    value = excluded.value,
    468                    value_sha256 = excluded.value_sha256",
    469             )
    470             .bind(projection_id.as_str())
    471             .bind(generation.as_bytes().as_slice())
    472             .bind(document.key())
    473             .bind(document.value())
    474             .bind(document.value_sha256().as_slice())
    475             .execute(self.pool())
    476             .await
    477             .map_err(map_backend)?;
    478             Ok(())
    479         })
    480     }
    481 
    482     fn projection_document(
    483         &self,
    484         projection_id: ProjectionId,
    485         generation: ProjectionGeneration,
    486         key: String,
    487     ) -> BoxFuture<'_, Result<Option<ProjectionDocument>, Error>> {
    488         Box::pin(async move {
    489             sqlx::query(
    490                 "SELECT document_key, value, value_sha256
    491                  FROM radroots_runtime_projection_documents
    492                  WHERE projection_id = ? AND generation = ? AND document_key = ?",
    493             )
    494             .bind(projection_id.as_str())
    495             .bind(generation.as_bytes().as_slice())
    496             .bind(key)
    497             .fetch_optional(self.pool())
    498             .await
    499             .map_err(map_backend)?
    500             .map(|row| {
    501                 ProjectionDocument::from_stored_parts(
    502                     row.try_get("document_key").map_err(map_corrupt)?,
    503                     row.try_get("value").map_err(map_corrupt)?,
    504                     array(row.try_get("value_sha256").map_err(map_corrupt)?)?,
    505                 )
    506             })
    507             .transpose()
    508         })
    509     }
    510 
    511     fn query_projection_documents(
    512         &self,
    513         query: radroots_storage::projection::document_query::ProjectionDocumentQuery,
    514     ) -> BoxFuture<
    515         '_,
    516         Result<radroots_storage::projection::document_query::ProjectionDocumentPage, Error>,
    517     > {
    518         Box::pin(document_query::page(self, query))
    519     }
    520 
    521     fn put_projection_snapshot(
    522         &self,
    523         snapshot: ProjectionSnapshot,
    524     ) -> BoxFuture<'_, Result<(), Error>> {
    525         Box::pin(async move {
    526             self.require_projection_writer()?;
    527             let result = sqlx::query(
    528                 "INSERT INTO radroots_runtime_projection_snapshots (
    529                    projection_id, snapshot_id, generation, created_at_unix_ms,
    530                    value, value_sha256
    531                  ) VALUES (?, ?, ?, ?, ?, ?)
    532                  ON CONFLICT(projection_id, snapshot_id) DO NOTHING",
    533             )
    534             .bind(snapshot.projection_id().as_str())
    535             .bind(snapshot.snapshot_id().as_slice())
    536             .bind(snapshot.generation().as_bytes().as_slice())
    537             .bind(i64_from_u64(snapshot.created_at_unix_ms())?)
    538             .bind(snapshot.value())
    539             .bind(snapshot.value_sha256().as_slice())
    540             .execute(self.pool())
    541             .await
    542             .map_err(map_backend)?;
    543             if result.rows_affected() == 1 {
    544                 return Ok(());
    545             }
    546             match self
    547                 .projection_snapshot(snapshot.projection_id().clone(), *snapshot.snapshot_id())
    548                 .await?
    549             {
    550                 Some(existing) if existing == snapshot => Ok(()),
    551                 Some(_) => Err(Error::CorruptProjectionDocument),
    552                 None => Err(Error::CorruptProjectionDocument),
    553             }
    554         })
    555     }
    556 
    557     fn projection_snapshot(
    558         &self,
    559         projection_id: ProjectionId,
    560         snapshot_id: [u8; 32],
    561     ) -> BoxFuture<'_, Result<Option<ProjectionSnapshot>, Error>> {
    562         Box::pin(async move {
    563             sqlx::query(
    564                 "SELECT projection_id, snapshot_id, generation, created_at_unix_ms,
    565                         value, value_sha256
    566                  FROM radroots_runtime_projection_snapshots
    567                  WHERE projection_id = ? AND snapshot_id = ?",
    568             )
    569             .bind(projection_id.as_str())
    570             .bind(snapshot_id.as_slice())
    571             .fetch_optional(self.pool())
    572             .await
    573             .map_err(map_backend)?
    574             .map(|row| {
    575                 ProjectionSnapshot::from_stored_parts(
    576                     ProjectionId::parse(
    577                         row.try_get::<String, _>("projection_id")
    578                             .map_err(map_corrupt)?,
    579                     )
    580                     .map_err(|_| Error::CorruptProjectionDocument)?,
    581                     array(row.try_get("snapshot_id").map_err(map_corrupt)?)?,
    582                     ProjectionGeneration::new(array(
    583                         row.try_get("generation").map_err(map_corrupt)?,
    584                     )?)
    585                     .map_err(|_| Error::CorruptProjectionDocument)?,
    586                     u64_from_i64(row.try_get("created_at_unix_ms").map_err(map_corrupt)?)?,
    587                     row.try_get("value").map_err(map_corrupt)?,
    588                     array(row.try_get("value_sha256").map_err(map_corrupt)?)?,
    589                 )
    590             })
    591             .transpose()
    592         })
    593     }
    594 }
    595 
    596 #[cfg_attr(coverage_nightly, coverage(off))]
    597 pub(crate) async fn checkpoint_transaction(
    598     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    599     checkpoint: ProjectionCheckpoint,
    600 ) -> Result<ProjectionStatus, Error> {
    601     let prior = sqlx::query(
    602         "SELECT * FROM radroots_runtime_projection_checkpoints WHERE projection_id = ?",
    603     )
    604     .bind(checkpoint.projection_id().as_str())
    605     .fetch_optional(&mut **transaction)
    606     .await
    607     .map_err(map_backend)?
    608     .as_ref()
    609     .map(decode_status)
    610     .transpose()?;
    611     if let Some(prior) = prior.as_ref() {
    612         if prior.generation() != checkpoint.generation() {
    613             return Err(Error::ProjectionCheckpointMismatch);
    614         }
    615         if prior
    616             .checkpoint()
    617             .is_some_and(|value| !checkpoint.advances(value))
    618         {
    619             return Err(Error::ProjectionCheckpointRegression);
    620         }
    621     }
    622     let status = ProjectionStatus::new(
    623         checkpoint.projection_id().clone(),
    624         checkpoint.generation(),
    625         ProjectionHealth::Ready,
    626         Some(checkpoint),
    627         None,
    628     )?;
    629     put_status_transaction(transaction, &status).await?;
    630     Ok(status)
    631 }
    632 
    633 impl SqliteStorage {
    634     fn require_projection_writer(&self) -> Result<(), Error> {
    635         if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
    636             return Err(Error::BackendUnavailable);
    637         }
    638         Ok(())
    639     }
    640 }
    641 
    642 #[cfg_attr(coverage_nightly, coverage(off))]
    643 async fn put_status_transaction(
    644     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    645     status: &ProjectionStatus,
    646 ) -> Result<(), Error> {
    647     let values = checkpoint_values(status.checkpoint())?;
    648     sqlx::query(STATUS_UPSERT)
    649         .bind(status.projection_id().as_str())
    650         .bind(status.generation().as_bytes().as_slice())
    651         .bind(health_name(status.health()))
    652         .bind(values.0)
    653         .bind(values.1)
    654         .bind(values.2)
    655         .bind(values.3)
    656         .bind(
    657             status
    658                 .active_rebuild()
    659                 .map(|ticket| ticket.as_bytes().to_vec()),
    660         )
    661         .execute(&mut **transaction)
    662         .await
    663         .map_err(map_backend)?;
    664     Ok(())
    665 }
    666 
    667 const STATUS_UPSERT: &str = "INSERT INTO radroots_runtime_projection_checkpoints (
    668        projection_id, projection_generation, health, source_generation, source_sequence,
    669        projected_rows, checkpoint_updated_at_unix_ms, active_rebuild
    670      ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
    671      ON CONFLICT(projection_id) DO UPDATE SET
    672        projection_generation = excluded.projection_generation,
    673        health = excluded.health,
    674        source_generation = excluded.source_generation,
    675        source_sequence = excluded.source_sequence,
    676        projected_rows = excluded.projected_rows,
    677        checkpoint_updated_at_unix_ms = excluded.checkpoint_updated_at_unix_ms,
    678        active_rebuild = excluded.active_rebuild";
    679 
    680 type CheckpointValues = (Option<Vec<u8>>, Option<i64>, Option<i64>, Option<i64>);
    681 
    682 fn checkpoint_values(checkpoint: Option<&ProjectionCheckpoint>) -> Result<CheckpointValues, Error> {
    683     let Some(checkpoint) = checkpoint else {
    684         return Ok((None, None, None, None));
    685     };
    686     let (generation, sequence) = checkpoint
    687         .source_position()
    688         .map_or((None, None), |position| {
    689             (
    690                 Some(position.generation().as_bytes().to_vec()),
    691                 Some(i64_from_u64(position.sequence().get())),
    692             )
    693         });
    694     Ok((
    695         generation,
    696         sequence.transpose()?,
    697         Some(i64_from_u64(checkpoint.projected_rows())?),
    698         Some(i64_from_u64(checkpoint.updated_at_unix_ms())?),
    699     ))
    700 }
    701 
    702 fn decode_status(row: &sqlx::sqlite::SqliteRow) -> Result<ProjectionStatus, Error> {
    703     let projection_id = ProjectionId::parse(
    704         row.try_get::<String, _>("projection_id")
    705             .map_err(map_corrupt)?,
    706     )
    707     .map_err(|_| Error::CorruptProjectionRecord)?;
    708     let generation = projection_generation(row, "projection_generation")?;
    709     let checkpoint = decode_checkpoint(row, projection_id.clone(), generation, "")?;
    710     let active = row
    711         .try_get::<Option<Vec<u8>>, _>("active_rebuild")
    712         .map_err(map_corrupt)?
    713         .map(|value| {
    714             RebuildTicketId::new(array(value)?).map_err(|_| Error::CorruptProjectionRecord)
    715         })
    716         .transpose()?;
    717     ProjectionStatus::new(
    718         projection_id,
    719         generation,
    720         health(
    721             row.try_get::<String, _>("health")
    722                 .map_err(map_corrupt)?
    723                 .as_str(),
    724         )?,
    725         checkpoint,
    726         active,
    727     )
    728 }
    729 
    730 fn decode_checkpoint(
    731     row: &sqlx::sqlite::SqliteRow,
    732     projection_id: ProjectionId,
    733     generation: ProjectionGeneration,
    734     prefix: &str,
    735 ) -> Result<Option<ProjectionCheckpoint>, Error> {
    736     let rows = row
    737         .try_get::<Option<i64>, _>(format!("{prefix}projected_rows").as_str())
    738         .map_err(map_corrupt)?;
    739     let updated = row
    740         .try_get::<Option<i64>, _>(format!("{prefix}updated_at_unix_ms").as_str())
    741         .or_else(|_| row.try_get(format!("{prefix}checkpoint_updated_at_unix_ms").as_str()))
    742         .map_err(map_corrupt)?;
    743     let source_generation = row
    744         .try_get::<Option<Vec<u8>>, _>(format!("{prefix}source_generation").as_str())
    745         .map_err(map_corrupt)?;
    746     let source_sequence = row
    747         .try_get::<Option<i64>, _>(format!("{prefix}source_sequence").as_str())
    748         .map_err(map_corrupt)?;
    749     match (rows, updated, source_generation, source_sequence) {
    750         (None, None, None, None) => Ok(None),
    751         (Some(rows), Some(updated), source_generation, source_sequence) => {
    752             let position = match (source_generation, source_sequence) {
    753                 (None, None) => None,
    754                 (Some(source_generation), Some(source_sequence)) => Some(EventPosition::new(
    755                     SourceGeneration::new(array(source_generation)?)
    756                         .map_err(|_| Error::CorruptProjectionRecord)?,
    757                     radroots_storage::event::EventSequence::new(u64_from_i64(source_sequence)?)
    758                         .map_err(|_| Error::CorruptProjectionRecord)?,
    759                 )),
    760                 _ => return Err(Error::CorruptProjectionRecord),
    761             };
    762             ProjectionCheckpoint::new(
    763                 projection_id,
    764                 generation,
    765                 position,
    766                 u64_from_i64(rows)?,
    767                 u64_from_i64(updated)?,
    768             )
    769             .map(Some)
    770             .map_err(|_| Error::CorruptProjectionRecord)
    771         }
    772         _ => Err(Error::CorruptProjectionRecord),
    773     }
    774 }
    775 
    776 #[cfg_attr(coverage_nightly, coverage(off))]
    777 async fn load_invalidation(
    778     connection: &mut SqliteConnection,
    779     projection_id: &ProjectionId,
    780     generation: ProjectionGeneration,
    781 ) -> Result<Option<ProjectionInvalidation>, Error> {
    782     sqlx::query(
    783         "SELECT * FROM radroots_runtime_projection_invalidations
    784          WHERE projection_id = ? AND invalid_generation = ?",
    785     )
    786     .bind(projection_id.as_str())
    787     .bind(generation.as_bytes().as_slice())
    788     .fetch_optional(&mut *connection)
    789     .await
    790     .map_err(map_backend)?
    791     .as_ref()
    792     .map(decode_invalidation)
    793     .transpose()
    794 }
    795 
    796 fn decode_invalidation(row: &sqlx::sqlite::SqliteRow) -> Result<ProjectionInvalidation, Error> {
    797     ProjectionInvalidation::new(
    798         ProjectionId::parse(
    799             row.try_get::<String, _>("projection_id")
    800                 .map_err(map_corrupt)?,
    801         )
    802         .map_err(|_| Error::CorruptProjectionRecord)?,
    803         projection_generation(row, "invalid_generation")?,
    804         projection_generation(row, "replacement_generation")?,
    805         reason(
    806             row.try_get::<String, _>("reason")
    807                 .map_err(map_corrupt)?
    808                 .as_str(),
    809         )?,
    810         u64_from_i64(row.try_get("invalidated_at_unix_ms").map_err(map_corrupt)?)?,
    811     )
    812     .map_err(|_| Error::CorruptProjectionRecord)
    813 }
    814 
    815 #[cfg_attr(coverage_nightly, coverage(off))]
    816 async fn insert_ticket(
    817     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    818     ticket: &RebuildTicket,
    819 ) -> Result<(), Error> {
    820     let checkpoint = checkpoint_values(ticket.checkpoint())?;
    821     sqlx::query(
    822         "INSERT INTO radroots_runtime_projection_rebuilds (
    823            ticket_id, projection_id, invalid_generation, replacement_generation, revision, stage,
    824            source_generation, source_sequence, source_digest,
    825            checkpoint_source_generation, checkpoint_source_sequence, checkpoint_projected_rows,
    826            checkpoint_updated_at_unix_ms, failure, requested_at_unix_ms, updated_at_unix_ms
    827          ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
    828     )
    829     .bind(ticket.ticket_id().as_bytes().as_slice())
    830     .bind(ticket.invalidation().projection_id().as_str())
    831     .bind(
    832         ticket
    833             .invalidation()
    834             .invalid_generation()
    835             .as_bytes()
    836             .as_slice(),
    837     )
    838     .bind(
    839         ticket
    840             .invalidation()
    841             .replacement_generation()
    842             .as_bytes()
    843             .as_slice(),
    844     )
    845     .bind(i64_from_u64(ticket.revision().get())?)
    846     .bind(rebuild_stage_name(ticket.stage()))
    847     .bind(ticket.source_generation().as_bytes().as_slice())
    848     .bind(
    849         ticket
    850             .source_high_water()
    851             .map(|position| i64_from_u64(position.sequence().get()))
    852             .transpose()?,
    853     )
    854     .bind(ticket.source_digest().as_bytes().as_slice())
    855     .bind(checkpoint.0)
    856     .bind(checkpoint.1)
    857     .bind(checkpoint.2)
    858     .bind(checkpoint.3)
    859     .bind(ticket.failure().map(rebuild_failure_name))
    860     .bind(i64_from_u64(ticket.requested_at_unix_ms())?)
    861     .bind(i64_from_u64(ticket.updated_at_unix_ms())?)
    862     .execute(&mut **transaction)
    863     .await
    864     .map_err(map_backend)?;
    865     Ok(())
    866 }
    867 
    868 #[cfg_attr(coverage_nightly, coverage(off))]
    869 async fn update_ticket(
    870     transaction: &mut sqlx::Transaction<'_, Sqlite>,
    871     ticket: &RebuildTicket,
    872     prior: ProjectionRevision,
    873 ) -> Result<(), Error> {
    874     let checkpoint = checkpoint_values(ticket.checkpoint())?;
    875     let result = sqlx::query(
    876         "UPDATE radroots_runtime_projection_rebuilds SET
    877            revision = ?, stage = ?, checkpoint_source_generation = ?,
    878            checkpoint_source_sequence = ?, checkpoint_projected_rows = ?,
    879            checkpoint_updated_at_unix_ms = ?, failure = ?, updated_at_unix_ms = ?
    880          WHERE ticket_id = ? AND revision = ?",
    881     )
    882     .bind(i64_from_u64(ticket.revision().get())?)
    883     .bind(rebuild_stage_name(ticket.stage()))
    884     .bind(checkpoint.0)
    885     .bind(checkpoint.1)
    886     .bind(checkpoint.2)
    887     .bind(checkpoint.3)
    888     .bind(ticket.failure().map(rebuild_failure_name))
    889     .bind(i64_from_u64(ticket.updated_at_unix_ms())?)
    890     .bind(ticket.ticket_id().as_bytes().as_slice())
    891     .bind(i64_from_u64(prior.get())?)
    892     .execute(&mut **transaction)
    893     .await
    894     .map_err(map_backend)?;
    895     if result.rows_affected() != 1 {
    896         return Err(Error::ProjectionRevisionConflict);
    897     }
    898     Ok(())
    899 }
    900 
    901 #[cfg_attr(coverage_nightly, coverage(off))]
    902 async fn decode_ticket(
    903     connection: &mut SqliteConnection,
    904     row: &sqlx::sqlite::SqliteRow,
    905 ) -> Result<RebuildTicket, Error> {
    906     let projection_id = ProjectionId::parse(
    907         row.try_get::<String, _>("projection_id")
    908             .map_err(map_corrupt)?,
    909     )
    910     .map_err(|_| Error::CorruptProjectionRecord)?;
    911     let invalid_generation = projection_generation(row, "invalid_generation")?;
    912     let invalidation = load_invalidation(connection, &projection_id, invalid_generation)
    913         .await?
    914         .ok_or(Error::CorruptProjectionRecord)?;
    915     if invalidation.replacement_generation()
    916         != projection_generation(row, "replacement_generation")?
    917     {
    918         return Err(Error::CorruptProjectionRecord);
    919     }
    920     let checkpoint = decode_checkpoint(
    921         row,
    922         projection_id,
    923         invalidation.replacement_generation(),
    924         "checkpoint_",
    925     )?;
    926     RebuildTicket::from_durable_parts(
    927         RebuildTicketId::new(array(
    928             row.try_get::<Vec<u8>, _>("ticket_id")
    929                 .map_err(map_corrupt)?,
    930         )?)
    931         .map_err(|_| Error::CorruptProjectionRecord)?,
    932         invalidation,
    933         ProjectionRevision::new(u64_from_i64(row.try_get("revision").map_err(map_corrupt)?)?)
    934             .map_err(|_| Error::CorruptProjectionRecord)?,
    935         rebuild_stage(
    936             row.try_get::<String, _>("stage")
    937                 .map_err(map_corrupt)?
    938                 .as_str(),
    939         )?,
    940         SourceGeneration::new(array(
    941             row.try_get::<Vec<u8>, _>("source_generation")
    942                 .map_err(map_corrupt)?,
    943         )?)
    944         .map_err(|_| Error::CorruptProjectionRecord)?,
    945         row.try_get::<Option<i64>, _>("source_sequence")
    946             .map_err(map_corrupt)?
    947             .map(|sequence| {
    948                 Ok(EventPosition::new(
    949                     SourceGeneration::new(array(
    950                         row.try_get::<Vec<u8>, _>("source_generation")
    951                             .map_err(map_corrupt)?,
    952                     )?)
    953                     .map_err(|_| Error::CorruptProjectionRecord)?,
    954                     radroots_storage::event::EventSequence::new(u64_from_i64(sequence)?)
    955                         .map_err(|_| Error::CorruptProjectionRecord)?,
    956                 ))
    957             })
    958             .transpose()?,
    959         RawSourceDigest::new(array(
    960             row.try_get::<Vec<u8>, _>("source_digest")
    961                 .map_err(map_corrupt)?,
    962         )?),
    963         checkpoint,
    964         row.try_get::<Option<String>, _>("failure")
    965             .map_err(map_corrupt)?
    966             .as_deref()
    967             .map(rebuild_failure)
    968             .transpose()?,
    969         u64_from_i64(row.try_get("requested_at_unix_ms").map_err(map_corrupt)?)?,
    970         u64_from_i64(row.try_get("updated_at_unix_ms").map_err(map_corrupt)?)?,
    971     )
    972 }
    973 
    974 #[cfg_attr(coverage_nightly, coverage(off))]
    975 async fn load_manifest(
    976     connection: &mut SqliteConnection,
    977     generation: ProjectionGeneration,
    978 ) -> Result<Option<EventIndexManifest>, Error> {
    979     let Some(row) = sqlx::query(
    980         "SELECT * FROM radroots_runtime_event_index_manifests
    981          WHERE projection_generation = ?",
    982     )
    983     .bind(generation.as_bytes().as_slice())
    984     .fetch_optional(&mut *connection)
    985     .await
    986     .map_err(map_backend)?
    987     else {
    988         return Ok(None);
    989     };
    990     let shards = sqlx::query(
    991         "SELECT * FROM radroots_runtime_event_index_shards
    992          WHERE projection_generation = ? ORDER BY ordinal",
    993     )
    994     .bind(generation.as_bytes().as_slice())
    995     .fetch_all(&mut *connection)
    996     .await
    997     .map_err(map_backend)?
    998     .iter()
    999     .enumerate()
   1000     .map(|(ordinal, row)| {
   1001         if row.try_get::<i64, _>("ordinal").map_err(map_corrupt)?
   1002             != i64::try_from(ordinal).map_err(|_| Error::CorruptProjectionRecord)?
   1003         {
   1004             return Err(Error::CorruptProjectionRecord);
   1005         }
   1006         EventIndexShard::new(
   1007             EventIndexShardId::parse(row.try_get::<String, _>("shard_id").map_err(map_corrupt)?)
   1008                 .map_err(|_| Error::CorruptProjectionRecord)?,
   1009             row.try_get::<String, _>("artifact_path")
   1010                 .map_err(map_corrupt)?,
   1011             u32::try_from(row.try_get::<i64, _>("event_count").map_err(map_corrupt)?)
   1012                 .map_err(|_| Error::CorruptProjectionRecord)?,
   1013             EventIdRange::new(
   1014                 event_id(row, "first_event_id")?,
   1015                 event_id(row, "last_event_id")?,
   1016             )
   1017             .map_err(|_| Error::CorruptProjectionRecord)?,
   1018             u64_from_i64(
   1019                 row.try_get("first_published_at_unix_s")
   1020                     .map_err(map_corrupt)?,
   1021             )?,
   1022             u64_from_i64(
   1023                 row.try_get("last_published_at_unix_s")
   1024                     .map_err(map_corrupt)?,
   1025             )?,
   1026             ArtifactDigest::new(array(
   1027                 row.try_get::<Vec<u8>, _>("artifact_digest")
   1028                     .map_err(map_corrupt)?,
   1029             )?),
   1030         )
   1031         .map_err(|_| Error::CorruptProjectionRecord)
   1032     })
   1033     .collect::<Result<Vec<_>, _>>()?;
   1034     EventIndexManifest::new(
   1035         generation,
   1036         u64_from_i64(row.try_get("total_events").map_err(map_corrupt)?)?,
   1037         u32::try_from(
   1038             row.try_get::<i64, _>("target_shard_size")
   1039                 .map_err(map_corrupt)?,
   1040         )
   1041         .map_err(|_| Error::CorruptProjectionRecord)?,
   1042         u64_from_i64(
   1043             row.try_get("first_published_at_unix_s")
   1044                 .map_err(map_corrupt)?,
   1045         )?,
   1046         u64_from_i64(
   1047             row.try_get("last_published_at_unix_s")
   1048                 .map_err(map_corrupt)?,
   1049         )?,
   1050         shards,
   1051     )
   1052     .map(Some)
   1053     .map_err(|_| Error::CorruptProjectionRecord)
   1054 }
   1055 
   1056 fn encode_index_checkpoint(checkpoint: &EventIndexCheckpoint) -> Result<Vec<u8>, Error> {
   1057     let mut bytes = vec![1];
   1058     let count =
   1059         u16::try_from(checkpoint.shards().len()).map_err(|_| Error::CorruptProjectionRecord)?;
   1060     bytes.extend_from_slice(&count.to_be_bytes());
   1061     for shard in checkpoint.shards() {
   1062         put_string(&mut bytes, shard.shard_id().as_str())?;
   1063         bytes.extend_from_slice(&shard.last_created_at_unix_s().to_be_bytes());
   1064         match shard.last_event_id() {
   1065             Some(event_id) => {
   1066                 bytes.push(1);
   1067                 bytes.extend_from_slice(event_id.as_bytes());
   1068             }
   1069             None => bytes.push(0),
   1070         }
   1071         match shard.cursor() {
   1072             Some(cursor) => {
   1073                 bytes.push(1);
   1074                 put_string(&mut bytes, cursor)?;
   1075             }
   1076             None => bytes.push(0),
   1077         }
   1078     }
   1079     Ok(bytes)
   1080 }
   1081 
   1082 fn decode_index_checkpoint(
   1083     row: &sqlx::sqlite::SqliteRow,
   1084     generation: ProjectionGeneration,
   1085 ) -> Result<EventIndexCheckpoint, Error> {
   1086     let bytes = row
   1087         .try_get::<Vec<u8>, _>("checkpoint")
   1088         .map_err(map_corrupt)?;
   1089     let mut cursor = Cursor::new(bytes.as_slice());
   1090     if cursor.byte()? != 1 {
   1091         return Err(Error::CorruptProjectionRecord);
   1092     }
   1093     let count = usize::from(cursor.u16()?);
   1094     if count > EVENT_INDEX_SHARDS_MAX {
   1095         return Err(Error::CorruptProjectionRecord);
   1096     }
   1097     let mut shards = Vec::with_capacity(count);
   1098     for _ in 0..count {
   1099         let shard_id = EventIndexShardId::parse(cursor.string()?.to_owned())
   1100             .map_err(|_| Error::CorruptProjectionRecord)?;
   1101         let last_created_at_unix_s = cursor.u64()?;
   1102         let last_event_id = match cursor.byte()? {
   1103             0 => None,
   1104             1 => Some(EventId::from_bytes(cursor.array()?)),
   1105             _ => return Err(Error::CorruptProjectionRecord),
   1106         };
   1107         let checkpoint_cursor = match cursor.byte()? {
   1108             0 => None,
   1109             1 => Some(cursor.string()?.to_owned()),
   1110             _ => return Err(Error::CorruptProjectionRecord),
   1111         };
   1112         shards.push(
   1113             EventIndexShardCheckpoint::new(
   1114                 shard_id,
   1115                 last_created_at_unix_s,
   1116                 last_event_id,
   1117                 checkpoint_cursor,
   1118             )
   1119             .map_err(|_| Error::CorruptProjectionRecord)?,
   1120         );
   1121     }
   1122     cursor.finish()?;
   1123     EventIndexCheckpoint::new(
   1124         generation,
   1125         u64_from_i64(row.try_get("generated_at_unix_ms").map_err(map_corrupt)?)?,
   1126         shards,
   1127     )
   1128     .map_err(|_| Error::CorruptProjectionRecord)
   1129 }
   1130 
   1131 pub(crate) fn encode_status_snapshot(status: &ProjectionStatus) -> Result<Vec<u8>, Error> {
   1132     let mut bytes = Vec::with_capacity(128);
   1133     bytes.push(1);
   1134     put_string(&mut bytes, status.projection_id().as_str())?;
   1135     bytes.extend_from_slice(status.generation().as_bytes());
   1136     bytes.push(match status.health() {
   1137         ProjectionHealth::Ready => 0,
   1138         ProjectionHealth::Invalidated => 1,
   1139         ProjectionHealth::Rebuilding => 2,
   1140         ProjectionHealth::Failed => 3,
   1141     });
   1142     match status.checkpoint() {
   1143         Some(checkpoint) => {
   1144             bytes.push(1);
   1145             match checkpoint.source_position() {
   1146                 Some(position) => {
   1147                     bytes.push(1);
   1148                     bytes.extend_from_slice(position.generation().as_bytes());
   1149                     bytes.extend_from_slice(&position.sequence().get().to_be_bytes());
   1150                 }
   1151                 None => bytes.push(0),
   1152             }
   1153             bytes.extend_from_slice(&checkpoint.projected_rows().to_be_bytes());
   1154             bytes.extend_from_slice(&checkpoint.updated_at_unix_ms().to_be_bytes());
   1155         }
   1156         None => bytes.push(0),
   1157     }
   1158     match status.active_rebuild() {
   1159         Some(ticket) => {
   1160             bytes.push(1);
   1161             bytes.extend_from_slice(ticket.as_bytes());
   1162         }
   1163         None => bytes.push(0),
   1164     }
   1165     Ok(bytes)
   1166 }
   1167 
   1168 pub(crate) fn decode_status_snapshot(bytes: &[u8]) -> Result<ProjectionStatus, Error> {
   1169     let mut cursor = Cursor::new(bytes);
   1170     if cursor.byte()? != 1 {
   1171         return Err(Error::CorruptProjectionRecord);
   1172     }
   1173     let projection_id =
   1174         ProjectionId::parse(cursor.string()?).map_err(|_| Error::CorruptProjectionRecord)?;
   1175     let generation =
   1176         ProjectionGeneration::new(cursor.array()?).map_err(|_| Error::CorruptProjectionRecord)?;
   1177     let health = match cursor.byte()? {
   1178         0 => ProjectionHealth::Ready,
   1179         1 => ProjectionHealth::Invalidated,
   1180         2 => ProjectionHealth::Rebuilding,
   1181         3 => ProjectionHealth::Failed,
   1182         _ => return Err(Error::CorruptProjectionRecord),
   1183     };
   1184     let checkpoint = match cursor.byte()? {
   1185         0 => None,
   1186         1 => {
   1187             let source_position = match cursor.byte()? {
   1188                 0 => None,
   1189                 1 => Some(EventPosition::new(
   1190                     SourceGeneration::new(cursor.array()?)
   1191                         .map_err(|_| Error::CorruptProjectionRecord)?,
   1192                     radroots_storage::event::EventSequence::new(cursor.u64()?)
   1193                         .map_err(|_| Error::CorruptProjectionRecord)?,
   1194                 )),
   1195                 _ => return Err(Error::CorruptProjectionRecord),
   1196             };
   1197             Some(
   1198                 ProjectionCheckpoint::new(
   1199                     projection_id.clone(),
   1200                     generation,
   1201                     source_position,
   1202                     cursor.u64()?,
   1203                     cursor.u64()?,
   1204                 )
   1205                 .map_err(|_| Error::CorruptProjectionRecord)?,
   1206             )
   1207         }
   1208         _ => return Err(Error::CorruptProjectionRecord),
   1209     };
   1210     let active_rebuild = match cursor.byte()? {
   1211         0 => None,
   1212         1 => Some(
   1213             RebuildTicketId::new(cursor.array()?).map_err(|_| Error::CorruptProjectionRecord)?,
   1214         ),
   1215         _ => return Err(Error::CorruptProjectionRecord),
   1216     };
   1217     cursor.finish()?;
   1218     ProjectionStatus::new(
   1219         projection_id,
   1220         generation,
   1221         health,
   1222         checkpoint,
   1223         active_rebuild,
   1224     )
   1225     .map_err(|_| Error::CorruptProjectionRecord)
   1226 }
   1227 
   1228 fn put_string(bytes: &mut Vec<u8>, value: &str) -> Result<(), Error> {
   1229     let length = u16::try_from(value.len()).map_err(|_| Error::CorruptProjectionRecord)?;
   1230     bytes.extend_from_slice(&length.to_be_bytes());
   1231     bytes.extend_from_slice(value.as_bytes());
   1232     Ok(())
   1233 }
   1234 
   1235 struct Cursor<'a> {
   1236     bytes: &'a [u8],
   1237     offset: usize,
   1238 }
   1239 
   1240 impl<'a> Cursor<'a> {
   1241     const fn new(bytes: &'a [u8]) -> Self {
   1242         Self { bytes, offset: 0 }
   1243     }
   1244     fn byte(&mut self) -> Result<u8, Error> {
   1245         let value = self
   1246             .bytes
   1247             .get(self.offset)
   1248             .copied()
   1249             .ok_or(Error::CorruptProjectionRecord)?;
   1250         self.offset += 1;
   1251         Ok(value)
   1252     }
   1253     fn u16(&mut self) -> Result<u16, Error> {
   1254         Ok(u16::from_be_bytes(self.array()?))
   1255     }
   1256     fn u64(&mut self) -> Result<u64, Error> {
   1257         Ok(u64::from_be_bytes(self.array()?))
   1258     }
   1259     fn string(&mut self) -> Result<&'a str, Error> {
   1260         let length = usize::from(self.u16()?);
   1261         core::str::from_utf8(self.take(length)?).map_err(|_| Error::CorruptProjectionRecord)
   1262     }
   1263     fn array<const N: usize>(&mut self) -> Result<[u8; N], Error> {
   1264         self.take(N)?
   1265             .try_into()
   1266             .map_err(|_| Error::CorruptProjectionRecord)
   1267     }
   1268     fn take(&mut self, length: usize) -> Result<&'a [u8], Error> {
   1269         let end = self
   1270             .offset
   1271             .checked_add(length)
   1272             .ok_or(Error::CorruptProjectionRecord)?;
   1273         let value = self
   1274             .bytes
   1275             .get(self.offset..end)
   1276             .ok_or(Error::CorruptProjectionRecord)?;
   1277         self.offset = end;
   1278         Ok(value)
   1279     }
   1280     fn finish(self) -> Result<(), Error> {
   1281         if self.offset == self.bytes.len() {
   1282             Ok(())
   1283         } else {
   1284             Err(Error::CorruptProjectionRecord)
   1285         }
   1286     }
   1287 }
   1288 
   1289 fn projection_generation(
   1290     row: &sqlx::sqlite::SqliteRow,
   1291     column: &str,
   1292 ) -> Result<ProjectionGeneration, Error> {
   1293     ProjectionGeneration::new(array(
   1294         row.try_get::<Vec<u8>, _>(column).map_err(map_corrupt)?,
   1295     )?)
   1296     .map_err(|_| Error::CorruptProjectionRecord)
   1297 }
   1298 
   1299 fn event_id(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<EventId, Error> {
   1300     Ok(EventId::from_bytes(array(
   1301         row.try_get::<Vec<u8>, _>(column).map_err(map_corrupt)?,
   1302     )?))
   1303 }
   1304 
   1305 const fn reason_name(value: InvalidationReason) -> &'static str {
   1306     match value {
   1307         InvalidationReason::SourceGenerationChanged => "source_generation_changed",
   1308         InvalidationReason::ProjectionGenerationChanged => "projection_generation_changed",
   1309         InvalidationReason::EventIndexManifestChanged => "event_index_manifest_changed",
   1310         InvalidationReason::IntegrityFailure => "integrity_failure",
   1311         InvalidationReason::OperatorRequested => "operator_requested",
   1312     }
   1313 }
   1314 
   1315 const fn reason(value: &str) -> Result<InvalidationReason, Error> {
   1316     match value.as_bytes() {
   1317         b"source_generation_changed" => Ok(InvalidationReason::SourceGenerationChanged),
   1318         b"projection_generation_changed" => Ok(InvalidationReason::ProjectionGenerationChanged),
   1319         b"event_index_manifest_changed" => Ok(InvalidationReason::EventIndexManifestChanged),
   1320         b"integrity_failure" => Ok(InvalidationReason::IntegrityFailure),
   1321         b"operator_requested" => Ok(InvalidationReason::OperatorRequested),
   1322         _ => Err(Error::CorruptProjectionRecord),
   1323     }
   1324 }
   1325 
   1326 const fn health_name(value: ProjectionHealth) -> &'static str {
   1327     match value {
   1328         ProjectionHealth::Ready => "ready",
   1329         ProjectionHealth::Invalidated => "invalidated",
   1330         ProjectionHealth::Rebuilding => "rebuilding",
   1331         ProjectionHealth::Failed => "failed",
   1332     }
   1333 }
   1334 
   1335 const fn health(value: &str) -> Result<ProjectionHealth, Error> {
   1336     match value.as_bytes() {
   1337         b"ready" => Ok(ProjectionHealth::Ready),
   1338         b"invalidated" => Ok(ProjectionHealth::Invalidated),
   1339         b"rebuilding" => Ok(ProjectionHealth::Rebuilding),
   1340         b"failed" => Ok(ProjectionHealth::Failed),
   1341         _ => Err(Error::CorruptProjectionRecord),
   1342     }
   1343 }
   1344 
   1345 const fn rebuild_stage_name(value: RebuildStage) -> &'static str {
   1346     match value {
   1347         RebuildStage::Requested => "requested",
   1348         RebuildStage::Running => "running",
   1349         RebuildStage::Completed => "completed",
   1350         RebuildStage::Failed => "failed",
   1351     }
   1352 }
   1353 
   1354 const fn rebuild_stage(value: &str) -> Result<RebuildStage, Error> {
   1355     match value.as_bytes() {
   1356         b"requested" => Ok(RebuildStage::Requested),
   1357         b"running" => Ok(RebuildStage::Running),
   1358         b"completed" => Ok(RebuildStage::Completed),
   1359         b"failed" => Ok(RebuildStage::Failed),
   1360         _ => Err(Error::CorruptProjectionRecord),
   1361     }
   1362 }
   1363 
   1364 const fn rebuild_failure_name(value: RebuildFailure) -> &'static str {
   1365     match value {
   1366         RebuildFailure::ReducerRejected => "reducer_rejected",
   1367         RebuildFailure::SourceChanged => "source_changed",
   1368         RebuildFailure::IntegrityFailure => "integrity_failure",
   1369         RebuildFailure::PromotionRejected => "promotion_rejected",
   1370     }
   1371 }
   1372 
   1373 const fn rebuild_failure(value: &str) -> Result<RebuildFailure, Error> {
   1374     match value.as_bytes() {
   1375         b"reducer_rejected" => Ok(RebuildFailure::ReducerRejected),
   1376         b"source_changed" => Ok(RebuildFailure::SourceChanged),
   1377         b"integrity_failure" => Ok(RebuildFailure::IntegrityFailure),
   1378         b"promotion_rejected" => Ok(RebuildFailure::PromotionRejected),
   1379         _ => Err(Error::CorruptProjectionRecord),
   1380     }
   1381 }
   1382 
   1383 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
   1384     bytes.try_into().map_err(|_| Error::CorruptProjectionRecord)
   1385 }
   1386 
   1387 fn i64_from_u64(value: u64) -> Result<i64, Error> {
   1388     i64::try_from(value).map_err(|_| Error::CorruptProjectionRecord)
   1389 }
   1390 
   1391 fn u64_from_i64(value: i64) -> Result<u64, Error> {
   1392     u64::try_from(value).map_err(|_| Error::CorruptProjectionRecord)
   1393 }
   1394 
   1395 fn map_corrupt(_: sqlx::Error) -> Error {
   1396     Error::CorruptProjectionRecord
   1397 }
   1398 
   1399 #[cfg(test)]
   1400 #[cfg_attr(coverage_nightly, coverage(off))]
   1401 mod tests {
   1402     use super::*;
   1403     use crate::migration::runtime::{MIGRATIONS, migration_sql};
   1404     use radroots_storage::{ProjectionStore, event::EventSequence, status::EventStoreMode};
   1405     use sqlx::sqlite::SqlitePoolOptions;
   1406 
   1407     async fn store(mode: EventStoreMode) -> SqliteStorage {
   1408         let pool = SqlitePoolOptions::new()
   1409             .max_connections(1)
   1410             .connect("sqlite::memory:")
   1411             .await
   1412             .expect("memory SQLite");
   1413         sqlx::query("PRAGMA foreign_keys = ON")
   1414             .execute(&pool)
   1415             .await
   1416             .expect("foreign keys");
   1417         for migration in MIGRATIONS {
   1418             sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL"))
   1419                 .execute(&pool)
   1420                 .await
   1421                 .expect("runtime migration");
   1422         }
   1423         sqlx::query(
   1424             "INSERT INTO radroots_runtime_source_generations (
   1425                generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms
   1426              ) VALUES (?, 0, 'active', 1, NULL)",
   1427         )
   1428         .bind([41_u8; 32].as_slice())
   1429         .execute(&pool)
   1430         .await
   1431         .expect("active source generation");
   1432         SqliteStorage::new(
   1433             pool,
   1434             SourceGeneration::new([41; 32]).expect("generation"),
   1435             mode,
   1436         )
   1437     }
   1438 
   1439     fn projection_id() -> ProjectionId {
   1440         ProjectionId::parse("trade_projection").expect("projection id")
   1441     }
   1442 
   1443     fn generation(byte: u8) -> ProjectionGeneration {
   1444         ProjectionGeneration::new([byte; 32]).expect("projection generation")
   1445     }
   1446 
   1447     fn checkpoint(
   1448         generation: ProjectionGeneration,
   1449         sequence: u64,
   1450         rows: u64,
   1451         at: u64,
   1452     ) -> ProjectionCheckpoint {
   1453         ProjectionCheckpoint::new(
   1454             projection_id(),
   1455             generation,
   1456             Some(EventPosition::new(
   1457                 SourceGeneration::new([41; 32]).expect("source generation"),
   1458                 EventSequence::new(sequence).expect("sequence"),
   1459             )),
   1460             rows,
   1461             at,
   1462         )
   1463         .expect("checkpoint")
   1464     }
   1465 
   1466     fn invalidation() -> ProjectionInvalidation {
   1467         ProjectionInvalidation::new(
   1468             projection_id(),
   1469             generation(1),
   1470             generation(2),
   1471             InvalidationReason::ProjectionGenerationChanged,
   1472             200,
   1473         )
   1474         .expect("invalidation")
   1475     }
   1476 
   1477     #[tokio::test]
   1478     async fn checkpoints_invalidation_and_rebuild_lifecycle_are_durable() {
   1479         let store = store(EventStoreMode::ReadWrite).await;
   1480         assert!(
   1481             store
   1482                 .status(projection_id())
   1483                 .await
   1484                 .expect("empty status")
   1485                 .is_none()
   1486         );
   1487         let first = store
   1488             .checkpoint(checkpoint(generation(1), 1, 10, 100))
   1489             .await
   1490             .expect("first checkpoint");
   1491         assert_eq!(first.health(), ProjectionHealth::Ready);
   1492         let advanced = store
   1493             .checkpoint(checkpoint(generation(1), 2, 20, 150))
   1494             .await
   1495             .expect("advanced checkpoint");
   1496         assert_eq!(
   1497             advanced.checkpoint().expect("checkpoint").projected_rows(),
   1498             20
   1499         );
   1500         assert_eq!(
   1501             store
   1502                 .checkpoint(checkpoint(generation(1), 1, 19, 140))
   1503                 .await,
   1504             Err(Error::ProjectionCheckpointRegression)
   1505         );
   1506 
   1507         let invalidation = invalidation();
   1508         let invalidated = store
   1509             .invalidate(invalidation.clone())
   1510             .await
   1511             .expect("invalidate");
   1512         assert_eq!(invalidated.health(), ProjectionHealth::Invalidated);
   1513         let ticket = RebuildTicket::requested(
   1514             RebuildTicketId::new([3; 16]).expect("ticket"),
   1515             invalidation,
   1516             SourceGeneration::new([41; 32]).expect("source generation"),
   1517             Some(EventPosition::new(
   1518                 SourceGeneration::new([41; 32]).expect("source generation"),
   1519                 EventSequence::new(3).expect("sequence"),
   1520             )),
   1521             RawSourceDigest::new([8; 32]),
   1522         )
   1523         .expect("ticket");
   1524         let requested = store
   1525             .request_rebuild(ticket.clone())
   1526             .await
   1527             .expect("request rebuild");
   1528         assert_eq!(
   1529             store.request_rebuild(ticket).await.expect("replay"),
   1530             requested
   1531         );
   1532         let running = store
   1533             .transition_rebuild(RebuildTransition::start(
   1534                 requested.ticket_id(),
   1535                 requested.revision(),
   1536                 210,
   1537             ))
   1538             .await
   1539             .expect("start");
   1540         let progress = store
   1541             .transition_rebuild(RebuildTransition::checkpoint(
   1542                 running.ticket_id(),
   1543                 running.revision(),
   1544                 220,
   1545                 checkpoint(generation(2), 2, 20, 220),
   1546             ))
   1547             .await
   1548             .expect("progress");
   1549         assert_eq!(
   1550             store
   1551                 .transition_rebuild(RebuildTransition::fail(
   1552                     progress.ticket_id(),
   1553                     running.revision(),
   1554                     230,
   1555                     RebuildFailure::IntegrityFailure,
   1556                 ))
   1557                 .await,
   1558             Err(Error::ProjectionRevisionConflict)
   1559         );
   1560         sqlx::query(
   1561             "UPDATE radroots_runtime_source_generations SET sequence_head = 3 WHERE state = 'active'",
   1562         )
   1563         .execute(store.pool())
   1564         .await
   1565         .expect("advance source high water");
   1566         let completed = store
   1567             .transition_rebuild(RebuildTransition::complete(
   1568                 progress.ticket_id(),
   1569                 progress.revision(),
   1570                 240,
   1571                 checkpoint(generation(2), 3, 30, 240),
   1572             ))
   1573             .await
   1574             .expect("complete");
   1575         assert_eq!(completed.stage(), RebuildStage::Completed);
   1576         let status = store
   1577             .status(projection_id())
   1578             .await
   1579             .expect("status")
   1580             .expect("projection");
   1581         assert_eq!(status.health(), ProjectionHealth::Ready);
   1582         assert_eq!(status.generation(), generation(2));
   1583         assert_eq!(status.active_rebuild(), None);
   1584     }
   1585 
   1586     fn manifest(generation: ProjectionGeneration, digest: u8) -> EventIndexManifest {
   1587         EventIndexManifest::new(
   1588             generation,
   1589             4,
   1590             2,
   1591             10,
   1592             40,
   1593             vec![
   1594                 EventIndexShard::new(
   1595                     EventIndexShardId::parse("shard_a").expect("shard id"),
   1596                     "index/shard_a.bin",
   1597                     2,
   1598                     EventIdRange::new(EventId::from_bytes([1; 32]), EventId::from_bytes([2; 32]))
   1599                         .expect("range"),
   1600                     10,
   1601                     20,
   1602                     ArtifactDigest::new([digest; 32]),
   1603                 )
   1604                 .expect("first shard"),
   1605                 EventIndexShard::new(
   1606                     EventIndexShardId::parse("shard_b").expect("shard id"),
   1607                     "index/shard_b.bin",
   1608                     2,
   1609                     EventIdRange::new(EventId::from_bytes([3; 32]), EventId::from_bytes([4; 32]))
   1610                         .expect("range"),
   1611                     30,
   1612                     40,
   1613                     ArtifactDigest::new([digest.wrapping_add(1); 32]),
   1614                 )
   1615                 .expect("second shard"),
   1616             ],
   1617         )
   1618         .expect("manifest")
   1619     }
   1620 
   1621     fn index_checkpoint(generation: ProjectionGeneration, at: u64) -> EventIndexCheckpoint {
   1622         EventIndexCheckpoint::new(
   1623             generation,
   1624             at,
   1625             vec![
   1626                 EventIndexShardCheckpoint::new(
   1627                     EventIndexShardId::parse("shard_b").expect("shard id"),
   1628                     40,
   1629                     Some(EventId::from_bytes([4; 32])),
   1630                     Some("cursor-b".to_owned()),
   1631                 )
   1632                 .expect("second checkpoint"),
   1633                 EventIndexShardCheckpoint::new(
   1634                     EventIndexShardId::parse("shard_a").expect("shard id"),
   1635                     20,
   1636                     Some(EventId::from_bytes([2; 32])),
   1637                     None,
   1638                 )
   1639                 .expect("first checkpoint"),
   1640             ],
   1641         )
   1642         .expect("index checkpoint")
   1643     }
   1644 
   1645     #[tokio::test]
   1646     async fn event_index_manifest_and_checkpoint_round_trip_and_reject_regression() {
   1647         let store = store(EventStoreMode::ReadWrite).await;
   1648         let generation = generation(7);
   1649         let expected_manifest = manifest(generation, 8);
   1650         store
   1651             .put_event_index_manifest(expected_manifest.clone())
   1652             .await
   1653             .expect("put manifest");
   1654         store
   1655             .put_event_index_manifest(expected_manifest.clone())
   1656             .await
   1657             .expect("manifest replay");
   1658         assert_eq!(
   1659             store
   1660                 .event_index_manifest(generation)
   1661                 .await
   1662                 .expect("manifest lookup")
   1663                 .expect("manifest"),
   1664             expected_manifest
   1665         );
   1666         assert_eq!(
   1667             store
   1668                 .put_event_index_manifest(manifest(generation, 9))
   1669                 .await,
   1670             Err(Error::CorruptProjectionRecord)
   1671         );
   1672 
   1673         let first = index_checkpoint(generation, 100);
   1674         store
   1675             .put_event_index_checkpoint(first.clone())
   1676             .await
   1677             .expect("put checkpoint");
   1678         assert_eq!(
   1679             store
   1680                 .event_index_checkpoint(generation)
   1681                 .await
   1682                 .expect("checkpoint lookup")
   1683                 .expect("checkpoint"),
   1684             first
   1685         );
   1686         store
   1687             .put_event_index_checkpoint(index_checkpoint(generation, 110))
   1688             .await
   1689             .expect("advance checkpoint");
   1690         assert_eq!(
   1691             store
   1692                 .put_event_index_checkpoint(index_checkpoint(generation, 109))
   1693                 .await,
   1694             Err(Error::InvalidEventIndexCheckpoint)
   1695         );
   1696         for corrupt in [&[0_u8][..], &[1_u8, 0xff, 0xff][..]] {
   1697             sqlx::query(
   1698                 "UPDATE radroots_runtime_event_index_checkpoints
   1699                  SET checkpoint = ? WHERE projection_generation = ?",
   1700             )
   1701             .bind(corrupt)
   1702             .bind(generation.as_bytes().as_slice())
   1703             .execute(store.pool())
   1704             .await
   1705             .expect("forge corrupt index checkpoint");
   1706             assert_eq!(
   1707                 store.event_index_checkpoint(generation).await,
   1708                 Err(Error::CorruptProjectionRecord)
   1709             );
   1710         }
   1711     }
   1712 
   1713     #[tokio::test]
   1714     async fn materialized_documents_and_frozen_snapshots_round_trip_and_fail_closed() {
   1715         let writer = store(EventStoreMode::ReadWrite).await;
   1716         let id = projection_id();
   1717         let generation = generation(12);
   1718         let first = ProjectionDocument::new("context.alpha".into(), b"{\"cards\":[]}".to_vec())
   1719             .expect("document");
   1720         writer
   1721             .put_projection_document(id.clone(), generation, first)
   1722             .await
   1723             .expect("put document");
   1724         assert_eq!(
   1725             writer
   1726                 .projection_document(id.clone(), generation, "context.alpha".into(),)
   1727                 .await
   1728                 .expect("document lookup")
   1729                 .expect("document")
   1730                 .value(),
   1731             b"{\"cards\":[]}"
   1732         );
   1733         writer
   1734             .put_projection_document(
   1735                 id.clone(),
   1736                 generation,
   1737                 ProjectionDocument::new("context.alpha".into(), b"{\"cards\":[1]}".to_vec())
   1738                     .expect("replacement"),
   1739             )
   1740             .await
   1741             .expect("replace document");
   1742         assert_eq!(
   1743             writer
   1744                 .projection_document(id.clone(), generation, "context.alpha".into(),)
   1745                 .await
   1746                 .expect("document lookup")
   1747                 .expect("document")
   1748                 .value(),
   1749             b"{\"cards\":[1]}"
   1750         );
   1751 
   1752         let snapshot = ProjectionSnapshot::new(
   1753             id.clone(),
   1754             [13; 32],
   1755             generation,
   1756             1_000,
   1757             b"{\"frozen\":true}".to_vec(),
   1758         )
   1759         .expect("snapshot");
   1760         writer
   1761             .put_projection_snapshot(snapshot.clone())
   1762             .await
   1763             .expect("put snapshot");
   1764         writer
   1765             .put_projection_snapshot(snapshot.clone())
   1766             .await
   1767             .expect("idempotent snapshot replay");
   1768         let concurrent_snapshot = ProjectionSnapshot::new(
   1769             id.clone(),
   1770             [14; 32],
   1771             generation,
   1772             1_001,
   1773             b"{\"frozen\":\"concurrent\"}".to_vec(),
   1774         )
   1775         .expect("concurrent snapshot");
   1776         let (left, right) = tokio::join!(
   1777             writer.put_projection_snapshot(concurrent_snapshot.clone()),
   1778             writer.put_projection_snapshot(concurrent_snapshot),
   1779         );
   1780         left.expect("concurrent left snapshot insert");
   1781         right.expect("concurrent right snapshot insert");
   1782         assert_eq!(
   1783             writer
   1784                 .projection_snapshot(id.clone(), [13; 32])
   1785                 .await
   1786                 .expect("snapshot lookup"),
   1787             Some(snapshot)
   1788         );
   1789         assert_eq!(
   1790             writer
   1791                 .put_projection_snapshot(
   1792                     ProjectionSnapshot::new(
   1793                         id.clone(),
   1794                         [13; 32],
   1795                         generation,
   1796                         1_000,
   1797                         b"{\"frozen\":false}".to_vec(),
   1798                     )
   1799                     .expect("conflicting snapshot"),
   1800                 )
   1801                 .await,
   1802             Err(Error::CorruptProjectionDocument)
   1803         );
   1804 
   1805         sqlx::query(
   1806             "UPDATE radroots_runtime_projection_documents
   1807              SET value = X'00' WHERE projection_id = ?",
   1808         )
   1809         .bind(id.as_str())
   1810         .execute(writer.pool())
   1811         .await
   1812         .expect("forge corrupt document");
   1813         assert_eq!(
   1814             writer
   1815                 .projection_document(id, generation, "context.alpha".into())
   1816                 .await,
   1817             Err(Error::CorruptProjectionDocument)
   1818         );
   1819 
   1820         let read_only = store(EventStoreMode::ReadOnly).await;
   1821         assert_eq!(
   1822             read_only
   1823                 .put_projection_document(
   1824                     projection_id(),
   1825                     generation,
   1826                     ProjectionDocument::new("context.alpha".into(), vec![1]).expect("document"),
   1827                 )
   1828                 .await,
   1829             Err(Error::BackendUnavailable)
   1830         );
   1831     }
   1832 
   1833     #[tokio::test]
   1834     async fn failed_rebuild_corruption_and_read_only_mode_fail_closed() {
   1835         let store = store(EventStoreMode::ReadWrite).await;
   1836         let initial = store
   1837             .checkpoint(checkpoint(generation(1), 1, 1, 100))
   1838             .await
   1839             .expect("checkpoint");
   1840         let encoded = encode_status_snapshot(&initial).expect("encode projection status");
   1841         for end in 0..encoded.len() {
   1842             let _ = decode_status_snapshot(&encoded[..end]);
   1843         }
   1844         let mut trailing = encoded.clone();
   1845         trailing.push(0);
   1846         assert_eq!(
   1847             decode_status_snapshot(&trailing),
   1848             Err(Error::CorruptProjectionRecord)
   1849         );
   1850         for index in 0..encoded.len() {
   1851             let mut corrupt = encoded.clone();
   1852             corrupt[index] ^= 0xff;
   1853             let _ = decode_status_snapshot(&corrupt);
   1854         }
   1855         let invalidation = invalidation();
   1856         store
   1857             .invalidate(invalidation.clone())
   1858             .await
   1859             .expect("invalidate");
   1860         let ticket = store
   1861             .request_rebuild(
   1862                 RebuildTicket::requested(
   1863                     RebuildTicketId::new([9; 16]).expect("ticket"),
   1864                     invalidation,
   1865                     SourceGeneration::new([41; 32]).expect("source generation"),
   1866                     None,
   1867                     RawSourceDigest::new([8; 32]),
   1868                 )
   1869                 .expect("ticket"),
   1870             )
   1871             .await
   1872             .expect("request rebuild");
   1873         let failed = store
   1874             .transition_rebuild(RebuildTransition::fail(
   1875                 ticket.ticket_id(),
   1876                 ticket.revision(),
   1877                 210,
   1878                 RebuildFailure::IntegrityFailure,
   1879             ))
   1880             .await
   1881             .expect("fail rebuild");
   1882         assert_eq!(failed.stage(), RebuildStage::Failed);
   1883         assert_eq!(
   1884             store
   1885                 .status(projection_id())
   1886                 .await
   1887                 .expect("status")
   1888                 .expect("projection")
   1889                 .health(),
   1890             ProjectionHealth::Ready
   1891         );
   1892 
   1893         sqlx::query("PRAGMA ignore_check_constraints = ON")
   1894             .execute(store.pool())
   1895             .await
   1896             .expect("disable constraints");
   1897         sqlx::query(
   1898             "UPDATE radroots_runtime_projection_checkpoints
   1899              SET projection_generation = X'01' WHERE projection_id = ?",
   1900         )
   1901         .bind(projection_id().as_str())
   1902         .execute(store.pool())
   1903         .await
   1904         .expect("forge corruption");
   1905         assert_eq!(
   1906             store.status(projection_id()).await,
   1907             Err(Error::CorruptProjectionRecord)
   1908         );
   1909 
   1910         let read_only = SqliteStorage::new(
   1911             store.pool().clone(),
   1912             SourceGeneration::new([41; 32]).expect("generation"),
   1913             EventStoreMode::ReadOnly,
   1914         );
   1915         assert_eq!(
   1916             read_only
   1917                 .checkpoint(checkpoint(generation(3), 1, 1, 300))
   1918                 .await,
   1919             Err(Error::BackendUnavailable)
   1920         );
   1921     }
   1922 }