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 }