lib

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

source_paging_tests.rs (9173B)


      1 use super::*;
      2 use crate::{Config, RelayUrlPolicy};
      3 use radroots_transport::{Target, TargetSet, source::FetchBounds};
      4 use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
      5 
      6 struct CappedRelay {
      7     history: Vec<String>,
      8     calls: AtomicUsize,
      9 }
     10 
     11 impl RelaySourceClient for CappedRelay {
     12     fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
     13         Box::pin(async move {
     14             self.calls.fetch_add(1, AtomicOrdering::SeqCst);
     15             let events = self
     16                 .history
     17                 .iter()
     18                 .filter(|raw| {
     19                     let event = radroots_event_codec::decode::signed_event(raw).unwrap();
     20                     query.selector.matches(&event)
     21                         && query
     22                             .until_unix_seconds
     23                             .is_none_or(|until| event.created_at() <= until)
     24                 })
     25                 .take(UPSTREAM_FETCH_LIMIT)
     26                 .cloned()
     27                 .collect::<Vec<_>>();
     28             query
     29                 .relays
     30                 .into_iter()
     31                 .map(|relay| RelayFetchBatch {
     32                     relay,
     33                     result: RelayFetchResult::Complete(events.clone()),
     34                 })
     35                 .collect()
     36         })
     37     }
     38 }
     39 
     40 // Envelope IDs are canonical; signature verification remains in the shared
     41 // ingest owner. This fixture exercises transport ordering only.
     42 fn raw_event(id: usize, time: u64) -> String {
     43     let pubkey = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
     44     let content = format!("page fixture {id}");
     45     let canonical = serde_json::json!([0, pubkey, time, 1, [], content]).to_string();
     46     let event_id = hex_encode(&Sha256::digest(canonical.as_bytes()));
     47     serde_json::json!({
     48         "id": event_id, "pubkey": pubkey,
     49         "created_at": time, "kind": 1, "tags": [], "content": content,
     50         "sig": "2".repeat(128)
     51     })
     52     .to_string()
     53 }
     54 
     55 fn fixture(
     56     ties: usize,
     57     time: u64,
     58     older: bool,
     59     duplicate_relay: bool,
     60 ) -> (NostrTransport, FetchRequest, Arc<CappedRelay>) {
     61     let mut history = (1..=ties)
     62         .rev()
     63         .map(|id| raw_event(id, time))
     64         .collect::<Vec<_>>();
     65     history.sort_by_cached_key(|raw| {
     66         std::cmp::Reverse(
     67             radroots_event_codec::decode::signed_event(raw)
     68                 .unwrap()
     69                 .id_str()
     70                 .to_owned(),
     71         )
     72     });
     73     if older {
     74         history.push(raw_event(ties + 1, time - 1));
     75     }
     76     let source = Arc::new(CappedRelay {
     77         history,
     78         calls: AtomicUsize::new(0),
     79     });
     80     let mut urls = vec!["wss://one.example"];
     81     if duplicate_relay {
     82         urls.push("wss://two.example");
     83     }
     84     let config = Config::from_profile(
     85         crate::profile::test_profile(
     86             crate::RelayProfileKind::Public,
     87             RelayUrlPolicy::Public,
     88             urls.clone(),
     89         )
     90         .unwrap(),
     91     );
     92     let transport = NostrTransport::with_source_client(config, source.clone());
     93     let request = FetchRequest::new(
     94         "capped-page",
     95         TargetSet::new(
     96             urls.into_iter()
     97                 .map(|url| Target::nostr_relay(url).unwrap())
     98                 .collect(),
     99         )
    100         .unwrap(),
    101         FetchBounds::new(500, u64::MAX).unwrap(),
    102     )
    103     .unwrap();
    104     (transport, request, source)
    105 }
    106 
    107 #[tokio::test]
    108 async fn ordinary_equal_time_pages_retain_all_peers_and_deduplicate_relays() {
    109     let (transport, request, source) = fixture(501, 100, true, true);
    110     let first = transport.fetch(request.clone()).await.unwrap();
    111     assert_eq!(first.events().len(), 500);
    112     let NextPage::Cursor(cursor) = first.next_page() else {
    113         panic!("next equal-time page")
    114     };
    115     let second = transport
    116         .fetch(request.with_cursor(cursor.clone()))
    117         .await
    118         .unwrap();
    119     assert_eq!(second.events().len(), 2);
    120     assert!(matches!(second.next_page(), NextPage::Complete));
    121     let ids = first
    122         .events()
    123         .iter()
    124         .chain(second.events())
    125         .map(|e| e.event().id_str())
    126         .collect::<BTreeSet<_>>();
    127     assert_eq!(ids.len(), 502);
    128     assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 2);
    129 }
    130 
    131 #[tokio::test]
    132 async fn capped_eose_yields_partial_coverage_and_explicit_older_backfill() {
    133     for ties in [1000, 1001] {
    134         let (transport, request, source) = fixture(ties, 100, true, true);
    135         let first = transport.fetch(request.clone()).await.unwrap();
    136         assert_eq!(first.events().len(), 500);
    137         assert!(
    138             first
    139                 .target_outcomes()
    140                 .iter()
    141                 .all(|outcome| outcome.state() == FetchTargetState::Partial)
    142         );
    143         let NextPage::Cursor(cursor) = first.next_page() else {
    144             panic!("collected peers remain")
    145         };
    146         let second = transport
    147             .fetch(request.clone().with_cursor(cursor.clone()))
    148             .await
    149             .unwrap();
    150         assert_eq!(second.events().len(), 500);
    151         let NextPage::Cancelled {
    152             resume_from: Some(older),
    153         } = second.next_page()
    154         else {
    155             panic!("a capped boundary must yield explicit older continuation")
    156         };
    157         let received = first
    158             .events()
    159             .iter()
    160             .chain(second.events())
    161             .map(|e| e.event().id_str())
    162             .collect::<BTreeSet<_>>();
    163         assert_eq!(received.len(), 1000);
    164         assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 2);
    165         let third = transport
    166             .fetch(request.with_cursor(older.clone()))
    167             .await
    168             .unwrap();
    169         assert_eq!(third.events().len(), 1);
    170         assert_eq!(third.events()[0].event().created_at(), 99);
    171         assert!(matches!(third.next_page(), NextPage::Complete));
    172         assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 3);
    173         // The 1001st same-time peer is outside the capped discovery window;
    174         // partial evidence, not an invented completeness claim, describes it.
    175     }
    176 }
    177 
    178 #[tokio::test]
    179 async fn zero_timestamp_yields_partial_without_fabricating_older_history() {
    180     let (transport, request, _) = fixture(1000, 0, false, false);
    181     let first = transport.fetch(request.clone()).await.unwrap();
    182     let NextPage::Cursor(cursor) = first.next_page() else {
    183         panic!("received peers")
    184     };
    185     let second = transport
    186         .fetch(request.with_cursor(cursor.clone()))
    187         .await
    188         .unwrap();
    189     assert_eq!(second.events().len(), 500);
    190     assert_eq!(
    191         second.target_outcomes()[0].state(),
    192         FetchTargetState::Partial
    193     );
    194     assert!(matches!(
    195         second.next_page(),
    196         NextPage::Cancelled { resume_from: None }
    197     ));
    198 }
    199 
    200 #[tokio::test]
    201 async fn older_window_rejects_changed_scope_before_access() {
    202     let (transport, request, source) = fixture(1000, 100, false, false);
    203     let older = window::before_boundary(100, &request_scope(&request)).unwrap();
    204     for selector in [
    205         radroots_transport::source::FetchSelector::all()
    206             .with_kinds(vec![1])
    207             .unwrap(),
    208         radroots_transport::source::FetchSelector::all()
    209             .with_since_unix_seconds(1)
    210             .unwrap(),
    211         radroots_transport::source::FetchSelector::all()
    212             .with_until_unix_seconds(100)
    213             .unwrap(),
    214     ] {
    215         assert!(
    216             transport
    217                 .fetch(
    218                     request
    219                         .clone()
    220                         .with_selector(selector)
    221                         .with_cursor(older.clone())
    222                 )
    223                 .await
    224                 .is_err()
    225         );
    226     }
    227     let (_, different, _) = fixture(1, 100, false, true);
    228     assert!(transport.fetch(different.with_cursor(older)).await.is_err());
    229     assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 0);
    230 }
    231 
    232 struct UnfilteredSource(Vec<String>);
    233 impl RelaySourceClient for UnfilteredSource {
    234     fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
    235         Box::pin(async move {
    236             query
    237                 .relays
    238                 .into_iter()
    239                 .map(|relay| RelayFetchBatch {
    240                     relay,
    241                     result: RelayFetchResult::Complete(self.0.clone()),
    242                 })
    243                 .collect()
    244         })
    245     }
    246 }
    247 
    248 #[tokio::test]
    249 async fn malformed_or_out_of_bound_caps_cannot_fabricate_a_backward_boundary() {
    250     for raw in ["{".to_owned(), raw_event(1, 101)] {
    251         let (base, request, _) = fixture(1, 100, false, false);
    252         let transport = NostrTransport::with_source_client(
    253             base.config().clone(),
    254             Arc::new(UnfilteredSource(vec![raw; 1000])),
    255         );
    256         let cursor = window::before_boundary(101, &request_scope(&request)).unwrap();
    257         let page = transport.fetch(request.with_cursor(cursor)).await.unwrap();
    258         assert!(page.events().is_empty());
    259         assert_eq!(page.target_outcomes()[0].state(), FetchTargetState::Partial);
    260         assert!(matches!(
    261             page.next_page(),
    262             NextPage::Cancelled { resume_from: None }
    263         ));
    264     }
    265 }