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