lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }