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 }