authored_delivery.rs (31233B)
1 //! Independent durable delivery-plan, attempt, retry, and evidence models. 2 3 mod facts; 4 pub use facts::AuthoredDeliveryFact; 5 mod history; 6 pub use history::{AuthoredDeliveryClaim, AuthoredDeliveryHistory}; 7 mod reconciliation; 8 9 use core::num::{NonZeroU32, NonZeroU64}; 10 use radroots_transport::{ 11 DeliveryReceipt, DeliveryRequest, SinkFailure, 12 outcome::Retryability, 13 policy::{SatisfactionClass, SatisfactionPolicy, SatisfactionState, evaluate_satisfaction}, 14 sink::{DeliveryPayload, DeliveryRequestId}, 15 target::TargetSet, 16 }; 17 use sha2::{Digest, Sha256}; 18 use std::vec::Vec; 19 20 use crate::{ 21 Error, 22 authored::{ 23 AuthoredArtifactId, FailureClass, RetrySchedule, WorkClaim, WorkFailure, WorkPhase, 24 }, 25 }; 26 27 pub const DELIVERY_PLAN_ATTEMPTS_MAX: u32 = 1_024; 28 29 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 30 #[cfg_attr(feature = "serde", serde(try_from = "[u8; 16]", into = "[u8; 16]"))] 31 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 32 pub struct AuthoredDeliveryPlanId([u8; 16]); 33 34 impl AuthoredDeliveryPlanId { 35 pub const fn new(value: [u8; 16]) -> Result<Self, Error> { 36 if all_zero(&value) { 37 Err(Error::InvalidAuthoredDeliveryPlan) 38 } else { 39 Ok(Self(value)) 40 } 41 } 42 43 pub const fn as_bytes(&self) -> &[u8; 16] { 44 &self.0 45 } 46 } 47 48 impl TryFrom<[u8; 16]> for AuthoredDeliveryPlanId { 49 type Error = Error; 50 fn try_from(value: [u8; 16]) -> Result<Self, Self::Error> { 51 Self::new(value) 52 } 53 } 54 55 impl From<AuthoredDeliveryPlanId> for [u8; 16] { 56 fn from(value: AuthoredDeliveryPlanId) -> Self { 57 value.0 58 } 59 } 60 61 /// Exact delivery intent persisted before a signed payload exists. 62 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 63 #[derive(Clone, Debug, Eq, PartialEq)] 64 pub struct AuthoredDeliveryIntent { 65 request_id: DeliveryRequestId, 66 target_set: TargetSet, 67 satisfaction: SatisfactionPolicy, 68 deadline_unix_ms: u64, 69 } 70 71 impl AuthoredDeliveryIntent { 72 pub fn new( 73 request_id: impl Into<String>, 74 target_set: TargetSet, 75 satisfaction: SatisfactionPolicy, 76 deadline_unix_ms: u64, 77 ) -> Result<Self, Error> { 78 if deadline_unix_ms == 0 || satisfaction.validate_for(&target_set).is_err() { 79 return Err(Error::InvalidAuthoredDeliveryPlan); 80 } 81 Ok(Self { 82 request_id: DeliveryRequestId::parse(request_id) 83 .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?, 84 target_set, 85 satisfaction, 86 deadline_unix_ms, 87 }) 88 } 89 90 pub fn from_request(request: &DeliveryRequest) -> Self { 91 Self { 92 request_id: request.request_id().clone(), 93 target_set: request.target_set().clone(), 94 satisfaction: request.satisfaction().clone(), 95 deadline_unix_ms: request.deadline_unix_ms(), 96 } 97 } 98 99 pub fn materialize(&self, payload: DeliveryPayload) -> Result<DeliveryRequest, Error> { 100 DeliveryRequest::new( 101 self.request_id.as_str(), 102 payload, 103 self.target_set.clone(), 104 self.satisfaction.clone(), 105 self.deadline_unix_ms, 106 ) 107 .map_err(|_| Error::InvalidAuthoredDeliveryPlan) 108 } 109 110 pub const fn request_id(&self) -> &DeliveryRequestId { 111 &self.request_id 112 } 113 pub const fn target_set(&self) -> &TargetSet { 114 &self.target_set 115 } 116 pub const fn satisfaction(&self) -> &SatisfactionPolicy { 117 &self.satisfaction 118 } 119 pub const fn deadline_unix_ms(&self) -> u64 { 120 self.deadline_unix_ms 121 } 122 } 123 124 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 125 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 126 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 127 pub enum AuthoredDeliveryState { 128 Pending, 129 Retryable, 130 Satisfied, 131 Exhausted, 132 FailedTerminal, 133 Cancelled, 134 } 135 136 impl AuthoredDeliveryState { 137 pub const fn is_terminal(self) -> bool { 138 matches!( 139 self, 140 Self::Satisfied | Self::Exhausted | Self::FailedTerminal | Self::Cancelled 141 ) 142 } 143 } 144 145 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 146 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 147 #[derive(Clone, Debug, Eq, PartialEq)] 148 pub enum DeliveryAttemptOutcome { 149 Receipt(DeliveryReceipt), 150 SinkFailure(SinkFailure), 151 } 152 153 impl DeliveryAttemptOutcome { 154 fn validate_for(&self, request: &DeliveryRequest) -> Result<(), Error> { 155 match self { 156 Self::Receipt(receipt) => receipt 157 .validate_for_request(request) 158 .map_err(|_| Error::InvalidAuthoredDeliveryPlan), 159 Self::SinkFailure(failure) => failure 160 .validate_for_request(request) 161 .map_err(|_| Error::InvalidAuthoredDeliveryPlan), 162 } 163 } 164 } 165 166 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 167 #[derive(Clone, Debug, Eq, PartialEq)] 168 pub struct AuthoredDeliveryAttempt { 169 attempt: NonZeroU32, 170 recorded_at_unix_ms: u64, 171 outcome: DeliveryAttemptOutcome, 172 satisfaction: SatisfactionState, 173 #[cfg_attr( 174 feature = "serde", 175 serde(default, skip_serializing_if = "Option::is_none") 176 )] 177 claim: Option<WorkClaim>, 178 } 179 180 impl AuthoredDeliveryAttempt { 181 pub fn reconstruct( 182 attempt: NonZeroU32, 183 recorded_at_unix_ms: u64, 184 outcome: DeliveryAttemptOutcome, 185 satisfaction: SatisfactionState, 186 ) -> Result<Self, Error> { 187 if recorded_at_unix_ms == 0 { 188 return Err(Error::InvalidAuthoredDeliveryPlan); 189 } 190 Ok(Self { 191 attempt, 192 recorded_at_unix_ms, 193 outcome, 194 satisfaction, 195 claim: None, 196 }) 197 } 198 199 pub const fn attempt(&self) -> NonZeroU32 { 200 self.attempt 201 } 202 pub const fn recorded_at_unix_ms(&self) -> u64 { 203 self.recorded_at_unix_ms 204 } 205 pub const fn outcome(&self) -> &DeliveryAttemptOutcome { 206 &self.outcome 207 } 208 pub const fn satisfaction(&self) -> SatisfactionState { 209 self.satisfaction 210 } 211 /// Exact fact provenance for new reconciliation; absent on historical attempts. 212 pub const fn claim_evidence(&self) -> Option<&WorkClaim> { 213 self.claim.as_ref() 214 } 215 } 216 217 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 218 #[cfg_attr( 219 feature = "serde", 220 serde( 221 try_from = "AuthoredDeliveryPlanWire", 222 into = "AuthoredDeliveryPlanWire" 223 ) 224 )] 225 #[derive(Clone, Debug, Eq, PartialEq)] 226 pub struct AuthoredDeliveryPlan { 227 plan_id: AuthoredDeliveryPlanId, 228 artifact_id: AuthoredArtifactId, 229 request_digest: [u8; 32], 230 intent: AuthoredDeliveryIntent, 231 request: Option<DeliveryRequest>, 232 state: AuthoredDeliveryState, 233 attempts: Vec<AuthoredDeliveryAttempt>, 234 delivery_facts: Vec<AuthoredDeliveryFact>, 235 stop_requested_at_unix_ms: Option<u64>, 236 attempt_count: u32, 237 retry: Option<RetrySchedule>, 238 claim: Option<WorkClaim>, 239 last_failure: Option<WorkFailure>, 240 created_at_unix_ms: u64, 241 updated_at_unix_ms: u64, 242 revision: NonZeroU64, 243 } 244 245 impl AuthoredDeliveryPlan { 246 pub fn new( 247 plan_id: AuthoredDeliveryPlanId, 248 artifact_id: AuthoredArtifactId, 249 intent: AuthoredDeliveryIntent, 250 created_at_unix_ms: u64, 251 ) -> Result<Self, Error> { 252 let request_digest = delivery_intent_digest(&intent); 253 Self::reconstruct(Self { 254 plan_id, 255 artifact_id, 256 request_digest, 257 intent, 258 request: None, 259 state: AuthoredDeliveryState::Pending, 260 attempts: Vec::new(), 261 delivery_facts: Vec::new(), 262 stop_requested_at_unix_ms: None, 263 attempt_count: 0, 264 retry: None, 265 claim: None, 266 last_failure: None, 267 created_at_unix_ms, 268 updated_at_unix_ms: created_at_unix_ms, 269 revision: NonZeroU64::MIN, 270 }) 271 } 272 273 pub fn new_bound( 274 plan_id: AuthoredDeliveryPlanId, 275 artifact_id: AuthoredArtifactId, 276 request: DeliveryRequest, 277 created_at_unix_ms: u64, 278 ) -> Result<Self, Error> { 279 let intent = AuthoredDeliveryIntent::from_request(&request); 280 let mut plan = Self::new(plan_id, artifact_id, intent, created_at_unix_ms)?; 281 plan.request = Some(request); 282 plan.validate()?; 283 Ok(plan) 284 } 285 286 pub fn reconstruct(value: Self) -> Result<Self, Error> { 287 value.validate()?; 288 Ok(value) 289 } 290 291 pub fn validate(&self) -> Result<(), Error> { 292 self.validate_facts()?; 293 if self.created_at_unix_ms == 0 294 || self.updated_at_unix_ms < self.created_at_unix_ms 295 || self.request_digest != delivery_intent_digest(&self.intent) 296 || self.attempt_count > DELIVERY_PLAN_ATTEMPTS_MAX 297 || usize::try_from(self.attempt_count).ok() != Some(self.attempts.len()) 298 || (matches!(self.state, AuthoredDeliveryState::Retryable) != self.retry.is_some()) 299 || (self.state.is_terminal() && self.claim.is_some()) 300 { 301 return Err(Error::InvalidAuthoredDeliveryPlan); 302 } 303 if self 304 .request 305 .as_ref() 306 .is_some_and(|request| AuthoredDeliveryIntent::from_request(request) != self.intent) 307 || (!self.attempts.is_empty() && self.request.is_none()) 308 { 309 return Err(Error::InvalidAuthoredDeliveryPlan); 310 } 311 for (index, attempt) in self.attempts.iter().enumerate() { 312 if attempt.attempt.get() != u32::try_from(index + 1).unwrap_or(u32::MAX) 313 || attempt.recorded_at_unix_ms < self.created_at_unix_ms 314 || self 315 .request 316 .as_ref() 317 .is_none_or(|request| attempt.outcome.validate_for(request).is_err()) 318 || self 319 .evaluate_outcomes(self.attempts[..=index].iter().map(|value| &value.outcome)) 320 .ok() 321 != Some(attempt.satisfaction) 322 { 323 return Err(Error::InvalidAuthoredDeliveryPlan); 324 } 325 } 326 if self 327 .attempts 328 .windows(2) 329 .any(|pair| pair[0].recorded_at_unix_ms > pair[1].recorded_at_unix_ms) 330 || self 331 .claim 332 .as_ref() 333 .is_some_and(|claim| claim.validate().is_err()) 334 || self.retry.as_ref().is_some_and(|retry| { 335 retry.failure().phase() != WorkPhase::Delivery 336 || retry.attempt().get() != self.attempt_count 337 || self.attempts.last().is_none_or(|attempt| { 338 retry.not_before_unix_ms() <= attempt.recorded_at_unix_ms 339 }) 340 }) 341 || self.claim.as_ref().is_some_and(|claim| { 342 claim.row_revision().get().checked_add(1) != Some(self.revision.get()) 343 || claim.acquired_at_unix_ms() != self.updated_at_unix_ms 344 }) 345 { 346 return Err(Error::InvalidAuthoredDeliveryPlan); 347 } 348 match self.state { 349 AuthoredDeliveryState::Pending 350 | AuthoredDeliveryState::Satisfied 351 | AuthoredDeliveryState::Cancelled => { 352 if self.retry.is_some() || self.last_failure.is_some() { 353 return Err(Error::InvalidAuthoredDeliveryPlan); 354 } 355 } 356 AuthoredDeliveryState::Exhausted => { 357 if self.retry.is_some() 358 || self.last_failure.as_ref().is_some_and(|failure| { 359 failure.phase() != WorkPhase::Delivery 360 || failure.class() != FailureClass::Terminal 361 }) 362 { 363 return Err(Error::InvalidAuthoredDeliveryPlan); 364 } 365 } 366 AuthoredDeliveryState::Retryable => { 367 if self.retry.as_ref().map(RetrySchedule::failure) != self.last_failure.as_ref() { 368 return Err(Error::InvalidAuthoredDeliveryPlan); 369 } 370 } 371 AuthoredDeliveryState::FailedTerminal => { 372 if !matches!( 373 self.last_failure.as_ref().map(WorkFailure::class), 374 Some(FailureClass::Terminal) 375 ) { 376 return Err(Error::InvalidAuthoredDeliveryPlan); 377 } 378 } 379 } 380 Ok(()) 381 } 382 383 pub fn claim(&mut self, claim: WorkClaim, now_unix_ms: u64) -> Result<(), Error> { 384 let existing_blocks = self.claim.as_ref().is_some_and(|existing| { 385 now_unix_ms < existing.expires_at_unix_ms() 386 || claim.generation() <= existing.generation() 387 }); 388 if self.state.is_terminal() 389 || self.stop_requested_at_unix_ms.is_some() 390 || self.delivery_facts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize 391 || self.request.is_none() 392 || existing_blocks 393 || claim.row_revision() != self.revision 394 || claim.acquired_at_unix_ms() != now_unix_ms 395 || self 396 .retry 397 .as_ref() 398 .is_some_and(|retry| now_unix_ms < retry.not_before_unix_ms()) 399 { 400 return Err(Error::DeliveryPlanClaimConflict); 401 } 402 claim.validate()?; 403 let previous = self.clone(); 404 self.claim = Some(claim); 405 if let Err(error) = self.advance(now_unix_ms) { 406 *self = previous; 407 return Err(error); 408 } 409 if let Err(error) = self.validate() { 410 *self = previous; 411 return Err(error); 412 } 413 Ok(()) 414 } 415 416 pub fn bind_signed_event( 417 &mut self, 418 event: radroots_event::SignedEvent, 419 bound_at_unix_ms: u64, 420 ) -> Result<(), Error> { 421 if self.request.is_some() || self.attempt_count != 0 || self.state.is_terminal() { 422 return Err(Error::InvalidAuthoredDeliveryPlan); 423 } 424 let previous = self.clone(); 425 self.request = Some(self.intent.materialize(DeliveryPayload::new(event))?); 426 if let Err(error) = self.advance(bound_at_unix_ms) { 427 *self = previous; 428 return Err(error); 429 } 430 if let Err(error) = self.validate() { 431 *self = previous; 432 return Err(error); 433 } 434 Ok(()) 435 } 436 437 pub fn apply_receipt( 438 &mut self, 439 token: &[u8; 16], 440 generation: NonZeroU64, 441 claim_revision: NonZeroU64, 442 receipt: DeliveryReceipt, 443 retry: Option<RetrySchedule>, 444 recorded_at_unix_ms: u64, 445 ) -> Result<(), Error> { 446 self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?; 447 let request = self 448 .request 449 .as_ref() 450 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 451 receipt 452 .validate_for_request(request) 453 .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; 454 let satisfaction = self.evaluate_with(DeliveryAttemptOutcome::Receipt(receipt.clone()))?; 455 let at_attempt_limit = self.attempt_count.saturating_add(1) == DELIVERY_PLAN_ATTEMPTS_MAX; 456 if satisfaction == SatisfactionState::Pending && at_attempt_limit && retry.is_some() { 457 return Err(Error::InvalidRetrySchedule); 458 } 459 let (state, last_failure) = match satisfaction { 460 SatisfactionState::Satisfied if retry.is_none() => { 461 (AuthoredDeliveryState::Satisfied, None) 462 } 463 SatisfactionState::Exhausted if retry.is_none() => { 464 (AuthoredDeliveryState::Exhausted, None) 465 } 466 SatisfactionState::Pending if at_attempt_limit && retry.is_none() => ( 467 AuthoredDeliveryState::Exhausted, 468 Some(WorkFailure::new( 469 "delivery_attempt_limit", 470 WorkPhase::Delivery, 471 FailureClass::Terminal, 472 None, 473 None, 474 )?), 475 ), 476 SatisfactionState::Pending => { 477 let schedule = retry.as_ref().ok_or(Error::InvalidRetrySchedule)?; 478 if schedule.failure().phase() != WorkPhase::Delivery { 479 return Err(Error::InvalidRetrySchedule); 480 } 481 ( 482 AuthoredDeliveryState::Retryable, 483 Some(schedule.failure().clone()), 484 ) 485 } 486 SatisfactionState::Satisfied | SatisfactionState::Exhausted => { 487 return Err(Error::InvalidRetrySchedule); 488 } 489 }; 490 self.apply_attempt( 491 DeliveryAttemptOutcome::Receipt(receipt), 492 satisfaction, 493 state, 494 retry, 495 last_failure, 496 recorded_at_unix_ms, 497 ) 498 } 499 500 pub fn apply_sink_failure( 501 &mut self, 502 token: &[u8; 16], 503 generation: NonZeroU64, 504 claim_revision: NonZeroU64, 505 failure: SinkFailure, 506 retry: Option<RetrySchedule>, 507 recorded_at_unix_ms: u64, 508 ) -> Result<(), Error> { 509 self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?; 510 let request = self 511 .request 512 .as_ref() 513 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 514 failure 515 .validate_for_request(request) 516 .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; 517 let outcome = DeliveryAttemptOutcome::SinkFailure(failure.clone()); 518 let satisfaction = self.evaluate_with(outcome.clone())?; 519 let at_attempt_limit = self.attempt_count.saturating_add(1) == DELIVERY_PLAN_ATTEMPTS_MAX; 520 if satisfaction == SatisfactionState::Pending && at_attempt_limit && retry.is_some() { 521 return Err(Error::InvalidRetrySchedule); 522 } 523 let typed = WorkFailure::new( 524 failure.code(), 525 WorkPhase::Delivery, 526 match failure.retryability() { 527 Retryability::Retryable => FailureClass::Retryable, 528 Retryability::Terminal | Retryability::NotApplicable => FailureClass::Terminal, 529 }, 530 failure.retry_after_unix_ms(), 531 failure.message().map(str::to_owned), 532 )?; 533 let (state, retry, last_failure) = if satisfaction == SatisfactionState::Satisfied { 534 if retry.is_some() { 535 return Err(Error::InvalidRetrySchedule); 536 } 537 (AuthoredDeliveryState::Satisfied, None, None) 538 } else if satisfaction == SatisfactionState::Exhausted { 539 if retry.is_some() { 540 return Err(Error::InvalidRetrySchedule); 541 } 542 (AuthoredDeliveryState::Exhausted, None, Some(typed)) 543 } else if at_attempt_limit && retry.is_none() { 544 ( 545 AuthoredDeliveryState::Exhausted, 546 None, 547 Some(WorkFailure::new( 548 "delivery_attempt_limit", 549 WorkPhase::Delivery, 550 FailureClass::Terminal, 551 None, 552 None, 553 )?), 554 ) 555 } else if failure.retryability() == Retryability::Retryable { 556 let schedule = retry.ok_or(Error::InvalidRetrySchedule)?; 557 if schedule.failure() != &typed { 558 return Err(Error::InvalidRetrySchedule); 559 } 560 ( 561 AuthoredDeliveryState::Retryable, 562 Some(schedule), 563 Some(typed), 564 ) 565 } else { 566 if retry.is_some() { 567 return Err(Error::InvalidRetrySchedule); 568 } 569 (AuthoredDeliveryState::FailedTerminal, None, Some(typed)) 570 }; 571 self.apply_attempt( 572 outcome, 573 satisfaction, 574 state, 575 retry, 576 last_failure, 577 recorded_at_unix_ms, 578 ) 579 } 580 581 pub fn cancel(&mut self, cancelled_at_unix_ms: u64) -> Result<(), Error> { 582 if self.state.is_terminal() { 583 return Err(Error::InvalidAuthoredDeliveryPlan); 584 } 585 self.request_stop(cancelled_at_unix_ms) 586 } 587 588 /// Retains the first stop intent, including when delivery already settled. 589 pub fn request_stop(&mut self, cancelled_at_unix_ms: u64) -> Result<(), Error> { 590 if cancelled_at_unix_ms < self.created_at_unix_ms { 591 return Err(Error::InvalidAuthoredDeliveryPlan); 592 } 593 if self.stop_requested_at_unix_ms.is_some() { 594 return Ok(()); 595 } 596 let previous = self.clone(); 597 self.stop_requested_at_unix_ms = Some(cancelled_at_unix_ms); 598 if !self.state.is_terminal() { 599 self.state = AuthoredDeliveryState::Cancelled; 600 self.last_failure = None; 601 } 602 self.claim = None; 603 self.retry = None; 604 if let Err(error) = self.advance(cancelled_at_unix_ms) { 605 *self = previous; 606 return Err(error); 607 } 608 if let Err(error) = self.validate() { 609 *self = previous; 610 return Err(error); 611 } 612 Ok(()) 613 } 614 615 fn apply_attempt( 616 &mut self, 617 outcome: DeliveryAttemptOutcome, 618 satisfaction: SatisfactionState, 619 state: AuthoredDeliveryState, 620 retry: Option<RetrySchedule>, 621 last_failure: Option<WorkFailure>, 622 recorded_at_unix_ms: u64, 623 ) -> Result<(), Error> { 624 let next = self 625 .attempt_count 626 .checked_add(1) 627 .filter(|attempt| *attempt <= DELIVERY_PLAN_ATTEMPTS_MAX) 628 .and_then(NonZeroU32::new) 629 .ok_or(Error::DeliveryAttemptOverflow)?; 630 let previous = self.clone(); 631 self.attempt_count = next.get(); 632 self.attempts.push(AuthoredDeliveryAttempt::reconstruct( 633 next, 634 recorded_at_unix_ms, 635 outcome, 636 satisfaction, 637 )?); 638 self.state = state; 639 self.retry = retry; 640 self.claim = None; 641 self.last_failure = last_failure; 642 if let Err(error) = self.advance(recorded_at_unix_ms) { 643 *self = previous; 644 return Err(error); 645 } 646 if let Err(error) = self.validate() { 647 *self = previous; 648 return Err(error); 649 } 650 Ok(()) 651 } 652 653 fn evaluate_with(&self, next: DeliveryAttemptOutcome) -> Result<SatisfactionState, Error> { 654 self.evaluate_outcomes( 655 self.attempts 656 .iter() 657 .map(|attempt| &attempt.outcome) 658 .chain(core::iter::once(&next)), 659 ) 660 } 661 662 fn evaluate_outcomes<'a, I>(&self, outcomes: I) -> Result<SatisfactionState, Error> 663 where 664 I: IntoIterator<Item = &'a DeliveryAttemptOutcome>, 665 { 666 let mut evidence = Vec::new(); 667 for outcome in outcomes { 668 match outcome { 669 DeliveryAttemptOutcome::Receipt(receipt) => { 670 evidence.extend( 671 receipt 672 .target_receipts() 673 .iter() 674 .map(|entry| (entry.target().fingerprint(), entry.outcome())), 675 ); 676 } 677 DeliveryAttemptOutcome::SinkFailure(failure) => { 678 evidence.extend( 679 failure 680 .partial_evidence() 681 .iter() 682 .map(|entry| (entry.target().fingerprint(), entry.outcome())), 683 ); 684 } 685 } 686 } 687 evaluate_satisfaction( 688 self.request 689 .as_ref() 690 .ok_or(Error::InvalidAuthoredDeliveryPlan)? 691 .satisfaction(), 692 self.request 693 .as_ref() 694 .ok_or(Error::InvalidAuthoredDeliveryPlan)? 695 .target_set(), 696 evidence, 697 ) 698 .map_err(|_| Error::InvalidAuthoredDeliveryPlan) 699 } 700 701 fn require_claim( 702 &self, 703 token: &[u8; 16], 704 generation: NonZeroU64, 705 claim_revision: NonZeroU64, 706 now_unix_ms: u64, 707 ) -> Result<(), Error> { 708 if !self.claim.as_ref().is_some_and(|claim| { 709 claim.matches_fence(token, generation, claim_revision, now_unix_ms) 710 }) { 711 return Err(Error::DeliveryPlanClaimConflict); 712 } 713 Ok(()) 714 } 715 716 fn advance(&mut self, at_unix_ms: u64) -> Result<(), Error> { 717 if at_unix_ms < self.updated_at_unix_ms { 718 return Err(Error::InvalidAuthoredDeliveryPlan); 719 } 720 self.revision = self 721 .revision 722 .get() 723 .checked_add(1) 724 .and_then(NonZeroU64::new) 725 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 726 self.updated_at_unix_ms = at_unix_ms; 727 Ok(()) 728 } 729 730 pub const fn plan_id(&self) -> AuthoredDeliveryPlanId { 731 self.plan_id 732 } 733 pub const fn artifact_id(&self) -> AuthoredArtifactId { 734 self.artifact_id 735 } 736 pub const fn request_digest(&self) -> &[u8; 32] { 737 &self.request_digest 738 } 739 pub const fn intent(&self) -> &AuthoredDeliveryIntent { 740 &self.intent 741 } 742 pub const fn request(&self) -> Option<&DeliveryRequest> { 743 self.request.as_ref() 744 } 745 pub const fn state(&self) -> AuthoredDeliveryState { 746 self.state 747 } 748 pub fn attempts(&self) -> &[AuthoredDeliveryAttempt] { 749 self.attempts.as_slice() 750 } 751 pub const fn attempt_count(&self) -> u32 { 752 self.attempt_count 753 } 754 pub const fn retry(&self) -> Option<&RetrySchedule> { 755 self.retry.as_ref() 756 } 757 pub const fn claim_evidence(&self) -> Option<&WorkClaim> { 758 self.claim.as_ref() 759 } 760 pub const fn last_failure(&self) -> Option<&WorkFailure> { 761 self.last_failure.as_ref() 762 } 763 pub const fn created_at_unix_ms(&self) -> u64 { 764 self.created_at_unix_ms 765 } 766 pub const fn updated_at_unix_ms(&self) -> u64 { 767 self.updated_at_unix_ms 768 } 769 pub const fn revision(&self) -> NonZeroU64 { 770 self.revision 771 } 772 773 /// Evaluates one prospective attempt against all durable prior evidence. 774 pub fn evaluate_next_attempt( 775 &self, 776 outcome: &DeliveryAttemptOutcome, 777 ) -> Result<SatisfactionState, Error> { 778 let request = self 779 .request 780 .as_ref() 781 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 782 outcome.validate_for(request)?; 783 self.evaluate_with(outcome.clone()) 784 } 785 } 786 787 #[cfg(feature = "serde")] 788 #[derive(serde::Serialize, serde::Deserialize)] 789 struct AuthoredDeliveryPlanWire { 790 plan_id: AuthoredDeliveryPlanId, 791 artifact_id: AuthoredArtifactId, 792 request_digest: [u8; 32], 793 intent: AuthoredDeliveryIntent, 794 request: Option<DeliveryRequest>, 795 state: AuthoredDeliveryState, 796 attempts: Vec<AuthoredDeliveryAttempt>, 797 #[serde(default)] 798 delivery_facts: Vec<AuthoredDeliveryFact>, 799 #[serde(default)] 800 stop_requested_at_unix_ms: Option<u64>, 801 attempt_count: u32, 802 retry: Option<RetrySchedule>, 803 claim: Option<WorkClaim>, 804 last_failure: Option<WorkFailure>, 805 created_at_unix_ms: u64, 806 updated_at_unix_ms: u64, 807 revision: NonZeroU64, 808 } 809 810 #[cfg(feature = "serde")] 811 impl TryFrom<AuthoredDeliveryPlanWire> for AuthoredDeliveryPlan { 812 type Error = Error; 813 fn try_from(value: AuthoredDeliveryPlanWire) -> Result<Self, Self::Error> { 814 Self::reconstruct(Self { 815 plan_id: value.plan_id, 816 artifact_id: value.artifact_id, 817 request_digest: value.request_digest, 818 intent: value.intent, 819 request: value.request, 820 state: value.state, 821 attempts: value.attempts, 822 delivery_facts: value.delivery_facts, 823 stop_requested_at_unix_ms: value.stop_requested_at_unix_ms.or_else(|| { 824 (value.state == AuthoredDeliveryState::Cancelled) 825 .then_some(value.updated_at_unix_ms) 826 }), 827 attempt_count: value.attempt_count, 828 retry: value.retry, 829 claim: value.claim, 830 last_failure: value.last_failure, 831 created_at_unix_ms: value.created_at_unix_ms, 832 updated_at_unix_ms: value.updated_at_unix_ms, 833 revision: value.revision, 834 }) 835 } 836 } 837 838 #[cfg(feature = "serde")] 839 impl From<AuthoredDeliveryPlan> for AuthoredDeliveryPlanWire { 840 fn from(value: AuthoredDeliveryPlan) -> Self { 841 Self { 842 plan_id: value.plan_id, 843 artifact_id: value.artifact_id, 844 request_digest: value.request_digest, 845 intent: value.intent, 846 request: value.request, 847 state: value.state, 848 attempts: value.attempts, 849 delivery_facts: value.delivery_facts, 850 stop_requested_at_unix_ms: value.stop_requested_at_unix_ms, 851 attempt_count: value.attempt_count, 852 retry: value.retry, 853 claim: value.claim, 854 last_failure: value.last_failure, 855 created_at_unix_ms: value.created_at_unix_ms, 856 updated_at_unix_ms: value.updated_at_unix_ms, 857 revision: value.revision, 858 } 859 } 860 } 861 862 fn delivery_intent_digest(intent: &AuthoredDeliveryIntent) -> [u8; 32] { 863 let mut hasher = Sha256::new(); 864 hash_field(&mut hasher, b"radroots.authored.delivery.v2"); 865 hash_field(&mut hasher, intent.request_id().as_str().as_bytes()); 866 for target in intent.target_set().targets() { 867 hash_field(&mut hasher, target.fingerprint().as_str().as_bytes()); 868 } 869 let policy = intent.satisfaction(); 870 hasher.update([match policy.class() { 871 SatisfactionClass::Accepted => 0, 872 SatisfactionClass::Delivered => 1, 873 }]); 874 hash_policy(&mut hasher, policy); 875 hasher.update(intent.deadline_unix_ms().to_be_bytes()); 876 hasher.finalize().into() 877 } 878 879 fn hash_policy(hasher: &mut Sha256, policy: &SatisfactionPolicy) { 880 let targets = policy.targets(); 881 if targets.is_any() { 882 hasher.update([0]); 883 } else if targets.is_all() { 884 hasher.update([1]); 885 } else if let Some(threshold) = targets.quorum_threshold() { 886 hasher.update([2]); 887 hasher.update(threshold.to_be_bytes()); 888 } else if let Some(required) = targets.required_targets() { 889 hasher.update([3]); 890 for target in required { 891 hash_field(hasher, target.as_str().as_bytes()); 892 } 893 } 894 } 895 896 fn hash_field(hasher: &mut Sha256, value: &[u8]) { 897 hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); 898 hasher.update(value); 899 } 900 901 const fn all_zero(bytes: &[u8; 16]) -> bool { 902 let mut index = 0; 903 while index < bytes.len() { 904 if bytes[index] != 0 { 905 return false; 906 } 907 index += 1; 908 } 909 true 910 }