lib

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

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 }