document_query.rs (4712B)
1 use super::{SqliteStorage, map_backend}; 2 use radroots_storage::{ 3 Error, 4 projection::{ 5 ProjectionDocument, ProjectionGeneration, 6 document_query::{ 7 PROJECTION_DOCUMENT_PAGE_BYTES_MAX, ProjectionDocumentGenerations, 8 ProjectionDocumentPage, ProjectionDocumentQuery, ProjectionDocumentRecord, 9 }, 10 }, 11 }; 12 use sqlx::Row; 13 14 pub(super) async fn page( 15 store: &SqliteStorage, 16 query: ProjectionDocumentQuery, 17 ) -> Result<ProjectionDocumentPage, Error> { 18 let mut transaction = store.pool().begin().await.map_err(map_backend)?; 19 let result = read_page(&mut transaction, &query).await; 20 let rollback = transaction.rollback().await.map_err(map_backend); 21 match result { 22 Ok(page) => { 23 rollback?; 24 Ok(page) 25 } 26 Err(error) => Err(error), 27 } 28 } 29 30 async fn read_page( 31 transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, 32 query: &ProjectionDocumentQuery, 33 ) -> Result<ProjectionDocumentPage, Error> { 34 let generation = match query.generations() { 35 ProjectionDocumentGenerations::All => None, 36 ProjectionDocumentGenerations::Exact(value) => Some(value.as_bytes().to_vec()), 37 }; 38 let after_generation = query 39 .after() 40 .map(|(generation, _)| generation.as_bytes().to_vec()); 41 let after_key = query.after().map(|(_, key)| key); 42 // CASE bounds every variable-width metadata value before SQLx allocates it. 43 // Payloads are loaded separately, only after charging this page's byte budget. 44 let rows = sqlx::query( 45 "SELECT CASE WHEN length(generation) = 32 THEN generation END AS generation, 46 CASE WHEN length(CAST(document_key AS BLOB)) BETWEEN 1 AND 512 THEN document_key END AS document_key, 47 CASE WHEN length(value_sha256) = 32 THEN value_sha256 END AS value_sha256, 48 length(value) AS value_bytes 49 FROM radroots_runtime_projection_documents AS documents 50 WHERE projection_id = ? AND (? IS NULL OR generation = ?) 51 AND (? IS NULL OR (generation, document_key) > (?, ?)) 52 ORDER BY documents.generation, documents.document_key LIMIT ?" 53 ).bind(query.projection_id().as_str()).bind(&generation).bind(generation) 54 .bind(&after_generation).bind(&after_generation).bind(after_key) 55 .bind(i64::from(query.limit()) + 1) 56 .fetch_all(&mut **transaction).await.map_err(map_backend)?; 57 let mut has_more = rows.len() > usize::from(query.limit()); 58 let mut records = Vec::new(); 59 let mut bytes = 0; 60 for row in rows.iter().take(usize::from(query.limit())) { 61 let generation: Vec<u8> = row 62 .try_get("generation") 63 .map_err(|_| Error::CorruptProjectionDocument)?; 64 let generation = ProjectionGeneration::new( 65 generation 66 .try_into() 67 .map_err(|_| Error::CorruptProjectionDocument)?, 68 ) 69 .map_err(|_| Error::CorruptProjectionDocument)?; 70 let key: String = row 71 .try_get("document_key") 72 .map_err(|_| Error::CorruptProjectionDocument)?; 73 let corrupt = ProjectionDocumentRecord::corrupt(generation, key.clone())?; 74 let digest: Option<Vec<u8>> = row 75 .try_get("value_sha256") 76 .map_err(|_| Error::CorruptProjectionDocument)?; 77 let size: i64 = row 78 .try_get("value_bytes") 79 .map_err(|_| Error::CorruptProjectionDocument)?; 80 let Some(digest) = digest else { 81 records.push(corrupt); 82 continue; 83 }; 84 let digest: [u8; 32] = digest 85 .try_into() 86 .map_err(|_| Error::CorruptProjectionDocument)?; 87 if size <= 0 || size > PROJECTION_DOCUMENT_PAGE_BYTES_MAX as i64 { 88 records.push(corrupt); 89 continue; 90 } 91 let size = size as usize; 92 if bytes + size > PROJECTION_DOCUMENT_PAGE_BYTES_MAX { 93 has_more = true; 94 break; 95 } 96 bytes += size; 97 let value: Vec<u8> = sqlx::query_scalar( 98 "SELECT value FROM radroots_runtime_projection_documents WHERE projection_id = ? AND generation = ? AND document_key = ?" 99 ).bind(query.projection_id().as_str()).bind(generation.as_bytes().as_slice()).bind(&key) 100 .fetch_one(&mut **transaction).await.map_err(map_backend)?; 101 records.push( 102 match ProjectionDocument::from_stored_parts(key, value, digest) { 103 Ok(document) => ProjectionDocumentRecord::new(generation, document)?, 104 Err(_) => corrupt, 105 }, 106 ); 107 } 108 ProjectionDocumentPage::new(query, records, has_more) 109 } 110 111 #[cfg(test)] 112 mod tests;