authored_draft_pair.rs (6605B)
1 use radroots_storage::{ 2 Error, 3 authored_draft::{ 4 AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore, 5 DraftAppendDisposition, 6 }, 7 authored_draft_pair::AuthoredDraftPair, 8 }; 9 fn draft(id: u8, payload: &[u8]) -> AuthoredDraft { 10 AuthoredDraft::initial( 11 AuthoredDraftId::new([id; 16]).unwrap(), 12 [9; 32], 13 "fixture.pair.v1", 14 payload.to_vec(), 15 AuthoredDraftStage::Draft, 16 None, 17 10, 18 ) 19 .unwrap() 20 } 21 fn pair(a: AuthoredDraft, b: AuthoredDraft) -> AuthoredDraftPair { 22 AuthoredDraftPair::new(a, None, b, None).unwrap() 23 } 24 async fn contract(store: &dyn AuthoredDraftStore) { 25 let a = draft(1, b"claim"); 26 let b = draft(2, b"captured"); 27 let request = pair(a.clone(), b.clone()); 28 let receipt = store 29 .append_authored_draft_pair(request.clone()) 30 .await 31 .unwrap(); 32 assert_eq!(receipt[0].draft(), &a); 33 assert_eq!(receipt[1].draft(), &b); 34 assert!( 35 receipt 36 .iter() 37 .all(|v| v.disposition() == DraftAppendDisposition::Inserted) 38 ); 39 let replay = store 40 .append_authored_draft_pair(request.clone()) 41 .await 42 .unwrap(); 43 assert!( 44 replay 45 .iter() 46 .all(|v| v.disposition() == DraftAppendDisposition::Replay) 47 ); 48 // First-member conflict and second-member conflict must retain no fresh row. 49 for candidate in [ 50 pair(draft(1, b"changed"), draft(3, b"new")), 51 pair(draft(3, b"new"), draft(2, b"changed")), 52 pair(a.clone(), draft(3, b"new")), 53 pair(draft(3, b"new"), b.clone()), 54 ] { 55 assert_eq!( 56 store.append_authored_draft_pair(candidate).await, 57 Err(Error::DraftRevisionConflict) 58 ); 59 assert!( 60 store 61 .authored_draft_head(draft(3, b"new").draft_id()) 62 .await 63 .unwrap() 64 .is_none() 65 ); 66 } 67 let a2 = a 68 .successor(b"next claim".to_vec(), AuthoredDraftStage::Draft, None, 11) 69 .unwrap(); 70 let b2 = b 71 .successor( 72 b"next capture".to_vec(), 73 AuthoredDraftStage::Draft, 74 None, 75 11, 76 ) 77 .unwrap(); 78 let incorrect = 79 AuthoredDraftPair::new(a2.clone(), Some(a.revision()), b2.clone(), None).unwrap(); 80 assert_eq!( 81 store.append_authored_draft_pair(incorrect).await, 82 Err(Error::DraftRevisionConflict) 83 ); 84 assert_eq!( 85 store.authored_draft_head(a.draft_id()).await.unwrap(), 86 Some(a.clone()) 87 ); 88 let incorrect = 89 AuthoredDraftPair::new(a2.clone(), None, b2.clone(), Some(b.revision())).unwrap(); 90 assert_eq!( 91 store.append_authored_draft_pair(incorrect).await, 92 Err(Error::DraftRevisionConflict) 93 ); 94 let next = AuthoredDraftPair::new( 95 a2.clone(), 96 Some(a.revision()), 97 b2.clone(), 98 Some(b.revision()), 99 ) 100 .unwrap(); 101 store.append_authored_draft_pair(next).await.unwrap(); 102 // Replay returns immutable history and never rewinds the current head. 103 let replay = store.append_authored_draft_pair(request).await.unwrap(); 104 assert_eq!(replay[0].draft(), &a); 105 assert_eq!( 106 store.authored_draft_head(a.draft_id()).await.unwrap(), 107 Some(a2) 108 ); 109 assert_eq!( 110 store.authored_draft_head(b.draft_id()).await.unwrap(), 111 Some(b2) 112 ); 113 // The untouched single-row API retains its original behavior. 114 let c = draft(4, b"single"); 115 store.append_authored_draft(c.clone(), None).await.unwrap(); 116 assert_eq!( 117 store 118 .append_authored_draft(c, None) 119 .await 120 .unwrap() 121 .disposition(), 122 DraftAppendDisposition::Replay 123 ); 124 } 125 126 use radroots_storage::event::SourceGeneration; 127 use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage}; 128 async fn open(root: &std::path::Path, mode: OpenMode) -> SqliteStorage { 129 let options = OpenOptions::new(Paths::from_directory(root).unwrap(), mode); 130 let options = if mode == OpenMode::Create { 131 options 132 .with_source_generation(SourceGeneration::new([7; 32]).unwrap(), 10) 133 .unwrap() 134 } else { 135 options 136 }; 137 SqliteStorage::open(options).await.unwrap() 138 } 139 #[tokio::test] 140 async fn sqlite_pair_contract_persists_across_reopen_and_read_only_refuses() { 141 let dir = tempfile::tempdir().unwrap(); 142 let store = open(dir.path(), OpenMode::Create).await; 143 contract(&store).await; 144 store.close().await.unwrap(); 145 let store = open(dir.path(), OpenMode::ReadOnly).await; 146 for id in [1, 2] { 147 assert_eq!( 148 store 149 .authored_draft_head(draft(id, b"lookup").draft_id()) 150 .await 151 .unwrap() 152 .unwrap() 153 .revision() 154 .get(), 155 2 156 ); 157 } 158 assert_eq!( 159 store 160 .append_authored_draft_pair(pair(draft(5, b"a"), draft(6, b"b"))) 161 .await, 162 Err(Error::BackendUnavailable) 163 ); 164 store.close().await.unwrap(); 165 let store = open(dir.path(), OpenMode::ReadWriteExisting).await; 166 let replay = store 167 .append_authored_draft_pair(pair(draft(1, b"claim"), draft(2, b"captured"))) 168 .await 169 .unwrap(); 170 assert!( 171 replay 172 .iter() 173 .all(|v| v.disposition() == DraftAppendDisposition::Replay) 174 ); 175 store.close().await.unwrap(); 176 } 177 #[tokio::test] 178 async fn concurrent_sqlite_reservations_keep_exactly_one_complete_pair() { 179 let dir = tempfile::tempdir().unwrap(); 180 let store = open(dir.path(), OpenMode::Create).await; 181 let (a, b) = tokio::join!( 182 store.append_authored_draft_pair(pair(draft(1, b"a"), draft(2, b"capture a"))), 183 store.append_authored_draft_pair(pair(draft(1, b"b"), draft(3, b"capture b"))) 184 ); 185 assert_ne!(a.is_ok(), b.is_ok()); 186 assert_eq!( 187 store 188 .authored_draft_head(draft(2, b"lookup").draft_id()) 189 .await 190 .unwrap() 191 .is_some(), 192 a.is_ok() 193 ); 194 assert_eq!( 195 store 196 .authored_draft_head(draft(3, b"lookup").draft_id()) 197 .await 198 .unwrap() 199 .is_some(), 200 b.is_ok() 201 ); 202 store.close().await.unwrap(); 203 let store = open(dir.path(), OpenMode::ReadWriteExisting).await; 204 assert_eq!( 205 store.authored_draft_heads([9; 32], 10).await.unwrap().len(), 206 2 207 ); 208 store.close().await.unwrap(); 209 }