facts.rs (4637B)
1 //! Immutable delivery observations, independent of the scheduling revision. 2 3 use super::{ 4 AuthoredDeliveryPlan, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, Error, 5 SatisfactionState, WorkClaim, 6 }; 7 8 /// One final sink result from an exactly identified durable claim. 9 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 10 #[derive(Clone, Debug, Eq, PartialEq)] 11 pub struct AuthoredDeliveryFact { 12 claim: WorkClaim, 13 outcome: DeliveryAttemptOutcome, 14 observed_at_unix_ms: u64, 15 } 16 17 impl AuthoredDeliveryFact { 18 pub const fn claim(&self) -> &WorkClaim { 19 &self.claim 20 } 21 pub const fn outcome(&self) -> &DeliveryAttemptOutcome { 22 &self.outcome 23 } 24 pub const fn observed_at_unix_ms(&self) -> u64 { 25 self.observed_at_unix_ms 26 } 27 } 28 29 impl AuthoredDeliveryPlan { 30 pub fn delivery_facts(&self) -> &[AuthoredDeliveryFact] { 31 &self.delivery_facts 32 } 33 34 pub const fn stop_requested_at_unix_ms(&self) -> Option<u64> { 35 self.stop_requested_at_unix_ms 36 } 37 38 /// Cumulative transport evidence, independent of cancellation or retry state. 39 pub fn delivery_satisfaction(&self) -> Result<SatisfactionState, Error> { 40 self.evaluate_outcomes( 41 self.attempts 42 .iter() 43 .map(|attempt| &attempt.outcome) 44 .chain(self.delivery_facts.iter().map(|fact| &fact.outcome)), 45 ) 46 } 47 48 pub(super) fn validate_facts(&self) -> Result<(), Error> { 49 if self.delivery_facts.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize 50 || self.stop_requested_at_unix_ms.is_some_and(|at| { 51 at < self.created_at_unix_ms 52 || at > self.updated_at_unix_ms 53 || !self.state.is_terminal() 54 }) 55 { 56 return Err(Error::InvalidAuthoredDeliveryPlan); 57 } 58 let mut claims = std::collections::BTreeSet::new(); 59 for fact in &self.delivery_facts { 60 if fact.claim.validate().is_err() 61 || fact.claim.acquired_at_unix_ms() < self.created_at_unix_ms 62 || fact.claim.acquired_at_unix_ms() > self.updated_at_unix_ms 63 || fact.claim.row_revision() >= self.revision 64 || fact.observed_at_unix_ms < fact.claim.acquired_at_unix_ms() 65 || self 66 .request 67 .as_ref() 68 .is_none_or(|request| fact.outcome.validate_for(request).is_err()) 69 || !claims.insert(( 70 fact.claim.token(), 71 fact.claim.owner(), 72 fact.claim.generation(), 73 fact.claim.row_revision(), 74 fact.claim.acquired_at_unix_ms(), 75 fact.claim.expires_at_unix_ms(), 76 )) 77 { 78 return Err(Error::InvalidAuthoredDeliveryPlan); 79 } 80 } 81 let mut reconciled = Vec::new(); 82 for attempt in &self.attempts { 83 if let Some(claim) = attempt.claim_evidence() { 84 if reconciled.contains(&claim) 85 || !self 86 .delivery_facts 87 .iter() 88 .any(|fact| fact.claim() == claim && fact.outcome() == attempt.outcome()) 89 { 90 return Err(Error::InvalidAuthoredDeliveryPlan); 91 } 92 reconciled.push(claim); 93 } 94 } 95 Ok(()) 96 } 97 98 /// Called only after the atomic owner has checked its original claim receipt. 99 pub(crate) fn record_delivery_fact( 100 &mut self, 101 claim: WorkClaim, 102 outcome: DeliveryAttemptOutcome, 103 observed_at_unix_ms: u64, 104 ) -> Result<(), Error> { 105 if let Some(prior) = self.delivery_facts.iter().find(|fact| fact.claim == claim) { 106 return if prior.outcome == outcome { 107 Ok(()) 108 } else { 109 Err(Error::AtomicWorkflowMismatch) 110 }; 111 } 112 if self.delivery_facts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize { 113 return Err(Error::DeliveryAttemptOverflow); 114 } 115 self.delivery_facts.push(AuthoredDeliveryFact { 116 claim, 117 outcome, 118 observed_at_unix_ms, 119 }); 120 // Facts have their own observation time and immutable atomic receipt. They 121 // do not mutate scheduling revision/time, fences, attempts or backoff. 122 if let Err(error) = self.validate() { 123 self.delivery_facts.pop(); 124 return Err(error); 125 } 126 Ok(()) 127 } 128 }