lib

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

delivery.rs (14772B)


      1 //! Delivery observations outlive the lease that admitted their external effect.
      2 
      3 use super::*;
      4 use radroots_storage::authored_atomic::{ReconcileDeliveryFacts, RecordDeliveryFact};
      5 
      6 impl Engine {
      7     /// Reconciles retained evidence or executes one bounded delivery attempt.
      8     ///
      9     /// A late result retains its original raw binding without authorizing a
     10     /// stale worker to change scheduling. Hosts must keep polling this future
     11     /// to deliver late results; dropping it leaves durable unresolved work.
     12     pub async fn deliver_push(
     13         &self,
     14         operation_id: SyncId,
     15     ) -> Result<DeliveryExecutionReceipt, Error> {
     16         self.deliver_push_inner(operation_id, None).await
     17     }
     18 
     19     /// Attempts only an exact nonempty subset of the frozen delivery targets.
     20     /// The full persisted request, claim and raw result bindings remain intact.
     21     /// Callers own selection policy; no ineligible target may be attempted.
     22     pub async fn deliver_push_selected(
     23         &self,
     24         operation_id: SyncId,
     25         selected: radroots_transport::TargetSet,
     26     ) -> Result<DeliveryExecutionReceipt, Error> {
     27         self.deliver_push_inner(operation_id, Some(selected)).await
     28     }
     29 
     30     async fn deliver_push_inner(
     31         &self,
     32         operation_id: SyncId,
     33         selected: Option<radroots_transport::TargetSet>,
     34     ) -> Result<DeliveryExecutionReceipt, Error> {
     35         let status = self.push_status(operation_id).await?.ok_or_else(|| {
     36             if self.sink.is_none() {
     37                 Error::MissingSink
     38             } else {
     39                 Error::StorageFailed
     40             }
     41         })?;
     42         let plan = status.delivery_plan();
     43         if plan.state().is_terminal() {
     44             return Ok(DeliveryExecutionReceipt {
     45                 plan: plan.clone(),
     46                 replay: true,
     47             });
     48         }
     49         if status.artifact.signing_state() != SigningState::Signed {
     50             return Err(Error::InvalidSignerOutput);
     51         }
     52         if !status.artifact.admission_state().is_admitted() {
     53             return Err(Error::AdmissionFailed);
     54         }
     55         let now = self.delivery_now()?.max(plan.updated_at_unix_ms());
     56         // This invocation is fresh scheduling authority, not a late callback.
     57         // Reconciliation is a complete bounded action; never also send a retry.
     58         if plan.pending_delivery_facts().next().is_some() {
     59             let plan = self
     60                 .reconcile_delivery(status.delivery_history(), None, now)
     61                 .await?;
     62             return Ok(DeliveryExecutionReceipt { plan, replay: true });
     63         }
     64         if plan
     65             .claim_evidence()
     66             .is_some_and(|claim| now < claim.expires_at_unix_ms())
     67         {
     68             return Err(Error::WorkClaimConflict);
     69         }
     70         if plan
     71             .retry()
     72             .is_some_and(|retry| now < retry.not_before_unix_ms())
     73         {
     74             return Err(Error::DeliveryDeferred);
     75         }
     76         let sink = self.sink.as_deref().ok_or(Error::MissingSink)?;
     77         let request = plan.request().cloned().ok_or(Error::InvalidSignerOutput)?;
     78         if let Some(selected) = &selected {
     79             request
     80                 .validate_target_selection(selected)
     81                 .map_err(|_| Error::InvalidDeliveryRequest)?;
     82         }
     83         let claimed = self.claim_delivery_plan(plan, now).await?;
     84         let claim = claimed
     85             .claim_evidence()
     86             .cloned()
     87             .ok_or(Error::StorageFailed)?;
     88         // Re-read after claim admission: an observed stop or newer worker must
     89         // prevent this callback from initiating an effect.
     90         let current = self.delivery_history(claimed.plan_id()).await?;
     91         if current.plan().state().is_terminal() {
     92             return Ok(DeliveryExecutionReceipt {
     93                 plan: current.plan().clone(),
     94                 replay: true,
     95             });
     96         }
     97         let execution_started_at = self
     98             .delivery_now()?
     99             .max(current.plan().updated_at_unix_ms());
    100         if current.plan().claim_evidence() != Some(&claim)
    101             || execution_started_at >= claim.expires_at_unix_ms()
    102         {
    103             return Err(Error::WorkClaimConflict);
    104         }
    105         if current.plan().pending_delivery_facts().next().is_some() {
    106             let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision())
    107                 .map_err(map_storage_error)?;
    108             let plan = self
    109                 .reconcile_delivery(&current, Some(fence), execution_started_at)
    110                 .await?;
    111             return Ok(DeliveryExecutionReceipt { plan, replay: true });
    112         }
    113         let outcome = if execution_started_at >= request.deadline_unix_ms() {
    114             DeliveryAttemptOutcome::SinkFailure(
    115                 SinkFailure::for_request(
    116                     &request,
    117                     "delivery_deadline_exceeded",
    118                     Retryability::Terminal,
    119                     None,
    120                     None,
    121                     Vec::new(),
    122                 )
    123                 .map_err(|_| Error::InvalidDeliveryRequest)?,
    124             )
    125         } else {
    126             let result = match selected.clone() {
    127                 Some(targets) => sink.deliver_selected(request.clone(), targets).await,
    128                 None => sink.deliver(request.clone()).await,
    129             };
    130             let allowed = |rows: &[radroots_transport::sink::DeliveryTargetReceipt]| {
    131                 selected.as_ref().is_none_or(|targets| {
    132                     rows.iter()
    133                         .all(|row| !row.was_attempted() || targets.targets().contains(row.target()))
    134                 })
    135             };
    136             match result {
    137                 Ok(receipt)
    138                     if receipt.validate_for_request(&request).is_ok()
    139                         && allowed(receipt.target_receipts()) =>
    140                 {
    141                     DeliveryAttemptOutcome::Receipt(receipt)
    142                 }
    143                 Err(failure)
    144                     if failure.validate_for_request(&request).is_ok()
    145                         && allowed(failure.partial_evidence()) =>
    146                 {
    147                     DeliveryAttemptOutcome::SinkFailure(failure)
    148                 }
    149                 Ok(_) | Err(_) => {
    150                     DeliveryAttemptOutcome::SinkFailure(SinkFailure::invalid_contract(&request))
    151                 }
    152             }
    153         };
    154         let observed = self.delivery_now().map(|at| at.max(execution_started_at));
    155         // A known pre-effect lower bound suffices for the non-expiring delivery
    156         // fact. It does not claim a measured response time or grant retry rights.
    157         // Strict signing/expiring authorization observation rules are separate.
    158         let command = AuthoredAtomicCommand::RecordDelivery(
    159             RecordDeliveryFact::new(
    160                 claimed.plan_id(),
    161                 claimed.artifact_id(),
    162                 claim.clone(),
    163                 outcome.clone(),
    164                 observed.unwrap_or(execution_started_at),
    165             )
    166             .map_err(map_storage_error)?,
    167         );
    168         let receipt = self
    169             .storage
    170             .execute_authored(command.clone())
    171             .await
    172             .map_err(map_storage_error)?;
    173         if !receipt.matches_command(&command) {
    174             return Err(Error::StorageFailed);
    175         }
    176         let current = self.delivery_history(claimed.plan_id()).await?;
    177         if !current
    178             .plan()
    179             .delivery_facts()
    180             .iter()
    181             .any(|fact| fact.claim() == &claim && fact.outcome() == &outcome)
    182         {
    183             return Err(Error::StorageFailed);
    184         }
    185         let at = observed?;
    186         let plan = if !current.plan().state().is_terminal()
    187             && current.plan().claim_evidence() == Some(&claim)
    188             && at < claim.expires_at_unix_ms()
    189         {
    190             let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision())
    191                 .map_err(map_storage_error)?;
    192             self.reconcile_delivery(&current, Some(fence), at).await?
    193         } else {
    194             // Acceptance remains visible even when scheduling was stopped,
    195             // expired or superseded. Only a fresh caller may reconcile later.
    196             current.plan().clone()
    197         };
    198         Ok(DeliveryExecutionReceipt {
    199             plan,
    200             replay: false,
    201         })
    202     }
    203 
    204     fn delivery_now(&self) -> Result<u64, Error> {
    205         match self.clock.now_unix_ms()? {
    206             0 => Err(Error::ClockUnavailable),
    207             at => Ok(at),
    208         }
    209     }
    210 
    211     async fn delivery_history(
    212         &self,
    213         id: AuthoredDeliveryPlanId,
    214     ) -> Result<AuthoredDeliveryHistory, Error> {
    215         let history = self
    216             .storage
    217             .authored_delivery_history(id)
    218             .await
    219             .map_err(map_storage_error)?
    220             .ok_or(Error::StorageFailed)?;
    221         history.validate().map_err(map_storage_error)?;
    222         if history.plan().plan_id() != id {
    223             return Err(Error::StorageFailed);
    224         }
    225         Ok(history)
    226     }
    227 
    228     async fn reconcile_delivery(
    229         &self,
    230         history: &AuthoredDeliveryHistory,
    231         fence: Option<WorkFence>,
    232         at: u64,
    233     ) -> Result<AuthoredDeliveryPlan, Error> {
    234         let plan = history.plan();
    235         let at = at.max(plan.updated_at_unix_ms());
    236         // Never advance host time to a future fact to bypass a competing lease.
    237         if plan
    238             .delivery_facts()
    239             .iter()
    240             .any(|fact| fact.observed_at_unix_ms() > at)
    241         {
    242             return Err(Error::ClockUnavailable);
    243         }
    244         let retry = reconciliation_retry(history, at)?;
    245         let reconcile =
    246             ReconcileDeliveryFacts::new(plan, fence, retry, at).map_err(map_storage_error)?;
    247         reconcile.apply_to(history).map_err(map_claim_error)?;
    248         let command = AuthoredAtomicCommand::ReconcileDelivery(reconcile);
    249         let receipt = self
    250             .storage
    251             .execute_authored(command.clone())
    252             .await
    253             .map_err(map_claim_error)?;
    254         if !receipt.matches_command(&command) {
    255             return Err(Error::StorageFailed);
    256         }
    257         // A receipt replay is historical; expose current stop/facts/authority.
    258         Ok(self.delivery_history(plan.plan_id()).await?.plan().clone())
    259     }
    260 
    261     async fn claim_delivery_plan(
    262         &self,
    263         plan: &AuthoredDeliveryPlan,
    264         acquired_at: u64,
    265     ) -> Result<AuthoredDeliveryPlan, Error> {
    266         let generation = plan
    267             .claim_evidence()
    268             .map_or(1, |claim| claim.generation().get().saturating_add(1));
    269         let expires_at = acquired_at
    270             .checked_add(self.deadlines.timeout_ms(OperationKind::Deliver))
    271             .ok_or(Error::DeadlineOverflow)?;
    272         let claim = WorkClaim::new(
    273             *self.ids.next_id(OperationKind::Deliver)?.as_bytes(),
    274             DELIVERY_CLAIM_OWNER,
    275             NonZeroU64::new(generation).ok_or(Error::StorageFailed)?,
    276             acquired_at,
    277             expires_at,
    278             plan.revision(),
    279         )
    280         .map_err(map_storage_error)?;
    281         let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    282             ClaimAuthoredTarget::DeliveryPlan(plan.plan_id()),
    283             claim.clone(),
    284         ));
    285         let receipt = self
    286             .storage
    287             .execute_authored(command.clone())
    288             .await
    289             .map_err(map_claim_error)?;
    290         if !receipt.matches_command(&command) {
    291             return Err(Error::StorageFailed);
    292         }
    293         let AuthoredAtomicOutcome::DeliveryPlan(claimed) = receipt.outcome() else {
    294             return Err(Error::StorageFailed);
    295         };
    296         if claimed.validate().is_err()
    297             || claimed.plan_id() != plan.plan_id()
    298             || claimed.artifact_id() != plan.artifact_id()
    299             || claimed.request() != plan.request()
    300             || claimed.created_at_unix_ms() != plan.created_at_unix_ms()
    301             || claimed.claim_evidence() != Some(&claim)
    302             || plan.revision().get().checked_add(1) != Some(claimed.revision().get())
    303             || claimed.updated_at_unix_ms() != acquired_at
    304             || receipt.committed_at_unix_ms() != acquired_at
    305             || claimed.attempts() != plan.attempts()
    306             || claimed.state() != plan.state()
    307             || claimed.retry() != plan.retry()
    308         {
    309             return Err(Error::StorageFailed);
    310         }
    311         Ok(claimed.clone())
    312     }
    313 }
    314 
    315 fn reconciliation_retry(
    316     history: &AuthoredDeliveryHistory,
    317     at: u64,
    318 ) -> Result<Option<RetrySchedule>, Error> {
    319     history
    320         .require_pending_fact_provenance()
    321         .map_err(map_storage_error)?;
    322     let plan = history.plan();
    323     let mut count = plan.attempt_count();
    324     let mut matched_legacy = Vec::new();
    325     let mut last = plan.attempts().last().map(|attempt| attempt.outcome());
    326     for fact in plan.pending_delivery_facts() {
    327         let original = history
    328             .claims()
    329             .iter()
    330             .find(|entry| entry.claim() == fact.claim())
    331             .ok_or(Error::StorageFailed)?;
    332         let legacy = plan.attempts().iter().find(|attempt| {
    333             !matched_legacy.contains(&attempt.attempt())
    334                 && attempt.claim_evidence().is_none()
    335                 && original.prior_attempt_count().checked_add(1) == Some(attempt.attempt().get())
    336                 && attempt.recorded_at_unix_ms() >= fact.claim().acquired_at_unix_ms()
    337                 && attempt.recorded_at_unix_ms() < fact.claim().expires_at_unix_ms()
    338         });
    339         if let Some(legacy) = legacy {
    340             if legacy.outcome() != fact.outcome() {
    341                 return Err(Error::StorageConflict);
    342             }
    343             matched_legacy.push(legacy.attempt());
    344         } else {
    345             count = count.checked_add(1).ok_or(Error::StorageFailed)?;
    346             last = Some(fact.outcome());
    347         }
    348     }
    349     if count > DELIVERY_PLAN_ATTEMPTS_MAX {
    350         return Err(Error::StorageFailed);
    351     }
    352     let satisfaction = plan.delivery_satisfaction().map_err(map_storage_error)?;
    353     if satisfaction != SatisfactionState::Pending || count == DELIVERY_PLAN_ATTEMPTS_MAX {
    354         return Ok(None);
    355     }
    356     let last = last.ok_or(Error::StorageFailed)?;
    357     let normalized;
    358     let last = if let DeliveryAttemptOutcome::SinkFailure(failure) = last {
    359         normalized = DeliveryAttemptOutcome::SinkFailure(normalize_sink_failure(
    360             plan.request().ok_or(Error::InvalidDeliveryRequest)?,
    361             failure.clone(),
    362             at,
    363         )?);
    364         &normalized
    365     } else {
    366         last
    367     };
    368     delivery_retry_schedule(
    369         NonZeroU32::new(count).ok_or(Error::StorageFailed)?,
    370         last,
    371         satisfaction,
    372         at,
    373     )
    374 }