lib

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

projection.rs (40769B)


      1 //! Projection checkpoint, event-index manifest, and rebuild contracts.
      2 //!
      3 //! Storage owns durable coordination metadata. Domain reducers and projected
      4 //! row representations remain in their domain packages.
      5 
      6 pub mod document_query;
      7 
      8 pub use radroots_event::EventId;
      9 pub use radroots_transport::BoxFuture;
     10 use sha2::{Digest, Sha256};
     11 use std::collections::BTreeSet;
     12 
     13 use crate::{
     14     Error,
     15     event::{EventPosition, SourceGeneration},
     16 };
     17 
     18 pub const PROJECTION_ID_MAX_BYTES: usize = 128;
     19 pub const EVENT_INDEX_SHARD_ID_MAX_BYTES: usize = 128;
     20 pub const EVENT_INDEX_ARTIFACT_PATH_MAX_BYTES: usize = 512;
     21 pub const EVENT_INDEX_CURSOR_MAX_BYTES: usize = 2_048;
     22 pub const EVENT_INDEX_SHARDS_MAX: usize = 4_096;
     23 
     24 /// Stable, backend-neutral projection identity.
     25 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     26 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     27 pub struct ProjectionId(String);
     28 
     29 impl ProjectionId {
     30     pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
     31         let value = value.into();
     32         if !valid_label(value.as_str(), PROJECTION_ID_MAX_BYTES) {
     33             return Err(Error::InvalidProjectionId);
     34         }
     35         Ok(Self(value))
     36     }
     37 
     38     pub fn as_str(&self) -> &str {
     39         self.0.as_str()
     40     }
     41 }
     42 
     43 /// Content-derived generation of a projection implementation and its inputs.
     44 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     45 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     46 pub struct ProjectionGeneration([u8; 32]);
     47 
     48 impl ProjectionGeneration {
     49     pub const fn new(bytes: [u8; 32]) -> Result<Self, Error> {
     50         if bytes32_are_zero(&bytes) {
     51             return Err(Error::InvalidProjectionGeneration);
     52         }
     53         Ok(Self(bytes))
     54     }
     55     pub const fn as_bytes(&self) -> &[u8; 32] {
     56         &self.0
     57     }
     58 }
     59 
     60 /// SHA-256 digest of one ordered, immutable canonical raw-event snapshot.
     61 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     62 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     63 pub struct RawSourceDigest([u8; 32]);
     64 
     65 impl RawSourceDigest {
     66     pub const fn new(bytes: [u8; 32]) -> Self {
     67         Self(bytes)
     68     }
     69 
     70     pub const fn as_bytes(&self) -> &[u8; 32] {
     71         &self.0
     72     }
     73 }
     74 
     75 /// Non-zero optimistic revision for projection coordination state.
     76 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     77 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
     78 pub struct ProjectionRevision(u64);
     79 
     80 impl ProjectionRevision {
     81     pub const INITIAL: Self = Self(1);
     82     pub const fn new(value: u64) -> Result<Self, Error> {
     83         if value == 0 {
     84             return Err(Error::InvalidProjectionRevision);
     85         }
     86         Ok(Self(value))
     87     }
     88     pub const fn get(self) -> u64 {
     89         self.0
     90     }
     91     fn next(self) -> Result<Self, Error> {
     92         self.0
     93             .checked_add(1)
     94             .map(Self)
     95             .ok_or(Error::CorruptProjectionRecord)
     96     }
     97 }
     98 
     99 /// Last canonical event incorporated by a projection generation.
    100 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    101 #[derive(Clone, Debug, Eq, PartialEq)]
    102 pub struct ProjectionCheckpoint {
    103     projection_id: ProjectionId,
    104     generation: ProjectionGeneration,
    105     source_position: Option<EventPosition>,
    106     projected_rows: u64,
    107     updated_at_unix_ms: u64,
    108 }
    109 
    110 impl ProjectionCheckpoint {
    111     pub fn new(
    112         projection_id: ProjectionId,
    113         generation: ProjectionGeneration,
    114         source_position: Option<EventPosition>,
    115         projected_rows: u64,
    116         updated_at_unix_ms: u64,
    117     ) -> Result<Self, Error> {
    118         if updated_at_unix_ms == 0 {
    119             return Err(Error::InvalidProjectionTimestamp);
    120         }
    121         Ok(Self {
    122             projection_id,
    123             generation,
    124             source_position,
    125             projected_rows,
    126             updated_at_unix_ms,
    127         })
    128     }
    129     pub const fn projection_id(&self) -> &ProjectionId {
    130         &self.projection_id
    131     }
    132     pub const fn generation(&self) -> ProjectionGeneration {
    133         self.generation
    134     }
    135     pub const fn source_position(&self) -> Option<EventPosition> {
    136         self.source_position
    137     }
    138     pub const fn projected_rows(&self) -> u64 {
    139         self.projected_rows
    140     }
    141     pub const fn updated_at_unix_ms(&self) -> u64 {
    142         self.updated_at_unix_ms
    143     }
    144 
    145     pub fn advances(&self, prior: &Self) -> bool {
    146         self.projection_id == prior.projection_id
    147             && self.generation == prior.generation
    148             && self.updated_at_unix_ms >= prior.updated_at_unix_ms
    149             && self.projected_rows >= prior.projected_rows
    150             && match (self.source_position, prior.source_position) {
    151                 (Some(next), Some(previous)) => {
    152                     next.generation() == previous.generation()
    153                         && next.sequence() >= previous.sequence()
    154                 }
    155                 (Some(_), None) | (None, None) => true,
    156                 (None, Some(_)) => false,
    157             }
    158     }
    159 }
    160 
    161 /// Stable reason a projection can no longer be trusted.
    162 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    163 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    164 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    165 pub enum InvalidationReason {
    166     SourceGenerationChanged,
    167     ProjectionGenerationChanged,
    168     EventIndexManifestChanged,
    169     IntegrityFailure,
    170     OperatorRequested,
    171 }
    172 
    173 /// Durable projection invalidation evidence.
    174 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    175 #[derive(Clone, Debug, Eq, PartialEq)]
    176 pub struct ProjectionInvalidation {
    177     projection_id: ProjectionId,
    178     invalid_generation: ProjectionGeneration,
    179     replacement_generation: ProjectionGeneration,
    180     reason: InvalidationReason,
    181     invalidated_at_unix_ms: u64,
    182 }
    183 
    184 impl ProjectionInvalidation {
    185     pub fn new(
    186         projection_id: ProjectionId,
    187         invalid_generation: ProjectionGeneration,
    188         replacement_generation: ProjectionGeneration,
    189         reason: InvalidationReason,
    190         invalidated_at_unix_ms: u64,
    191     ) -> Result<Self, Error> {
    192         if invalid_generation.0 == replacement_generation.0 || invalidated_at_unix_ms == 0 {
    193             return Err(Error::InvalidProjectionInvalidation);
    194         }
    195         Ok(Self {
    196             projection_id,
    197             invalid_generation,
    198             replacement_generation,
    199             reason,
    200             invalidated_at_unix_ms,
    201         })
    202     }
    203     pub const fn projection_id(&self) -> &ProjectionId {
    204         &self.projection_id
    205     }
    206     pub const fn invalid_generation(&self) -> ProjectionGeneration {
    207         self.invalid_generation
    208     }
    209     pub const fn replacement_generation(&self) -> ProjectionGeneration {
    210         self.replacement_generation
    211     }
    212     pub const fn reason(&self) -> InvalidationReason {
    213         self.reason
    214     }
    215     pub const fn invalidated_at_unix_ms(&self) -> u64 {
    216         self.invalidated_at_unix_ms
    217     }
    218 }
    219 
    220 /// Stable identity of a projection rebuild execution.
    221 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    222 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    223 pub struct RebuildTicketId([u8; 16]);
    224 
    225 impl RebuildTicketId {
    226     pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
    227         if bytes16_are_zero(&bytes) {
    228             return Err(Error::InvalidRebuildTicketId);
    229         }
    230         Ok(Self(bytes))
    231     }
    232     pub const fn as_bytes(&self) -> &[u8; 16] {
    233         &self.0
    234     }
    235 }
    236 
    237 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    238 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    239 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    240 pub enum RebuildStage {
    241     Requested,
    242     Running,
    243     Completed,
    244     Failed,
    245 }
    246 
    247 /// Stable, secret-safe classification retained for a failed rebuild.
    248 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    249 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    250 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    251 pub enum RebuildFailure {
    252     ReducerRejected,
    253     SourceChanged,
    254     IntegrityFailure,
    255     PromotionRejected,
    256 }
    257 
    258 /// Optimistic, monotonic projection rebuild state.
    259 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    260 #[derive(Clone, Debug, Eq, PartialEq)]
    261 pub struct RebuildTicket {
    262     ticket_id: RebuildTicketId,
    263     invalidation: ProjectionInvalidation,
    264     revision: ProjectionRevision,
    265     stage: RebuildStage,
    266     source_generation: SourceGeneration,
    267     source_high_water: Option<EventPosition>,
    268     source_digest: RawSourceDigest,
    269     checkpoint: Option<ProjectionCheckpoint>,
    270     failure: Option<RebuildFailure>,
    271     requested_at_unix_ms: u64,
    272     updated_at_unix_ms: u64,
    273 }
    274 
    275 impl RebuildTicket {
    276     pub fn requested(
    277         ticket_id: RebuildTicketId,
    278         invalidation: ProjectionInvalidation,
    279         source_generation: SourceGeneration,
    280         source_high_water: Option<EventPosition>,
    281         source_digest: RawSourceDigest,
    282     ) -> Result<Self, Error> {
    283         if source_high_water.is_some_and(|position| position.generation() != source_generation) {
    284             return Err(Error::SourceGenerationChanged);
    285         }
    286         let at = invalidation.invalidated_at_unix_ms();
    287         Ok(Self {
    288             ticket_id,
    289             invalidation,
    290             revision: ProjectionRevision::INITIAL,
    291             stage: RebuildStage::Requested,
    292             source_generation,
    293             source_high_water,
    294             source_digest,
    295             checkpoint: None,
    296             failure: None,
    297             requested_at_unix_ms: at,
    298             updated_at_unix_ms: at,
    299         })
    300     }
    301 
    302     /// Reconstructs and validates one durable rebuild ticket.
    303     #[allow(clippy::too_many_arguments)]
    304     pub fn from_durable_parts(
    305         ticket_id: RebuildTicketId,
    306         invalidation: ProjectionInvalidation,
    307         revision: ProjectionRevision,
    308         stage: RebuildStage,
    309         source_generation: SourceGeneration,
    310         source_high_water: Option<EventPosition>,
    311         source_digest: RawSourceDigest,
    312         checkpoint: Option<ProjectionCheckpoint>,
    313         failure: Option<RebuildFailure>,
    314         requested_at_unix_ms: u64,
    315         updated_at_unix_ms: u64,
    316     ) -> Result<Self, Error> {
    317         if source_high_water.is_some_and(|position| position.generation() != source_generation)
    318             || requested_at_unix_ms != invalidation.invalidated_at_unix_ms()
    319             || updated_at_unix_ms < requested_at_unix_ms
    320             || matches!(stage, RebuildStage::Requested)
    321                 && (revision != ProjectionRevision::INITIAL
    322                     || updated_at_unix_ms != requested_at_unix_ms)
    323             || !matches!(stage, RebuildStage::Requested) && revision == ProjectionRevision::INITIAL
    324             || matches!(stage, RebuildStage::Requested) && checkpoint.is_some()
    325             || matches!(stage, RebuildStage::Completed) && checkpoint.is_none()
    326             || matches!(stage, RebuildStage::Failed) != failure.is_some()
    327             || !matches!(stage, RebuildStage::Failed) && failure.is_some()
    328             || checkpoint.as_ref().is_some_and(|checkpoint| {
    329                 checkpoint.projection_id() != invalidation.projection_id()
    330                     || checkpoint.generation() != invalidation.replacement_generation()
    331                     || checkpoint
    332                         .source_position()
    333                         .is_some_and(|position| position.generation() != source_generation)
    334                     || checkpoint.updated_at_unix_ms() > updated_at_unix_ms
    335             })
    336         {
    337             return Err(Error::CorruptProjectionRecord);
    338         }
    339         Ok(Self {
    340             ticket_id,
    341             invalidation,
    342             revision,
    343             stage,
    344             source_generation,
    345             source_high_water,
    346             source_digest,
    347             checkpoint,
    348             failure,
    349             requested_at_unix_ms,
    350             updated_at_unix_ms,
    351         })
    352     }
    353     pub const fn ticket_id(&self) -> RebuildTicketId {
    354         self.ticket_id
    355     }
    356     pub const fn invalidation(&self) -> &ProjectionInvalidation {
    357         &self.invalidation
    358     }
    359     pub const fn revision(&self) -> ProjectionRevision {
    360         self.revision
    361     }
    362     pub const fn stage(&self) -> RebuildStage {
    363         self.stage
    364     }
    365     pub const fn source_generation(&self) -> SourceGeneration {
    366         self.source_generation
    367     }
    368     pub const fn source_high_water(&self) -> Option<EventPosition> {
    369         self.source_high_water
    370     }
    371     pub const fn source_digest(&self) -> RawSourceDigest {
    372         self.source_digest
    373     }
    374     pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
    375         self.checkpoint.as_ref()
    376     }
    377     pub const fn requested_at_unix_ms(&self) -> u64 {
    378         self.requested_at_unix_ms
    379     }
    380     pub const fn updated_at_unix_ms(&self) -> u64 {
    381         self.updated_at_unix_ms
    382     }
    383     pub const fn failure(&self) -> Option<RebuildFailure> {
    384         self.failure
    385     }
    386 
    387     pub fn transition(&self, transition: RebuildTransition) -> Result<Self, Error> {
    388         if transition.ticket_id != self.ticket_id || transition.expected_revision != self.revision {
    389             return Err(Error::ProjectionRevisionConflict);
    390         }
    391         if transition.at_unix_ms < self.updated_at_unix_ms {
    392             return Err(Error::InvalidProjectionTimestamp);
    393         }
    394         let (stage, checkpoint, failure) = match (&self.stage, transition.kind) {
    395             (RebuildStage::Requested, RebuildTransitionKind::Start) => {
    396                 (RebuildStage::Running, None, None)
    397             }
    398             (RebuildStage::Running, RebuildTransitionKind::Checkpoint(checkpoint)) => {
    399                 self.validate_checkpoint(&checkpoint)?;
    400                 if self
    401                     .checkpoint
    402                     .as_ref()
    403                     .is_some_and(|prior| !checkpoint.advances(prior))
    404                 {
    405                     return Err(Error::ProjectionCheckpointRegression);
    406                 }
    407                 (RebuildStage::Running, Some(checkpoint), None)
    408             }
    409             (RebuildStage::Running, RebuildTransitionKind::Complete(checkpoint)) => {
    410                 self.validate_checkpoint(&checkpoint)?;
    411                 if self
    412                     .checkpoint
    413                     .as_ref()
    414                     .is_some_and(|prior| !checkpoint.advances(prior))
    415                 {
    416                     return Err(Error::ProjectionCheckpointRegression);
    417                 }
    418                 (RebuildStage::Completed, Some(checkpoint), None)
    419             }
    420             (
    421                 RebuildStage::Requested | RebuildStage::Running,
    422                 RebuildTransitionKind::Fail(failure),
    423             ) => (RebuildStage::Failed, self.checkpoint.clone(), Some(failure)),
    424             (RebuildStage::Completed | RebuildStage::Failed, _) => {
    425                 return Err(Error::RebuildTicketTerminal);
    426             }
    427             _ => return Err(Error::InvalidRebuildTransition),
    428         };
    429         Ok(Self {
    430             ticket_id: self.ticket_id,
    431             invalidation: self.invalidation.clone(),
    432             revision: self.revision.next()?,
    433             stage,
    434             source_generation: self.source_generation,
    435             source_high_water: self.source_high_water,
    436             source_digest: self.source_digest,
    437             checkpoint,
    438             failure,
    439             requested_at_unix_ms: self.requested_at_unix_ms,
    440             updated_at_unix_ms: transition.at_unix_ms,
    441         })
    442     }
    443 
    444     fn validate_checkpoint(&self, checkpoint: &ProjectionCheckpoint) -> Result<(), Error> {
    445         if checkpoint.projection_id() != self.invalidation.projection_id()
    446             || checkpoint.generation() != self.invalidation.replacement_generation()
    447             || checkpoint
    448                 .source_position()
    449                 .is_some_and(|position| position.generation() != self.source_generation)
    450         {
    451             return Err(Error::ProjectionCheckpointMismatch);
    452         }
    453         Ok(())
    454     }
    455 }
    456 
    457 #[derive(Clone, Debug, Eq, PartialEq)]
    458 pub struct RebuildTransition {
    459     ticket_id: RebuildTicketId,
    460     expected_revision: ProjectionRevision,
    461     at_unix_ms: u64,
    462     kind: RebuildTransitionKind,
    463 }
    464 
    465 #[derive(Clone, Debug, Eq, PartialEq)]
    466 enum RebuildTransitionKind {
    467     Start,
    468     Checkpoint(ProjectionCheckpoint),
    469     Complete(ProjectionCheckpoint),
    470     Fail(RebuildFailure),
    471 }
    472 
    473 impl RebuildTransition {
    474     pub const fn start(
    475         ticket_id: RebuildTicketId,
    476         expected_revision: ProjectionRevision,
    477         at_unix_ms: u64,
    478     ) -> Self {
    479         Self {
    480             ticket_id,
    481             expected_revision,
    482             at_unix_ms,
    483             kind: RebuildTransitionKind::Start,
    484         }
    485     }
    486     pub const fn checkpoint(
    487         ticket_id: RebuildTicketId,
    488         expected_revision: ProjectionRevision,
    489         at_unix_ms: u64,
    490         checkpoint: ProjectionCheckpoint,
    491     ) -> Self {
    492         Self {
    493             ticket_id,
    494             expected_revision,
    495             at_unix_ms,
    496             kind: RebuildTransitionKind::Checkpoint(checkpoint),
    497         }
    498     }
    499     pub const fn complete(
    500         ticket_id: RebuildTicketId,
    501         expected_revision: ProjectionRevision,
    502         at_unix_ms: u64,
    503         checkpoint: ProjectionCheckpoint,
    504     ) -> Self {
    505         Self {
    506             ticket_id,
    507             expected_revision,
    508             at_unix_ms,
    509             kind: RebuildTransitionKind::Complete(checkpoint),
    510         }
    511     }
    512     pub const fn fail(
    513         ticket_id: RebuildTicketId,
    514         expected_revision: ProjectionRevision,
    515         at_unix_ms: u64,
    516         failure: RebuildFailure,
    517     ) -> Self {
    518         Self {
    519             ticket_id,
    520             expected_revision,
    521             at_unix_ms,
    522             kind: RebuildTransitionKind::Fail(failure),
    523         }
    524     }
    525     pub const fn ticket_id(&self) -> RebuildTicketId {
    526         self.ticket_id
    527     }
    528 }
    529 
    530 /// SHA-256 digest of an immutable event-index shard artifact.
    531 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    532 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    533 pub struct ArtifactDigest([u8; 32]);
    534 
    535 impl ArtifactDigest {
    536     pub const fn new(bytes: [u8; 32]) -> Self {
    537         Self(bytes)
    538     }
    539     pub const fn as_bytes(&self) -> &[u8; 32] {
    540         &self.0
    541     }
    542 }
    543 
    544 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    545 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    546 pub struct EventIndexShardId(String);
    547 
    548 impl EventIndexShardId {
    549     pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
    550         let value = value.into();
    551         if !valid_label(value.as_str(), EVENT_INDEX_SHARD_ID_MAX_BYTES) {
    552             return Err(Error::InvalidEventIndexShardId);
    553         }
    554         Ok(Self(value))
    555     }
    556     pub fn as_str(&self) -> &str {
    557         self.0.as_str()
    558     }
    559 }
    560 
    561 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    562 #[derive(Clone, Debug, Eq, PartialEq)]
    563 pub struct EventIdRange {
    564     first: EventId,
    565     last: EventId,
    566 }
    567 
    568 impl EventIdRange {
    569     pub fn new(first: EventId, last: EventId) -> Result<Self, Error> {
    570         if first > last {
    571             return Err(Error::InvalidEventIndexRange);
    572         }
    573         Ok(Self { first, last })
    574     }
    575     pub const fn first(&self) -> &EventId {
    576         &self.first
    577     }
    578     pub const fn last(&self) -> &EventId {
    579         &self.last
    580     }
    581 }
    582 
    583 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    584 #[derive(Clone, Debug, Eq, PartialEq)]
    585 pub struct EventIndexShard {
    586     shard_id: EventIndexShardId,
    587     artifact_path: String,
    588     event_count: u32,
    589     event_ids: EventIdRange,
    590     first_published_at_unix_s: u64,
    591     last_published_at_unix_s: u64,
    592     sha256: ArtifactDigest,
    593 }
    594 
    595 impl EventIndexShard {
    596     #[allow(clippy::too_many_arguments)]
    597     pub fn new(
    598         shard_id: EventIndexShardId,
    599         artifact_path: impl Into<String>,
    600         event_count: u32,
    601         event_ids: EventIdRange,
    602         first_published_at_unix_s: u64,
    603         last_published_at_unix_s: u64,
    604         sha256: ArtifactDigest,
    605     ) -> Result<Self, Error> {
    606         let artifact_path = artifact_path.into();
    607         if !valid_artifact_path(artifact_path.as_str()) {
    608             return Err(Error::InvalidEventIndexArtifactPath);
    609         }
    610         if event_count == 0 {
    611             return Err(Error::InvalidEventIndexShardCount);
    612         }
    613         if first_published_at_unix_s == 0 || last_published_at_unix_s < first_published_at_unix_s {
    614             return Err(Error::InvalidEventIndexTimestamp);
    615         }
    616         Ok(Self {
    617             shard_id,
    618             artifact_path,
    619             event_count,
    620             event_ids,
    621             first_published_at_unix_s,
    622             last_published_at_unix_s,
    623             sha256,
    624         })
    625     }
    626     pub const fn shard_id(&self) -> &EventIndexShardId {
    627         &self.shard_id
    628     }
    629     pub fn artifact_path(&self) -> &str {
    630         self.artifact_path.as_str()
    631     }
    632     pub const fn event_count(&self) -> u32 {
    633         self.event_count
    634     }
    635     pub const fn event_ids(&self) -> &EventIdRange {
    636         &self.event_ids
    637     }
    638     pub const fn first_published_at_unix_s(&self) -> u64 {
    639         self.first_published_at_unix_s
    640     }
    641     pub const fn last_published_at_unix_s(&self) -> u64 {
    642         self.last_published_at_unix_s
    643     }
    644     pub const fn sha256(&self) -> ArtifactDigest {
    645         self.sha256
    646     }
    647 }
    648 
    649 /// Validated, immutable event-index artifact inventory.
    650 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    651 #[derive(Clone, Debug, Eq, PartialEq)]
    652 pub struct EventIndexManifest {
    653     generation: ProjectionGeneration,
    654     total_events: u64,
    655     target_shard_size: u32,
    656     first_published_at_unix_s: u64,
    657     last_published_at_unix_s: u64,
    658     shards: Vec<EventIndexShard>,
    659 }
    660 
    661 impl EventIndexManifest {
    662     pub fn new(
    663         generation: ProjectionGeneration,
    664         total_events: u64,
    665         target_shard_size: u32,
    666         first_published_at_unix_s: u64,
    667         last_published_at_unix_s: u64,
    668         shards: Vec<EventIndexShard>,
    669     ) -> Result<Self, Error> {
    670         if shards.is_empty() || shards.len() > EVENT_INDEX_SHARDS_MAX {
    671             return Err(Error::InvalidEventIndexShardCount);
    672         }
    673         if target_shard_size == 0 || total_events == 0 {
    674             return Err(Error::InvalidEventIndexManifest);
    675         }
    676         let sum = shards.iter().try_fold(0_u64, |sum, shard| {
    677             if shard.event_count() > target_shard_size {
    678                 return Err(Error::InvalidEventIndexManifest);
    679             }
    680             sum.checked_add(u64::from(shard.event_count()))
    681                 .ok_or(Error::InvalidEventIndexManifest)
    682         })?;
    683         if sum != total_events
    684             || first_published_at_unix_s != shards[0].first_published_at_unix_s()
    685             || last_published_at_unix_s != shards[shards.len() - 1].last_published_at_unix_s()
    686         {
    687             return Err(Error::InvalidEventIndexManifest);
    688         }
    689         let mut shard_ids = BTreeSet::new();
    690         let mut artifact_paths = BTreeSet::new();
    691         if shards.iter().any(|shard| {
    692             !shard_ids.insert(shard.shard_id()) || !artifact_paths.insert(shard.artifact_path())
    693         }) {
    694             return Err(Error::InvalidEventIndexManifest);
    695         }
    696         for pair in shards.windows(2) {
    697             if pair[0].shard_id() >= pair[1].shard_id()
    698                 || pair[0].event_ids().last() >= pair[1].event_ids().first()
    699                 || pair[0].last_published_at_unix_s() > pair[1].first_published_at_unix_s()
    700             {
    701                 return Err(Error::InvalidEventIndexManifest);
    702             }
    703         }
    704         Ok(Self {
    705             generation,
    706             total_events,
    707             target_shard_size,
    708             first_published_at_unix_s,
    709             last_published_at_unix_s,
    710             shards,
    711         })
    712     }
    713     pub const fn generation(&self) -> ProjectionGeneration {
    714         self.generation
    715     }
    716     pub const fn total_events(&self) -> u64 {
    717         self.total_events
    718     }
    719     pub const fn target_shard_size(&self) -> u32 {
    720         self.target_shard_size
    721     }
    722     pub const fn first_published_at_unix_s(&self) -> u64 {
    723         self.first_published_at_unix_s
    724     }
    725     pub const fn last_published_at_unix_s(&self) -> u64 {
    726         self.last_published_at_unix_s
    727     }
    728     pub fn shards(&self) -> &[EventIndexShard] {
    729         self.shards.as_slice()
    730     }
    731 }
    732 
    733 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    734 #[derive(Clone, Debug, Eq, PartialEq)]
    735 pub struct EventIndexShardCheckpoint {
    736     shard_id: EventIndexShardId,
    737     last_created_at_unix_s: u64,
    738     last_event_id: Option<EventId>,
    739     cursor: Option<String>,
    740 }
    741 
    742 impl EventIndexShardCheckpoint {
    743     pub fn new(
    744         shard_id: EventIndexShardId,
    745         last_created_at_unix_s: u64,
    746         last_event_id: Option<EventId>,
    747         cursor: Option<String>,
    748     ) -> Result<Self, Error> {
    749         if last_created_at_unix_s == 0 {
    750             return Err(Error::InvalidEventIndexTimestamp);
    751         }
    752         if let Some(value) = cursor.as_deref()
    753             && (value.is_empty()
    754                 || value.len() > EVENT_INDEX_CURSOR_MAX_BYTES
    755                 || value != value.trim()
    756                 || value.chars().any(char::is_control))
    757         {
    758             return Err(Error::InvalidEventIndexCursor);
    759         }
    760         Ok(Self {
    761             shard_id,
    762             last_created_at_unix_s,
    763             last_event_id,
    764             cursor,
    765         })
    766     }
    767     pub const fn shard_id(&self) -> &EventIndexShardId {
    768         &self.shard_id
    769     }
    770     pub const fn last_created_at_unix_s(&self) -> u64 {
    771         self.last_created_at_unix_s
    772     }
    773     pub const fn last_event_id(&self) -> Option<&EventId> {
    774         self.last_event_id.as_ref()
    775     }
    776     pub fn cursor(&self) -> Option<&str> {
    777         self.cursor.as_deref()
    778     }
    779 }
    780 
    781 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    782 #[derive(Clone, Debug, Eq, PartialEq)]
    783 pub struct EventIndexCheckpoint {
    784     generation: ProjectionGeneration,
    785     generated_at_unix_ms: u64,
    786     shards: Vec<EventIndexShardCheckpoint>,
    787 }
    788 
    789 impl EventIndexCheckpoint {
    790     pub fn new(
    791         generation: ProjectionGeneration,
    792         generated_at_unix_ms: u64,
    793         mut shards: Vec<EventIndexShardCheckpoint>,
    794     ) -> Result<Self, Error> {
    795         if generated_at_unix_ms == 0 || shards.len() > EVENT_INDEX_SHARDS_MAX {
    796             return Err(Error::InvalidEventIndexCheckpoint);
    797         }
    798         shards.sort_by(|left, right| left.shard_id().cmp(right.shard_id()));
    799         if shards
    800             .windows(2)
    801             .any(|pair| pair[0].shard_id() == pair[1].shard_id())
    802         {
    803             return Err(Error::DuplicateEventIndexShard);
    804         }
    805         Ok(Self {
    806             generation,
    807             generated_at_unix_ms,
    808             shards,
    809         })
    810     }
    811     pub const fn generation(&self) -> ProjectionGeneration {
    812         self.generation
    813     }
    814     pub const fn generated_at_unix_ms(&self) -> u64 {
    815         self.generated_at_unix_ms
    816     }
    817     pub fn shards(&self) -> &[EventIndexShardCheckpoint] {
    818         self.shards.as_slice()
    819     }
    820     pub fn shard(&self, id: &EventIndexShardId) -> Option<&EventIndexShardCheckpoint> {
    821         self.shards
    822             .binary_search_by(|candidate| candidate.shard_id().cmp(id))
    823             .ok()
    824             .map(|index| &self.shards[index])
    825     }
    826 }
    827 
    828 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    829 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    830 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    831 pub enum ProjectionHealth {
    832     Ready,
    833     Invalidated,
    834     Rebuilding,
    835     Failed,
    836 }
    837 
    838 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    839 #[derive(Clone, Debug, Eq, PartialEq)]
    840 pub struct ProjectionStatus {
    841     projection_id: ProjectionId,
    842     generation: ProjectionGeneration,
    843     health: ProjectionHealth,
    844     checkpoint: Option<ProjectionCheckpoint>,
    845     active_rebuild: Option<RebuildTicketId>,
    846 }
    847 
    848 /// Maximum UTF-8 bytes in one materialized projection document key.
    849 pub const PROJECTION_DOCUMENT_KEY_MAX_BYTES: usize = 512;
    850 /// Maximum bytes in one materialized projection document value.
    851 pub const PROJECTION_DOCUMENT_VALUE_MAX_BYTES: usize = 16 * 1024 * 1024;
    852 
    853 /// One backend-neutral, opaque materialized projection document.
    854 ///
    855 /// Projection owners define the value encoding. Storage verifies its digest
    856 /// and treats the bytes as opaque so product-specific DTOs do not leak into
    857 /// the generic persistence boundary.
    858 #[derive(Clone, Debug, Eq, PartialEq)]
    859 pub struct ProjectionDocument {
    860     key: String,
    861     value: Vec<u8>,
    862     value_sha256: [u8; 32],
    863 }
    864 
    865 impl ProjectionDocument {
    866     pub fn new(key: String, value: Vec<u8>) -> Result<Self, Error> {
    867         if !valid_document_key(&key)
    868             || value.is_empty()
    869             || value.len() > PROJECTION_DOCUMENT_VALUE_MAX_BYTES
    870         {
    871             return Err(Error::InvalidProjectionDocument);
    872         }
    873         let value_sha256 = Sha256::digest(&value).into();
    874         Ok(Self {
    875             key,
    876             value,
    877             value_sha256,
    878         })
    879     }
    880 
    881     pub fn from_stored_parts(
    882         key: String,
    883         value: Vec<u8>,
    884         value_sha256: [u8; 32],
    885     ) -> Result<Self, Error> {
    886         let document = Self::new(key, value)?;
    887         if document.value_sha256 != value_sha256 {
    888             return Err(Error::CorruptProjectionDocument);
    889         }
    890         Ok(document)
    891     }
    892 
    893     pub fn key(&self) -> &str {
    894         &self.key
    895     }
    896 
    897     pub fn value(&self) -> &[u8] {
    898         &self.value
    899     }
    900 
    901     pub const fn value_sha256(&self) -> &[u8; 32] {
    902         &self.value_sha256
    903     }
    904 }
    905 
    906 /// One immutable, durable frozen-query snapshot.
    907 #[derive(Clone, Debug, Eq, PartialEq)]
    908 pub struct ProjectionSnapshot {
    909     projection_id: ProjectionId,
    910     snapshot_id: [u8; 32],
    911     generation: ProjectionGeneration,
    912     created_at_unix_ms: u64,
    913     value: Vec<u8>,
    914     value_sha256: [u8; 32],
    915 }
    916 
    917 impl ProjectionSnapshot {
    918     pub fn new(
    919         projection_id: ProjectionId,
    920         snapshot_id: [u8; 32],
    921         generation: ProjectionGeneration,
    922         created_at_unix_ms: u64,
    923         value: Vec<u8>,
    924     ) -> Result<Self, Error> {
    925         if snapshot_id.iter().all(|byte| *byte == 0)
    926             || created_at_unix_ms == 0
    927             || value.is_empty()
    928             || value.len() > PROJECTION_DOCUMENT_VALUE_MAX_BYTES
    929         {
    930             return Err(Error::InvalidProjectionSnapshot);
    931         }
    932         let value_sha256 = Sha256::digest(&value).into();
    933         Ok(Self {
    934             projection_id,
    935             snapshot_id,
    936             generation,
    937             created_at_unix_ms,
    938             value,
    939             value_sha256,
    940         })
    941     }
    942 
    943     pub fn from_stored_parts(
    944         projection_id: ProjectionId,
    945         snapshot_id: [u8; 32],
    946         generation: ProjectionGeneration,
    947         created_at_unix_ms: u64,
    948         value: Vec<u8>,
    949         value_sha256: [u8; 32],
    950     ) -> Result<Self, Error> {
    951         let snapshot = Self::new(
    952             projection_id,
    953             snapshot_id,
    954             generation,
    955             created_at_unix_ms,
    956             value,
    957         )?;
    958         if snapshot.value_sha256 != value_sha256 {
    959             return Err(Error::CorruptProjectionDocument);
    960         }
    961         Ok(snapshot)
    962     }
    963 
    964     pub const fn projection_id(&self) -> &ProjectionId {
    965         &self.projection_id
    966     }
    967 
    968     pub const fn snapshot_id(&self) -> &[u8; 32] {
    969         &self.snapshot_id
    970     }
    971 
    972     pub const fn generation(&self) -> ProjectionGeneration {
    973         self.generation
    974     }
    975 
    976     pub const fn created_at_unix_ms(&self) -> u64 {
    977         self.created_at_unix_ms
    978     }
    979 
    980     pub fn value(&self) -> &[u8] {
    981         &self.value
    982     }
    983 
    984     pub const fn value_sha256(&self) -> &[u8; 32] {
    985         &self.value_sha256
    986     }
    987 }
    988 
    989 impl ProjectionStatus {
    990     pub fn new(
    991         projection_id: ProjectionId,
    992         generation: ProjectionGeneration,
    993         health: ProjectionHealth,
    994         checkpoint: Option<ProjectionCheckpoint>,
    995         active_rebuild: Option<RebuildTicketId>,
    996     ) -> Result<Self, Error> {
    997         if checkpoint.as_ref().is_some_and(|value| {
    998             value.projection_id() != &projection_id || value.generation() != generation
    999         }) || (health == ProjectionHealth::Rebuilding) != active_rebuild.is_some()
   1000         {
   1001             return Err(Error::CorruptProjectionRecord);
   1002         }
   1003         Ok(Self {
   1004             projection_id,
   1005             generation,
   1006             health,
   1007             checkpoint,
   1008             active_rebuild,
   1009         })
   1010     }
   1011     pub const fn projection_id(&self) -> &ProjectionId {
   1012         &self.projection_id
   1013     }
   1014     pub const fn generation(&self) -> ProjectionGeneration {
   1015         self.generation
   1016     }
   1017     pub const fn health(&self) -> ProjectionHealth {
   1018         self.health
   1019     }
   1020     pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
   1021         self.checkpoint.as_ref()
   1022     }
   1023     pub const fn active_rebuild(&self) -> Option<RebuildTicketId> {
   1024         self.active_rebuild
   1025     }
   1026 }
   1027 
   1028 /// Backend-neutral projection coordination SPI.
   1029 pub trait ProjectionStore: Send + Sync {
   1030     fn status(
   1031         &self,
   1032         projection_id: ProjectionId,
   1033     ) -> BoxFuture<'_, Result<Option<ProjectionStatus>, Error>>;
   1034     fn checkpoint(
   1035         &self,
   1036         checkpoint: ProjectionCheckpoint,
   1037     ) -> BoxFuture<'_, Result<ProjectionStatus, Error>>;
   1038     fn invalidate(
   1039         &self,
   1040         invalidation: ProjectionInvalidation,
   1041     ) -> BoxFuture<'_, Result<ProjectionStatus, Error>>;
   1042     /// Returns the latest durable invalidation selecting a replacement generation.
   1043     fn invalidation(
   1044         &self,
   1045         projection_id: ProjectionId,
   1046         replacement_generation: ProjectionGeneration,
   1047     ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>>;
   1048     fn request_rebuild(&self, ticket: RebuildTicket)
   1049     -> BoxFuture<'_, Result<RebuildTicket, Error>>;
   1050     /// Returns one durable rebuild execution by identity.
   1051     fn rebuild(
   1052         &self,
   1053         ticket_id: RebuildTicketId,
   1054     ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>>;
   1055     fn transition_rebuild(
   1056         &self,
   1057         transition: RebuildTransition,
   1058     ) -> BoxFuture<'_, Result<RebuildTicket, Error>>;
   1059     fn event_index_manifest(
   1060         &self,
   1061         generation: ProjectionGeneration,
   1062     ) -> BoxFuture<'_, Result<Option<EventIndexManifest>, Error>>;
   1063     fn put_event_index_manifest(
   1064         &self,
   1065         manifest: EventIndexManifest,
   1066     ) -> BoxFuture<'_, Result<(), Error>>;
   1067     fn event_index_checkpoint(
   1068         &self,
   1069         generation: ProjectionGeneration,
   1070     ) -> BoxFuture<'_, Result<Option<EventIndexCheckpoint>, Error>>;
   1071     fn put_event_index_checkpoint(
   1072         &self,
   1073         checkpoint: EventIndexCheckpoint,
   1074     ) -> BoxFuture<'_, Result<(), Error>>;
   1075     /// Replaces one named materialized document for a projection generation.
   1076     fn put_projection_document(
   1077         &self,
   1078         projection_id: ProjectionId,
   1079         generation: ProjectionGeneration,
   1080         document: ProjectionDocument,
   1081     ) -> BoxFuture<'_, Result<(), Error>>;
   1082     /// Loads one named materialized document for an exact generation.
   1083     fn projection_document(
   1084         &self,
   1085         projection_id: ProjectionId,
   1086         generation: ProjectionGeneration,
   1087         key: String,
   1088     ) -> BoxFuture<'_, Result<Option<ProjectionDocument>, Error>>;
   1089     /// Inventories are live bounded scans. Unsupported backends fail explicitly.
   1090     fn query_projection_documents(
   1091         &self,
   1092         _query: document_query::ProjectionDocumentQuery,
   1093     ) -> BoxFuture<'_, Result<document_query::ProjectionDocumentPage, Error>> {
   1094         Box::pin(async { Err(Error::BackendUnavailable) })
   1095     }
   1096     /// Persists one immutable frozen-query snapshot idempotently.
   1097     fn put_projection_snapshot(
   1098         &self,
   1099         snapshot: ProjectionSnapshot,
   1100     ) -> BoxFuture<'_, Result<(), Error>>;
   1101     /// Loads one immutable frozen-query snapshot by exact identity.
   1102     fn projection_snapshot(
   1103         &self,
   1104         projection_id: ProjectionId,
   1105         snapshot_id: [u8; 32],
   1106     ) -> BoxFuture<'_, Result<Option<ProjectionSnapshot>, Error>>;
   1107 }
   1108 
   1109 fn valid_label(value: &str, max: usize) -> bool {
   1110     !value.is_empty()
   1111         && value.len() <= max
   1112         && value == value.trim()
   1113         && value.bytes().all(|byte| {
   1114             byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-' | b'.')
   1115         })
   1116 }
   1117 
   1118 fn valid_document_key(value: &str) -> bool {
   1119     !value.is_empty()
   1120         && value.len() <= PROJECTION_DOCUMENT_KEY_MAX_BYTES
   1121         && value == value.trim()
   1122         && !value.chars().any(char::is_control)
   1123 }
   1124 
   1125 fn valid_artifact_path(value: &str) -> bool {
   1126     !value.is_empty()
   1127         && value.len() <= EVENT_INDEX_ARTIFACT_PATH_MAX_BYTES
   1128         && value == value.trim()
   1129         && !value.starts_with('/')
   1130         && !value.contains('\\')
   1131         && value.split('/').all(|part| {
   1132             !part.is_empty() && part != "." && part != ".." && !part.chars().any(char::is_control)
   1133         })
   1134 }
   1135 
   1136 const fn bytes16_are_zero(bytes: &[u8; 16]) -> bool {
   1137     let mut index = 0;
   1138     while index < bytes.len() {
   1139         if bytes[index] != 0 {
   1140             return false;
   1141         }
   1142         index += 1;
   1143     }
   1144     true
   1145 }
   1146 
   1147 const fn bytes32_are_zero(bytes: &[u8; 32]) -> bool {
   1148     let mut index = 0;
   1149     while index < bytes.len() {
   1150         if bytes[index] != 0 {
   1151             return false;
   1152         }
   1153         index += 1;
   1154     }
   1155     true
   1156 }
   1157 
   1158 #[cfg(test)]
   1159 mod materialized_tests {
   1160     use super::*;
   1161 
   1162     #[test]
   1163     fn materialized_document_and_snapshot_bounds_and_digests_fail_closed() {
   1164         assert_eq!(
   1165             ProjectionDocument::new(String::new(), vec![1]),
   1166             Err(Error::InvalidProjectionDocument)
   1167         );
   1168         assert_eq!(
   1169             ProjectionDocument::new("key".into(), Vec::new()),
   1170             Err(Error::InvalidProjectionDocument)
   1171         );
   1172         assert_eq!(
   1173             ProjectionDocument::new("k".repeat(PROJECTION_DOCUMENT_KEY_MAX_BYTES + 1), vec![1],),
   1174             Err(Error::InvalidProjectionDocument)
   1175         );
   1176         assert_eq!(
   1177             ProjectionDocument::new(" key".into(), vec![1]),
   1178             Err(Error::InvalidProjectionDocument)
   1179         );
   1180         assert_eq!(
   1181             ProjectionDocument::new("key\npart".into(), vec![1]),
   1182             Err(Error::InvalidProjectionDocument)
   1183         );
   1184         assert_eq!(
   1185             ProjectionDocument::new(
   1186                 "key".into(),
   1187                 vec![0; PROJECTION_DOCUMENT_VALUE_MAX_BYTES + 1],
   1188             ),
   1189             Err(Error::InvalidProjectionDocument)
   1190         );
   1191         let document = ProjectionDocument::new("key".into(), vec![1, 2]).unwrap();
   1192         assert_eq!(document.key(), "key");
   1193         assert_eq!(document.value(), [1, 2]);
   1194         assert_eq!(document.value_sha256().len(), 32);
   1195         assert_eq!(
   1196             ProjectionDocument::from_stored_parts("key".into(), vec![1, 2], [9; 32]),
   1197             Err(Error::CorruptProjectionDocument)
   1198         );
   1199 
   1200         let projection_id = ProjectionId::parse("today").unwrap();
   1201         let generation = ProjectionGeneration::new([1; 32]).unwrap();
   1202         assert_eq!(
   1203             ProjectionSnapshot::new(projection_id.clone(), [0; 32], generation, 1, vec![1]),
   1204             Err(Error::InvalidProjectionSnapshot)
   1205         );
   1206         assert_eq!(
   1207             ProjectionSnapshot::new(projection_id.clone(), [2; 32], generation, 0, vec![1]),
   1208             Err(Error::InvalidProjectionSnapshot)
   1209         );
   1210         assert_eq!(
   1211             ProjectionSnapshot::new(projection_id.clone(), [2; 32], generation, 1, Vec::new()),
   1212             Err(Error::InvalidProjectionSnapshot)
   1213         );
   1214         assert_eq!(
   1215             ProjectionSnapshot::new(
   1216                 projection_id.clone(),
   1217                 [2; 32],
   1218                 generation,
   1219                 1,
   1220                 vec![0; PROJECTION_DOCUMENT_VALUE_MAX_BYTES + 1],
   1221             ),
   1222             Err(Error::InvalidProjectionSnapshot)
   1223         );
   1224         let snapshot =
   1225             ProjectionSnapshot::new(projection_id.clone(), [2; 32], generation, 1, vec![1])
   1226                 .unwrap();
   1227         assert_eq!(snapshot.projection_id(), &projection_id);
   1228         assert_eq!(snapshot.snapshot_id(), &[2; 32]);
   1229         assert_eq!(snapshot.generation(), generation);
   1230         assert_eq!(snapshot.created_at_unix_ms(), 1);
   1231         assert_eq!(snapshot.value(), [1]);
   1232         assert_eq!(snapshot.value_sha256().len(), 32);
   1233         assert_eq!(
   1234             ProjectionSnapshot::from_stored_parts(
   1235                 projection_id,
   1236                 [2; 32],
   1237                 generation,
   1238                 1,
   1239                 vec![1],
   1240                 [9; 32],
   1241             ),
   1242             Err(Error::CorruptProjectionDocument)
   1243         );
   1244     }
   1245 }