lib

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

authored_durability_tests.rs (11473B)


      1 use super::*;
      2 use crate::{OpenMode, OpenOptions, Paths};
      3 use radroots_storage::event::SourceGeneration;
      4 use std::path::Path;
      5 use tempfile::TempDir;
      6 
      7 #[path = "authored_durability_crash_tests.rs"]
      8 mod crash;
      9 
     10 async fn open(directory: &Path, mode: OpenMode) -> SqliteStorage {
     11     let mut options = OpenOptions::new(Paths::from_directory(directory).unwrap(), mode);
     12     if matches!(mode, OpenMode::Create) {
     13         options = options
     14             .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9)
     15             .unwrap();
     16     }
     17     SqliteStorage::open(options).await.unwrap()
     18 }
     19 
     20 fn first() -> AuthoredDraft {
     21     AuthoredDraft::initial(
     22         AuthoredDraftId::new([41; 16]).unwrap(),
     23         [7; 32],
     24         "radroots.durability-fixture.v1",
     25         b"acknowledged baseline".to_vec(),
     26         AuthoredDraftStage::Draft,
     27         None,
     28         10,
     29     )
     30     .unwrap()
     31 }
     32 
     33 fn next(previous: &AuthoredDraft, payload: Vec<u8>) -> AuthoredDraft {
     34     previous
     35         .successor(payload, AuthoredDraftStage::Draft, None, 11)
     36         .unwrap()
     37 }
     38 
     39 async fn assert_head(store: &SqliteStorage, expected: &AuthoredDraft) {
     40     assert_eq!(
     41         store
     42             .authored_draft_head(expected.draft_id())
     43             .await
     44             .unwrap(),
     45         Some(expected.clone())
     46     );
     47 }
     48 
     49 #[tokio::test]
     50 async fn authored_durability_commit_fault_never_acknowledges_and_exact_retry_recovers() {
     51     let temp = TempDir::new().unwrap();
     52     let store = open(temp.path(), OpenMode::Create).await;
     53     let baseline = first();
     54     let pending = next(&baseline, b"pending complete revision".to_vec());
     55     store
     56         .append_authored_draft(baseline.clone(), None)
     57         .await
     58         .unwrap();
     59     sqlx::query("CREATE TABLE authored_commit_parent (id INTEGER PRIMARY KEY)")
     60         .execute(store.pool())
     61         .await
     62         .unwrap();
     63     sqlx::query("CREATE TABLE authored_commit_fault (id INTEGER REFERENCES authored_commit_parent(id) DEFERRABLE INITIALLY DEFERRED)")
     64         .execute(store.pool()).await.unwrap();
     65     sqlx::query("CREATE TRIGGER authored_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_draft_revisions BEGIN INSERT INTO authored_commit_fault VALUES (99); END")
     66         .execute(store.pool()).await.unwrap();
     67 
     68     // Establish that the insertion succeeds and the actual COMMIT is the fault.
     69     let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap();
     70     insert_draft_tx(&mut transaction, &pending).await.unwrap();
     71     let failure = transaction.commit().await.unwrap_err();
     72     assert_eq!(
     73         failure.as_database_error().unwrap().kind(),
     74         sqlx::error::ErrorKind::ForeignKeyViolation
     75     );
     76     assert_eq!(
     77         store
     78             .append_authored_draft(pending.clone(), Some(baseline.revision()))
     79             .await,
     80         Err(Error::BackendUnavailable)
     81     );
     82     assert_head(&store, &baseline).await;
     83     assert!(
     84         store
     85             .authored_draft_revision(pending.draft_id(), pending.revision())
     86             .await
     87             .unwrap()
     88             .is_none()
     89     );
     90 
     91     for statement in [
     92         "DROP TRIGGER authored_commit_fault_trigger",
     93         "DROP TABLE authored_commit_fault",
     94         "DROP TABLE authored_commit_parent",
     95     ] {
     96         sqlx::query(statement).execute(store.pool()).await.unwrap();
     97     }
     98     let receipt = store
     99         .append_authored_draft(pending.clone(), Some(baseline.revision()))
    100         .await
    101         .unwrap();
    102     assert_eq!(receipt.disposition(), DraftAppendDisposition::Inserted);
    103     store.close().await.unwrap();
    104     let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await;
    105     assert_head(&reopened, &pending).await;
    106     assert_eq!(
    107         reopened
    108             .append_authored_draft(pending, Some(baseline.revision()))
    109             .await
    110             .unwrap()
    111             .disposition(),
    112         DraftAppendDisposition::Replay
    113     );
    114     reopened.close().await.unwrap();
    115 }
    116 
    117 #[tokio::test]
    118 async fn authored_durability_sqlite_capacity_failure_preserves_the_acknowledged_head() {
    119     let temp = TempDir::new().unwrap();
    120     let store = open(temp.path(), OpenMode::Create).await;
    121     let baseline = first();
    122     store
    123         .append_authored_draft(baseline.clone(), None)
    124         .await
    125         .unwrap();
    126     let mut connections = Vec::new();
    127     for _ in 0..4 {
    128         connections.push(store.pool().acquire().await.unwrap());
    129     }
    130     let mut original_limits = Vec::new();
    131     for connection in &mut connections {
    132         original_limits.push(
    133             sqlx::query_scalar::<_, i64>("PRAGMA max_page_count")
    134                 .fetch_one(&mut **connection)
    135                 .await
    136                 .unwrap(),
    137         );
    138         let pages = sqlx::query_scalar::<_, i64>("PRAGMA page_count")
    139             .fetch_one(&mut **connection)
    140             .await
    141             .unwrap();
    142         // PRAGMA assignments do not accept bind parameters; only an i64 is interpolated.
    143         let limit = sqlx::query_scalar::<_, i64>(sqlx::AssertSqlSafe(format!(
    144             "PRAGMA max_page_count = {}",
    145             pages + 2
    146         )))
    147         .fetch_one(&mut **connection)
    148         .await
    149         .unwrap();
    150         assert_eq!(limit, pages + 2);
    151     }
    152     drop(connections);
    153     // This bounded allocation proves SQLITE_FULL without filling the host disk.
    154     let failure =
    155         sqlx::query("CREATE TABLE authored_capacity_probe AS SELECT zeroblob(1048576) AS data")
    156             .execute(store.pool())
    157             .await
    158             .unwrap_err();
    159     assert_eq!(
    160         failure.as_database_error().unwrap().code().as_deref(),
    161         Some("13")
    162     );
    163     let pending = next(&baseline, vec![42; 256 * 1024]);
    164     assert_eq!(
    165         store
    166             .append_authored_draft(pending.clone(), Some(baseline.revision()))
    167             .await,
    168         Err(Error::SpaceInsufficient)
    169     );
    170     assert_head(&store, &baseline).await;
    171 
    172     let mut connections = Vec::new();
    173     for _ in 0..4 {
    174         connections.push(store.pool().acquire().await.unwrap());
    175     }
    176     for (connection, limit) in connections.iter_mut().zip(original_limits) {
    177         sqlx::query(sqlx::AssertSqlSafe(format!(
    178             "PRAGMA max_page_count = {limit}"
    179         )))
    180         .execute(&mut **connection)
    181         .await
    182         .unwrap();
    183     }
    184     drop(connections);
    185     store.close().await.unwrap();
    186     let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await;
    187     assert_head(&reopened, &baseline).await;
    188     assert_eq!(
    189         reopened
    190             .append_authored_draft(pending.clone(), Some(baseline.revision()))
    191             .await
    192             .unwrap()
    193             .disposition(),
    194         DraftAppendDisposition::Inserted
    195     );
    196     assert_head(&reopened, &pending).await;
    197     reopened.close().await.unwrap();
    198 }
    199 
    200 #[tokio::test]
    201 async fn authored_durability_denied_writes_have_no_receipt_or_head_advance() {
    202     let temp = TempDir::new().unwrap();
    203     let store = open(temp.path(), OpenMode::Create).await;
    204     let baseline = first();
    205     let pending = next(&baseline, b"denied revision".to_vec());
    206     store
    207         .append_authored_draft(baseline.clone(), None)
    208         .await
    209         .unwrap();
    210     let mut connections = Vec::new();
    211     for _ in 0..4 {
    212         connections.push(store.pool().acquire().await.unwrap());
    213     }
    214     for connection in &mut connections {
    215         sqlx::query("PRAGMA query_only = ON")
    216             .execute(&mut **connection)
    217             .await
    218             .unwrap();
    219     }
    220     drop(connections);
    221     assert_eq!(
    222         store
    223             .append_authored_draft(pending, Some(baseline.revision()))
    224             .await,
    225         Err(Error::BackendUnavailable)
    226     );
    227     assert_head(&store, &baseline).await;
    228     store.close().await.unwrap();
    229     let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await;
    230     assert_head(&reopened, &baseline).await;
    231     reopened.close().await.unwrap();
    232 }
    233 
    234 #[tokio::test]
    235 async fn authored_durability_busy_writer_never_returns_a_success_receipt() {
    236     let temp = TempDir::new().unwrap();
    237     let store = SqliteStorage::open(
    238         OpenOptions::new(
    239             Paths::from_directory(temp.path()).unwrap(),
    240             OpenMode::Create,
    241         )
    242         .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9)
    243         .unwrap()
    244         .with_busy_timeout(std::time::Duration::from_millis(10))
    245         .unwrap(),
    246     )
    247     .await
    248     .unwrap();
    249     let baseline = first();
    250     let pending = next(&baseline, b"blocked complete revision".to_vec());
    251     store
    252         .append_authored_draft(baseline.clone(), None)
    253         .await
    254         .unwrap();
    255     let transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap();
    256     assert_eq!(
    257         store
    258             .append_authored_draft(pending.clone(), Some(baseline.revision()))
    259             .await,
    260         Err(Error::BackendUnavailable)
    261     );
    262     assert_head(&store, &baseline).await;
    263     transaction.rollback().await.unwrap();
    264     assert_eq!(
    265         store
    266             .append_authored_draft(pending.clone(), Some(baseline.revision()))
    267             .await
    268             .unwrap()
    269             .disposition(),
    270         DraftAppendDisposition::Inserted
    271     );
    272     assert_head(&store, &pending).await;
    273     store.close().await.unwrap();
    274 }
    275 
    276 #[tokio::test]
    277 async fn authored_durability_read_only_and_closed_stores_cannot_acknowledge() {
    278     let temp = TempDir::new().unwrap();
    279     let store = open(temp.path(), OpenMode::Create).await;
    280     let baseline = first();
    281     let pending = next(&baseline, b"not writable".to_vec());
    282     store
    283         .append_authored_draft(baseline.clone(), None)
    284         .await
    285         .unwrap();
    286     store.close().await.unwrap();
    287     assert_eq!(
    288         store
    289             .append_authored_draft(pending.clone(), Some(baseline.revision()))
    290             .await,
    291         Err(Error::BackendUnavailable)
    292     );
    293     let read_only = open(temp.path(), OpenMode::ReadOnly).await;
    294     assert_head(&read_only, &baseline).await;
    295     assert_eq!(
    296         read_only
    297             .append_authored_draft(pending, Some(baseline.revision()))
    298             .await,
    299         Err(Error::BackendUnavailable)
    300     );
    301     read_only.close().await.unwrap();
    302 }
    303 
    304 #[test]
    305 fn authored_durability_contract_distinguishes_crash_and_power_loss() {
    306     let policy: toml::Value = toml::from_str(include_str!(
    307         "../../../contracts/storage/failure_injection_policy_v1.toml"
    308     ))
    309     .unwrap();
    310     let authored = &policy["authored_write"];
    311     assert_eq!(
    312         authored["acknowledgment"].as_str(),
    313         Some("only_after_successful_commit_or_exact_committed_replay")
    314     );
    315     assert_eq!(
    316         authored["commit_fault"].as_str(),
    317         Some("deferred_foreign_key_at_actual_commit")
    318     );
    319     assert_eq!(
    320         authored["capacity_fault"].as_str(),
    321         Some("bounded_sqlite_max_page_count")
    322     );
    323     assert_eq!(
    324         authored["write_denied_fault"].as_str(),
    325         Some("owned_connection_query_only")
    326     );
    327     assert_eq!(
    328         authored["busy_fault"].as_str(),
    329         Some("owned_begin_immediate_with_bounded_timeout")
    330     );
    331     assert_eq!(
    332         authored["process_termination_points"]
    333             .as_array()
    334             .unwrap()
    335             .iter()
    336             .map(|value| value.as_str().unwrap())
    337             .collect::<Vec<_>>(),
    338         ["after_acknowledgment", "during_uncommitted_insert"]
    339     );
    340     assert_eq!(authored["power_loss_qualified"].as_bool(), Some(false));
    341     assert_eq!(
    342         authored["protected_data_policy_owner"].as_str(),
    343         Some("native_host")
    344     );
    345 }