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 }