lib

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

fetch_windows.rs (5769B)


      1 use futures::{SinkExt, StreamExt};
      2 use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Timestamp};
      3 use radroots_transport::{
      4     EventSource, FetchRequest, TargetSet,
      5     outcome::FetchTargetState,
      6     source::{FetchBounds, NextPage},
      7 };
      8 use radroots_transport_nostr::{
      9     Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind,
     10     RelayUrlPolicy,
     11 };
     12 use serde_json::Value;
     13 use std::{
     14     collections::BTreeSet,
     15     sync::{
     16         Arc,
     17         atomic::{AtomicUsize, Ordering},
     18     },
     19     time::{Duration, SystemTime, UNIX_EPOCH},
     20 };
     21 use tokio::net::TcpListener;
     22 use tokio_tungstenite::{accept_async, tungstenite::Message};
     23 
     24 #[tokio::test(flavor = "multi_thread")]
     25 async fn real_capped_eose_preserves_received_ties_and_yields_explicit_older_history() {
     26     let keys =
     27         Keys::parse("0000000000000000000000000000000000000000000000000000000000000001").unwrap();
     28     let mut history = (0..1001)
     29         .map(|index| {
     30             let event = EventBuilder::text_note(format!("equal-time {index}"))
     31                 .custom_created_at(Timestamp::from_secs(100))
     32                 .sign_with_keys(&keys)
     33                 .unwrap();
     34             serde_json::from_str::<Value>(&event.as_json()).unwrap()
     35         })
     36         .collect::<Vec<_>>();
     37     history.sort_by(|left, right| right["id"].as_str().cmp(&left["id"].as_str()));
     38     let expected = history[..1000]
     39         .iter()
     40         .map(|event| event["id"].as_str().unwrap().to_owned())
     41         .collect::<BTreeSet<_>>();
     42     let older = EventBuilder::text_note("older")
     43         .custom_created_at(Timestamp::from_secs(99))
     44         .sign_with_keys(&keys)
     45         .unwrap();
     46     history.push(serde_json::from_str(&older.as_json()).unwrap());
     47     let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
     48     let url = format!("ws://{}", listener.local_addr().unwrap());
     49     let calls = Arc::new(AtomicUsize::new(0));
     50     let server_calls = calls.clone();
     51     let server = tokio::spawn(async move {
     52         let (stream, _) = listener.accept().await.unwrap();
     53         let mut socket = accept_async(stream).await.unwrap();
     54         while let Some(Ok(message)) = socket.next().await {
     55             let Message::Text(message) = message else {
     56                 continue;
     57             };
     58             let request: Value = serde_json::from_str(&message).unwrap();
     59             if request[0] != "REQ" {
     60                 continue;
     61             }
     62             server_calls.fetch_add(1, Ordering::SeqCst);
     63             assert_eq!(request[2]["limit"], 1000);
     64             let until = request[2]["until"].as_u64().unwrap_or(u64::MAX);
     65             for event in history
     66                 .iter()
     67                 .filter(|event| event["created_at"].as_u64().unwrap() <= until)
     68                 .take(1000)
     69             {
     70                 socket
     71                     .send(Message::Text(
     72                         serde_json::to_string(&("EVENT", &request[1], event))
     73                             .unwrap()
     74                             .into(),
     75                     ))
     76                     .await
     77                     .unwrap();
     78             }
     79             socket
     80                 .send(Message::Text(
     81                     serde_json::to_string(&("EOSE", &request[1]))
     82                         .unwrap()
     83                         .into(),
     84                 ))
     85                 .await
     86                 .unwrap();
     87         }
     88     });
     89     let profile = RelayProfile::explicit(
     90         RelayProfileKind::Simulator,
     91         [RelayEndpoint::new(&url, RelayUrlPolicy::Local, RelayAccess::ReadOnly).unwrap()],
     92     )
     93     .unwrap();
     94     let config = Config::from_profile(profile)
     95         .with_timeouts(5000, 5000, 500)
     96         .unwrap();
     97     let targets = TargetSet::new(
     98         config
     99             .read_relays()
    100             .map(|relay| relay.to_target().unwrap())
    101             .collect(),
    102     )
    103     .unwrap();
    104     let transport = NostrTransport::new(config);
    105     let deadline = SystemTime::now()
    106         .duration_since(UNIX_EPOCH)
    107         .unwrap()
    108         .as_millis() as u64
    109         + 20000;
    110     let request = FetchRequest::new(
    111         "live-capped-window",
    112         targets,
    113         FetchBounds::new(500, deadline).unwrap(),
    114     )
    115     .unwrap();
    116     tokio::time::timeout(Duration::from_secs(15), async {
    117         let first = transport.fetch(request.clone()).await.unwrap();
    118         assert_eq!(first.events().len(), 500);
    119         assert_eq!(
    120             first.target_outcomes()[0].state(),
    121             FetchTargetState::Partial
    122         );
    123         let NextPage::Cursor(cursor) = first.next_page() else {
    124             panic!("received peers remain")
    125         };
    126         let second = transport
    127             .fetch(request.clone().with_cursor(cursor.clone()))
    128             .await
    129             .unwrap();
    130         assert_eq!(second.events().len(), 500);
    131         assert_eq!(
    132             second.target_outcomes()[0].state(),
    133             FetchTargetState::Partial
    134         );
    135         assert_eq!(calls.load(Ordering::SeqCst), 2);
    136         let received = first
    137             .events()
    138             .iter()
    139             .chain(second.events())
    140             .map(|e| e.event().id_str().to_owned())
    141             .collect::<BTreeSet<_>>();
    142         assert_eq!(received, expected);
    143         let NextPage::Cancelled {
    144             resume_from: Some(older),
    145         } = second.next_page()
    146         else {
    147             panic!("explicit partial window yield")
    148         };
    149         let third = transport
    150             .fetch(request.with_cursor(older.clone()))
    151             .await
    152             .unwrap();
    153         assert_eq!(third.events().len(), 1);
    154         assert_eq!(third.events()[0].event().created_at(), 99);
    155         assert!(matches!(third.next_page(), NextPage::Complete));
    156         assert_eq!(calls.load(Ordering::SeqCst), 3);
    157     })
    158     .await
    159     .unwrap();
    160     server.abort();
    161 }