lib

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

tests.rs (9294B)


      1 use super::*;
      2 use crate::{OpenMode, OpenOptions, Paths};
      3 use radroots_storage::{
      4     ProjectionStore, event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId,
      5 };
      6 use tempfile::TempDir;
      7 
      8 async fn open(temp: &TempDir) -> SqliteStorage {
      9     SqliteStorage::open(
     10         OpenOptions::new(
     11             Paths::from_directory(temp.path()).unwrap(),
     12             OpenMode::Create,
     13         )
     14         .with_source_generation(SourceGeneration::new([9; 32]).unwrap(), 9)
     15         .unwrap(),
     16     )
     17     .await
     18     .unwrap()
     19 }
     20 fn generation(byte: u8) -> ProjectionGeneration {
     21     ProjectionGeneration::new([byte; 32]).unwrap()
     22 }
     23 fn query(selection: ProjectionDocumentGenerations, limit: u16) -> ProjectionDocumentQuery {
     24     ProjectionDocumentQuery::new(ProjectionId::parse("fixture").unwrap(), selection, limit).unwrap()
     25 }
     26 async fn put(store: &dyn ProjectionStore, generation: u8, key: &str, value: Vec<u8>) {
     27     store
     28         .put_projection_document(
     29             ProjectionId::parse("fixture").unwrap(),
     30             self::generation(generation),
     31             ProjectionDocument::new(key.into(), value).unwrap(),
     32         )
     33         .await
     34         .unwrap();
     35 }
     36 
     37 #[tokio::test]
     38 async fn inventory_thousand_records_matches_memory_across_generations_updates_and_reopen() {
     39     let temp = TempDir::new().unwrap();
     40     let store = open(&temp).await;
     41     let memory = MemoryStorage::new(SourceGeneration::new([9; 32]).unwrap());
     42     for i in (0..1000).rev() {
     43         for backend in [&store as &dyn ProjectionStore, &memory] {
     44             put(backend, 1 + (i % 2) as u8, &format!("key.{i:04}"), vec![1]).await;
     45         }
     46     }
     47     store
     48         .put_projection_document(
     49             ProjectionId::parse("foreign").unwrap(),
     50             generation(1),
     51             ProjectionDocument::new("key.0000".into(), vec![3]).unwrap(),
     52         )
     53         .await
     54         .unwrap();
     55     let q = query(ProjectionDocumentGenerations::All, 37);
     56     let first = store.query_projection_documents(q.clone()).await.unwrap();
     57     assert_eq!(
     58         first,
     59         memory.query_projection_documents(q.clone()).await.unwrap()
     60     );
     61     for backend in [&store as &dyn ProjectionStore, &memory] {
     62         put(backend, 1, "key.0000", vec![2]).await;
     63         put(backend, 2, "key.0999", vec![2]).await;
     64     }
     65     store.close().await.unwrap();
     66     let store = open(&temp).await;
     67     let mut count = first.records().len();
     68     let mut cursor = first.next_cursor().cloned();
     69     while let Some(current) = cursor {
     70         let query = q.clone().with_cursor(&current).unwrap();
     71         let page = store
     72             .query_projection_documents(query.clone())
     73             .await
     74             .unwrap();
     75         assert_eq!(
     76             page,
     77             memory.query_projection_documents(query).await.unwrap()
     78         );
     79         count += page.records().len();
     80         cursor = page.next_cursor().cloned();
     81     }
     82     assert_eq!(count, 1000);
     83     for byte in [1, 2, 3] {
     84         let q = query(ProjectionDocumentGenerations::Exact(generation(byte)), 256);
     85         let page = store.query_projection_documents(q.clone()).await.unwrap();
     86         assert_eq!(
     87             page,
     88             memory.query_projection_documents(q.clone()).await.unwrap()
     89         );
     90         assert!(
     91             page.records()
     92                 .iter()
     93                 .all(|r| r.generation() == generation(byte))
     94         );
     95         if byte == 3 {
     96             assert!(page.records().is_empty());
     97         } else {
     98             let last = store
     99                 .query_projection_documents(q.with_cursor(page.next_cursor().unwrap()).unwrap())
    100                 .await
    101                 .unwrap();
    102             assert_eq!(last.records().len() + page.records().len(), 500);
    103             assert!(last.next_cursor().is_none());
    104         }
    105     }
    106     store.close().await.unwrap();
    107     assert_eq!(
    108         store.query_projection_documents(q).await,
    109         Err(Error::BackendUnavailable)
    110     );
    111 }
    112 
    113 #[tokio::test]
    114 async fn inventory_readonly_preserves_identical_keys_and_byte_order() {
    115     let temp = TempDir::new().unwrap();
    116     let store = open(&temp).await;
    117     for byte in [1, 2] {
    118         for key in ["same", "z", "é"] {
    119             put(&store, byte, key, vec![byte]).await;
    120         }
    121     }
    122     store.close().await.unwrap();
    123     let store = SqliteStorage::open(OpenOptions::new(
    124         Paths::from_directory(temp.path()).unwrap(),
    125         OpenMode::ReadOnly,
    126     ))
    127     .await
    128     .unwrap();
    129     let q = query(ProjectionDocumentGenerations::All, 1);
    130     let mut current = q.clone();
    131     let mut found = Vec::new();
    132     loop {
    133         let page = store.query_projection_documents(current).await.unwrap();
    134         found.extend(
    135             page.records()
    136                 .iter()
    137                 .map(|r| (r.generation(), r.key().to_owned())),
    138         );
    139         let Some(cursor) = page.next_cursor() else {
    140             break;
    141         };
    142         current = q.clone().with_cursor(cursor).unwrap();
    143     }
    144     assert_eq!(
    145         found,
    146         [1, 2]
    147             .into_iter()
    148             .flat_map(|byte| ["same", "z", "é"].map(|key| (generation(byte), key.to_owned())))
    149             .collect::<Vec<_>>()
    150     );
    151     let exact = query(ProjectionDocumentGenerations::Exact(generation(1)), 1);
    152     let page = store.query_projection_documents(exact).await.unwrap();
    153     assert!(
    154         query(ProjectionDocumentGenerations::Exact(generation(2)), 1)
    155             .with_cursor(page.next_cursor().unwrap())
    156             .is_err()
    157     );
    158     store.close().await.unwrap();
    159 }
    160 
    161 #[tokio::test]
    162 async fn inventory_bounds_bytes_including_corrupt_payloads_without_losing_continuation() {
    163     let temp = TempDir::new().unwrap();
    164     let store = open(&temp).await;
    165     put(&store, 1, "a", vec![7; PROJECTION_DOCUMENT_PAGE_BYTES_MAX]).await;
    166     put(&store, 1, "b", vec![8]).await;
    167     let q = query(ProjectionDocumentGenerations::All, 256);
    168     for corrupt in [false, true] {
    169         if corrupt {
    170             sqlx::query("UPDATE radroots_runtime_projection_documents SET value_sha256 = zeroblob(32) WHERE document_key = 'a'")
    171                 .execute(store.pool()).await.unwrap();
    172         }
    173         let first = store.query_projection_documents(q.clone()).await.unwrap();
    174         assert_eq!(first.records().len(), 1);
    175         assert_eq!(first.records()[0].document().is_none(), corrupt);
    176         let next = store
    177             .query_projection_documents(
    178                 q.clone().with_cursor(first.next_cursor().unwrap()).unwrap(),
    179             )
    180             .await
    181             .unwrap();
    182         assert_eq!(next.records()[0].key(), "b");
    183         assert!(next.next_cursor().is_none());
    184     }
    185     store.close().await.unwrap();
    186 }
    187 
    188 #[tokio::test]
    189 async fn inventory_exposes_bounded_corruption_and_rejects_unpageable_keys() {
    190     let temp = TempDir::new().unwrap();
    191     let store = open(&temp).await;
    192     // One held connection makes disabled-check injection local to this fixture.
    193     let mut connection = store.pool().acquire().await.unwrap();
    194     sqlx::query("PRAGMA ignore_check_constraints = ON")
    195         .execute(&mut *connection)
    196         .await
    197         .unwrap();
    198     for (key, size, digest) in [("a", 0, 32), ("b", 16777217, 32), ("c", 1, 33)] {
    199         sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', ?, ?, zeroblob(?), zeroblob(?))")
    200             .bind(generation(1).as_bytes().as_slice()).bind(key).bind(size).bind(digest)
    201             .execute(&mut *connection).await.unwrap();
    202     }
    203     sqlx::query("PRAGMA ignore_check_constraints = OFF")
    204         .execute(&mut *connection)
    205         .await
    206         .unwrap();
    207     drop(connection);
    208     put(&store, 1, "d", vec![1]).await;
    209     let q = query(ProjectionDocumentGenerations::All, 2);
    210     let page = store.query_projection_documents(q.clone()).await.unwrap();
    211     assert!(page.records().iter().all(|r| r.document().is_none()));
    212     let next = store
    213         .query_projection_documents(q.with_cursor(page.next_cursor().unwrap()).unwrap())
    214         .await
    215         .unwrap();
    216     assert!(next.records()[0].document().is_none());
    217     assert!(next.records()[1].document().is_some());
    218     assert!(next.next_cursor().is_none());
    219     for key in ["\n".to_owned(), "é".repeat(512)] {
    220         sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', ?, ?, x'01', zeroblob(32))")
    221             .bind(generation(2).as_bytes().as_slice()).bind(&key).execute(store.pool()).await.unwrap();
    222         assert_eq!(
    223             store
    224                 .query_projection_documents(query(
    225                     ProjectionDocumentGenerations::Exact(generation(2)),
    226                     1
    227                 ))
    228                 .await,
    229             Err(Error::CorruptProjectionDocument)
    230         );
    231         sqlx::query("DELETE FROM radroots_runtime_projection_documents WHERE generation = ?")
    232             .bind(generation(2).as_bytes().as_slice())
    233             .execute(store.pool())
    234             .await
    235             .unwrap();
    236     }
    237     sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', zeroblob(32), 'zero', x'01', zeroblob(32))")
    238         .execute(store.pool()).await.unwrap();
    239     assert_eq!(
    240         store
    241             .query_projection_documents(query(ProjectionDocumentGenerations::All, 1))
    242             .await,
    243         Err(Error::CorruptProjectionDocument)
    244     );
    245     store.close().await.unwrap();
    246 }