lib

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

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 }