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(¤t).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 }