Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions nostr-sdk/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,14 @@
### Breaking change

- Redesign `Relay::count_events` API
- Return `NotificationStream<T>` from client and relay `notifications()` methods instead of a boxed stream.

### Added

- Add `LocalRelay::connections_left` (https://github.com/nostrdevkit/nostr/pull/1459)
- Add `LocalRelayBuilderNip42::relay_url` to make it possible to configure a custom `relay_url` for the nip42 challenge validation (https://github.com/nostrdevkit/nostr/pull/1476)
- Add `LocalRelayBuilder::new_event_channel_size` for customizing the size of the channel used to notify new received events
- Add opt-in receiver gap reporting with `NotificationStream::with_gaps()`.

### Fixed

Expand Down
55 changes: 33 additions & 22 deletions nostr-sdk/src/client/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,11 @@

use std::borrow::Cow;
use std::collections::HashMap;
use std::pin::Pin;
use std::sync::{Arc, Weak};
use std::time::Duration;

use futures::{Stream, StreamExt};
use nostr::prelude::*;
use nostr_database::prelude::*;
use tokio::sync::oneshot;

mod api;
mod builder;
Expand Down Expand Up @@ -194,32 +191,45 @@ impl Client {
///
/// The stream terminates when the client shutdowns.
///
/// Notifications lost when this receiver falls behind are skipped.
/// Use [`NotificationStream::with_gaps`] to observe such losses.
///
/// <div class="warning">When you call this method, you subscribe to the notifications channel from that precise moment. Anything received by relay/s before that moment is not included in the channel!</div>
///
/// # Examples
///
/// ## Report notification gaps
///
/// ```rust,no_run
/// # use nostr_sdk::prelude::*;
/// # async fn example(client: &Client) -> Result<(), Box<dyn std::error::Error>> {
/// // Create the stream before subscribing so it can observe the response.
/// let mut notifications = client.notifications().with_gaps();
/// client.subscribe(Filter::new().kind(Kind::TextNote)).await?;
///
/// while let Some(update) = notifications.next().await {
/// match update {
/// Ok(ClientNotification::Event { event, .. }) => println!("{}", event.id),
/// Ok(ClientNotification::Shutdown) => break,
/// Ok(_) => {}
/// Err(gap) => {
/// eprintln!("Skipped {} notifications", gap.skipped);
/// // Reacquire any subscription coverage the application requires.
/// }
/// }
/// }
/// # Ok(())
/// # }
/// ```
#[inline]
pub fn notifications(&self) -> Pin<Box<dyn Stream<Item = ClientNotification> + Send>> {
pub fn notifications(&self) -> NotificationStream<ClientNotification> {
if self.is_shutdown() {
return Box::pin(futures::stream::empty());
return NotificationStream::empty();
}

// Subscribe to notifications
let rx = self.pool().notifications();

// Create a oneshot channel
let (tx, rx_done) = oneshot::channel();
let mut tx: Option<oneshot::Sender<()>> = Some(tx);

Box::pin(
NotificationStream::new(rx)
.inspect(move |notification| {
if let ClientNotification::Shutdown = &notification {
// Take the sender and send the oneshot notification
if let Some(tx) = tx.take() {
let _ = tx.send(());
}
}
})
.take_until(rx_done),
)
NotificationStream::new(rx, |notification| notification.is_shutdown())
}

/// Get relays from the relay pool.
Expand Down Expand Up @@ -1190,6 +1200,7 @@ impl Client {

#[cfg(test)]
mod tests {
use futures::StreamExt;
use nostr_gossip_memory::prelude::*;

use super::*;
Expand Down
7 changes: 7 additions & 0 deletions nostr-sdk/src/client/notification.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,3 +39,10 @@ pub enum ClientNotification {
/// This notification variant is sent after [`Client::shutdown`](super::Client::shutdown) method is called and all connections have been closed.
Shutdown,
}

impl ClientNotification {
#[inline]
pub(super) fn is_shutdown(&self) -> bool {
matches!(self, Self::Shutdown)
}
}
2 changes: 1 addition & 1 deletion nostr-sdk/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ pub mod prelude;
pub mod proxy;
pub mod relay;
mod shared;
mod stream;
pub mod stream;
#[cfg(test)]
mod test_utils;
pub mod transport;
1 change: 1 addition & 0 deletions nostr-sdk/src/prelude.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,5 @@ pub use crate::policy::*;
#[cfg(not(target_arch = "wasm32"))]
pub use crate::proxy::{self, *};
pub use crate::relay::{self, *};
pub use crate::stream::{self, *};
pub use crate::*;
32 changes: 10 additions & 22 deletions nostr-sdk/src/relay/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,13 @@ use std::cmp;
use std::collections::HashMap;
#[cfg(not(target_arch = "wasm32"))]
use std::net::SocketAddr;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;

use async_utility::time;
use futures::{Stream, StreamExt};
use nostr_database::prelude::*;
use tokio::sync::{broadcast, oneshot};
use tokio::sync::broadcast;

mod api;
mod builder;
Expand Down Expand Up @@ -211,35 +209,24 @@ impl Relay {
///
/// The stream terminates when the relay shutdowns or is banned.
///
/// Notifications lost when this receiver falls behind are skipped.
/// Use [`NotificationStream::with_gaps`] to observe such losses.
///
/// <div class="warning">When you call this method, you subscribe to the notifications channel from that precise moment. Anything received by relay/s before that moment is not included in the channel!</div>
#[inline]
pub fn notifications(&self) -> Pin<Box<dyn Stream<Item = RelayNotification> + Send>> {
pub fn notifications(&self) -> NotificationStream<RelayNotification> {
// If the relay is permanently unusable, return an empty stream
let status: RelayStatus = self.status();
if status.is_banned() || status.is_shutdown() {
return Box::pin(futures::stream::empty());
return NotificationStream::empty();
}

// Subscribe to notifications
let rx = self.inner.internal_notification_sender.subscribe();

// Create a oneshot channel
let (tx, rx_done) = oneshot::channel();
let mut tx: Option<oneshot::Sender<()>> = Some(tx);

Box::pin(
NotificationStream::new(rx)
.inspect(move |notification| {
if let RelayNotification::RelayStatus { status } = &notification {
if status.is_banned() || status.is_shutdown() {
// Take the sender and send the oneshot notification
if let Some(tx) = tx.take() {
let _ = tx.send(());
}
}
}
})
.take_until(rx_done),
NotificationStream::new(
rx,
|notification| matches!(notification, RelayNotification::RelayStatus { status } if status.is_banned() || status.is_shutdown()),
)
}

Expand Down Expand Up @@ -432,6 +419,7 @@ mod tests {
use std::sync::Arc;

use async_utility::time;
use futures::StreamExt;

use super::*;
use crate::error::{Error, ErrorKind};
Expand Down
61 changes: 0 additions & 61 deletions nostr-sdk/src/stream.rs

This file was deleted.

30 changes: 30 additions & 0 deletions nostr-sdk/src/stream/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
//! Streams

use std::pin::Pin;
use std::task::{Context, Poll};

use futures::Stream;
use tokio::sync::mpsc::Receiver;

mod notification;

pub use self::notification::*;

pub(crate) struct ReceiverStream<T> {
inner: Receiver<T>,
}

impl<T> ReceiverStream<T> {
#[inline]
pub(crate) fn new(recv: Receiver<T>) -> Self {
Self { inner: recv }
}
}

impl<T> Stream for ReceiverStream<T> {
type Item = T;

fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.inner.poll_recv(cx)
}
}
Loading
Loading