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 }