lib

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

authored_draft_query.rs (4721B)


      1 use super::{SqliteStorage, bounded_row, decode_row, map_backend};
      2 use radroots_storage::{
      3     Error,
      4     authored_draft::AuthoredDraftRevision,
      5     authored_draft_query::{
      6         AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES,
      7         AuthoredDraftPage, AuthoredDraftQuery, AuthoredDraftQueryRecord,
      8     },
      9 };
     10 use sqlx::Row;
     11 
     12 pub(super) async fn page(
     13     store: &SqliteStorage,
     14     query: AuthoredDraftQuery,
     15 ) -> Result<AuthoredDraftPage, Error> {
     16     let mut transaction = store.pool().begin().await.map_err(map_backend)?;
     17     let result = read_page(&mut transaction, &query).await;
     18     // The read snapshot is always released before returning any continuation.
     19     let rollback = transaction.rollback().await.map_err(map_backend);
     20     match result {
     21         Ok(page) => {
     22             rollback?;
     23             Ok(page)
     24         }
     25         Err(error) => Err(error),
     26     }
     27 }
     28 
     29 async fn read_page(
     30     transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
     31     query: &AuthoredDraftQuery,
     32 ) -> Result<AuthoredDraftPage, Error> {
     33     let after = query.after().map(|value| value.to_vec());
     34     let rows = sqlx::query(
     35         "SELECT CASE WHEN typeof(revisions.draft_id) = 'blob' AND octet_length(revisions.draft_id) = 16 THEN revisions.draft_id END AS draft_id,
     36                 CASE WHEN typeof(revisions.revision) = 'integer' THEN revisions.revision END AS revision,
     37                 octet_length(revisions.snapshot) AS snapshot_bytes,
     38                 CASE WHEN typeof(revisions.payload_schema) = 'text' AND octet_length(revisions.payload_schema) <= 128
     39                      THEN revisions.payload_schema = '' ELSE 1 END AS unknown_schema
     40          FROM radroots_runtime_authored_draft_revisions AS revisions
     41          WHERE revisions.author = ?
     42            AND (? OR (revisions.payload_schema = ? AND (? OR revisions.payload_scope IS ?))
     43                 OR revisions.payload_schema = '')
     44            AND (? IS NULL OR revisions.draft_id > ?)
     45            AND revisions.revision = (SELECT MAX(head.revision)
     46              FROM radroots_runtime_authored_draft_revisions AS head WHERE head.draft_id = revisions.draft_id)
     47          ORDER BY revisions.draft_id LIMIT ?"
     48     ).bind(query.author().as_slice()).bind(query.payload_schema().is_none()).bind(query.payload_schema())
     49         .bind(query.is_author_wide())
     50         .bind(query.scope().map(|value| value.as_bytes().to_vec()))
     51         .bind(&after).bind(after).bind(i64::from(query.limit()) + 1)
     52         .fetch_all(&mut **transaction).await.map_err(map_backend)?;
     53     let mut records = Vec::new();
     54     let mut snapshot_bytes = 0usize;
     55     let mut payload_bytes = 0usize;
     56     let mut has_more = rows.len() > usize::from(query.limit());
     57     for row in rows.iter().take(usize::from(query.limit())) {
     58         // Keys and positive revisions are enforced by the STRICT table. A zero
     59         // key is still a valid scan/repair position for a corrupt domain ID.
     60         let key: [u8; 16] = row
     61             .try_get::<Vec<u8>, _>("draft_id")
     62             .map_err(map_backend)?
     63             .try_into()
     64             .map_err(|_| Error::CorruptAuthoredDraft)?;
     65         let revision = row.try_get::<i64, _>("revision").map_err(map_backend)?;
     66         let revision = AuthoredDraftRevision::new(
     67             u64::try_from(revision).map_err(|_| Error::CorruptAuthoredDraft)?,
     68         )?;
     69         let size = row
     70             .try_get::<i64, _>("snapshot_bytes")
     71             .map_err(map_backend)?;
     72         let unknown = row
     73             .try_get::<bool, _>("unknown_schema")
     74             .map_err(map_backend)?;
     75         let corrupt = AuthoredDraftQueryRecord::Corrupt {
     76             draft_key: key,
     77             revision,
     78         };
     79         if unknown || size <= 0 || size as u64 > AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES as u64 {
     80             records.push(corrupt);
     81             continue;
     82         }
     83         let size = size as usize;
     84         if snapshot_bytes + size > AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES {
     85             has_more = true;
     86             break;
     87         }
     88         snapshot_bytes += size;
     89         let row = bounded_row::load(&mut **transaction, &key, Some(revision))
     90             .await?
     91             .ok_or(Error::CorruptAuthoredDraft)?;
     92         match decode_row(&row) {
     93             Ok(draft) if query.matches(&draft) => {
     94                 if payload_bytes + draft.payload().len() > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES {
     95                     has_more = true;
     96                     break;
     97                 }
     98                 payload_bytes += draft.payload().len();
     99                 records.push(AuthoredDraftQueryRecord::Draft(draft));
    100             }
    101             Ok(_) | Err(_) => records.push(corrupt),
    102         }
    103     }
    104     AuthoredDraftPage::new(query, records, has_more)
    105 }