lib

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

reconciliation.rs (7346B)


      1 //! Reconcile immutable facts only under current scheduling authority.
      2 
      3 use super::{
      4     AuthoredDeliveryAttempt, AuthoredDeliveryClaim, AuthoredDeliveryFact, AuthoredDeliveryPlan,
      5     AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, Error, FailureClass,
      6     NonZeroU32, RetrySchedule, Retryability, SatisfactionState, WorkFailure, WorkPhase,
      7 };
      8 use crate::authored_atomic::WorkFence;
      9 
     10 impl AuthoredDeliveryPlan {
     11     pub fn pending_delivery_facts(&self) -> impl Iterator<Item = &AuthoredDeliveryFact> {
     12         self.delivery_facts.iter().filter(|fact| {
     13             !self
     14                 .attempts
     15                 .iter()
     16                 .any(|attempt| attempt.claim_evidence() == Some(fact.claim()))
     17         })
     18     }
     19 
     20     pub(crate) fn reconcile_delivery_facts(
     21         &mut self,
     22         fence: Option<&WorkFence>,
     23         retry: Option<RetrySchedule>,
     24         at: u64,
     25         claims: &[AuthoredDeliveryClaim],
     26     ) -> Result<(), Error> {
     27         self.validate()?;
     28         if self.stop_requested_at_unix_ms.is_some() || self.state.is_terminal() {
     29             return Err(Error::InvalidAuthoredTransition);
     30         }
     31         match fence {
     32             Some(fence) => {
     33                 self.require_claim(fence.token(), fence.generation(), fence.row_revision(), at)?
     34             }
     35             None if self
     36                 .claim
     37                 .as_ref()
     38                 .is_some_and(|claim| at < claim.expires_at_unix_ms()) =>
     39             {
     40                 return Err(Error::DeliveryPlanClaimConflict);
     41             }
     42             None => {}
     43         }
     44         let pending: Vec<_> = self.pending_delivery_facts().cloned().collect();
     45         if pending.is_empty()
     46             || pending.iter().any(|fact| at < fact.observed_at_unix_ms())
     47             || at < self.updated_at_unix_ms
     48         {
     49             return Err(Error::AtomicWorkflowMismatch);
     50         }
     51         let mut candidate = self.clone();
     52         for fact in &pending {
     53             // A legacy fenced application may already have persisted this exact
     54             // result. Both its original ordinal and valid lease interval must
     55             // match: reconciling another fact can clear a live lease early.
     56             let original = claims
     57                 .iter()
     58                 .find(|entry| entry.claim() == fact.claim())
     59                 .ok_or(Error::AtomicWorkflowMismatch)?;
     60             if let Some(attempt) = candidate.attempts.iter_mut().find(|attempt| {
     61                 attempt.claim_evidence().is_none()
     62                     && original.prior_attempt_count().checked_add(1)
     63                         == Some(attempt.attempt().get())
     64                     && attempt.recorded_at_unix_ms() >= fact.claim().acquired_at_unix_ms()
     65                     && attempt.recorded_at_unix_ms() < fact.claim().expires_at_unix_ms()
     66             }) {
     67                 if attempt.outcome() != fact.outcome() {
     68                     return Err(Error::AtomicWorkflowMismatch);
     69                 }
     70                 attempt.claim = Some(fact.claim().clone());
     71                 continue;
     72             }
     73             if candidate.attempts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize {
     74                 return Err(Error::DeliveryAttemptOverflow);
     75             }
     76             let satisfaction = candidate.evaluate_with(fact.outcome().clone())?;
     77             let next = u32::try_from(candidate.attempts.len() + 1)
     78                 .ok()
     79                 .and_then(NonZeroU32::new)
     80                 .ok_or(Error::DeliveryAttemptOverflow)?;
     81             candidate.attempts.push(AuthoredDeliveryAttempt {
     82                 attempt: next,
     83                 recorded_at_unix_ms: at,
     84                 outcome: fact.outcome().clone(),
     85                 satisfaction,
     86                 claim: Some(fact.claim().clone()),
     87             });
     88         }
     89         candidate.attempt_count =
     90             u32::try_from(candidate.attempts.len()).map_err(|_| Error::DeliveryAttemptOverflow)?;
     91         let last = candidate
     92             .attempts
     93             .last()
     94             .ok_or(Error::AtomicWorkflowMismatch)?;
     95         let (state, failure) = reconciliation_state(
     96             last.outcome(),
     97             last.satisfaction(),
     98             candidate.attempt_count,
     99             retry.as_ref(),
    100         )?;
    101         candidate.state = state;
    102         candidate.last_failure = failure;
    103         candidate.retry = retry;
    104         candidate.claim = None;
    105         candidate.advance(at)?;
    106         candidate.validate()?;
    107         *self = candidate;
    108         Ok(())
    109     }
    110 }
    111 
    112 fn reconciliation_state(
    113     outcome: &DeliveryAttemptOutcome,
    114     satisfaction: SatisfactionState,
    115     attempt_count: u32,
    116     retry: Option<&RetrySchedule>,
    117 ) -> Result<(AuthoredDeliveryState, Option<WorkFailure>), Error> {
    118     let failure = match outcome {
    119         DeliveryAttemptOutcome::Receipt(_) => None,
    120         DeliveryAttemptOutcome::SinkFailure(failure) => Some(WorkFailure::new(
    121             failure.code(),
    122             WorkPhase::Delivery,
    123             if failure.retryability() == Retryability::Retryable {
    124                 FailureClass::Retryable
    125             } else {
    126                 FailureClass::Terminal
    127             },
    128             failure.retry_after_unix_ms(),
    129             failure.message().map(str::to_owned),
    130         )?),
    131     };
    132     match satisfaction {
    133         SatisfactionState::Satisfied if retry.is_none() => {
    134             return Ok((AuthoredDeliveryState::Satisfied, None));
    135         }
    136         SatisfactionState::Exhausted if retry.is_none() => {
    137             // Exhausted acceptance is settled even when the last adapter failure
    138             // was retryable; retain only a compatible terminal diagnostic.
    139             return Ok((
    140                 AuthoredDeliveryState::Exhausted,
    141                 failure.filter(|failure| failure.class() == FailureClass::Terminal),
    142             ));
    143         }
    144         SatisfactionState::Satisfied | SatisfactionState::Exhausted => {
    145             return Err(Error::InvalidRetrySchedule);
    146         }
    147         SatisfactionState::Pending => {}
    148     }
    149     if attempt_count == DELIVERY_PLAN_ATTEMPTS_MAX {
    150         if retry.is_some() {
    151             return Err(Error::InvalidRetrySchedule);
    152         }
    153         return Ok((
    154             AuthoredDeliveryState::Exhausted,
    155             Some(WorkFailure::new(
    156                 "delivery_attempt_limit",
    157                 WorkPhase::Delivery,
    158                 FailureClass::Terminal,
    159                 None,
    160                 None,
    161             )?),
    162         ));
    163     }
    164     if failure
    165         .as_ref()
    166         .is_some_and(|failure| failure.class() == FailureClass::Terminal)
    167     {
    168         if retry.is_some() {
    169             return Err(Error::InvalidRetrySchedule);
    170         }
    171         return Ok((AuthoredDeliveryState::FailedTerminal, failure));
    172     }
    173     let retry = retry.ok_or(Error::InvalidRetrySchedule)?;
    174     if retry.attempt().get() != attempt_count
    175         || retry.failure().phase() != WorkPhase::Delivery
    176         || failure.as_ref().is_some_and(|failure| {
    177             retry.failure().code() != failure.code()
    178                 || retry.failure().diagnostic() != failure.diagnostic()
    179                 || failure
    180                     .retry_after_unix_ms()
    181                     .is_some_and(|at| retry.not_before_unix_ms() < at)
    182         })
    183     {
    184         return Err(Error::InvalidRetrySchedule);
    185     }
    186     Ok((
    187         AuthoredDeliveryState::Retryable,
    188         Some(retry.failure().clone()),
    189     ))
    190 }