lib

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

authored_durability_crash_tests.rs (5087B)


      1 use super::*;
      2 use std::{
      3     fs,
      4     process::{Child, Command, Stdio},
      5     time::{Duration, Instant},
      6 };
      7 
      8 const CHILD_TEST: &str = "authored_draft::durability_tests::crash::authored_durability_child";
      9 const CHILD_ROOT: &str = "RADROOTS_AUTHORED_CRASH_FIXTURE_ROOT";
     10 const CHILD_PHASE: &str = "RADROOTS_AUTHORED_CRASH_FIXTURE_PHASE";
     11 
     12 struct ChildGuard(Child);
     13 
     14 impl Drop for ChildGuard {
     15     fn drop(&mut self) {
     16         let _ = self.0.kill();
     17         let _ = self.0.wait();
     18     }
     19 }
     20 
     21 #[tokio::test]
     22 async fn authored_durability_child() {
     23     let Some(directory) = std::env::var_os(CHILD_ROOT) else {
     24         return;
     25     };
     26     let directory = Path::new(&directory);
     27     assert!(directory.is_absolute() && directory.is_dir());
     28     let phase = std::env::var(CHILD_PHASE).unwrap();
     29     assert!(matches!(
     30         phase.as_str(),
     31         "after_acknowledgment" | "during_uncommitted_insert"
     32     ));
     33     let store = open(directory, OpenMode::ReadWriteExisting).await;
     34     let baseline = first();
     35     let pending = next(&baseline, b"complete child revision".to_vec());
     36     if phase == "after_acknowledgment" {
     37         let receipt = store
     38             .append_authored_draft(pending.clone(), Some(baseline.revision()))
     39             .await
     40             .unwrap();
     41         assert_eq!(receipt.disposition(), DraftAppendDisposition::Inserted);
     42         assert_head(&store, &pending).await;
     43         fs::write(directory.join("child-ready"), phase.as_bytes()).unwrap();
     44         std::future::pending::<()>().await;
     45     } else {
     46         let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap();
     47         insert_draft_tx(&mut transaction, &pending).await.unwrap();
     48         assert_eq!(
     49             load_head_tx(&mut transaction, baseline.draft_id())
     50                 .await
     51                 .unwrap(),
     52             Some(pending)
     53         );
     54         fs::write(directory.join("child-ready"), phase.as_bytes()).unwrap();
     55         std::future::pending::<()>().await;
     56         // Keep the actual uncommitted transaction alive until process termination.
     57         transaction.rollback().await.unwrap();
     58     }
     59 }
     60 
     61 async fn kill_and_reopen(phase: &str) {
     62     let temp = TempDir::new().unwrap();
     63     let store = open(temp.path(), OpenMode::Create).await;
     64     let baseline = first();
     65     store
     66         .append_authored_draft(baseline.clone(), None)
     67         .await
     68         .unwrap();
     69     store.close().await.unwrap();
     70     let output = fs::File::create(temp.path().join("child-output")).unwrap();
     71     let mut child = ChildGuard(
     72         Command::new(std::env::current_exe().unwrap())
     73             .args(["--exact", CHILD_TEST, "--nocapture", "--test-threads=1"])
     74             .env(CHILD_ROOT, temp.path())
     75             .env(CHILD_PHASE, phase)
     76             .stdin(Stdio::null())
     77             .stdout(output.try_clone().unwrap())
     78             .stderr(output)
     79             .spawn()
     80             .unwrap(),
     81     );
     82     let deadline = Instant::now() + Duration::from_secs(30);
     83     loop {
     84         assert!(
     85             child.0.try_wait().unwrap().is_none(),
     86             "child exited before ready: {}",
     87             fs::read_to_string(temp.path().join("child-output")).unwrap()
     88         );
     89         if fs::read(temp.path().join("child-ready")).ok().as_deref() == Some(phase.as_bytes()) {
     90             break;
     91         }
     92         assert!(
     93             Instant::now() < deadline,
     94             "child did not reach exact write boundary"
     95         );
     96         std::thread::sleep(Duration::from_millis(10));
     97     }
     98     assert!(temp.path().join("runtime.sqlite-wal").is_file());
     99     child.0.kill().unwrap();
    100     assert!(!child.0.wait().unwrap().success());
    101     let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await;
    102     let pending = next(&baseline, b"complete child revision".to_vec());
    103     let expected = if phase == "after_acknowledgment" {
    104         &pending
    105     } else {
    106         &baseline
    107     };
    108     assert_head(&reopened, expected).await;
    109     assert_eq!(
    110         reopened
    111             .authored_draft_revision(baseline.draft_id(), baseline.revision())
    112             .await
    113             .unwrap(),
    114         Some(baseline.clone())
    115     );
    116     if phase == "during_uncommitted_insert" {
    117         assert!(
    118             reopened
    119                 .authored_draft_revision(pending.draft_id(), pending.revision())
    120                 .await
    121                 .unwrap()
    122                 .is_none()
    123         );
    124     }
    125     let receipt = reopened
    126         .append_authored_draft(pending.clone(), Some(baseline.revision()))
    127         .await
    128         .unwrap();
    129     assert_eq!(
    130         receipt.disposition(),
    131         if phase == "after_acknowledgment" {
    132             DraftAppendDisposition::Replay
    133         } else {
    134             DraftAppendDisposition::Inserted
    135         }
    136     );
    137     assert_head(&reopened, &pending).await;
    138     reopened.close().await.unwrap();
    139 }
    140 
    141 #[tokio::test]
    142 async fn authored_durability_acknowledged_revision_survives_process_kill() {
    143     kill_and_reopen("after_acknowledgment").await;
    144 }
    145 
    146 #[tokio::test]
    147 async fn authored_durability_uncommitted_revision_is_absent_after_process_kill() {
    148     kill_and_reopen("during_uncommitted_insert").await;
    149 }