lib

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

authored_draft_pair_tests.rs (5126B)


      1 use super::*;
      2 use crate::{OpenMode, OpenOptions, Paths};
      3 use radroots_storage::{authored_draft_pair::AuthoredDraftPair, event::SourceGeneration};
      4 use std::{future::poll_fn, task::Poll};
      5 fn draft(id: u8, schema: &str) -> AuthoredDraft {
      6     AuthoredDraft::initial(
      7         AuthoredDraftId::new([id; 16]).unwrap(),
      8         [9; 32],
      9         schema,
     10         b"fixture".to_vec(),
     11         AuthoredDraftStage::Draft,
     12         None,
     13         10,
     14     )
     15     .unwrap()
     16 }
     17 async fn open(path: &std::path::Path, mode: OpenMode) -> SqliteStorage {
     18     let options = OpenOptions::new(Paths::from_directory(path).unwrap(), mode);
     19     let options = if mode == OpenMode::Create {
     20         options
     21             .with_source_generation(SourceGeneration::new([7; 32]).unwrap(), 10)
     22             .unwrap()
     23     } else {
     24         options
     25     };
     26     SqliteStorage::open(options).await.unwrap()
     27 }
     28 #[tokio::test]
     29 async fn second_sql_insert_failure_rolls_back_first_revision() {
     30     let dir = tempfile::tempdir().unwrap();
     31     let store = open(dir.path(), OpenMode::Create).await;
     32     sqlx::query("CREATE TRIGGER fixture_pair_abort BEFORE INSERT ON radroots_runtime_authored_draft_revisions WHEN NEW.payload_schema='fixture.reject.v1' BEGIN SELECT RAISE(ABORT,'fixture pair rejection'); END").execute(store.pool()).await.unwrap();
     33     let request = AuthoredDraftPair::new(
     34         draft(1, "fixture.pair.v1"),
     35         None,
     36         draft(2, "fixture.reject.v1"),
     37         None,
     38     )
     39     .unwrap();
     40     assert!(store.append_authored_draft_pair(request).await.is_err());
     41     assert!(
     42         store
     43             .authored_draft_heads([9; 32], 10)
     44             .await
     45             .unwrap()
     46             .is_empty()
     47     );
     48     sqlx::query("DROP TRIGGER fixture_pair_abort")
     49         .execute(store.pool())
     50         .await
     51         .unwrap();
     52     store.close().await.unwrap();
     53     let store = open(dir.path(), OpenMode::ReadWriteExisting).await;
     54     assert!(
     55         store
     56             .authored_draft_heads([9; 32], 10)
     57             .await
     58             .unwrap()
     59             .is_empty()
     60     );
     61     store.close().await.unwrap();
     62 }
     63 #[tokio::test]
     64 async fn cancelling_real_pair_futures_never_leaves_one_member_after_reopen() {
     65     let dir = tempfile::tempdir().unwrap();
     66     let store = open(dir.path(), OpenMode::Create).await;
     67     let mut cancelled = 0;
     68     let mut completed = 0;
     69     let mut durable = Vec::new();
     70     // Cancel at successive real Pending boundaries, including caller loss near
     71     // commit. No test hook or surrogate transaction replaces the public method.
     72     for boundary in 1..=64u8 {
     73         let a = draft(boundary * 2, "fixture.pair.v1");
     74         let b = draft(boundary * 2 + 1, "fixture.pair.v1");
     75         let mut future = store.append_authored_draft_pair(
     76             AuthoredDraftPair::new(a.clone(), None, b.clone(), None).unwrap(),
     77         );
     78         let mut pending = 0;
     79         let acknowledged = poll_fn(|cx| match future.as_mut().poll(cx) {
     80             Poll::Ready(result) => {
     81                 result.unwrap();
     82                 Poll::Ready(true)
     83             }
     84             Poll::Pending => {
     85                 pending += 1;
     86                 if pending == boundary {
     87                     Poll::Ready(false)
     88                 } else {
     89                     Poll::Pending
     90                 }
     91             }
     92         })
     93         .await;
     94         drop(future);
     95         if acknowledged {
     96             completed += 1;
     97         } else {
     98             cancelled += 1;
     99         }
    100         // Drain any queued commit/rollback before observing both heads.
    101         store
    102             .pool()
    103             .begin_with("BEGIN IMMEDIATE")
    104             .await
    105             .unwrap()
    106             .rollback()
    107             .await
    108             .unwrap();
    109         let first = store.authored_draft_head(a.draft_id()).await.unwrap();
    110         let second = store.authored_draft_head(b.draft_id()).await.unwrap();
    111         assert_eq!(first.is_some(), second.is_some(), "boundary {boundary}");
    112         if acknowledged {
    113             assert_eq!(first, Some(a.clone()));
    114         }
    115         durable.push((a, b, first.is_some()));
    116     }
    117     assert!(
    118         cancelled > 0 && completed > 0,
    119         "cancelled={cancelled}, completed={completed}"
    120     );
    121     store.close().await.unwrap();
    122     let store = open(dir.path(), OpenMode::ReadWriteExisting).await;
    123     for (a, b, present) in durable {
    124         assert_eq!(
    125             store
    126                 .authored_draft_head(a.draft_id())
    127                 .await
    128                 .unwrap()
    129                 .is_some(),
    130             present
    131         );
    132         assert_eq!(
    133             store
    134                 .authored_draft_head(b.draft_id())
    135                 .await
    136                 .unwrap()
    137                 .is_some(),
    138             present
    139         );
    140         if present {
    141             let receipts = store
    142                 .append_authored_draft_pair(AuthoredDraftPair::new(a, None, b, None).unwrap())
    143                 .await
    144                 .unwrap();
    145             assert!(
    146                 receipts
    147                     .iter()
    148                     .all(|r| r.disposition() == DraftAppendDisposition::Replay)
    149             );
    150         }
    151     }
    152     store.close().await.unwrap();
    153 }