commit 69c3ce2ded69b08e20993620d48391be7b93c29c
parent 7ef2e7dc401aa2019583e480f8e589d18ef3d558
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 06:32:40 +0000
transport-nostr: implement the Nostr event source
- implement EventSource with bounded relay filters, absolute deadlines, and opaque stable cursors
- decode wire events into canonical observations with exact relay provenance and deterministic deduplication
- report malformed events and relay failures explicitly without ingestion, projection, or persistence
- verify pagination, cancellation, package, Clippy, architecture, boundary, and source-maintenance gates
Diffstat:
2 files changed, 445 insertions(+), 1 deletion(-)
diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs
@@ -126,6 +126,7 @@ fn validate_timeout(field: &'static str, value_ms: u64) -> Result<(), Error> {
pub struct NostrTransport {
config: Config,
pub(crate) client: Arc<dyn crate::sink::RelayClient>,
+ pub(crate) source_client: Arc<dyn crate::source::RelaySourceClient>,
}
impl NostrTransport {
@@ -134,6 +135,7 @@ impl NostrTransport {
Self {
config,
client: Arc::new(crate::sink::LiveRelayClient),
+ source_client: Arc::new(crate::source::LiveRelaySourceClient),
}
}
@@ -144,7 +146,23 @@ impl NostrTransport {
#[cfg(test)]
pub(crate) fn with_client(config: Config, client: Arc<dyn crate::sink::RelayClient>) -> Self {
- Self { config, client }
+ Self {
+ config,
+ client,
+ source_client: Arc::new(crate::source::LiveRelaySourceClient),
+ }
+ }
+
+ #[cfg(test)]
+ pub(crate) fn with_source_client(
+ config: Config,
+ source_client: Arc<dyn crate::source::RelaySourceClient>,
+ ) -> Self {
+ Self {
+ config,
+ client: Arc::new(crate::sink::LiveRelayClient),
+ source_client,
+ }
}
}
diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs
@@ -1 +1,427 @@
//! Nostr implementation of the transport event source.
+
+use crate::{NostrTransport, RelayUrl};
+use core::cmp::Ordering;
+use core::time::Duration;
+use nostr_sdk::prelude::{Filter, JsonUtil, Timestamp};
+use radroots_transport::{
+ BoxFuture, EventSource, FetchPage, FetchRequest,
+ capability::{Availability, Maturity, SourceCapabilities},
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ source::{EventProvenance, FetchCursor, NextPage, ObservedEvent, SourceStatus},
+};
+use std::collections::{BTreeMap, BTreeSet};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+const UPSTREAM_FETCH_LIMIT: usize = 1_000;
+const CURSOR_PREFIX: &str = "nostr-v1";
+
+#[derive(Clone, Debug)]
+pub(crate) struct SourceQuery {
+ relays: Vec<RelayUrl>,
+ until_unix_seconds: Option<u64>,
+ timeout: Duration,
+}
+
+#[derive(Clone, Debug)]
+pub(crate) struct RelayFetchBatch {
+ relay: RelayUrl,
+ result: Result<Vec<String>, String>,
+}
+
+pub(crate) trait RelaySourceClient: Send + Sync {
+ fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>>;
+}
+
+#[derive(Debug)]
+pub(crate) struct LiveRelaySourceClient;
+
+impl RelaySourceClient for LiveRelaySourceClient {
+ fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
+ Box::pin(async move {
+ let client = nostr_sdk::Client::default();
+ let mut batches = Vec::with_capacity(query.relays.len());
+ for relay in query.relays {
+ let url = relay.as_str().to_owned();
+ let result = async {
+ client.add_relay(url.as_str()).await?;
+ client
+ .try_connect_relay(url.as_str(), query.timeout)
+ .await?;
+ let mut filter = Filter::new().limit(UPSTREAM_FETCH_LIMIT);
+ if let Some(until) = query.until_unix_seconds {
+ filter = filter.until(Timestamp::from_secs(until));
+ }
+ client
+ .fetch_events_from([url.as_str()], filter, query.timeout)
+ .await
+ .map(|events| events.iter().map(JsonUtil::as_json).collect())
+ }
+ .await
+ .map_err(|error: nostr_sdk::client::Error| error.to_string());
+ batches.push(RelayFetchBatch { relay, result });
+ }
+ batches
+ })
+ }
+}
+
+#[derive(Debug)]
+struct Candidate {
+ relay: RelayUrl,
+ raw: String,
+ created_at: u64,
+ event_id: String,
+}
+
+impl EventSource for NostrTransport {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
+ let configured = !self.config().relays().is_empty();
+ Box::pin(async move {
+ Ok(SourceStatus::new(
+ radroots_transport::TransportId::NOSTR,
+ configured,
+ Maturity::Preview,
+ Availability::Available,
+ SourceCapabilities::FETCH,
+ "bounded Nostr event source configured",
+ ))
+ })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
+ Box::pin(async move {
+ let cursor = request.cursor().map(parse_cursor).transpose()?.flatten();
+ let now_ms = unix_time_ms();
+ let remaining_ms = request.bounds().deadline_unix_ms().saturating_sub(now_ms);
+ let timeout_ms = remaining_ms.min(self.config().request_timeout_ms());
+
+ let mut targets = BTreeMap::new();
+ let mut outcomes = Vec::new();
+ for target in request.target_set().targets() {
+ match RelayUrl::from_target(target, self.config().relay_url_policy()) {
+ Ok(relay) if self.config().relays().contains(&relay) => {
+ targets.insert(relay, target.clone());
+ }
+ _ => outcomes.push(
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::FailedTerminal,
+ )
+ .with_message("target is not configured for this source"),
+ ),
+ }
+ }
+
+ if timeout_ms == 0 {
+ outcomes.extend(targets.values().map(|target| {
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::FailedRetryable,
+ )
+ .with_message("fetch deadline elapsed before relay access")
+ }));
+ return FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete);
+ }
+
+ let batches = self
+ .source_client
+ .fetch(SourceQuery {
+ relays: targets.keys().cloned().collect(),
+ until_unix_seconds: cursor.as_ref().map(|cursor| cursor.created_at),
+ timeout: Duration::from_millis(timeout_ms),
+ })
+ .await;
+ let mut candidates = Vec::new();
+ let mut malformed_by_relay = BTreeMap::<RelayUrl, usize>::new();
+ let mut reported = BTreeSet::new();
+ for batch in batches {
+ let Some(target) = targets.get(&batch.relay) else {
+ return Err(radroots_transport::Error::UnexpectedFetchTargetOutcome);
+ };
+ if !reported.insert(batch.relay.clone()) {
+ return Err(radroots_transport::Error::DuplicateFetchTargetOutcome);
+ }
+ match batch.result {
+ Ok(raw_events) => {
+ for raw in raw_events {
+ match radroots_event_codec::decode::signed_event(raw.as_str()) {
+ Ok(event) => candidates.push(Candidate {
+ relay: batch.relay.clone(),
+ created_at: event.created_at(),
+ event_id: event.id_str().to_owned(),
+ raw,
+ }),
+ Err(_) => {
+ *malformed_by_relay.entry(batch.relay.clone()).or_default() +=
+ 1;
+ }
+ }
+ }
+ let malformed = malformed_by_relay
+ .get(&batch.relay)
+ .copied()
+ .unwrap_or_default();
+ let outcome = if malformed == 0 {
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::Complete,
+ )
+ } else {
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::Partial,
+ )
+ .with_message(format!("ignored {malformed} malformed relay event(s)"))
+ };
+ outcomes.push(outcome);
+ }
+ Err(message) => outcomes.push(
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::FailedRetryable,
+ )
+ .with_message(safe_message(message)),
+ ),
+ }
+ }
+ for (relay, target) in &targets {
+ if !reported.contains(relay) {
+ outcomes.push(
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::FailedRetryable,
+ )
+ .with_message("relay returned no fetch result"),
+ );
+ }
+ }
+
+ candidates.sort_by(compare_candidate);
+ if let Some(cursor) = &cursor {
+ candidates.retain(|candidate| candidate_is_after_cursor(candidate, cursor));
+ }
+ let mut seen = BTreeSet::new();
+ candidates.retain(|candidate| seen.insert(candidate.event_id.clone()));
+
+ let has_more = candidates.len() > usize::from(request.bounds().limit());
+ candidates.truncate(usize::from(request.bounds().limit()));
+ let next_page = if has_more {
+ let last = candidates.last().expect("non-empty bounded page");
+ NextPage::Cursor(FetchCursor::parse(format!(
+ "{CURSOR_PREFIX}:{}:{}",
+ last.created_at, last.event_id
+ ))?)
+ } else {
+ NextPage::Complete
+ };
+ let observed_at = unix_time_ms().max(1);
+ let mut events = Vec::with_capacity(candidates.len());
+ for candidate in candidates {
+ let target = targets
+ .get(&candidate.relay)
+ .expect("candidate relay has requested target");
+ let mut provenance = EventProvenance::new(
+ radroots_transport::TransportId::NOSTR,
+ target.fingerprint().clone(),
+ observed_at,
+ )?;
+ if let Some(request_cursor) = request.cursor().cloned() {
+ provenance = provenance.with_cursor(request_cursor);
+ }
+ let event = radroots_event_codec::decode::signed_event(candidate.raw.as_str())
+ .map_err(|_| radroots_transport::Error::UnexpectedFetchProvenance)?;
+ events.push(ObservedEvent::new(event, provenance));
+ }
+ FetchPage::for_request(&request, events, outcomes, next_page)
+ })
+ }
+}
+
+#[derive(Clone, Debug)]
+struct CursorPosition {
+ created_at: u64,
+ event_id: String,
+}
+
+fn parse_cursor(cursor: &FetchCursor) -> Result<Option<CursorPosition>, radroots_transport::Error> {
+ let mut parts = 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();
+ if !valid || parts.next().is_some() {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ }
+ let (Some(created_at), Some(event_id)) = (created_at, event_id) else {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ };
+ if event_id.len() != 64 || !event_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
+ return Err(radroots_transport::Error::InvalidFetchCursor);
+ }
+ Ok(Some(CursorPosition {
+ created_at,
+ event_id: event_id.to_ascii_lowercase(),
+ }))
+}
+
+fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering {
+ right
+ .created_at
+ .cmp(&left.created_at)
+ .then_with(|| right.event_id.cmp(&left.event_id))
+ .then_with(|| left.relay.cmp(&right.relay))
+}
+
+fn candidate_is_after_cursor(candidate: &Candidate, cursor: &CursorPosition) -> bool {
+ candidate.created_at < cursor.created_at
+ || candidate.created_at == cursor.created_at && candidate.event_id < cursor.event_id
+}
+
+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()
+}
+
+fn safe_message(message: String) -> String {
+ let message = message
+ .chars()
+ .filter(|character| !character.is_control())
+ .take(1_024)
+ .collect::<String>();
+ if message.trim().is_empty() {
+ "relay fetch failed".to_owned()
+ } else {
+ message.trim().to_owned()
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::{Config, RelayUrlPolicy};
+ use radroots_transport::{
+ FetchRequest, Target, TargetSet,
+ source::{FetchBounds, NextPage},
+ };
+ use std::sync::{
+ Arc,
+ atomic::{AtomicUsize, Ordering as AtomicOrdering},
+ };
+
+ const FIRST: &str = r#"{"id":"762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1800000100,"kind":0,"tags":[],"content":"{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}","sig":"4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"}"#;
+ const SECOND: &str = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
+
+ #[derive(Debug)]
+ struct MockSourceClient;
+
+ impl RelaySourceClient for MockSourceClient {
+ fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
+ Box::pin(async move {
+ query
+ .relays
+ .into_iter()
+ .map(|relay| RelayFetchBatch {
+ relay,
+ result: Ok(vec![FIRST.to_owned(), SECOND.to_owned(), "{".to_owned()]),
+ })
+ .collect()
+ })
+ }
+ }
+
+ fn transport() -> NostrTransport {
+ let config = Config::new(
+ RelayUrlPolicy::Public,
+ ["wss://one.example", "wss://two.example"],
+ )
+ .expect("config");
+ NostrTransport::with_source_client(config, Arc::new(MockSourceClient))
+ }
+
+ fn request(limit: u16) -> FetchRequest {
+ FetchRequest::new(
+ "nostr-fetch",
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("one"),
+ Target::nostr_relay("wss://two.example").expect("two"),
+ ])
+ .expect("targets"),
+ FetchBounds::new(limit, u64::MAX).expect("bounds"),
+ )
+ .expect("request")
+ }
+
+ #[test]
+ fn source_deduplicates_relays_reports_malformed_and_paginates() {
+ let transport = transport();
+ let first_request = request(1);
+ let first =
+ futures::executor::block_on(transport.fetch(first_request)).expect("first page");
+ assert_eq!(first.events().len(), 1);
+ assert!(
+ first
+ .target_outcomes()
+ .iter()
+ .all(|outcome| outcome.state() == FetchTargetState::Partial)
+ );
+ let NextPage::Cursor(cursor) = first.next_page() else {
+ panic!("cursor expected");
+ };
+
+ let second_request = request(2).with_cursor(cursor.clone());
+ let second =
+ futures::executor::block_on(transport.fetch(second_request)).expect("second page");
+ assert_eq!(second.events().len(), 1);
+ assert!(matches!(second.next_page(), NextPage::Complete));
+ assert_ne!(
+ first.events()[0].event().id(),
+ second.events()[0].event().id()
+ );
+ }
+
+ #[test]
+ fn malformed_cursor_fails_before_relay_access() {
+ let request = request(1).with_cursor(FetchCursor::parse("other:1:value").expect("opaque"));
+ let error = futures::executor::block_on(transport().fetch(request)).expect_err("cursor");
+ assert_eq!(error, radroots_transport::Error::InvalidFetchCursor);
+ }
+
+ #[test]
+ fn dropping_an_unpolled_fetch_performs_no_relay_work() {
+ #[derive(Debug)]
+ struct CountingSourceClient(Arc<AtomicUsize>);
+
+ impl RelaySourceClient for CountingSourceClient {
+ fn fetch<'a>(&'a self, _query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
+ self.0.fetch_add(1, AtomicOrdering::SeqCst);
+ Box::pin(async { Vec::new() })
+ }
+ }
+
+ let calls = Arc::new(AtomicUsize::new(0));
+ let config = Config::new(RelayUrlPolicy::Public, ["wss://one.example"]).expect("config");
+ let transport = NostrTransport::with_source_client(
+ config,
+ Arc::new(CountingSourceClient(Arc::clone(&calls))),
+ );
+ let target_set = TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("target"),
+ ])
+ .expect("targets");
+ let fetch = transport.fetch(
+ FetchRequest::new(
+ "cancel-before-poll",
+ target_set,
+ FetchBounds::new(1, u64::MAX).expect("bounds"),
+ )
+ .expect("request"),
+ );
+ drop(fetch);
+ assert_eq!(calls.load(AtomicOrdering::SeqCst), 0);
+ }
+}