lib

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

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 }