memory.rs (84091B)
1 //! Deterministic in-memory reference storage backend. 2 3 #[path = "memory_delivery_history.rs"] 4 mod delivery_history; 5 6 use radroots_event::EventId; 7 use radroots_protocol::runtime::v1::OperationId; 8 use radroots_transport::{BoxFuture, source::EventProvenance}; 9 use std::sync::{Mutex, MutexGuard}; 10 11 use crate::{ 12 Error, EventStore, Journal, Outbox, ProjectionStore, 13 atomic::{ 14 AtomicCommit, AtomicCommitDisposition, AtomicCommitId, AtomicCommitOutcome, 15 AtomicCommitReceipt, AtomicStorage, AtomicWorkflow, 16 }, 17 authored::{AdmissionState, FailureClass, WorkFailure, WorkPhase}, 18 authored_atomic::{ 19 AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, AuthoredAtomicStorage, 20 AuthoredWorkTarget, CancelAuthoredTarget, ClaimAuthoredTarget, 21 }, 22 authored_delivery::DeliveryAttemptOutcome, 23 authored_draft::{ 24 AUTHORED_DRAFT_QUERY_LIMIT_MAX, AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, 25 AuthoredDraftStore, DraftAppendDisposition, DraftAppendReceipt, 26 }, 27 backup::{ 28 BackupId, BackupOperation, BackupPlan, BackupTransition, ReliabilityRevision, 29 RestoreOperation, RestorePlan, RestoreTransition, StorageReliability, 30 }, 31 event::{ 32 AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage, 33 EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration, 34 StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent, 35 VisibilityEvaluation, VisibilityInput, VisibilitySnapshot, evaluate_visibility, 36 }, 37 journal::{ 38 IdempotencyKey, JournalStage, JournalTransition, OperationInstanceId, OperationRecord, 39 PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX, 40 }, 41 outbox::{ 42 ClaimOutboxItems, ClaimedOutboxItem, DeliveryAttemptEvidence, EnqueueDisposition, 43 EnqueueOutboxItem, EnqueueReceipt, LeaseId, OutboxItemId, OutboxLease, OutboxRecord, 44 OutboxRevision, OutboxStage, OutboxStatus, 45 }, 46 private_artifact::{ 47 DeletionReason, EXPIRED_ARTIFACT_QUERY_LIMIT_MAX, PrivateArtifactId, 48 PrivateArtifactMetadata, PrivateArtifactResealReceipt, PrivateArtifactResealRequest, 49 PrivateArtifactRevision, PrivateArtifactStage, PrivateArtifactStatus, PrivateArtifactStore, 50 }, 51 projection::{ 52 EventIndexCheckpoint, EventIndexManifest, ProjectionCheckpoint, ProjectionDocument, 53 ProjectionGeneration, ProjectionHealth, ProjectionId, ProjectionInvalidation, 54 ProjectionSnapshot, ProjectionStatus, RebuildStage, RebuildTicket, RebuildTicketId, 55 RebuildTransition, 56 }, 57 status::{ 58 EventStoreHealth, EventStoreMode, EventStoreStatus, IntegrityHealth, IntegrityStatus, 59 ShutdownState, StorageBackend, StorageOpenMode, StorageStatus, StorageStatusProvider, 60 WriterPolicy, 61 }, 62 }; 63 64 #[derive(Clone)] 65 struct EventEntry { 66 position: EventPosition, 67 admission: EventAdmission, 68 provenance: Vec<EventProvenance>, 69 } 70 71 #[derive(Clone, Default)] 72 struct State { 73 events: Vec<EventEntry>, 74 journal: Vec<OperationRecord>, 75 outbox: Vec<OutboxRecord>, 76 projections: Vec<ProjectionStatus>, 77 projection_invalidations: Vec<ProjectionInvalidation>, 78 rebuilds: Vec<RebuildTicket>, 79 event_index_manifests: Vec<EventIndexManifest>, 80 event_index_checkpoints: Vec<EventIndexCheckpoint>, 81 projection_documents: Vec<(ProjectionId, ProjectionGeneration, ProjectionDocument)>, 82 projection_snapshots: Vec<ProjectionSnapshot>, 83 private_artifacts: Vec<PrivateArtifactMetadata>, 84 private_artifact_reseals: Vec<PrivateArtifactResealReceipt>, 85 backups: Vec<BackupOperation>, 86 restores: Vec<RestoreOperation>, 87 atomic_receipts: Vec<AtomicCommitReceipt>, 88 authored_operations: Vec<crate::authored::AuthoredOperation>, 89 authored_artifacts: Vec<crate::authored::AuthoredArtifact>, 90 authored_delivery_plans: Vec<crate::authored_delivery::AuthoredDeliveryPlan>, 91 authored_atomic_receipts: Vec<AuthoredAtomicReceipt>, 92 authored_delivery_history: std::collections::BTreeMap< 93 crate::authored_delivery::AuthoredDeliveryPlanId, 94 delivery_history::Entry, 95 >, 96 authored_drafts: Vec<AuthoredDraft>, 97 closed: bool, 98 } 99 100 /// Bounded deterministic reference backend with no hidden tasks or globals. 101 pub struct MemoryStorage { 102 generation: SourceGeneration, 103 state: Mutex<State>, 104 } 105 106 impl MemoryStorage { 107 pub const fn new(generation: SourceGeneration) -> Self { 108 Self { 109 generation, 110 state: Mutex::new(State { 111 events: Vec::new(), 112 journal: Vec::new(), 113 outbox: Vec::new(), 114 projections: Vec::new(), 115 projection_invalidations: Vec::new(), 116 rebuilds: Vec::new(), 117 event_index_manifests: Vec::new(), 118 event_index_checkpoints: Vec::new(), 119 projection_documents: Vec::new(), 120 projection_snapshots: Vec::new(), 121 private_artifacts: Vec::new(), 122 private_artifact_reseals: Vec::new(), 123 backups: Vec::new(), 124 restores: Vec::new(), 125 atomic_receipts: Vec::new(), 126 authored_operations: Vec::new(), 127 authored_artifacts: Vec::new(), 128 authored_delivery_plans: Vec::new(), 129 authored_atomic_receipts: Vec::new(), 130 authored_delivery_history: std::collections::BTreeMap::new(), 131 authored_drafts: Vec::new(), 132 closed: false, 133 }), 134 } 135 } 136 137 pub const fn generation(&self) -> SourceGeneration { 138 self.generation 139 } 140 141 fn state(&self) -> Result<MutexGuard<'_, State>, Error> { 142 let state = self.state_any()?; 143 if state.closed { 144 return Err(Error::BackendUnavailable); 145 } 146 Ok(state) 147 } 148 149 fn state_any(&self) -> Result<MutexGuard<'_, State>, Error> { 150 self.state.lock().map_err(|_| Error::BackendUnavailable) 151 } 152 153 fn selected( 154 &self, 155 state: &State, 156 query: &EventQuery, 157 mut eligible: impl FnMut(&EventEntry) -> bool, 158 ) -> Result<(Vec<EventEntry>, Option<EventPosition>), Error> { 159 if query 160 .bounds() 161 .cursor() 162 .is_some_and(|cursor| cursor.generation() != self.generation) 163 { 164 return Err(Error::SourceGenerationChanged); 165 } 166 let after = query 167 .bounds() 168 .cursor() 169 .map_or(0, |cursor| cursor.sequence().get()); 170 let mut selected = state 171 .events 172 .iter() 173 .filter(|entry| { 174 entry.position.sequence().get() > after 175 && query.selects(entry.admission.event_id()) 176 && eligible(entry) 177 }) 178 .take(usize::from(query.bounds().limit()) + 1) 179 .cloned() 180 .collect::<Vec<_>>(); 181 let next = if selected.len() > usize::from(query.bounds().limit()) { 182 selected.truncate(usize::from(query.bounds().limit())); 183 selected.last().map(|entry| entry.position) 184 } else { 185 None 186 }; 187 Ok((selected, next)) 188 } 189 190 fn visibility_locked(&self, state: &State) -> Result<VisibilityEvaluation, Error> { 191 evaluate_visibility( 192 self.generation, 193 state.events.iter().map(|entry| { 194 VisibilityInput::new( 195 entry.position, 196 entry.admission.event(), 197 entry.admission.stage(), 198 ) 199 }), 200 ) 201 } 202 203 fn admit_locked( 204 &self, 205 state: &mut State, 206 admission: EventAdmission, 207 ) -> Result<AdmissionReceipt, Error> { 208 if let Some(entry) = state 209 .events 210 .iter_mut() 211 .find(|entry| entry.admission.event_id() == admission.event_id()) 212 { 213 if entry.admission.event() != admission.event() { 214 return Err(Error::EventConflict); 215 } 216 if admission.stage() < entry.admission.stage() { 217 return Err(Error::AdmissionRegression); 218 } 219 let disposition = if admission.stage() == entry.admission.stage() { 220 AdmissionDisposition::Duplicate 221 } else { 222 AdmissionDisposition::Advanced 223 }; 224 if !entry.provenance.contains(admission.provenance()) { 225 entry.provenance.push(admission.provenance().clone()); 226 } 227 entry.admission = admission; 228 return Ok(AdmissionReceipt::new( 229 *entry.admission.event_id(), 230 entry.position, 231 entry.admission.stage(), 232 disposition, 233 )); 234 } 235 let next = u64::try_from(state.events.len()) 236 .map_err(|_| Error::CorruptStoredEvent)? 237 .checked_add(1) 238 .ok_or(Error::CorruptStoredEvent)?; 239 let position = EventPosition::new(self.generation, EventSequence::new(next)?); 240 let receipt = AdmissionReceipt::new( 241 *admission.event_id(), 242 position, 243 admission.stage(), 244 AdmissionDisposition::Inserted, 245 ); 246 let provenance = vec![admission.provenance().clone()]; 247 state.events.push(EventEntry { 248 position, 249 admission, 250 provenance, 251 }); 252 Ok(receipt) 253 } 254 255 fn prepare_locked( 256 state: &mut State, 257 operation: PrepareOperation, 258 ) -> Result<PrepareReceipt, Error> { 259 if let Some(record) = state 260 .journal 261 .iter() 262 .find(|record| record.idempotency_key() == operation.idempotency_key()) 263 { 264 if record.operation_id() != operation.operation_id() 265 || record.input_digest() != operation.input_digest() 266 || record.instance_id() != operation.instance_id() 267 { 268 return Err(Error::IdempotencyConflict); 269 } 270 return Ok(PrepareReceipt::new( 271 PrepareDisposition::Replay, 272 record.clone(), 273 )); 274 } 275 if state 276 .journal 277 .iter() 278 .any(|record| record.instance_id() == operation.instance_id()) 279 { 280 return Err(Error::OperationIdentityMismatch); 281 } 282 let record = operation.into_record()?; 283 state.journal.push(record.clone()); 284 Ok(PrepareReceipt::new(PrepareDisposition::Created, record)) 285 } 286 287 fn transition_locked( 288 state: &mut State, 289 transition: JournalTransition, 290 ) -> Result<OperationRecord, Error> { 291 let record = state 292 .journal 293 .iter_mut() 294 .find(|record| record.instance_id() == transition.instance_id()) 295 .ok_or(Error::OperationNotFound)?; 296 let next = record.transition(&transition)?; 297 *record = next.clone(); 298 Ok(next) 299 } 300 301 fn enqueue_locked(state: &mut State, item: EnqueueOutboxItem) -> Result<EnqueueReceipt, Error> { 302 if let Some(record) = state 303 .outbox 304 .iter() 305 .find(|record| record.item_id() == item.item_id()) 306 { 307 let candidate = item.into_record(); 308 if record.operation_instance_id() != candidate.operation_instance_id() 309 || record.plan_digest() != candidate.plan_digest() 310 || record.request() != candidate.request() 311 || record.created_at_unix_ms() != candidate.created_at_unix_ms() 312 { 313 return Err(Error::OutboxPlanConflict); 314 } 315 return Ok(EnqueueReceipt::new( 316 EnqueueDisposition::Replay, 317 record.clone(), 318 )); 319 } 320 if state.outbox.iter().any(|record| { 321 record.operation_instance_id() == item.operation_instance_id() 322 && record.item_id() != item.item_id() 323 }) { 324 return Err(Error::OutboxPlanConflict); 325 } 326 let record = item.into_record(); 327 state.outbox.push(record.clone()); 328 Ok(EnqueueReceipt::new(EnqueueDisposition::Created, record)) 329 } 330 331 fn checkpoint_locked( 332 state: &mut State, 333 checkpoint: ProjectionCheckpoint, 334 ) -> Result<ProjectionStatus, Error> { 335 if let Some(status) = state 336 .projections 337 .iter_mut() 338 .find(|status| status.projection_id() == checkpoint.projection_id()) 339 { 340 if status.generation() != checkpoint.generation() { 341 return Err(Error::ProjectionCheckpointMismatch); 342 } 343 if status 344 .checkpoint() 345 .is_some_and(|prior| !checkpoint.advances(prior)) 346 { 347 return Err(Error::ProjectionCheckpointRegression); 348 } 349 let next = ProjectionStatus::new( 350 checkpoint.projection_id().clone(), 351 checkpoint.generation(), 352 ProjectionHealth::Ready, 353 Some(checkpoint), 354 None, 355 )?; 356 *status = next.clone(); 357 return Ok(next); 358 } 359 let status = ProjectionStatus::new( 360 checkpoint.projection_id().clone(), 361 checkpoint.generation(), 362 ProjectionHealth::Ready, 363 Some(checkpoint), 364 None, 365 )?; 366 state.projections.push(status.clone()); 367 Ok(status) 368 } 369 370 fn integrity_locked(state: &State) -> Result<IntegrityStatus, Error> { 371 let members = state 372 .events 373 .len() 374 .checked_add(state.journal.len()) 375 .and_then(|count| count.checked_add(state.outbox.len())) 376 .and_then(|count| count.checked_add(state.projections.len())) 377 .and_then(|count| count.checked_add(state.private_artifacts.len())) 378 .ok_or(Error::InvalidIntegrityStatus)?; 379 IntegrityStatus::new( 380 IntegrityHealth::Healthy, 381 None, 382 u32::try_from(members).map_err(|_| Error::InvalidIntegrityStatus)?, 383 0, 384 ) 385 } 386 387 fn status_locked(state: &State) -> Result<StorageStatus, Error> { 388 StorageStatus::new( 389 StorageBackend::Memory, 390 StorageOpenMode::Create, 391 WriterPolicy::NoWriter, 392 if state.closed { 393 ShutdownState::Closed 394 } else { 395 ShutdownState::Open 396 }, 397 Self::integrity_locked(state)?, 398 false, 399 0, 400 ) 401 } 402 } 403 404 impl Default for MemoryStorage { 405 fn default() -> Self { 406 Self::new(SourceGeneration::new([1; 32]).expect("fixed non-zero memory generation")) 407 } 408 } 409 410 impl EventStore for MemoryStorage { 411 fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> { 412 Box::pin(async move { 413 let state = self.state()?; 414 let raw = u64::try_from(state.events.len()).map_err(|_| Error::CorruptStoredEvent)?; 415 let verified = u64::try_from( 416 state 417 .events 418 .iter() 419 .filter(|entry| entry.admission.stage() >= AdmissionStage::Verified) 420 .count(), 421 ) 422 .map_err(|_| Error::CorruptStoredEvent)?; 423 let visible = u64::try_from( 424 self.visibility_locked(&state)? 425 .snapshot() 426 .visible_event_ids() 427 .len(), 428 ) 429 .map_err(|_| Error::CorruptStoredEvent)?; 430 EventStoreStatus::new( 431 self.generation, 432 EventStoreMode::ReadWrite, 433 EventStoreHealth::Available, 434 raw, 435 verified, 436 visible, 437 ) 438 }) 439 } 440 441 fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> { 442 Box::pin(async move { 443 let mut state = self.state()?; 444 self.admit_locked(&mut state, admission) 445 }) 446 } 447 448 fn query_raw( 449 &self, 450 query: EventQuery, 451 ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> { 452 Box::pin(async move { 453 let state = self.state()?; 454 let (entries, next) = self.selected(&state, &query, |_| true)?; 455 let items = entries 456 .into_iter() 457 .map(|entry| { 458 StoredRawEvent::new( 459 entry.position, 460 entry.admission.event().clone(), 461 entry.admission.stage(), 462 ) 463 }) 464 .collect(); 465 EventPage::new(self.generation, items, next, query.bounds()) 466 }) 467 } 468 469 fn query_verified( 470 &self, 471 query: EventQuery, 472 ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> { 473 Box::pin(async move { 474 let state = self.state()?; 475 let (entries, next) = self.selected(&state, &query, |entry| { 476 entry.admission.stage() >= AdmissionStage::Verified 477 })?; 478 let items = entries 479 .into_iter() 480 .map(|entry| { 481 StoredVerifiedEvent::new(entry.position, entry.admission.event().clone()) 482 }) 483 .collect(); 484 EventPage::new(self.generation, items, next, query.bounds()) 485 }) 486 } 487 488 fn query_visible( 489 &self, 490 query: EventQuery, 491 ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> { 492 Box::pin(async move { 493 let state = self.state()?; 494 let visibility = self.visibility_locked(&state)?; 495 let (entries, next) = self.selected(&state, &query, |entry| { 496 visibility.is_visible(entry.admission.event_id()) 497 })?; 498 let items = entries 499 .into_iter() 500 .map(|entry| { 501 StoredVisibleEvent::new(entry.position, entry.admission.event().clone()) 502 }) 503 .collect(); 504 EventPage::new(self.generation, items, next, query.bounds()) 505 }) 506 } 507 508 fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> { 509 Box::pin(async move { 510 let state = self.state()?; 511 Ok(self.visibility_locked(&state)?.into_snapshot()) 512 }) 513 } 514 515 fn query_provenance( 516 &self, 517 event_id: EventId, 518 bounds: EventQueryBounds, 519 ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> { 520 Box::pin(async move { 521 if bounds 522 .cursor() 523 .is_some_and(|cursor| cursor.generation() != self.generation) 524 { 525 return Err(Error::SourceGenerationChanged); 526 } 527 let state = self.state()?; 528 let entry = state 529 .events 530 .iter() 531 .find(|entry| entry.admission.event_id() == &event_id) 532 .ok_or(Error::EventNotFound)?; 533 let after = bounds.cursor().map_or(0, |cursor| cursor.sequence().get()); 534 let items = if entry.position.sequence().get() > after { 535 entry 536 .provenance 537 .iter() 538 .take(usize::from(bounds.limit())) 539 .cloned() 540 .map(|provenance| StoredEventProvenance::new(entry.position, provenance)) 541 .collect() 542 } else { 543 Vec::new() 544 }; 545 EventPage::new(self.generation, items, None, bounds) 546 }) 547 } 548 } 549 550 impl Journal for MemoryStorage { 551 fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> { 552 Box::pin(async move { 553 let mut state = self.state()?; 554 Self::prepare_locked(&mut state, operation) 555 }) 556 } 557 558 fn operation( 559 &self, 560 instance_id: OperationInstanceId, 561 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 562 Box::pin(async move { 563 Ok(self 564 .state()? 565 .journal 566 .iter() 567 .find(|record| record.instance_id() == instance_id) 568 .cloned()) 569 }) 570 } 571 572 fn by_idempotency_key( 573 &self, 574 operation_id: OperationId, 575 idempotency_key: IdempotencyKey, 576 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 577 Box::pin(async move { 578 Ok(self 579 .state()? 580 .journal 581 .iter() 582 .find(|record| { 583 record.operation_id() == operation_id 584 && record.idempotency_key() == &idempotency_key 585 }) 586 .cloned()) 587 }) 588 } 589 590 fn transition( 591 &self, 592 transition: JournalTransition, 593 ) -> BoxFuture<'_, Result<OperationRecord, Error>> { 594 Box::pin(async move { 595 let mut state = self.state()?; 596 Self::transition_locked(&mut state, transition) 597 }) 598 } 599 600 fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> { 601 Box::pin(async move { 602 if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX { 603 return Err(Error::InvalidJournalQueryLimit); 604 } 605 Ok(self 606 .state()? 607 .journal 608 .iter() 609 .filter(|record| record.state().stage() == JournalStage::Recoverable) 610 .take(usize::from(limit)) 611 .cloned() 612 .collect()) 613 }) 614 } 615 } 616 617 impl Outbox for MemoryStorage { 618 fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>> { 619 Box::pin(async move { 620 let mut state = self.state()?; 621 Self::enqueue_locked(&mut state, item) 622 }) 623 } 624 625 fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>> { 626 Box::pin(async move { 627 Ok(self 628 .state()? 629 .outbox 630 .iter() 631 .find(|record| record.item_id() == item_id) 632 .cloned()) 633 }) 634 } 635 636 fn claim( 637 &self, 638 request: ClaimOutboxItems, 639 ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>> { 640 Box::pin(async move { 641 let mut state = self.state()?; 642 let mut claimed = Vec::new(); 643 for record in &mut state.outbox { 644 if claimed.len() >= usize::from(request.limit()) || record.stage().is_terminal() { 645 continue; 646 } 647 if record 648 .retry_not_before_unix_ms() 649 .is_some_and(|at| request.now_unix_ms() < at) 650 || record 651 .lease() 652 .is_some_and(|lease| lease.is_active_at(request.now_unix_ms())) 653 { 654 continue; 655 } 656 let lease = OutboxLease::new( 657 request.lease_id_for(record.item_id()), 658 request.owner().clone(), 659 request.now_unix_ms(), 660 request.lease_expires_at_unix_ms(), 661 )?; 662 record.claim(lease.clone())?; 663 claimed.push(ClaimedOutboxItem::new(record.clone(), lease)); 664 } 665 Ok(claimed) 666 }) 667 } 668 669 fn record_attempt( 670 &self, 671 evidence: DeliveryAttemptEvidence, 672 ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { 673 Box::pin(async move { 674 let mut state = self.state()?; 675 let record = state 676 .outbox 677 .iter_mut() 678 .find(|record| record.item_id() == evidence.item_id()) 679 .ok_or(Error::OutboxItemNotFound)?; 680 record.record_attempt(evidence)?; 681 Ok(record.clone()) 682 }) 683 } 684 685 fn release( 686 &self, 687 item_id: OutboxItemId, 688 lease_id: LeaseId, 689 expected_revision: OutboxRevision, 690 released_at_unix_ms: u64, 691 retry_not_before_unix_ms: Option<u64>, 692 ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { 693 Box::pin(async move { 694 let mut state = self.state()?; 695 let record = state 696 .outbox 697 .iter_mut() 698 .find(|record| record.item_id() == item_id) 699 .ok_or(Error::OutboxItemNotFound)?; 700 record.release( 701 lease_id, 702 expected_revision, 703 released_at_unix_ms, 704 retry_not_before_unix_ms, 705 )?; 706 Ok(record.clone()) 707 }) 708 } 709 710 fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>> { 711 Box::pin(async move { 712 let state = self.state()?; 713 let mut status = OutboxStatus { 714 pending: 0, 715 leased: 0, 716 retryable: 0, 717 satisfied: 0, 718 exhausted: 0, 719 }; 720 for record in &state.outbox { 721 let count = match record.stage() { 722 OutboxStage::Pending => &mut status.pending, 723 OutboxStage::Leased => &mut status.leased, 724 OutboxStage::Retryable => &mut status.retryable, 725 OutboxStage::Satisfied => &mut status.satisfied, 726 OutboxStage::Exhausted => &mut status.exhausted, 727 }; 728 *count = count.checked_add(1).ok_or(Error::CorruptOutboxRecord)?; 729 } 730 Ok(status) 731 }) 732 } 733 } 734 735 impl ProjectionStore for MemoryStorage { 736 fn status( 737 &self, 738 projection_id: ProjectionId, 739 ) -> BoxFuture<'_, Result<Option<ProjectionStatus>, Error>> { 740 Box::pin(async move { 741 Ok(self 742 .state()? 743 .projections 744 .iter() 745 .find(|status| status.projection_id() == &projection_id) 746 .cloned()) 747 }) 748 } 749 750 fn checkpoint( 751 &self, 752 checkpoint: ProjectionCheckpoint, 753 ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> { 754 Box::pin(async move { 755 let mut state = self.state()?; 756 Self::checkpoint_locked(&mut state, checkpoint) 757 }) 758 } 759 760 fn invalidate( 761 &self, 762 invalidation: ProjectionInvalidation, 763 ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> { 764 Box::pin(async move { 765 let mut state = self.state()?; 766 let status_index = state 767 .projections 768 .iter() 769 .position(|status| status.projection_id() == invalidation.projection_id()) 770 .ok_or(Error::ProjectionCheckpointMismatch)?; 771 if state.projections[status_index].generation() != invalidation.invalid_generation() { 772 return Err(Error::ProjectionCheckpointMismatch); 773 } 774 if let Some(existing) = state.projection_invalidations.iter().find(|existing| { 775 existing.projection_id() == invalidation.projection_id() 776 && existing.invalid_generation() == invalidation.invalid_generation() 777 }) && existing != &invalidation 778 { 779 return Err(Error::ProjectionRevisionConflict); 780 } 781 let next = ProjectionStatus::new( 782 invalidation.projection_id().clone(), 783 state.projections[status_index].generation(), 784 ProjectionHealth::Invalidated, 785 state.projections[status_index].checkpoint().cloned(), 786 None, 787 )?; 788 if !state 789 .projection_invalidations 790 .iter() 791 .any(|existing| existing == &invalidation) 792 { 793 state.projection_invalidations.push(invalidation); 794 } 795 state.projections[status_index] = next.clone(); 796 Ok(next) 797 }) 798 } 799 800 fn invalidation( 801 &self, 802 projection_id: ProjectionId, 803 replacement_generation: ProjectionGeneration, 804 ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>> { 805 Box::pin(async move { 806 Ok(self 807 .state()? 808 .projection_invalidations 809 .iter() 810 .rev() 811 .find(|invalidation| { 812 invalidation.projection_id() == &projection_id 813 && invalidation.replacement_generation() == replacement_generation 814 }) 815 .cloned()) 816 }) 817 } 818 819 fn request_rebuild( 820 &self, 821 ticket: RebuildTicket, 822 ) -> BoxFuture<'_, Result<RebuildTicket, Error>> { 823 Box::pin(async move { 824 let mut state = self.state()?; 825 if let Some(existing) = state 826 .rebuilds 827 .iter() 828 .find(|existing| existing.ticket_id() == ticket.ticket_id()) 829 { 830 return if existing == &ticket { 831 Ok(existing.clone()) 832 } else { 833 Err(Error::ProjectionRevisionConflict) 834 }; 835 } 836 let projection_id = ticket.invalidation().projection_id(); 837 let has_invalidation = state 838 .projection_invalidations 839 .iter() 840 .any(|invalidation| invalidation == ticket.invalidation()); 841 let status = state 842 .projections 843 .iter_mut() 844 .find(|status| status.projection_id() == projection_id) 845 .ok_or(Error::ProjectionCheckpointMismatch)?; 846 if status.generation() != ticket.invalidation().invalid_generation() 847 || status.health() != ProjectionHealth::Invalidated 848 || !has_invalidation 849 { 850 return Err(Error::ProjectionCheckpointMismatch); 851 } 852 *status = ProjectionStatus::new( 853 projection_id.clone(), 854 status.generation(), 855 ProjectionHealth::Rebuilding, 856 status.checkpoint().cloned(), 857 Some(ticket.ticket_id()), 858 )?; 859 state.rebuilds.push(ticket.clone()); 860 Ok(ticket) 861 }) 862 } 863 864 fn rebuild( 865 &self, 866 ticket_id: RebuildTicketId, 867 ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>> { 868 Box::pin(async move { 869 Ok(self 870 .state()? 871 .rebuilds 872 .iter() 873 .find(|ticket| ticket.ticket_id() == ticket_id) 874 .cloned()) 875 }) 876 } 877 878 fn transition_rebuild( 879 &self, 880 transition: RebuildTransition, 881 ) -> BoxFuture<'_, Result<RebuildTicket, Error>> { 882 Box::pin(async move { 883 let mut state = self.state()?; 884 let index = state 885 .rebuilds 886 .iter() 887 .position(|ticket| ticket.ticket_id() == transition.ticket_id()) 888 .ok_or(Error::ProjectionRevisionConflict)?; 889 let next = state.rebuilds[index].transition(transition)?; 890 if next.stage() == RebuildStage::Completed 891 && (next.source_generation() != self.generation 892 || next 893 .source_high_water() 894 .map_or(0, |position| position.sequence().get()) 895 != u64::try_from(state.events.len()) 896 .map_err(|_| Error::CorruptProjectionRecord)?) 897 { 898 return Err(Error::SourceGenerationChanged); 899 } 900 let projection_id = next.invalidation().projection_id(); 901 let status = state 902 .projections 903 .iter_mut() 904 .find(|status| status.projection_id() == projection_id) 905 .ok_or(Error::CorruptProjectionRecord)?; 906 if status.generation() != next.invalidation().invalid_generation() 907 || status.health() != ProjectionHealth::Rebuilding 908 || status.active_rebuild() != Some(next.ticket_id()) 909 { 910 return Err(Error::CorruptProjectionRecord); 911 } 912 let (generation, health, checkpoint, active_rebuild) = match next.stage() { 913 RebuildStage::Requested | RebuildStage::Running => ( 914 status.generation(), 915 ProjectionHealth::Rebuilding, 916 status.checkpoint().cloned(), 917 Some(next.ticket_id()), 918 ), 919 RebuildStage::Completed => ( 920 next.invalidation().replacement_generation(), 921 ProjectionHealth::Ready, 922 next.checkpoint().cloned(), 923 None, 924 ), 925 RebuildStage::Failed => ( 926 status.generation(), 927 ProjectionHealth::Ready, 928 status.checkpoint().cloned(), 929 None, 930 ), 931 }; 932 *status = ProjectionStatus::new( 933 projection_id.clone(), 934 generation, 935 health, 936 checkpoint, 937 active_rebuild, 938 )?; 939 state.rebuilds[index] = next.clone(); 940 Ok(next) 941 }) 942 } 943 944 fn event_index_manifest( 945 &self, 946 generation: ProjectionGeneration, 947 ) -> BoxFuture<'_, Result<Option<EventIndexManifest>, Error>> { 948 Box::pin(async move { 949 Ok(self 950 .state()? 951 .event_index_manifests 952 .iter() 953 .find(|manifest| manifest.generation() == generation) 954 .cloned()) 955 }) 956 } 957 958 fn put_event_index_manifest( 959 &self, 960 manifest: EventIndexManifest, 961 ) -> BoxFuture<'_, Result<(), Error>> { 962 Box::pin(async move { 963 let mut state = self.state()?; 964 if let Some(existing) = state 965 .event_index_manifests 966 .iter() 967 .find(|existing| existing.generation() == manifest.generation()) 968 { 969 return if existing == &manifest { 970 Ok(()) 971 } else { 972 Err(Error::CorruptProjectionRecord) 973 }; 974 } 975 state.event_index_manifests.push(manifest); 976 Ok(()) 977 }) 978 } 979 980 fn event_index_checkpoint( 981 &self, 982 generation: ProjectionGeneration, 983 ) -> BoxFuture<'_, Result<Option<EventIndexCheckpoint>, Error>> { 984 Box::pin(async move { 985 Ok(self 986 .state()? 987 .event_index_checkpoints 988 .iter() 989 .find(|checkpoint| checkpoint.generation() == generation) 990 .cloned()) 991 }) 992 } 993 994 fn put_event_index_checkpoint( 995 &self, 996 checkpoint: EventIndexCheckpoint, 997 ) -> BoxFuture<'_, Result<(), Error>> { 998 Box::pin(async move { 999 let mut state = self.state()?; 1000 if let Some(existing) = state 1001 .event_index_checkpoints 1002 .iter_mut() 1003 .find(|existing| existing.generation() == checkpoint.generation()) 1004 { 1005 if checkpoint.generated_at_unix_ms() < existing.generated_at_unix_ms() { 1006 return Err(Error::InvalidEventIndexCheckpoint); 1007 } 1008 *existing = checkpoint; 1009 } else { 1010 state.event_index_checkpoints.push(checkpoint); 1011 } 1012 Ok(()) 1013 }) 1014 } 1015 1016 fn put_projection_document( 1017 &self, 1018 projection_id: ProjectionId, 1019 generation: ProjectionGeneration, 1020 document: ProjectionDocument, 1021 ) -> BoxFuture<'_, Result<(), Error>> { 1022 Box::pin(async move { 1023 let mut state = self.state()?; 1024 if let Some((_, _, existing)) = state.projection_documents.iter_mut().find( 1025 |(existing_id, existing_generation, existing)| { 1026 existing_id == &projection_id 1027 && *existing_generation == generation 1028 && existing.key() == document.key() 1029 }, 1030 ) { 1031 *existing = document; 1032 } else { 1033 state 1034 .projection_documents 1035 .push((projection_id, generation, document)); 1036 } 1037 Ok(()) 1038 }) 1039 } 1040 1041 fn query_projection_documents( 1042 &self, 1043 query: crate::projection::document_query::ProjectionDocumentQuery, 1044 ) -> BoxFuture<'_, Result<crate::projection::document_query::ProjectionDocumentPage, Error>> 1045 { 1046 use crate::projection::document_query::{ 1047 PROJECTION_DOCUMENT_PAGE_BYTES_MAX, ProjectionDocumentPage, ProjectionDocumentRecord, 1048 }; 1049 Box::pin(async move { 1050 let state = self.state()?; 1051 let mut selected = std::collections::BTreeMap::new(); 1052 let maximum = usize::from(query.limit()) + 1; 1053 for (id, generation, document) in &state.projection_documents { 1054 if id == query.projection_id() && query.matches(*generation, document.key()) { 1055 selected.insert((*generation, document.key()), document); 1056 if selected.len() > maximum { 1057 selected.pop_last(); 1058 } 1059 } 1060 } 1061 let mut has_more = selected.len() > usize::from(query.limit()); 1062 let mut records = Vec::new(); 1063 let mut bytes = 0; 1064 for ((generation, _), document) in selected.into_iter().take(usize::from(query.limit())) 1065 { 1066 if bytes + document.value().len() > PROJECTION_DOCUMENT_PAGE_BYTES_MAX { 1067 has_more = true; 1068 break; 1069 } 1070 bytes += document.value().len(); 1071 records.push(ProjectionDocumentRecord::new(generation, document.clone())?); 1072 } 1073 ProjectionDocumentPage::new(&query, records, has_more) 1074 }) 1075 } 1076 1077 fn projection_document( 1078 &self, 1079 projection_id: ProjectionId, 1080 generation: ProjectionGeneration, 1081 key: String, 1082 ) -> BoxFuture<'_, Result<Option<ProjectionDocument>, Error>> { 1083 Box::pin(async move { 1084 Ok(self 1085 .state()? 1086 .projection_documents 1087 .iter() 1088 .find(|(existing_id, existing_generation, existing)| { 1089 existing_id == &projection_id 1090 && *existing_generation == generation 1091 && existing.key() == key 1092 }) 1093 .map(|(_, _, document)| document.clone())) 1094 }) 1095 } 1096 1097 fn put_projection_snapshot( 1098 &self, 1099 snapshot: ProjectionSnapshot, 1100 ) -> BoxFuture<'_, Result<(), Error>> { 1101 Box::pin(async move { 1102 let mut state = self.state()?; 1103 if let Some(existing) = state.projection_snapshots.iter().find(|existing| { 1104 existing.projection_id() == snapshot.projection_id() 1105 && existing.snapshot_id() == snapshot.snapshot_id() 1106 }) { 1107 return if existing == &snapshot { 1108 Ok(()) 1109 } else { 1110 Err(Error::CorruptProjectionDocument) 1111 }; 1112 } 1113 state.projection_snapshots.push(snapshot); 1114 Ok(()) 1115 }) 1116 } 1117 1118 fn projection_snapshot( 1119 &self, 1120 projection_id: ProjectionId, 1121 snapshot_id: [u8; 32], 1122 ) -> BoxFuture<'_, Result<Option<ProjectionSnapshot>, Error>> { 1123 Box::pin(async move { 1124 Ok(self 1125 .state()? 1126 .projection_snapshots 1127 .iter() 1128 .find(|snapshot| { 1129 snapshot.projection_id() == &projection_id 1130 && snapshot.snapshot_id() == &snapshot_id 1131 }) 1132 .cloned()) 1133 }) 1134 } 1135 } 1136 1137 impl PrivateArtifactStore for MemoryStorage { 1138 fn put_metadata( 1139 &self, 1140 metadata: PrivateArtifactMetadata, 1141 ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> { 1142 Box::pin(async move { 1143 let mut state = self.state()?; 1144 if let Some(existing) = state 1145 .private_artifacts 1146 .iter() 1147 .find(|existing| existing.artifact_id() == metadata.artifact_id()) 1148 { 1149 return if existing == &metadata { 1150 Ok(existing.clone()) 1151 } else { 1152 Err(Error::PrivateArtifactConflict) 1153 }; 1154 } 1155 state.private_artifacts.push(metadata.clone()); 1156 Ok(metadata) 1157 }) 1158 } 1159 1160 fn metadata( 1161 &self, 1162 artifact_id: PrivateArtifactId, 1163 ) -> BoxFuture<'_, Result<Option<PrivateArtifactMetadata>, Error>> { 1164 Box::pin(async move { 1165 Ok(self 1166 .state()? 1167 .private_artifacts 1168 .iter() 1169 .find(|metadata| metadata.artifact_id() == artifact_id) 1170 .cloned()) 1171 }) 1172 } 1173 1174 fn reseal_metadata( 1175 &self, 1176 request: PrivateArtifactResealRequest, 1177 ) -> BoxFuture<'_, Result<PrivateArtifactResealReceipt, Error>> { 1178 Box::pin(async move { 1179 let mut state = self.state()?; 1180 if let Some(receipt) = state 1181 .private_artifact_reseals 1182 .iter() 1183 .find(|receipt| receipt.reseal_id() == request.reseal_id()) 1184 { 1185 return receipt.replay(&request); 1186 } 1187 let metadata = state 1188 .private_artifacts 1189 .iter_mut() 1190 .find(|metadata| metadata.artifact_id() == request.artifact_id()) 1191 .ok_or(Error::PrivateArtifactNotFound)?; 1192 let next = metadata.resealed(&request)?; 1193 let receipt = PrivateArtifactResealReceipt::committed(&request, next.revision()); 1194 *metadata = next; 1195 state.private_artifact_reseals.push(receipt); 1196 Ok(receipt) 1197 }) 1198 } 1199 1200 fn mark_expired( 1201 &self, 1202 artifact_id: PrivateArtifactId, 1203 expected_revision: PrivateArtifactRevision, 1204 at_unix_ms: u64, 1205 ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> { 1206 Box::pin(async move { 1207 let mut state = self.state()?; 1208 let metadata = state 1209 .private_artifacts 1210 .iter_mut() 1211 .find(|metadata| metadata.artifact_id() == artifact_id) 1212 .ok_or(Error::PrivateArtifactNotFound)?; 1213 let next = metadata.mark_expired(expected_revision, at_unix_ms)?; 1214 *metadata = next.clone(); 1215 Ok(next) 1216 }) 1217 } 1218 1219 fn tombstone( 1220 &self, 1221 artifact_id: PrivateArtifactId, 1222 expected_revision: PrivateArtifactRevision, 1223 at_unix_ms: u64, 1224 reason: DeletionReason, 1225 ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> { 1226 Box::pin(async move { 1227 let mut state = self.state()?; 1228 let metadata = state 1229 .private_artifacts 1230 .iter_mut() 1231 .find(|metadata| metadata.artifact_id() == artifact_id) 1232 .ok_or(Error::PrivateArtifactNotFound)?; 1233 let next = metadata.tombstone(expected_revision, at_unix_ms, reason)?; 1234 *metadata = next.clone(); 1235 Ok(next) 1236 }) 1237 } 1238 1239 fn expired( 1240 &self, 1241 at_unix_ms: u64, 1242 limit: u16, 1243 ) -> BoxFuture<'_, Result<Vec<PrivateArtifactMetadata>, Error>> { 1244 Box::pin(async move { 1245 if at_unix_ms == 0 || limit == 0 || limit > EXPIRED_ARTIFACT_QUERY_LIMIT_MAX { 1246 return Err(Error::InvalidExpiredArtifactQueryLimit); 1247 } 1248 Ok(self 1249 .state()? 1250 .private_artifacts 1251 .iter() 1252 .filter(|metadata| { 1253 metadata.stage() == PrivateArtifactStage::Active 1254 && metadata.retention().is_expired_at(at_unix_ms) 1255 }) 1256 .take(usize::from(limit)) 1257 .cloned() 1258 .collect()) 1259 }) 1260 } 1261 1262 fn status(&self) -> BoxFuture<'_, Result<PrivateArtifactStatus, Error>> { 1263 Box::pin(async move { 1264 let state = self.state()?; 1265 let mut status = PrivateArtifactStatus { 1266 active: 0, 1267 expired: 0, 1268 tombstoned: 0, 1269 }; 1270 for metadata in &state.private_artifacts { 1271 let count = match metadata.stage() { 1272 PrivateArtifactStage::Active => &mut status.active, 1273 PrivateArtifactStage::Expired => &mut status.expired, 1274 PrivateArtifactStage::Tombstoned => &mut status.tombstoned, 1275 }; 1276 *count = count 1277 .checked_add(1) 1278 .ok_or(Error::CorruptPrivateArtifactMetadata)?; 1279 } 1280 Ok(status) 1281 }) 1282 } 1283 } 1284 1285 impl StorageReliability for MemoryStorage { 1286 fn begin_backup(&self, plan: BackupPlan) -> BoxFuture<'_, Result<BackupOperation, Error>> { 1287 Box::pin(async move { 1288 let mut state = self.state()?; 1289 if let Some(existing) = state 1290 .backups 1291 .iter() 1292 .find(|operation| operation.plan().backup_id() == plan.backup_id()) 1293 { 1294 return if existing.plan() == &plan { 1295 Ok(existing.clone()) 1296 } else { 1297 Err(Error::ReliabilityRevisionConflict) 1298 }; 1299 } 1300 let operation = BackupOperation::planned(plan); 1301 state.backups.push(operation.clone()); 1302 Ok(operation) 1303 }) 1304 } 1305 1306 fn transition_backup( 1307 &self, 1308 backup_id: BackupId, 1309 expected_revision: ReliabilityRevision, 1310 transition: BackupTransition, 1311 at_unix_ms: u64, 1312 ) -> BoxFuture<'_, Result<BackupOperation, Error>> { 1313 Box::pin(async move { 1314 let mut state = self.state()?; 1315 let operation = state 1316 .backups 1317 .iter_mut() 1318 .find(|operation| operation.plan().backup_id() == backup_id) 1319 .ok_or(Error::CorruptReliabilityOperation)?; 1320 let next = operation.transition(expected_revision, transition, at_unix_ms)?; 1321 *operation = next.clone(); 1322 Ok(next) 1323 }) 1324 } 1325 1326 fn begin_restore(&self, plan: RestorePlan) -> BoxFuture<'_, Result<RestoreOperation, Error>> { 1327 Box::pin(async move { 1328 let mut state = self.state()?; 1329 let backup_id = plan.manifest().backup_id(); 1330 if let Some(existing) = state 1331 .restores 1332 .iter() 1333 .find(|operation| operation.plan().manifest().backup_id() == backup_id) 1334 { 1335 return if existing.plan() == &plan { 1336 Ok(existing.clone()) 1337 } else { 1338 Err(Error::ReliabilityRevisionConflict) 1339 }; 1340 } 1341 let operation = RestoreOperation::staging(plan); 1342 state.restores.push(operation.clone()); 1343 Ok(operation) 1344 }) 1345 } 1346 1347 fn transition_restore( 1348 &self, 1349 backup_id: BackupId, 1350 expected_revision: ReliabilityRevision, 1351 transition: RestoreTransition, 1352 at_unix_ms: u64, 1353 ) -> BoxFuture<'_, Result<RestoreOperation, Error>> { 1354 Box::pin(async move { 1355 let mut state = self.state()?; 1356 let operation = state 1357 .restores 1358 .iter_mut() 1359 .find(|operation| operation.plan().manifest().backup_id() == backup_id) 1360 .ok_or(Error::CorruptReliabilityOperation)?; 1361 let next = operation.transition(expected_revision, transition, at_unix_ms)?; 1362 *operation = next.clone(); 1363 Ok(next) 1364 }) 1365 } 1366 1367 fn integrity(&self) -> BoxFuture<'_, Result<IntegrityStatus, Error>> { 1368 Box::pin(async move { 1369 let state = self.state_any()?; 1370 Self::integrity_locked(&state) 1371 }) 1372 } 1373 1374 fn status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { 1375 Box::pin(async move { 1376 let state = self.state_any()?; 1377 Self::status_locked(&state) 1378 }) 1379 } 1380 1381 fn close(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { 1382 Box::pin(async move { 1383 let mut state = self.state_any()?; 1384 state.closed = true; 1385 Self::status_locked(&state) 1386 }) 1387 } 1388 } 1389 1390 impl StorageStatusProvider for MemoryStorage { 1391 fn storage_status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { 1392 Box::pin(async move { 1393 let state = self.state_any()?; 1394 Self::status_locked(&state) 1395 }) 1396 } 1397 } 1398 1399 impl AtomicStorage for MemoryStorage { 1400 fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> { 1401 Box::pin(async move { 1402 let mut state = self.state()?; 1403 if let Some(existing) = state 1404 .atomic_receipts 1405 .iter() 1406 .find(|receipt| receipt.commit_id() == request.commit_id()) 1407 { 1408 if existing.digest() != request.digest() 1409 || existing.outcome().kind() != request.workflow().kind() 1410 { 1411 return Err(Error::AtomicCommitConflict); 1412 } 1413 return AtomicCommitReceipt::new( 1414 &request, 1415 AtomicCommitDisposition::Replay, 1416 existing.committed_at_unix_ms(), 1417 existing.outcome().clone(), 1418 ); 1419 } 1420 let mut candidate = state.clone(); 1421 let outcome = match request.workflow().clone() { 1422 AtomicWorkflow::Ingested(ingested) => { 1423 let admission = 1424 self.admit_locked(&mut candidate, ingested.admission().clone())?; 1425 let projection = ingested 1426 .projection() 1427 .cloned() 1428 .map(|checkpoint| Self::checkpoint_locked(&mut candidate, checkpoint)) 1429 .transpose()? 1430 .map(Box::new); 1431 AtomicCommitOutcome::Ingested { 1432 admission, 1433 projection, 1434 } 1435 } 1436 }; 1437 let receipt = AtomicCommitReceipt::new( 1438 &request, 1439 AtomicCommitDisposition::Committed, 1440 request.requested_at_unix_ms(), 1441 outcome, 1442 )?; 1443 candidate.atomic_receipts.push(receipt.clone()); 1444 *state = candidate; 1445 Ok(receipt) 1446 }) 1447 } 1448 1449 fn receipt( 1450 &self, 1451 commit_id: AtomicCommitId, 1452 ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>> { 1453 Box::pin(async move { 1454 Ok(self 1455 .state()? 1456 .atomic_receipts 1457 .iter() 1458 .find(|receipt| receipt.commit_id() == commit_id) 1459 .cloned()) 1460 }) 1461 } 1462 } 1463 1464 fn prepare_authored_memory( 1465 candidate: &mut State, 1466 value: crate::authored_atomic::PrepareAuthoredOperation, 1467 ) -> Result<AuthoredAtomicOutcome, Error> { 1468 if candidate 1469 .authored_operations 1470 .iter() 1471 .any(|operation| operation.operation_id() == value.operation().operation_id()) 1472 || value.artifacts().iter().any(|artifact| { 1473 candidate 1474 .authored_artifacts 1475 .iter() 1476 .any(|existing| existing.artifact_id() == artifact.artifact_id()) 1477 }) 1478 || value.delivery_plans().iter().any(|plan| { 1479 candidate 1480 .authored_delivery_plans 1481 .iter() 1482 .any(|existing| existing.plan_id() == plan.plan_id()) 1483 }) 1484 { 1485 return Err(Error::AtomicCommitConflict); 1486 } 1487 candidate 1488 .authored_operations 1489 .push(value.operation().clone()); 1490 candidate 1491 .authored_artifacts 1492 .extend(value.artifacts().iter().cloned()); 1493 candidate 1494 .authored_delivery_plans 1495 .extend(value.delivery_plans().iter().cloned()); 1496 Ok(AuthoredAtomicOutcome::Prepared { 1497 operation: value.operation().clone(), 1498 artifacts: value.artifacts().to_vec(), 1499 delivery_plans: value.delivery_plans().to_vec(), 1500 }) 1501 } 1502 1503 impl AuthoredAtomicStorage for MemoryStorage { 1504 fn authored_delivery_history( 1505 &self, 1506 plan_id: crate::authored_delivery::AuthoredDeliveryPlanId, 1507 ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryHistory>, Error>> 1508 { 1509 Box::pin(async move { 1510 let state = self.state()?; 1511 delivery_history::history(&state, plan_id) 1512 }) 1513 } 1514 fn execute_authored( 1515 &self, 1516 command: AuthoredAtomicCommand, 1517 ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>> { 1518 Box::pin(async move { 1519 let mut state = self.state()?; 1520 if let Some(existing) = state 1521 .authored_atomic_receipts 1522 .iter() 1523 .find(|receipt| receipt.commit_id() == command.commit_id()) 1524 { 1525 if !existing.matches_command(&command) { 1526 return Err(Error::AtomicCommitConflict); 1527 } 1528 return AuthoredAtomicReceipt::from_durable_parts( 1529 existing.commit_id(), 1530 existing.digest(), 1531 AtomicCommitDisposition::Replay, 1532 existing.committed_at_unix_ms(), 1533 existing.outcome().clone(), 1534 ); 1535 } 1536 1537 let mut candidate = state.clone(); 1538 let outcome = match command.clone() { 1539 AuthoredAtomicCommand::Prepare(value) => { 1540 prepare_authored_memory(&mut candidate, value)? 1541 } 1542 AuthoredAtomicCommand::PrepareFromDraft(value) => { 1543 value.validate()?; 1544 let head = candidate 1545 .authored_drafts 1546 .iter() 1547 .filter(|draft| draft.draft_id() == value.source().draft_id()) 1548 .max_by_key(|draft| draft.revision()); 1549 if !head.is_some_and(|head| value.source().matches(head)) { 1550 return Err(Error::DraftRevisionConflict); 1551 } 1552 if candidate 1553 .authored_drafts 1554 .iter() 1555 .any(|draft| draft.draft_id() == value.intent().draft_id()) 1556 { 1557 return Err(Error::DraftRevisionConflict); 1558 } 1559 let ordinary = AuthoredAtomicCommand::Prepare(value.preparation().clone()); 1560 if candidate 1561 .authored_atomic_receipts 1562 .iter() 1563 .any(|receipt| receipt.commit_id() == ordinary.commit_id()) 1564 { 1565 return Err(Error::AtomicCommitConflict); 1566 } 1567 let prepared = 1568 prepare_authored_memory(&mut candidate, value.preparation().clone())?; 1569 candidate.authored_drafts.push(value.intent().clone()); 1570 let ordinary_index = candidate.authored_atomic_receipts.len(); 1571 delivery_history::register(&mut candidate, &ordinary, ordinary_index)?; 1572 candidate 1573 .authored_atomic_receipts 1574 .push(AuthoredAtomicReceipt::new( 1575 &ordinary, 1576 AtomicCommitDisposition::Committed, 1577 ordinary.requested_at_unix_ms(), 1578 prepared, 1579 )?); 1580 AuthoredAtomicOutcome::Submitted(value) 1581 } 1582 AuthoredAtomicCommand::Claim(value) => match value.target() { 1583 ClaimAuthoredTarget::ArtifactSigning(artifact_id) => { 1584 let artifact = candidate 1585 .authored_artifacts 1586 .iter_mut() 1587 .find(|artifact| artifact.artifact_id() == *artifact_id) 1588 .ok_or(Error::InvalidAuthoredArtifact)?; 1589 artifact.set_signing_claim( 1590 value.claim().clone(), 1591 value.claim().acquired_at_unix_ms(), 1592 )?; 1593 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1594 } 1595 ClaimAuthoredTarget::ArtifactAdmission(artifact_id) => { 1596 let artifact = candidate 1597 .authored_artifacts 1598 .iter_mut() 1599 .find(|artifact| artifact.artifact_id() == *artifact_id) 1600 .ok_or(Error::InvalidAuthoredArtifact)?; 1601 artifact.set_admission_claim( 1602 value.claim().clone(), 1603 value.claim().acquired_at_unix_ms(), 1604 )?; 1605 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1606 } 1607 ClaimAuthoredTarget::DeliveryPlan(plan_id) => { 1608 let artifact_id = candidate 1609 .authored_delivery_plans 1610 .iter() 1611 .find(|plan| plan.plan_id() == *plan_id) 1612 .ok_or(Error::InvalidAuthoredDeliveryPlan)? 1613 .artifact_id(); 1614 let artifact = candidate 1615 .authored_artifacts 1616 .iter() 1617 .find(|artifact| artifact.artifact_id() == artifact_id) 1618 .ok_or(Error::InvalidAuthoredArtifact)?; 1619 if artifact.signing_state() != crate::authored::SigningState::Signed { 1620 return Err(Error::InvalidAuthoredTransition); 1621 } 1622 let plan = candidate 1623 .authored_delivery_plans 1624 .iter_mut() 1625 .find(|plan| plan.plan_id() == *plan_id) 1626 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 1627 plan.claim(value.claim().clone(), value.claim().acquired_at_unix_ms())?; 1628 AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) 1629 } 1630 }, 1631 AuthoredAtomicCommand::ApplySigned(value) => { 1632 let artifact = candidate 1633 .authored_artifacts 1634 .iter_mut() 1635 .find(|artifact| artifact.artifact_id() == value.artifact_id()) 1636 .ok_or(Error::InvalidAuthoredArtifact)?; 1637 require_artifact_claim( 1638 artifact.signing_claim(), 1639 value.fence(), 1640 value.applied_at_unix_ms(), 1641 )?; 1642 artifact.record_signed(value.event().clone(), value.applied_at_unix_ms())?; 1643 let artifact = artifact.clone(); 1644 for plan in candidate 1645 .authored_delivery_plans 1646 .iter_mut() 1647 .filter(|plan| plan.artifact_id() == value.artifact_id()) 1648 { 1649 plan.bind_signed_event(value.event().clone(), value.applied_at_unix_ms())?; 1650 } 1651 AuthoredAtomicOutcome::Artifact(artifact) 1652 } 1653 AuthoredAtomicCommand::ApplyAdmission(value) => { 1654 let artifact = candidate 1655 .authored_artifacts 1656 .iter_mut() 1657 .find(|artifact| artifact.artifact_id() == value.artifact_id()) 1658 .ok_or(Error::InvalidAuthoredArtifact)?; 1659 require_artifact_claim( 1660 artifact.admission_claim(), 1661 value.fence(), 1662 value.applied_at_unix_ms(), 1663 )?; 1664 artifact.record_admission( 1665 value.state(), 1666 value.failure().cloned(), 1667 value.retry().cloned(), 1668 value.applied_at_unix_ms(), 1669 )?; 1670 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1671 } 1672 AuthoredAtomicCommand::RecordSigned(value) => { 1673 let claim_id = value.claim_command().commit_id(); 1674 let original = candidate 1675 .authored_atomic_receipts 1676 .iter() 1677 .find(|receipt| receipt.commit_id() == claim_id) 1678 .ok_or(Error::AtomicWorkflowMismatch)?; 1679 let artifact = candidate 1680 .authored_artifacts 1681 .iter_mut() 1682 .find(|artifact| artifact.artifact_id() == value.artifact_id()) 1683 .ok_or(Error::InvalidAuthoredArtifact)?; 1684 let already_signed = artifact.signed().is_some(); 1685 value.apply_to(artifact, original)?; 1686 let artifact = artifact.clone(); 1687 if !already_signed 1688 && artifact.signing_state() == crate::authored::SigningState::Signed 1689 { 1690 for plan in candidate.authored_delivery_plans.iter_mut().filter(|plan| { 1691 plan.artifact_id() == value.artifact_id() && !plan.state().is_terminal() 1692 }) { 1693 plan.bind_signed_event( 1694 value.event().clone(), 1695 value.observed_at_unix_ms().max(plan.updated_at_unix_ms()), 1696 )?; 1697 } 1698 } 1699 AuthoredAtomicOutcome::Artifact(artifact) 1700 } 1701 AuthoredAtomicCommand::RecordDelivery(value) => { 1702 let claim_id = value.claim_command().commit_id(); 1703 let original = candidate 1704 .authored_atomic_receipts 1705 .iter() 1706 .find(|receipt| receipt.commit_id() == claim_id) 1707 .ok_or(Error::AtomicWorkflowMismatch)?; 1708 let plan = candidate 1709 .authored_delivery_plans 1710 .iter_mut() 1711 .find(|plan| plan.plan_id() == value.plan_id()) 1712 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 1713 value.apply_to(plan, original)?; 1714 AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) 1715 } 1716 AuthoredAtomicCommand::ReconcileDelivery(value) => { 1717 let plan = delivery_history::reconcile(&mut candidate, &value)?; 1718 AuthoredAtomicOutcome::DeliveryPlan(plan) 1719 } 1720 AuthoredAtomicCommand::ApplyDelivery(value) => { 1721 let plan = candidate 1722 .authored_delivery_plans 1723 .iter_mut() 1724 .find(|plan| plan.plan_id() == value.plan_id()) 1725 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 1726 match value.outcome().clone() { 1727 DeliveryAttemptOutcome::Receipt(receipt) => plan.apply_receipt( 1728 value.fence().token(), 1729 value.fence().generation(), 1730 value.fence().row_revision(), 1731 receipt, 1732 value.retry().cloned(), 1733 value.applied_at_unix_ms(), 1734 )?, 1735 DeliveryAttemptOutcome::SinkFailure(failure) => plan.apply_sink_failure( 1736 value.fence().token(), 1737 value.fence().generation(), 1738 value.fence().row_revision(), 1739 failure, 1740 value.retry().cloned(), 1741 value.applied_at_unix_ms(), 1742 )?, 1743 } 1744 AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) 1745 } 1746 AuthoredAtomicCommand::ApplyFailure(value) => match value.target() { 1747 AuthoredWorkTarget::Artifact(artifact_id) => { 1748 let artifact = candidate 1749 .authored_artifacts 1750 .iter_mut() 1751 .find(|artifact| artifact.artifact_id() == *artifact_id) 1752 .ok_or(Error::InvalidAuthoredArtifact)?; 1753 match value.failure().phase() { 1754 WorkPhase::Signing => { 1755 require_artifact_claim( 1756 artifact.signing_claim(), 1757 value.fence(), 1758 value.applied_at_unix_ms(), 1759 )?; 1760 artifact.record_signing_failure( 1761 value.failure().clone(), 1762 value.retry().cloned(), 1763 value.applied_at_unix_ms(), 1764 )?; 1765 } 1766 WorkPhase::Admission => { 1767 require_artifact_claim( 1768 artifact.admission_claim(), 1769 value.fence(), 1770 value.applied_at_unix_ms(), 1771 )?; 1772 let state = match value.failure().class() { 1773 FailureClass::Retryable => AdmissionState::Retryable, 1774 FailureClass::Terminal => AdmissionState::Rejected, 1775 FailureClass::Indeterminate => { 1776 return Err(Error::InvalidAuthoredTransition); 1777 } 1778 }; 1779 artifact.record_admission( 1780 state, 1781 Some(value.failure().clone()), 1782 value.retry().cloned(), 1783 value.applied_at_unix_ms(), 1784 )?; 1785 } 1786 WorkPhase::Delivery => { 1787 return Err(Error::AtomicWorkflowMismatch); 1788 } 1789 } 1790 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1791 } 1792 AuthoredWorkTarget::DeliveryPlan(plan_id) => { 1793 let plan = candidate 1794 .authored_delivery_plans 1795 .iter_mut() 1796 .find(|plan| plan.plan_id() == *plan_id) 1797 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 1798 if value.failure().phase() != WorkPhase::Delivery 1799 || value.failure().class() == FailureClass::Indeterminate 1800 { 1801 return Err(Error::AtomicWorkflowMismatch); 1802 } 1803 let retryability = match value.failure().class() { 1804 FailureClass::Retryable => { 1805 radroots_transport::outcome::Retryability::Retryable 1806 } 1807 FailureClass::Terminal => { 1808 radroots_transport::outcome::Retryability::Terminal 1809 } 1810 FailureClass::Indeterminate => unreachable!(), 1811 }; 1812 let failure = radroots_transport::SinkFailure::for_request( 1813 plan.request().ok_or(Error::InvalidAuthoredDeliveryPlan)?, 1814 value.failure().code(), 1815 retryability, 1816 value.failure().retry_after_unix_ms(), 1817 value.failure().diagnostic().map(str::to_owned), 1818 Vec::new(), 1819 ) 1820 .map_err(|_| Error::AtomicWorkflowMismatch)?; 1821 plan.apply_sink_failure( 1822 value.fence().token(), 1823 value.fence().generation(), 1824 value.fence().row_revision(), 1825 failure, 1826 value.retry().cloned(), 1827 value.applied_at_unix_ms(), 1828 )?; 1829 AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) 1830 } 1831 }, 1832 AuthoredAtomicCommand::Cancel(value) => match value.target() { 1833 CancelAuthoredTarget::ArtifactSigning(artifact_id) => { 1834 let artifact = candidate 1835 .authored_artifacts 1836 .iter_mut() 1837 .find(|artifact| artifact.artifact_id() == *artifact_id) 1838 .ok_or(Error::InvalidAuthoredArtifact)?; 1839 if artifact.revision() != value.expected_revision() { 1840 return Err(Error::InvalidAuthoredTransition); 1841 } 1842 artifact.cancel_signing(value.cancelled_at_unix_ms())?; 1843 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1844 } 1845 CancelAuthoredTarget::ArtifactAdmission(artifact_id) => { 1846 let artifact = candidate 1847 .authored_artifacts 1848 .iter_mut() 1849 .find(|artifact| artifact.artifact_id() == *artifact_id) 1850 .ok_or(Error::InvalidAuthoredArtifact)?; 1851 if artifact.revision() != value.expected_revision() { 1852 return Err(Error::InvalidAuthoredTransition); 1853 } 1854 let failure = WorkFailure::new( 1855 "cancelled", 1856 WorkPhase::Admission, 1857 FailureClass::Terminal, 1858 None, 1859 None, 1860 )?; 1861 artifact.record_admission( 1862 AdmissionState::Cancelled, 1863 Some(failure), 1864 None, 1865 value.cancelled_at_unix_ms(), 1866 )?; 1867 AuthoredAtomicOutcome::Artifact(artifact.clone()) 1868 } 1869 CancelAuthoredTarget::DeliveryPlan(plan_id) => { 1870 let plan = candidate 1871 .authored_delivery_plans 1872 .iter_mut() 1873 .find(|plan| plan.plan_id() == *plan_id) 1874 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 1875 if plan.revision() != value.expected_revision() { 1876 return Err(Error::InvalidAuthoredDeliveryPlan); 1877 } 1878 plan.request_stop(value.cancelled_at_unix_ms())?; 1879 AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) 1880 } 1881 }, 1882 }; 1883 let committed_at = match (&command, &outcome) { 1884 ( 1885 AuthoredAtomicCommand::RecordDelivery(_), 1886 AuthoredAtomicOutcome::DeliveryPlan(plan), 1887 ) => command 1888 .requested_at_unix_ms() 1889 .max(plan.updated_at_unix_ms()), 1890 ( 1891 AuthoredAtomicCommand::RecordSigned(_), 1892 AuthoredAtomicOutcome::Artifact(artifact), 1893 ) => command 1894 .requested_at_unix_ms() 1895 .max(artifact.updated_at_unix_ms()), 1896 _ => command.requested_at_unix_ms(), 1897 }; 1898 let receipt = AuthoredAtomicReceipt::new( 1899 &command, 1900 AtomicCommitDisposition::Committed, 1901 committed_at, 1902 outcome, 1903 )?; 1904 let receipt_index = candidate.authored_atomic_receipts.len(); 1905 delivery_history::register(&mut candidate, &command, receipt_index)?; 1906 candidate.authored_atomic_receipts.push(receipt.clone()); 1907 *state = candidate; 1908 Ok(receipt) 1909 }) 1910 } 1911 1912 fn authored_receipt( 1913 &self, 1914 commit_id: AtomicCommitId, 1915 ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>> { 1916 Box::pin(async move { 1917 Ok(self 1918 .state()? 1919 .authored_atomic_receipts 1920 .iter() 1921 .find(|receipt| receipt.commit_id() == commit_id) 1922 .cloned()) 1923 }) 1924 } 1925 1926 fn authored_operation( 1927 &self, 1928 operation_id: OperationInstanceId, 1929 ) -> BoxFuture<'_, Result<Option<crate::authored::AuthoredOperation>, Error>> { 1930 Box::pin(async move { 1931 Ok(self 1932 .state()? 1933 .authored_operations 1934 .iter() 1935 .find(|operation| operation.operation_id() == operation_id) 1936 .cloned()) 1937 }) 1938 } 1939 1940 fn authored_artifact( 1941 &self, 1942 artifact_id: crate::authored::AuthoredArtifactId, 1943 ) -> BoxFuture<'_, Result<Option<crate::authored::AuthoredArtifact>, Error>> { 1944 Box::pin(async move { 1945 Ok(self 1946 .state()? 1947 .authored_artifacts 1948 .iter() 1949 .find(|artifact| artifact.artifact_id() == artifact_id) 1950 .cloned()) 1951 }) 1952 } 1953 1954 fn authored_delivery_plan( 1955 &self, 1956 plan_id: crate::authored_delivery::AuthoredDeliveryPlanId, 1957 ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryPlan>, Error>> { 1958 Box::pin(async move { 1959 Ok(self 1960 .state()? 1961 .authored_delivery_plans 1962 .iter() 1963 .find(|plan| plan.plan_id() == plan_id) 1964 .cloned()) 1965 }) 1966 } 1967 } 1968 1969 impl AuthoredDraftStore for MemoryStorage { 1970 fn append_authored_draft_pair( 1971 &self, 1972 pair: crate::authored_draft_pair::AuthoredDraftPair, 1973 ) -> BoxFuture<'_, Result<[DraftAppendReceipt; 2], Error>> { 1974 Box::pin(async move { 1975 let mut state = self.state()?; 1976 let [first, second] = pair.drafts(); 1977 let [first_expected, second_expected] = *pair.expected_heads(); 1978 let a = draft_append_disposition(&state.authored_drafts, first, first_expected)?; 1979 let b = draft_append_disposition(&state.authored_drafts, second, second_expected)?; 1980 if a != b { 1981 return Err(Error::DraftRevisionConflict); 1982 } 1983 if a == DraftAppendDisposition::Inserted { 1984 state 1985 .authored_drafts 1986 .extend([first.clone(), second.clone()]); 1987 } 1988 Ok([ 1989 DraftAppendReceipt::new(first.clone(), a), 1990 DraftAppendReceipt::new(second.clone(), b), 1991 ]) 1992 }) 1993 } 1994 1995 fn query_authored_drafts( 1996 &self, 1997 query: crate::authored_draft_query::AuthoredDraftQuery, 1998 ) -> BoxFuture<'_, Result<crate::authored_draft_query::AuthoredDraftPage, Error>> { 1999 Box::pin(async move { 2000 use crate::authored_draft_query::{ 2001 AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AuthoredDraftPage, AuthoredDraftQueryRecord, 2002 }; 2003 let state = self.state()?; 2004 let mut heads: std::collections::BTreeMap<AuthoredDraftId, &AuthoredDraft> = 2005 std::collections::BTreeMap::new(); 2006 let capacity = usize::from(query.limit()) + 1; 2007 for draft in &state.authored_drafts { 2008 if !query.matches(draft) 2009 || query 2010 .after() 2011 .is_some_and(|after| *draft.draft_id().as_bytes() <= after) 2012 { 2013 continue; 2014 } 2015 if let Some(head) = heads.get_mut(&draft.draft_id()) { 2016 if draft.revision() > head.revision() { 2017 *head = draft; 2018 } 2019 continue; 2020 } 2021 // Retain only the smallest requested IDs and one lookahead. 2022 // Draft identity metadata is immutable across revisions. 2023 if heads.len() == capacity { 2024 if heads 2025 .last_key_value() 2026 .is_some_and(|(last, _)| draft.draft_id() >= *last) 2027 { 2028 continue; 2029 } 2030 heads.pop_last(); 2031 } 2032 heads.insert(draft.draft_id(), draft); 2033 } 2034 let mut records = Vec::new(); 2035 let mut bytes = 0usize; 2036 let mut has_more = false; 2037 for draft in heads.values() { 2038 if records.len() == usize::from(query.limit()) 2039 || bytes + draft.payload().len() > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES 2040 { 2041 has_more = true; 2042 break; 2043 } 2044 bytes += draft.payload().len(); 2045 records.push(AuthoredDraftQueryRecord::Draft((*draft).clone())); 2046 } 2047 AuthoredDraftPage::new(&query, records, has_more) 2048 }) 2049 } 2050 2051 fn append_authored_draft( 2052 &self, 2053 draft: AuthoredDraft, 2054 expected_head: Option<AuthoredDraftRevision>, 2055 ) -> BoxFuture<'_, Result<DraftAppendReceipt, Error>> { 2056 Box::pin(async move { 2057 draft.validate()?; 2058 let mut state = self.state()?; 2059 let disposition = 2060 draft_append_disposition(&state.authored_drafts, &draft, expected_head)?; 2061 if disposition == DraftAppendDisposition::Replay { 2062 return Ok(DraftAppendReceipt::new(draft, disposition)); 2063 } 2064 state.authored_drafts.push(draft.clone()); 2065 Ok(DraftAppendReceipt::new( 2066 draft, 2067 DraftAppendDisposition::Inserted, 2068 )) 2069 }) 2070 } 2071 2072 fn authored_draft_head( 2073 &self, 2074 draft_id: AuthoredDraftId, 2075 ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { 2076 Box::pin(async move { 2077 Ok(self 2078 .state()? 2079 .authored_drafts 2080 .iter() 2081 .filter(|draft| draft.draft_id() == draft_id) 2082 .max_by_key(|draft| draft.revision()) 2083 .cloned()) 2084 }) 2085 } 2086 2087 fn authored_draft_revision( 2088 &self, 2089 draft_id: AuthoredDraftId, 2090 revision: AuthoredDraftRevision, 2091 ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { 2092 Box::pin(async move { 2093 Ok(self 2094 .state()? 2095 .authored_drafts 2096 .iter() 2097 .find(|draft| draft.draft_id() == draft_id && draft.revision() == revision) 2098 .cloned()) 2099 }) 2100 } 2101 2102 fn authored_draft_heads( 2103 &self, 2104 author: [u8; 32], 2105 limit: u16, 2106 ) -> BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>> { 2107 Box::pin(async move { 2108 if author.iter().all(|byte| *byte == 0) 2109 || limit == 0 2110 || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX 2111 { 2112 return Err(Error::InvalidAuthoredDraft); 2113 } 2114 let state = self.state()?; 2115 let mut heads = state 2116 .authored_drafts 2117 .iter() 2118 .filter(|draft| draft.author() == &author) 2119 .fold(Vec::<AuthoredDraft>::new(), |mut heads, draft| { 2120 match heads 2121 .iter_mut() 2122 .find(|head| head.draft_id() == draft.draft_id()) 2123 { 2124 Some(head) if draft.revision() > head.revision() => *head = draft.clone(), 2125 None => heads.push(draft.clone()), 2126 Some(_) => {} 2127 } 2128 heads 2129 }); 2130 heads.sort_by(|left, right| { 2131 right 2132 .updated_at_unix_ms() 2133 .cmp(&left.updated_at_unix_ms()) 2134 .then_with(|| left.draft_id().cmp(&right.draft_id())) 2135 }); 2136 heads.truncate(usize::from(limit)); 2137 Ok(heads) 2138 }) 2139 } 2140 } 2141 2142 fn require_artifact_claim( 2143 claim: Option<&crate::authored::WorkClaim>, 2144 fence: &crate::authored_atomic::WorkFence, 2145 now_unix_ms: u64, 2146 ) -> Result<(), Error> { 2147 if !claim.is_some_and(|claim| { 2148 claim.matches_fence( 2149 fence.token(), 2150 fence.generation(), 2151 fence.row_revision(), 2152 now_unix_ms, 2153 ) 2154 }) { 2155 return Err(Error::DeliveryPlanClaimConflict); 2156 } 2157 Ok(()) 2158 } 2159 2160 fn draft_append_disposition( 2161 drafts: &[AuthoredDraft], 2162 draft: &AuthoredDraft, 2163 expected_head: Option<AuthoredDraftRevision>, 2164 ) -> Result<DraftAppendDisposition, Error> { 2165 if let Some(existing) = drafts.iter().find(|existing| { 2166 existing.draft_id() == draft.draft_id() && existing.revision() == draft.revision() 2167 }) { 2168 return if existing == draft { 2169 Ok(DraftAppendDisposition::Replay) 2170 } else { 2171 Err(Error::DraftRevisionConflict) 2172 }; 2173 } 2174 let head = drafts 2175 .iter() 2176 .filter(|existing| existing.draft_id() == draft.draft_id()) 2177 .max_by_key(|existing| existing.revision()); 2178 match (head, expected_head) { 2179 (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} 2180 (Some(previous), Some(expected)) if previous.revision() == expected => { 2181 draft.validate_successor_of(previous)?; 2182 } 2183 _ => return Err(Error::DraftRevisionConflict), 2184 } 2185 Ok(DraftAppendDisposition::Inserted) 2186 }