lib

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

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;