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 }