authored_draft_query.rs (13918B)
1 //! Bounded, independently scoped pages of current opaque draft revisions. 2 use crate::{ 3 Error, 4 authored_draft::{ 5 AUTHORED_DRAFT_QUERY_LIMIT_MAX, AUTHORED_DRAFT_SCHEMA_MAX_BYTES, AuthoredDraft, 6 AuthoredDraftId, AuthoredDraftRevision, 7 }, 8 }; 9 10 /// Maximum decoded payload bytes retained by one query page. 11 pub const AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES: usize = 4 * 1024 * 1024; 12 /// Maximum serialized snapshot bytes read by one native query page. 13 pub const AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024; 14 pub const AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION: u16 = 1; 15 /// Explicit author/schema traversal has a distinct continuation authority. 16 pub const AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION: u16 = 2; 17 /// Explicit author-wide traversal of every payload schema has separate authority. 18 pub const AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION: u16 = 3; 19 20 /// An opaque, stable application-selected scope digest; never a credential. 21 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 22 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 23 #[cfg_attr(feature = "serde", serde(try_from = "[u8; 32]", into = "[u8; 32]"))] 24 pub struct AuthoredDraftScope([u8; 32]); 25 26 impl AuthoredDraftScope { 27 pub fn new(bytes: [u8; 32]) -> Result<Self, Error> { 28 if bytes.iter().all(|byte| *byte == 0) { 29 return Err(Error::InvalidAuthoredDraft); 30 } 31 Ok(Self(bytes)) 32 } 33 pub const fn as_bytes(&self) -> &[u8; 32] { 34 &self.0 35 } 36 } 37 impl TryFrom<[u8; 32]> for AuthoredDraftScope { 38 type Error = Error; 39 fn try_from(value: [u8; 32]) -> Result<Self, Error> { 40 Self::new(value) 41 } 42 } 43 impl From<AuthoredDraftScope> for [u8; 32] { 44 fn from(value: AuthoredDraftScope) -> Self { 45 value.0 46 } 47 } 48 49 /// Scope is selected independently on every call, including continuations. 50 #[derive(Clone, Debug, Eq, PartialEq)] 51 pub struct AuthoredDraftQuery { 52 author: [u8; 32], 53 payload_schema: Option<String>, 54 scope: Option<AuthoredDraftScope>, 55 author_wide: bool, 56 limit: u16, 57 after: Option<[u8; 16]>, 58 } 59 impl AuthoredDraftQuery { 60 pub fn new( 61 author: [u8; 32], 62 payload_schema: impl AsRef<str>, 63 scope: Option<AuthoredDraftScope>, 64 limit: u16, 65 ) -> Result<Self, Error> { 66 let schema = payload_schema.as_ref(); 67 if author.iter().all(|byte| *byte == 0) 68 || schema.is_empty() 69 || schema.len() > AUTHORED_DRAFT_SCHEMA_MAX_BYTES 70 || schema != schema.trim() 71 || schema.chars().any(char::is_control) 72 || limit == 0 73 || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX 74 { 75 return Err(Error::InvalidAuthoredDraft); 76 } 77 Ok(Self { 78 author, 79 payload_schema: Some(schema.to_owned()), 80 scope, 81 author_wide: false, 82 limit, 83 after: None, 84 }) 85 } 86 /// Selects this author's exact payload schema across all application scopes. 87 /// This does not broaden `new(..., None, ...)`, which remains unscoped only. 88 pub fn for_author( 89 author: [u8; 32], 90 payload_schema: impl AsRef<str>, 91 limit: u16, 92 ) -> Result<Self, Error> { 93 let mut query = Self::new(author, payload_schema, None, limit)?; 94 query.author_wide = true; 95 Ok(query) 96 } 97 98 /// Selects every current payload schema for this author across all scopes. 99 /// Unknown application schemas remain opaque records, never absent owners. 100 pub fn for_author_all_schemas(author: [u8; 32], limit: u16) -> Result<Self, Error> { 101 if author.iter().all(|byte| *byte == 0) 102 || limit == 0 103 || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX 104 { 105 return Err(Error::InvalidAuthoredDraft); 106 } 107 Ok(Self { 108 author, 109 payload_schema: None, 110 scope: None, 111 author_wide: true, 112 limit, 113 after: None, 114 }) 115 } 116 117 /// Whether a named author-wide constructor explicitly omitted scope filtering. 118 pub const fn is_author_wide(&self) -> bool { 119 self.author_wide 120 } 121 122 const fn cursor_version(&self) -> u16 { 123 if self.payload_schema.is_none() { 124 AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION 125 } else if self.author_wide { 126 AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION 127 } else { 128 AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION 129 } 130 } 131 132 pub fn with_cursor(mut self, cursor: &AuthoredDraftCursor) -> Result<Self, Error> { 133 if cursor.schema_version != self.cursor_version() 134 || cursor.author != self.author 135 || cursor.payload_schema != self.payload_schema 136 || cursor.scope != self.scope 137 { 138 return Err(Error::InvalidAuthoredDraft); 139 } 140 self.after = Some(cursor.after_id); 141 Ok(self) 142 } 143 pub const fn author(&self) -> &[u8; 32] { 144 &self.author 145 } 146 /// Exact schema, or explicit all-schema selection. An empty string is never a wildcard. 147 pub fn payload_schema(&self) -> Option<&str> { 148 self.payload_schema.as_deref() 149 } 150 /// Exact optional scope for ordinary queries. Author-wide queries return 151 /// `None`; backends must also honor `is_author_wide` when selecting rows. 152 pub const fn scope(&self) -> Option<AuthoredDraftScope> { 153 self.scope 154 } 155 pub const fn limit(&self) -> u16 { 156 self.limit 157 } 158 pub const fn after(&self) -> Option<[u8; 16]> { 159 self.after 160 } 161 pub fn matches(&self, draft: &AuthoredDraft) -> bool { 162 draft.author() == &self.author 163 && self 164 .payload_schema() 165 .is_none_or(|schema| draft.payload_schema() == schema) 166 && (self.author_wide || draft.scope() == self.scope) 167 } 168 pub fn cursor_after(&self, after_id: [u8; 16]) -> AuthoredDraftCursor { 169 AuthoredDraftCursor { 170 schema_version: self.cursor_version(), 171 author: self.author, 172 payload_schema: self.payload_schema.clone(), 173 scope: self.scope, 174 after_id, 175 } 176 } 177 } 178 179 /// Stable ID ordering does not move when a draft receives another revision. 180 /// This is a bounded scan, not an immutable cross-page database snapshot. 181 #[derive(Clone, Debug, Eq, PartialEq)] 182 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 183 #[cfg_attr(feature = "serde", serde(try_from = "CursorWire", into = "CursorWire"))] 184 pub struct AuthoredDraftCursor { 185 schema_version: u16, 186 author: [u8; 32], 187 payload_schema: Option<String>, 188 scope: Option<AuthoredDraftScope>, 189 after_id: [u8; 16], 190 } 191 #[cfg(feature = "serde")] 192 #[derive(serde::Serialize, serde::Deserialize)] 193 #[serde(untagged)] 194 enum CursorWire { 195 Exact(ExactCursorWire), 196 Author(AuthorCursorWire), 197 AllSchemas(AllSchemasCursorWire), 198 } 199 200 #[cfg(feature = "serde")] 201 #[derive(serde::Serialize, serde::Deserialize)] 202 #[serde(deny_unknown_fields)] 203 struct ExactCursorWire { 204 schema_version: u16, 205 author: [u8; 32], 206 payload_schema: String, 207 scope: Option<AuthoredDraftScope>, 208 after_id: [u8; 16], 209 } 210 #[cfg(feature = "serde")] 211 #[derive(serde::Serialize, serde::Deserialize)] 212 #[serde(deny_unknown_fields)] 213 struct AuthorCursorWire { 214 schema_version: u16, 215 author: [u8; 32], 216 payload_schema: String, 217 scope: (), 218 after_id: [u8; 16], 219 selection: CursorSelection, 220 } 221 222 #[cfg(feature = "serde")] 223 #[derive(serde::Serialize, serde::Deserialize)] 224 #[serde(deny_unknown_fields)] 225 struct AllSchemasCursorWire { 226 schema_version: u16, 227 author: [u8; 32], 228 payload_schema: (), 229 scope: (), 230 after_id: [u8; 16], 231 selection: CursorSelection, 232 } 233 #[cfg(feature = "serde")] 234 #[derive(serde::Serialize, serde::Deserialize)] 235 enum CursorSelection { 236 #[serde(rename = "author_schema")] 237 AuthorSchema, 238 #[serde(rename = "author_all_schemas")] 239 AuthorAllSchemas, 240 } 241 #[cfg(feature = "serde")] 242 impl TryFrom<CursorWire> for AuthoredDraftCursor { 243 type Error = Error; 244 fn try_from(value: CursorWire) -> Result<Self, Error> { 245 let (query, after) = match value { 246 CursorWire::Exact(value) 247 if value.schema_version == AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION => 248 { 249 ( 250 AuthoredDraftQuery::new(value.author, value.payload_schema, value.scope, 1)?, 251 value.after_id, 252 ) 253 } 254 CursorWire::Author(value) 255 if value.schema_version == AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION 256 && matches!(value.selection, CursorSelection::AuthorSchema) => 257 { 258 ( 259 AuthoredDraftQuery::for_author(value.author, value.payload_schema, 1)?, 260 value.after_id, 261 ) 262 } 263 CursorWire::AllSchemas(value) 264 if value.schema_version == AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION 265 && matches!(value.selection, CursorSelection::AuthorAllSchemas) => 266 { 267 ( 268 AuthoredDraftQuery::for_author_all_schemas(value.author, 1)?, 269 value.after_id, 270 ) 271 } 272 _ => return Err(Error::InvalidAuthoredDraft), 273 }; 274 Ok(query.cursor_after(after)) 275 } 276 } 277 #[cfg(feature = "serde")] 278 impl From<AuthoredDraftCursor> for CursorWire { 279 fn from(value: AuthoredDraftCursor) -> Self { 280 let Some(payload_schema) = value.payload_schema else { 281 return Self::AllSchemas(AllSchemasCursorWire { 282 schema_version: value.schema_version, 283 author: value.author, 284 payload_schema: (), 285 scope: (), 286 after_id: value.after_id, 287 selection: CursorSelection::AuthorAllSchemas, 288 }); 289 }; 290 if value.schema_version == AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION { 291 return Self::Author(AuthorCursorWire { 292 schema_version: value.schema_version, 293 author: value.author, 294 payload_schema, 295 scope: (), 296 after_id: value.after_id, 297 selection: CursorSelection::AuthorSchema, 298 }); 299 } 300 Self::Exact(ExactCursorWire { 301 schema_version: value.schema_version, 302 author: value.author, 303 payload_schema, 304 scope: value.scope, 305 after_id: value.after_id, 306 }) 307 } 308 } 309 310 /// Corrupt rows expose only their local position, never untrusted payloads. 311 /// A raw position can identify a malformed all-zero draft ID for repair. 312 #[derive(Clone, Eq, PartialEq)] 313 pub enum AuthoredDraftQueryRecord { 314 Draft(AuthoredDraft), 315 Corrupt { 316 draft_key: [u8; 16], 317 revision: AuthoredDraftRevision, 318 }, 319 } 320 impl core::fmt::Debug for AuthoredDraftQueryRecord { 321 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { 322 f.debug_struct(match self { 323 Self::Draft(_) => "Draft", 324 Self::Corrupt { .. } => "Corrupt", 325 }) 326 .field("draft_key", &self.draft_key()) 327 .field("revision", &self.revision()) 328 .finish_non_exhaustive() 329 } 330 } 331 332 impl AuthoredDraftQueryRecord { 333 pub fn draft_key(&self) -> [u8; 16] { 334 match self { 335 Self::Draft(draft) => *draft.draft_id().as_bytes(), 336 Self::Corrupt { draft_key, .. } => *draft_key, 337 } 338 } 339 pub const fn revision(&self) -> AuthoredDraftRevision { 340 match self { 341 Self::Draft(draft) => draft.revision(), 342 Self::Corrupt { revision, .. } => *revision, 343 } 344 } 345 pub fn draft_id(&self) -> Result<AuthoredDraftId, Error> { 346 AuthoredDraftId::new(self.draft_key()) 347 } 348 } 349 350 #[derive(Clone, Debug, Eq, PartialEq)] 351 pub struct AuthoredDraftPage { 352 records: Vec<AuthoredDraftQueryRecord>, 353 next_cursor: Option<AuthoredDraftCursor>, 354 } 355 impl AuthoredDraftPage { 356 /// Backends return an ordered, bounded page after applying independent scope. 357 pub fn new( 358 query: &AuthoredDraftQuery, 359 records: Vec<AuthoredDraftQueryRecord>, 360 has_more: bool, 361 ) -> Result<Self, Error> { 362 if records.len() > usize::from(query.limit) || (has_more && records.is_empty()) { 363 return Err(Error::InvalidAuthoredDraft); 364 } 365 let mut previous = query.after; 366 let mut payload_bytes = 0usize; 367 for record in &records { 368 let key = record.draft_key(); 369 if previous.is_some_and(|previous| key <= previous) { 370 return Err(Error::InvalidAuthoredDraft); 371 } 372 previous = Some(key); 373 if let AuthoredDraftQueryRecord::Draft(draft) = record { 374 if !query.matches(draft) { 375 return Err(Error::InvalidAuthoredDraft); 376 } 377 draft.validate()?; 378 payload_bytes = payload_bytes 379 .checked_add(draft.payload().len()) 380 .ok_or(Error::InvalidAuthoredDraft)?; 381 if payload_bytes > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES { 382 return Err(Error::InvalidAuthoredDraft); 383 } 384 } 385 } 386 let next_cursor = if has_more { 387 previous.map(|key| query.cursor_after(key)) 388 } else { 389 None 390 }; 391 Ok(Self { 392 records, 393 next_cursor, 394 }) 395 } 396 pub fn records(&self) -> &[AuthoredDraftQueryRecord] { 397 &self.records 398 } 399 pub const fn next_cursor(&self) -> Option<&AuthoredDraftCursor> { 400 self.next_cursor.as_ref() 401 } 402 pub fn into_records(self) -> Vec<AuthoredDraftQueryRecord> { 403 self.records 404 } 405 }