lib

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

delivery_selection.rs (11481B)


      1 use super::*;
      2 use futures::channel::oneshot;
      3 use radroots_storage::authored_delivery::DeliveryAttemptOutcome;
      4 use radroots_transport::policy::SatisfactionState;
      5 
      6 type Pending = (
      7     DeliveryRequest,
      8     TargetSet,
      9     oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>,
     10 );
     11 
     12 #[derive(Default)]
     13 struct SelectedSink {
     14     pending: Mutex<VecDeque<Pending>>,
     15     calls: AtomicUsize,
     16 }
     17 
     18 impl EventSink for SelectedSink {
     19     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
     20         Box::pin(async { panic!("no status probe") })
     21     }
     22     fn deliver(
     23         &self,
     24         _: DeliveryRequest,
     25     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     26         Box::pin(async { panic!("selected attempts must not fall back") })
     27     }
     28     fn deliver_selected(
     29         &self,
     30         request: DeliveryRequest,
     31         selected: TargetSet,
     32     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     33         Box::pin(async move {
     34             request.validate_target_selection(&selected).unwrap();
     35             self.calls.fetch_add(1, Ordering::SeqCst);
     36             let (send, receive) = oneshot::channel();
     37             self.pending
     38                 .lock()
     39                 .unwrap()
     40                 .push_back((request, selected, send));
     41             receive.await.unwrap()
     42         })
     43     }
     44 }
     45 
     46 fn selected_receipt(request: &DeliveryRequest, selected: &TargetSet) -> DeliveryReceipt {
     47     DeliveryReceipt::for_request(
     48         request,
     49         request
     50             .target_set()
     51             .targets()
     52             .iter()
     53             .map(|target| {
     54                 if selected.targets().contains(target) {
     55                     DeliveryTargetReceipt::attempted(target.clone(), DeliveryOutcome::accepted())
     56                 } else {
     57                     DeliveryTargetReceipt::skipped(target.clone(), DeliveryOutcome::unavailable())
     58                         .unwrap()
     59                 }
     60             })
     61             .collect(),
     62     )
     63     .unwrap()
     64 }
     65 
     66 fn poll_pending(future: &mut (impl std::future::Future + Unpin)) {
     67     assert!(
     68         future
     69             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
     70             .is_pending()
     71     );
     72 }
     73 
     74 fn engine(storage: Arc<dyn SyncStorage>, clock: Arc<TestClock>, sink: Arc<SelectedSink>) -> Engine {
     75     Engine::builder(
     76         storage,
     77         clock,
     78         Arc::new(TestIds(AtomicU64::new(10))),
     79         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
     80     )
     81     .sink(sink)
     82     .signer(Arc::new(MockSigner::new(SignBehavior::Success {
     83         completed_at_unix_ms: 1_800_000_200_500,
     84     })))
     85     .build()
     86     .unwrap()
     87 }
     88 
     89 fn push() -> PushRequest {
     90     request_with_policy(
     91         211,
     92         &["wss://one.example", "wss://two.example"],
     93         SatisfactionClass::Accepted,
     94         TargetPolicy::all(),
     95     )
     96 }
     97 
     98 #[test]
     99 fn invalid_selection_cannot_claim_and_out_of_selection_results_fail_closed() {
    100     for failure in [false, true] {
    101         let storage = Arc::new(MemoryStorage::new(
    102             SourceGeneration::new([211; 32]).unwrap(),
    103         ));
    104         let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
    105         let sink = Arc::new(SelectedSink::default());
    106         let engine = engine(storage, clock, sink.clone());
    107         let push = push();
    108         execute_to_admitted(&engine, &push);
    109         let before = block_on(engine.push_status(push.operation_id()))
    110             .unwrap()
    111             .unwrap();
    112         let original = before.delivery_plan().request().unwrap();
    113         let foreign =
    114             TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap();
    115         assert!(matches!(
    116             block_on(engine.deliver_push_selected(push.operation_id(), foreign)),
    117             Err(Error::InvalidDeliveryRequest)
    118         ));
    119         let unchanged = block_on(engine.push_status(push.operation_id()))
    120             .unwrap()
    121             .unwrap();
    122         assert_eq!(unchanged.delivery_plan(), before.delivery_plan());
    123         assert_eq!(sink.calls.load(Ordering::SeqCst), 0);
    124         let selected = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap();
    125         let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected));
    126         poll_pending(&mut future);
    127         let (request, _, sender) = sink.pending.lock().unwrap().pop_front().unwrap();
    128         assert_eq!(&request, original);
    129         let forbidden = DeliveryTargetReceipt::attempted(
    130             request.target_set().targets()[1].clone(),
    131             DeliveryOutcome::accepted(),
    132         );
    133         let result = if failure {
    134             Err(SinkFailure::for_request(
    135                 &request,
    136                 "upstream_failure",
    137                 Retryability::Retryable,
    138                 None,
    139                 None,
    140                 vec![forbidden],
    141             )
    142             .unwrap())
    143         } else {
    144             Ok(receipt(&request, vec![DeliveryOutcome::accepted(); 2]).unwrap())
    145         };
    146         sender.send(result).unwrap();
    147         let result = block_on(future).unwrap();
    148         assert_eq!(result.plan().request(), Some(original));
    149         let DeliveryAttemptOutcome::SinkFailure(failure) =
    150             result.plan().delivery_facts()[0].outcome()
    151         else {
    152             panic!("invalid adapter evidence is never acceptance")
    153         };
    154         assert_eq!(failure.code(), "invalid_transport_contract");
    155         assert!(failure.partial_evidence().is_empty());
    156         assert_ne!(
    157             result.plan().delivery_satisfaction().unwrap(),
    158             SatisfactionState::Satisfied
    159         );
    160     }
    161 }
    162 
    163 #[test]
    164 fn selected_late_acceptance_survives_stop_without_scheduling_new_targets() {
    165     let storage = Arc::new(MemoryStorage::new(
    166         SourceGeneration::new([212; 32]).unwrap(),
    167     ));
    168     let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
    169     let sink = Arc::new(SelectedSink::default());
    170     let engine = engine(storage.clone(), clock.clone(), sink.clone());
    171     let push = push();
    172     execute_to_admitted(&engine, &push);
    173     let before = block_on(engine.push_status(push.operation_id()))
    174         .unwrap()
    175         .unwrap();
    176     let request = before.delivery_plan().request().unwrap();
    177     let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap();
    178     let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected.clone()));
    179     poll_pending(&mut future);
    180     block_on(engine.cancel_push(push.operation_id())).unwrap();
    181     let (captured, captured_selection, sender) = sink.pending.lock().unwrap().pop_front().unwrap();
    182     assert_eq!(&captured, request);
    183     assert_eq!(captured_selection, selected);
    184     let expected = selected_receipt(&captured, &selected);
    185     sender.send(Ok(expected.clone())).unwrap();
    186     let late = block_on(future).unwrap();
    187     assert_eq!(late.plan().state(), AuthoredDeliveryState::Cancelled);
    188     assert_eq!(late.plan().attempt_count(), 0);
    189     assert_eq!(
    190         late.plan().delivery_facts()[0].outcome(),
    191         &DeliveryAttemptOutcome::Receipt(expected)
    192     );
    193     let recovery = delivery_evidence::source_only(storage, clock);
    194     let replay = block_on(recovery.deliver_push_selected(push.operation_id(), selected)).unwrap();
    195     assert!(replay.is_replay());
    196     assert_eq!(replay.plan(), late.plan());
    197     assert_eq!(sink.calls.load(Ordering::SeqCst), 1);
    198 }
    199 
    200 #[tokio::test]
    201 async fn sqlite_selected_facts_survive_reopen_and_next_target_finishes_same_request() {
    202     let directory = tempfile::tempdir().unwrap();
    203     let paths = Paths::from_directory(directory.path()).unwrap();
    204     let storage = Arc::new(
    205         SqliteStorage::open(
    206             OpenOptions::new(paths.clone(), OpenMode::Create)
    207                 .with_source_generation(SourceGeneration::new([213; 32]).unwrap(), 1)
    208                 .unwrap(),
    209         )
    210         .await
    211         .unwrap(),
    212     );
    213     let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
    214     let sink = Arc::new(SelectedSink::default());
    215     let first_engine = engine(storage.clone(), clock.clone(), sink.clone());
    216     let push = push();
    217     first_engine.sign_prepared(push.clone()).await.unwrap();
    218     first_engine
    219         .admit_signed(push.operation_id())
    220         .await
    221         .unwrap();
    222     let before = first_engine
    223         .push_status(push.operation_id())
    224         .await
    225         .unwrap()
    226         .unwrap();
    227     let original = before.delivery_plan().request().unwrap().clone();
    228     let selected_a = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap();
    229     let first = {
    230         let mut future =
    231             Box::pin(first_engine.deliver_push_selected(push.operation_id(), selected_a.clone()));
    232         // SQLite storage needs an executor turn before the sink is reached.
    233         let responder = async {
    234             let pending = loop {
    235                 if let Some(pending) = sink.pending.lock().unwrap().pop_front() {
    236                     break pending;
    237                 }
    238                 tokio::task::yield_now().await;
    239             };
    240             assert_eq!(pending.0, original);
    241             assert_eq!(pending.1, selected_a);
    242             pending
    243                 .2
    244                 .send(Ok(selected_receipt(&pending.0, &pending.1)))
    245                 .unwrap();
    246         };
    247         tokio::time::timeout(std::time::Duration::from_secs(10), async {
    248             let (result, ()) = tokio::join!(&mut future, responder);
    249             result.unwrap()
    250         })
    251         .await
    252         .unwrap()
    253     };
    254     assert_ne!(first.plan().state(), AuthoredDeliveryState::Satisfied);
    255     assert_eq!(first.plan().attempt_count(), 1);
    256     clock.0.store(
    257         first.plan().retry().unwrap().not_before_unix_ms(),
    258         Ordering::SeqCst,
    259     );
    260     drop(first_engine);
    261     storage.close().await.unwrap();
    262     let storage = Arc::new(
    263         SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
    264             .await
    265             .unwrap(),
    266     );
    267     let second_engine = engine(storage.clone(), clock.clone(), sink.clone());
    268     let recovered = second_engine
    269         .push_status(push.operation_id())
    270         .await
    271         .unwrap()
    272         .unwrap();
    273     assert_eq!(
    274         recovered.delivery_plan().delivery_facts(),
    275         first.plan().delivery_facts()
    276     );
    277     let selected_b = TargetSet::new(vec![original.target_set().targets()[1].clone()]).unwrap();
    278     let responder = async {
    279         let pending = loop {
    280             if let Some(pending) = sink.pending.lock().unwrap().pop_front() {
    281                 break pending;
    282             }
    283             tokio::task::yield_now().await;
    284         };
    285         assert_eq!(pending.0, original);
    286         assert_eq!(pending.1, selected_b);
    287         pending
    288             .2
    289             .send(Ok(selected_receipt(&pending.0, &pending.1)))
    290             .unwrap();
    291     };
    292     let finished = tokio::time::timeout(std::time::Duration::from_secs(10), async {
    293         let (result, ()) = tokio::join!(
    294             second_engine.deliver_push_selected(push.operation_id(), selected_b.clone()),
    295             responder
    296         );
    297         result.unwrap()
    298     })
    299     .await
    300     .unwrap();
    301     assert_eq!(finished.plan().request(), Some(&original));
    302     assert_eq!(finished.plan().state(), AuthoredDeliveryState::Satisfied);
    303     assert_eq!(finished.plan().attempt_count(), 2);
    304     assert_eq!(finished.plan().delivery_facts().len(), 2);
    305     assert_eq!(sink.calls.load(Ordering::SeqCst), 2);
    306     drop(second_engine);
    307     storage.close().await.unwrap();
    308 }