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 }