lib

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

authored_delivery_fact_tests.rs (10126B)


      1 use super::signed_fact_fixture::{RAW, claim, ids, record};
      2 use super::*;
      3 use crate::OpenMode;
      4 use radroots_storage::{
      5     authored::WorkClaim,
      6     authored_atomic::{CancelAuthoredWork, ClaimAuthoredWork, RecordDeliveryFact},
      7 };
      8 use radroots_transport::{DeliveryReceipt, outcome::DeliveryOutcome, sink::DeliveryTargetReceipt};
      9 use tempfile::TempDir;
     10 
     11 async fn prepared(temp: &TempDir) -> (SqliteStorage, WorkClaim, AuthoredAtomicReceipt) {
     12     let (store, event, signing) = super::signed_fact_tests::prepared(temp).await;
     13     store
     14         .execute_authored(record(event, signing, 12))
     15         .await
     16         .unwrap();
     17     let plan = store
     18         .authored_delivery_plan(ids().2)
     19         .await
     20         .unwrap()
     21         .unwrap();
     22     let active = claim(plan.revision(), 5, 13);
     23     let original = store
     24         .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
     25             ClaimAuthoredTarget::DeliveryPlan(ids().2),
     26             active.clone(),
     27         )))
     28         .await
     29         .unwrap();
     30     (store, active, original)
     31 }
     32 
     33 pub(super) fn fact(plan: &AuthoredDeliveryPlan, active: WorkClaim) -> AuthoredAtomicCommand {
     34     let request = plan.request().unwrap();
     35     let receipt = DeliveryReceipt::for_request(
     36         request,
     37         request
     38             .target_set()
     39             .targets()
     40             .iter()
     41             .cloned()
     42             .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()))
     43             .collect(),
     44     )
     45     .unwrap();
     46     AuthoredAtomicCommand::RecordDelivery(
     47         RecordDeliveryFact::new(
     48             ids().2,
     49             ids().1,
     50             active,
     51             DeliveryAttemptOutcome::Receipt(receipt),
     52             50,
     53         )
     54         .unwrap(),
     55     )
     56 }
     57 
     58 #[tokio::test]
     59 async fn stopped_late_delivery_reopens_and_preserves_original_receipt_and_exact_raw() {
     60     let temp = TempDir::new().unwrap();
     61     let (store, active, original) = prepared(&temp).await;
     62     let plan = store
     63         .authored_delivery_plan(ids().2)
     64         .await
     65         .unwrap()
     66         .unwrap();
     67     let command = fact(&plan, active);
     68     store
     69         .execute_authored(AuthoredAtomicCommand::Cancel(
     70             CancelAuthoredWork::new(
     71                 CancelAuthoredTarget::DeliveryPlan(ids().2),
     72                 plan.revision(),
     73                 20,
     74             )
     75             .unwrap(),
     76         ))
     77         .await
     78         .unwrap();
     79     let receipt = store.execute_authored(command.clone()).await.unwrap();
     80     let retained = store
     81         .authored_delivery_plan(ids().2)
     82         .await
     83         .unwrap()
     84         .unwrap();
     85     assert_eq!(retained.state(), AuthoredDeliveryState::Cancelled);
     86     assert_eq!(retained.stop_requested_at_unix_ms(), Some(20));
     87     assert_eq!(
     88         retained.delivery_satisfaction().unwrap(),
     89         SatisfactionState::Satisfied
     90     );
     91     assert_eq!(
     92         retained.request().unwrap().payload().event().raw_json(),
     93         RAW
     94     );
     95     assert_eq!(retained.delivery_facts().len(), 1);
     96     for sql in [
     97         "UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99",
     98         "DELETE FROM radroots_runtime_authored_delivery_facts",
     99         "UPDATE radroots_runtime_authored_delivery_plans SET stop_requested_at_unix_ms = NULL",
    100     ] {
    101         assert!(sqlx::query(sql).execute(store.pool()).await.is_err());
    102     }
    103     store.close().await.unwrap();
    104     let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await;
    105     assert_eq!(
    106         store
    107             .authored_delivery_plan(ids().2)
    108             .await
    109             .unwrap()
    110             .unwrap(),
    111         retained
    112     );
    113     assert_eq!(
    114         store
    115             .authored_receipt(original.commit_id())
    116             .await
    117             .unwrap()
    118             .unwrap(),
    119         original
    120     );
    121     let replay = store.execute_authored(command).await.unwrap();
    122     assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
    123     assert_eq!(replay.outcome(), receipt.outcome());
    124     assert_eq!(
    125         store
    126             .authored_delivery_plan(ids().2)
    127             .await
    128             .unwrap()
    129             .unwrap(),
    130         retained
    131     );
    132     store.close().await.unwrap();
    133 }
    134 
    135 #[tokio::test]
    136 async fn real_delivery_commit_failure_rolls_back_facts_snapshot_and_receipt() {
    137     let temp = TempDir::new().unwrap();
    138     let (store, active, _) = prepared(&temp).await;
    139     let before = store
    140         .authored_delivery_plan(ids().2)
    141         .await
    142         .unwrap()
    143         .unwrap();
    144     let command = fact(&before, active);
    145     sqlx::query("CREATE TABLE delivery_fact_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)").execute(store.pool()).await.unwrap();
    146     sqlx::query("CREATE TRIGGER delivery_fact_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_atomic_commits WHEN NEW.phase IN ('delivery', 'cancel') BEGIN INSERT INTO delivery_fact_commit_fault VALUES (x'99999999999999999999999999999999'); END").execute(store.pool()).await.unwrap();
    147     let stop = AuthoredAtomicCommand::Cancel(
    148         CancelAuthoredWork::new(
    149             CancelAuthoredTarget::DeliveryPlan(ids().2),
    150             before.revision(),
    151             20,
    152         )
    153         .unwrap(),
    154     );
    155     assert!(store.execute_authored(stop.clone()).await.is_err());
    156     assert_eq!(
    157         store
    158             .authored_delivery_plan(ids().2)
    159             .await
    160             .unwrap()
    161             .unwrap(),
    162         before
    163     );
    164     assert!(
    165         store
    166             .authored_receipt(stop.commit_id())
    167             .await
    168             .unwrap()
    169             .is_none()
    170     );
    171     assert!(store.execute_authored(command.clone()).await.is_err());
    172     assert_eq!(
    173         store
    174             .authored_delivery_plan(ids().2)
    175             .await
    176             .unwrap()
    177             .unwrap(),
    178         before
    179     );
    180     assert!(
    181         store
    182             .authored_receipt(command.commit_id())
    183             .await
    184             .unwrap()
    185             .is_none()
    186     );
    187     assert_eq!(
    188         sqlx::query_scalar::<_, i64>(
    189             "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_facts"
    190         )
    191         .fetch_one(store.pool())
    192         .await
    193         .unwrap(),
    194         0
    195     );
    196     assert_eq!(
    197         sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_fact_commit_fault")
    198             .fetch_one(store.pool())
    199             .await
    200             .unwrap(),
    201         0
    202     );
    203     sqlx::query("DROP TRIGGER delivery_fact_commit_fault_trigger")
    204         .execute(store.pool())
    205         .await
    206         .unwrap();
    207     sqlx::query("DROP TABLE delivery_fact_commit_fault")
    208         .execute(store.pool())
    209         .await
    210         .unwrap();
    211     store.close().await.unwrap();
    212     let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await;
    213     assert_eq!(
    214         store
    215             .authored_delivery_plan(ids().2)
    216             .await
    217             .unwrap()
    218             .unwrap(),
    219         before
    220     );
    221     store.execute_authored(command.clone()).await.unwrap();
    222     assert_eq!(
    223         store
    224             .authored_delivery_plan(ids().2)
    225             .await
    226             .unwrap()
    227             .unwrap()
    228             .delivery_satisfaction()
    229             .unwrap(),
    230         SatisfactionState::Satisfied
    231     );
    232     store.close().await.unwrap();
    233 }
    234 
    235 #[tokio::test]
    236 async fn forged_claim_and_normalized_fact_corruption_fail_closed() {
    237     let temp = TempDir::new().unwrap();
    238     let (store, active, _) = prepared(&temp).await;
    239     let before = store
    240         .authored_delivery_plan(ids().2)
    241         .await
    242         .unwrap()
    243         .unwrap();
    244     let forged = WorkClaim::new(
    245         *active.token(),
    246         "forged-owner",
    247         active.generation(),
    248         active.acquired_at_unix_ms(),
    249         active.expires_at_unix_ms(),
    250         active.row_revision(),
    251     )
    252     .unwrap();
    253     assert_eq!(
    254         store.execute_authored(fact(&before, forged)).await,
    255         Err(Error::AtomicWorkflowMismatch)
    256     );
    257     assert_eq!(
    258         store
    259             .authored_delivery_plan(ids().2)
    260             .await
    261             .unwrap()
    262             .unwrap(),
    263         before
    264     );
    265     store.execute_authored(fact(&before, active)).await.unwrap();
    266     sqlx::query("DROP TRIGGER radroots_runtime_authored_delivery_facts_update_guard")
    267         .execute(store.pool())
    268         .await
    269         .unwrap();
    270     sqlx::query("UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99")
    271         .execute(store.pool())
    272         .await
    273         .unwrap();
    274     assert_eq!(
    275         store.authored_delivery_plan(ids().2).await,
    276         Err(Error::InvalidAuthoredDeliveryPlan)
    277     );
    278     store.close().await.unwrap();
    279 }
    280 
    281 #[tokio::test]
    282 async fn delivery_read_snapshot_remains_consistent_across_a_late_commit() {
    283     let temp = TempDir::new().unwrap();
    284     let (store, active, _) = prepared(&temp).await;
    285     let before = store
    286         .authored_delivery_plan(ids().2)
    287         .await
    288         .unwrap()
    289         .unwrap();
    290     let mut read = store.pool().begin().await.unwrap();
    291     let frozen = load_optional_plan_tx(&mut read, ids().2)
    292         .await
    293         .unwrap()
    294         .unwrap();
    295     assert_eq!(frozen, before);
    296     store.execute_authored(fact(&before, active)).await.unwrap();
    297     validate_plan_children_tx(&mut read, &frozen).await.unwrap();
    298     assert_eq!(
    299         load_optional_plan_tx(&mut read, ids().2)
    300             .await
    301             .unwrap()
    302             .unwrap(),
    303         frozen
    304     );
    305     read.rollback().await.unwrap();
    306     let current = store
    307         .authored_delivery_plan(ids().2)
    308         .await
    309         .unwrap()
    310         .unwrap();
    311     assert_eq!(current.delivery_facts().len(), 1);
    312     assert_eq!(current.revision(), before.revision());
    313     let mut stale = store.pool().begin().await.unwrap();
    314     assert_eq!(
    315         delivery_facts::persist(&mut stale, &before).await,
    316         Err(Error::InvalidAuthoredDeliveryPlan)
    317     );
    318     stale.rollback().await.unwrap();
    319     assert_eq!(
    320         store
    321             .authored_delivery_plan(ids().2)
    322             .await
    323             .unwrap()
    324             .unwrap(),
    325         current
    326     );
    327     store.close().await.unwrap();
    328 }