migration_delivery_reconciliation_tests.rs (8232B)
1 use super::{ 2 tests::{connection, establish_runtime_version, pragma}, 3 *, 4 }; 5 use crate::authored::signed_fact_fixture::{claim, prepare}; 6 use radroots_storage::{ 7 atomic::AtomicCommitDisposition, 8 authored_atomic::{ 9 AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, ClaimAuthoredTarget, 10 ClaimAuthoredWork, 11 }, 12 }; 13 14 const FAILING_V17: &str = concat!( 15 include_str!("migration/runtime/0017_authored_delivery_reconciliation.up.sql"), 16 "\nINSERT INTO missing_fixture_table VALUES (1);" 17 ); 18 19 fn plan(current: u32, fail: bool) -> MigrationPlan { 20 MigrationPlan { 21 database: RUNTIME_DATABASE, 22 application_id: RUNTIME_APPLICATION_ID, 23 set_application_id_sql: SET_RUNTIME_APPLICATION_ID, 24 minimum_version: runtime::MINIMUM_VERSION, 25 current_version: current, 26 steps: runtime::MIGRATIONS 27 .iter() 28 .take(current as usize) 29 .map(|step| MigrationStep { 30 version: step.version(), 31 sql: if fail && step.version() == 17 { 32 FAILING_V17 33 } else { 34 runtime::migration_sql(step.version()).unwrap() 35 }, 36 owned_objects: step.owned_objects(), 37 }) 38 .collect(), 39 } 40 } 41 42 async fn seed(connection: &mut SqliteConnection) { 43 super::signed_facts_tests::seed(connection).await; 44 let (command, event) = prepare(); 45 let AuthoredAtomicCommand::Prepare(prepared) = command else { 46 unreachable!() 47 }; 48 let mut delivery = prepared.delivery_plans()[0].clone(); 49 delivery.bind_signed_event(event, 12).unwrap(); 50 for (token, at) in [(5, 13), (6, 40)] { 51 let active = claim(delivery.revision(), token, at); 52 delivery.claim(active.clone(), at).unwrap(); 53 let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 54 ClaimAuthoredTarget::DeliveryPlan(delivery.plan_id()), 55 active, 56 )); 57 let receipt = AuthoredAtomicReceipt::new( 58 &command, 59 AtomicCommitDisposition::Committed, 60 at, 61 AuthoredAtomicOutcome::DeliveryPlan(delivery.clone()), 62 ) 63 .unwrap(); 64 sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, ?, ?, ?)") 65 .bind(receipt.commit_id().as_bytes().as_slice()).bind(receipt.digest().as_bytes().as_slice()) 66 .bind(delivery.plan_id().as_bytes().as_slice()).bind(at as i64).bind(at as i64) 67 .bind(serde_json::to_vec(&serde_json::json!({"outcome": receipt.outcome()})).unwrap()) 68 .execute(&mut *connection).await.unwrap(); 69 } 70 delivery.request_stop(61).unwrap(); 71 let mut transaction = connection.begin().await.unwrap(); 72 crate::authored::persist_plan_v11(&mut transaction, &delivery) 73 .await 74 .unwrap(); 75 transaction.commit().await.unwrap(); 76 } 77 78 async fn snapshots(connection: &mut SqliteConnection) -> (Vec<u8>, Vec<(Vec<u8>, Vec<u8>)>) { 79 let plan = sqlx::query_scalar("SELECT snapshot FROM radroots_runtime_authored_delivery_plans") 80 .fetch_one(&mut *connection) 81 .await 82 .unwrap(); 83 let receipts = sqlx::query_as("SELECT commit_id, receipt FROM radroots_runtime_authored_atomic_commits ORDER BY commit_id").fetch_all(&mut *connection).await.unwrap(); 84 (plan, receipts) 85 } 86 87 #[tokio::test] 88 async fn v17_backfills_stopped_issued_history_without_rewriting_old_receipts() { 89 let mut connection = connection().await; 90 establish_runtime_version(&mut connection, 16).await; 91 seed(&mut connection).await; 92 let before = snapshots(&mut connection).await; 93 assert!(matches!( 94 migrate(&mut connection, OpenMode::ReadOnly, &plan(17, false)).await, 95 Err(Error::SchemaMigrationRequired { 96 actual: 16, 97 current: 17, 98 .. 99 }) 100 )); 101 assert_eq!( 102 migrate( 103 &mut connection, 104 OpenMode::ReadWriteExisting, 105 &plan(17, false) 106 ) 107 .await 108 .unwrap() 109 .applied(), 110 1 111 ); 112 assert_eq!(snapshots(&mut connection).await, before); 113 let projected: Vec<Vec<u8>> = sqlx::query_scalar( 114 "SELECT claim_id FROM radroots_runtime_authored_delivery_claims ORDER BY claim_id", 115 ) 116 .fetch_all(&mut connection) 117 .await 118 .unwrap(); 119 let original: Vec<Vec<u8>> = sqlx::query_scalar("SELECT commit_id FROM radroots_runtime_authored_atomic_commits WHERE phase = 'claim' ORDER BY commit_id").fetch_all(&mut connection).await.unwrap(); 120 assert_eq!(projected.len(), 2); 121 assert_eq!(projected, original); 122 for mode in [OpenMode::ReadOnly, OpenMode::ReadWriteExisting] { 123 assert!(matches!( 124 migrate(&mut connection, mode, &plan(16, false)).await, 125 Err(Error::SchemaTooNew { 126 actual: 17, 127 supported: 16, 128 .. 129 }) 130 )); 131 } 132 assert_eq!( 133 migrate(&mut connection, OpenMode::ReadOnly, &plan(17, false)) 134 .await 135 .unwrap() 136 .applied(), 137 0 138 ); 139 connection.close().await.unwrap(); 140 } 141 142 #[tokio::test] 143 async fn failed_v17_rolls_back_all_objects_and_preserves_original_claim_bytes() { 144 let mut connection = connection().await; 145 establish_runtime_version(&mut connection, 16).await; 146 seed(&mut connection).await; 147 let before = snapshots(&mut connection).await; 148 assert!( 149 migrate( 150 &mut connection, 151 OpenMode::ReadWriteExisting, 152 &plan(17, true) 153 ) 154 .await 155 .is_err() 156 ); 157 assert_eq!(pragma(&mut connection, "user_version").await, 16); 158 assert_eq!(snapshots(&mut connection).await, before); 159 assert_eq!(sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM sqlite_schema WHERE name LIKE 'radroots_runtime_authored_delivery_claims%' OR name LIKE 'radroots_runtime_authored_delivery_reconciliations%' OR name = 'radroots_runtime_authored_atomic_target_phase_idx'").fetch_one(&mut connection).await.unwrap(), 0); 160 migrate( 161 &mut connection, 162 OpenMode::ReadWriteExisting, 163 &plan(17, false), 164 ) 165 .await 166 .unwrap(); 167 assert_eq!(snapshots(&mut connection).await, before); 168 connection.close().await.unwrap(); 169 } 170 171 #[tokio::test] 172 async fn malformed_claim_namespaces_cannot_silently_disappear_during_backfill() { 173 for wire in [ 174 "{}", 175 "{!", 176 r#"{"outcome":{"delivery_plan":null}}"#, 177 r#"{"outcome":{"delivery_plan":{},"artifact":{}}}"#, 178 ] { 179 let mut connection = connection().await; 180 establish_runtime_version(&mut connection, 16).await; 181 super::signed_facts_tests::seed(&mut connection).await; 182 sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, 13, 13, ?)") 183 .bind([9u8; 16].as_slice()).bind([8u8; 32].as_slice()).bind([3u8; 16].as_slice()).bind(wire.as_bytes()).execute(&mut connection).await.unwrap(); 184 let before = snapshots(&mut connection).await; 185 assert!( 186 migrate( 187 &mut connection, 188 OpenMode::ReadWriteExisting, 189 &plan(17, false) 190 ) 191 .await 192 .is_err(), 193 "{wire}" 194 ); 195 assert_eq!(pragma(&mut connection, "user_version").await, 16); 196 assert_eq!(snapshots(&mut connection).await, before); 197 connection.close().await.unwrap(); 198 } 199 } 200 201 #[test] 202 fn delivery_reconciliation_decision_binds_exact_forward_migration() { 203 let decision: serde_json::Value = serde_json::from_str(include_str!( 204 "../../../contracts/architecture/decisions/authored_delivery_reconciliation.v1.json" 205 )) 206 .unwrap(); 207 let migration = runtime::MIGRATIONS[16]; 208 assert_eq!(decision["migration"]["version"], migration.version()); 209 assert_eq!(decision["migration"]["name"], migration.name()); 210 assert_eq!(decision["migration"]["sha256"], migration.up_sha256()); 211 }