pull.rs (9770B)
1 //! Bounded source pagination and ingestion. 2 3 mod summary; 4 #[cfg(feature = "serde")] 5 mod wire; 6 7 pub use summary::PullTargetSummary; 8 9 use radroots_transport::{ 10 FetchRequest, 11 outcome::FetchTargetOutcome, 12 source::{FETCH_PAGE_MAX_EVENTS, FetchBounds, FetchCursor, FetchSelector, NextPage}, 13 target::TargetSet, 14 }; 15 16 use crate::{ 17 Engine, 18 ingest::{AdmissionPolicy, IngestReceipt}, 19 policy::{Error, OperationKind, SyncId}, 20 }; 21 22 /// Hard maximum number of source pages in one explicit pull call. 23 pub const PULL_MAX_PAGES: u16 = 1_000; 24 25 /// Caller-owned bounds and continuation state for one pull operation. 26 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 27 #[derive(Clone, Debug, Eq, PartialEq)] 28 pub struct PullRequest { 29 targets: TargetSet, 30 page_limit: u16, 31 max_pages: u16, 32 cursor: Option<FetchCursor>, 33 selector: FetchSelector, 34 } 35 36 impl PullRequest { 37 pub fn new(targets: TargetSet, page_limit: u16, max_pages: u16) -> Result<Self, Error> { 38 if page_limit == 0 39 || page_limit > FETCH_PAGE_MAX_EVENTS 40 || max_pages == 0 41 || max_pages > PULL_MAX_PAGES 42 { 43 return Err(Error::InvalidPullRequest); 44 } 45 Ok(Self { 46 targets, 47 page_limit, 48 max_pages, 49 cursor: None, 50 selector: FetchSelector::all(), 51 }) 52 } 53 54 #[must_use] 55 pub fn with_cursor(mut self, cursor: FetchCursor) -> Self { 56 self.cursor = Some(cursor); 57 self 58 } 59 60 pub const fn targets(&self) -> &TargetSet { 61 &self.targets 62 } 63 64 pub const fn page_limit(&self) -> u16 { 65 self.page_limit 66 } 67 68 pub const fn max_pages(&self) -> u16 { 69 self.max_pages 70 } 71 72 pub const fn cursor(&self) -> Option<&FetchCursor> { 73 self.cursor.as_ref() 74 } 75 76 /// Applies explicit transport-neutral event constraints to every page. 77 #[must_use] 78 pub fn with_selector(mut self, selector: FetchSelector) -> Self { 79 self.selector = selector; 80 self 81 } 82 83 /// Returns the exact event constraints for this pull. 84 pub const fn selector(&self) -> &FetchSelector { 85 &self.selector 86 } 87 } 88 89 /// Deterministic reason why a bounded pull returned control to its caller. 90 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 91 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 92 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 93 pub enum PullTermination { 94 Complete, 95 PageLimit, 96 Deadline, 97 Cancelled, 98 SourceFailed, 99 } 100 101 /// Normalized receipt retaining ingest outcomes, final states and bounded page evidence. 102 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 103 #[derive(Clone, Debug, Eq, PartialEq)] 104 pub struct PullReceipt { 105 sync_id: SyncId, 106 deadline_unix_ms: u64, 107 pages_fetched: u16, 108 events_observed: usize, 109 ingest_outcomes: Vec<Result<IngestReceipt, Error>>, 110 target_outcomes: Vec<FetchTargetOutcome>, 111 target_summaries: Option<summary::PullTargetSummaries>, 112 termination: PullTermination, 113 resume_from: Option<FetchCursor>, 114 } 115 116 impl PullReceipt { 117 pub const fn sync_id(&self) -> SyncId { 118 self.sync_id 119 } 120 121 pub const fn deadline_unix_ms(&self) -> u64 { 122 self.deadline_unix_ms 123 } 124 125 pub const fn pages_fetched(&self) -> u16 { 126 self.pages_fetched 127 } 128 129 pub const fn events_observed(&self) -> usize { 130 self.events_observed 131 } 132 133 pub fn ingest_outcomes(&self) -> &[Result<IngestReceipt, Error>] { 134 self.ingest_outcomes.as_slice() 135 } 136 137 pub fn target_outcomes(&self) -> &[FetchTargetOutcome] { 138 self.target_outcomes.as_slice() 139 } 140 141 /// Cumulative evidence in request-target order, or `None` for a legacy receipt. 142 /// 143 /// A missing summary inventory is unknown evidence. Even complete summaries 144 /// must be interpreted with the termination and exact request scope; they 145 /// do not prove complete global history. 146 pub fn target_summaries(&self) -> Option<&[PullTargetSummary]> { 147 self.target_summaries.as_ref().map(|value| value.as_slice()) 148 } 149 150 pub const fn termination(&self) -> PullTermination { 151 self.termination 152 } 153 154 pub const fn resume_from(&self) -> Option<&FetchCursor> { 155 self.resume_from.as_ref() 156 } 157 } 158 159 impl Engine { 160 /// Fetches and ingests at most the caller-bounded number of source pages. 161 pub async fn pull( 162 &self, 163 request: PullRequest, 164 admission: &dyn AdmissionPolicy, 165 ) -> Result<PullReceipt, Error> { 166 let source = self.source.as_deref().ok_or(Error::MissingSource)?; 167 let sync_id = self.ids.next_id(OperationKind::Pull)?; 168 let started_at = self.clock.now_unix_ms()?; 169 let deadline_unix_ms = self 170 .deadlines 171 .deadline_unix_ms(OperationKind::Pull, started_at)?; 172 let bounds = FetchBounds::new(request.page_limit, deadline_unix_ms) 173 .map_err(|_| Error::InvalidPullRequest)?; 174 let mut receipt = PullReceipt { 175 sync_id, 176 deadline_unix_ms, 177 pages_fetched: 0, 178 events_observed: 0, 179 ingest_outcomes: Vec::new(), 180 target_outcomes: Vec::new(), 181 target_summaries: Some(summary::PullTargetSummaries::new(&request.targets)), 182 termination: PullTermination::Complete, 183 resume_from: request.cursor.clone(), 184 }; 185 let mut cursor = request.cursor; 186 187 for page_index in 0..request.max_pages { 188 if page_index != 0 && self.clock.now_unix_ms()? >= deadline_unix_ms { 189 receipt.termination = PullTermination::Deadline; 190 receipt.resume_from = cursor; 191 return Ok(receipt); 192 } 193 let mut fetch = FetchRequest::new( 194 fetch_request_id(sync_id, page_index), 195 request.targets.clone(), 196 bounds, 197 ) 198 .map_err(|_| Error::InvalidPullRequest)? 199 .with_selector(request.selector.clone()); 200 if let Some(current) = cursor.clone() { 201 fetch = fetch.with_cursor(current); 202 } 203 let page = match source.fetch(fetch.clone()).await { 204 Ok(page) => page, 205 Err(_) => { 206 receipt.termination = PullTermination::SourceFailed; 207 receipt.resume_from = cursor; 208 return Ok(receipt); 209 } 210 }; 211 page.validate_for_request(&fetch) 212 .map_err(|_| Error::InvalidSourcePage)?; 213 receipt.pages_fetched += 1; 214 receipt.events_observed += page.events().len(); 215 merge_target_outcomes(&mut receipt.target_outcomes, page.target_outcomes()); 216 if let Some(summaries) = &mut receipt.target_summaries { 217 summaries.observe(page.target_outcomes()); 218 } 219 let outcomes = self.ingest_batch(page.events().to_vec(), admission).await; 220 receipt 221 .ingest_outcomes 222 .extend_from_slice(outcomes.outcomes()); 223 224 match page.next_page() { 225 NextPage::Complete => { 226 receipt.termination = PullTermination::Complete; 227 receipt.resume_from = None; 228 return Ok(receipt); 229 } 230 NextPage::Cancelled { resume_from } => { 231 receipt.termination = PullTermination::Cancelled; 232 receipt.resume_from = resume_from.clone(); 233 return Ok(receipt); 234 } 235 NextPage::Cursor(next) => { 236 cursor = Some(next.clone()); 237 if page_index + 1 == request.max_pages { 238 receipt.termination = PullTermination::PageLimit; 239 receipt.resume_from = cursor; 240 return Ok(receipt); 241 } 242 } 243 } 244 } 245 unreachable!("validated pull requests execute at least one page") 246 } 247 } 248 249 fn merge_target_outcomes(current: &mut Vec<FetchTargetOutcome>, page: &[FetchTargetOutcome]) { 250 for outcome in page { 251 if let Some(existing) = current 252 .iter_mut() 253 .find(|existing| existing.target() == outcome.target()) 254 { 255 *existing = outcome.clone(); 256 } else { 257 current.push(outcome.clone()); 258 } 259 } 260 } 261 262 fn fetch_request_id(sync_id: SyncId, page_index: u16) -> String { 263 const HEX: &[u8; 16] = b"0123456789abcdef"; 264 let mut value = String::with_capacity(48); 265 value.push_str("sync-"); 266 for byte in sync_id.as_bytes() { 267 value.push(HEX[(byte >> 4) as usize] as char); 268 value.push(HEX[(byte & 0x0f) as usize] as char); 269 } 270 value.push('-'); 271 value.push_str(page_index.to_string().as_str()); 272 value 273 } 274 275 #[cfg(feature = "serde")] 276 impl<'de> serde::Deserialize<'de> for PullRequest { 277 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 278 where 279 D: serde::Deserializer<'de>, 280 { 281 #[derive(serde::Deserialize)] 282 #[serde(deny_unknown_fields)] 283 struct Wire { 284 targets: TargetSet, 285 page_limit: u16, 286 max_pages: u16, 287 cursor: Option<FetchCursor>, 288 #[serde(default)] 289 selector: FetchSelector, 290 } 291 292 let wire = Wire::deserialize(deserializer)?; 293 let mut request = Self::new(wire.targets, wire.page_limit, wire.max_pages) 294 .map_err(serde::de::Error::custom)?; 295 if let Some(cursor) = wire.cursor { 296 request = request.with_cursor(cursor); 297 } 298 Ok(request.with_selector(wire.selector)) 299 } 300 }