lib

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

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 }