lib

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

subscription.rs (52157B)


      1 //! Bounded Nostr live-subscription adapter.
      2 
      3 use crate::{NostrTransport, RelayCursor, RelayUrl};
      4 use nostr_sdk::prelude::{
      5     ClientMessage, Filter, JsonUtil, Kind, RelayMessage, RelayPoolNotification, ReqExitPolicy,
      6     SubscribeAutoCloseOptions, SubscribeOptions, SubscriptionId, Timestamp,
      7 };
      8 use radroots_transport::{
      9     BoxFuture, BoxSubscription, EventSubscriber, EventSubscription, SubscriptionEnd,
     10     SubscriptionEndReason, SubscriptionEvent, SubscriptionNext, SubscriptionRequest,
     11     source::{EventProvenance, FetchCursor, ObservedEvent, SubscriptionCheckpoint},
     12 };
     13 use sha2::{Digest, Sha256};
     14 use std::{
     15     collections::{BTreeMap, BTreeSet},
     16     sync::{
     17         Arc,
     18         atomic::{AtomicBool, Ordering},
     19     },
     20     time::{Duration, SystemTime, UNIX_EPOCH},
     21 };
     22 
     23 const CURSOR_PREFIX: &str = "nostr-live-v1";
     24 const CURSOR_SCOPE_DOMAIN: &[u8] = b"radroots.transport-nostr.subscription-cursor.v1\0";
     25 const SUBSCRIPTION_ID_DOMAIN: &[u8] = b"radroots.transport-nostr.subscription-id.v1\0";
     26 
     27 #[derive(Clone, Debug)]
     28 pub(crate) struct RelaySubscriptionQuery {
     29     id: String,
     30     targets: Vec<RelaySubscriptionTarget>,
     31     selector: radroots_transport::source::FetchSelector,
     32     connect_timeout: Duration,
     33     timeout: Duration,
     34 }
     35 
     36 #[derive(Clone, Debug)]
     37 struct RelaySubscriptionTarget {
     38     relay: RelayUrl,
     39     since_unix_seconds: Option<u64>,
     40 }
     41 
     42 #[derive(Clone, Debug, Eq, PartialEq)]
     43 pub(crate) enum RelaySubscriptionItem {
     44     Event { relay: RelayUrl, raw: String },
     45     Closed { relay: RelayUrl },
     46     Shutdown,
     47 }
     48 
     49 pub(crate) trait RelaySubscriptionSession: Send {
     50     fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>>;
     51     fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>>;
     52 }
     53 
     54 pub(crate) trait RelaySubscriptionClient: Send + Sync {
     55     fn subscribe(
     56         &self,
     57         query: RelaySubscriptionQuery,
     58     ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>>;
     59 }
     60 
     61 #[derive(Clone, Debug)]
     62 pub(crate) struct LiveRelaySubscriptionClient {
     63     client: nostr_sdk::Client,
     64 }
     65 
     66 impl LiveRelaySubscriptionClient {
     67     pub(crate) const fn new(client: nostr_sdk::Client) -> Self {
     68         Self { client }
     69     }
     70 
     71     #[cfg(test)]
     72     pub(crate) fn isolated() -> Self {
     73         let client = nostr_sdk::Client::default();
     74         client.automatic_authentication(false);
     75         Self::new(client)
     76     }
     77 }
     78 
     79 impl RelaySubscriptionClient for LiveRelaySubscriptionClient {
     80     #[cfg_attr(coverage_nightly, coverage(off))]
     81     fn subscribe(
     82         &self,
     83         query: RelaySubscriptionQuery,
     84     ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>> {
     85         Box::pin(async move {
     86             let notifications = self.client.notifications();
     87             let subscription_id = SubscriptionId::new(query.id);
     88             let mut targeted = Vec::with_capacity(query.targets.len());
     89             let mut relay_lookup = BTreeMap::new();
     90 
     91             for target in query.targets {
     92                 let url = target.relay.as_str().to_owned();
     93                 let filter = subscription_filter(&query.selector, target.since_unix_seconds)?;
     94                 self.client.add_relay(url.as_str()).await.map_err(|_| ())?;
     95                 self.client
     96                     .try_connect_relay(url.as_str(), query.connect_timeout)
     97                     .await
     98                     .map_err(|_| ())?;
     99                 relay_lookup.insert(url.clone(), target.relay);
    100                 targeted.push((url, filter));
    101             }
    102 
    103             let auto_close = SubscribeAutoCloseOptions::default()
    104                 .exit_policy(ReqExitPolicy::WaitDurationAfterEOSE(query.timeout))
    105                 .timeout(Some(query.timeout));
    106             let output = self
    107                 .client
    108                 .subscribe_targeted(
    109                     subscription_id.clone(),
    110                     targeted,
    111                     SubscribeOptions::default().close_on(Some(auto_close)),
    112                 )
    113                 .await
    114                 .map_err(|_| ())?;
    115             if !output.failed.is_empty() || output.success.len() != relay_lookup.len() {
    116                 let _ = self
    117                     .client
    118                     .send_msg_to(
    119                         relay_lookup.keys().map(String::as_str),
    120                         ClientMessage::close(subscription_id.clone()),
    121                     )
    122                     .await;
    123                 self.client.unsubscribe(&subscription_id).await;
    124                 return Err(());
    125             }
    126 
    127             // Keep the receiver created before REQ publication so an immediate
    128             // relay event cannot race ahead of local observation.
    129             let session = LiveRelaySubscriptionSession {
    130                 client: self.client.clone(),
    131                 subscription_id,
    132                 notifications: Some(notifications),
    133                 relay_lookup,
    134                 cancelled: false,
    135             };
    136             Ok(Box::new(session) as Box<dyn RelaySubscriptionSession>)
    137         })
    138     }
    139 }
    140 
    141 struct LiveRelaySubscriptionSession {
    142     client: nostr_sdk::Client,
    143     subscription_id: SubscriptionId,
    144     notifications: Option<tokio::sync::broadcast::Receiver<RelayPoolNotification>>,
    145     relay_lookup: BTreeMap<String, RelayUrl>,
    146     cancelled: bool,
    147 }
    148 
    149 impl RelaySubscriptionSession for LiveRelaySubscriptionSession {
    150     #[cfg_attr(coverage_nightly, coverage(off))]
    151     fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>> {
    152         Box::pin(async move {
    153             let Some(notifications) = self.notifications.as_mut() else {
    154                 return Ok(RelaySubscriptionItem::Shutdown);
    155             };
    156             loop {
    157                 let notification = match notifications.recv().await {
    158                     Ok(notification) => notification,
    159                     Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => return Err(()),
    160                     Err(tokio::sync::broadcast::error::RecvError::Closed) => {
    161                         return Ok(RelaySubscriptionItem::Shutdown);
    162                     }
    163                 };
    164                 match notification {
    165                     RelayPoolNotification::Message { relay_url, message } => {
    166                         let Some(relay) = self.relay_lookup.get(relay_url.as_str()).cloned() else {
    167                             continue;
    168                         };
    169                         match message {
    170                             RelayMessage::Event {
    171                                 subscription_id,
    172                                 event,
    173                             } if subscription_id.as_ref() == &self.subscription_id => {
    174                                 return Ok(RelaySubscriptionItem::Event {
    175                                     relay,
    176                                     raw: event.as_json(),
    177                                 });
    178                             }
    179                             RelayMessage::Closed {
    180                                 subscription_id, ..
    181                             } if subscription_id.as_ref() == &self.subscription_id => {
    182                                 return Ok(RelaySubscriptionItem::Closed { relay });
    183                             }
    184                             _ => {}
    185                         }
    186                     }
    187                     RelayPoolNotification::Shutdown => {
    188                         return Ok(RelaySubscriptionItem::Shutdown);
    189                     }
    190                     RelayPoolNotification::Event { .. } => {}
    191                 }
    192             }
    193         })
    194     }
    195 
    196     #[cfg_attr(coverage_nightly, coverage(off))]
    197     fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>> {
    198         Box::pin(async move {
    199             if !self.cancelled {
    200                 let output = self
    201                     .client
    202                     .send_msg_to(
    203                         self.relay_lookup.keys().map(String::as_str),
    204                         ClientMessage::close(self.subscription_id.clone()),
    205                     )
    206                     .await
    207                     .map_err(|_| ())?;
    208                 if !output.failed.is_empty() || output.success.len() != self.relay_lookup.len() {
    209                     return Err(());
    210                 }
    211                 self.client.unsubscribe(&self.subscription_id).await;
    212                 self.cancelled = true;
    213                 self.notifications = None;
    214             }
    215             Ok(())
    216         })
    217     }
    218 }
    219 
    220 struct RelayEventSubscription {
    221     request: SubscriptionRequest,
    222     session: Option<Box<dyn RelaySubscriptionSession>>,
    223     targets: BTreeMap<RelayUrl, radroots_transport::Target>,
    224     active_relays: BTreeSet<RelayUrl>,
    225     resume_cursors: BTreeMap<radroots_transport::target::TargetFingerprint, RelayCursor>,
    226     cursors: BTreeMap<radroots_transport::target::TargetFingerprint, RelayCursor>,
    227     checkpoints: BTreeMap<radroots_transport::target::TargetFingerprint, SubscriptionCheckpoint>,
    228     seen_event_ids: BTreeSet<String>,
    229     event_count: u16,
    230     terminal: Option<SubscriptionEnd>,
    231     cancellation_requested: Arc<AtomicBool>,
    232     status: Arc<crate::status::StatusTracker>,
    233 }
    234 
    235 impl RelayEventSubscription {
    236     fn ended(
    237         request: SubscriptionRequest,
    238         reason: SubscriptionEndReason,
    239         status: Arc<crate::status::StatusTracker>,
    240     ) -> Result<Self, ()> {
    241         let terminal = SubscriptionEnd::for_request(&request, 0, [], reason).map_err(|_| ())?;
    242         Ok(Self {
    243             request,
    244             session: None,
    245             targets: BTreeMap::new(),
    246             active_relays: BTreeSet::new(),
    247             resume_cursors: BTreeMap::new(),
    248             cursors: BTreeMap::new(),
    249             checkpoints: BTreeMap::new(),
    250             seen_event_ids: BTreeSet::new(),
    251             event_count: 0,
    252             terminal: Some(terminal),
    253             cancellation_requested: Arc::new(AtomicBool::new(false)),
    254             status,
    255         })
    256     }
    257 
    258     async fn next_inner(&mut self) -> Result<SubscriptionNext, radroots_transport::Error> {
    259         if let Some(terminal) = &self.terminal {
    260             return Ok(SubscriptionNext::End(terminal.clone()));
    261         }
    262         if self.cancellation_requested.swap(false, Ordering::SeqCst) {
    263             return self
    264                 .terminate(SubscriptionEndReason::Cancelled)
    265                 .await
    266                 .map(SubscriptionNext::End);
    267         }
    268         if self.event_count >= self.request.bounds().event_limit() {
    269             return self
    270                 .terminate(SubscriptionEndReason::EventLimit)
    271                 .await
    272                 .map(SubscriptionNext::End);
    273         }
    274 
    275         loop {
    276             let remaining = remaining_duration(self.request.bounds().deadline_unix_ms());
    277             if remaining.is_zero() {
    278                 return self
    279                     .terminate(SubscriptionEndReason::Deadline)
    280                     .await
    281                     .map(SubscriptionNext::End);
    282             }
    283             let item = {
    284                 let session = self
    285                     .session
    286                     .as_mut()
    287                     .ok_or(radroots_transport::Error::SubscriptionUnavailable)?;
    288                 tokio::time::timeout(remaining, session.next()).await
    289             };
    290             let item = match item {
    291                 Ok(Ok(item)) => item,
    292                 Ok(Err(())) => return Err(radroots_transport::Error::SubscriptionUnavailable),
    293                 Err(_) => {
    294                     return self
    295                         .terminate(SubscriptionEndReason::Deadline)
    296                         .await
    297                         .map(SubscriptionNext::End);
    298                 }
    299             };
    300             match item {
    301                 RelaySubscriptionItem::Event { relay, raw } => {
    302                     if let Some(event) = self.admit_event(&relay, raw.as_str())? {
    303                         return Ok(SubscriptionNext::Event(Box::new(event)));
    304                     }
    305                 }
    306                 RelaySubscriptionItem::Closed { relay } => {
    307                     if !self.active_relays.remove(&relay) {
    308                         return Err(radroots_transport::Error::SubscriptionUnavailable);
    309                     }
    310                     if self.active_relays.is_empty() {
    311                         return self
    312                             .terminate(SubscriptionEndReason::SourceClosed)
    313                             .await
    314                             .map(SubscriptionNext::End);
    315                     }
    316                 }
    317                 RelaySubscriptionItem::Shutdown => {
    318                     return self
    319                         .terminate(SubscriptionEndReason::SourceClosed)
    320                         .await
    321                         .map(SubscriptionNext::End);
    322                 }
    323             }
    324         }
    325     }
    326 
    327     fn admit_event(
    328         &mut self,
    329         relay: &RelayUrl,
    330         raw: &str,
    331     ) -> Result<Option<SubscriptionEvent>, radroots_transport::Error> {
    332         let target = self
    333             .targets
    334             .get(relay)
    335             .ok_or(radroots_transport::Error::UnexpectedSubscriptionEvent)?;
    336         let event = radroots_event_codec::decode::signed_event(raw)
    337             .map_err(|_| radroots_transport::Error::UnexpectedSubscriptionEvent)?;
    338         if !self.request.selector().matches(&event) {
    339             return Err(radroots_transport::Error::UnexpectedSubscriptionEvent);
    340         }
    341         let event_id = event.id_str().to_owned();
    342         let created_at = event.created_at();
    343         if self.seen_event_ids.contains(event_id.as_str()) {
    344             return Ok(None);
    345         }
    346         if self
    347             .resume_cursors
    348             .get(target.fingerprint())
    349             .is_some_and(|cursor| {
    350                 created_at < cursor.created_at_unix_s()
    351                     || (created_at == cursor.created_at_unix_s()
    352                         && event_id.as_str() == cursor.event_id())
    353             })
    354         {
    355             return Ok(None);
    356         }
    357 
    358         let cursor = RelayCursor::new(created_at, event_id.clone())
    359             .map_err(|_| radroots_transport::Error::UnexpectedSubscriptionEvent)?;
    360         let cursor_advances = self
    361             .cursors
    362             .get(target.fingerprint())
    363             .is_none_or(|current| cursor > *current);
    364         if cursor_advances {
    365             self.cursors
    366                 .insert(target.fingerprint().clone(), cursor.clone());
    367             let opaque = encode_cursor(&self.request, target.fingerprint(), &cursor)?;
    368             self.checkpoints.insert(
    369                 target.fingerprint().clone(),
    370                 SubscriptionCheckpoint::new(target.fingerprint().clone(), opaque),
    371             );
    372         }
    373         let checkpoint = self
    374             .checkpoints
    375             .get(target.fingerprint())
    376             .cloned()
    377             .ok_or(radroots_transport::Error::UnexpectedSubscriptionEvent)?;
    378         let provenance = EventProvenance::new(
    379             radroots_transport::TransportId::NOSTR,
    380             target.fingerprint().clone(),
    381             unix_time_ms().max(1),
    382         )?
    383         .with_cursor(checkpoint.cursor().clone());
    384         let observed = ObservedEvent::new(event, provenance);
    385         let subscription_event =
    386             SubscriptionEvent::for_request(&self.request, observed, checkpoint.clone())?;
    387 
    388         self.seen_event_ids.insert(event_id);
    389         self.event_count = self.event_count.saturating_add(1);
    390         self.status
    391             .record_read(relay, true, false, unix_time_ms().max(1));
    392         Ok(Some(subscription_event))
    393     }
    394 
    395     async fn terminate(
    396         &mut self,
    397         mut reason: SubscriptionEndReason,
    398     ) -> Result<SubscriptionEnd, radroots_transport::Error> {
    399         if let Some(terminal) = &self.terminal {
    400             return Ok(terminal.clone());
    401         }
    402         if matches!(
    403             reason,
    404             SubscriptionEndReason::EventLimit | SubscriptionEndReason::Cancelled
    405         ) {
    406             let remaining = remaining_duration(self.request.bounds().deadline_unix_ms());
    407             if remaining.is_zero() {
    408                 reason = SubscriptionEndReason::Deadline;
    409             } else if let Some(session) = self.session.as_mut() {
    410                 match tokio::time::timeout(remaining, session.cancel()).await {
    411                     Ok(Ok(())) => {}
    412                     Ok(Err(())) => {
    413                         return Err(radroots_transport::Error::SubscriptionUnavailable);
    414                     }
    415                     Err(_) => reason = SubscriptionEndReason::Deadline,
    416                 }
    417             }
    418         }
    419         self.session = None;
    420         let terminal = SubscriptionEnd::for_request(
    421             &self.request,
    422             self.event_count,
    423             self.checkpoints.values().cloned(),
    424             reason,
    425         )?;
    426         self.terminal = Some(terminal.clone());
    427         Ok(terminal)
    428     }
    429 }
    430 
    431 impl EventSubscription for RelayEventSubscription {
    432     fn request(&self) -> &SubscriptionRequest {
    433         &self.request
    434     }
    435 
    436     fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, radroots_transport::Error>> {
    437         let cancellation = CancellationOnDrop::new(Arc::clone(&self.cancellation_requested));
    438         Box::pin(async move {
    439             let result = self.next_inner().await;
    440             cancellation.complete();
    441             result
    442         })
    443     }
    444 
    445     fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, radroots_transport::Error>> {
    446         let cancellation = CancellationOnDrop::new(Arc::clone(&self.cancellation_requested));
    447         Box::pin(async move {
    448             let result = self.terminate(SubscriptionEndReason::Cancelled).await;
    449             cancellation.complete();
    450             result
    451         })
    452     }
    453 }
    454 
    455 impl EventSubscriber for NostrTransport {
    456     fn subscribe(
    457         &self,
    458         request: SubscriptionRequest,
    459     ) -> BoxFuture<'_, Result<BoxSubscription, radroots_transport::Error>> {
    460         Box::pin(async move {
    461             let remaining = remaining_duration(request.bounds().deadline_unix_ms());
    462             if remaining.is_zero() {
    463                 return RelayEventSubscription::ended(
    464                     request,
    465                     SubscriptionEndReason::Deadline,
    466                     Arc::clone(&self.status),
    467                 )
    468                 .map(|subscription| Box::new(subscription) as BoxSubscription)
    469                 .map_err(|()| radroots_transport::Error::SubscriptionUnavailable);
    470             }
    471             if !selector_is_representable(request.selector()) {
    472                 return RelayEventSubscription::ended(
    473                     request,
    474                     SubscriptionEndReason::SourceClosed,
    475                     Arc::clone(&self.status),
    476                 )
    477                 .map(|subscription| Box::new(subscription) as BoxSubscription)
    478                 .map_err(|()| radroots_transport::Error::SubscriptionUnavailable);
    479             }
    480 
    481             let now_ms = unix_time_ms();
    482             let mut targets = BTreeMap::new();
    483             let mut active_relays = BTreeSet::new();
    484             for target in request.target_set().targets() {
    485                 let endpoint = self
    486                     .config()
    487                     .endpoint_for_target(target)
    488                     .filter(|endpoint| endpoint.access().can_read())
    489                     .ok_or(radroots_transport::Error::SubscriptionUnavailable)?;
    490                 if !self.status.may_read(endpoint.url(), now_ms)
    491                     || targets
    492                         .insert(endpoint.url().clone(), target.clone())
    493                         .is_some()
    494                 {
    495                     return Err(radroots_transport::Error::SubscriptionUnavailable);
    496                 }
    497                 active_relays.insert(endpoint.url().clone());
    498             }
    499 
    500             let mut cursors = BTreeMap::new();
    501             let mut checkpoints = BTreeMap::new();
    502             let mut target_queries = Vec::with_capacity(targets.len());
    503             for (relay, target) in &targets {
    504                 let checkpoint = request
    505                     .checkpoints()
    506                     .iter()
    507                     .find(|checkpoint| checkpoint.target() == target.fingerprint());
    508                 let cursor = checkpoint
    509                     .map(|checkpoint| parse_cursor(&request, checkpoint))
    510                     .transpose()?;
    511                 let since = cursor
    512                     .as_ref()
    513                     .map(RelayCursor::created_at_unix_s)
    514                     .or(request.selector().since_unix_seconds());
    515                 if let Some(cursor) = cursor {
    516                     cursors.insert(target.fingerprint().clone(), cursor);
    517                 }
    518                 if let Some(checkpoint) = checkpoint {
    519                     checkpoints.insert(target.fingerprint().clone(), checkpoint.clone());
    520                 }
    521                 target_queries.push(RelaySubscriptionTarget {
    522                     relay: relay.clone(),
    523                     since_unix_seconds: since,
    524                 });
    525             }
    526 
    527             let timeout = remaining.min(Duration::from_millis(self.config().request_timeout_ms()));
    528             let id = subscription_id(&request, self.next_subscription_sequence()?);
    529             for relay in &active_relays {
    530                 self.status.begin_read(relay, now_ms);
    531             }
    532             let query = RelaySubscriptionQuery {
    533                 id,
    534                 targets: target_queries,
    535                 selector: request.selector().clone(),
    536                 connect_timeout: timeout
    537                     .min(Duration::from_millis(self.config().connect_timeout_ms())),
    538                 timeout: remaining,
    539             };
    540             let session = match tokio::time::timeout(
    541                 timeout,
    542                 self.subscription_client.subscribe(query),
    543             )
    544             .await
    545             {
    546                 Ok(Ok(session)) => session,
    547                 Ok(Err(())) | Err(_) => {
    548                     let observed_at = unix_time_ms().max(1);
    549                     for relay in &active_relays {
    550                         self.status.record_read(relay, false, true, observed_at);
    551                     }
    552                     return Err(radroots_transport::Error::SubscriptionUnavailable);
    553                 }
    554             };
    555 
    556             let resume_cursors = cursors.clone();
    557             Ok(Box::new(RelayEventSubscription {
    558                 request,
    559                 session: Some(session),
    560                 targets,
    561                 active_relays,
    562                 resume_cursors,
    563                 cursors,
    564                 checkpoints,
    565                 seen_event_ids: BTreeSet::new(),
    566                 event_count: 0,
    567                 terminal: None,
    568                 cancellation_requested: Arc::new(AtomicBool::new(false)),
    569                 status: Arc::clone(&self.status),
    570             }) as BoxSubscription)
    571         })
    572     }
    573 }
    574 
    575 struct CancellationOnDrop {
    576     requested: Arc<AtomicBool>,
    577     completed: AtomicBool,
    578 }
    579 
    580 impl CancellationOnDrop {
    581     fn new(requested: Arc<AtomicBool>) -> Arc<Self> {
    582         Arc::new(Self {
    583             requested,
    584             completed: AtomicBool::new(false),
    585         })
    586     }
    587 
    588     fn complete(&self) {
    589         self.completed.store(true, Ordering::SeqCst);
    590     }
    591 }
    592 
    593 impl Drop for CancellationOnDrop {
    594     fn drop(&mut self) {
    595         if !self.completed.load(Ordering::SeqCst) {
    596             self.requested.store(true, Ordering::SeqCst);
    597         }
    598     }
    599 }
    600 
    601 fn subscription_filter(
    602     selector: &radroots_transport::source::FetchSelector,
    603     since: Option<u64>,
    604 ) -> Result<Filter, ()> {
    605     let kinds = selector
    606         .kinds()
    607         .iter()
    608         .filter_map(|kind| u16::try_from(*kind).ok())
    609         .map(Kind::from)
    610         .collect::<Vec<_>>();
    611     if !selector.kinds().is_empty() && kinds.is_empty() {
    612         return Err(());
    613     }
    614     let authors = selector
    615         .authors()
    616         .iter()
    617         .map(|author| radroots_nostr::key::public_key_to_nostr(*author).map_err(|_| ()))
    618         .collect::<Result<Vec<_>, _>>()?;
    619     let mut filter = Filter::new();
    620     if !kinds.is_empty() {
    621         filter = filter.kinds(kinds);
    622     }
    623     if !authors.is_empty() {
    624         filter = filter.authors(authors);
    625     }
    626     filter = crate::source::apply_exact_tag_filters(filter, selector)?;
    627     if let Some(since) = since {
    628         filter = filter.since(Timestamp::from_secs(since));
    629     }
    630     if let Some(until) = selector.until_unix_seconds() {
    631         filter = filter.until(Timestamp::from_secs(until));
    632     }
    633     Ok(filter)
    634 }
    635 
    636 fn selector_is_representable(selector: &radroots_transport::source::FetchSelector) -> bool {
    637     selector.kinds().is_empty()
    638         || selector
    639             .kinds()
    640             .iter()
    641             .any(|kind| u16::try_from(*kind).is_ok())
    642 }
    643 
    644 fn subscription_id(request: &SubscriptionRequest, sequence: u64) -> String {
    645     let mut hasher = Sha256::new();
    646     hasher.update(SUBSCRIPTION_ID_DOMAIN);
    647     hasher.update(request.request_id().as_str().as_bytes());
    648     hasher.update([0]);
    649     hasher.update(sequence.to_be_bytes());
    650     hex_encode(&hasher.finalize())
    651 }
    652 
    653 fn encode_cursor(
    654     request: &SubscriptionRequest,
    655     target: &radroots_transport::target::TargetFingerprint,
    656     cursor: &RelayCursor,
    657 ) -> Result<FetchCursor, radroots_transport::Error> {
    658     FetchCursor::parse(format!(
    659         "{CURSOR_PREFIX}:{}:{}:{}",
    660         cursor.created_at_unix_s(),
    661         cursor.event_id(),
    662         cursor_scope(request, target),
    663     ))
    664 }
    665 
    666 fn parse_cursor(
    667     request: &SubscriptionRequest,
    668     checkpoint: &SubscriptionCheckpoint,
    669 ) -> Result<RelayCursor, radroots_transport::Error> {
    670     let mut parts = checkpoint.cursor().as_str().split(':');
    671     let valid = parts.next() == Some(CURSOR_PREFIX);
    672     let created_at = parts.next().and_then(|value| value.parse::<u64>().ok());
    673     let event_id = parts.next();
    674     let scope = parts.next();
    675     if !valid || parts.next().is_some() {
    676         return Err(radroots_transport::Error::InvalidFetchCursor);
    677     }
    678     let (Some(created_at), Some(event_id), Some(scope)) = (created_at, event_id, scope) else {
    679         return Err(radroots_transport::Error::InvalidFetchCursor);
    680     };
    681     if scope != cursor_scope(request, checkpoint.target()) {
    682         return Err(radroots_transport::Error::InvalidFetchCursor);
    683     }
    684     RelayCursor::new(created_at, event_id)
    685         .map_err(|_| radroots_transport::Error::InvalidFetchCursor)
    686 }
    687 
    688 fn cursor_scope(
    689     request: &SubscriptionRequest,
    690     target: &radroots_transport::target::TargetFingerprint,
    691 ) -> String {
    692     let mut hasher = Sha256::new();
    693     hasher.update(CURSOR_SCOPE_DOMAIN);
    694     hasher.update(target.as_str().as_bytes());
    695     hasher.update([0]);
    696     for kind in request.selector().kinds() {
    697         hasher.update(kind.to_be_bytes());
    698     }
    699     hasher.update([0]);
    700     for author in request.selector().authors() {
    701         hasher.update(author.as_bytes());
    702     }
    703     hasher.update([0]);
    704     crate::source::hash_exact_tag_filters(&mut hasher, request.selector());
    705     hash_optional_u64(&mut hasher, request.selector().since_unix_seconds());
    706     hash_optional_u64(&mut hasher, request.selector().until_unix_seconds());
    707     hex_encode(&hasher.finalize())
    708 }
    709 
    710 fn hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) {
    711     match value {
    712         Some(value) => {
    713             hasher.update([1]);
    714             hasher.update(value.to_be_bytes());
    715         }
    716         None => hasher.update([0]),
    717     }
    718 }
    719 
    720 fn remaining_duration(deadline_unix_ms: u64) -> Duration {
    721     Duration::from_millis(deadline_unix_ms.saturating_sub(unix_time_ms()))
    722 }
    723 
    724 fn hex_encode(bytes: &[u8]) -> String {
    725     const HEX: &[u8; 16] = b"0123456789abcdef";
    726     let mut encoded = String::with_capacity(bytes.len() * 2);
    727     for byte in bytes {
    728         encoded.push(HEX[(byte >> 4) as usize] as char);
    729         encoded.push(HEX[(byte & 0x0f) as usize] as char);
    730     }
    731     encoded
    732 }
    733 
    734 #[cfg_attr(coverage_nightly, coverage(off))]
    735 fn unix_time_ms() -> u64 {
    736     SystemTime::now()
    737         .duration_since(UNIX_EPOCH)
    738         .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
    739         .unwrap_or_default()
    740 }
    741 
    742 #[cfg(test)]
    743 mod tests {
    744     use super::*;
    745     use crate::{
    746         Config, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind, RelayUrlPolicy,
    747     };
    748     use nostr_sdk::prelude::{EventBuilder, Keys};
    749     use radroots_transport::{
    750         EventSubscriber, Target, TargetSet,
    751         source::{FetchSelector, SubscriptionBounds, SubscriptionCheckpoint},
    752     };
    753     use std::{
    754         collections::VecDeque,
    755         sync::{
    756             Arc, Mutex,
    757             atomic::{AtomicBool, AtomicUsize, Ordering as AtomicOrdering},
    758         },
    759     };
    760 
    761     const FIXTURE_SECRET_KEY: &str =
    762         "0000000000000000000000000000000000000000000000000000000000000001";
    763 
    764     #[derive(Clone, Debug)]
    765     enum ScriptedItem {
    766         Item(RelaySubscriptionItem),
    767         Error,
    768     }
    769 
    770     #[derive(Debug)]
    771     struct MockSubscriptionSession {
    772         items: VecDeque<ScriptedItem>,
    773         cancellations: Arc<AtomicUsize>,
    774     }
    775 
    776     impl RelaySubscriptionSession for MockSubscriptionSession {
    777         fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>> {
    778             match self.items.pop_front() {
    779                 Some(ScriptedItem::Item(item)) => Box::pin(async move { Ok(item) }),
    780                 Some(ScriptedItem::Error) => Box::pin(async { Err(()) }),
    781                 None => Box::pin(core::future::pending()),
    782             }
    783         }
    784 
    785         fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>> {
    786             self.cancellations.fetch_add(1, AtomicOrdering::SeqCst);
    787             Box::pin(async { Ok(()) })
    788         }
    789     }
    790 
    791     #[derive(Debug)]
    792     struct MockSubscriptionClient {
    793         queries: Mutex<Vec<RelaySubscriptionQuery>>,
    794         items: Mutex<Option<VecDeque<ScriptedItem>>>,
    795         cancellations: Arc<AtomicUsize>,
    796         fail: AtomicBool,
    797     }
    798 
    799     impl MockSubscriptionClient {
    800         fn new(items: impl IntoIterator<Item = ScriptedItem>) -> Self {
    801             Self {
    802                 queries: Mutex::new(Vec::new()),
    803                 items: Mutex::new(Some(items.into_iter().collect())),
    804                 cancellations: Arc::new(AtomicUsize::new(0)),
    805                 fail: AtomicBool::new(false),
    806             }
    807         }
    808 
    809         fn failing() -> Self {
    810             let client = Self::new([]);
    811             client.fail.store(true, AtomicOrdering::SeqCst);
    812             client
    813         }
    814 
    815         fn query(&self) -> RelaySubscriptionQuery {
    816             self.queries.lock().expect("queries")[0].clone()
    817         }
    818     }
    819 
    820     impl RelaySubscriptionClient for MockSubscriptionClient {
    821         fn subscribe(
    822             &self,
    823             query: RelaySubscriptionQuery,
    824         ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>> {
    825             self.queries.lock().expect("queries").push(query);
    826             if self.fail.load(AtomicOrdering::SeqCst) {
    827                 return Box::pin(async { Err(()) });
    828             }
    829             let items = self.items.lock().expect("items").take().unwrap_or_default();
    830             let session = MockSubscriptionSession {
    831                 items,
    832                 cancellations: Arc::clone(&self.cancellations),
    833             };
    834             Box::pin(async move { Ok(Box::new(session) as Box<dyn RelaySubscriptionSession>) })
    835         }
    836     }
    837 
    838     fn configured(relays: &[&str]) -> Config {
    839         let endpoints = relays
    840             .iter()
    841             .map(|relay| {
    842                 RelayEndpoint::new(relay, RelayUrlPolicy::Public, RelayAccess::ReadWrite)
    843                     .expect("endpoint")
    844             })
    845             .collect::<Vec<_>>();
    846         Config::from_profile(
    847             RelayProfile::explicit(RelayProfileKind::Public, endpoints).expect("profile"),
    848         )
    849     }
    850 
    851     fn target_set(relays: &[&str]) -> TargetSet {
    852         TargetSet::new(
    853             relays
    854                 .iter()
    855                 .map(|relay| Target::nostr_relay(relay).expect("target"))
    856                 .collect::<Vec<_>>(),
    857         )
    858         .expect("target set")
    859     }
    860 
    861     fn request(relays: &[&str], event_limit: u16) -> SubscriptionRequest {
    862         SubscriptionRequest::new(
    863             "nostr-live",
    864             target_set(relays),
    865             SubscriptionBounds::new(event_limit, unix_time_ms() + 60_000).expect("bounds"),
    866         )
    867         .expect("request")
    868     }
    869 
    870     fn signed_event(content: &str, created_at: u64) -> String {
    871         EventBuilder::text_note(content)
    872             .custom_created_at(Timestamp::from_secs(created_at))
    873             .sign_with_keys(&Keys::parse(FIXTURE_SECRET_KEY).expect("fixture keys"))
    874             .expect("signed event")
    875             .as_json()
    876     }
    877 
    878     fn event_id(raw: &str) -> String {
    879         radroots_event_codec::decode::signed_event(raw)
    880             .expect("signed event")
    881             .id_str()
    882             .to_owned()
    883     }
    884 
    885     #[tokio::test]
    886     async fn subscription_translates_targets_selector_and_unique_ids() {
    887         let relay = "wss://one.example";
    888         let client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
    889             RelaySubscriptionItem::Shutdown,
    890         )]));
    891         let transport =
    892             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
    893         let selector = FetchSelector::all()
    894             .with_kinds(vec![1])
    895             .expect("kind")
    896             .with_authors(vec![
    897                 *radroots_event_codec::decode::signed_event(&signed_event(
    898                     "subscription-author",
    899                     1_800_000_000,
    900                 ))
    901                 .expect("signed author event")
    902                 .pubkey(),
    903             ])
    904             .expect("author")
    905             .with_exact_tag_value('d', "trade-1")
    906             .expect("tag")
    907             .with_since_unix_seconds(1_700_000_000)
    908             .expect("since")
    909             .with_until_unix_seconds(1_800_000_000)
    910             .expect("until");
    911         let first = transport
    912             .subscribe(request(&[relay], 2).with_selector(selector.clone()))
    913             .await
    914             .expect("first subscription");
    915         drop(first);
    916         let second = transport
    917             .subscribe(request(&[relay], 2).with_selector(selector.clone()))
    918             .await
    919             .expect("second subscription");
    920         drop(second);
    921 
    922         let queries = client.queries.lock().expect("queries");
    923         assert_eq!(queries.len(), 2);
    924         assert_ne!(queries[0].id, queries[1].id);
    925         assert_eq!(queries[0].targets.len(), 1);
    926         assert_eq!(queries[0].targets[0].relay.as_str(), relay);
    927         assert_eq!(
    928             queries[0].targets[0].since_unix_seconds,
    929             Some(1_700_000_000)
    930         );
    931         assert_eq!(queries[0].selector, selector);
    932         assert!(queries[0].connect_timeout <= queries[0].timeout);
    933 
    934         let filter = subscription_filter(&queries[0].selector, Some(1_700_000_000))
    935             .expect("filter")
    936             .as_json();
    937         assert!(filter.contains("\"kinds\":[1]"));
    938         assert!(filter.contains("\"#d\":[\"trade-1\"]"));
    939         assert!(filter.contains("\"since\":1700000000"));
    940         assert!(filter.contains("\"until\":1800000000"));
    941     }
    942 
    943     #[tokio::test]
    944     async fn reconnect_checkpoint_is_equal_timestamp_safe_and_request_bound() {
    945         let relay = "wss://one.example";
    946         let mut events = [
    947             signed_event("equal-a", 1_800_000_000),
    948             signed_event("equal-b", 1_800_000_000),
    949             signed_event("equal-c", 1_800_000_000),
    950         ];
    951         events.sort_by_key(|event| event_id(event));
    952         let base_request = request(&[relay], 3);
    953         let target = base_request.target_set().targets()[0].fingerprint().clone();
    954         let middle_cursor = RelayCursor::new(1_800_000_000, event_id(&events[1])).expect("cursor");
    955         let checkpoint = SubscriptionCheckpoint::new(
    956             target,
    957             encode_cursor(
    958                 &base_request,
    959                 base_request.target_set().targets()[0].fingerprint(),
    960                 &middle_cursor,
    961             )
    962             .expect("opaque cursor"),
    963         );
    964         let resumed = base_request
    965             .with_checkpoints([checkpoint])
    966             .expect("checkpointed request");
    967         let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay");
    968         let older = signed_event("older", 1_799_999_999);
    969         let later = signed_event("later", 1_800_000_001);
    970         let client = Arc::new(MockSubscriptionClient::new([
    971             ScriptedItem::Item(RelaySubscriptionItem::Event {
    972                 relay: relay_url.clone(),
    973                 raw: older,
    974             }),
    975             ScriptedItem::Item(RelaySubscriptionItem::Event {
    976                 relay: relay_url.clone(),
    977                 raw: events[0].clone(),
    978             }),
    979             ScriptedItem::Item(RelaySubscriptionItem::Event {
    980                 relay: relay_url.clone(),
    981                 raw: events[1].clone(),
    982             }),
    983             ScriptedItem::Item(RelaySubscriptionItem::Event {
    984                 relay: relay_url.clone(),
    985                 raw: events[2].clone(),
    986             }),
    987             ScriptedItem::Item(RelaySubscriptionItem::Event {
    988                 relay: relay_url,
    989                 raw: later.clone(),
    990             }),
    991         ]));
    992         let transport =
    993             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
    994         let mut subscription = transport.subscribe(resumed).await.expect("subscription");
    995         let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") else {
    996             panic!("event expected");
    997         };
    998         assert_eq!(event.observed().event().id_str(), event_id(&events[0]));
    999         assert_eq!(
   1000             event.checkpoint().cursor().as_str().split(':').nth(2),
   1001             Some(event_id(&events[1]).as_str())
   1002         );
   1003         let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") else {
   1004             panic!("event expected");
   1005         };
   1006         assert_eq!(event.observed().event().id_str(), event_id(&events[2]));
   1007         assert_eq!(
   1008             event.checkpoint().cursor().as_str().split(':').nth(2),
   1009             Some(event_id(&events[2]).as_str())
   1010         );
   1011         let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") else {
   1012             panic!("event expected");
   1013         };
   1014         assert_eq!(event.observed().event().id_str(), event_id(&later));
   1015         assert_eq!(
   1016             client.query().targets[0].since_unix_seconds,
   1017             Some(1_800_000_000)
   1018         );
   1019     }
   1020 
   1021     #[tokio::test]
   1022     async fn live_subscription_accepts_same_second_events_in_relay_arrival_order() {
   1023         let relay = "wss://one.example";
   1024         let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay");
   1025         let mut events = [
   1026             signed_event("same-second-a", 1_800_000_000),
   1027             signed_event("same-second-b", 1_800_000_000),
   1028         ];
   1029         events.sort_by_key(|event| event_id(event));
   1030         let client = Arc::new(MockSubscriptionClient::new([
   1031             ScriptedItem::Item(RelaySubscriptionItem::Event {
   1032                 relay: relay_url.clone(),
   1033                 raw: events[1].clone(),
   1034             }),
   1035             ScriptedItem::Item(RelaySubscriptionItem::Event {
   1036                 relay: relay_url.clone(),
   1037                 raw: events[1].clone(),
   1038             }),
   1039             ScriptedItem::Item(RelaySubscriptionItem::Event {
   1040                 relay: relay_url,
   1041                 raw: events[0].clone(),
   1042             }),
   1043         ]));
   1044         let transport = NostrTransport::with_subscription_client(configured(&[relay]), client);
   1045         let mut subscription = transport
   1046             .subscribe(request(&[relay], 2))
   1047             .await
   1048             .expect("subscription");
   1049 
   1050         for expected in [&events[1], &events[0]] {
   1051             let SubscriptionNext::Event(event) = subscription.next().await.expect("next event")
   1052             else {
   1053                 panic!("event expected");
   1054             };
   1055             assert_eq!(event.observed().event().id_str(), event_id(expected));
   1056             assert_eq!(
   1057                 event.checkpoint().cursor().as_str().split(':').nth(2),
   1058                 Some(event_id(&events[1]).as_str())
   1059             );
   1060         }
   1061     }
   1062 
   1063     #[tokio::test]
   1064     async fn malformed_or_mismatched_checkpoint_fails_before_backend_work() {
   1065         let relay = "wss://one.example";
   1066         let client = Arc::new(MockSubscriptionClient::new([]));
   1067         let transport =
   1068             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
   1069         let base = request(&[relay], 1);
   1070         let target = base.target_set().targets()[0].fingerprint().clone();
   1071         let malformed = base
   1072             .clone()
   1073             .with_checkpoints([SubscriptionCheckpoint::new(
   1074                 target.clone(),
   1075                 FetchCursor::parse("not-a-live-cursor").expect("opaque cursor"),
   1076             )])
   1077             .expect("request checkpoint");
   1078         assert_eq!(
   1079             transport.subscribe(malformed).await.err(),
   1080             Some(radroots_transport::Error::InvalidFetchCursor)
   1081         );
   1082 
   1083         let other_selector = FetchSelector::all()
   1084             .with_exact_tag_value('d', "trade-1")
   1085             .expect("selector");
   1086         let scoped = encode_cursor(
   1087             &base,
   1088             &target,
   1089             &RelayCursor::new(10, "a".repeat(64)).expect("cursor"),
   1090         )
   1091         .expect("scoped cursor");
   1092         for malformed_cursor in [
   1093             FetchCursor::parse("nostr-live-v1:10").expect("opaque missing-field cursor"),
   1094             FetchCursor::parse(format!("{}:extra", scoped.as_str()))
   1095                 .expect("opaque extra-field cursor"),
   1096         ] {
   1097             let checkpoint = SubscriptionCheckpoint::new(target.clone(), malformed_cursor);
   1098             assert_eq!(
   1099                 parse_cursor(&base, &checkpoint),
   1100                 Err(radroots_transport::Error::InvalidFetchCursor)
   1101             );
   1102         }
   1103         let mismatched = base
   1104             .with_selector(other_selector)
   1105             .with_checkpoints([SubscriptionCheckpoint::new(target, scoped)])
   1106             .expect("request checkpoint");
   1107         assert_eq!(
   1108             transport.subscribe(mismatched).await.err(),
   1109             Some(radroots_transport::Error::InvalidFetchCursor)
   1110         );
   1111         assert!(client.queries.lock().expect("queries").is_empty());
   1112     }
   1113 
   1114     #[test]
   1115     fn dropping_an_unpolled_subscription_start_performs_no_backend_work() {
   1116         let relay = "wss://one.example";
   1117         let client = Arc::new(MockSubscriptionClient::new([]));
   1118         let transport =
   1119             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
   1120         let future = transport.subscribe(request(&[relay], 1));
   1121         drop(future);
   1122         assert!(client.queries.lock().expect("queries").is_empty());
   1123     }
   1124 
   1125     #[tokio::test]
   1126     async fn event_limit_and_explicit_cancellation_are_stable_and_idempotent() {
   1127         let relay = "wss://one.example";
   1128         let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay");
   1129         let client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
   1130             RelaySubscriptionItem::Event {
   1131                 relay: relay_url,
   1132                 raw: signed_event("bounded", 1_800_000_001),
   1133             },
   1134         )]));
   1135         let transport =
   1136             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
   1137         let mut subscription = transport
   1138             .subscribe(request(&[relay], 1))
   1139             .await
   1140             .expect("subscription");
   1141         assert!(matches!(
   1142             subscription.next().await.expect("event"),
   1143             SubscriptionNext::Event(_)
   1144         ));
   1145         let SubscriptionNext::End(limit) = subscription.next().await.expect("limit") else {
   1146             panic!("terminal expected");
   1147         };
   1148         assert_eq!(limit.reason(), SubscriptionEndReason::EventLimit);
   1149         assert_eq!(limit.event_count(), 1);
   1150         assert_eq!(limit.checkpoints().len(), 1);
   1151         assert_eq!(subscription.cancel().await.expect("stable cancel"), limit);
   1152         assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 1);
   1153 
   1154         let cancel_client = Arc::new(MockSubscriptionClient::new([]));
   1155         let cancel_transport =
   1156             NostrTransport::with_subscription_client(configured(&[relay]), cancel_client.clone());
   1157         let mut cancelled = cancel_transport
   1158             .subscribe(request(&[relay], 2))
   1159             .await
   1160             .expect("subscription");
   1161         let terminal = cancelled.cancel().await.expect("cancel");
   1162         assert_eq!(terminal.reason(), SubscriptionEndReason::Cancelled);
   1163         assert_eq!(cancelled.cancel().await.expect("repeat cancel"), terminal);
   1164         assert_eq!(cancel_client.cancellations.load(AtomicOrdering::SeqCst), 1);
   1165     }
   1166 
   1167     #[tokio::test]
   1168     async fn dropping_pending_next_requests_cancellation_on_the_next_call() {
   1169         let relay = "wss://one.example";
   1170         let client = Arc::new(MockSubscriptionClient::new([]));
   1171         let transport =
   1172             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
   1173         let mut subscription = transport
   1174             .subscribe(request(&[relay], 2))
   1175             .await
   1176             .expect("subscription");
   1177         let mut pending = subscription.next();
   1178         assert!(futures::poll!(pending.as_mut()).is_pending());
   1179         drop(pending);
   1180         let SubscriptionNext::End(terminal) = subscription.next().await.expect("cancelled") else {
   1181             panic!("terminal expected");
   1182         };
   1183         assert_eq!(terminal.reason(), SubscriptionEndReason::Cancelled);
   1184         assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 1);
   1185     }
   1186 
   1187     #[tokio::test]
   1188     async fn an_admitted_subscription_that_expires_before_next_is_deadline_bounded() {
   1189         let relay = "wss://one.example";
   1190         let client = Arc::new(MockSubscriptionClient::new([]));
   1191         let transport = NostrTransport::with_subscription_client(configured(&[relay]), client);
   1192         let request = SubscriptionRequest::new(
   1193             "expires-after-admission",
   1194             target_set(&[relay]),
   1195             SubscriptionBounds::new(1, unix_time_ms() + 50).expect("bounds"),
   1196         )
   1197         .expect("request");
   1198         let mut subscription = transport.subscribe(request).await.expect("subscription");
   1199         tokio::time::sleep(Duration::from_millis(75)).await;
   1200 
   1201         let SubscriptionNext::End(terminal) = subscription.next().await.expect("deadline") else {
   1202             panic!("terminal expected");
   1203         };
   1204         assert_eq!(terminal.reason(), SubscriptionEndReason::Deadline);
   1205     }
   1206 
   1207     #[tokio::test]
   1208     async fn expired_unrepresentable_closed_and_failed_sources_are_bounded() {
   1209         let relay = "wss://one.example";
   1210         let client = Arc::new(MockSubscriptionClient::new([]));
   1211         let transport =
   1212             NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
   1213         let expired = SubscriptionRequest::new(
   1214             "expired",
   1215             target_set(&[relay]),
   1216             SubscriptionBounds::new(1, 1).expect("bounds"),
   1217         )
   1218         .expect("request");
   1219         let mut expired = transport
   1220             .subscribe(expired)
   1221             .await
   1222             .expect("expired capability");
   1223         let SubscriptionNext::End(terminal) = expired.next().await.expect("deadline") else {
   1224             panic!("terminal expected");
   1225         };
   1226         assert_eq!(terminal.reason(), SubscriptionEndReason::Deadline);
   1227 
   1228         let unsupported = request(&[relay], 1).with_selector(
   1229             FetchSelector::all()
   1230                 .with_kinds(vec![u32::MAX])
   1231                 .expect("selector"),
   1232         );
   1233         let mut unsupported = transport
   1234             .subscribe(unsupported)
   1235             .await
   1236             .expect("closed capability");
   1237         let SubscriptionNext::End(terminal) = unsupported.next().await.expect("source closed")
   1238         else {
   1239             panic!("terminal expected");
   1240         };
   1241         assert_eq!(terminal.reason(), SubscriptionEndReason::SourceClosed);
   1242         assert!(client.queries.lock().expect("queries").is_empty());
   1243 
   1244         let deadline_request = SubscriptionRequest::new(
   1245             "pending-deadline",
   1246             target_set(&[relay]),
   1247             SubscriptionBounds::new(1, unix_time_ms() + 100).expect("bounds"),
   1248         )
   1249         .expect("request");
   1250         let mut deadline = transport
   1251             .subscribe(deadline_request)
   1252             .await
   1253             .expect("deadline capability");
   1254         let SubscriptionNext::End(terminal) = deadline.next().await.expect("deadline") else {
   1255             panic!("terminal expected");
   1256         };
   1257         assert_eq!(terminal.reason(), SubscriptionEndReason::Deadline);
   1258         assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 0);
   1259 
   1260         let failed = Arc::new(MockSubscriptionClient::failing());
   1261         let failed_transport =
   1262             NostrTransport::with_subscription_client(configured(&[relay]), failed.clone());
   1263         assert_eq!(
   1264             failed_transport.subscribe(request(&[relay], 1)).await.err(),
   1265             Some(radroots_transport::Error::SubscriptionUnavailable)
   1266         );
   1267         assert_eq!(failed.queries.lock().expect("queries").len(), 1);
   1268     }
   1269 
   1270     #[tokio::test]
   1271     async fn source_closure_backend_failure_and_event_admission_are_explicit() {
   1272         let relays = ["wss://one.example", "wss://two.example"];
   1273         let one = RelayUrl::parse(relays[0], RelayUrlPolicy::Public).expect("one");
   1274         let two = RelayUrl::parse(relays[1], RelayUrlPolicy::Public).expect("two");
   1275         let client = Arc::new(MockSubscriptionClient::new([
   1276             ScriptedItem::Item(RelaySubscriptionItem::Closed { relay: one }),
   1277             ScriptedItem::Item(RelaySubscriptionItem::Closed { relay: two }),
   1278         ]));
   1279         let transport =
   1280             NostrTransport::with_subscription_client(configured(&relays), client.clone());
   1281         let mut subscription = transport
   1282             .subscribe(request(&relays, 2))
   1283             .await
   1284             .expect("subscription");
   1285         let SubscriptionNext::End(terminal) = subscription.next().await.expect("closed") else {
   1286             panic!("terminal expected");
   1287         };
   1288         assert_eq!(terminal.reason(), SubscriptionEndReason::SourceClosed);
   1289         assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 0);
   1290 
   1291         let unexpected_close = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
   1292             RelaySubscriptionItem::Closed {
   1293                 relay: RelayUrl::parse("wss://unexpected.example", RelayUrlPolicy::Public)
   1294                     .expect("unexpected relay"),
   1295             },
   1296         )]));
   1297         let unexpected_transport =
   1298             NostrTransport::with_subscription_client(configured(&[relays[0]]), unexpected_close);
   1299         let mut unexpected = unexpected_transport
   1300             .subscribe(request(&[relays[0]], 1))
   1301             .await
   1302             .expect("subscription");
   1303         assert_eq!(
   1304             unexpected.next().await,
   1305             Err(radroots_transport::Error::SubscriptionUnavailable)
   1306         );
   1307 
   1308         let error_client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Error]));
   1309         let error_transport =
   1310             NostrTransport::with_subscription_client(configured(&[relays[0]]), error_client);
   1311         let mut subscription = error_transport
   1312             .subscribe(request(&[relays[0]], 1))
   1313             .await
   1314             .expect("subscription");
   1315         assert_eq!(
   1316             subscription.next().await,
   1317             Err(radroots_transport::Error::SubscriptionUnavailable)
   1318         );
   1319 
   1320         let mismatch_client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
   1321             RelaySubscriptionItem::Event {
   1322                 relay: RelayUrl::parse(relays[0], RelayUrlPolicy::Public).expect("relay"),
   1323                 raw: signed_event("wrong-kind", 1_800_000_002),
   1324             },
   1325         )]));
   1326         let mismatch_transport =
   1327             NostrTransport::with_subscription_client(configured(&[relays[0]]), mismatch_client);
   1328         let selector = FetchSelector::all().with_kinds(vec![2]).expect("selector");
   1329         let mut subscription = mismatch_transport
   1330             .subscribe(request(&[relays[0]], 1).with_selector(selector))
   1331             .await
   1332             .expect("subscription");
   1333         assert_eq!(
   1334             subscription.next().await,
   1335             Err(radroots_transport::Error::UnexpectedSubscriptionEvent)
   1336         );
   1337     }
   1338 }