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 }