lib

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

commit d287d41c2cd97cd0e455445da90f22180029f089
parent 066cbd4f741deb21674859d16b127b19820260eb
Author: triesap <tyson@radroots.org>
Date:   Tue, 25 Aug 2026 03:57:20 +0000

transport-nostr: preserve same-second live events

Diffstat:
Mcrates/transport_nostr/README.md | 10++++++----
Mcrates/transport_nostr/src/cursor.rs | 11++++++-----
Mcrates/transport_nostr/src/subscription.rs | 102++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mcrates/transport_nostr/tests/package_boundary.rs | 8++++++--
4 files changed, 109 insertions(+), 22 deletions(-)

diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md @@ -144,10 +144,12 @@ rewritten as completion even when it carries admissible events. 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. +the adapter reconnects with Nostr's inclusive `since` timestamp, suppresses +older timestamps and the exact checkpoint event, and permits at-least-once +replay of other events from the checkpoint second. Within one subscription it +deduplicates exact event IDs, accepts same-second events in relay arrival order, +and never regresses the canonical target checkpoint. Each emitted event carries +exact relay provenance and that current checkpoint. Event limits, absolute deadlines, explicit cancellation, source closure, and stable repeated terminal results follow the generic subscription contract. diff --git a/crates/transport_nostr/src/cursor.rs b/crates/transport_nostr/src/cursor.rs @@ -5,8 +5,8 @@ use crate::Error; /// Stable total-order position for one Nostr event. /// /// Relay timestamps are only second-granular. The canonical lowercase event id -/// is therefore a required tie-breaker for both descending fetch pages and -/// inclusive reconnect catch-up after a subscription interruption. +/// provides a deterministic tie-breaker for descending fetch pages and a +/// non-regressing checkpoint position for inclusive subscription reconnects. #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] pub struct RelayCursor { created_at_unix_s: u64, @@ -42,9 +42,10 @@ impl RelayCursor { self.event_id.as_str() } - /// Returns whether a candidate follows this cursor in ascending reconnect - /// order. Equal timestamps are resolved by event id, preventing loss when - /// a reconnect query uses an inclusive `since` timestamp. + /// Returns whether a candidate follows this cursor in ascending total + /// order. Equal timestamps are resolved by event id. Live subscriptions + /// additionally admit out-of-order peers from the checkpoint second so + /// this total-order helper is not itself used as a lossless admission gate. #[must_use] pub fn precedes(&self, created_at_unix_s: u64, event_id: &str) -> bool { created_at_unix_s > self.created_at_unix_s diff --git a/crates/transport_nostr/src/subscription.rs b/crates/transport_nostr/src/subscription.rs @@ -222,8 +222,10 @@ struct RelayEventSubscription { session: Option<Box<dyn RelaySubscriptionSession>>, targets: BTreeMap<RelayUrl, radroots_transport::Target>, active_relays: BTreeSet<RelayUrl>, + resume_cursors: BTreeMap<radroots_transport::target::TargetFingerprint, RelayCursor>, cursors: BTreeMap<radroots_transport::target::TargetFingerprint, RelayCursor>, checkpoints: BTreeMap<radroots_transport::target::TargetFingerprint, SubscriptionCheckpoint>, + seen_event_ids: BTreeSet<String>, event_count: u16, terminal: Option<SubscriptionEnd>, cancellation_requested: Arc<AtomicBool>, @@ -242,8 +244,10 @@ impl RelayEventSubscription { session: None, targets: BTreeMap::new(), active_relays: BTreeSet::new(), + resume_cursors: BTreeMap::new(), cursors: BTreeMap::new(), checkpoints: BTreeMap::new(), + seen_event_ids: BTreeSet::new(), event_count: 0, terminal: Some(terminal), cancellation_requested: Arc::new(AtomicBool::new(false)), @@ -334,31 +338,54 @@ impl RelayEventSubscription { if !self.request.selector().matches(&event) { return Err(radroots_transport::Error::UnexpectedSubscriptionEvent); } + let event_id = event.id_str().to_owned(); + let created_at = event.created_at(); + if self.seen_event_ids.contains(event_id.as_str()) { + return Ok(None); + } if self - .cursors + .resume_cursors .get(target.fingerprint()) - .is_some_and(|cursor| !cursor.precedes(event.created_at(), event.id_str())) + .is_some_and(|cursor| { + created_at < cursor.created_at_unix_s() + || (created_at == cursor.created_at_unix_s() + && event_id.as_str() == cursor.event_id()) + }) { return Ok(None); } - let cursor = RelayCursor::new(event.created_at(), event.id_str()) + let cursor = RelayCursor::new(created_at, event_id.clone()) .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 cursor_advances = self + .cursors + .get(target.fingerprint()) + .is_none_or(|current| cursor > *current); + if cursor_advances { + self.cursors + .insert(target.fingerprint().clone(), cursor.clone()); + let opaque = encode_cursor(&self.request, target.fingerprint(), &cursor)?; + self.checkpoints.insert( + target.fingerprint().clone(), + SubscriptionCheckpoint::new(target.fingerprint().clone(), opaque), + ); + } + let checkpoint = self + .checkpoints + .get(target.fingerprint()) + .cloned() + .ok_or(radroots_transport::Error::UnexpectedSubscriptionEvent)?; let provenance = EventProvenance::new( radroots_transport::TransportId::NOSTR, target.fingerprint().clone(), unix_time_ms().max(1), )? - .with_cursor(opaque); + .with_cursor(checkpoint.cursor().clone()); 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.seen_event_ids.insert(event_id); self.event_count = self.event_count.saturating_add(1); self.status .record_read(relay, true, false, unix_time_ms().max(1)); @@ -526,13 +553,16 @@ impl EventSubscriber for NostrTransport { } }; + let resume_cursors = cursors.clone(); Ok(Box::new(RelayEventSubscription { request, session: Some(session), targets, active_relays, + resume_cursors, cursors, checkpoints, + seen_event_ids: BTreeSet::new(), event_count: 0, terminal: None, cancellation_requested: Arc::new(AtomicBool::new(false)), @@ -932,6 +962,10 @@ mod tests { raw: events[0].clone(), }), ScriptedItem::Item(RelaySubscriptionItem::Event { + relay: relay_url.clone(), + raw: events[1].clone(), + }), + ScriptedItem::Item(RelaySubscriptionItem::Event { relay: relay_url, raw: events[2].clone(), }), @@ -942,10 +976,18 @@ mod tests { let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") else { panic!("event expected"); }; + assert_eq!(event.observed().event().id_str(), event_id(&events[0])); + assert_eq!( + event.checkpoint().cursor().as_str().split(':').nth(2), + Some(event_id(&events[1]).as_str()) + ); + 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") + event.checkpoint().cursor().as_str().split(':').nth(2), + Some(event_id(&events[2]).as_str()) ); assert_eq!( client.query().targets[0].since_unix_seconds, @@ -954,6 +996,44 @@ mod tests { } #[tokio::test] + async fn live_subscription_accepts_same_second_events_in_relay_arrival_order() { + let relay = "wss://one.example"; + let relay_url = RelayUrl::parse(relay, RelayUrlPolicy::Public).expect("relay"); + let mut events = [ + signed_event("same-second-a", 1_800_000_000), + signed_event("same-second-b", 1_800_000_000), + ]; + events.sort_by_key(|event| event_id(event)); + let client = Arc::new(MockSubscriptionClient::new([ + ScriptedItem::Item(RelaySubscriptionItem::Event { + relay: relay_url.clone(), + raw: events[1].clone(), + }), + ScriptedItem::Item(RelaySubscriptionItem::Event { + relay: relay_url, + raw: events[0].clone(), + }), + ])); + let transport = NostrTransport::with_subscription_client(configured(&[relay]), client); + let mut subscription = transport + .subscribe(request(&[relay], 2)) + .await + .expect("subscription"); + + for expected in [&events[1], &events[0]] { + let SubscriptionNext::Event(event) = subscription.next().await.expect("next event") + else { + panic!("event expected"); + }; + assert_eq!(event.observed().event().id_str(), event_id(expected)); + assert_eq!( + event.checkpoint().cursor().as_str().split(':').nth(2), + Some(event_id(&events[1]).as_str()) + ); + } + } + + #[tokio::test] async fn malformed_or_mismatched_checkpoint_fails_before_backend_work() { let relay = "wss://one.example"; let client = Arc::new(MockSubscriptionClient::new([])); diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs @@ -83,7 +83,9 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() { "Deadline expiry is `Cancelled`", "exceeds the bounded inventory is `Partial`", "inclusive `since` timestamp", - "event-ID tie breaker", + "permits at-least-once", + "same-second events in relay arrival order", + "never regresses the canonical target checkpoint", "upstream auto-close deadline", "adapter-owned worker", "validates the exact request, writable\nrelay bindings, and signed-event conversion without reading a clock, polling\nstatus, or performing relay I/O", @@ -299,7 +301,9 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() { "impl EventSubscriber for NostrTransport", "SubscribeAutoCloseOptions::default()", "ReqExitPolicy::WaitDurationAfterEOSE(query.timeout)", - "cursor.precedes(event.created_at(), event.id_str())", + "self.seen_event_ids.contains(event_id.as_str())", + "resume_cursors", + "let cursor_advances = self", "self.terminate(SubscriptionEndReason::Cancelled)", ] { assert!(