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 }