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(¤t, 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(¤t, 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 }