authored_draft_query_tests.rs (12956B)
1 use super::{SqliteStorage, tests::open_store}; 2 use radroots_storage::{ 3 Error, 4 authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, 5 authored_draft_query::{AuthoredDraftQuery, AuthoredDraftQueryRecord, AuthoredDraftScope}, 6 }; 7 use tempfile::TempDir; 8 9 pub(super) fn draft( 10 id: u8, 11 schema: &str, 12 scope: Option<AuthoredDraftScope>, 13 payload: Vec<u8>, 14 ) -> AuthoredDraft { 15 let draft = AuthoredDraft::initial( 16 AuthoredDraftId::new([id; 16]).unwrap(), 17 [7; 32], 18 schema, 19 payload, 20 AuthoredDraftStage::Draft, 21 None, 22 10, 23 ) 24 .unwrap(); 25 match scope { 26 Some(scope) => draft.with_scope(scope).unwrap(), 27 None => draft, 28 } 29 } 30 fn query(scope: Option<AuthoredDraftScope>, limit: u16) -> AuthoredDraftQuery { 31 AuthoredDraftQuery::new([7; 32], "fixture.composer.v1", scope, limit).unwrap() 32 } 33 34 #[tokio::test] 35 async fn author_wide_sqlite_pages_preserve_scope_isolation_corruption_and_restart() { 36 let temp = TempDir::new().unwrap(); 37 let store = open_store(&temp).await; 38 let scope = AuthoredDraftScope::new([9; 32]).unwrap(); 39 for id in 1_u128..=1000 { 40 let value = AuthoredDraft::initial( 41 AuthoredDraftId::new(id.to_be_bytes()).unwrap(), 42 [7; 32], 43 "fixture.composer.v1", 44 vec![1], 45 AuthoredDraftStage::Draft, 46 None, 47 10, 48 ) 49 .unwrap(); 50 let value = if id % 2 == 0 { 51 value.with_scope(scope).unwrap() 52 } else { 53 value 54 }; 55 store.append_authored_draft(value, None).await.unwrap(); 56 } 57 corrupt(&store, 0, 7, "fixture.composer.v1", Some(scope)).await; 58 corrupt(&store, 2, 8, "fixture.composer.v1", Some(scope)).await; 59 corrupt(&store, 3, 7, "fixture.other.v1", Some(scope)).await; 60 let q = AuthoredDraftQuery::for_author([7; 32], "fixture.composer.v1", 37).unwrap(); 61 let first = store.query_authored_drafts(q.clone()).await.unwrap(); 62 assert_eq!(first.records().len(), 37); 63 assert!(matches!( 64 first.records()[0], 65 AuthoredDraftQueryRecord::Corrupt { 66 draft_key: [0, ..], 67 .. 68 } 69 )); 70 for (index, record) in first.records().iter().enumerate().skip(1) { 71 assert_eq!(record.draft_key(), (index as u128).to_be_bytes()); 72 } 73 let mut cursor = first.next_cursor().cloned(); 74 // Revisions of both a visited and an unvisited ID do not move position. 75 for id in [1_u128, 999] { 76 let original = store 77 .authored_draft_head(AuthoredDraftId::new(id.to_be_bytes()).unwrap()) 78 .await 79 .unwrap() 80 .unwrap(); 81 let next = original 82 .successor(vec![2], AuthoredDraftStage::Draft, None, 11) 83 .unwrap(); 84 store 85 .append_authored_draft(next, Some(original.revision())) 86 .await 87 .unwrap(); 88 } 89 store.close().await.unwrap(); 90 let store = open_store(&temp).await; 91 let mut expected = 37_u128; 92 while let Some(current) = cursor { 93 let bytes = serde_json::to_vec(¤t).unwrap(); 94 let decoded = serde_json::from_slice(&bytes).unwrap(); 95 let page = store 96 .query_authored_drafts(q.clone().with_cursor(&decoded).unwrap()) 97 .await 98 .unwrap(); 99 assert!(page.records().len() <= 37); 100 for record in page.records() { 101 assert_eq!(record.draft_key(), expected.to_be_bytes()); 102 if expected == 999 { 103 assert_eq!(record.revision().get(), 2); 104 } 105 expected += 1; 106 } 107 cursor = page.next_cursor().cloned(); 108 } 109 assert_eq!(expected, 1001); 110 let unscoped = store.query_authored_drafts(query(None, 256)).await.unwrap(); 111 assert!( 112 unscoped 113 .records() 114 .iter() 115 .all(|record| u128::from_be_bytes(record.draft_key()) % 2 == 1) 116 ); 117 let fresh = store.query_authored_drafts(q).await.unwrap(); 118 assert_eq!(fresh.records()[1].revision().get(), 2); 119 store.close().await.unwrap(); 120 } 121 122 #[tokio::test] 123 async fn author_wide_sqlite_pages_preserve_snapshot_and_payload_budgets() { 124 for byte in [0, 99] { 125 let temp = TempDir::new().unwrap(); 126 let store = open_store(&temp).await; 127 let scope = AuthoredDraftScope::new([9; 32]).unwrap(); 128 for (id, scope) in [(1, None), (2, Some(scope))] { 129 store 130 .append_authored_draft( 131 draft( 132 id, 133 "fixture.composer.v1", 134 scope, 135 vec![byte; 3 * 1024 * 1024], 136 ), 137 None, 138 ) 139 .await 140 .unwrap(); 141 } 142 let q = AuthoredDraftQuery::for_author([7; 32], "fixture.composer.v1", 256).unwrap(); 143 let first = store.query_authored_drafts(q.clone()).await.unwrap(); 144 assert_eq!(first.records().len(), 1); 145 let next = store 146 .query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap()) 147 .await 148 .unwrap(); 149 assert_eq!(next.records().len(), 1); 150 assert_eq!(next.records()[0].draft_key(), [2; 16]); 151 assert!(matches!( 152 next.records()[0], 153 AuthoredDraftQueryRecord::Draft(_) 154 )); 155 assert!(next.next_cursor().is_none()); 156 store.close().await.unwrap(); 157 } 158 } 159 pub(super) async fn corrupt( 160 store: &SqliteStorage, 161 id: u8, 162 author: u8, 163 schema: &str, 164 scope: Option<AuthoredDraftScope>, 165 ) { 166 sqlx::query("INSERT INTO radroots_runtime_authored_draft_revisions 167 (draft_id, revision, author, stage, payload_sha256, created_at_unix_ms, updated_at_unix_ms, snapshot, payload_schema, payload_scope) 168 VALUES (?, 1, ?, 0, ?, 10, 10, ?, ?, ?)") 169 .bind([id; 16].as_slice()).bind([author; 32].as_slice()).bind([3; 32].as_slice()) 170 .bind(b"{malformed private payload".as_slice()).bind(schema).bind(scope.map(|scope| scope.as_bytes().to_vec())) 171 .execute(store.pool()).await.unwrap(); 172 } 173 174 #[tokio::test] 175 async fn scoped_sqlite_pages_isolate_corruption_foreign_schemas_and_contexts() { 176 let temp = TempDir::new().unwrap(); 177 let store = open_store(&temp).await; 178 let scope = AuthoredDraftScope::new([9; 32]).unwrap(); 179 let first = draft( 180 2, 181 "fixture.composer.v1", 182 Some(scope), 183 b"incomplete 0.".to_vec(), 184 ); 185 let second = draft(6, "fixture.composer.v1", Some(scope), b"2026-".to_vec()); 186 for value in [ 187 first.clone(), 188 second.clone(), 189 draft(1, "fixture.profile.v1", Some(scope), b"profile".to_vec()), 190 draft(3, "fixture.composer.v1", None, b"unscoped".to_vec()), 191 ] { 192 store.append_authored_draft(value, None).await.unwrap(); 193 } 194 corrupt(&store, 4, 7, "fixture.composer.v1", Some(scope)).await; 195 // Unknown historical schema yields only an author-bound corruption locator. 196 corrupt(&store, 5, 7, "", None).await; 197 corrupt(&store, 7, 8, "", None).await; 198 corrupt(&store, 8, 7, "fixture.profile.v1", Some(scope)).await; 199 let q = query(Some(scope), 1); 200 let first_page = store.query_authored_drafts(q.clone()).await.unwrap(); 201 assert_eq!( 202 first_page.records(), 203 [AuthoredDraftQueryRecord::Draft(first.clone())] 204 ); 205 let changed = first 206 .successor( 207 b"later saved edit".to_vec(), 208 AuthoredDraftStage::Draft, 209 None, 210 11, 211 ) 212 .unwrap(); 213 store 214 .append_authored_draft(changed, Some(first.revision())) 215 .await 216 .unwrap(); 217 let mut cursor = first_page.next_cursor().cloned(); 218 let mut records = first_page.into_records(); 219 while let Some(next) = cursor { 220 let page = store 221 .query_authored_drafts(q.clone().with_cursor(&next).unwrap()) 222 .await 223 .unwrap(); 224 cursor = page.next_cursor().cloned(); 225 records.extend(page.into_records()); 226 } 227 assert_eq!( 228 records 229 .iter() 230 .map(AuthoredDraftQueryRecord::draft_key) 231 .collect::<Vec<_>>(), 232 [[2; 16], [4; 16], [5; 16], [6; 16]] 233 ); 234 assert!(matches!( 235 records[1], 236 AuthoredDraftQueryRecord::Corrupt { .. } 237 )); 238 assert!(matches!( 239 records[2], 240 AuthoredDraftQueryRecord::Corrupt { .. } 241 )); 242 assert_eq!(records[3], AuthoredDraftQueryRecord::Draft(second)); 243 assert!(!format!("{records:?}").contains("private payload")); 244 assert_eq!( 245 store.authored_draft_heads([7; 32], 10).await, 246 Err(Error::CorruptAuthoredDraft) 247 ); 248 store.close().await.unwrap(); 249 assert!(store.query_authored_drafts(q.clone()).await.is_err()); 250 let reopened = open_store(&temp).await; 251 let page = reopened 252 .query_authored_drafts(query(Some(scope), 10)) 253 .await 254 .unwrap(); 255 assert_eq!(page.records().len(), 4); 256 assert!(page.next_cursor().is_none()); 257 reopened.close().await.unwrap(); 258 } 259 260 #[tokio::test] 261 async fn sqlite_snapshot_budget_keeps_large_valid_records_on_later_pages() { 262 // Small numeric bytes exhaust decoded payload first; larger ones exhaust serialized snapshots first. 263 for payload_byte in [0, 99] { 264 let temp = TempDir::new().unwrap(); 265 let store = open_store(&temp).await; 266 for id in [1, 2] { 267 store 268 .append_authored_draft( 269 draft( 270 id, 271 "fixture.composer.v1", 272 None, 273 vec![payload_byte; 3 * 1024 * 1024], 274 ), 275 None, 276 ) 277 .await 278 .unwrap(); 279 } 280 let q = query(None, 256); 281 let first = store.query_authored_drafts(q.clone()).await.unwrap(); 282 assert_eq!(first.records().len(), 1); 283 let next = store 284 .query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap()) 285 .await 286 .unwrap(); 287 assert_eq!(next.records()[0].draft_key(), [2; 16]); 288 assert!(matches!( 289 next.records()[0], 290 AuthoredDraftQueryRecord::Draft(_) 291 )); 292 assert!(next.next_cursor().is_none()); 293 store.close().await.unwrap(); 294 } 295 } 296 297 #[tokio::test] 298 async fn zero_domain_identity_is_an_isolated_repair_position() { 299 let temp = TempDir::new().unwrap(); 300 let store = open_store(&temp).await; 301 corrupt(&store, 0, 7, "", None).await; 302 let good = draft(1, "fixture.composer.v1", None, b"partial".to_vec()); 303 store 304 .append_authored_draft(good.clone(), None) 305 .await 306 .unwrap(); 307 let q = query(None, 1); 308 let page = store.query_authored_drafts(q.clone()).await.unwrap(); 309 assert!(page.records()[0].draft_id().is_err()); 310 let next = store 311 .query_authored_drafts(q.with_cursor(page.next_cursor().unwrap()).unwrap()) 312 .await 313 .unwrap(); 314 assert_eq!(next.records(), [AuthoredDraftQueryRecord::Draft(good)]); 315 store.close().await.unwrap(); 316 } 317 318 #[tokio::test] 319 async fn inconsistent_query_metadata_is_isolated_and_never_returns_foreign_payload() { 320 for corrupt_scope in [false, true] { 321 let temp = TempDir::new().unwrap(); 322 let store = open_store(&temp).await; 323 let value = draft(1, "fixture.composer.v1", None, b"private original".to_vec()); 324 store 325 .append_authored_draft(value.clone(), None) 326 .await 327 .unwrap(); 328 sqlx::query("DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard") 329 .execute(store.pool()) 330 .await 331 .unwrap(); 332 let (schema, scope) = if corrupt_scope { 333 ( 334 "fixture.composer.v1", 335 Some(AuthoredDraftScope::new([4; 32]).unwrap()), 336 ) 337 } else { 338 ("fixture.other.v1", None) 339 }; 340 sqlx::query("UPDATE radroots_runtime_authored_draft_revisions SET payload_schema = ?, payload_scope = ?") 341 .bind(schema).bind(scope.map(|v|v.as_bytes().to_vec())).execute(store.pool()).await.unwrap(); 342 assert_eq!( 343 store.authored_draft_head(value.draft_id()).await, 344 Err(Error::CorruptAuthoredDraft) 345 ); 346 let page = store 347 .query_authored_drafts(AuthoredDraftQuery::new([7; 32], schema, scope, 1).unwrap()) 348 .await 349 .unwrap(); 350 assert_eq!( 351 page.records(), 352 [AuthoredDraftQueryRecord::Corrupt { 353 draft_key: *value.draft_id().as_bytes(), 354 revision: value.revision() 355 }] 356 ); 357 assert!(!format!("{page:?}").contains("private original")); 358 store.close().await.unwrap(); 359 } 360 }