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 }