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