source.rs (76041B)
1 //! Nostr implementation of the transport event source. 2 3 use crate::{NostrTransport, RelayCursor, RelayUrl, status}; 4 use core::cmp::Ordering; 5 use core::time::Duration; 6 use futures::{StreamExt, stream}; 7 use nostr_sdk::prelude::{ 8 Filter, JsonUtil, Kind, RelayMessage, RelayNotification, ReqExitPolicy, SingleLetterTag, 9 SubscribeAutoCloseOptions, SubscribeOptions, SubscriptionId, Timestamp, 10 }; 11 use radroots_transport::{ 12 BoxFuture, EventSource, FetchPage, FetchRequest, 13 outcome::{FetchTargetOutcome, FetchTargetState}, 14 source::{EventProvenance, FetchCursor, NextPage, ObservedEvent, SourceStatus}, 15 }; 16 use sha2::{Digest, Sha256}; 17 use std::collections::{BTreeMap, BTreeSet}; 18 use std::sync::Arc; 19 use std::time::{SystemTime, UNIX_EPOCH}; 20 21 #[path = "source_budget.rs"] 22 pub(crate) mod budget; 23 use budget::FetchBudget; 24 25 #[path = "source_window.rs"] 26 mod window; 27 28 #[cfg(test)] 29 #[path = "source_paging_tests.rs"] 30 mod paging_tests; 31 32 const UPSTREAM_FETCH_LIMIT: usize = 1_000; 33 const CURSOR_PREFIX: &str = "nostr-v2"; 34 const CURSOR_SCOPE_DOMAIN: &[u8] = b"radroots.transport-nostr.fetch-cursor.v2\0"; 35 36 #[derive(Clone, Debug)] 37 pub(crate) struct SourceQuery { 38 relays: Vec<RelayUrl>, 39 selector: radroots_transport::source::FetchSelector, 40 until_unix_seconds: Option<u64>, 41 connect_timeout: Duration, 42 deadline: tokio::time::Instant, 43 max_connections: usize, 44 } 45 46 #[derive(Clone, Debug)] 47 pub(crate) struct RelayFetchBatch { 48 relay: RelayUrl, 49 result: RelayFetchResult, 50 } 51 52 #[derive(Clone, Debug, PartialEq, Eq)] 53 pub(crate) enum RelayFetchResult { 54 Complete(Vec<String>), 55 Timeout(Vec<String>), 56 ResourceLimit(Vec<String>), 57 Failed(String), 58 } 59 60 pub(crate) trait RelaySourceClient: Send + Sync { 61 fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>>; 62 } 63 64 #[derive(Clone, Debug)] 65 pub(crate) struct LiveRelaySourceClient { 66 client: nostr_sdk::Client, 67 ingress: crate::source_ingress::IngressRegistry, 68 } 69 70 impl LiveRelaySourceClient { 71 pub(crate) const fn new( 72 client: nostr_sdk::Client, 73 ingress: crate::source_ingress::IngressRegistry, 74 ) -> Self { 75 Self { client, ingress } 76 } 77 78 #[cfg(test)] 79 pub(crate) fn isolated() -> Self { 80 let client = nostr_sdk::Client::default(); 81 client.automatic_authentication(false); 82 Self::new(client, crate::source_ingress::IngressRegistry::default()) 83 } 84 } 85 86 impl RelaySourceClient for LiveRelaySourceClient { 87 // The live SDK loop is exercised by bounded loopback regressions; injected 88 // source tests independently cover selection and normalized results. 89 #[cfg_attr(coverage_nightly, coverage(off))] 90 fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { 91 Box::pin(async move { 92 let SourceQuery { 93 relays, 94 selector, 95 until_unix_seconds, 96 connect_timeout, 97 deadline, 98 max_connections, 99 } = query; 100 let budget = Arc::new(FetchBudget::default()); 101 // An unencodable selector has no network work or ingress to admit. 102 if !selector.kinds().is_empty() 103 && selector 104 .kinds() 105 .iter() 106 .all(|kind| u16::try_from(*kind).is_err()) 107 { 108 return relays 109 .into_iter() 110 .map(|relay| RelayFetchBatch { 111 relay, 112 result: RelayFetchResult::Complete(Vec::new()), 113 }) 114 .collect(); 115 } 116 let connections = match admit_fetch_connections(&self.client, &relays, deadline).await { 117 Ok(connections) => connections, 118 Err(result) => { 119 return relays 120 .into_iter() 121 .map(|relay| RelayFetchBatch { 122 relay, 123 result: result.clone(), 124 }) 125 .collect(); 126 } 127 }; 128 let registration = self.ingress.register( 129 relays.iter().map(|relay| relay.as_str().to_owned()), 130 Arc::clone(&budget), 131 ); 132 let Some(registration) = registration else { 133 // No ingress or query I/O was admitted, so existing shared 134 // connections retain their previous operation ownership. 135 connections.finish(); 136 return relays 137 .into_iter() 138 .map(|relay| RelayFetchBatch { 139 relay, 140 result: RelayFetchResult::ResourceLimit(Vec::new()), 141 }) 142 .collect(); 143 }; 144 let mut batches: Vec<RelayFetchBatch> = stream::iter(relays.into_iter().map(|relay| { 145 let selector = selector.clone(); 146 let budget = Arc::clone(&budget); 147 async move { 148 let url = relay.as_str().to_owned(); 149 let result = async { 150 if budget.exhausted() { 151 return Ok(RelayFetchResult::ResourceLimit(Vec::new())); 152 } 153 if tokio::time::Instant::now() >= deadline { 154 return Ok(RelayFetchResult::Timeout(Vec::new())); 155 } 156 let kinds = selector 157 .kinds() 158 .iter() 159 .filter_map(|kind| u16::try_from(*kind).ok()) 160 .map(Kind::from) 161 .collect::<Vec<_>>(); 162 let authors = selector 163 .authors() 164 .iter() 165 .filter_map(|author| { 166 radroots_nostr::key::public_key_to_nostr(*author).ok() 167 }) 168 .collect::<Vec<_>>(); 169 if authors.len() != selector.authors().len() { 170 return Ok(RelayFetchResult::Complete(Vec::new())); 171 } 172 let mut filter = Filter::new().limit(UPSTREAM_FETCH_LIMIT); 173 if !kinds.is_empty() { 174 filter = filter.kinds(kinds); 175 } 176 if !authors.is_empty() { 177 filter = filter.authors(authors); 178 } 179 filter = apply_exact_tag_filters(filter, &selector).map_err(|()| { 180 String::from("validated indexed tag cannot be encoded") 181 })?; 182 if let Some(since) = selector.since_unix_seconds() { 183 filter = filter.since(Timestamp::from_secs(since)); 184 } 185 if let Some(until) = until_unix_seconds { 186 filter = filter.until(Timestamp::from_secs(until)); 187 } 188 let connected = tokio::time::timeout_at(deadline, async { 189 self.client 190 .try_connect_relay( 191 url.as_str(), 192 connect_timeout.min( 193 deadline 194 .saturating_duration_since(tokio::time::Instant::now()), 195 ), 196 ) 197 .await 198 .map_err(|error| error.to_string()) 199 }) 200 .await; 201 match connected { 202 Ok(result) => result?, 203 Err(_) => return Ok(RelayFetchResult::Timeout(Vec::new())), 204 } 205 let relay = self 206 .client 207 .relay(url.as_str()) 208 .await 209 .map_err(|error| error.to_string())?; 210 let mut notifications = relay.notifications(); 211 let subscription_id = SubscriptionId::generate(); 212 let remaining = 213 deadline.saturating_duration_since(tokio::time::Instant::now()); 214 if remaining.is_zero() { 215 return Ok(RelayFetchResult::Timeout(Vec::new())); 216 } 217 let options = SubscribeOptions::default().close_on(Some( 218 SubscribeAutoCloseOptions::default() 219 .exit_policy(ReqExitPolicy::ExitOnEOSE) 220 .timeout(Some(remaining)), 221 )); 222 match tokio::time::timeout_at( 223 deadline, 224 relay.subscribe_with_id(subscription_id.clone(), filter, options), 225 ) 226 .await 227 { 228 Ok(result) => { 229 result.map_err(|error| error.to_string())?; 230 } 231 Err(_) => return Ok(RelayFetchResult::Timeout(Vec::new())), 232 } 233 let result = collect_until_eose( 234 &mut notifications, 235 &subscription_id, 236 deadline, 237 &budget, 238 ) 239 .await; 240 Ok(result) 241 } 242 .await; 243 RelayFetchBatch { 244 relay, 245 result: match result { 246 Err(_) if budget.exhausted() => { 247 RelayFetchResult::ResourceLimit(Vec::new()) 248 } 249 Ok(RelayFetchResult::Timeout(events)) if budget.exhausted() => { 250 RelayFetchResult::ResourceLimit(events) 251 } 252 result => result.unwrap_or_else(RelayFetchResult::Failed), 253 }, 254 } 255 } 256 })) 257 .buffered(max_connections) 258 .collect() 259 .await; 260 if finish_fetch(&mut batches, registration) { 261 connections.finish(); 262 } 263 batches 264 }) 265 } 266 } 267 268 async fn admit_fetch_connections( 269 client: &nostr_sdk::Client, 270 relays: &[RelayUrl], 271 deadline: tokio::time::Instant, 272 ) -> Result<FetchConnections, RelayFetchResult> { 273 let connections = FetchConnections::default(); 274 let admitted = tokio::time::timeout_at(deadline, async { 275 for relay in relays { 276 if tokio::time::Instant::now() >= deadline { 277 return Err(RelayFetchResult::Timeout(Vec::new())); 278 } 279 // Adding a relay is inert SDK metadata admission, not connection 280 // work. Own even preexisting queued sockets before ingress can be 281 // revoked by this fetch's registration. 282 client 283 .add_relay(relay.as_str()) 284 .await 285 .map_err(|error| RelayFetchResult::Failed(error.to_string()))?; 286 connections 287 .retain( 288 client 289 .relay(relay.as_str()) 290 .await 291 .map_err(|error| RelayFetchResult::Failed(error.to_string()))?, 292 ) 293 .map_err(RelayFetchResult::Failed)?; 294 } 295 if tokio::time::Instant::now() >= deadline { 296 Err(RelayFetchResult::Timeout(Vec::new())) 297 } else { 298 Ok(()) 299 } 300 }) 301 .await; 302 match admitted { 303 Ok(Ok(())) => Ok(connections), 304 Ok(Err(result)) => Err(result), 305 Err(_) => Err(RelayFetchResult::Timeout(Vec::new())), 306 } 307 } 308 309 fn finish_fetch( 310 batches: &mut [RelayFetchBatch], 311 registration: crate::source_ingress::FetchRegistration, 312 ) -> bool { 313 if !batches 314 .iter() 315 .all(|batch| matches!(batch.result, RelayFetchResult::Complete(_))) 316 { 317 return false; 318 } 319 if registration.finish() { 320 return true; 321 } 322 // Earlier EOSE evidence remains useful, but finalization cannot declare 323 // an entirely complete fetch after ingress exhausted its shared budget. 324 if let Some(last) = batches.last_mut() 325 && let RelayFetchResult::Complete(events) = &mut last.result 326 { 327 last.result = RelayFetchResult::ResourceLimit(std::mem::take(events)); 328 } 329 false 330 } 331 332 // The SDK owns its socket tasks. Its synchronous disconnect signal disables 333 // reconnect and tears down those tasks, including when the fetch future drops. 334 #[derive(Default)] 335 struct FetchConnections(std::sync::Mutex<Vec<nostr_sdk::Relay>>); 336 337 impl FetchConnections { 338 fn retain(&self, relay: nostr_sdk::Relay) -> Result<(), String> { 339 self.0 340 .lock() 341 .map_err(|_| String::from("fetch connection state unavailable"))? 342 .push(relay); 343 Ok(()) 344 } 345 346 fn finish(&self) { 347 if let Ok(mut relays) = self.0.lock() { 348 relays.clear(); 349 } 350 } 351 } 352 353 impl Drop for FetchConnections { 354 fn drop(&mut self) { 355 let relays = self 356 .0 357 .get_mut() 358 .unwrap_or_else(std::sync::PoisonError::into_inner); 359 for relay in relays { 360 relay.disconnect(); 361 } 362 } 363 } 364 365 async fn collect_until_eose( 366 notifications: &mut tokio::sync::broadcast::Receiver<RelayNotification>, 367 subscription_id: &SubscriptionId, 368 deadline: tokio::time::Instant, 369 budget: &FetchBudget, 370 ) -> RelayFetchResult { 371 let mut events = Vec::new(); 372 loop { 373 if budget.exhausted() { 374 return RelayFetchResult::ResourceLimit(events); 375 } 376 if tokio::time::Instant::now() >= deadline { 377 return RelayFetchResult::Timeout(events); 378 } 379 let received = tokio::select! { 380 biased; 381 _ = budget.wait_exhausted() => return RelayFetchResult::ResourceLimit(events), 382 received = tokio::time::timeout_at(deadline, notifications.recv()) => received, 383 }; 384 let notification = match received { 385 Ok(Ok(notification)) => notification, 386 Ok(Err(error)) => return RelayFetchResult::Failed(error.to_string()), 387 Err(_) => return RelayFetchResult::Timeout(events), 388 }; 389 if !budget.notification() { 390 return RelayFetchResult::ResourceLimit(events); 391 } 392 match notification { 393 RelayNotification::Message { 394 message: 395 RelayMessage::Event { 396 subscription_id: observed_subscription, 397 event, 398 }, 399 } if observed_subscription.as_ref() == subscription_id => { 400 if events.len() >= UPSTREAM_FETCH_LIMIT { 401 return RelayFetchResult::ResourceLimit(events); 402 } 403 let raw = event.as_ref().as_json(); 404 if !budget.event(raw.len()) { 405 return RelayFetchResult::ResourceLimit(events); 406 } 407 events.push(raw); 408 } 409 RelayNotification::Message { 410 message: RelayMessage::EndOfStoredEvents(observed_subscription), 411 } if observed_subscription.as_ref() == subscription_id => { 412 return if budget.exhausted() { 413 RelayFetchResult::ResourceLimit(events) 414 } else if tokio::time::Instant::now() < deadline { 415 RelayFetchResult::Complete(events) 416 } else { 417 RelayFetchResult::Timeout(events) 418 }; 419 } 420 RelayNotification::Message { 421 message: 422 RelayMessage::Closed { 423 subscription_id: observed_subscription, 424 message, 425 }, 426 } if observed_subscription.as_ref() == subscription_id => { 427 return RelayFetchResult::Failed(message.into_owned()); 428 } 429 RelayNotification::RelayStatus { 430 status: 431 nostr_sdk::prelude::RelayStatus::Disconnected 432 | nostr_sdk::prelude::RelayStatus::Terminated 433 | nostr_sdk::prelude::RelayStatus::Banned, 434 } => { 435 return RelayFetchResult::Failed(String::from("relay disconnected")); 436 } 437 RelayNotification::AuthenticationFailed => { 438 return RelayFetchResult::Failed(String::from("relay authentication failed")); 439 } 440 RelayNotification::Shutdown => { 441 return RelayFetchResult::Failed(String::from("relay shutdown")); 442 } 443 _ => {} 444 } 445 } 446 } 447 448 #[derive(Debug)] 449 struct Candidate { 450 relay: RelayUrl, 451 raw: String, 452 created_at: u64, 453 event_id: String, 454 } 455 456 impl EventSource for NostrTransport { 457 fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> { 458 Box::pin(async move { Ok(status::source_status(&self.status)) }) 459 } 460 461 fn fetch( 462 &self, 463 request: FetchRequest, 464 ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> { 465 Box::pin(async move { 466 let cursor_scope = request_scope(&request); 467 let cursor = request 468 .cursor() 469 .map(|cursor| window::Position::parse(cursor, cursor_scope.as_str())) 470 .transpose()?; 471 let effective_until = 472 window::effective_until(request.selector().until_unix_seconds(), cursor.as_ref()); 473 let now_ms = unix_time_ms(); 474 let remaining_ms = request.bounds().deadline_unix_ms().saturating_sub(now_ms); 475 let timeout_ms = remaining_ms.min(self.config().request_timeout_ms()); 476 let deadline = tokio::time::Instant::now() + Duration::from_millis(timeout_ms); 477 478 let mut targets = BTreeMap::new(); 479 let mut outcomes = Vec::new(); 480 for target in request.target_set().targets() { 481 match self.config().endpoint_for_target(target) { 482 Some(endpoint) if endpoint.access().can_read() => { 483 let relay = endpoint.url().clone(); 484 if self.status.may_read(&relay, now_ms) { 485 targets.insert(relay, target.clone()); 486 } else { 487 outcomes.push( 488 FetchTargetOutcome::new( 489 target.fingerprint().clone(), 490 FetchTargetState::FailedRetryable, 491 ) 492 .with_message("relay reconnect backoff is active"), 493 ); 494 } 495 } 496 None | Some(_) => outcomes.push( 497 FetchTargetOutcome::new( 498 target.fingerprint().clone(), 499 FetchTargetState::FailedTerminal, 500 ) 501 .with_message("target is not configured for this source"), 502 ), 503 } 504 } 505 506 if timeout_ms == 0 { 507 for relay in targets.keys() { 508 self.status.record_read(relay, false, true, now_ms); 509 } 510 outcomes.extend(targets.values().map(|target| { 511 FetchTargetOutcome::new( 512 target.fingerprint().clone(), 513 FetchTargetState::Cancelled, 514 ) 515 .with_message("fetch deadline elapsed before relay access") 516 })); 517 return FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete); 518 } 519 520 for relay in targets.keys() { 521 self.status.begin_read(relay, now_ms); 522 } 523 let batches = self 524 .source_client 525 .fetch(SourceQuery { 526 relays: targets.keys().cloned().collect(), 527 selector: request.selector().clone(), 528 until_unix_seconds: effective_until, 529 connect_timeout: Duration::from_millis( 530 timeout_ms.min(self.config().connect_timeout_ms()), 531 ), 532 deadline, 533 max_connections: self.config().max_connections(), 534 }) 535 .await; 536 let mut candidates = Vec::new(); 537 let mut capped_window = false; 538 let mut older_boundary = None; 539 let parse_budget = FetchBudget::default(); 540 let mut malformed_by_relay = BTreeMap::<RelayUrl, usize>::new(); 541 let mut reported = BTreeSet::new(); 542 let observed_at_unix_ms = unix_time_ms().max(now_ms); 543 for batch in batches { 544 let RelayFetchBatch { relay, result } = batch; 545 let Some(target) = targets.get(&relay) else { 546 return Err(radroots_transport::Error::UnexpectedFetchTargetOutcome); 547 }; 548 if !reported.insert(relay.clone()) { 549 return Err(radroots_transport::Error::DuplicateFetchTargetOutcome); 550 } 551 let (raw_events, mut terminal) = match result { 552 RelayFetchResult::Complete(events) => (events, FetchTargetState::Complete), 553 RelayFetchResult::Timeout(events) => (events, FetchTargetState::Cancelled), 554 RelayFetchResult::ResourceLimit(events) => (events, FetchTargetState::Partial), 555 RelayFetchResult::Failed(message) => { 556 let (state, safe) = status::fetch_failure(message.as_str()); 557 self.status.record_read( 558 &relay, 559 false, 560 state.is_retryable(), 561 observed_at_unix_ms, 562 ); 563 outcomes.push( 564 FetchTargetOutcome::new(target.fingerprint().clone(), state) 565 .with_message(safe), 566 ); 567 continue; 568 } 569 }; 570 let capped_eose = terminal == FetchTargetState::Complete 571 && raw_events.len() >= UPSTREAM_FETCH_LIMIT; 572 capped_window |= capped_eose; 573 let mut oldest_matching = None; 574 for (index, raw) in raw_events.into_iter().enumerate() { 575 if index >= UPSTREAM_FETCH_LIMIT || !parse_budget.event(raw.len()) { 576 if terminal != FetchTargetState::Cancelled { 577 terminal = FetchTargetState::Partial; 578 } 579 break; 580 } 581 match radroots_event_codec::decode::signed_event(raw.as_str()) { 582 Ok(event) 583 if request.selector().matches(&event) 584 && effective_until 585 .is_none_or(|until| event.created_at() <= until) => 586 { 587 oldest_matching = 588 Some(oldest_matching.map_or(event.created_at(), |oldest: u64| { 589 oldest.min(event.created_at()) 590 })); 591 candidates.push(Candidate { 592 relay: relay.clone(), 593 created_at: event.created_at(), 594 event_id: event.id_str().to_owned(), 595 raw, 596 }) 597 } 598 Ok(_) => {} 599 Err(_) => { 600 *malformed_by_relay.entry(relay.clone()).or_default() += 1; 601 } 602 } 603 } 604 if capped_eose && let Some(oldest) = oldest_matching { 605 older_boundary = 606 Some(older_boundary.map_or(oldest, |boundary: u64| boundary.max(oldest))); 607 } 608 // A capped EOSE proves relay availability, while the page's 609 // coverage stays partial. It must not start reconnect backoff. 610 self.status.record_read( 611 &relay, 612 terminal == FetchTargetState::Complete, 613 terminal != FetchTargetState::Complete, 614 observed_at_unix_ms, 615 ); 616 let malformed = malformed_by_relay.get(&relay).copied().unwrap_or_default(); 617 let outcome = if terminal == FetchTargetState::Cancelled { 618 FetchTargetOutcome::new( 619 target.fingerprint().clone(), 620 FetchTargetState::Cancelled, 621 ) 622 .with_message("relay fetch deadline elapsed before EOSE") 623 } else if terminal == FetchTargetState::Partial || capped_eose { 624 FetchTargetOutcome::new(target.fingerprint().clone(), FetchTargetState::Partial) 625 .with_message("relay result reached the bounded fetch inventory") 626 } else if malformed == 0 { 627 FetchTargetOutcome::new( 628 target.fingerprint().clone(), 629 FetchTargetState::Complete, 630 ) 631 } else { 632 FetchTargetOutcome::new(target.fingerprint().clone(), FetchTargetState::Partial) 633 .with_message(format!("ignored {malformed} malformed relay event(s)")) 634 }; 635 outcomes.push(outcome); 636 } 637 for (relay, target) in &targets { 638 if !reported.contains(relay) { 639 self.status 640 .record_read(relay, false, true, observed_at_unix_ms); 641 outcomes.push( 642 FetchTargetOutcome::new( 643 target.fingerprint().clone(), 644 FetchTargetState::FailedRetryable, 645 ) 646 .with_message("relay returned no fetch result"), 647 ); 648 } 649 } 650 candidates.sort_by(compare_candidate); 651 if let Some(cursor) = &cursor { 652 candidates.retain(|candidate| cursor.includes(candidate)); 653 } 654 let mut seen = BTreeSet::new(); 655 candidates.retain(|candidate| seen.insert(candidate.event_id.clone())); 656 657 let has_more = candidates.len() > usize::from(request.bounds().limit()); 658 candidates.truncate(usize::from(request.bounds().limit())); 659 let next_page = if has_more { 660 let last = candidates.last().expect("non-empty bounded page"); 661 NextPage::Cursor(FetchCursor::parse(format!( 662 "{CURSOR_PREFIX}:{}:{}:{cursor_scope}", 663 last.created_at, last.event_id, 664 ))?) 665 } else if capped_window { 666 NextPage::Cancelled { 667 resume_from: older_boundary 668 .and_then(|boundary| window::before_boundary(boundary, &cursor_scope)), 669 } 670 } else { 671 NextPage::Complete 672 }; 673 let observed_at = unix_time_ms().max(1); 674 let mut events = Vec::with_capacity(candidates.len()); 675 for candidate in candidates { 676 let target = targets 677 .get(&candidate.relay) 678 .expect("candidate relay has requested target"); 679 let mut provenance = EventProvenance::new( 680 radroots_transport::TransportId::NOSTR, 681 target.fingerprint().clone(), 682 observed_at, 683 )?; 684 if let Some(request_cursor) = request.cursor().cloned() { 685 provenance = provenance.with_cursor(request_cursor); 686 } 687 let event = radroots_event_codec::decode::signed_event(candidate.raw.as_str()) 688 .map_err(|_| radroots_transport::Error::UnexpectedFetchProvenance)?; 689 events.push(ObservedEvent::new(event, provenance)); 690 } 691 FetchPage::for_request(&request, events, outcomes, next_page) 692 }) 693 } 694 } 695 696 fn parse_cursor( 697 cursor: &FetchCursor, 698 expected_scope: &str, 699 ) -> Result<RelayCursor, radroots_transport::Error> { 700 let mut parts = cursor.as_str().split(':'); 701 let valid = parts.next() == Some(CURSOR_PREFIX); 702 let created_at = parts.next().and_then(|value| value.parse::<u64>().ok()); 703 let event_id = parts.next(); 704 let scope = parts.next(); 705 if !valid || parts.next().is_some() { 706 return Err(radroots_transport::Error::InvalidFetchCursor); 707 } 708 let (Some(created_at), Some(event_id), Some(scope)) = (created_at, event_id, scope) else { 709 return Err(radroots_transport::Error::InvalidFetchCursor); 710 }; 711 if scope != expected_scope { 712 return Err(radroots_transport::Error::InvalidFetchCursor); 713 } 714 RelayCursor::new(created_at, event_id) 715 .map_err(|_| radroots_transport::Error::InvalidFetchCursor) 716 } 717 718 fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering { 719 right 720 .created_at 721 .cmp(&left.created_at) 722 .then_with(|| right.event_id.cmp(&left.event_id)) 723 .then_with(|| left.relay.cmp(&right.relay)) 724 } 725 726 fn candidate_is_after_cursor(candidate: &Candidate, cursor: &RelayCursor) -> bool { 727 cursor.page_precedes(candidate.created_at, candidate.event_id.as_str()) 728 } 729 730 fn request_scope(request: &FetchRequest) -> String { 731 let mut hasher = Sha256::new(); 732 hasher.update(CURSOR_SCOPE_DOMAIN); 733 for target in request.target_set().targets() { 734 hasher.update(target.fingerprint().as_str().as_bytes()); 735 hasher.update([0]); 736 } 737 for kind in request.selector().kinds() { 738 hasher.update(kind.to_be_bytes()); 739 } 740 hasher.update([0]); 741 for author in request.selector().authors() { 742 hasher.update(author.as_bytes()); 743 } 744 hasher.update([0]); 745 hash_exact_tag_filters(&mut hasher, request.selector()); 746 hash_optional_u64(&mut hasher, request.selector().since_unix_seconds()); 747 hash_optional_u64(&mut hasher, request.selector().until_unix_seconds()); 748 hex_encode(&hasher.finalize()) 749 } 750 751 pub(crate) fn hash_exact_tag_filters( 752 hasher: &mut Sha256, 753 selector: &radroots_transport::source::FetchSelector, 754 ) { 755 let mut filters = selector.exact_tag_filters().peekable(); 756 if filters.peek().is_none() { 757 return; 758 } 759 hasher.update([2]); 760 for (key, values) in filters { 761 hasher.update([u8::try_from(key).expect("validated ASCII tag key")]); 762 hasher.update( 763 u64::try_from(values.len()) 764 .expect("bounded tag-value count") 765 .to_be_bytes(), 766 ); 767 for value in values { 768 hasher.update( 769 u64::try_from(value.len()) 770 .expect("bounded tag-value length") 771 .to_be_bytes(), 772 ); 773 hasher.update(value.as_bytes()); 774 } 775 } 776 hasher.update([0]); 777 } 778 779 pub(crate) fn apply_exact_tag_filters( 780 mut filter: Filter, 781 selector: &radroots_transport::source::FetchSelector, 782 ) -> Result<Filter, ()> { 783 for (key, values) in selector.exact_tag_filters() { 784 let tag = SingleLetterTag::from_char(key).map_err(|_| ())?; 785 filter = filter.custom_tags(tag, values.iter().map(String::as_str)); 786 } 787 Ok(filter) 788 } 789 790 fn hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) { 791 match value { 792 Some(value) => { 793 hasher.update([1]); 794 hasher.update(value.to_be_bytes()); 795 } 796 None => hasher.update([0]), 797 } 798 } 799 800 fn hex_encode(bytes: &[u8]) -> String { 801 const HEX: &[u8; 16] = b"0123456789abcdef"; 802 let mut encoded = String::with_capacity(bytes.len() * 2); 803 for byte in bytes { 804 encoded.push(HEX[(byte >> 4) as usize] as char); 805 encoded.push(HEX[(byte & 0x0f) as usize] as char); 806 } 807 encoded 808 } 809 810 #[cfg_attr(coverage_nightly, coverage(off))] 811 fn unix_time_ms() -> u64 { 812 SystemTime::now() 813 .duration_since(UNIX_EPOCH) 814 .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)) 815 .unwrap_or_default() 816 } 817 818 #[cfg(test)] 819 mod tests { 820 use super::*; 821 use crate::{Config, RelayUrlPolicy}; 822 use radroots_transport::{ 823 FetchRequest, Target, TargetSet, 824 source::{FetchBounds, FetchSelector, NextPage}, 825 }; 826 use std::borrow::Cow; 827 use std::sync::{ 828 Arc, 829 atomic::{AtomicUsize, Ordering as AtomicOrdering}, 830 }; 831 832 const FIRST: &str = r#"{"id":"762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1800000100,"kind":0,"tags":[],"content":"{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}","sig":"4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"}"#; 833 const SECOND: &str = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#; 834 835 #[derive(Debug)] 836 struct MockSourceClient; 837 838 impl RelaySourceClient for MockSourceClient { 839 fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { 840 Box::pin(async move { 841 query 842 .relays 843 .into_iter() 844 .map(|relay| RelayFetchBatch { 845 relay, 846 result: RelayFetchResult::Complete(vec![ 847 FIRST.to_owned(), 848 SECOND.to_owned(), 849 "{".to_owned(), 850 ]), 851 }) 852 .collect() 853 }) 854 } 855 } 856 857 fn transport() -> NostrTransport { 858 let config = Config::from_profile( 859 crate::profile::test_profile( 860 crate::RelayProfileKind::Public, 861 RelayUrlPolicy::Public, 862 ["wss://one.example", "wss://two.example"], 863 ) 864 .expect("profile"), 865 ); 866 NostrTransport::with_source_client(config, Arc::new(MockSourceClient)) 867 } 868 869 fn request(limit: u16) -> FetchRequest { 870 FetchRequest::new( 871 "nostr-fetch", 872 TargetSet::new(vec![ 873 Target::nostr_relay("wss://one.example").expect("one"), 874 Target::nostr_relay("wss://two.example").expect("two"), 875 ]) 876 .expect("targets"), 877 FetchBounds::new(limit, u64::MAX).expect("bounds"), 878 ) 879 .expect("request") 880 } 881 882 fn single_request(limit: u16) -> FetchRequest { 883 FetchRequest::new( 884 "nostr-fetch-one", 885 TargetSet::new(vec![Target::nostr_relay("wss://one.example").expect("one")]) 886 .expect("targets"), 887 FetchBounds::new(limit, u64::MAX).expect("bounds"), 888 ) 889 .expect("request") 890 } 891 892 #[test] 893 fn source_deduplicates_relays_reports_malformed_and_paginates() { 894 let transport = transport(); 895 let first_request = request(1); 896 let first = 897 futures::executor::block_on(transport.fetch(first_request)).expect("first page"); 898 assert_eq!(first.events().len(), 1); 899 assert!( 900 first 901 .target_outcomes() 902 .iter() 903 .all(|outcome| outcome.state() == FetchTargetState::Partial) 904 ); 905 let NextPage::Cursor(cursor) = first.next_page() else { 906 panic!("cursor expected"); 907 }; 908 909 let second_request = request(2).with_cursor(cursor.clone()); 910 let second = 911 futures::executor::block_on(transport.fetch(second_request)).expect("second page"); 912 assert_eq!(second.events().len(), 1); 913 assert!(matches!(second.next_page(), NextPage::Complete)); 914 assert_ne!( 915 first.events()[0].event().id(), 916 second.events()[0].event().id() 917 ); 918 } 919 920 #[test] 921 fn source_applies_kind_author_and_time_selectors_before_page_bounds() { 922 let event = radroots_event_codec::decode::signed_event(FIRST).expect("fixture event"); 923 let selector = FetchSelector::all() 924 .with_kinds(vec![0]) 925 .expect("kind") 926 .with_authors(vec![*event.pubkey()]) 927 .expect("author") 928 .with_since_unix_seconds(1_750_000_000) 929 .expect("since"); 930 let selected = 931 futures::executor::block_on(transport().fetch(request(10).with_selector(selector))) 932 .expect("selected page"); 933 assert_eq!(selected.events().len(), 1); 934 assert_eq!(selected.events()[0].event().id_str(), event.id_str()); 935 936 let excluded = FetchSelector::all() 937 .with_kinds(vec![1]) 938 .expect("excluded kind"); 939 let selected = 940 futures::executor::block_on(transport().fetch(request(10).with_selector(excluded))) 941 .expect("empty selected page"); 942 assert!(selected.events().is_empty()); 943 } 944 945 #[test] 946 fn malformed_cursor_fails_before_relay_access() { 947 let request = request(1).with_cursor(FetchCursor::parse("other:1:value").expect("opaque")); 948 let error = futures::executor::block_on(transport().fetch(request)).expect_err("cursor"); 949 assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); 950 } 951 952 #[test] 953 fn dropping_an_unpolled_fetch_performs_no_relay_work() { 954 #[derive(Debug)] 955 struct CountingSourceClient(Arc<AtomicUsize>); 956 957 impl RelaySourceClient for CountingSourceClient { 958 fn fetch<'a>(&'a self, _query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { 959 self.0.fetch_add(1, AtomicOrdering::SeqCst); 960 Box::pin(async { Vec::new() }) 961 } 962 } 963 964 let calls = Arc::new(AtomicUsize::new(0)); 965 let config = Config::from_profile( 966 crate::profile::test_profile( 967 crate::RelayProfileKind::Public, 968 RelayUrlPolicy::Public, 969 ["wss://one.example"], 970 ) 971 .expect("profile"), 972 ); 973 let transport = NostrTransport::with_source_client( 974 config, 975 Arc::new(CountingSourceClient(Arc::clone(&calls))), 976 ); 977 let target_set = TargetSet::new(vec![ 978 Target::nostr_relay("wss://one.example").expect("target"), 979 ]) 980 .expect("targets"); 981 let fetch = transport.fetch( 982 FetchRequest::new( 983 "cancel-before-poll", 984 target_set, 985 FetchBounds::new(1, u64::MAX).expect("bounds"), 986 ) 987 .expect("request"), 988 ); 989 drop(fetch); 990 assert_eq!(calls.load(AtomicOrdering::SeqCst), 0); 991 } 992 993 #[test] 994 fn adapter_and_generic_page_limits_are_identical() { 995 assert_eq!(UPSTREAM_FETCH_LIMIT, 1_000); 996 assert_eq!( 997 UPSTREAM_FETCH_LIMIT, 998 usize::from(radroots_transport::source::FETCH_PAGE_MAX_EVENTS) 999 ); 1000 assert!(FetchBounds::new(1_001, u64::MAX).is_err()); 1001 } 1002 1003 #[derive(Debug)] 1004 struct ScriptedSourceClient(Vec<RelayFetchBatch>); 1005 1006 impl RelaySourceClient for ScriptedSourceClient { 1007 fn fetch<'a>(&'a self, _query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { 1008 Box::pin(async move { self.0.clone() }) 1009 } 1010 } 1011 1012 fn scripted(batches: Vec<RelayFetchBatch>) -> NostrTransport { 1013 let config = Config::from_profile( 1014 crate::profile::test_profile( 1015 crate::RelayProfileKind::Public, 1016 RelayUrlPolicy::Public, 1017 ["wss://one.example", "wss://two.example"], 1018 ) 1019 .expect("profile"), 1020 ); 1021 NostrTransport::with_source_client(config, Arc::new(ScriptedSourceClient(batches))) 1022 } 1023 1024 #[test] 1025 fn source_handles_failed_missing_duplicate_and_unexpected_batches() { 1026 let one = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("one"); 1027 let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).expect("two"); 1028 let failed = futures::executor::block_on( 1029 scripted(vec![RelayFetchBatch { 1030 relay: one.clone(), 1031 result: RelayFetchResult::Failed("connection timeout".into()), 1032 }]) 1033 .fetch(request(10)), 1034 ) 1035 .expect("failed page"); 1036 assert_eq!(failed.target_outcomes().len(), 2); 1037 assert!(failed.events().is_empty()); 1038 1039 let duplicate = scripted(vec![ 1040 RelayFetchBatch { 1041 relay: one.clone(), 1042 result: RelayFetchResult::Complete(vec![]), 1043 }, 1044 RelayFetchBatch { 1045 relay: one, 1046 result: RelayFetchResult::Complete(vec![]), 1047 }, 1048 ]); 1049 assert_eq!( 1050 futures::executor::block_on(duplicate.fetch(request(10))), 1051 Err(radroots_transport::Error::DuplicateFetchTargetOutcome) 1052 ); 1053 1054 let other = RelayUrl::parse("wss://other.example", RelayUrlPolicy::Public).expect("other"); 1055 let unexpected = scripted(vec![RelayFetchBatch { 1056 relay: other, 1057 result: RelayFetchResult::Complete(vec![]), 1058 }]); 1059 assert_eq!( 1060 futures::executor::block_on(unexpected.fetch(request(10))), 1061 Err(radroots_transport::Error::UnexpectedFetchTargetOutcome) 1062 ); 1063 1064 let complete = futures::executor::block_on( 1065 scripted(vec![RelayFetchBatch { 1066 relay: two, 1067 result: RelayFetchResult::Complete(vec![]), 1068 }]) 1069 .fetch(request(10)), 1070 ) 1071 .expect("partial reporting"); 1072 assert_eq!(complete.target_outcomes().len(), 2); 1073 } 1074 1075 #[test] 1076 fn eose_timeout_and_resource_exhaustion_remain_distinct() { 1077 let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay"); 1078 let fetch = |result| { 1079 futures::executor::block_on( 1080 scripted(vec![RelayFetchBatch { 1081 relay: relay.clone(), 1082 result, 1083 }]) 1084 .fetch(single_request(10)), 1085 ) 1086 .expect("page") 1087 }; 1088 1089 let complete = fetch(RelayFetchResult::Complete(vec![FIRST.to_owned()])); 1090 assert_eq!( 1091 complete.target_outcomes()[0].state(), 1092 FetchTargetState::Complete 1093 ); 1094 assert!(matches!(complete.next_page(), NextPage::Complete)); 1095 1096 let timeout = fetch(RelayFetchResult::Timeout(vec![FIRST.to_owned()])); 1097 assert_eq!(timeout.events().len(), 1); 1098 assert_eq!( 1099 timeout.target_outcomes()[0].state(), 1100 FetchTargetState::Cancelled 1101 ); 1102 1103 let exhausted = fetch(RelayFetchResult::ResourceLimit(vec![FIRST.to_owned()])); 1104 assert_eq!(exhausted.events().len(), 1); 1105 assert_eq!( 1106 exhausted.target_outcomes()[0].state(), 1107 FetchTargetState::Partial 1108 ); 1109 } 1110 1111 #[tokio::test] 1112 async fn live_completion_requires_eose_strictly_before_deadline() { 1113 let subscription_id = SubscriptionId::generate(); 1114 let (sender, mut receiver) = tokio::sync::broadcast::channel(2); 1115 sender 1116 .send(RelayNotification::Message { 1117 message: RelayMessage::EndOfStoredEvents(Cow::Owned(subscription_id.clone())), 1118 }) 1119 .expect("EOSE"); 1120 assert_eq!( 1121 collect_until_eose( 1122 &mut receiver, 1123 &subscription_id, 1124 tokio::time::Instant::now() + Duration::from_secs(1), 1125 &FetchBudget::default(), 1126 ) 1127 .await, 1128 RelayFetchResult::Complete(Vec::new()) 1129 ); 1130 1131 let (_sender, mut receiver) = tokio::sync::broadcast::channel(2); 1132 assert_eq!( 1133 collect_until_eose( 1134 &mut receiver, 1135 &subscription_id, 1136 tokio::time::Instant::now(), 1137 &FetchBudget::default() 1138 ) 1139 .await, 1140 RelayFetchResult::Timeout(Vec::new()) 1141 ); 1142 } 1143 1144 #[tokio::test] 1145 async fn exhausted_ingress_wins_over_queued_eose() { 1146 let subscription_id = SubscriptionId::generate(); 1147 let (sender, mut receiver) = tokio::sync::broadcast::channel(1); 1148 sender 1149 .send(RelayNotification::Message { 1150 message: RelayMessage::EndOfStoredEvents(Cow::Owned(subscription_id.clone())), 1151 }) 1152 .unwrap(); 1153 let budget = FetchBudget::default(); 1154 assert!(!budget.wire(budget::MAX_FETCH_BYTES + 1, 0, 0)); 1155 assert_eq!( 1156 collect_until_eose( 1157 &mut receiver, 1158 &subscription_id, 1159 tokio::time::Instant::now() + Duration::from_secs(1), 1160 &budget 1161 ) 1162 .await, 1163 RelayFetchResult::ResourceLimit(Vec::new()) 1164 ); 1165 } 1166 1167 #[tokio::test] 1168 async fn exhaustion_after_eose_before_finalization_retains_events_and_earlier_evidence() { 1169 let relays = ["wss://one.example", "wss://two.example"]; 1170 let registry = 1171 crate::source_ingress::IngressRegistry::new(relays.into_iter().map(str::to_owned)); 1172 let budget = Arc::new(FetchBudget::default()); 1173 let registration = registry 1174 .register(relays.into_iter().map(str::to_owned), Arc::clone(&budget)) 1175 .unwrap(); 1176 let connection = registry.connection(relays[1]).unwrap(); 1177 let mut batches = Vec::new(); 1178 for url in relays { 1179 let id = SubscriptionId::generate(); 1180 let (sender, mut receiver) = tokio::sync::broadcast::channel(2); 1181 sender 1182 .send(RelayNotification::Message { 1183 message: RelayMessage::Event { 1184 subscription_id: Cow::Owned(id.clone()), 1185 event: Cow::Owned(nostr_sdk::prelude::Event::from_json(FIRST).unwrap()), 1186 }, 1187 }) 1188 .unwrap(); 1189 sender 1190 .send(RelayNotification::Message { 1191 message: RelayMessage::EndOfStoredEvents(Cow::Owned(id.clone())), 1192 }) 1193 .unwrap(); 1194 batches.push(RelayFetchBatch { 1195 relay: RelayUrl::parse(url, RelayUrlPolicy::Public).unwrap(), 1196 result: collect_until_eose( 1197 &mut receiver, 1198 &id, 1199 tokio::time::Instant::now() + Duration::from_secs(1), 1200 &budget, 1201 ) 1202 .await, 1203 }); 1204 } 1205 let before = batches.clone(); 1206 assert!(before.iter().all( 1207 |batch| matches!(&batch.result, RelayFetchResult::Complete(events) if events.len() == 1) 1208 )); 1209 assert!(!connection.charge(budget::MAX_FETCH_BYTES + 1, 0, 0)); 1210 assert!(!finish_fetch(&mut batches, registration)); 1211 assert_eq!(batches[0].result, before[0].result); 1212 let RelayFetchResult::Complete(expected) = &before[1].result else { 1213 panic!("completed collection"); 1214 }; 1215 assert_eq!( 1216 batches[1].result, 1217 RelayFetchResult::ResourceLimit(expected.clone()) 1218 ); 1219 } 1220 1221 #[test] 1222 fn finalization_preserves_success_partial_and_empty_results() { 1223 for (result, exhausted, expected_complete) in [ 1224 ( 1225 Some(RelayFetchResult::Complete(vec![FIRST.to_owned()])), 1226 false, 1227 true, 1228 ), 1229 ( 1230 Some(RelayFetchResult::Timeout(vec![FIRST.to_owned()])), 1231 false, 1232 false, 1233 ), 1234 (None, false, true), 1235 (None, true, false), 1236 ] { 1237 let registry = crate::source_ingress::IngressRegistry::new( 1238 ["wss://one.example".to_owned()].into_iter(), 1239 ); 1240 let budget = Arc::new(FetchBudget::default()); 1241 let registration = registry 1242 .register( 1243 ["wss://one.example".to_owned()].into_iter(), 1244 Arc::clone(&budget), 1245 ) 1246 .unwrap(); 1247 let connection = registry.connection("wss://one.example").unwrap(); 1248 let mut batches = result 1249 .map(|result| RelayFetchBatch { 1250 relay: RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).unwrap(), 1251 result, 1252 }) 1253 .into_iter() 1254 .collect::<Vec<_>>(); 1255 let before = batches 1256 .iter() 1257 .map(|batch| batch.result.clone()) 1258 .collect::<Vec<_>>(); 1259 if exhausted { 1260 assert!(!connection.charge(budget::MAX_FETCH_BYTES + 1, 0, 0)); 1261 } 1262 assert_eq!(finish_fetch(&mut batches, registration), expected_complete); 1263 assert_eq!( 1264 batches 1265 .iter() 1266 .map(|batch| batch.result.clone()) 1267 .collect::<Vec<_>>(), 1268 before 1269 ); 1270 assert_eq!( 1271 connection.admitted(&std::task::Context::from_waker( 1272 futures::task::noop_waker_ref() 1273 )), 1274 expected_complete 1275 ); 1276 } 1277 } 1278 1279 #[tokio::test] 1280 async fn completed_connection_ownership_is_retained_until_the_whole_query_finishes() { 1281 let client = nostr_sdk::Client::default(); 1282 client.add_relay("wss://one.example").await.unwrap(); 1283 let relay = client.relay("wss://one.example").await.unwrap(); 1284 let connections = FetchConnections::default(); 1285 connections.retain(relay.clone()).unwrap(); 1286 assert_eq!(connections.0.lock().unwrap().len(), 1); 1287 connections.finish(); 1288 assert!(connections.0.lock().unwrap().is_empty()); 1289 connections.retain(relay).unwrap(); 1290 drop(connections); 1291 } 1292 1293 #[tokio::test] 1294 async fn complete_connection_inventory_is_inert_and_preserves_shared_sockets_on_success() { 1295 let client = nostr_sdk::Client::default(); 1296 let relays = ["wss://one.example", "wss://two.example"] 1297 .map(|url| RelayUrl::parse(url, RelayUrlPolicy::Public).unwrap()); 1298 let connections = admit_fetch_connections( 1299 &client, 1300 &relays, 1301 tokio::time::Instant::now() + Duration::from_secs(1), 1302 ) 1303 .await 1304 .unwrap(); 1305 assert_eq!(connections.0.lock().unwrap().len(), 2); 1306 for relay in &relays { 1307 assert_eq!( 1308 client.relay(relay.as_str()).await.unwrap().status(), 1309 nostr_sdk::prelude::RelayStatus::Initialized 1310 ); 1311 } 1312 connections.finish(); 1313 drop(connections); 1314 for relay in &relays { 1315 assert_eq!( 1316 client.relay(relay.as_str()).await.unwrap().status(), 1317 nostr_sdk::prelude::RelayStatus::Initialized 1318 ); 1319 } 1320 } 1321 1322 #[tokio::test] 1323 async fn expired_inventory_performs_no_sdk_admission() { 1324 let client = nostr_sdk::Client::default(); 1325 let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).unwrap(); 1326 for relays in [vec![relay], Vec::new()] { 1327 assert!(matches!( 1328 admit_fetch_connections(&client, &relays, tokio::time::Instant::now()).await, 1329 Err(RelayFetchResult::Timeout(_)) 1330 )); 1331 assert!(client.relays().await.is_empty()); 1332 } 1333 } 1334 1335 #[tokio::test] 1336 async fn failed_inventory_releases_every_handle_already_admitted() { 1337 let client = nostr_sdk::Client::builder() 1338 .opts( 1339 nostr_sdk::ClientOptions::default() 1340 .pool(nostr_sdk::prelude::RelayPoolOptions::default().max_relays(Some(1))), 1341 ) 1342 .build(); 1343 let relays = ["wss://one.example", "wss://two.example"] 1344 .map(|url| RelayUrl::parse(url, RelayUrlPolicy::Public).unwrap()); 1345 assert!(matches!( 1346 admit_fetch_connections( 1347 &client, 1348 &relays, 1349 tokio::time::Instant::now() + Duration::from_secs(1) 1350 ) 1351 .await, 1352 Err(RelayFetchResult::Failed(_)) 1353 )); 1354 assert_eq!(client.relays().await.len(), 1); 1355 assert_eq!( 1356 client.relay(relays[0].as_str()).await.unwrap().status(), 1357 nostr_sdk::prelude::RelayStatus::Terminated 1358 ); 1359 } 1360 1361 #[derive(Clone, Copy)] 1362 enum QueuedCleanup { 1363 Cancel, 1364 Deadline, 1365 Resource, 1366 } 1367 1368 async fn queued_connection_cleanup(mode: QueuedCleanup) { 1369 use futures::SinkExt; 1370 use tokio_tungstenite::{accept_async, tungstenite::Message}; 1371 let first = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 1372 let second = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); 1373 let urls = [ 1374 format!("ws://{}", first.local_addr().unwrap()), 1375 format!("ws://{}", second.local_addr().unwrap()), 1376 ]; 1377 let endpoints = urls 1378 .iter() 1379 .map(|url| { 1380 crate::RelayEndpoint::new(url, RelayUrlPolicy::Local, crate::RelayAccess::ReadOnly) 1381 .unwrap() 1382 }) 1383 .collect::<Vec<_>>(); 1384 let connector = crate::relay::HardenedWebsocketTransport::new(&endpoints); 1385 let ingress = connector.ingress.clone(); 1386 let client = nostr_sdk::Client::builder() 1387 .websocket_transport(connector) 1388 .build(); 1389 client.automatic_authentication(false); 1390 let (started, observed) = tokio::sync::oneshot::channel(); 1391 let first_server = tokio::spawn(async move { 1392 let (tcp, _) = first.accept().await.unwrap(); 1393 let mut socket = accept_async(tcp).await.unwrap(); 1394 let mut started = Some(started); 1395 while let Some(message) = socket.next().await { 1396 match message { 1397 Ok(Message::Text(text)) => { 1398 let value: serde_json::Value = serde_json::from_str(&text).unwrap(); 1399 if value[0] == "REQ" { 1400 if let Some(started) = started.take() { 1401 let _ = started.send(()); 1402 } 1403 if matches!(mode, QueuedCleanup::Resource) { 1404 let malformed = serde_json::to_string(&( 1405 "EVENT", 1406 &value[1], 1407 serde_json::json!({"invalid": "x".repeat(384 * 1024)}), 1408 )) 1409 .unwrap(); 1410 for _ in 0..24 { 1411 if socket 1412 .send(Message::Text(malformed.clone().into())) 1413 .await 1414 .is_err() 1415 { 1416 return; 1417 } 1418 } 1419 } 1420 } 1421 } 1422 Ok(Message::Close(_)) | Err(_) => return, 1423 _ => { 1424 if socket.flush().await.is_err() { 1425 return; 1426 } 1427 } 1428 } 1429 } 1430 }); 1431 let queued_requests = Arc::new(AtomicUsize::new(0)); 1432 let requests = Arc::clone(&queued_requests); 1433 let (closed, closure) = tokio::sync::oneshot::channel(); 1434 let second_server = tokio::spawn(async move { 1435 let (tcp, _) = second.accept().await.unwrap(); 1436 let mut socket = accept_async(tcp).await.unwrap(); 1437 while let Some(message) = socket.next().await { 1438 match message { 1439 Ok(Message::Text(text)) => { 1440 let value: serde_json::Value = serde_json::from_str(&text).unwrap(); 1441 if value[0] == "REQ" { 1442 requests.fetch_add(1, AtomicOrdering::SeqCst); 1443 } 1444 } 1445 Ok(Message::Close(_)) | Err(_) => break, 1446 _ => { 1447 if socket.flush().await.is_err() { 1448 break; 1449 } 1450 } 1451 } 1452 } 1453 let _ = closed.send(()); 1454 }); 1455 client.add_relay(urls[1].as_str()).await.unwrap(); 1456 client 1457 .try_connect_relay(urls[1].as_str(), Duration::from_secs(1)) 1458 .await 1459 .unwrap(); 1460 let queued = client.relay(urls[1].as_str()).await.unwrap(); 1461 let source = LiveRelaySourceClient::new(client, ingress); 1462 let query = SourceQuery { 1463 relays: urls 1464 .iter() 1465 .map(|url| RelayUrl::parse(url, RelayUrlPolicy::Local).unwrap()) 1466 .collect(), 1467 selector: radroots_transport::source::FetchSelector::all(), 1468 until_unix_seconds: None, 1469 connect_timeout: Duration::from_secs(1), 1470 deadline: tokio::time::Instant::now() 1471 + if matches!(mode, QueuedCleanup::Deadline) { 1472 Duration::from_millis(500) 1473 } else { 1474 Duration::from_secs(5) 1475 }, 1476 max_connections: 1, 1477 }; 1478 let mut fetch = source.fetch(query); 1479 let started = tokio::time::timeout(Duration::from_secs(10), async { 1480 tokio::select! { 1481 biased; 1482 result = observed => result.is_ok(), 1483 _ = &mut fetch => false, 1484 } 1485 }) 1486 .await; 1487 let batches = if matches!(mode, QueuedCleanup::Cancel) { 1488 None 1489 } else { 1490 Some(tokio::time::timeout(Duration::from_secs(10), &mut fetch).await) 1491 }; 1492 drop(fetch); 1493 let status = queued.status(); 1494 let closed = tokio::time::timeout(Duration::from_secs(2), closure).await; 1495 for server in [first_server, second_server] { 1496 server.abort(); 1497 let result = server.await; 1498 assert!(result.is_ok() || result.unwrap_err().is_cancelled()); 1499 } 1500 assert!( 1501 started.unwrap(), 1502 "first relay must own the only scheduled batch" 1503 ); 1504 assert_eq!(queued_requests.load(AtomicOrdering::SeqCst), 0); 1505 assert_eq!( 1506 status, 1507 nostr_sdk::prelude::RelayStatus::Terminated, 1508 "queued preexisting socket must lose automatic reconnect authority" 1509 ); 1510 closed.unwrap().unwrap(); 1511 if let Some(batches) = batches { 1512 let batches = batches.unwrap(); 1513 assert_eq!(batches.len(), 2); 1514 assert!(batches.iter().all(|batch| match mode { 1515 QueuedCleanup::Deadline => matches!(batch.result, RelayFetchResult::Timeout(_)), 1516 QueuedCleanup::Resource => 1517 matches!(batch.result, RelayFetchResult::ResourceLimit(_)), 1518 QueuedCleanup::Cancel => false, 1519 })); 1520 } 1521 } 1522 1523 #[tokio::test(flavor = "multi_thread")] 1524 async fn cancellation_terminates_preexisting_queued_relay_reconnect_authority() { 1525 queued_connection_cleanup(QueuedCleanup::Cancel).await; 1526 } 1527 1528 #[tokio::test(flavor = "multi_thread")] 1529 async fn deadline_terminates_preexisting_queued_relay_reconnect_authority() { 1530 queued_connection_cleanup(QueuedCleanup::Deadline).await; 1531 } 1532 1533 #[tokio::test(flavor = "multi_thread")] 1534 async fn exhaustion_terminates_preexisting_queued_relay_reconnect_authority() { 1535 queued_connection_cleanup(QueuedCleanup::Resource).await; 1536 } 1537 1538 #[test] 1539 fn source_rejects_unconfigured_targets_and_expired_deadlines() { 1540 let unconfigured_targets = TargetSet::new(vec![ 1541 Target::nostr_relay("wss://other.example").expect("other"), 1542 ]) 1543 .expect("targets"); 1544 let unconfigured = FetchRequest::new( 1545 "unconfigured", 1546 unconfigured_targets, 1547 FetchBounds::new(1, u64::MAX).expect("bounds"), 1548 ) 1549 .expect("request"); 1550 let page = futures::executor::block_on(transport().fetch(unconfigured)).expect("page"); 1551 assert_eq!( 1552 page.target_outcomes()[0].state(), 1553 FetchTargetState::FailedTerminal 1554 ); 1555 1556 let expired = FetchRequest::new( 1557 "expired", 1558 single_request(1).target_set().clone(), 1559 FetchBounds::new(1, 1).expect("bounds"), 1560 ) 1561 .expect("request"); 1562 let page = futures::executor::block_on(transport().fetch(expired)).expect("page"); 1563 assert!(page.events().is_empty()); 1564 assert_eq!(page.target_outcomes().len(), 1); 1565 assert_eq!( 1566 page.target_outcomes()[0].state(), 1567 FetchTargetState::Cancelled 1568 ); 1569 assert!(futures::executor::block_on(transport().status()).is_ok()); 1570 } 1571 1572 #[test] 1573 fn cursor_and_candidate_ordering_cover_boundaries() { 1574 for invalid in [ 1575 "nostr-v1", 1576 "nostr-v2:not-a-time:id:scope", 1577 "nostr-v2:1", 1578 "nostr-v2:1:id:scope:extra", 1579 "nostr-v2:1:abc:scope", 1580 "nostr-v2:1:gggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggg:scope", 1581 ] { 1582 let cursor = FetchCursor::parse(invalid).expect("opaque cursor"); 1583 assert!(matches!( 1584 parse_cursor(&cursor, &"a".repeat(64)), 1585 Err(radroots_transport::Error::InvalidFetchCursor) 1586 )); 1587 } 1588 let cursor = RelayCursor::new(10, "b".repeat(64)).expect("cursor"); 1589 let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay"); 1590 let older = Candidate { 1591 relay: relay.clone(), 1592 raw: String::new(), 1593 created_at: 9, 1594 event_id: "f".repeat(64), 1595 }; 1596 let earlier_id = Candidate { 1597 relay: relay.clone(), 1598 raw: String::new(), 1599 created_at: 10, 1600 event_id: "a".repeat(64), 1601 }; 1602 let later = Candidate { 1603 relay, 1604 raw: String::new(), 1605 created_at: 11, 1606 event_id: "0".repeat(64), 1607 }; 1608 assert!(candidate_is_after_cursor(&older, &cursor)); 1609 assert!(candidate_is_after_cursor(&earlier_id, &cursor)); 1610 assert!(!candidate_is_after_cursor(&later, &cursor)); 1611 assert_ne!(compare_candidate(&older, &later), Ordering::Equal); 1612 } 1613 1614 #[test] 1615 fn continuation_cursor_is_bound_to_exact_targets_and_selector() { 1616 let first = futures::executor::block_on(transport().fetch(request(1))).expect("first page"); 1617 let NextPage::Cursor(cursor) = first.next_page() else { 1618 panic!("cursor expected"); 1619 }; 1620 let different_selector = FetchSelector::all().with_kinds(vec![1]).expect("selector"); 1621 let error = futures::executor::block_on( 1622 transport().fetch( 1623 request(1) 1624 .with_selector(different_selector) 1625 .with_cursor(cursor.clone()), 1626 ), 1627 ) 1628 .expect_err("scope mismatch"); 1629 assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); 1630 1631 let tagged_selector = FetchSelector::all() 1632 .with_exact_tag_value('d', "trade-1") 1633 .expect("tag selector"); 1634 let error = futures::executor::block_on( 1635 transport().fetch( 1636 request(1) 1637 .with_selector(tagged_selector) 1638 .with_cursor(cursor.clone()), 1639 ), 1640 ) 1641 .expect_err("tag scope mismatch"); 1642 assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); 1643 1644 let other_targets = 1645 TargetSet::new(vec![Target::nostr_relay("wss://one.example").expect("one")]) 1646 .expect("targets"); 1647 let error = futures::executor::block_on( 1648 transport().fetch( 1649 FetchRequest::new( 1650 "different-targets", 1651 other_targets, 1652 FetchBounds::new(1, u64::MAX).expect("bounds"), 1653 ) 1654 .expect("request") 1655 .with_cursor(cursor.clone()), 1656 ), 1657 ) 1658 .expect_err("target scope mismatch"); 1659 assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); 1660 } 1661 1662 #[test] 1663 fn exact_tag_selector_maps_to_nostr_filter_and_is_defensively_enforced() { 1664 let selector = FetchSelector::all() 1665 .with_exact_tag_value('d', "trade-2") 1666 .and_then(|selector| selector.with_exact_tag_value('d', "trade-1")) 1667 .expect("tag selector"); 1668 let encoded = apply_exact_tag_filters(Filter::new(), &selector) 1669 .expect("Nostr filter") 1670 .as_json(); 1671 assert!(encoded.contains("\"#d\":[\"trade-1\",\"trade-2\"]")); 1672 1673 let selected = 1674 futures::executor::block_on(transport().fetch(request(10).with_selector(selector))) 1675 .expect("selected page"); 1676 assert!(selected.events().is_empty()); 1677 } 1678 1679 #[test] 1680 fn live_source_short_circuits_selectors_that_cannot_be_encoded() { 1681 let client = LiveRelaySourceClient::isolated(); 1682 let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay"); 1683 let selector = FetchSelector::all() 1684 .with_kinds(vec![u32::MAX]) 1685 .expect("kind"); 1686 let batches = futures::executor::block_on(client.fetch(SourceQuery { 1687 relays: vec![relay], 1688 selector, 1689 until_unix_seconds: None, 1690 connect_timeout: Duration::from_millis(1), 1691 deadline: tokio::time::Instant::now() + Duration::from_millis(1), 1692 max_connections: 1, 1693 })); 1694 assert_eq!(batches.len(), 1); 1695 assert_eq!(batches[0].result, RelayFetchResult::Complete(vec![])); 1696 } 1697 1698 #[tokio::test] 1699 async fn continuous_notifications_terminate_at_the_shared_work_limit() { 1700 let id = SubscriptionId::generate(); 1701 let (sender, mut receiver) = 1702 tokio::sync::broadcast::channel(budget::MAX_FETCH_NOTIFICATIONS + 1); 1703 for _ in 0..=budget::MAX_FETCH_NOTIFICATIONS { 1704 sender 1705 .send(RelayNotification::RelayStatus { 1706 status: nostr_sdk::prelude::RelayStatus::Connected, 1707 }) 1708 .unwrap(); 1709 } 1710 assert_eq!( 1711 collect_until_eose( 1712 &mut receiver, 1713 &id, 1714 tokio::time::Instant::now() + Duration::from_secs(10), 1715 &FetchBudget::default() 1716 ) 1717 .await, 1718 RelayFetchResult::ResourceLimit(Vec::new()) 1719 ); 1720 } 1721 1722 #[tokio::test] 1723 async fn repeated_events_consume_inventory_before_deduplication() { 1724 let id = SubscriptionId::generate(); 1725 let event = nostr_sdk::prelude::Event::from_json(FIRST).unwrap(); 1726 let (sender, mut receiver) = tokio::sync::broadcast::channel(UPSTREAM_FETCH_LIMIT + 1); 1727 for _ in 0..=UPSTREAM_FETCH_LIMIT { 1728 sender 1729 .send(RelayNotification::Message { 1730 message: RelayMessage::Event { 1731 subscription_id: Cow::Owned(id.clone()), 1732 event: Cow::Owned(event.clone()), 1733 }, 1734 }) 1735 .unwrap(); 1736 } 1737 let result = collect_until_eose( 1738 &mut receiver, 1739 &id, 1740 tokio::time::Instant::now() + Duration::from_secs(10), 1741 &FetchBudget::default(), 1742 ) 1743 .await; 1744 let RelayFetchResult::ResourceLimit(events) = result else { 1745 panic!("bounded inventory"); 1746 }; 1747 assert_eq!(events.len(), UPSTREAM_FETCH_LIMIT); 1748 } 1749 1750 #[test] 1751 fn defensive_parse_budget_preserves_earlier_events_and_refuses_excess_work() { 1752 let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).unwrap(); 1753 let mut records = vec![FIRST.to_owned()]; 1754 records.extend((1..UPSTREAM_FETCH_LIMIT).map(|_| "{".to_owned())); 1755 records.push(SECOND.to_owned()); 1756 for records in [ 1757 records, 1758 vec![ 1759 FIRST.to_owned(), 1760 "x".repeat(budget::MAX_EVENT_BYTES + 1), 1761 SECOND.to_owned(), 1762 ], 1763 ] { 1764 let page = futures::executor::block_on( 1765 scripted(vec![RelayFetchBatch { 1766 relay: relay.clone(), 1767 result: RelayFetchResult::Complete(records), 1768 }]) 1769 .fetch(single_request(10)), 1770 ) 1771 .unwrap(); 1772 assert_eq!(page.events().len(), 1); 1773 assert_eq!( 1774 page.events()[0].event().id_str(), 1775 radroots_event_codec::decode::signed_event(FIRST) 1776 .unwrap() 1777 .id_str() 1778 ); 1779 assert_eq!(page.target_outcomes()[0].state(), FetchTargetState::Partial); 1780 } 1781 } 1782 1783 #[tokio::test] 1784 async fn completed_relay_batches_do_not_refund_the_shared_byte_budget() { 1785 let mut wire: serde_json::Value = serde_json::from_str(FIRST).unwrap(); 1786 wire["content"] = serde_json::json!(""); 1787 let empty = nostr_sdk::prelude::Event::from_json(wire.to_string()).unwrap(); 1788 wire["content"] = 1789 serde_json::json!("x".repeat(budget::MAX_EVENT_BYTES - empty.as_json().len())); 1790 // Collection bounds precede canonical event admission; the fixture only 1791 // needs the upstream event structure and an exact serialized size. 1792 let event = nostr_sdk::prelude::Event::from_json(wire.to_string()).unwrap(); 1793 assert_eq!(event.as_json().len(), budget::MAX_EVENT_BYTES); 1794 let shared = FetchBudget::default(); 1795 let mut retained_bytes = 0; 1796 for batch in 0..2 { 1797 let id = SubscriptionId::generate(); 1798 let (sender, mut receiver) = tokio::sync::broadcast::channel(18); 1799 for _ in 0..17 { 1800 sender 1801 .send(RelayNotification::Message { 1802 message: RelayMessage::Event { 1803 subscription_id: Cow::Owned(id.clone()), 1804 event: Cow::Owned(event.clone()), 1805 }, 1806 }) 1807 .unwrap(); 1808 } 1809 sender 1810 .send(RelayNotification::Message { 1811 message: RelayMessage::EndOfStoredEvents(Cow::Owned(id.clone())), 1812 }) 1813 .unwrap(); 1814 let result = collect_until_eose( 1815 &mut receiver, 1816 &id, 1817 tokio::time::Instant::now() + Duration::from_secs(10), 1818 &shared, 1819 ) 1820 .await; 1821 let events = match (batch, result) { 1822 (0, RelayFetchResult::Complete(events)) => { 1823 assert_eq!(events.len(), 17); 1824 events 1825 } 1826 (1, RelayFetchResult::ResourceLimit(events)) => { 1827 assert_eq!(events.len(), 15); 1828 events 1829 } 1830 _ => panic!("the second relay must exhaust the shared byte budget"), 1831 }; 1832 retained_bytes += events.iter().map(String::len).sum::<usize>(); 1833 } 1834 assert_eq!(retained_bytes, budget::MAX_FETCH_BYTES); 1835 assert!(!shared.event(1)); 1836 } 1837 1838 #[test] 1839 fn duplicate_inventory_is_bounded_across_all_relay_batches_before_deduplication() { 1840 let urls = (0..5) 1841 .map(|index| format!("wss://relay{index}.example")) 1842 .collect::<Vec<_>>(); 1843 let config = Config::from_profile( 1844 crate::profile::test_profile( 1845 crate::RelayProfileKind::Public, 1846 RelayUrlPolicy::Public, 1847 urls.iter().map(String::as_str), 1848 ) 1849 .unwrap(), 1850 ); 1851 let targets = TargetSet::new( 1852 config 1853 .read_relays() 1854 .map(|relay| relay.to_target().unwrap()) 1855 .collect(), 1856 ) 1857 .unwrap(); 1858 let request = FetchRequest::new( 1859 "aggregate-duplicates", 1860 targets, 1861 FetchBounds::new(10, unix_time_ms() + 10_000).unwrap(), 1862 ) 1863 .unwrap(); 1864 let batches = urls 1865 .iter() 1866 .map(|url| RelayFetchBatch { 1867 relay: RelayUrl::parse(url, RelayUrlPolicy::Public).unwrap(), 1868 result: RelayFetchResult::Complete(vec![FIRST.to_owned(); UPSTREAM_FETCH_LIMIT]), 1869 }) 1870 .collect(); 1871 let transport = 1872 NostrTransport::with_source_client(config, Arc::new(ScriptedSourceClient(batches))); 1873 let page = futures::executor::block_on(transport.fetch(request)).unwrap(); 1874 assert_eq!(page.events().len(), 1); 1875 assert_eq!(page.target_outcomes().len(), 5); 1876 assert_eq!( 1877 page.target_outcomes() 1878 .iter() 1879 .filter(|outcome| outcome.state() == FetchTargetState::Complete) 1880 .count(), 1881 0 1882 ); 1883 assert_eq!( 1884 page.target_outcomes() 1885 .iter() 1886 .filter(|outcome| outcome.state() == FetchTargetState::Partial) 1887 .count(), 1888 5 1889 ); 1890 let report = transport.relay_status(); 1891 assert_eq!( 1892 report 1893 .relays() 1894 .iter() 1895 .filter(|relay| relay.read().state() == crate::RelayEvidenceState::Available) 1896 .count(), 1897 4 1898 ); 1899 assert_eq!( 1900 report 1901 .relays() 1902 .iter() 1903 .filter(|relay| relay.read().state() == crate::RelayEvidenceState::Unavailable) 1904 .count(), 1905 1 1906 ); 1907 } 1908 }