suite.rs (9321B)
1 use core::{future::Future, pin::Pin, task::Context}; 2 use futures::{executor::block_on, task::noop_waker_ref}; 3 use radroots_event::{SignedEvent, wire::v1::Nip01EventWire}; 4 use radroots_transport::{ 5 DeliveryRequest, Error, EventSink, EventSource, FetchRequest, SinkStatus, SourceStatus, Target, 6 TargetSet, 7 outcome::FetchTargetState, 8 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 9 sink::{DELIVERY_REQUEST_ID_MAX_BYTES, DeliveryPayload}, 10 source::{FETCH_PAGE_MAX_EVENTS, FETCH_REQUEST_ID_MAX_BYTES, FetchBounds, FetchCursor}, 11 }; 12 13 pub(crate) const NOW_UNIX_MS: u64 = 1_700_000_000_000; 14 15 pub(crate) trait SourceConformanceHarness { 16 fn source(&self) -> &dyn EventSource; 17 fn expected_status(&self) -> SourceStatus; 18 fn target_set(&self) -> TargetSet; 19 fn now_unix_ms(&self) -> u64; 20 fn captured_request(&self) -> Option<FetchRequest>; 21 fn published(&self) -> bool; 22 fn cancelled_after_publish(&self) -> bool; 23 } 24 25 pub(crate) trait SinkConformanceHarness { 26 fn sink(&self) -> &dyn EventSink; 27 fn expected_status(&self) -> SinkStatus; 28 fn target_set(&self) -> TargetSet; 29 fn now_unix_ms(&self) -> u64; 30 fn captured_request(&self) -> Option<DeliveryRequest>; 31 fn published(&self) -> bool; 32 fn cancelled_after_publish(&self) -> bool; 33 } 34 35 fn fetch_request(id: &str, targets: TargetSet, deadline: u64) -> FetchRequest { 36 FetchRequest::new( 37 id, 38 targets, 39 FetchBounds::new(2, deadline).expect("fetch bounds"), 40 ) 41 .expect("fetch request") 42 .with_cursor(FetchCursor::parse("opaque-cursor").expect("cursor")) 43 } 44 45 fn delivery_request(id: &str, targets: TargetSet, deadline: u64) -> DeliveryRequest { 46 DeliveryRequest::new( 47 id, 48 DeliveryPayload::new(signed_event()), 49 targets, 50 SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::any()), 51 deadline, 52 ) 53 .expect("delivery request") 54 } 55 56 pub(crate) fn assert_source_conformance(harness: &impl SourceConformanceHarness) { 57 let source = harness.source(); 58 let status = block_on(source.status()).expect("source status"); 59 assert_eq!(status, harness.expected_status()); 60 assert!(status.is_configured()); 61 assert!(status.capabilities().can_fetch()); 62 assert!(!status.message().is_empty()); 63 64 let request = fetch_request( 65 "source-conformance", 66 harness.target_set(), 67 harness.now_unix_ms() + 100, 68 ); 69 let page = block_on(source.fetch(request.clone())).expect("fetch page"); 70 page.validate_for_request(&request) 71 .expect("request binding"); 72 assert_eq!(page.request_id().as_str(), request.request_id().as_str()); 73 assert!(page.events().len() <= usize::from(request.bounds().limit())); 74 assert_eq!(page.target_outcomes().len(), request.target_set().len()); 75 assert_eq!( 76 page.target_outcomes()[0].target(), 77 request.target_set().targets()[0].fingerprint() 78 ); 79 assert_eq!( 80 page.target_outcomes()[0].state(), 81 FetchTargetState::Complete 82 ); 83 assert_eq!( 84 page.target_outcomes()[1].target(), 85 request.target_set().targets()[1].fingerprint() 86 ); 87 assert!(page.target_outcomes()[1].state().is_retryable()); 88 assert_eq!(harness.captured_request().as_ref(), Some(&request)); 89 90 let expired = fetch_request( 91 "source-expired", 92 harness.target_set(), 93 harness.now_unix_ms(), 94 ); 95 assert_eq!( 96 block_on(source.fetch(expired)).expect_err("expired fetch"), 97 Error::InvalidFetchDeadline 98 ); 99 } 100 101 pub(crate) fn assert_sink_conformance(harness: &impl SinkConformanceHarness) { 102 let sink = harness.sink(); 103 let status = block_on(sink.status()).expect("sink status"); 104 assert_eq!(status, harness.expected_status()); 105 assert!(status.is_configured()); 106 assert!(status.capabilities().can_deliver()); 107 assert!(!status.message().is_empty()); 108 109 let request = delivery_request( 110 "sink-conformance", 111 harness.target_set(), 112 harness.now_unix_ms() + 100, 113 ); 114 let receipt = block_on(sink.deliver(request.clone())).expect("delivery receipt"); 115 receipt 116 .validate_for_request(&request) 117 .expect("request binding"); 118 assert_eq!(receipt.request_id().as_str(), request.request_id().as_str()); 119 assert_eq!(receipt.target_receipts().len(), request.target_set().len()); 120 for (target_receipt, requested_target) in receipt 121 .target_receipts() 122 .iter() 123 .zip(request.target_set().targets()) 124 { 125 assert_eq!(target_receipt.target(), requested_target); 126 } 127 assert!( 128 receipt.target_receipts()[0] 129 .outcome() 130 .satisfies(SatisfactionClass::Delivered) 131 ); 132 assert!(receipt.target_receipts()[1].outcome().is_retryable()); 133 assert!(receipt.is_satisfied(&request).expect("satisfaction")); 134 assert_eq!(harness.captured_request().as_ref(), Some(&request)); 135 136 let expired = delivery_request("sink-expired", harness.target_set(), harness.now_unix_ms()); 137 let failure = block_on(sink.deliver(expired)).expect_err("expired delivery"); 138 assert_eq!(failure.code(), "invalid_transport_contract"); 139 } 140 141 pub(crate) fn assert_request_boundaries() { 142 assert_eq!( 143 FetchBounds::new(0, 1).expect_err("zero fetch bound"), 144 Error::InvalidFetchLimit 145 ); 146 assert_eq!( 147 FetchBounds::new(FETCH_PAGE_MAX_EVENTS + 1, 1).expect_err("oversized fetch bound"), 148 Error::InvalidFetchLimit 149 ); 150 assert_eq!( 151 FetchRequest::new( 152 "x".repeat(FETCH_REQUEST_ID_MAX_BYTES + 1), 153 target_set(), 154 FetchBounds::new(1, 1).expect("fetch bounds"), 155 ) 156 .expect_err("oversized fetch request id"), 157 Error::InvalidFetchRequestId 158 ); 159 assert_eq!( 160 DeliveryRequest::new( 161 "x".repeat(DELIVERY_REQUEST_ID_MAX_BYTES + 1), 162 DeliveryPayload::new(signed_event()), 163 target_set(), 164 SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all()), 165 1, 166 ) 167 .expect_err("oversized delivery request id"), 168 Error::InvalidDeliveryRequestId 169 ); 170 } 171 172 pub(crate) fn assert_source_error(harness: &impl SourceConformanceHarness, expected: Error) { 173 let request = fetch_request( 174 "source-error", 175 harness.target_set(), 176 harness.now_unix_ms() + 100, 177 ); 178 assert_eq!( 179 block_on(harness.source().fetch(request)).expect_err("source error"), 180 expected 181 ); 182 } 183 184 pub(crate) fn assert_sink_error(harness: &impl SinkConformanceHarness, _expected: Error) { 185 let request = delivery_request( 186 "sink-error", 187 harness.target_set(), 188 harness.now_unix_ms() + 100, 189 ); 190 let failure = block_on(harness.sink().deliver(request)).expect_err("sink error"); 191 assert_eq!(failure.code(), "invalid_transport_contract"); 192 assert!(matches!( 193 failure.retryability(), 194 radroots_transport::outcome::Retryability::Terminal 195 )); 196 } 197 198 pub(crate) fn assert_source_cancellation(harness: &impl SourceConformanceHarness) { 199 let source = harness.source(); 200 let unpolled = source.fetch(fetch_request( 201 "source-unpolled", 202 harness.target_set(), 203 harness.now_unix_ms() + 100, 204 )); 205 drop(unpolled); 206 assert!(!harness.published()); 207 assert!(!harness.cancelled_after_publish()); 208 209 let mut published = source.fetch(fetch_request( 210 "source-published", 211 harness.target_set(), 212 harness.now_unix_ms() + 100, 213 )); 214 let mut context = Context::from_waker(noop_waker_ref()); 215 assert!(Pin::new(&mut published).poll(&mut context).is_pending()); 216 assert!(harness.published()); 217 drop(published); 218 assert!(harness.cancelled_after_publish()); 219 } 220 221 pub(crate) fn assert_sink_cancellation(harness: &impl SinkConformanceHarness) { 222 let sink = harness.sink(); 223 let unpolled = sink.deliver(delivery_request( 224 "sink-unpolled", 225 harness.target_set(), 226 harness.now_unix_ms() + 100, 227 )); 228 drop(unpolled); 229 assert!(!harness.published()); 230 assert!(!harness.cancelled_after_publish()); 231 232 let mut published = sink.deliver(delivery_request( 233 "sink-published", 234 harness.target_set(), 235 harness.now_unix_ms() + 100, 236 )); 237 let mut context = Context::from_waker(noop_waker_ref()); 238 assert!(Pin::new(&mut published).poll(&mut context).is_pending()); 239 assert!(harness.published()); 240 drop(published); 241 assert!(harness.cancelled_after_publish()); 242 } 243 244 fn target_set() -> TargetSet { 245 TargetSet::new(vec![ 246 Target::local("local:conformance-a").expect("first target"), 247 Target::local("local:conformance-b").expect("second target"), 248 ]) 249 .expect("target set") 250 } 251 252 fn signed_event() -> SignedEvent { 253 let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#; 254 let wire = Nip01EventWire::parse_json(raw).expect("wire event"); 255 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 256 }