lib

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

document_query.rs (6367B)


      1 //! Bounded opaque document inventory. A continuation is not a frozen snapshot.
      2 #[cfg(test)]
      3 mod tests;
      4 use super::{Error, ProjectionDocument, ProjectionGeneration, ProjectionId, valid_document_key};
      5 
      6 pub const PROJECTION_DOCUMENT_QUERY_LIMIT_MAX: u16 = 256;
      7 pub const PROJECTION_DOCUMENT_PAGE_BYTES_MAX: usize = super::PROJECTION_DOCUMENT_VALUE_MAX_BYTES;
      8 
      9 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     10 pub enum ProjectionDocumentGenerations {
     11     Exact(ProjectionGeneration),
     12     All,
     13 }
     14 
     15 #[derive(Clone, Debug, Eq, PartialEq)]
     16 pub struct ProjectionDocumentQuery {
     17     projection_id: ProjectionId,
     18     generations: ProjectionDocumentGenerations,
     19     limit: u16,
     20     after: Option<(ProjectionGeneration, String)>,
     21 }
     22 
     23 impl ProjectionDocumentQuery {
     24     pub fn new(
     25         projection_id: ProjectionId,
     26         generations: ProjectionDocumentGenerations,
     27         limit: u16,
     28     ) -> Result<Self, Error> {
     29         ProjectionId::parse(projection_id.as_str())?;
     30         if let ProjectionDocumentGenerations::Exact(generation) = generations {
     31             ProjectionGeneration::new(*generation.as_bytes())?;
     32         }
     33         if limit == 0 || limit > PROJECTION_DOCUMENT_QUERY_LIMIT_MAX {
     34             return Err(Error::InvalidProjectionDocument);
     35         }
     36         Ok(Self {
     37             projection_id,
     38             generations,
     39             limit,
     40             after: None,
     41         })
     42     }
     43     pub fn with_cursor(mut self, cursor: &ProjectionDocumentCursor) -> Result<Self, Error> {
     44         if self.projection_id != cursor.projection_id || self.generations != cursor.generations {
     45             return Err(Error::InvalidProjectionDocument);
     46         }
     47         self.after = Some(cursor.after.clone());
     48         Ok(self)
     49     }
     50     pub fn projection_id(&self) -> &ProjectionId {
     51         &self.projection_id
     52     }
     53     pub const fn generations(&self) -> ProjectionDocumentGenerations {
     54         self.generations
     55     }
     56     pub const fn limit(&self) -> u16 {
     57         self.limit
     58     }
     59     pub fn after(&self) -> Option<(ProjectionGeneration, &str)> {
     60         self.after
     61             .as_ref()
     62             .map(|(generation, key)| (*generation, key.as_str()))
     63     }
     64     pub fn matches(&self, generation: ProjectionGeneration, key: &str) -> bool {
     65         (match self.generations {
     66             ProjectionDocumentGenerations::Exact(expected) => generation == expected,
     67             ProjectionDocumentGenerations::All => true,
     68         }) && self.after().is_none_or(|after| (generation, key) > after)
     69     }
     70 }
     71 
     72 /// In-process continuation; the next call independently supplies its scope.
     73 #[derive(Clone, Debug, Eq, PartialEq)]
     74 pub struct ProjectionDocumentCursor {
     75     projection_id: ProjectionId,
     76     generations: ProjectionDocumentGenerations,
     77     after: (ProjectionGeneration, String),
     78 }
     79 
     80 /// An absent document is explicit corruption requiring caller reconciliation.
     81 #[derive(Clone, Eq, PartialEq)]
     82 pub struct ProjectionDocumentRecord {
     83     generation: ProjectionGeneration,
     84     key: String,
     85     document: Option<ProjectionDocument>,
     86 }
     87 impl ProjectionDocumentRecord {
     88     pub fn new(
     89         generation: ProjectionGeneration,
     90         document: ProjectionDocument,
     91     ) -> Result<Self, Error> {
     92         ProjectionGeneration::new(*generation.as_bytes())?;
     93         Ok(Self {
     94             generation,
     95             key: document.key().to_owned(),
     96             document: Some(document),
     97         })
     98     }
     99     pub fn corrupt(generation: ProjectionGeneration, key: String) -> Result<Self, Error> {
    100         ProjectionGeneration::new(*generation.as_bytes())?;
    101         if !valid_document_key(&key) {
    102             return Err(Error::CorruptProjectionDocument);
    103         }
    104         Ok(Self {
    105             generation,
    106             key,
    107             document: None,
    108         })
    109     }
    110     pub const fn generation(&self) -> ProjectionGeneration {
    111         self.generation
    112     }
    113     pub fn key(&self) -> &str {
    114         &self.key
    115     }
    116     pub const fn document(&self) -> Option<&ProjectionDocument> {
    117         self.document.as_ref()
    118     }
    119 }
    120 impl std::fmt::Debug for ProjectionDocumentRecord {
    121     fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
    122         formatter
    123             .debug_struct("ProjectionDocumentRecord")
    124             .field("generation", &self.generation)
    125             .field("key", &self.key)
    126             .field("corrupt", &self.document.is_none())
    127             .finish()
    128     }
    129 }
    130 
    131 #[derive(Clone, Debug, Eq, PartialEq)]
    132 pub struct ProjectionDocumentPage {
    133     records: Vec<ProjectionDocumentRecord>,
    134     next_cursor: Option<ProjectionDocumentCursor>,
    135 }
    136 impl ProjectionDocumentPage {
    137     /// Backends must bound bytes read, including invalid payloads, before constructing a page.
    138     pub fn new(
    139         query: &ProjectionDocumentQuery,
    140         records: Vec<ProjectionDocumentRecord>,
    141         has_more: bool,
    142     ) -> Result<Self, Error> {
    143         if records.len() > usize::from(query.limit) || (has_more && records.is_empty()) {
    144             return Err(Error::InvalidProjectionDocument);
    145         }
    146         let mut previous = query.after();
    147         let mut bytes = 0;
    148         for record in &records {
    149             let position = (record.generation, record.key());
    150             if !query.matches(position.0, position.1)
    151                 || previous.is_some_and(|last| position <= last)
    152             {
    153                 return Err(Error::InvalidProjectionDocument);
    154             }
    155             bytes += record
    156                 .document()
    157                 .map_or(0, |document| document.value().len());
    158             if bytes > PROJECTION_DOCUMENT_PAGE_BYTES_MAX {
    159                 return Err(Error::InvalidProjectionDocument);
    160             }
    161             previous = Some(position);
    162         }
    163         let next_cursor = if has_more {
    164             previous.map(|(generation, key)| ProjectionDocumentCursor {
    165                 projection_id: query.projection_id.clone(),
    166                 generations: query.generations,
    167                 after: (generation, key.to_owned()),
    168             })
    169         } else {
    170             None
    171         };
    172         Ok(Self {
    173             records,
    174             next_cursor,
    175         })
    176     }
    177     pub fn records(&self) -> &[ProjectionDocumentRecord] {
    178         &self.records
    179     }
    180     pub const fn next_cursor(&self) -> Option<&ProjectionDocumentCursor> {
    181         self.next_cursor.as_ref()
    182     }
    183     pub fn into_records(self) -> Vec<ProjectionDocumentRecord> {
    184         self.records
    185     }
    186 }