projection.rs (24675B)
1 use radroots_event::EventId; 2 use radroots_storage::{ 3 Error, ProjectionStore, 4 event::{EventPosition, EventSequence, SourceGeneration}, 5 projection::{ 6 ArtifactDigest, EventIdRange, EventIndexCheckpoint, EventIndexManifest, EventIndexShard, 7 EventIndexShardCheckpoint, EventIndexShardId, InvalidationReason, ProjectionCheckpoint, 8 ProjectionGeneration, ProjectionHealth, ProjectionId, ProjectionInvalidation, 9 ProjectionRevision, ProjectionStatus, RawSourceDigest, RebuildFailure, RebuildStage, 10 RebuildTicket, RebuildTicketId, RebuildTransition, 11 }, 12 }; 13 14 fn projection_id() -> ProjectionId { 15 ProjectionId::parse("food_availability.v1").expect("projection id") 16 } 17 18 fn generation(byte: u8) -> ProjectionGeneration { 19 ProjectionGeneration::new([byte; 32]).expect("projection generation") 20 } 21 22 fn event_id(character: char) -> EventId { 23 EventId::parse(character.to_string().repeat(64)).expect("event id") 24 } 25 26 fn position(source_byte: u8, sequence: u64) -> EventPosition { 27 EventPosition::new( 28 SourceGeneration::new([source_byte; 32]).expect("source generation"), 29 EventSequence::new(sequence).expect("event sequence"), 30 ) 31 } 32 33 fn checkpoint( 34 generation: ProjectionGeneration, 35 sequence: u64, 36 rows: u64, 37 at: u64, 38 ) -> ProjectionCheckpoint { 39 ProjectionCheckpoint::new( 40 projection_id(), 41 generation, 42 Some(position(9, sequence)), 43 rows, 44 at, 45 ) 46 .expect("projection checkpoint") 47 } 48 49 fn requested_ticket( 50 ticket_id: RebuildTicketId, 51 invalidation: ProjectionInvalidation, 52 ) -> RebuildTicket { 53 RebuildTicket::requested( 54 ticket_id, 55 invalidation, 56 SourceGeneration::new([9; 32]).expect("source generation"), 57 Some(position(9, 11)), 58 RawSourceDigest::new([8; 32]), 59 ) 60 .expect("requested ticket") 61 } 62 63 fn shard( 64 id: &str, 65 path: &str, 66 first: char, 67 last: char, 68 first_at: u64, 69 last_at: u64, 70 ) -> EventIndexShard { 71 EventIndexShard::new( 72 EventIndexShardId::parse(id).expect("shard id"), 73 path, 74 2, 75 EventIdRange::new(event_id(first), event_id(last)).expect("event range"), 76 first_at, 77 last_at, 78 ArtifactDigest::new([id.as_bytes()[0]; 32]), 79 ) 80 .expect("event-index shard") 81 } 82 83 #[test] 84 fn checkpoints_require_monotonic_source_and_row_progress() { 85 let prior = checkpoint(generation(1), 10, 4, 100); 86 assert!(checkpoint(generation(1), 11, 5, 101).advances(&prior)); 87 assert!(!checkpoint(generation(1), 9, 5, 101).advances(&prior)); 88 assert!(!checkpoint(generation(1), 11, 3, 101).advances(&prior)); 89 let other_source = ProjectionCheckpoint::new( 90 projection_id(), 91 generation(1), 92 Some(position(8, 11)), 93 5, 94 101, 95 ) 96 .expect("checkpoint"); 97 assert!(!other_source.advances(&prior)); 98 assert_eq!( 99 ProjectionCheckpoint::new(projection_id(), generation(1), None, 0, 0), 100 Err(Error::InvalidProjectionTimestamp) 101 ); 102 } 103 104 #[test] 105 fn manifests_validate_typed_ranges_totals_order_paths_and_bounds() { 106 let first = shard("a", "index/a.json", '0', '3', 10, 20); 107 let second = shard("b", "index/b.json", '4', '7', 20, 30); 108 let manifest = EventIndexManifest::new( 109 generation(2), 110 4, 111 2, 112 10, 113 30, 114 vec![first.clone(), second.clone()], 115 ) 116 .expect("manifest"); 117 assert_eq!(manifest.total_events(), 4); 118 assert_eq!(manifest.shards().len(), 2); 119 120 assert_eq!( 121 EventIndexManifest::new( 122 generation(2), 123 5, 124 2, 125 10, 126 30, 127 vec![first.clone(), second.clone()] 128 ), 129 Err(Error::InvalidEventIndexManifest) 130 ); 131 let overlap = shard("b", "index/b.json", '3', '7', 20, 30); 132 assert_eq!( 133 EventIndexManifest::new(generation(2), 4, 2, 10, 30, vec![first.clone(), overlap]), 134 Err(Error::InvalidEventIndexManifest) 135 ); 136 let duplicate_path = shard("b", "index/a.json", '4', '7', 20, 30); 137 assert_eq!( 138 EventIndexManifest::new(generation(2), 4, 2, 10, 30, vec![first, duplicate_path]), 139 Err(Error::InvalidEventIndexManifest) 140 ); 141 assert_eq!( 142 EventIndexShard::new( 143 EventIndexShardId::parse("unsafe").expect("id"), 144 "../escape.json", 145 1, 146 EventIdRange::new(event_id('0'), event_id('1')).expect("range"), 147 1, 148 2, 149 ArtifactDigest::new([1; 32]), 150 ), 151 Err(Error::InvalidEventIndexArtifactPath) 152 ); 153 } 154 155 #[test] 156 fn event_index_checkpoints_sort_lookup_and_reject_duplicates() { 157 let first = EventIndexShardCheckpoint::new( 158 EventIndexShardId::parse("a").expect("id"), 159 10, 160 Some(event_id('1')), 161 Some("cursor-a".to_owned()), 162 ) 163 .expect("checkpoint"); 164 let second = EventIndexShardCheckpoint::new( 165 EventIndexShardId::parse("b").expect("id"), 166 20, 167 Some(event_id('2')), 168 None, 169 ) 170 .expect("checkpoint"); 171 let checkpoint = EventIndexCheckpoint::new(generation(2), 100, vec![second, first.clone()]) 172 .expect("index checkpoint"); 173 assert_eq!( 174 checkpoint.shard(&EventIndexShardId::parse("a").expect("id")), 175 Some(&first) 176 ); 177 assert_eq!( 178 EventIndexCheckpoint::new(generation(2), 100, vec![first.clone(), first]), 179 Err(Error::DuplicateEventIndexShard) 180 ); 181 assert_eq!( 182 EventIndexShardCheckpoint::new( 183 EventIndexShardId::parse("a").expect("id"), 184 10, 185 None, 186 Some("x".repeat(2_049)), 187 ), 188 Err(Error::InvalidEventIndexCursor) 189 ); 190 } 191 192 #[test] 193 fn invalidation_and_rebuild_lifecycle_is_optimistic_and_terminal() { 194 let invalidation = radroots_storage::projection::ProjectionInvalidation::new( 195 projection_id(), 196 generation(1), 197 generation(2), 198 InvalidationReason::ProjectionGenerationChanged, 199 100, 200 ) 201 .expect("invalidation"); 202 let ticket_id = RebuildTicketId::new([7; 16]).expect("ticket id"); 203 let requested = requested_ticket(ticket_id, invalidation); 204 assert_eq!(requested.stage(), RebuildStage::Requested); 205 206 let running = requested 207 .transition(RebuildTransition::start( 208 ticket_id, 209 ProjectionRevision::INITIAL, 210 110, 211 )) 212 .expect("start rebuild"); 213 assert_eq!(running.stage(), RebuildStage::Running); 214 assert_eq!( 215 running.transition(RebuildTransition::checkpoint( 216 ticket_id, 217 ProjectionRevision::INITIAL, 218 120, 219 checkpoint(generation(2), 10, 4, 120), 220 )), 221 Err(Error::ProjectionRevisionConflict) 222 ); 223 let progressed = running 224 .transition(RebuildTransition::checkpoint( 225 ticket_id, 226 running.revision(), 227 120, 228 checkpoint(generation(2), 10, 4, 120), 229 )) 230 .expect("checkpoint rebuild"); 231 assert_eq!( 232 progressed.transition(RebuildTransition::checkpoint( 233 ticket_id, 234 progressed.revision(), 235 130, 236 checkpoint(generation(2), 9, 5, 130), 237 )), 238 Err(Error::ProjectionCheckpointRegression) 239 ); 240 let completed = progressed 241 .transition(RebuildTransition::complete( 242 ticket_id, 243 progressed.revision(), 244 140, 245 checkpoint(generation(2), 11, 5, 140), 246 )) 247 .expect("complete rebuild"); 248 assert_eq!(completed.stage(), RebuildStage::Completed); 249 assert_eq!( 250 completed.transition(RebuildTransition::fail( 251 ticket_id, 252 completed.revision(), 253 150, 254 RebuildFailure::IntegrityFailure, 255 )), 256 Err(Error::RebuildTicketTerminal) 257 ); 258 259 let status = ProjectionStatus::new( 260 projection_id(), 261 generation(2), 262 ProjectionHealth::Ready, 263 completed.checkpoint().cloned(), 264 None, 265 ) 266 .expect("ready status"); 267 assert_eq!(status.health(), ProjectionHealth::Ready); 268 } 269 270 #[test] 271 fn projection_spi_is_dyn_compatible_and_validated_identifiers_fail_closed() { 272 fn accepts_dyn(_: Option<&dyn ProjectionStore>) {} 273 accepts_dyn(None); 274 assert_eq!( 275 ProjectionId::parse("Uppercase"), 276 Err(Error::InvalidProjectionId) 277 ); 278 assert_eq!( 279 ProjectionGeneration::new([0; 32]), 280 Err(Error::InvalidProjectionGeneration) 281 ); 282 assert_eq!( 283 RebuildTicketId::new([0; 16]), 284 Err(Error::InvalidRebuildTicketId) 285 ); 286 assert_eq!( 287 radroots_storage::projection::ProjectionInvalidation::new( 288 projection_id(), 289 generation(1), 290 generation(1), 291 InvalidationReason::OperatorRequested, 292 1, 293 ), 294 Err(Error::InvalidProjectionInvalidation) 295 ); 296 } 297 298 #[test] 299 fn projection_models_cover_all_accessors_and_validation_bounds() { 300 let id = projection_id(); 301 let generation_one = generation(1); 302 assert_eq!(id.as_str(), "food_availability.v1"); 303 assert_eq!(generation_one.as_bytes(), &[1; 32]); 304 for invalid in ["", "Uppercase", " leading", "trailing ", "bad/slash"] { 305 assert_eq!( 306 ProjectionId::parse(invalid), 307 Err(Error::InvalidProjectionId) 308 ); 309 assert_eq!( 310 EventIndexShardId::parse(invalid), 311 Err(Error::InvalidEventIndexShardId) 312 ); 313 } 314 assert_eq!( 315 ProjectionId::parse("x".repeat(radroots_storage::projection::PROJECTION_ID_MAX_BYTES + 1)), 316 Err(Error::InvalidProjectionId) 317 ); 318 assert_eq!( 319 ProjectionRevision::new(0), 320 Err(Error::InvalidProjectionRevision) 321 ); 322 assert_eq!(ProjectionRevision::new(2).unwrap().get(), 2); 323 324 let checkpoint = checkpoint(generation_one, 10, 4, 100); 325 assert_eq!(checkpoint.projection_id(), &id); 326 assert_eq!(checkpoint.generation(), generation_one); 327 assert_eq!(checkpoint.source_position(), Some(position(9, 10))); 328 assert_eq!(checkpoint.projected_rows(), 4); 329 assert_eq!(checkpoint.updated_at_unix_ms(), 100); 330 let empty = ProjectionCheckpoint::new(id.clone(), generation_one, None, 0, 100).unwrap(); 331 assert!(empty.advances(&empty)); 332 assert!(checkpoint.advances(&empty)); 333 assert!(!empty.advances(&checkpoint)); 334 assert!( 335 !ProjectionCheckpoint::new(id.clone(), generation(2), None, 0, 101) 336 .unwrap() 337 .advances(&empty) 338 ); 339 assert!( 340 !ProjectionCheckpoint::new( 341 ProjectionId::parse("other").unwrap(), 342 generation_one, 343 None, 344 0, 345 101, 346 ) 347 .unwrap() 348 .advances(&empty) 349 ); 350 assert!( 351 !ProjectionCheckpoint::new(id.clone(), generation_one, None, 0, 99) 352 .unwrap() 353 .advances(&empty) 354 ); 355 356 assert_eq!( 357 ProjectionInvalidation::new( 358 id.clone(), 359 generation_one, 360 generation(2), 361 InvalidationReason::OperatorRequested, 362 0, 363 ), 364 Err(Error::InvalidProjectionInvalidation) 365 ); 366 let invalidation = ProjectionInvalidation::new( 367 id.clone(), 368 generation_one, 369 generation(2), 370 InvalidationReason::IntegrityFailure, 371 100, 372 ) 373 .unwrap(); 374 assert_eq!(invalidation.projection_id(), &id); 375 assert_eq!(invalidation.invalid_generation(), generation_one); 376 assert_eq!(invalidation.replacement_generation(), generation(2)); 377 assert_eq!(invalidation.reason(), InvalidationReason::IntegrityFailure); 378 assert_eq!(invalidation.invalidated_at_unix_ms(), 100); 379 let ticket_id = RebuildTicketId::new([3; 16]).unwrap(); 380 assert_eq!(ticket_id.as_bytes(), &[3; 16]); 381 let ticket = requested_ticket(ticket_id, invalidation); 382 assert_eq!(ticket.ticket_id(), ticket_id); 383 assert_eq!(ticket.revision(), ProjectionRevision::INITIAL); 384 assert_eq!(ticket.stage(), RebuildStage::Requested); 385 assert!(ticket.checkpoint().is_none()); 386 assert_eq!(ticket.requested_at_unix_ms(), 100); 387 assert_eq!(ticket.updated_at_unix_ms(), 100); 388 assert_eq!( 389 RebuildTransition::start(ticket_id, ticket.revision(), 101).ticket_id(), 390 ticket_id 391 ); 392 } 393 394 #[test] 395 fn durable_rebuild_matrix_rejects_every_inconsistent_shape() { 396 let invalidation = ProjectionInvalidation::new( 397 projection_id(), 398 generation(1), 399 generation(2), 400 InvalidationReason::ProjectionGenerationChanged, 401 100, 402 ) 403 .unwrap(); 404 let ticket_id = RebuildTicketId::new([3; 16]).unwrap(); 405 let replacement = checkpoint(generation(2), 1, 1, 105); 406 let durable = |revision, stage, checkpoint, requested, updated| { 407 RebuildTicket::from_durable_parts( 408 ticket_id, 409 invalidation.clone(), 410 revision, 411 stage, 412 SourceGeneration::new([9; 32]).unwrap(), 413 Some(position(9, 11)), 414 RawSourceDigest::new([8; 32]), 415 checkpoint, 416 None, 417 requested, 418 updated, 419 ) 420 }; 421 assert!( 422 durable( 423 ProjectionRevision::INITIAL, 424 RebuildStage::Requested, 425 None, 426 100, 427 100 428 ) 429 .is_ok() 430 ); 431 assert!( 432 durable( 433 ProjectionRevision::new(2).unwrap(), 434 RebuildStage::Running, 435 Some(replacement.clone()), 436 100, 437 105 438 ) 439 .is_ok() 440 ); 441 assert!( 442 durable( 443 ProjectionRevision::new(2).unwrap(), 444 RebuildStage::Completed, 445 Some(replacement.clone()), 446 100, 447 105 448 ) 449 .is_ok() 450 ); 451 for result in [ 452 durable( 453 ProjectionRevision::INITIAL, 454 RebuildStage::Requested, 455 None, 456 99, 457 100, 458 ), 459 durable( 460 ProjectionRevision::INITIAL, 461 RebuildStage::Requested, 462 None, 463 100, 464 99, 465 ), 466 durable( 467 ProjectionRevision::new(2).unwrap(), 468 RebuildStage::Requested, 469 None, 470 100, 471 100, 472 ), 473 durable( 474 ProjectionRevision::INITIAL, 475 RebuildStage::Requested, 476 None, 477 100, 478 101, 479 ), 480 durable( 481 ProjectionRevision::INITIAL, 482 RebuildStage::Running, 483 None, 484 100, 485 101, 486 ), 487 durable( 488 ProjectionRevision::INITIAL, 489 RebuildStage::Requested, 490 Some(replacement.clone()), 491 100, 492 100, 493 ), 494 durable( 495 ProjectionRevision::new(2).unwrap(), 496 RebuildStage::Completed, 497 None, 498 100, 499 105, 500 ), 501 durable( 502 ProjectionRevision::new(2).unwrap(), 503 RebuildStage::Running, 504 Some(checkpoint(generation(1), 1, 1, 105)), 505 100, 506 105, 507 ), 508 durable( 509 ProjectionRevision::new(2).unwrap(), 510 RebuildStage::Running, 511 Some(checkpoint(generation(2), 1, 1, 106)), 512 100, 513 105, 514 ), 515 ] { 516 assert_eq!(result, Err(Error::CorruptProjectionRecord)); 517 } 518 519 let requested = requested_ticket(ticket_id, invalidation.clone()); 520 assert_eq!( 521 requested.transition(RebuildTransition::start( 522 ticket_id, 523 ProjectionRevision::new(2).unwrap(), 524 101 525 )), 526 Err(Error::ProjectionRevisionConflict) 527 ); 528 assert_eq!( 529 requested.transition(RebuildTransition::start( 530 RebuildTicketId::new([4; 16]).unwrap(), 531 requested.revision(), 532 101, 533 )), 534 Err(Error::ProjectionRevisionConflict) 535 ); 536 assert_eq!( 537 requested.transition(RebuildTransition::start( 538 ticket_id, 539 requested.revision(), 540 99 541 )), 542 Err(Error::InvalidProjectionTimestamp) 543 ); 544 assert_eq!( 545 requested.transition(RebuildTransition::checkpoint( 546 ticket_id, 547 requested.revision(), 548 101, 549 replacement.clone() 550 )), 551 Err(Error::InvalidRebuildTransition) 552 ); 553 let running = requested 554 .transition(RebuildTransition::start( 555 ticket_id, 556 requested.revision(), 557 101, 558 )) 559 .unwrap(); 560 assert_eq!( 561 running.transition(RebuildTransition::checkpoint( 562 ticket_id, 563 running.revision(), 564 102, 565 checkpoint(generation(1), 1, 1, 102), 566 )), 567 Err(Error::ProjectionCheckpointMismatch) 568 ); 569 let progressed = running 570 .transition(RebuildTransition::checkpoint( 571 ticket_id, 572 running.revision(), 573 103, 574 checkpoint(generation(2), 2, 2, 103), 575 )) 576 .unwrap(); 577 assert_eq!( 578 progressed.transition(RebuildTransition::complete( 579 ticket_id, 580 progressed.revision(), 581 104, 582 checkpoint(generation(2), 1, 2, 104), 583 )), 584 Err(Error::ProjectionCheckpointRegression) 585 ); 586 let failed = running 587 .transition(RebuildTransition::fail( 588 ticket_id, 589 running.revision(), 590 102, 591 RebuildFailure::ReducerRejected, 592 )) 593 .unwrap(); 594 assert_eq!(failed.stage(), RebuildStage::Failed); 595 assert_eq!( 596 failed.transition(RebuildTransition::fail( 597 ticket_id, 598 failed.revision(), 599 103, 600 RebuildFailure::ReducerRejected, 601 )), 602 Err(Error::RebuildTicketTerminal) 603 ); 604 } 605 606 #[test] 607 fn event_index_models_cover_manifest_and_checkpoint_edges() { 608 let first_id = event_id('1'); 609 let last_id = event_id('2'); 610 assert_eq!( 611 EventIdRange::new(last_id, first_id), 612 Err(Error::InvalidEventIndexRange) 613 ); 614 let range = EventIdRange::new(first_id, last_id).unwrap(); 615 assert_eq!(range.first(), &first_id); 616 assert_eq!(range.last(), &last_id); 617 let digest = ArtifactDigest::new([5; 32]); 618 assert_eq!(digest.as_bytes(), &[5; 32]); 619 let shard_id = EventIndexShardId::parse("a").unwrap(); 620 assert_eq!(shard_id.as_str(), "a"); 621 for path in ["", "/absolute", "../escape", "a/../b", "a//b", "a\\b", " a"] { 622 assert_eq!( 623 EventIndexShard::new(shard_id.clone(), path, 1, range.clone(), 1, 2, digest), 624 Err(Error::InvalidEventIndexArtifactPath) 625 ); 626 } 627 assert_eq!( 628 EventIndexShard::new( 629 shard_id.clone(), 630 "x".repeat(radroots_storage::projection::EVENT_INDEX_ARTIFACT_PATH_MAX_BYTES + 1), 631 1, 632 range.clone(), 633 1, 634 2, 635 digest, 636 ), 637 Err(Error::InvalidEventIndexArtifactPath) 638 ); 639 assert_eq!( 640 EventIndexShard::new(shard_id.clone(), "a.json", 0, range.clone(), 1, 2, digest), 641 Err(Error::InvalidEventIndexShardCount) 642 ); 643 assert_eq!( 644 EventIndexShard::new(shard_id.clone(), "a.json", 1, range.clone(), 0, 2, digest), 645 Err(Error::InvalidEventIndexTimestamp) 646 ); 647 assert_eq!( 648 EventIndexShard::new(shard_id.clone(), "a.json", 1, range.clone(), 2, 1, digest), 649 Err(Error::InvalidEventIndexTimestamp) 650 ); 651 let shard = EventIndexShard::new(shard_id.clone(), "a.json", 1, range, 1, 2, digest).unwrap(); 652 assert_eq!(shard.shard_id(), &shard_id); 653 assert_eq!(shard.artifact_path(), "a.json"); 654 assert_eq!(shard.event_count(), 1); 655 assert_eq!(shard.first_published_at_unix_s(), 1); 656 assert_eq!(shard.last_published_at_unix_s(), 2); 657 assert_eq!(shard.sha256(), digest); 658 assert_eq!( 659 EventIndexManifest::new(generation(1), 1, 1, 1, 2, vec![]), 660 Err(Error::InvalidEventIndexShardCount) 661 ); 662 assert_eq!( 663 EventIndexManifest::new(generation(1), 1, 0, 1, 2, vec![shard.clone()]), 664 Err(Error::InvalidEventIndexManifest) 665 ); 666 assert_eq!( 667 EventIndexManifest::new(generation(1), 0, 1, 1, 2, vec![shard.clone()]), 668 Err(Error::InvalidEventIndexManifest) 669 ); 670 assert_eq!( 671 EventIndexManifest::new(generation(1), 1, 1, 0, 2, vec![shard.clone()]), 672 Err(Error::InvalidEventIndexManifest) 673 ); 674 assert_eq!( 675 EventIndexManifest::new(generation(1), 1, 1, 1, 3, vec![shard.clone()]), 676 Err(Error::InvalidEventIndexManifest) 677 ); 678 assert_eq!( 679 EventIndexManifest::new( 680 generation(1), 681 1, 682 1, 683 1, 684 2, 685 vec![shard.clone(); radroots_storage::projection::EVENT_INDEX_SHARDS_MAX + 1], 686 ), 687 Err(Error::InvalidEventIndexShardCount) 688 ); 689 let manifest = EventIndexManifest::new(generation(1), 1, 1, 1, 2, vec![shard.clone()]).unwrap(); 690 assert_eq!(manifest.generation(), generation(1)); 691 assert_eq!(manifest.target_shard_size(), 1); 692 assert_eq!(manifest.first_published_at_unix_s(), 1); 693 assert_eq!(manifest.last_published_at_unix_s(), 2); 694 695 for cursor in [ 696 Some(String::new()), 697 Some(" leading".to_owned()), 698 Some("bad\nvalue".to_owned()), 699 ] { 700 assert_eq!( 701 EventIndexShardCheckpoint::new(shard_id.clone(), 1, None, cursor), 702 Err(Error::InvalidEventIndexCursor) 703 ); 704 } 705 assert_eq!( 706 EventIndexShardCheckpoint::new(shard_id.clone(), 0, None, None), 707 Err(Error::InvalidEventIndexTimestamp) 708 ); 709 let shard_checkpoint = EventIndexShardCheckpoint::new( 710 shard_id.clone(), 711 2, 712 Some(last_id), 713 Some("cursor".to_owned()), 714 ) 715 .unwrap(); 716 assert_eq!(shard_checkpoint.shard_id(), &shard_id); 717 assert_eq!(shard_checkpoint.last_created_at_unix_s(), 2); 718 assert_eq!(shard_checkpoint.last_event_id(), Some(&last_id)); 719 assert_eq!(shard_checkpoint.cursor(), Some("cursor")); 720 assert_eq!( 721 EventIndexCheckpoint::new(generation(1), 0, vec![]), 722 Err(Error::InvalidEventIndexCheckpoint) 723 ); 724 assert_eq!( 725 EventIndexCheckpoint::new( 726 generation(1), 727 1, 728 vec![ 729 shard_checkpoint.clone(); 730 radroots_storage::projection::EVENT_INDEX_SHARDS_MAX + 1 731 ], 732 ), 733 Err(Error::InvalidEventIndexCheckpoint) 734 ); 735 let index = EventIndexCheckpoint::new(generation(1), 3, vec![shard_checkpoint]).unwrap(); 736 assert_eq!(index.generation(), generation(1)); 737 assert_eq!(index.generated_at_unix_ms(), 3); 738 assert_eq!(index.shards().len(), 1); 739 assert!(index.shard(&shard_id).is_some()); 740 assert!( 741 index 742 .shard(&EventIndexShardId::parse("missing").unwrap()) 743 .is_none() 744 ); 745 746 let status = ProjectionStatus::new( 747 projection_id(), 748 generation(1), 749 ProjectionHealth::Ready, 750 None, 751 None, 752 ) 753 .unwrap(); 754 assert_eq!(status.projection_id(), &projection_id()); 755 assert_eq!(status.generation(), generation(1)); 756 assert!(status.checkpoint().is_none()); 757 assert!(status.active_rebuild().is_none()); 758 assert_eq!( 759 ProjectionStatus::new( 760 projection_id(), 761 generation(1), 762 ProjectionHealth::Rebuilding, 763 None, 764 None 765 ), 766 Err(Error::CorruptProjectionRecord) 767 ); 768 assert_eq!( 769 ProjectionStatus::new( 770 projection_id(), 771 generation(1), 772 ProjectionHealth::Ready, 773 None, 774 Some(RebuildTicketId::new([1; 16]).unwrap()) 775 ), 776 Err(Error::CorruptProjectionRecord) 777 ); 778 assert_eq!( 779 ProjectionStatus::new( 780 projection_id(), 781 generation(1), 782 ProjectionHealth::Ready, 783 Some(checkpoint(generation(2), 1, 1, 2)), 784 None 785 ), 786 Err(Error::CorruptProjectionRecord) 787 ); 788 assert_eq!( 789 ProjectionStatus::new( 790 projection_id(), 791 generation(1), 792 ProjectionHealth::Ready, 793 Some( 794 ProjectionCheckpoint::new( 795 ProjectionId::parse("other").unwrap(), 796 generation(1), 797 None, 798 1, 799 2, 800 ) 801 .unwrap(), 802 ), 803 None, 804 ), 805 Err(Error::CorruptProjectionRecord) 806 ); 807 }