lib

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

spi.rs (5785B)


      1 use futures::executor::block_on;
      2 use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
      3 use radroots_transport::{
      4     BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, EventSource, FetchPage, FetchRequest,
      5     SinkStatus, SourceStatus, Target, TargetSet, TransportId,
      6     capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
      7     outcome::DeliveryOutcome,
      8     policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
      9     sink::{DeliveryPayload, DeliveryTargetReceipt},
     10     source::{FetchBounds, NextPage},
     11 };
     12 
     13 struct SourceOnly;
     14 
     15 impl EventSource for SourceOnly {
     16     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
     17         Box::pin(async { Ok(source_status()) })
     18     }
     19 
     20     fn fetch(
     21         &self,
     22         request: FetchRequest,
     23     ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
     24         Box::pin(async move {
     25             FetchPage::for_request(&request, Vec::new(), Vec::new(), NextPage::Complete)
     26         })
     27     }
     28 }
     29 
     30 struct SinkOnly;
     31 
     32 impl EventSink for SinkOnly {
     33     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
     34         Box::pin(async { Ok(sink_status()) })
     35     }
     36 
     37     fn deliver(
     38         &self,
     39         request: DeliveryRequest,
     40     ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>> {
     41         Box::pin(async move {
     42             let receipts = request
     43                 .target_set()
     44                 .targets()
     45                 .iter()
     46                 .cloned()
     47                 .map(|target| {
     48                     DeliveryTargetReceipt::attempted(target, DeliveryOutcome::delivered())
     49                 })
     50                 .collect();
     51             DeliveryReceipt::for_request(&request, receipts)
     52                 .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))
     53         })
     54     }
     55 }
     56 
     57 struct Bidirectional {
     58     source: SourceOnly,
     59     sink: SinkOnly,
     60 }
     61 
     62 impl EventSource for Bidirectional {
     63     fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
     64         EventSource::status(&self.source)
     65     }
     66 
     67     fn fetch(
     68         &self,
     69         request: FetchRequest,
     70     ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
     71         self.source.fetch(request)
     72     }
     73 }
     74 
     75 impl EventSink for Bidirectional {
     76     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
     77         EventSink::status(&self.sink)
     78     }
     79 
     80     fn deliver(
     81         &self,
     82         request: DeliveryRequest,
     83     ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>> {
     84         self.sink.deliver(request)
     85     }
     86 }
     87 
     88 fn source_status() -> SourceStatus {
     89     SourceStatus::new(
     90         TransportId::LOCAL,
     91         true,
     92         Maturity::Stable,
     93         Availability::Available,
     94         SourceCapabilities::FETCH,
     95         "source ready",
     96     )
     97 }
     98 
     99 fn sink_status() -> SinkStatus {
    100     SinkStatus::new(
    101         TransportId::LOCAL,
    102         true,
    103         Maturity::Stable,
    104         Availability::Available,
    105         SinkCapabilities::DELIVER,
    106         "sink ready",
    107     )
    108 }
    109 
    110 fn target_set() -> TargetSet {
    111     TargetSet::new(vec![Target::local("local:spi").expect("local target")]).expect("target set")
    112 }
    113 
    114 fn delivery_payload() -> DeliveryPayload {
    115     let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
    116     let wire = Nip01EventWire::parse_json(raw).expect("wire event");
    117     DeliveryPayload::new(
    118         SignedEvent::from_wire_verified_id(wire, raw).expect("signed delivery event"),
    119     )
    120 }
    121 
    122 fn assert_source_dyn_compatible(_: &dyn EventSource) {}
    123 fn assert_sink_dyn_compatible(_: &dyn EventSink) {}
    124 
    125 #[test]
    126 fn source_only_and_sink_only_implementations_are_independently_dispatchable() {
    127     let source = SourceOnly;
    128     let sink = SinkOnly;
    129     assert_source_dyn_compatible(&source);
    130     assert_sink_dyn_compatible(&sink);
    131 
    132     let source_status = block_on(EventSource::status(&source)).expect("source status");
    133     assert!(source_status.capabilities().can_fetch());
    134     let request = FetchRequest::new(
    135         "fetch-1",
    136         target_set(),
    137         FetchBounds::new(10, 1_700_000_000_000).expect("fetch bounds"),
    138     )
    139     .expect("fetch request");
    140     let page = block_on(source.fetch(request)).expect("fetch page");
    141     assert_eq!(page.request_id().as_str(), "fetch-1");
    142 
    143     let sink_status = block_on(EventSink::status(&sink)).expect("sink status");
    144     assert!(sink_status.capabilities().can_deliver());
    145     let receipt = block_on(
    146         sink.deliver(
    147             DeliveryRequest::new(
    148                 "deliver-1",
    149                 delivery_payload(),
    150                 target_set(),
    151                 SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all()),
    152                 1_700_000_100_000,
    153             )
    154             .expect("delivery request"),
    155         ),
    156     )
    157     .expect("delivery receipt");
    158     assert_eq!(receipt.request_id().as_str(), "deliver-1");
    159 }
    160 
    161 #[test]
    162 fn a_bidirectional_adapter_exposes_both_dyn_contracts() {
    163     let adapter = Bidirectional {
    164         source: SourceOnly,
    165         sink: SinkOnly,
    166     };
    167     assert_source_dyn_compatible(&adapter);
    168     assert_sink_dyn_compatible(&adapter);
    169 
    170     assert!(
    171         block_on(EventSource::status(&adapter))
    172             .expect("source status")
    173             .capabilities()
    174             .can_fetch()
    175     );
    176     assert!(
    177         block_on(EventSink::status(&adapter))
    178             .expect("sink status")
    179             .capabilities()
    180             .can_deliver()
    181     );
    182 }