lib

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

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 }