lib

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

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 }