authored_draft_pair.rs (7722B)
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 #[test] 127 fn memory_atomic_pair_conflict_replay_and_successor_contract() { 128 futures_executor::block_on(contract(&radroots_storage::memory::MemoryStorage::default())); 129 } 130 #[test] 131 fn pair_rejects_duplicate_keys_and_foreign_author() { 132 let a = draft(1, b"a"); 133 assert_eq!( 134 AuthoredDraftPair::new(a.clone(), None, a.clone(), None), 135 Err(Error::InvalidAuthoredDraft) 136 ); 137 let foreign = AuthoredDraft::initial( 138 AuthoredDraftId::new([2; 16]).unwrap(), 139 [8; 32], 140 "fixture.pair.v1", 141 b"foreign".to_vec(), 142 AuthoredDraftStage::Draft, 143 None, 144 10, 145 ) 146 .unwrap(); 147 assert_eq!( 148 AuthoredDraftPair::new(a, None, foreign, None), 149 Err(Error::InvalidAuthoredDraft) 150 ); 151 } 152 #[test] 153 fn concurrent_memory_reservations_commit_one_complete_pair() { 154 use std::sync::{Arc, Barrier}; 155 let store = Arc::new(radroots_storage::memory::MemoryStorage::default()); 156 let barrier = Arc::new(Barrier::new(2)); 157 let tasks: Vec<_> = 158 (2..=3) 159 .map(|id| { 160 let store = store.clone(); 161 let barrier = barrier.clone(); 162 std::thread::spawn(move || { 163 barrier.wait(); 164 ( 165 id, 166 futures_executor::block_on(store.append_authored_draft_pair(pair( 167 draft(1, &[id]), 168 draft(id, b"capture"), 169 ))), 170 ) 171 }) 172 }) 173 .collect(); 174 let results: Vec<_> = tasks.into_iter().map(|t| t.join().unwrap()).collect(); 175 assert_eq!(results.iter().filter(|(_, v)| v.is_ok()).count(), 1); 176 for (id, result) in results { 177 assert_eq!( 178 futures_executor::block_on(store.authored_draft_head(draft(id, b"lookup").draft_id())) 179 .unwrap() 180 .is_some(), 181 result.is_ok() 182 ); 183 } 184 } 185 186 struct Unsupported(radroots_storage::memory::MemoryStorage); 187 impl AuthoredDraftStore for Unsupported { 188 fn query_authored_drafts( 189 &self, 190 q: radroots_storage::authored_draft_query::AuthoredDraftQuery, 191 ) -> radroots_transport::BoxFuture< 192 '_, 193 Result<radroots_storage::authored_draft_query::AuthoredDraftPage, Error>, 194 > { 195 self.0.query_authored_drafts(q) 196 } 197 fn append_authored_draft( 198 &self, 199 d: AuthoredDraft, 200 e: Option<radroots_storage::authored_draft::AuthoredDraftRevision>, 201 ) -> radroots_transport::BoxFuture< 202 '_, 203 Result<radroots_storage::authored_draft::DraftAppendReceipt, Error>, 204 > { 205 self.0.append_authored_draft(d, e) 206 } 207 fn authored_draft_head( 208 &self, 209 id: AuthoredDraftId, 210 ) -> radroots_transport::BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { 211 self.0.authored_draft_head(id) 212 } 213 fn authored_draft_revision( 214 &self, 215 id: AuthoredDraftId, 216 r: radroots_storage::authored_draft::AuthoredDraftRevision, 217 ) -> radroots_transport::BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { 218 self.0.authored_draft_revision(id, r) 219 } 220 fn authored_draft_heads( 221 &self, 222 a: [u8; 32], 223 n: u16, 224 ) -> radroots_transport::BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>> { 225 self.0.authored_draft_heads(a, n) 226 } 227 } 228 #[test] 229 fn unsupported_backend_never_falls_back_to_separate_writes() { 230 let store = Unsupported(radroots_storage::memory::MemoryStorage::default()); 231 assert_eq!( 232 futures_executor::block_on( 233 store.append_authored_draft_pair(pair(draft(1, b"a"), draft(2, b"b"))) 234 ), 235 Err(Error::BackendUnavailable) 236 ); 237 assert!( 238 futures_executor::block_on(store.authored_draft_heads([9; 32], 10)) 239 .unwrap() 240 .is_empty() 241 ); 242 }