commit de88ddf56a389f1424c704d73f3322aa926cecc8
parent ed3016d3620d98cfde1b53901e12eebdd73f6b85
Author: triesap <tyson@radroots.org>
Date: Sun, 16 Aug 2026 21:30:41 +0000
transport-nostr: implement bounded live subscriptions
- Intent: Implement explicit bounded Nostr live subscriptions with per-target provenance and checkpoints.
- Why: The transport-neutral subscription SPI requires an adapter with deadline, cancellation, and equal-timestamp reconnect semantics.
- Verification: Extbuild package tests, loopback I/O, Clippy, Rustdoc, workspace check, contract validation, and API baseline comparison passed.
- Risk: Upstream relay cleanup is best effort on explicit cancellation and remains independently bounded by the request auto-close deadline.
Diffstat:
12 files changed, 1478 insertions(+), 41 deletions(-)
diff --git a/contracts/api_baselines/radroots_transport.txt b/contracts/api_baselines/radroots_transport.txt
@@ -148,6 +148,7 @@ pub radroots_transport::error::Error::SubscriptionCheckpointSetTooLarge
pub radroots_transport::error::Error::SubscriptionEndLimitExceeded
pub radroots_transport::error::Error::SubscriptionEndRequestMismatch
pub radroots_transport::error::Error::SubscriptionEventCheckpointMismatch
+pub radroots_transport::error::Error::SubscriptionUnavailable
pub radroots_transport::error::Error::TargetSetTooLarge
pub radroots_transport::error::Error::TransportOutcomeStatusMismatch
pub radroots_transport::error::Error::UnexpectedDeliveryTargetReceipt
@@ -626,6 +627,7 @@ pub radroots_transport::Error::SubscriptionCheckpointSetTooLarge
pub radroots_transport::Error::SubscriptionEndLimitExceeded
pub radroots_transport::Error::SubscriptionEndRequestMismatch
pub radroots_transport::Error::SubscriptionEventCheckpointMismatch
+pub radroots_transport::Error::SubscriptionUnavailable
pub radroots_transport::Error::TargetSetTooLarge
pub radroots_transport::Error::TransportOutcomeStatusMismatch
pub radroots_transport::Error::UnexpectedDeliveryTargetReceipt
diff --git a/contracts/api_baselines/radroots_transport_nostr.txt b/contracts/api_baselines/radroots_transport_nostr.txt
@@ -106,6 +106,8 @@ pub fn radroots_transport_nostr::NostrTransport::status(&self) -> radroots_trans
impl radroots_transport::source::EventSource for radroots_transport_nostr::NostrTransport
pub fn radroots_transport_nostr::NostrTransport::fetch(&self, radroots_transport::source::FetchRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::FetchPage, radroots_transport::error::Error>>
pub fn radroots_transport_nostr::NostrTransport::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::status::SourceStatus, radroots_transport::error::Error>>
+impl radroots_transport::source::EventSubscriber for radroots_transport_nostr::NostrTransport
+pub fn radroots_transport_nostr::NostrTransport::subscribe(&self, radroots_transport::source::SubscriptionRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::BoxSubscription, radroots_transport::error::Error>>
pub struct radroots_transport_nostr::ReconnectBackoff
impl radroots_transport_nostr::ReconnectBackoff
pub const fn radroots_transport_nostr::ReconnectBackoff::initial_delay_ms(self) -> u64
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -47,6 +47,7 @@ pub enum Error {
SubscriptionEndLimitExceeded,
InvalidSubscriptionEnd,
SubscriptionEndRequestMismatch,
+ SubscriptionUnavailable,
InvalidSatisfactionPolicy,
EmptyRequiredTargetSet,
DuplicateRequiredTargetFingerprint,
@@ -167,6 +168,7 @@ impl fmt::Display for Error {
Self::SubscriptionEndRequestMismatch => {
f.write_str("transport subscription result does not match its request")
}
+ Self::SubscriptionUnavailable => f.write_str("transport subscription is unavailable"),
Self::InvalidSatisfactionPolicy => {
f.write_str("transport satisfaction policy is invalid")
}
diff --git a/crates/transport/tests/error_contract.rs b/crates/transport/tests/error_contract.rs
@@ -145,6 +145,10 @@ fn every_transport_error_has_stable_operator_facing_text() {
"transport subscription result does not match its request",
),
(
+ Error::SubscriptionUnavailable,
+ "transport subscription is unavailable",
+ ),
+ (
Error::InvalidSatisfactionPolicy,
"transport satisfaction policy is invalid",
),
diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md
@@ -1,11 +1,12 @@
# radroots_transport_nostr
`radroots_transport_nostr` is the concrete, signer-free Nostr implementation
-of the generic [`radroots_transport::EventSource`] and
-[`radroots_transport::EventSink`] interfaces. It validates relay configuration
-and network policy, performs bounded fetch and delivery attempts, exposes
-explicit host-mediated NIP-42 authentication, and normalizes relay outcomes
-and passive status.
+of the generic [`radroots_transport::EventSource`],
+[`radroots_transport::EventSubscriber`], and [`radroots_transport::EventSink`]
+interfaces. It validates relay configuration and network policy, performs
+bounded fetch, live-subscription, and delivery attempts, exposes explicit
+host-mediated NIP-42 authentication, and normalizes relay outcomes and passive
+status.
The crate does not own event ingestion, persistence, outbox claiming,
projection refresh, durable retry scheduling, or a process runtime. Those
@@ -23,7 +24,7 @@ Configuration is explicit, validated, and inert. Constructing
[`NostrTransport`] creates no socket and performs no DNS lookup:
```rust
-use radroots_transport::{EventSink, EventSource};
+use radroots_transport::{EventSink, EventSource, EventSubscriber};
use radroots_transport_nostr::{
Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile,
RelayProfileKind, RelayUrlPolicy,
@@ -39,17 +40,20 @@ let config = Config::from_profile(profile).with_timeouts(5_000, 20_000, 2_000)?;
let transport = NostrTransport::new(config);
let source: &dyn EventSource = &transport;
+let subscriber: &dyn EventSubscriber = &transport;
let sink: &dyn EventSink = &transport;
drop(source.status());
+let _ = subscriber;
drop(sink.status());
# Ok::<(), Box<dyn std::error::Error>>(())
```
A runnable version is available at
[`examples/configure_transport.rs`](examples/configure_transport.rs).
-The composing host constructs bounded `FetchRequest` and `DeliveryRequest`
-values from `radroots_transport`, polls the returned futures on its executor,
-and applies any retry or scheduling policy outside this crate.
+The composing host constructs bounded `FetchRequest`, `SubscriptionRequest`,
+and `DeliveryRequest` values from `radroots_transport`, polls the returned
+futures on its executor, and applies any retry or scheduling policy outside
+this crate.
## Public surface
@@ -63,7 +67,7 @@ and applies any retry or scheduling policy outside this crate.
trusted private-network destination rules.
- [`RelayCursor`] provides the equal-timestamp-safe event ordering primitive
used by scoped fetch continuation cursors.
-- [`NostrTransport`] implements both transport SPIs, exposes passive typed
+- [`NostrTransport`] implements all three transport SPIs, exposes passive typed
per-relay evidence, and provides explicit NIP-42 challenge lifecycle methods.
- [`Error`] contains only package-owned validation and authentication errors;
upstream failures are normalized before crossing the public boundary.
@@ -96,16 +100,25 @@ Callers that resolve addresses outside the adapter may use
[`RelayUrl::validate_resolved_addresses`] to apply the same destination-class
check before handing control to another network boundary.
-## Fetch, delivery, and outcome behavior
-
-Fetch accepts only configured readable Nostr targets, translates transport-neutral kind,
-author, and event-time selectors into Nostr filters, reapplies those selectors
-defensively, applies the request page bound, deduplicates events by event ID,
-preserves per-relay provenance, and emits an opaque versioned cursor bound to
-the exact target set and selector when more results remain. Equal timestamps
-are ordered by event ID so overlap-safe reconnect pagination cannot skip peers.
-Malformed relay events are ignored and reported as a partial target outcome
-rather than admitted.
+## Fetch, live subscription, delivery, and outcome behavior
+
+Fetch accepts only configured readable Nostr targets, translates
+transport-neutral kind, author, and event-time selectors into Nostr filters,
+reapplies those selectors defensively, applies the request page bound,
+deduplicates events by event ID, preserves per-relay provenance, and emits an
+opaque versioned cursor bound to the exact target set and selector when more
+results remain. Equal timestamps are ordered by event ID so overlap-safe
+reconnect pagination cannot skip peers. Malformed relay events are ignored and
+reported as a partial target outcome rather than admitted.
+
+Live subscriptions use the same explicit readable targets and selector
+translation. A caller checkpoint is scoped to one exact target and selector;
+the adapter reconnects with Nostr's inclusive `since` timestamp and suppresses
+only events at or before that checkpoint's event-ID tie breaker. This preserves
+every later event sharing the same second-granular timestamp. Each emitted
+event carries exact relay provenance and a new canonical target checkpoint.
+Event limits, absolute deadlines, explicit cancellation, source closure, and
+stable repeated terminal results follow the generic subscription contract.
Delivery converts an already validated signed Radroots event to Nostr, attempts
each configured writable target once, and returns one normalized receipt entry per
@@ -128,12 +141,21 @@ The absolute deadline in each generic request bounds the complete operation.
The configured connection and request timeouts are upper bounds within that
remaining budget. An already-expired request performs no relay work.
-Dropping an unpolled fetch or delivery future performs no I/O. Once polled,
-cancellation is best effort at the socket boundary. For delivery, submission
-to a relay is the remote commit point: after a relay accepts the event, dropping
-the future cannot retract it. A missing final response is therefore reported
-as unavailable or unknown evidence, never as proof that no publication
-occurred. Fetch is observational and has no local durable commit point.
+Dropping an unpolled fetch, subscription-start, or delivery future performs no
+I/O. Once polled, cancellation is best effort at the socket boundary. For
+delivery, submission to a relay is the remote commit point: after a relay
+accepts the event, dropping the future cannot retract it. A missing final
+response is therefore reported as unavailable or unknown evidence, never as
+proof that no publication occurred. Fetch and live observation have no local
+durable commit point.
+
+Dropping a pending subscription `next` or `cancel` future records a
+cancellation request in the retained capability; its next operation awaits
+relay unsubscription and returns the stable cancelled terminal result.
+Dropping the capability itself cannot await network cleanup. Every published
+relay subscription therefore also carries an upstream auto-close deadline
+bounded by the request's absolute deadline, so remote work cannot continue
+indefinitely.
NIP-42 authentication is explicit. [`NostrTransport::begin_authentication`]
records one bounded relay challenge and returns the exact host signing input.
@@ -160,15 +182,17 @@ event JSON.
The package has no Cargo features. It is a standard-library native adapter;
Tokio and the upstream Nostr relay client are private implementation choices.
-The crate never creates an executor, installs a runtime, spawns a background
-worker, installs a tracing subscriber, or owns process lifecycle. The host must
-poll operations from a compatible executor and provide clock/deadline policy
-through the generic requests.
+The crate never creates an executor, installs a runtime, launches an
+adapter-owned worker, installs a tracing subscriber, or owns process lifecycle.
+The host must poll operations from a compatible executor and provide
+clock/deadline policy through the generic requests. After an explicit operation
+begins, the private upstream relay client owns its ordinary socket tasks and
+the bounded auto-close timer for live relay subscriptions.
## Intended consumers
-- `radroots_sync` composes this source and sink with verification, storage,
- projection, outbox, and explicit retry decisions.
+- `radroots_sync` composes this source, subscriber, and sink with verification,
+ storage, projection, outbox, and explicit retry decisions.
- `radroots_sdk` selects and configures the adapter for advanced applications.
- Native services may compose it directly behind the generic transport SPIs.
diff --git a/crates/transport_nostr/examples/configure_transport.rs b/crates/transport_nostr/examples/configure_transport.rs
@@ -1,4 +1,4 @@
-use radroots_transport::{EventSink, EventSource};
+use radroots_transport::{EventSink, EventSource, EventSubscriber};
use radroots_transport_nostr::{
Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind,
RelayUrlPolicy,
@@ -15,10 +15,12 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
let transport = NostrTransport::new(config);
let source: &dyn EventSource = &transport;
+ let subscriber: &dyn EventSubscriber = &transport;
let sink: &dyn EventSink = &transport;
drop(source.status());
+ let _ = subscriber;
drop(sink.status());
- println!("configured a Nostr source and sink without opening a relay connection");
+ println!("configured a Nostr source, subscriber, and sink without opening a relay connection");
Ok(())
}
diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs
@@ -2,7 +2,10 @@
use crate::{Error, RelayEndpoint, RelayProfile, RelayProfileKind, RelayStatusReport, RelayUrl};
use core::fmt;
-use std::sync::Arc;
+use std::sync::{
+ Arc,
+ atomic::{AtomicU64, Ordering},
+};
/// Maximum relay targets accepted by one transport instance.
pub(crate) const MAX_RELAYS: usize = 64;
@@ -216,8 +219,10 @@ pub struct NostrTransport {
config: Config,
pub(crate) client: Arc<dyn crate::sink::RelayClient>,
pub(crate) source_client: Arc<dyn crate::source::RelaySourceClient>,
+ pub(crate) subscription_client: Arc<dyn crate::subscription::RelaySubscriptionClient>,
pub(crate) auth: Arc<crate::auth::AuthFlow>,
pub(crate) status: Arc<crate::status::StatusTracker>,
+ subscription_sequence: Arc<AtomicU64>,
}
impl NostrTransport {
@@ -234,10 +239,14 @@ impl NostrTransport {
config,
client: Arc::new(crate::sink::LiveRelayClient::new(client.clone())),
source_client: Arc::new(crate::source::LiveRelaySourceClient::new(client.clone())),
+ subscription_client: Arc::new(crate::subscription::LiveRelaySubscriptionClient::new(
+ client.clone(),
+ )),
auth: Arc::new(crate::auth::AuthFlow::new(Arc::new(
crate::auth::LiveAuthClient::new(client),
))),
status,
+ subscription_sequence: Arc::new(AtomicU64::new(0)),
}
}
@@ -259,8 +268,12 @@ impl NostrTransport {
config,
client,
source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()),
+ subscription_client: Arc::new(
+ crate::subscription::LiveRelaySubscriptionClient::isolated(),
+ ),
auth: Arc::new(crate::auth::AuthFlow::isolated()),
status,
+ subscription_sequence: Arc::new(AtomicU64::new(0)),
}
}
@@ -274,10 +287,40 @@ impl NostrTransport {
config,
client: Arc::new(crate::sink::LiveRelayClient::isolated()),
source_client,
+ subscription_client: Arc::new(
+ crate::subscription::LiveRelaySubscriptionClient::isolated(),
+ ),
auth: Arc::new(crate::auth::AuthFlow::isolated()),
status,
+ subscription_sequence: Arc::new(AtomicU64::new(0)),
}
}
+
+ #[cfg(test)]
+ pub(crate) fn with_subscription_client(
+ config: Config,
+ subscription_client: Arc<dyn crate::subscription::RelaySubscriptionClient>,
+ ) -> Self {
+ let status = Arc::new(crate::status::StatusTracker::new(&config));
+ Self {
+ config,
+ client: Arc::new(crate::sink::LiveRelayClient::isolated()),
+ source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()),
+ subscription_client,
+ auth: Arc::new(crate::auth::AuthFlow::isolated()),
+ status,
+ subscription_sequence: Arc::new(AtomicU64::new(0)),
+ }
+ }
+
+ pub(crate) fn next_subscription_sequence(&self) -> Result<u64, radroots_transport::Error> {
+ self.subscription_sequence
+ .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |current| {
+ current.checked_add(1)
+ })
+ .map(|previous| previous + 1)
+ .map_err(|_| radroots_transport::Error::SubscriptionUnavailable)
+ }
}
impl fmt::Debug for NostrTransport {
@@ -384,4 +427,26 @@ mod tests {
assert!(debug.contains("NostrTransport"));
assert!(!debug.contains("client"));
}
+
+ #[test]
+ fn subscription_sequence_is_monotonic_and_fails_closed_at_overflow() {
+ let config = Config::from_profile(
+ crate::profile::test_profile(
+ RelayProfileKind::Public,
+ crate::RelayUrlPolicy::Public,
+ ["wss://relay.example.com"],
+ )
+ .expect("profile"),
+ );
+ let transport = NostrTransport::new(config);
+ assert_eq!(transport.next_subscription_sequence(), Ok(1));
+ assert_eq!(transport.next_subscription_sequence(), Ok(2));
+ transport
+ .subscription_sequence
+ .store(u64::MAX, Ordering::SeqCst);
+ assert_eq!(
+ transport.next_subscription_sequence(),
+ Err(radroots_transport::Error::SubscriptionUnavailable)
+ );
+ }
}
diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs
@@ -11,6 +11,7 @@ mod relay;
mod sink;
mod source;
mod status;
+mod subscription;
pub use client::{Config, NostrTransport, ReconnectBackoff};
pub use cursor::RelayCursor;
diff --git a/crates/transport_nostr/src/subscription.rs b/crates/transport_nostr/src/subscription.rs
@@ -0,0 +1,1176 @@
+//! Bounded Nostr live-subscription adapter.
+
+use crate::{NostrTransport, RelayCursor, RelayUrl};
+use nostr_sdk::prelude::{
+ ClientMessage, Filter, JsonUtil, Kind, RelayMessage, RelayPoolNotification, ReqExitPolicy,
+ SubscribeAutoCloseOptions, SubscribeOptions, SubscriptionId, Timestamp,
+};
+use radroots_transport::{
+ BoxFuture, BoxSubscription, EventSubscriber, EventSubscription, SubscriptionEnd,
+ SubscriptionEndReason, SubscriptionEvent, SubscriptionNext, SubscriptionRequest,
+ source::{EventProvenance, FetchCursor, ObservedEvent, SubscriptionCheckpoint},
+};
+use sha2::{Digest, Sha256};
+use std::{
+ collections::{BTreeMap, BTreeSet},
+ sync::{
+ Arc,
+ atomic::{AtomicBool, Ordering},
+ },
+ time::{Duration, SystemTime, UNIX_EPOCH},
+};
+
+const CURSOR_PREFIX: &str = "nostr-live-v1";
+const CURSOR_SCOPE_DOMAIN: &[u8] = b"radroots.transport-nostr.subscription-cursor.v1\0";
+const SUBSCRIPTION_ID_DOMAIN: &[u8] = b"radroots.transport-nostr.subscription-id.v1\0";
+
+#[derive(Clone, Debug)]
+pub(crate) struct RelaySubscriptionQuery {
+ id: String,
+ targets: Vec<RelaySubscriptionTarget>,
+ selector: radroots_transport::source::FetchSelector,
+ connect_timeout: Duration,
+ timeout: Duration,
+}
+
+#[derive(Clone, Debug)]
+struct RelaySubscriptionTarget {
+ relay: RelayUrl,
+ since_unix_seconds: Option<u64>,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub(crate) enum RelaySubscriptionItem {
+ Event { relay: RelayUrl, raw: String },
+ Closed { relay: RelayUrl },
+ Shutdown,
+}
+
+pub(crate) trait RelaySubscriptionSession: Send {
+ fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>>;
+ fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>>;
+}
+
+pub(crate) trait RelaySubscriptionClient: Send + Sync {
+ fn subscribe(
+ &self,
+ query: RelaySubscriptionQuery,
+ ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>>;
+}
+
+#[derive(Clone, Debug)]
+pub(crate) struct LiveRelaySubscriptionClient {
+ client: nostr_sdk::Client,
+}
+
+impl LiveRelaySubscriptionClient {
+ pub(crate) const fn new(client: nostr_sdk::Client) -> Self {
+ Self { client }
+ }
+
+ #[cfg(test)]
+ pub(crate) fn isolated() -> Self {
+ let client = nostr_sdk::Client::default();
+ client.automatic_authentication(false);
+ Self::new(client)
+ }
+}
+
+impl RelaySubscriptionClient for LiveRelaySubscriptionClient {
+ #[cfg_attr(coverage_nightly, coverage(off))]
+ fn subscribe(
+ &self,
+ query: RelaySubscriptionQuery,
+ ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>> {
+ Box::pin(async move {
+ let notifications = self.client.notifications();
+ let subscription_id = SubscriptionId::new(query.id);
+ let mut targeted = Vec::with_capacity(query.targets.len());
+ let mut relay_lookup = BTreeMap::new();
+
+ for target in query.targets {
+ let url = target.relay.as_str().to_owned();
+ self.client.add_relay(url.as_str()).await.map_err(|_| ())?;
+ self.client
+ .try_connect_relay(url.as_str(), query.connect_timeout)
+ .await
+ .map_err(|_| ())?;
+ let filter = subscription_filter(&query.selector, target.since_unix_seconds)?;
+ relay_lookup.insert(url.clone(), target.relay);
+ targeted.push((url, filter));
+ }
+
+ let auto_close = SubscribeAutoCloseOptions::default()
+ .exit_policy(ReqExitPolicy::WaitDurationAfterEOSE(query.timeout))
+ .timeout(Some(query.timeout));
+ let output = self
+ .client
+ .subscribe_targeted(
+ subscription_id.clone(),
+ targeted,
+ SubscribeOptions::default().close_on(Some(auto_close)),
+ )
+ .await
+ .map_err(|_| ())?;
+ if !output.failed.is_empty() || output.success.len() != relay_lookup.len() {
+ let _ = self
+ .client
+ .send_msg_to(
+ relay_lookup.keys().map(String::as_str),
+ ClientMessage::close(subscription_id.clone()),
+ )
+ .await;
+ self.client.unsubscribe(&subscription_id).await;
+ return Err(());
+ }
+
+ // Keep the receiver created before REQ publication so an immediate
+ // relay event cannot race ahead of local observation.
+ let session = LiveRelaySubscriptionSession {
+ client: self.client.clone(),
+ subscription_id,
+ notifications: Some(notifications),
+ relay_lookup,
+ cancelled: false,
+ };
+ Ok(Box::new(session) as Box<dyn RelaySubscriptionSession>)
+ })
+ }
+}
+
+struct LiveRelaySubscriptionSession {
+ client: nostr_sdk::Client,
+ subscription_id: SubscriptionId,
+ notifications: Option<tokio::sync::broadcast::Receiver<RelayPoolNotification>>,
+ relay_lookup: BTreeMap<String, RelayUrl>,
+ cancelled: bool,
+}
+
+impl RelaySubscriptionSession for LiveRelaySubscriptionSession {
+ #[cfg_attr(coverage_nightly, coverage(off))]
+ fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>> {
+ Box::pin(async move {
+ let Some(notifications) = self.notifications.as_mut() else {
+ return Ok(RelaySubscriptionItem::Shutdown);
+ };
+ loop {
+ let notification = match notifications.recv().await {
+ Ok(notification) => notification,
+ Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => return Err(()),
+ Err(tokio::sync::broadcast::error::RecvError::Closed) => {
+ return Ok(RelaySubscriptionItem::Shutdown);
+ }
+ };
+ match notification {
+ RelayPoolNotification::Message { relay_url, message } => {
+ let Some(relay) = self.relay_lookup.get(relay_url.as_str()).cloned() else {
+ continue;
+ };
+ match message {
+ RelayMessage::Event {
+ subscription_id,
+ event,
+ } if subscription_id.as_ref() == &self.subscription_id => {
+ return Ok(RelaySubscriptionItem::Event {
+ relay,
+ raw: event.as_json(),
+ });
+ }
+ RelayMessage::Closed {
+ subscription_id, ..
+ } if subscription_id.as_ref() == &self.subscription_id => {
+ return Ok(RelaySubscriptionItem::Closed { relay });
+ }
+ _ => {}
+ }
+ }
+ RelayPoolNotification::Shutdown => {
+ return Ok(RelaySubscriptionItem::Shutdown);
+ }
+ RelayPoolNotification::Event { .. } => {}
+ }
+ }
+ })
+ }
+
+ #[cfg_attr(coverage_nightly, coverage(off))]
+ fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>> {
+ Box::pin(async move {
+ if !self.cancelled {
+ let output = self
+ .client
+ .send_msg_to(
+ self.relay_lookup.keys().map(String::as_str),
+ ClientMessage::close(self.subscription_id.clone()),
+ )
+ .await
+ .map_err(|_| ())?;
+ if !output.failed.is_empty() || output.success.len() != self.relay_lookup.len() {
+ return Err(());
+ }
+ self.client.unsubscribe(&self.subscription_id).await;
+ self.cancelled = true;
+ self.notifications = None;
+ }
+ Ok(())
+ })
+ }
+}
+
+struct RelayEventSubscription {
+ request: SubscriptionRequest,
+ session: Option<Box<dyn RelaySubscriptionSession>>,
+ targets: BTreeMap<RelayUrl, radroots_transport::Target>,
+ active_relays: BTreeSet<RelayUrl>,
+ cursors: BTreeMap<radroots_transport::target::TargetFingerprint, RelayCursor>,
+ checkpoints: BTreeMap<radroots_transport::target::TargetFingerprint, SubscriptionCheckpoint>,
+ event_count: u16,
+ terminal: Option<SubscriptionEnd>,
+ cancellation_requested: Arc<AtomicBool>,
+ status: Arc<crate::status::StatusTracker>,
+}
+
+impl RelayEventSubscription {
+ fn ended(
+ request: SubscriptionRequest,
+ reason: SubscriptionEndReason,
+ status: Arc<crate::status::StatusTracker>,
+ ) -> Result<Self, ()> {
+ let terminal = SubscriptionEnd::for_request(&request, 0, [], reason).map_err(|_| ())?;
+ Ok(Self {
+ request,
+ session: None,
+ targets: BTreeMap::new(),
+ active_relays: BTreeSet::new(),
+ cursors: BTreeMap::new(),
+ checkpoints: BTreeMap::new(),
+ event_count: 0,
+ terminal: Some(terminal),
+ cancellation_requested: Arc::new(AtomicBool::new(false)),
+ status,
+ })
+ }
+
+ async fn next_inner(&mut self) -> Result<SubscriptionNext, radroots_transport::Error> {
+ if let Some(terminal) = &self.terminal {
+ return Ok(SubscriptionNext::End(terminal.clone()));
+ }
+ if self.cancellation_requested.swap(false, Ordering::SeqCst) {
+ return self
+ .terminate(SubscriptionEndReason::Cancelled)
+ .await
+ .map(SubscriptionNext::End);
+ }
+ if self.event_count >= self.request.bounds().event_limit() {
+ return self
+ .terminate(SubscriptionEndReason::EventLimit)
+ .await
+ .map(SubscriptionNext::End);
+ }
+
+ loop {
+ let remaining = remaining_duration(self.request.bounds().deadline_unix_ms());
+ if remaining.is_zero() {
+ return self
+ .terminate(SubscriptionEndReason::Deadline)
+ .await
+ .map(SubscriptionNext::End);
+ }
+ let item = {
+ let session = self
+ .session
+ .as_mut()
+ .ok_or(radroots_transport::Error::SubscriptionUnavailable)?;
+ tokio::time::timeout(remaining, session.next()).await
+ };
+ let item = match item {
+ Ok(Ok(item)) => item,
+ Ok(Err(())) => return Err(radroots_transport::Error::SubscriptionUnavailable),
+ Err(_) => {
+ return self
+ .terminate(SubscriptionEndReason::Deadline)
+ .await
+ .map(SubscriptionNext::End);
+ }
+ };
+ match item {
+ RelaySubscriptionItem::Event { relay, raw } => {
+ if let Some(event) = self.admit_event(&relay, raw.as_str())? {
+ return Ok(SubscriptionNext::Event(Box::new(event)));
+ }
+ }
+ RelaySubscriptionItem::Closed { relay } => {
+ if !self.active_relays.remove(&relay) {
+ return Err(radroots_transport::Error::SubscriptionUnavailable);
+ }
+ if self.active_relays.is_empty() {
+ return self
+ .terminate(SubscriptionEndReason::SourceClosed)
+ .await
+ .map(SubscriptionNext::End);
+ }
+ }
+ RelaySubscriptionItem::Shutdown => {
+ return self
+ .terminate(SubscriptionEndReason::SourceClosed)
+ .await
+ .map(SubscriptionNext::End);
+ }
+ }
+ }
+ }
+
+ fn admit_event(
+ &mut self,
+ relay: &RelayUrl,
+ raw: &str,
+ ) -> Result<Option<SubscriptionEvent>, radroots_transport::Error> {
+ let target = self
+ .targets
+ .get(relay)
+ .ok_or(radroots_transport::Error::UnexpectedSubscriptionEvent)?;
+ let event = radroots_event_codec::decode::signed_event(raw)
+ .map_err(|_| radroots_transport::Error::UnexpectedSubscriptionEvent)?;
+ if !self.request.selector().matches(&event) {
+ return Err(radroots_transport::Error::UnexpectedSubscriptionEvent);
+ }
+ if self
+ .cursors
+ .get(target.fingerprint())
+ .is_some_and(|cursor| !cursor.precedes(event.created_at(), event.id_str()))
+ {
+ return Ok(None);
+ }
+
+ let cursor = RelayCursor::new(event.created_at(), event.id_str())
+ .map_err(|_| radroots_transport::Error::UnexpectedSubscriptionEvent)?;
+ let opaque = encode_cursor(&self.request, target.fingerprint(), &cursor)?;
+ let checkpoint = SubscriptionCheckpoint::new(target.fingerprint().clone(), opaque.clone());
+ let provenance = EventProvenance::new(
+ radroots_transport::TransportId::NOSTR,
+ target.fingerprint().clone(),
+ unix_time_ms().max(1),
+ )?
+ .with_cursor(opaque);
+ let observed = ObservedEvent::new(event, provenance);
+ let subscription_event =
+ SubscriptionEvent::for_request(&self.request, observed, checkpoint.clone())?;
+
+ self.cursors.insert(target.fingerprint().clone(), cursor);
+ self.checkpoints
+ .insert(target.fingerprint().clone(), checkpoint);
+ self.event_count = self.event_count.saturating_add(1);
+ self.status
+ .record_read(relay, true, false, unix_time_ms().max(1));
+ Ok(Some(subscription_event))
+ }
+
+ async fn terminate(
+ &mut self,
+ mut reason: SubscriptionEndReason,
+ ) -> Result<SubscriptionEnd, radroots_transport::Error> {
+ if let Some(terminal) = &self.terminal {
+ return Ok(terminal.clone());
+ }
+ if matches!(
+ reason,
+ SubscriptionEndReason::EventLimit | SubscriptionEndReason::Cancelled
+ ) {
+ let remaining = remaining_duration(self.request.bounds().deadline_unix_ms());
+ if remaining.is_zero() {
+ reason = SubscriptionEndReason::Deadline;
+ } else if let Some(session) = self.session.as_mut() {
+ match tokio::time::timeout(remaining, session.cancel()).await {
+ Ok(Ok(())) => {}
+ Ok(Err(())) => {
+ return Err(radroots_transport::Error::SubscriptionUnavailable);
+ }
+ Err(_) => reason = SubscriptionEndReason::Deadline,
+ }
+ }
+ }
+ self.session = None;
+ let terminal = SubscriptionEnd::for_request(
+ &self.request,
+ self.event_count,
+ self.checkpoints.values().cloned(),
+ reason,
+ )?;
+ self.terminal = Some(terminal.clone());
+ Ok(terminal)
+ }
+}
+
+impl EventSubscription for RelayEventSubscription {
+ fn request(&self) -> &SubscriptionRequest {
+ &self.request
+ }
+
+ fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, radroots_transport::Error>> {
+ let cancellation = CancellationOnDrop::new(Arc::clone(&self.cancellation_requested));
+ Box::pin(async move {
+ let result = self.next_inner().await;
+ cancellation.complete();
+ result
+ })
+ }
+
+ fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, radroots_transport::Error>> {
+ let cancellation = CancellationOnDrop::new(Arc::clone(&self.cancellation_requested));
+ Box::pin(async move {
+ let result = self.terminate(SubscriptionEndReason::Cancelled).await;
+ cancellation.complete();
+ result
+ })
+ }
+}
+
+impl EventSubscriber for NostrTransport {
+ fn subscribe(
+ &self,
+ request: SubscriptionRequest,
+ ) -> BoxFuture<'_, Result<BoxSubscription, radroots_transport::Error>> {
+ Box::pin(async move {
+ let remaining = remaining_duration(request.bounds().deadline_unix_ms());
+ if remaining.is_zero() {
+ return RelayEventSubscription::ended(
+ request,
+ SubscriptionEndReason::Deadline,
+ Arc::clone(&self.status),
+ )
+ .map(|subscription| Box::new(subscription) as BoxSubscription)
+ .map_err(|()| radroots_transport::Error::SubscriptionUnavailable);
+ }
+ if !selector_is_representable(request.selector()) {
+ return RelayEventSubscription::ended(
+ request,
+ SubscriptionEndReason::SourceClosed,
+ Arc::clone(&self.status),
+ )
+ .map(|subscription| Box::new(subscription) as BoxSubscription)
+ .map_err(|()| radroots_transport::Error::SubscriptionUnavailable);
+ }
+
+ let now_ms = unix_time_ms();
+ let mut targets = BTreeMap::new();
+ let mut active_relays = BTreeSet::new();
+ for target in request.target_set().targets() {
+ let endpoint = self
+ .config()
+ .endpoint_for_target(target)
+ .filter(|endpoint| endpoint.access().can_read())
+ .ok_or(radroots_transport::Error::SubscriptionUnavailable)?;
+ if !self.status.may_read(endpoint.url(), now_ms)
+ || targets
+ .insert(endpoint.url().clone(), target.clone())
+ .is_some()
+ {
+ return Err(radroots_transport::Error::SubscriptionUnavailable);
+ }
+ active_relays.insert(endpoint.url().clone());
+ }
+
+ let mut cursors = BTreeMap::new();
+ let mut checkpoints = BTreeMap::new();
+ let mut target_queries = Vec::with_capacity(targets.len());
+ for (relay, target) in &targets {
+ let checkpoint = request
+ .checkpoints()
+ .iter()
+ .find(|checkpoint| checkpoint.target() == target.fingerprint());
+ let cursor = checkpoint
+ .map(|checkpoint| parse_cursor(&request, checkpoint))
+ .transpose()?;
+ let since = cursor
+ .as_ref()
+ .map(RelayCursor::created_at_unix_s)
+ .or(request.selector().since_unix_seconds());
+ if let Some(cursor) = cursor {
+ cursors.insert(target.fingerprint().clone(), cursor);
+ }
+ if let Some(checkpoint) = checkpoint {
+ checkpoints.insert(target.fingerprint().clone(), checkpoint.clone());
+ }
+ target_queries.push(RelaySubscriptionTarget {
+ relay: relay.clone(),
+ since_unix_seconds: since,
+ });
+ }
+
+ let timeout = remaining.min(Duration::from_millis(self.config().request_timeout_ms()));
+ let id = subscription_id(&request, self.next_subscription_sequence()?);
+ for relay in &active_relays {
+ self.status.begin_read(relay, now_ms);
+ }
+ let query = RelaySubscriptionQuery {
+ id,
+ targets: target_queries,
+ selector: request.selector().clone(),
+ connect_timeout: timeout
+ .min(Duration::from_millis(self.config().connect_timeout_ms())),
+ timeout: remaining,
+ };
+ let session = match tokio::time::timeout(
+ timeout,
+ self.subscription_client.subscribe(query),
+ )
+ .await
+ {
+ Ok(Ok(session)) => session,
+ Ok(Err(())) | Err(_) => {
+ let observed_at = unix_time_ms().max(1);
+ for relay in &active_relays {
+ self.status.record_read(relay, false, true, observed_at);
+ }
+ return Err(radroots_transport::Error::SubscriptionUnavailable);
+ }
+ };
+
+ Ok(Box::new(RelayEventSubscription {
+ request,
+ session: Some(session),
+ targets,
+ active_relays,
+ cursors,
+ checkpoints,
+ event_count: 0,
+ terminal: None,
+ cancellation_requested: Arc::new(AtomicBool::new(false)),
+ status: Arc::clone(&self.status),
+ }) as BoxSubscription)
+ })
+ }
+}
+
+struct CancellationOnDrop {
+ requested: Arc<AtomicBool>,
+ completed: AtomicBool,
+}
+
+impl CancellationOnDrop {
+ fn new(requested: Arc<AtomicBool>) -> Arc<Self> {
+ Arc::new(Self {
+ requested,
+ completed: AtomicBool::new(false),
+ })
+ }
+
+ fn complete(&self) {
+ self.completed.store(true, Ordering::SeqCst);
+ }
+}
+
+impl Drop for CancellationOnDrop {
+ fn drop(&mut self) {
+ if !self.completed.load(Ordering::SeqCst) {
+ self.requested.store(true, Ordering::SeqCst);
+ }
+ }
+}
+
+fn subscription_filter(
+ selector: &radroots_transport::source::FetchSelector,
+ since: Option<u64>,
+) -> Result<Filter, ()> {
+ let kinds = selector
+ .kinds()
+ .iter()
+ .filter_map(|kind| u16::try_from(*kind).ok())
+ .map(Kind::from)
+ .collect::<Vec<_>>();
+ if !selector.kinds().is_empty() && kinds.is_empty() {
+ return Err(());
+ }
+ let authors = selector
+ .authors()
+ .iter()
+ .map(|author| radroots_nostr::key::public_key_to_nostr(*author).map_err(|_| ()))
+ .collect::<Result<Vec<_>, _>>()?;
+ let mut filter = Filter::new();
+ if !kinds.is_empty() {
+ filter = filter.kinds(kinds);
+ }
+ if !authors.is_empty() {
+ filter = filter.authors(authors);
+ }
+ if let Some(since) = since {
+ filter = filter.since(Timestamp::from_secs(since));
+ }
+ if let Some(until) = selector.until_unix_seconds() {
+ filter = filter.until(Timestamp::from_secs(until));
+ }
+ Ok(filter)
+}
+
+fn selector_is_representable(selector: &radroots_transport::source::FetchSelector) -> bool {
+ selector.kinds().is_empty()
+ || selector
+ .kinds()
+ .iter()
+ .any(|kind| u16::try_from(*kind).is_ok())
+}
+
+fn subscription_id(request: &SubscriptionRequest, sequence: u64) -> String {
+ let mut hasher = Sha256::new();
+ hasher.update(SUBSCRIPTION_ID_DOMAIN);
+ hasher.update(request.request_id().as_str().as_bytes());
+ hasher.update([0]);
+ hasher.update(sequence.to_be_bytes());
+ hex_encode(&hasher.finalize())
+}
+
+fn encode_cursor(
+ request: &SubscriptionRequest,
+ target: &radroots_transport::target::TargetFingerprint,
+ cursor: &RelayCursor,
+) -> Result<FetchCursor, radroots_transport::Error> {
+ FetchCursor::parse(format!(
+ "{CURSOR_PREFIX}:{}:{}:{}",
+ cursor.created_at_unix_s(),
+ cursor.event_id(),
+ cursor_scope(request, target),
+ ))
+}
+
+fn parse_cursor(
+ request: &SubscriptionRequest,
+ checkpoint: &SubscriptionCheckpoint,
+) -> Result<RelayCursor, radroots_transport::Error> {
+ let mut parts = checkpoint.cursor().as_str().split(':');
+ let valid = parts.next() == Some(CURSOR_PREFIX);
+ let created_at = parts.next().and_then(|value| value.parse::<u64>().ok());
+ let event_id = parts.next();
+ let scope = parts.next();
+ if !valid || parts.next().is_some() {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ }
+ let (Some(created_at), Some(event_id), Some(scope)) = (created_at, event_id, scope) else {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ };
+ if scope != cursor_scope(request, checkpoint.target()) {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ }
+ RelayCursor::new(created_at, event_id)
+ .map_err(|_| radroots_transport::Error::InvalidFetchCursor)
+}
+
+fn cursor_scope(
+ request: &SubscriptionRequest,
+ target: &radroots_transport::target::TargetFingerprint,
+) -> String {
+ let mut hasher = Sha256::new();
+ hasher.update(CURSOR_SCOPE_DOMAIN);
+ hasher.update(target.as_str().as_bytes());
+ hasher.update([0]);
+ for kind in request.selector().kinds() {
+ hasher.update(kind.to_be_bytes());
+ }
+ hasher.update([0]);
+ for author in request.selector().authors() {
+ hasher.update(author.as_bytes());
+ }
+ hasher.update([0]);
+ hash_optional_u64(&mut hasher, request.selector().since_unix_seconds());
+ hash_optional_u64(&mut hasher, request.selector().until_unix_seconds());
+ hex_encode(&hasher.finalize())
+}
+
+fn hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) {
+ match value {
+ Some(value) => {
+ hasher.update([1]);
+ hasher.update(value.to_be_bytes());
+ }
+ None => hasher.update([0]),
+ }
+}
+
+fn remaining_duration(deadline_unix_ms: u64) -> Duration {
+ Duration::from_millis(deadline_unix_ms.saturating_sub(unix_time_ms()))
+}
+
+fn hex_encode(bytes: &[u8]) -> String {
+ const HEX: &[u8; 16] = b"0123456789abcdef";
+ let mut encoded = String::with_capacity(bytes.len() * 2);
+ for byte in bytes {
+ encoded.push(HEX[(byte >> 4) as usize] as char);
+ encoded.push(HEX[(byte & 0x0f) as usize] as char);
+ }
+ encoded
+}
+
+#[cfg_attr(coverage_nightly, coverage(off))]
+fn unix_time_ms() -> u64 {
+ SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
+ .unwrap_or_default()
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::{
+ Config, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind, RelayUrlPolicy,
+ };
+ use nostr_sdk::prelude::{EventBuilder, Keys};
+ use radroots_transport::{
+ EventSubscriber, Target, TargetSet,
+ source::{FetchSelector, SubscriptionBounds, SubscriptionCheckpoint},
+ };
+ use std::{
+ collections::VecDeque,
+ sync::{
+ Arc, Mutex,
+ atomic::{AtomicBool, AtomicUsize, Ordering as AtomicOrdering},
+ },
+ };
+
+ const FIXTURE_SECRET_KEY: &str =
+ "0000000000000000000000000000000000000000000000000000000000000001";
+
+ #[derive(Clone, Debug)]
+ enum ScriptedItem {
+ Item(RelaySubscriptionItem),
+ Error,
+ }
+
+ #[derive(Debug)]
+ struct MockSubscriptionSession {
+ items: VecDeque<ScriptedItem>,
+ cancellations: Arc<AtomicUsize>,
+ }
+
+ impl RelaySubscriptionSession for MockSubscriptionSession {
+ fn next(&mut self) -> BoxFuture<'_, Result<RelaySubscriptionItem, ()>> {
+ match self.items.pop_front() {
+ Some(ScriptedItem::Item(item)) => Box::pin(async move { Ok(item) }),
+ Some(ScriptedItem::Error) => Box::pin(async { Err(()) }),
+ None => Box::pin(core::future::pending()),
+ }
+ }
+
+ fn cancel(&mut self) -> BoxFuture<'_, Result<(), ()>> {
+ self.cancellations.fetch_add(1, AtomicOrdering::SeqCst);
+ Box::pin(async { Ok(()) })
+ }
+ }
+
+ #[derive(Debug)]
+ struct MockSubscriptionClient {
+ queries: Mutex<Vec<RelaySubscriptionQuery>>,
+ items: Mutex<Option<VecDeque<ScriptedItem>>>,
+ cancellations: Arc<AtomicUsize>,
+ fail: AtomicBool,
+ }
+
+ impl MockSubscriptionClient {
+ fn new(items: impl IntoIterator<Item = ScriptedItem>) -> Self {
+ Self {
+ queries: Mutex::new(Vec::new()),
+ items: Mutex::new(Some(items.into_iter().collect())),
+ cancellations: Arc::new(AtomicUsize::new(0)),
+ fail: AtomicBool::new(false),
+ }
+ }
+
+ fn failing() -> Self {
+ let client = Self::new([]);
+ client.fail.store(true, AtomicOrdering::SeqCst);
+ client
+ }
+
+ fn query(&self) -> RelaySubscriptionQuery {
+ self.queries.lock().expect("queries")[0].clone()
+ }
+ }
+
+ impl RelaySubscriptionClient for MockSubscriptionClient {
+ fn subscribe(
+ &self,
+ query: RelaySubscriptionQuery,
+ ) -> BoxFuture<'_, Result<Box<dyn RelaySubscriptionSession>, ()>> {
+ self.queries.lock().expect("queries").push(query);
+ if self.fail.load(AtomicOrdering::SeqCst) {
+ return Box::pin(async { Err(()) });
+ }
+ let items = self.items.lock().expect("items").take().unwrap_or_default();
+ let session = MockSubscriptionSession {
+ items,
+ cancellations: Arc::clone(&self.cancellations),
+ };
+ Box::pin(async move { Ok(Box::new(session) as Box<dyn RelaySubscriptionSession>) })
+ }
+ }
+
+ fn configured(relays: &[&str]) -> Config {
+ let endpoints = relays
+ .iter()
+ .map(|relay| {
+ RelayEndpoint::new(relay, RelayUrlPolicy::Public, RelayAccess::ReadWrite)
+ .expect("endpoint")
+ })
+ .collect::<Vec<_>>();
+ Config::from_profile(
+ RelayProfile::explicit(RelayProfileKind::Public, endpoints).expect("profile"),
+ )
+ }
+
+ fn target_set(relays: &[&str]) -> TargetSet {
+ TargetSet::new(
+ relays
+ .iter()
+ .map(|relay| Target::nostr_relay(relay).expect("target"))
+ .collect::<Vec<_>>(),
+ )
+ .expect("target set")
+ }
+
+ fn request(relays: &[&str], event_limit: u16) -> SubscriptionRequest {
+ SubscriptionRequest::new(
+ "nostr-live",
+ target_set(relays),
+ SubscriptionBounds::new(event_limit, unix_time_ms() + 60_000).expect("bounds"),
+ )
+ .expect("request")
+ }
+
+ fn signed_event(content: &str, created_at: u64) -> String {
+ EventBuilder::text_note(content)
+ .custom_created_at(Timestamp::from_secs(created_at))
+ .sign_with_keys(&Keys::parse(FIXTURE_SECRET_KEY).expect("fixture keys"))
+ .expect("signed event")
+ .as_json()
+ }
+
+ fn event_id(raw: &str) -> String {
+ radroots_event_codec::decode::signed_event(raw)
+ .expect("signed event")
+ .id_str()
+ .to_owned()
+ }
+
+ #[tokio::test]
+ async fn subscription_translates_targets_selector_and_unique_ids() {
+ let relay = "wss://one.example";
+ let client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
+ RelaySubscriptionItem::Shutdown,
+ )]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let selector = FetchSelector::all()
+ .with_kinds(vec![1])
+ .expect("kind")
+ .with_since_unix_seconds(1_700_000_000)
+ .expect("since")
+ .with_until_unix_seconds(1_800_000_000)
+ .expect("until");
+ let first = transport
+ .subscribe(request(&[relay], 2).with_selector(selector.clone()))
+ .await
+ .expect("first subscription");
+ drop(first);
+ let second = transport
+ .subscribe(request(&[relay], 2).with_selector(selector.clone()))
+ .await
+ .expect("second subscription");
+ drop(second);
+
+ let queries = client.queries.lock().expect("queries");
+ assert_eq!(queries.len(), 2);
+ assert_ne!(queries[0].id, queries[1].id);
+ assert_eq!(queries[0].targets.len(), 1);
+ assert_eq!(queries[0].targets[0].relay.as_str(), relay);
+ assert_eq!(
+ queries[0].targets[0].since_unix_seconds,
+ Some(1_700_000_000)
+ );
+ assert_eq!(queries[0].selector, selector);
+ assert!(queries[0].connect_timeout <= queries[0].timeout);
+
+ let filter = subscription_filter(&queries[0].selector, Some(1_700_000_000))
+ .expect("filter")
+ .as_json();
+ assert!(filter.contains("\"kinds\":[1]"));
+ assert!(filter.contains("\"since\":1700000000"));
+ assert!(filter.contains("\"until\":1800000000"));
+ }
+
+ #[tokio::test]
+ async fn reconnect_checkpoint_is_equal_timestamp_safe_and_request_bound() {
+ let relay = "wss://one.example";
+ let mut events = [
+ signed_event("equal-a", 1_800_000_000),
+ signed_event("equal-b", 1_800_000_000),
+ signed_event("equal-c", 1_800_000_000),
+ ];
+ events.sort_by_key(|event| event_id(event));
+ let base_request = request(&[relay], 2);
+ let target = base_request.target_set().targets()[0].fingerprint().clone();
+ let middle_cursor = RelayCursor::new(1_800_000_000, event_id(&events[1])).expect("cursor");
+ let checkpoint = SubscriptionCheckpoint::new(
+ target,
+ encode_cursor(
+ &base_request,
+ base_request.target_set().targets()[0].fingerprint(),
+ &middle_cursor,
+ )
+ .expect("opaque cursor"),
+ );
+ let resumed = base_request
+ .with_checkpoints([checkpoint])
+ .expect("checkpointed request");
+ let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay");
+ let client = Arc::new(MockSubscriptionClient::new([
+ ScriptedItem::Item(RelaySubscriptionItem::Event {
+ relay: relay_url.clone(),
+ raw: events[0].clone(),
+ }),
+ ScriptedItem::Item(RelaySubscriptionItem::Event {
+ relay: relay_url,
+ raw: events[2].clone(),
+ }),
+ ]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let mut subscription = transport.subscribe(resumed).await.expect("subscription");
+ let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") else {
+ panic!("event expected");
+ };
+ assert_eq!(event.observed().event().id_str(), event_id(&events[2]));
+ assert_eq!(
+ event.checkpoint().cursor().as_str().split(':').nth(1),
+ Some("1800000000")
+ );
+ assert_eq!(
+ client.query().targets[0].since_unix_seconds,
+ Some(1_800_000_000)
+ );
+ }
+
+ #[tokio::test]
+ async fn malformed_or_mismatched_checkpoint_fails_before_backend_work() {
+ let relay = "wss://one.example";
+ let client = Arc::new(MockSubscriptionClient::new([]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let base = request(&[relay], 1);
+ let target = base.target_set().targets()[0].fingerprint().clone();
+ let malformed = base
+ .clone()
+ .with_checkpoints([SubscriptionCheckpoint::new(
+ target.clone(),
+ FetchCursor::parse("not-a-live-cursor").expect("opaque cursor"),
+ )])
+ .expect("request checkpoint");
+ assert_eq!(
+ transport.subscribe(malformed).await.err(),
+ Some(radroots_transport::Error::InvalidFetchCursor)
+ );
+
+ let other_selector = FetchSelector::all().with_kinds(vec![2]).expect("selector");
+ let scoped = encode_cursor(
+ &base,
+ &target,
+ &RelayCursor::new(10, "a".repeat(64)).expect("cursor"),
+ )
+ .expect("scoped cursor");
+ let mismatched = base
+ .with_selector(other_selector)
+ .with_checkpoints([SubscriptionCheckpoint::new(target, scoped)])
+ .expect("request checkpoint");
+ assert_eq!(
+ transport.subscribe(mismatched).await.err(),
+ Some(radroots_transport::Error::InvalidFetchCursor)
+ );
+ assert!(client.queries.lock().expect("queries").is_empty());
+ }
+
+ #[test]
+ fn dropping_an_unpolled_subscription_start_performs_no_backend_work() {
+ let relay = "wss://one.example";
+ let client = Arc::new(MockSubscriptionClient::new([]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let future = transport.subscribe(request(&[relay], 1));
+ drop(future);
+ assert!(client.queries.lock().expect("queries").is_empty());
+ }
+
+ #[tokio::test]
+ async fn event_limit_and_explicit_cancellation_are_stable_and_idempotent() {
+ let relay = "wss://one.example";
+ let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay");
+ let client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
+ RelaySubscriptionItem::Event {
+ relay: relay_url,
+ raw: signed_event("bounded", 1_800_000_001),
+ },
+ )]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let mut subscription = transport
+ .subscribe(request(&[relay], 1))
+ .await
+ .expect("subscription");
+ assert!(matches!(
+ subscription.next().await.expect("event"),
+ SubscriptionNext::Event(_)
+ ));
+ let SubscriptionNext::End(limit) = subscription.next().await.expect("limit") else {
+ panic!("terminal expected");
+ };
+ assert_eq!(limit.reason(), SubscriptionEndReason::EventLimit);
+ assert_eq!(limit.event_count(), 1);
+ assert_eq!(limit.checkpoints().len(), 1);
+ assert_eq!(subscription.cancel().await.expect("stable cancel"), limit);
+ assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 1);
+
+ let cancel_client = Arc::new(MockSubscriptionClient::new([]));
+ let cancel_transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), cancel_client.clone());
+ let mut cancelled = cancel_transport
+ .subscribe(request(&[relay], 2))
+ .await
+ .expect("subscription");
+ let terminal = cancelled.cancel().await.expect("cancel");
+ assert_eq!(terminal.reason(), SubscriptionEndReason::Cancelled);
+ assert_eq!(cancelled.cancel().await.expect("repeat cancel"), terminal);
+ assert_eq!(cancel_client.cancellations.load(AtomicOrdering::SeqCst), 1);
+ }
+
+ #[tokio::test]
+ async fn dropping_pending_next_requests_cancellation_on_the_next_call() {
+ let relay = "wss://one.example";
+ let client = Arc::new(MockSubscriptionClient::new([]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let mut subscription = transport
+ .subscribe(request(&[relay], 2))
+ .await
+ .expect("subscription");
+ let mut pending = subscription.next();
+ assert!(futures::poll!(pending.as_mut()).is_pending());
+ drop(pending);
+ let SubscriptionNext::End(terminal) = subscription.next().await.expect("cancelled") else {
+ panic!("terminal expected");
+ };
+ assert_eq!(terminal.reason(), SubscriptionEndReason::Cancelled);
+ assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 1);
+ }
+
+ #[tokio::test]
+ async fn expired_unrepresentable_closed_and_failed_sources_are_bounded() {
+ let relay = "wss://one.example";
+ let client = Arc::new(MockSubscriptionClient::new([]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), client.clone());
+ let expired = SubscriptionRequest::new(
+ "expired",
+ target_set(&[relay]),
+ SubscriptionBounds::new(1, 1).expect("bounds"),
+ )
+ .expect("request");
+ let mut expired = transport
+ .subscribe(expired)
+ .await
+ .expect("expired capability");
+ let SubscriptionNext::End(terminal) = expired.next().await.expect("deadline") else {
+ panic!("terminal expected");
+ };
+ assert_eq!(terminal.reason(), SubscriptionEndReason::Deadline);
+
+ let unsupported = request(&[relay], 1).with_selector(
+ FetchSelector::all()
+ .with_kinds(vec![u32::MAX])
+ .expect("selector"),
+ );
+ let mut unsupported = transport
+ .subscribe(unsupported)
+ .await
+ .expect("closed capability");
+ let SubscriptionNext::End(terminal) = unsupported.next().await.expect("source closed")
+ else {
+ panic!("terminal expected");
+ };
+ assert_eq!(terminal.reason(), SubscriptionEndReason::SourceClosed);
+ assert!(client.queries.lock().expect("queries").is_empty());
+
+ let deadline_request = SubscriptionRequest::new(
+ "pending-deadline",
+ target_set(&[relay]),
+ SubscriptionBounds::new(1, unix_time_ms() + 100).expect("bounds"),
+ )
+ .expect("request");
+ let mut deadline = transport
+ .subscribe(deadline_request)
+ .await
+ .expect("deadline capability");
+ let SubscriptionNext::End(terminal) = deadline.next().await.expect("deadline") else {
+ panic!("terminal expected");
+ };
+ assert_eq!(terminal.reason(), SubscriptionEndReason::Deadline);
+ assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 0);
+
+ let failed = Arc::new(MockSubscriptionClient::failing());
+ let failed_transport =
+ NostrTransport::with_subscription_client(configured(&[relay]), failed.clone());
+ assert_eq!(
+ failed_transport.subscribe(request(&[relay], 1)).await.err(),
+ Some(radroots_transport::Error::SubscriptionUnavailable)
+ );
+ assert_eq!(failed.queries.lock().expect("queries").len(), 1);
+ }
+
+ #[tokio::test]
+ async fn source_closure_backend_failure_and_event_admission_are_explicit() {
+ let relays = ["wss://one.example", "wss://two.example"];
+ let one = RelayUrl::parse(relays[0], RelayUrlPolicy::Public).expect("one");
+ let two = RelayUrl::parse(relays[1], RelayUrlPolicy::Public).expect("two");
+ let client = Arc::new(MockSubscriptionClient::new([
+ ScriptedItem::Item(RelaySubscriptionItem::Closed { relay: one }),
+ ScriptedItem::Item(RelaySubscriptionItem::Closed { relay: two }),
+ ]));
+ let transport =
+ NostrTransport::with_subscription_client(configured(&relays), client.clone());
+ let mut subscription = transport
+ .subscribe(request(&relays, 2))
+ .await
+ .expect("subscription");
+ let SubscriptionNext::End(terminal) = subscription.next().await.expect("closed") else {
+ panic!("terminal expected");
+ };
+ assert_eq!(terminal.reason(), SubscriptionEndReason::SourceClosed);
+ assert_eq!(client.cancellations.load(AtomicOrdering::SeqCst), 0);
+
+ let error_client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Error]));
+ let error_transport =
+ NostrTransport::with_subscription_client(configured(&[relays[0]]), error_client);
+ let mut subscription = error_transport
+ .subscribe(request(&[relays[0]], 1))
+ .await
+ .expect("subscription");
+ assert_eq!(
+ subscription.next().await,
+ Err(radroots_transport::Error::SubscriptionUnavailable)
+ );
+
+ let mismatch_client = Arc::new(MockSubscriptionClient::new([ScriptedItem::Item(
+ RelaySubscriptionItem::Event {
+ relay: RelayUrl::parse(relays[0], RelayUrlPolicy::Public).expect("relay"),
+ raw: signed_event("wrong-kind", 1_800_000_002),
+ },
+ )]));
+ let mismatch_transport =
+ NostrTransport::with_subscription_client(configured(&[relays[0]]), mismatch_client);
+ let selector = FetchSelector::all().with_kinds(vec![2]).expect("selector");
+ let mut subscription = mismatch_transport
+ .subscribe(request(&[relays[0]], 1).with_selector(selector))
+ .await
+ .expect("subscription");
+ assert_eq!(
+ subscription.next().await,
+ Err(radroots_transport::Error::UnexpectedSubscriptionEvent)
+ );
+ }
+}
diff --git a/crates/transport_nostr/tests/conformance.rs b/crates/transport_nostr/tests/conformance.rs
@@ -1,5 +1,5 @@
use core::fmt::Debug;
-use radroots_transport::{EventSink, EventSource};
+use radroots_transport::{EventSink, EventSource, EventSubscriber};
use radroots_transport_nostr::{Config, Error, NostrTransport, RelayUrl, RelayUrlPolicy};
const MANIFEST: &str = include_str!("../Cargo.toml");
@@ -9,10 +9,11 @@ const RELAY: &str = include_str!("../src/relay.rs");
const SINK: &str = include_str!("../src/sink.rs");
const SOURCE: &str = include_str!("../src/source.rs");
const STATUS: &str = include_str!("../src/status.rs");
+const SUBSCRIPTION: &str = include_str!("../src/subscription.rs");
fn assert_adapter_contract<T>()
where
- T: EventSource + EventSink + Clone + Debug + Send + Sync,
+ T: EventSource + EventSubscriber + EventSink + Clone + Debug + Send + Sync,
{
}
@@ -74,6 +75,18 @@ fn complete_mocked_conformance_matrix_remains_reachable() {
),
(SOURCE, "dropping_an_unpolled_fetch_performs_no_relay_work"),
(
+ SUBSCRIPTION,
+ "reconnect_checkpoint_is_equal_timestamp_safe_and_request_bound",
+ ),
+ (
+ SUBSCRIPTION,
+ "dropping_pending_next_requests_cancellation_on_the_next_call",
+ ),
+ (
+ SUBSCRIPTION,
+ "dropping_an_unpolled_subscription_start_performs_no_backend_work",
+ ),
+ (
STATUS,
"every_upstream_class_maps_to_stable_secret_safe_output",
),
diff --git a/crates/transport_nostr/tests/local_io.rs b/crates/transport_nostr/tests/local_io.rs
@@ -1,7 +1,9 @@
use futures::{SinkExt, StreamExt};
use radroots_transport::{
- EventSource, FetchRequest, TargetSet, capability::Availability, outcome::FetchTargetState,
- source::FetchBounds,
+ EventSource, EventSubscriber, FetchRequest, SubscriptionNext, SubscriptionRequest, TargetSet,
+ capability::Availability,
+ outcome::FetchTargetState,
+ source::{FetchBounds, SubscriptionBounds, SubscriptionEndReason},
};
use radroots_transport_nostr::{
Config, NostrTransport, RelayAccess, RelayAggregateState, RelayEndpoint, RelayEvidenceState,
@@ -12,6 +14,8 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tokio::net::TcpListener;
use tokio_tungstenite::{accept_async, tungstenite::Message};
+const FIXTURE_SECRET_KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001";
+
#[tokio::test(flavor = "multi_thread")]
async fn simulator_profile_proves_read_capability_against_a_real_loopback_socket() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind relay");
@@ -97,6 +101,104 @@ async fn simulator_profile_proves_read_capability_against_a_real_loopback_socket
server.await.expect("relay task");
}
+#[tokio::test(flavor = "multi_thread")]
+async fn simulator_profile_streams_and_cancels_a_real_live_subscription() {
+ use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Timestamp};
+
+ let event = EventBuilder::text_note("bounded live event")
+ .custom_created_at(Timestamp::from_secs(1_800_000_100))
+ .sign_with_keys(&Keys::parse(FIXTURE_SECRET_KEY).expect("fixture keys"))
+ .expect("signed event");
+ let expected_id = event.id.to_hex();
+ let event: Value = serde_json::from_str(event.as_json().as_str()).expect("event JSON");
+ let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind relay");
+ let address = listener.local_addr().expect("relay address");
+ let relay_url = format!("ws://{address}");
+ let server = tokio::spawn(async move {
+ let (stream, _) = listener.accept().await.expect("accept relay client");
+ let mut websocket = accept_async(stream).await.expect("websocket handshake");
+ let mut subscription_id = None;
+ while let Some(message) = websocket.next().await {
+ let Message::Text(message) = message.expect("client message") else {
+ continue;
+ };
+ let parsed: Value = serde_json::from_str(message.as_str()).expect("Nostr message");
+ let Some(values) = parsed.as_array() else {
+ continue;
+ };
+ match values.as_slice() {
+ [Value::String(kind), Value::String(subscription), ..] if kind == "REQ" => {
+ subscription_id = Some(subscription.clone());
+ websocket
+ .send(Message::Text(
+ serde_json::to_string(&("EVENT", subscription, &event))
+ .expect("EVENT message")
+ .into(),
+ ))
+ .await
+ .expect("send EVENT");
+ }
+ [Value::String(kind), Value::String(subscription)]
+ if kind == "CLOSE" && subscription_id.as_ref() == Some(subscription) =>
+ {
+ return;
+ }
+ _ => {}
+ }
+ }
+ panic!("client disconnected without CLOSE");
+ });
+
+ let endpoint = RelayEndpoint::new(
+ relay_url.as_str(),
+ RelayUrlPolicy::Local,
+ RelayAccess::ReadWrite,
+ )
+ .expect("loopback endpoint");
+ let profile =
+ RelayProfile::explicit(RelayProfileKind::Simulator, [endpoint]).expect("simulator profile");
+ let config = Config::from_profile(profile)
+ .with_timeouts(1_000, 2_000, 500)
+ .expect("timeouts");
+ let targets = TargetSet::new(
+ config
+ .read_relays()
+ .map(|relay| relay.to_target())
+ .collect::<Result<Vec<_>, _>>()
+ .expect("targets"),
+ )
+ .expect("target set");
+ let request = SubscriptionRequest::new(
+ "local-live-io",
+ targets,
+ SubscriptionBounds::new(2, unix_time_ms() + 5_000).expect("bounds"),
+ )
+ .expect("request");
+ let transport = NostrTransport::new(config);
+ let mut subscription =
+ tokio::time::timeout(Duration::from_secs(5), transport.subscribe(request))
+ .await
+ .expect("bounded subscribe")
+ .expect("subscription");
+ let SubscriptionNext::Event(observed) =
+ tokio::time::timeout(Duration::from_secs(5), subscription.next())
+ .await
+ .expect("bounded event")
+ .expect("event")
+ else {
+ panic!("event expected");
+ };
+ assert_eq!(observed.observed().event().id_str(), expected_id);
+ assert_eq!(
+ subscription.cancel().await.expect("cancel").reason(),
+ SubscriptionEndReason::Cancelled
+ );
+ tokio::time::timeout(Duration::from_secs(5), server)
+ .await
+ .expect("server deadline")
+ .expect("relay task");
+}
+
fn unix_time_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs
@@ -9,6 +9,7 @@ const PUBLIC_API: &str =
include_str!("../../../contracts/api_baselines/radroots_transport_nostr.txt");
const ROOT: &str = include_str!("../src/lib.rs");
const PROFILE: &str = include_str!("../src/profile.rs");
+const SUBSCRIPTION: &str = include_str!("../src/subscription.rs");
#[test]
fn manifest_and_root_match_the_governed_transport_boundary() {
@@ -36,7 +37,16 @@ fn manifest_and_root_match_the_governed_transport_boundary() {
assert_eq!(
private_modules(ROOT),
BTreeSet::from([
- "auth", "client", "cursor", "error", "profile", "relay", "sink", "source", "status"
+ "auth",
+ "client",
+ "cursor",
+ "error",
+ "profile",
+ "relay",
+ "sink",
+ "source",
+ "status",
+ "subscription"
])
);
for export in [
@@ -56,7 +66,7 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
"## Configure without connecting",
"## Public surface",
"## Relay and network security",
- "## Fetch, delivery, and outcome behavior",
+ "## Fetch, live subscription, delivery, and outcome behavior",
"## Deadlines, cancellation, and commit points",
"## Serialization and diagnostics",
"## Features and runtime requirements",
@@ -64,6 +74,11 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
"radroots_crates_release_v1.toml",
"examples/configure_transport.rs",
"contracts/api_baselines/radroots_transport_nostr.txt",
+ "Live subscriptions use the same explicit readable targets",
+ "inclusive `since` timestamp",
+ "event-ID tie breaker",
+ "upstream auto-close deadline",
+ "adapter-owned worker",
] {
assert!(README.contains(required), "README is missing `{required}`");
}
@@ -73,8 +88,10 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
"RelayProfile::explicit(",
"NostrTransport::new(config)",
"let source: &dyn EventSource",
+ "let subscriber: &dyn EventSubscriber",
"let sink: &dyn EventSink",
"drop(source.status())",
+ "let _ = subscriber",
"drop(sink.status())",
] {
assert!(
@@ -95,6 +112,7 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
"pub enum radroots_transport_nostr::Error",
"impl radroots_transport::sink::EventSink for radroots_transport_nostr::NostrTransport",
"impl radroots_transport::source::EventSource for radroots_transport_nostr::NostrTransport",
+ "impl radroots_transport::source::EventSubscriber for radroots_transport_nostr::NostrTransport",
"NostrTransport::begin_authentication",
"NostrTransport::complete_authentication",
"NostrTransport::reject_authentication",
@@ -212,6 +230,32 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() {
"sink.rs".to_owned(),
"source.rs".to_owned(),
"status.rs".to_owned(),
+ "subscription.rs".to_owned(),
])
);
+
+ for required in [
+ "impl EventSubscriber for NostrTransport",
+ "SubscribeAutoCloseOptions::default()",
+ "ReqExitPolicy::WaitDurationAfterEOSE(query.timeout)",
+ "cursor.precedes(event.created_at(), event.id_str())",
+ "self.terminate(SubscriptionEndReason::Cancelled)",
+ ] {
+ assert!(
+ SUBSCRIPTION.contains(required),
+ "subscription adapter is missing `{required}`"
+ );
+ }
+ for forbidden in [
+ "tokio::spawn",
+ "spawn_blocking",
+ "std::thread",
+ "process::",
+ "global_default",
+ ] {
+ assert!(
+ !SUBSCRIPTION.contains(forbidden),
+ "subscription adapter contains forbidden authority `{forbidden}`"
+ );
+ }
}