memory_delivery_history.rs (3646B)
1 //! Instance-owned indexes into the reference backend's immutable receipts. 2 3 use super::State; 4 use crate::{ 5 Error, 6 authored_atomic::{AuthoredAtomicCommand, ClaimAuthoredTarget, ReconcileDeliveryFacts}, 7 authored_delivery::{ 8 AuthoredDeliveryHistory, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, 9 DELIVERY_PLAN_ATTEMPTS_MAX, 10 }, 11 }; 12 13 #[derive(Clone)] 14 pub(super) struct Entry { 15 preparation: Option<usize>, 16 claims: Vec<usize>, 17 } 18 19 pub(super) fn register( 20 state: &mut State, 21 command: &AuthoredAtomicCommand, 22 receipt_index: usize, 23 ) -> Result<(), Error> { 24 match command { 25 AuthoredAtomicCommand::Prepare(value) => { 26 for plan in value.delivery_plans() { 27 if state 28 .authored_delivery_history 29 .insert( 30 plan.plan_id(), 31 Entry { 32 preparation: Some(receipt_index), 33 claims: Vec::new(), 34 }, 35 ) 36 .is_some() 37 { 38 return Err(Error::AtomicCommitConflict); 39 } 40 } 41 } 42 AuthoredAtomicCommand::Claim(value) => { 43 if let ClaimAuthoredTarget::DeliveryPlan(plan_id) = value.target() { 44 let entry = state 45 .authored_delivery_history 46 .entry(*plan_id) 47 .or_insert(Entry { 48 preparation: None, 49 claims: Vec::new(), 50 }); 51 if entry.claims.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize { 52 return Err(Error::DeliveryAttemptOverflow); 53 } 54 entry.claims.push(receipt_index); 55 } 56 } 57 _ => {} 58 } 59 Ok(()) 60 } 61 62 pub(super) fn history( 63 state: &State, 64 plan_id: AuthoredDeliveryPlanId, 65 ) -> Result<Option<AuthoredDeliveryHistory>, Error> { 66 let Some(plan) = state 67 .authored_delivery_plans 68 .iter() 69 .find(|plan| plan.plan_id() == plan_id) 70 else { 71 return Ok(None); 72 }; 73 let entry = state.authored_delivery_history.get(&plan_id); 74 let preparation = entry 75 .and_then(|entry| entry.preparation) 76 .map(|index| { 77 state 78 .authored_atomic_receipts 79 .get(index) 80 .ok_or(Error::AtomicWorkflowMismatch) 81 }) 82 .transpose()?; 83 let mut history = AuthoredDeliveryHistory::new(plan.clone(), preparation)?; 84 if let Some(entry) = entry { 85 if entry.claims.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize { 86 history.mark_truncated(); 87 } 88 for index in entry 89 .claims 90 .iter() 91 .take(DELIVERY_PLAN_ATTEMPTS_MAX as usize) 92 { 93 history.push_claim( 94 state 95 .authored_atomic_receipts 96 .get(*index) 97 .ok_or(Error::AtomicWorkflowMismatch)?, 98 )?; 99 } 100 } 101 history.validate()?; 102 Ok(Some(history)) 103 } 104 105 pub(super) fn reconcile( 106 state: &mut State, 107 command: &ReconcileDeliveryFacts, 108 ) -> Result<AuthoredDeliveryPlan, Error> { 109 let history = history(state, command.plan_id())?.ok_or(Error::InvalidAuthoredDeliveryPlan)?; 110 let plan = command.apply_to(&history)?; 111 let stored = state 112 .authored_delivery_plans 113 .iter_mut() 114 .find(|value| value.plan_id() == command.plan_id()) 115 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 116 *stored = plan.clone(); 117 Ok(plan) 118 }