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 }