lib

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

support.rs (11082B)


      1 use std::sync::{
      2     Arc, Mutex,
      3     atomic::{AtomicBool, Ordering},
      4 };
      5 
      6 use futures::future;
      7 use radroots_transport::{
      8     BoxFuture, DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource, FetchPage,
      9     FetchRequest, SinkFailure, SinkStatus, SourceStatus, TransportId,
     10     capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
     11     outcome::{DeliveryOutcome, FetchTargetOutcome, FetchTargetState},
     12     sink::DeliveryTargetReceipt,
     13     source::NextPage,
     14 };
     15 
     16 use crate::suite::{NOW_UNIX_MS, SinkConformanceHarness, SourceConformanceHarness};
     17 
     18 #[derive(Clone)]
     19 enum Mode {
     20     Success,
     21     Fail(Error),
     22     Pending,
     23 }
     24 
     25 #[derive(Default)]
     26 pub(crate) struct SourceState {
     27     request: Mutex<Option<FetchRequest>>,
     28     pub(crate) published: AtomicBool,
     29     pub(crate) cancelled_after_publish: AtomicBool,
     30 }
     31 
     32 impl SourceState {
     33     pub(crate) fn request(&self) -> Option<FetchRequest> {
     34         self.request.lock().expect("source request lock").clone()
     35     }
     36 }
     37 
     38 #[derive(Default)]
     39 pub(crate) struct SinkState {
     40     request: Mutex<Option<DeliveryRequest>>,
     41     pub(crate) published: AtomicBool,
     42     pub(crate) cancelled_after_publish: AtomicBool,
     43 }
     44 
     45 impl SinkState {
     46     pub(crate) fn request(&self) -> Option<DeliveryRequest> {
     47         self.request.lock().expect("sink request lock").clone()
     48     }
     49 }
     50 
     51 pub(crate) struct MockSource {
     52     mode: Mode,
     53     state: Arc<SourceState>,
     54 }
     55 
     56 impl MockSource {
     57     pub(crate) fn successful() -> Self {
     58         Self::new(Mode::Success)
     59     }
     60 
     61     pub(crate) fn failing(error: Error) -> Self {
     62         Self::new(Mode::Fail(error))
     63     }
     64 
     65     pub(crate) fn pending() -> Self {
     66         Self::new(Mode::Pending)
     67     }
     68 
     69     fn new(mode: Mode) -> Self {
     70         Self {
     71             mode,
     72             state: Arc::new(SourceState::default()),
     73         }
     74     }
     75 }
     76 
     77 impl EventSource for MockSource {
     78     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> {
     79         Box::pin(async { Ok(source_status()) })
     80     }
     81 
     82     fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> {
     83         let state = Arc::clone(&self.state);
     84         let mode = self.mode.clone();
     85         Box::pin(async move {
     86             *state.request.lock().expect("source request lock") = Some(request.clone());
     87             if request.bounds().deadline_unix_ms() <= NOW_UNIX_MS {
     88                 return Err(Error::InvalidFetchDeadline);
     89             }
     90             match mode {
     91                 Mode::Fail(error) => Err(error),
     92                 Mode::Pending => {
     93                     state.published.store(true, Ordering::SeqCst);
     94                     let state_for_drop = Arc::clone(&state);
     95                     let _publish_guard = ScopeGuard::new(move || {
     96                         state_for_drop
     97                             .cancelled_after_publish
     98                             .store(true, Ordering::SeqCst);
     99                     });
    100                     future::pending().await
    101                 }
    102                 Mode::Success => {
    103                     let outcomes = request
    104                         .target_set()
    105                         .targets()
    106                         .iter()
    107                         .enumerate()
    108                         .map(|(index, target)| {
    109                             FetchTargetOutcome::new(
    110                                 target.fingerprint().clone(),
    111                                 if index == 0 {
    112                                     FetchTargetState::Complete
    113                                 } else {
    114                                     FetchTargetState::FailedRetryable
    115                                 },
    116                             )
    117                         })
    118                         .collect();
    119                     FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete)
    120                 }
    121             }
    122         })
    123     }
    124 }
    125 
    126 impl SourceConformanceHarness for MockSource {
    127     fn source(&self) -> &dyn EventSource {
    128         self
    129     }
    130 
    131     fn expected_status(&self) -> SourceStatus {
    132         source_status()
    133     }
    134 
    135     fn target_set(&self) -> radroots_transport::TargetSet {
    136         target_set()
    137     }
    138 
    139     fn now_unix_ms(&self) -> u64 {
    140         NOW_UNIX_MS
    141     }
    142 
    143     fn captured_request(&self) -> Option<FetchRequest> {
    144         self.state.request()
    145     }
    146 
    147     fn published(&self) -> bool {
    148         self.state.published.load(Ordering::SeqCst)
    149     }
    150 
    151     fn cancelled_after_publish(&self) -> bool {
    152         self.state.cancelled_after_publish.load(Ordering::SeqCst)
    153     }
    154 }
    155 
    156 pub(crate) struct MockSink {
    157     mode: Mode,
    158     state: Arc<SinkState>,
    159 }
    160 
    161 impl MockSink {
    162     pub(crate) fn successful() -> Self {
    163         Self::new(Mode::Success)
    164     }
    165 
    166     pub(crate) fn failing(error: Error) -> Self {
    167         Self::new(Mode::Fail(error))
    168     }
    169 
    170     pub(crate) fn pending() -> Self {
    171         Self::new(Mode::Pending)
    172     }
    173 
    174     fn new(mode: Mode) -> Self {
    175         Self {
    176             mode,
    177             state: Arc::new(SinkState::default()),
    178         }
    179     }
    180 }
    181 
    182 impl EventSink for MockSink {
    183     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
    184         Box::pin(async { Ok(sink_status()) })
    185     }
    186 
    187     fn deliver(
    188         &self,
    189         request: DeliveryRequest,
    190     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    191         let state = Arc::clone(&self.state);
    192         let mode = self.mode.clone();
    193         Box::pin(async move {
    194             *state.request.lock().expect("sink request lock") = Some(request.clone());
    195             if request.deadline_unix_ms() <= NOW_UNIX_MS {
    196                 return Err(SinkFailure::invalid_contract(&request));
    197             }
    198             match mode {
    199                 Mode::Fail(_) => Err(SinkFailure::invalid_contract(&request)),
    200                 Mode::Pending => {
    201                     state.published.store(true, Ordering::SeqCst);
    202                     let state_for_drop = Arc::clone(&state);
    203                     let _publish_guard = ScopeGuard::new(move || {
    204                         state_for_drop
    205                             .cancelled_after_publish
    206                             .store(true, Ordering::SeqCst);
    207                     });
    208                     future::pending().await
    209                 }
    210                 Mode::Success => {
    211                     let receipts = request
    212                         .target_set()
    213                         .targets()
    214                         .iter()
    215                         .enumerate()
    216                         .map(|(index, target)| {
    217                             DeliveryTargetReceipt::attempted(
    218                                 target.clone(),
    219                                 if index == 0 {
    220                                     DeliveryOutcome::delivered()
    221                                 } else {
    222                                     DeliveryOutcome::unavailable()
    223                                 },
    224                             )
    225                         })
    226                         .collect();
    227                     DeliveryReceipt::for_request(&request, receipts)
    228                         .map_err(|_| SinkFailure::invalid_contract(&request))
    229                 }
    230             }
    231         })
    232     }
    233 }
    234 
    235 impl SinkConformanceHarness for MockSink {
    236     fn sink(&self) -> &dyn EventSink {
    237         self
    238     }
    239 
    240     fn expected_status(&self) -> SinkStatus {
    241         sink_status()
    242     }
    243 
    244     fn target_set(&self) -> radroots_transport::TargetSet {
    245         target_set()
    246     }
    247 
    248     fn now_unix_ms(&self) -> u64 {
    249         NOW_UNIX_MS
    250     }
    251 
    252     fn captured_request(&self) -> Option<DeliveryRequest> {
    253         self.state.request()
    254     }
    255 
    256     fn published(&self) -> bool {
    257         self.state.published.load(Ordering::SeqCst)
    258     }
    259 
    260     fn cancelled_after_publish(&self) -> bool {
    261         self.state.cancelled_after_publish.load(Ordering::SeqCst)
    262     }
    263 }
    264 
    265 struct ScopeGuard<F: FnOnce()>(Option<F>);
    266 
    267 impl<F: FnOnce()> ScopeGuard<F> {
    268     fn new(callback: F) -> Self {
    269         Self(Some(callback))
    270     }
    271 }
    272 
    273 impl<F: FnOnce()> Drop for ScopeGuard<F> {
    274     fn drop(&mut self) {
    275         if let Some(callback) = self.0.take() {
    276             callback();
    277         }
    278     }
    279 }
    280 
    281 pub(crate) struct CombinedAdapter {
    282     source: MockSource,
    283     sink: MockSink,
    284 }
    285 
    286 impl CombinedAdapter {
    287     pub(crate) fn successful() -> Self {
    288         Self {
    289             source: MockSource::successful(),
    290             sink: MockSink::successful(),
    291         }
    292     }
    293 }
    294 
    295 impl EventSource for CombinedAdapter {
    296     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> {
    297         EventSource::status(&self.source)
    298     }
    299 
    300     fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> {
    301         self.source.fetch(request)
    302     }
    303 }
    304 
    305 impl EventSink for CombinedAdapter {
    306     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
    307         EventSink::status(&self.sink)
    308     }
    309 
    310     fn deliver(
    311         &self,
    312         request: DeliveryRequest,
    313     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    314         self.sink.deliver(request)
    315     }
    316 }
    317 
    318 impl SourceConformanceHarness for CombinedAdapter {
    319     fn source(&self) -> &dyn EventSource {
    320         self
    321     }
    322 
    323     fn expected_status(&self) -> SourceStatus {
    324         source_status()
    325     }
    326 
    327     fn target_set(&self) -> radroots_transport::TargetSet {
    328         target_set()
    329     }
    330 
    331     fn now_unix_ms(&self) -> u64 {
    332         NOW_UNIX_MS
    333     }
    334 
    335     fn captured_request(&self) -> Option<FetchRequest> {
    336         self.source.state.request()
    337     }
    338 
    339     fn published(&self) -> bool {
    340         self.source.state.published.load(Ordering::SeqCst)
    341     }
    342 
    343     fn cancelled_after_publish(&self) -> bool {
    344         self.source
    345             .state
    346             .cancelled_after_publish
    347             .load(Ordering::SeqCst)
    348     }
    349 }
    350 
    351 impl SinkConformanceHarness for CombinedAdapter {
    352     fn sink(&self) -> &dyn EventSink {
    353         self
    354     }
    355 
    356     fn expected_status(&self) -> SinkStatus {
    357         sink_status()
    358     }
    359 
    360     fn target_set(&self) -> radroots_transport::TargetSet {
    361         target_set()
    362     }
    363 
    364     fn now_unix_ms(&self) -> u64 {
    365         NOW_UNIX_MS
    366     }
    367 
    368     fn captured_request(&self) -> Option<DeliveryRequest> {
    369         self.sink.state.request()
    370     }
    371 
    372     fn published(&self) -> bool {
    373         self.sink.state.published.load(Ordering::SeqCst)
    374     }
    375 
    376     fn cancelled_after_publish(&self) -> bool {
    377         self.sink
    378             .state
    379             .cancelled_after_publish
    380             .load(Ordering::SeqCst)
    381     }
    382 }
    383 
    384 fn target_set() -> radroots_transport::TargetSet {
    385     radroots_transport::TargetSet::new(vec![
    386         radroots_transport::Target::local("local:conformance-a").expect("first target"),
    387         radroots_transport::Target::local("local:conformance-b").expect("second target"),
    388     ])
    389     .expect("target set")
    390 }
    391 
    392 fn source_status() -> SourceStatus {
    393     SourceStatus::new(
    394         TransportId::LOCAL,
    395         true,
    396         Maturity::Stable,
    397         Availability::Available,
    398         SourceCapabilities::FETCH,
    399         "mock source ready",
    400     )
    401 }
    402 
    403 fn sink_status() -> SinkStatus {
    404     SinkStatus::new(
    405         TransportId::LOCAL,
    406         true,
    407         Maturity::Stable,
    408         Availability::Available,
    409         SinkCapabilities::DELIVER,
    410         "mock sink ready",
    411     )
    412 }