lib

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

selected.rs (5256B)


      1 use super::*;
      2 use radroots_transport::{
      3     DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure,
      4     outcome::DeliveryOutcome,
      5     sink::{DeliveryTargetReceipt, SinkStatus},
      6 };
      7 
      8 #[derive(Default)]
      9 struct SelectedSink(std::sync::Mutex<Vec<(DeliveryRequest, TargetSet)>>);
     10 
     11 impl EventSink for SelectedSink {
     12     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
     13         Box::pin(async { panic!("selected delivery does not probe status") })
     14     }
     15 
     16     fn deliver(
     17         &self,
     18         _: DeliveryRequest,
     19     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     20         Box::pin(async { panic!("selected delivery must not widen to ordinary delivery") })
     21     }
     22 
     23     fn deliver_selected(
     24         &self,
     25         request: DeliveryRequest,
     26         selected: TargetSet,
     27     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     28         Box::pin(async move {
     29             request.validate_target_selection(&selected).unwrap();
     30             self.0
     31                 .lock()
     32                 .unwrap()
     33                 .push((request.clone(), selected.clone()));
     34             Ok(DeliveryReceipt::for_request(
     35                 &request,
     36                 request
     37                     .target_set()
     38                     .targets()
     39                     .iter()
     40                     .map(|target| {
     41                         if selected.targets().contains(target) {
     42                             DeliveryTargetReceipt::attempted(
     43                                 target.clone(),
     44                                 DeliveryOutcome::accepted(),
     45                             )
     46                         } else {
     47                             DeliveryTargetReceipt::skipped(
     48                                 target.clone(),
     49                                 DeliveryOutcome::unavailable(),
     50                             )
     51                             .unwrap()
     52                         }
     53                     })
     54                     .collect(),
     55             )
     56             .unwrap())
     57         })
     58     }
     59 }
     60 
     61 #[tokio::test]
     62 async fn sdk_selected_delivery_preserves_subset_and_full_durable_request() {
     63     let storage = Arc::new(MemoryStorage::new(
     64         SourceGeneration::new([211; 32]).unwrap(),
     65     ));
     66     let signer = radroots_nostr::signing::LocalSigner::new(
     67         radroots_nostr::key::SecretKey::parse(
     68             "0000000000000000000000000000000000000000000000000000000000000001",
     69         )
     70         .unwrap(),
     71     )
     72     .unwrap();
     73     let author = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798";
     74     let (clock, ids, deadlines) = HostPolicy::default().composition();
     75     let now = clock.now_unix_ms().unwrap();
     76     let sink = Arc::new(SelectedSink::default());
     77     let engine = Engine::builder(storage.clone(), clock, ids, deadlines)
     78         .signer(Arc::new(signer))
     79         .sink(sink.clone())
     80         .build()
     81         .unwrap();
     82     let client = ClientBuilder::new()
     83         .storage(storage)
     84         .sync_engine(engine)
     85         .build()
     86         .unwrap();
     87     let operations = client.sync().unwrap().unwrap();
     88     let targets = TargetSet::new(vec![
     89         target(),
     90         Target::nostr_relay("wss://held.example").unwrap(),
     91     ])
     92     .unwrap();
     93     let selected = TargetSet::new(vec![targets.targets()[0].clone()]).unwrap();
     94     let request = PushRequest::new(
     95         SyncId::new([211; 16]).unwrap(),
     96         IdempotencyKey::parse("sdk-selected-delivery").unwrap(),
     97         Actor::new(
     98             PublicKey::from_hex(author).unwrap(),
     99             ActorSource::ExplicitPublicKey,
    100             [AuthorRole::Any],
    101         )
    102         .unwrap(),
    103         AuthoredEventPlan::from_generic(
    104             GenericEventDraft::new(
    105                 "radroots.social.geochat.v1",
    106                 20_000,
    107                 now / 1_000,
    108                 Vec::new(),
    109                 "selected",
    110                 author,
    111             )
    112             .unwrap(),
    113         )
    114         .unwrap(),
    115         targets.clone(),
    116         SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    117         now + 60_000,
    118         CancellationPolicy::LocalCooperative,
    119     )
    120     .unwrap();
    121     let id = request.operation_id();
    122     operations.submit_push(request).await.unwrap();
    123     let original = operations.push_status(id).await.unwrap().unwrap();
    124     let original_request = original.delivery_plan().request().unwrap();
    125     assert_eq!(
    126         operations
    127             .deliver_push_selected(
    128                 id,
    129                 TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()])
    130                     .unwrap(),
    131             )
    132             .await
    133             .unwrap_err(),
    134         Error::InvalidDeliveryRequest
    135     );
    136     assert!(sink.0.lock().unwrap().is_empty());
    137     operations
    138         .deliver_push_selected(id, selected.clone())
    139         .await
    140         .unwrap();
    141     let after = operations.push_status(id).await.unwrap().unwrap();
    142     assert_eq!(after.delivery_plan().request(), Some(original_request));
    143     assert_eq!(
    144         sink.0.lock().unwrap().as_slice(),
    145         &[(original_request.clone(), selected)]
    146     );
    147     assert!(!after.delivery_plan().state().is_terminal());
    148     assert_eq!(after.delivery_plan().delivery_facts().len(), 1);
    149     client.close().await.unwrap();
    150 }