lib

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

projection.rs (22585B)


      1 //! Projection refresh and rebuild orchestration.
      2 
      3 use radroots_storage::{
      4     Error as StorageError, ProjectionStore,
      5     event::{
      6         AdmissionStage, EVENT_QUERY_LIMIT_MAX, EventPosition, EventQuery, EventQueryBounds,
      7         SourceGeneration, StoredVisibleEvent,
      8     },
      9     projection::{
     10         InvalidationReason, ProjectionCheckpoint, ProjectionGeneration, ProjectionHealth,
     11         ProjectionId, ProjectionInvalidation, ProjectionRevision, ProjectionStatus,
     12         RawSourceDigest, RebuildFailure, RebuildStage, RebuildTicket, RebuildTicketId,
     13         RebuildTransition,
     14     },
     15 };
     16 use sha2::{Digest, Sha256};
     17 
     18 use crate::{
     19     Engine,
     20     policy::{Error, OperationKind},
     21 };
     22 
     23 /// Maximum number of reducer batches in one explicit refresh call.
     24 pub const PROJECTION_REFRESH_MAX_BATCHES: u16 = 1_000;
     25 /// Maximum canonical raw events included in one rebuild source preflight.
     26 pub const PROJECTION_RAW_SOURCE_MAX_EVENTS: u64 = 1_000_000;
     27 const RAW_SOURCE_DIGEST_DOMAIN: &[u8] = b"radroots:projection:raw-source:v1\0";
     28 
     29 /// Bounded refresh request for one exact reducer generation.
     30 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     31 #[derive(Clone, Debug, Eq, PartialEq)]
     32 pub struct RefreshRequest {
     33     projection_id: ProjectionId,
     34     generation: ProjectionGeneration,
     35     batch_limit: u16,
     36     max_batches: u16,
     37 }
     38 
     39 impl RefreshRequest {
     40     pub fn new(
     41         projection_id: ProjectionId,
     42         generation: ProjectionGeneration,
     43         batch_limit: u16,
     44         max_batches: u16,
     45     ) -> Result<Self, Error> {
     46         if batch_limit == 0
     47             || batch_limit > EVENT_QUERY_LIMIT_MAX
     48             || max_batches == 0
     49             || max_batches > PROJECTION_REFRESH_MAX_BATCHES
     50         {
     51             return Err(Error::InvalidProjectionRequest);
     52         }
     53         Ok(Self {
     54             projection_id,
     55             generation,
     56             batch_limit,
     57             max_batches,
     58         })
     59     }
     60 
     61     pub const fn projection_id(&self) -> &ProjectionId {
     62         &self.projection_id
     63     }
     64 
     65     pub const fn generation(&self) -> ProjectionGeneration {
     66         self.generation
     67     }
     68 
     69     pub const fn batch_limit(&self) -> u16 {
     70         self.batch_limit
     71     }
     72 
     73     pub const fn max_batches(&self) -> u16 {
     74         self.max_batches
     75     }
     76 }
     77 
     78 /// Owning-domain deterministic reducer capability.
     79 ///
     80 /// Reducers own domain semantics and projected row calculation. They receive
     81 /// canonical visible events in storage order and must perform no durable
     82 /// metadata mutation; sync owns the checkpoint/rebuild coordination boundary.
     83 pub trait Reducer: Send + Sync {
     84     fn projection_id(&self) -> &ProjectionId;
     85     fn generation(&self) -> ProjectionGeneration;
     86     /// Opens an isolated replacement generation. Existing readers must remain
     87     /// bound to the active generation until storage promotes the ticket.
     88     fn begin_rebuild(
     89         &self,
     90         ticket_id: RebuildTicketId,
     91         source_generation: SourceGeneration,
     92         source_digest: RawSourceDigest,
     93     ) -> Result<(), ReducerError>;
     94     fn reduce(
     95         &self,
     96         events: &[StoredVisibleEvent],
     97         prior_projected_rows: u64,
     98         rebuild_ticket: Option<RebuildTicketId>,
     99     ) -> Result<u64, ReducerError>;
    100     /// Discards an isolated replacement generation after durable failure.
    101     fn abort_rebuild(
    102         &self,
    103         ticket_id: RebuildTicketId,
    104         failure: RebuildFailure,
    105     ) -> Result<(), ReducerError>;
    106 }
    107 
    108 /// Secret-safe reducer rejection normalized at the orchestration boundary.
    109 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    110 pub struct ReducerError;
    111 
    112 /// Refresh execution class.
    113 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    114 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    115 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    116 pub enum RefreshKind {
    117     Incremental,
    118     Rebuild,
    119 }
    120 
    121 /// Deterministic state returned to the host scheduler.
    122 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    123 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    124 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    125 pub enum RefreshState {
    126     Complete,
    127     Partial,
    128     Failed,
    129 }
    130 
    131 /// Normalized projection progress after one bounded refresh call.
    132 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    133 #[derive(Clone, Debug, Eq, PartialEq)]
    134 pub struct RefreshReceipt {
    135     kind: RefreshKind,
    136     state: RefreshState,
    137     batches: u16,
    138     events_reduced: usize,
    139     checkpoint: Option<ProjectionCheckpoint>,
    140     rebuild_ticket: Option<RebuildTicketId>,
    141 }
    142 
    143 impl RefreshReceipt {
    144     pub const fn kind(&self) -> RefreshKind {
    145         self.kind
    146     }
    147     pub const fn state(&self) -> RefreshState {
    148         self.state
    149     }
    150     pub const fn batches(&self) -> u16 {
    151         self.batches
    152     }
    153     pub const fn events_reduced(&self) -> usize {
    154         self.events_reduced
    155     }
    156     pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
    157         self.checkpoint.as_ref()
    158     }
    159     pub const fn rebuild_ticket(&self) -> Option<RebuildTicketId> {
    160         self.rebuild_ticket
    161     }
    162 }
    163 
    164 impl Engine {
    165     /// Runs at most the requested number of deterministic reducer batches.
    166     pub async fn refresh_projection(
    167         &self,
    168         request: RefreshRequest,
    169         reducer: &dyn Reducer,
    170     ) -> Result<RefreshReceipt, Error> {
    171         if reducer.projection_id() != request.projection_id()
    172             || reducer.generation() != request.generation()
    173         {
    174             return Err(Error::InvalidProjectionRequest);
    175         }
    176         let status = ProjectionStore::status(self.storage.as_ref(), request.projection_id.clone())
    177             .await
    178             .map_err(map_storage_error)?;
    179         let source = if status.as_ref().is_some_and(|status| {
    180             status.generation() != request.generation
    181                 || status.health() == ProjectionHealth::Rebuilding
    182         }) {
    183             Some(self.raw_source_snapshot().await?)
    184         } else {
    185             None
    186         };
    187         let mut coordination = self
    188             .projection_coordination(&request, status, source.as_ref())
    189             .await?;
    190         let kind = if coordination.ticket.is_some() {
    191             RefreshKind::Rebuild
    192         } else {
    193             RefreshKind::Incremental
    194         };
    195         let mut receipt = RefreshReceipt {
    196             kind,
    197             state: RefreshState::Partial,
    198             batches: 0,
    199             events_reduced: 0,
    200             checkpoint: coordination.checkpoint.clone(),
    201             rebuild_ticket: coordination.ticket.as_ref().map(RebuildTicket::ticket_id),
    202         };
    203 
    204         if let Some(ticket) = coordination.ticket.as_ref()
    205             && !source.is_some_and(|source| source.matches_ticket(ticket))
    206         {
    207             self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
    208                 .await?;
    209             receipt.state = RefreshState::Failed;
    210             return Ok(receipt);
    211         }
    212 
    213         if coordination.started
    214             && let Some(ticket) = coordination.ticket.as_ref()
    215             && reducer
    216                 .begin_rebuild(
    217                     ticket.ticket_id(),
    218                     ticket.source_generation(),
    219                     ticket.source_digest(),
    220                 )
    221                 .is_err()
    222         {
    223             self.fail_rebuild(ticket, reducer, RebuildFailure::ReducerRejected)
    224                 .await?;
    225             receipt.state = RefreshState::Failed;
    226             return Ok(receipt);
    227         }
    228 
    229         for batch_index in 0..request.max_batches {
    230             let mut bounds =
    231                 EventQueryBounds::first(request.batch_limit).map_err(map_storage_error)?;
    232             if let Some(position) = coordination
    233                 .checkpoint
    234                 .as_ref()
    235                 .and_then(ProjectionCheckpoint::source_position)
    236             {
    237                 bounds = bounds.after(position);
    238             }
    239             let page = self
    240                 .storage
    241                 .query_visible(EventQuery::all(bounds))
    242                 .await
    243                 .map_err(map_storage_error)?;
    244             let prior_rows = coordination
    245                 .checkpoint
    246                 .as_ref()
    247                 .map_or(0, ProjectionCheckpoint::projected_rows);
    248             let projected_rows = if page.items().is_empty() {
    249                 prior_rows
    250             } else {
    251                 match reducer.reduce(
    252                     page.items(),
    253                     prior_rows,
    254                     coordination.ticket.as_ref().map(RebuildTicket::ticket_id),
    255                 ) {
    256                     Ok(rows) if rows >= prior_rows => rows,
    257                     Ok(_) => return Err(Error::InvalidReducerOutput),
    258                     Err(_) => {
    259                         if let Some(ticket) = coordination.ticket.as_ref() {
    260                             self.fail_rebuild(ticket, reducer, RebuildFailure::ReducerRejected)
    261                                 .await?;
    262                         }
    263                         receipt.state = RefreshState::Failed;
    264                         return Ok(receipt);
    265                     }
    266                 }
    267             };
    268             let source_position = page
    269                 .items()
    270                 .last()
    271                 .map(StoredVisibleEvent::position)
    272                 .or_else(|| {
    273                     coordination
    274                         .checkpoint
    275                         .as_ref()
    276                         .and_then(ProjectionCheckpoint::source_position)
    277                 });
    278             let checkpoint = ProjectionCheckpoint::new(
    279                 request.projection_id.clone(),
    280                 request.generation,
    281                 source_position,
    282                 projected_rows,
    283                 self.clock.now_unix_ms()?,
    284             )
    285             .map_err(map_storage_error)?;
    286             let complete = page.items().len() < usize::from(request.batch_limit);
    287             if let Some(ticket) = coordination.ticket.as_mut() {
    288                 if complete {
    289                     let current_source = self.raw_source_snapshot().await?;
    290                     if !current_source.matches_ticket(ticket) {
    291                         self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
    292                             .await?;
    293                         receipt.state = RefreshState::Failed;
    294                         return Ok(receipt);
    295                     }
    296                 }
    297                 let transition = if complete {
    298                     RebuildTransition::complete(
    299                         ticket.ticket_id(),
    300                         ticket.revision(),
    301                         checkpoint.updated_at_unix_ms(),
    302                         checkpoint.clone(),
    303                     )
    304                 } else {
    305                     RebuildTransition::checkpoint(
    306                         ticket.ticket_id(),
    307                         ticket.revision(),
    308                         checkpoint.updated_at_unix_ms(),
    309                         checkpoint.clone(),
    310                     )
    311                 };
    312                 match self.storage.transition_rebuild(transition).await {
    313                     Ok(next) => *ticket = next,
    314                     Err(StorageError::SourceGenerationChanged) if complete => {
    315                         self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
    316                             .await?;
    317                         receipt.state = RefreshState::Failed;
    318                         return Ok(receipt);
    319                     }
    320                     Err(error) => return Err(map_storage_error(error)),
    321                 }
    322             } else {
    323                 self.storage
    324                     .checkpoint(checkpoint.clone())
    325                     .await
    326                     .map_err(map_storage_error)?;
    327             }
    328             receipt.batches += 1;
    329             receipt.events_reduced += page.items().len();
    330             receipt.checkpoint = Some(checkpoint.clone());
    331             coordination.checkpoint = Some(checkpoint);
    332             if complete {
    333                 receipt.state = RefreshState::Complete;
    334                 return Ok(receipt);
    335             }
    336             if batch_index + 1 == request.max_batches {
    337                 receipt.state = RefreshState::Partial;
    338                 return Ok(receipt);
    339             }
    340         }
    341         unreachable!("validated refresh requests execute at least one batch")
    342     }
    343 
    344     async fn projection_coordination(
    345         &self,
    346         request: &RefreshRequest,
    347         status: Option<ProjectionStatus>,
    348         source: Option<&RawSourceSnapshot>,
    349     ) -> Result<ProjectionCoordination, Error> {
    350         let Some(status) = status else {
    351             return Ok(ProjectionCoordination::default());
    352         };
    353         if status.generation() == request.generation && status.health() == ProjectionHealth::Ready {
    354             return Ok(ProjectionCoordination {
    355                 checkpoint: status.checkpoint().cloned(),
    356                 ticket: None,
    357                 started: false,
    358             });
    359         }
    360         if status.health() == ProjectionHealth::Rebuilding {
    361             let ticket_id = status.active_rebuild().ok_or(Error::StorageFailed)?;
    362             let ticket = self
    363                 .storage
    364                 .rebuild(ticket_id)
    365                 .await
    366                 .map_err(map_storage_error)?
    367                 .ok_or(Error::StorageFailed)?;
    368             if ticket.invalidation().replacement_generation() != request.generation {
    369                 return Err(Error::StorageConflict);
    370             }
    371             return Ok(ProjectionCoordination {
    372                 checkpoint: ticket.checkpoint().cloned(),
    373                 ticket: Some(ticket),
    374                 started: false,
    375             });
    376         }
    377 
    378         let invalidation = if status.generation() != request.generation {
    379             if status.health() != ProjectionHealth::Ready {
    380                 return Err(Error::StorageConflict);
    381             }
    382             let invalidation = match self
    383                 .storage
    384                 .invalidation(request.projection_id.clone(), request.generation)
    385                 .await
    386                 .map_err(map_storage_error)?
    387             {
    388                 Some(existing) if existing.invalid_generation() == status.generation() => existing,
    389                 Some(_) => return Err(Error::StorageConflict),
    390                 None => ProjectionInvalidation::new(
    391                     request.projection_id.clone(),
    392                     status.generation(),
    393                     request.generation,
    394                     InvalidationReason::ProjectionGenerationChanged,
    395                     self.clock.now_unix_ms()?,
    396                 )
    397                 .map_err(map_storage_error)?,
    398             };
    399             self.storage
    400                 .invalidate(invalidation.clone())
    401                 .await
    402                 .map_err(map_storage_error)?;
    403             invalidation
    404         } else if status.health() == ProjectionHealth::Invalidated {
    405             self.storage
    406                 .invalidation(request.projection_id.clone(), request.generation)
    407                 .await
    408                 .map_err(map_storage_error)?
    409                 .ok_or(Error::StorageFailed)?
    410         } else {
    411             return Err(Error::StorageConflict);
    412         };
    413         let sync_id = self.ids.next_id(OperationKind::Projection)?;
    414         let source = source.ok_or(Error::StorageFailed)?;
    415         let ticket = RebuildTicket::requested(
    416             RebuildTicketId::new(*sync_id.as_bytes()).map_err(map_storage_error)?,
    417             invalidation,
    418             source.generation,
    419             source.high_water,
    420             source.digest,
    421         )
    422         .map_err(map_storage_error)?;
    423         let requested = self
    424             .storage
    425             .request_rebuild(ticket)
    426             .await
    427             .map_err(map_storage_error)?;
    428         let running = self
    429             .storage
    430             .transition_rebuild(RebuildTransition::start(
    431                 requested.ticket_id(),
    432                 ProjectionRevision::INITIAL,
    433                 self.clock.now_unix_ms()?,
    434             ))
    435             .await
    436             .map_err(map_storage_error)?;
    437         Ok(ProjectionCoordination {
    438             checkpoint: None,
    439             ticket: Some(running),
    440             started: true,
    441         })
    442     }
    443 
    444     async fn raw_source_snapshot(&self) -> Result<RawSourceSnapshot, Error> {
    445         let mut hasher = Sha256::new();
    446         hasher.update(RAW_SOURCE_DIGEST_DOMAIN);
    447         let mut cursor = None;
    448         let mut count = 0_u64;
    449         let mut generation = None;
    450         let mut high_water = None;
    451         loop {
    452             let mut bounds =
    453                 EventQueryBounds::first(EVENT_QUERY_LIMIT_MAX).map_err(map_storage_error)?;
    454             if let Some(position) = cursor {
    455                 bounds = bounds.after(position);
    456             }
    457             let page = self
    458                 .storage
    459                 .query_raw(EventQuery::all(bounds))
    460                 .await
    461                 .map_err(map_storage_error)?;
    462             if generation
    463                 .replace(page.generation())
    464                 .is_some_and(|prior| prior != page.generation())
    465             {
    466                 return Err(Error::StorageConflict);
    467             }
    468             hasher.update(page.generation().as_bytes());
    469             for event in page.items() {
    470                 count = count.checked_add(1).ok_or(Error::StorageFailed)?;
    471                 if count > PROJECTION_RAW_SOURCE_MAX_EVENTS {
    472                     return Err(Error::InvalidProjectionRequest);
    473                 }
    474                 let position = event.position();
    475                 hasher.update(position.sequence().get().to_be_bytes());
    476                 hasher.update([admission_stage_byte(event.stage())]);
    477                 let raw = event.event().raw_json().as_bytes();
    478                 hasher.update(
    479                     u64::try_from(raw.len())
    480                         .map_err(|_| Error::StorageFailed)?
    481                         .to_be_bytes(),
    482                 );
    483                 hasher.update(raw);
    484                 high_water = Some(position);
    485             }
    486             cursor = page.next_cursor();
    487             if cursor.is_none() {
    488                 break;
    489             }
    490         }
    491         Ok(RawSourceSnapshot {
    492             generation: generation.ok_or(Error::StorageFailed)?,
    493             high_water,
    494             digest: RawSourceDigest::new(hasher.finalize().into()),
    495         })
    496     }
    497 
    498     async fn fail_rebuild(
    499         &self,
    500         ticket: &RebuildTicket,
    501         reducer: &dyn Reducer,
    502         failure: RebuildFailure,
    503     ) -> Result<(), Error> {
    504         let failed = self
    505             .storage
    506             .transition_rebuild(RebuildTransition::fail(
    507                 ticket.ticket_id(),
    508                 ticket.revision(),
    509                 self.clock.now_unix_ms()?,
    510                 failure,
    511             ))
    512             .await
    513             .map_err(map_storage_error)?;
    514         debug_assert_eq!(failed.stage(), RebuildStage::Failed);
    515         reducer
    516             .abort_rebuild(ticket.ticket_id(), failure)
    517             .map_err(|_| Error::InvalidReducerOutput)
    518     }
    519 }
    520 
    521 #[derive(Default)]
    522 struct ProjectionCoordination {
    523     checkpoint: Option<ProjectionCheckpoint>,
    524     ticket: Option<RebuildTicket>,
    525     started: bool,
    526 }
    527 
    528 #[derive(Clone, Copy)]
    529 struct RawSourceSnapshot {
    530     generation: SourceGeneration,
    531     high_water: Option<EventPosition>,
    532     digest: RawSourceDigest,
    533 }
    534 
    535 impl RawSourceSnapshot {
    536     fn matches_ticket(self, ticket: &RebuildTicket) -> bool {
    537         self.generation == ticket.source_generation()
    538             && self.high_water == ticket.source_high_water()
    539             && self.digest == ticket.source_digest()
    540     }
    541 }
    542 
    543 const fn admission_stage_byte(stage: AdmissionStage) -> u8 {
    544     match stage {
    545         AdmissionStage::Raw => 0,
    546         AdmissionStage::Verified => 1,
    547         AdmissionStage::Visible => 2,
    548     }
    549 }
    550 
    551 fn map_storage_error(error: StorageError) -> Error {
    552     match error {
    553         StorageError::SpaceInsufficient => Error::StorageSpaceInsufficient,
    554         StorageError::ProjectionCheckpointMismatch
    555         | StorageError::ProjectionCheckpointRegression
    556         | StorageError::ProjectionRevisionConflict
    557         | StorageError::SourceGenerationChanged => Error::StorageConflict,
    558         _ => Error::StorageFailed,
    559     }
    560 }
    561 
    562 #[cfg(test)]
    563 #[cfg_attr(coverage_nightly, coverage(off))]
    564 mod tests {
    565     use super::*;
    566 
    567     #[test]
    568     fn raw_source_identity_stage_encoding_and_error_mapping_are_exact() {
    569         assert_eq!(
    570             map_storage_error(StorageError::SpaceInsufficient),
    571             Error::StorageSpaceInsufficient
    572         );
    573         let invalidation = ProjectionInvalidation::new(
    574             ProjectionId::parse("projection-helper").unwrap(),
    575             ProjectionGeneration::new([1; 32]).unwrap(),
    576             ProjectionGeneration::new([2; 32]).unwrap(),
    577             InvalidationReason::ProjectionGenerationChanged,
    578             1,
    579         )
    580         .unwrap();
    581         let source_generation = SourceGeneration::new([11; 32]).unwrap();
    582         let digest = RawSourceDigest::new([12; 32]);
    583         let ticket = RebuildTicket::requested(
    584             RebuildTicketId::new([13; 16]).unwrap(),
    585             invalidation,
    586             source_generation,
    587             None,
    588             digest,
    589         )
    590         .unwrap();
    591         let matching = RawSourceSnapshot {
    592             generation: source_generation,
    593             high_water: None,
    594             digest,
    595         };
    596         assert!(matching.matches_ticket(&ticket));
    597         assert!(
    598             !RawSourceSnapshot {
    599                 generation: SourceGeneration::new([14; 32]).unwrap(),
    600                 ..matching
    601             }
    602             .matches_ticket(&ticket)
    603         );
    604         assert!(
    605             !RawSourceSnapshot {
    606                 high_water: Some(EventPosition::new(
    607                     source_generation,
    608                     radroots_storage::event::EventSequence::new(1).unwrap(),
    609                 )),
    610                 ..matching
    611             }
    612             .matches_ticket(&ticket)
    613         );
    614         assert!(
    615             !RawSourceSnapshot {
    616                 digest: RawSourceDigest::new([15; 32]),
    617                 ..matching
    618             }
    619             .matches_ticket(&ticket)
    620         );
    621         assert_eq!(admission_stage_byte(AdmissionStage::Raw), 0);
    622         assert_eq!(admission_stage_byte(AdmissionStage::Verified), 1);
    623         assert_eq!(admission_stage_byte(AdmissionStage::Visible), 2);
    624         assert_eq!(
    625             map_storage_error(StorageError::ProjectionRevisionConflict),
    626             Error::StorageConflict
    627         );
    628         assert_eq!(
    629             map_storage_error(StorageError::BackendUnavailable),
    630             Error::StorageFailed
    631         );
    632     }
    633 }