authored_draft_all_schemas_tests.rs (6720B)
1 use super::{ 2 query_tests::{corrupt, draft}, 3 tests::open_store, 4 }; 5 use radroots_storage::{ 6 authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, 7 authored_draft_query::{ 8 AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES, AuthoredDraftQuery, AuthoredDraftQueryRecord, 9 AuthoredDraftScope, 10 }, 11 }; 12 use tempfile::TempDir; 13 14 #[tokio::test] 15 async fn all_schema_sqlite_pages_include_unknown_and_corrupt_owners_across_reopen() { 16 let temp = TempDir::new().unwrap(); 17 let store = open_store(&temp).await; 18 let scope = AuthoredDraftScope::new([3; 32]).unwrap(); 19 for id in 1_u128..=1001 { 20 let value = AuthoredDraft::initial( 21 AuthoredDraftId::new(id.to_be_bytes()).unwrap(), 22 [7; 32], 23 if id % 3 == 0 { 24 "future.unknown.v999" 25 } else { 26 "known.v1" 27 }, 28 vec![1], 29 AuthoredDraftStage::Draft, 30 None, 31 10, 32 ) 33 .unwrap(); 34 let value = if id % 2 == 0 { 35 value.with_scope(scope).unwrap() 36 } else { 37 value 38 }; 39 store.append_authored_draft(value, None).await.unwrap(); 40 } 41 corrupt(&store, 0, 7, "future.unknown.v999", Some(scope)).await; 42 corrupt(&store, 2, 8, "future.unknown.v999", None).await; 43 corrupt(&store, 3, 7, "", None).await; 44 corrupt(&store, 4, 7, "future.oversized.v999", None).await; 45 let update_guard: String = sqlx::query_scalar( 46 "SELECT sql FROM sqlite_schema WHERE type = 'trigger' AND name = 'radroots_runtime_authored_draft_revisions_update_guard'" 47 ).fetch_one(store.pool()).await.unwrap(); 48 sqlx::query("DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard") 49 .execute(store.pool()) 50 .await 51 .unwrap(); 52 // The production CHECK rejects oversized writes. Simulate a corrupted file 53 // on one held fixture connection, then restore enforcement before querying. 54 let mut connection = store.pool().acquire().await.unwrap(); 55 sqlx::query("PRAGMA ignore_check_constraints = ON") 56 .execute(&mut *connection) 57 .await 58 .unwrap(); 59 sqlx::query("UPDATE radroots_runtime_authored_draft_revisions SET snapshot = zeroblob(?) WHERE draft_id = ?") 60 .bind((AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES + 1) as i64).bind([4_u8; 16].as_slice()).execute(&mut *connection).await.unwrap(); 61 sqlx::query("PRAGMA ignore_check_constraints = OFF") 62 .execute(&mut *connection) 63 .await 64 .unwrap(); 65 // Exact SQL captured from this fresh fixture's governed migration, with no 66 // external input or interpolation. Restore its immutable catalog verbatim. 67 sqlx::query(sqlx::AssertSqlSafe(update_guard.as_str())) 68 .execute(&mut *connection) 69 .await 70 .unwrap(); 71 drop(connection); 72 let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 37).unwrap(); 73 let first = store.query_authored_drafts(q.clone()).await.unwrap(); 74 assert_eq!(first.records().len(), 37); 75 assert!(matches!( 76 first.records()[0], 77 AuthoredDraftQueryRecord::Corrupt { .. } 78 )); 79 assert!(first.records()[0].draft_id().is_err()); 80 let mut keys: Vec<_> = first 81 .records() 82 .iter() 83 .map(AuthoredDraftQueryRecord::draft_key) 84 .collect(); 85 let mut next = first.next_cursor().cloned(); 86 for id in [1_u128, 1001] { 87 let old = store 88 .authored_draft_head(AuthoredDraftId::new(id.to_be_bytes()).unwrap()) 89 .await 90 .unwrap() 91 .unwrap(); 92 let revised = old 93 .successor(vec![2], AuthoredDraftStage::Draft, None, 11) 94 .unwrap(); 95 store 96 .append_authored_draft(revised, Some(old.revision())) 97 .await 98 .unwrap(); 99 } 100 store.close().await.unwrap(); 101 let store = open_store(&temp).await; 102 while let Some(cursor) = next { 103 let cursor = serde_json::from_slice(&serde_json::to_vec(&cursor).unwrap()).unwrap(); 104 let page = store 105 .query_authored_drafts(q.clone().with_cursor(&cursor).unwrap()) 106 .await 107 .unwrap(); 108 assert!(page.records().len() <= 37); 109 for row in page.records() { 110 keys.push(row.draft_key()); 111 if row.draft_key() == 1001_u128.to_be_bytes() { 112 assert_eq!(row.revision().get(), 2); 113 } 114 if row.draft_key() == [3; 16] || row.draft_key() == [4; 16] { 115 assert!(matches!(row, AuthoredDraftQueryRecord::Corrupt { .. })); 116 } 117 } 118 next = page.next_cursor().cloned(); 119 } 120 let mut expected = vec![[0; 16]]; 121 expected.extend((1_u128..=1001).map(u128::to_be_bytes)); 122 expected.extend([[3; 16], [4; 16]]); 123 assert_eq!(keys, expected); 124 let fresh = store.query_authored_drafts(q).await.unwrap(); 125 assert_eq!(fresh.records()[1].revision().get(), 2); 126 let exact = store 127 .query_authored_drafts(AuthoredDraftQuery::for_author([7; 32], "known.v1", 37).unwrap()) 128 .await 129 .unwrap(); 130 assert!(exact.records().iter().all( 131 |r| matches!(r, AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "known.v1") 132 )); 133 store.close().await.unwrap(); 134 } 135 136 #[tokio::test] 137 async fn all_schema_sqlite_continuation_preserves_snapshot_and_payload_budgets() { 138 for byte in [0, 99] { 139 let temp = TempDir::new().unwrap(); 140 let store = open_store(&temp).await; 141 for (id, schema, scope) in [ 142 (1, "known.v1", None), 143 ( 144 2, 145 "future.v999", 146 Some(AuthoredDraftScope::new([3; 32]).unwrap()), 147 ), 148 ] { 149 store 150 .append_authored_draft(draft(id, schema, scope, vec![byte; 3 * 1024 * 1024]), None) 151 .await 152 .unwrap(); 153 } 154 let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 256).unwrap(); 155 let first = store.query_authored_drafts(q.clone()).await.unwrap(); 156 assert_eq!(first.records().len(), 1); 157 assert!( 158 matches!(&first.records()[0], AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "known.v1") 159 ); 160 let next = store 161 .query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap()) 162 .await 163 .unwrap(); 164 assert_eq!(next.records().len(), 1); 165 assert!( 166 matches!(&next.records()[0], AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "future.v999") 167 ); 168 assert!(next.next_cursor().is_none()); 169 store.close().await.unwrap(); 170 } 171 }