lib

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

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 }