lib

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

wire.rs (3409B)


      1 use super::{PullReceipt, PullTermination, summary::PullTargetSummaries};
      2 use crate::{
      3     ingest::IngestReceipt,
      4     policy::{Error, SyncId},
      5 };
      6 use radroots_transport::{
      7     outcome::{FetchTargetOutcome, FetchTargetState},
      8     source::FetchCursor,
      9 };
     10 
     11 impl<'de> serde::Deserialize<'de> for PullReceipt {
     12     fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
     13     where
     14         D: serde::Deserializer<'de>,
     15     {
     16         // Preserve the existing receipt's legacy fields and unknown-field policy.
     17         #[derive(serde::Deserialize)]
     18         struct Wire {
     19             sync_id: SyncId,
     20             deadline_unix_ms: u64,
     21             pages_fetched: u16,
     22             events_observed: usize,
     23             ingest_outcomes: Vec<Result<IngestReceipt, Error>>,
     24             target_outcomes: Vec<FetchTargetOutcome>,
     25             #[serde(default)]
     26             target_summaries: Option<PullTargetSummaries>,
     27             termination: PullTermination,
     28             resume_from: Option<FetchCursor>,
     29         }
     30         let wire = Wire::deserialize(deserializer)?;
     31         if let Some(summaries) = &wire.target_summaries {
     32             let summaries = summaries.as_slice();
     33             if summaries
     34                 .iter()
     35                 .any(|summary| summary.pages_observed() != wire.pages_fetched)
     36             {
     37                 return Err(serde::de::Error::custom("pull summary page count differs"));
     38             }
     39             let mut targets = std::collections::BTreeSet::new();
     40             for outcome in &wire.target_outcomes {
     41                 let Some(summary) = summaries
     42                     .iter()
     43                     .find(|summary| summary.target() == outcome.target())
     44                 else {
     45                     return Err(serde::de::Error::custom(
     46                         "pull summary target inventory differs",
     47                     ));
     48                 };
     49                 if !targets.insert(outcome.target()) {
     50                     return Err(serde::de::Error::custom("duplicate pull final target"));
     51                 }
     52                 let consistent = match outcome.state() {
     53                     FetchTargetState::Complete => {
     54                         summary.pages_observed()
     55                             > summary.incomplete_pages() + summary.missing_outcome_pages()
     56                     }
     57                     state => summary.last_incomplete() == Some(state),
     58                 };
     59                 if !consistent {
     60                     return Err(serde::de::Error::custom("pull summary final state differs"));
     61                 }
     62             }
     63             for summary in summaries {
     64                 let observed_outcome = summary.pages_observed() > summary.missing_outcome_pages();
     65                 if targets.contains(summary.target()) != observed_outcome {
     66                     return Err(serde::de::Error::custom(
     67                         "pull summary outcome evidence differs",
     68                     ));
     69                 }
     70             }
     71         }
     72         Ok(Self {
     73             sync_id: wire.sync_id,
     74             deadline_unix_ms: wire.deadline_unix_ms,
     75             pages_fetched: wire.pages_fetched,
     76             events_observed: wire.events_observed,
     77             ingest_outcomes: wire.ingest_outcomes,
     78             target_outcomes: wire.target_outcomes,
     79             target_summaries: wire.target_summaries,
     80             termination: wire.termination,
     81             resume_from: wire.resume_from,
     82         })
     83     }
     84 }