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 }