lib

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

commit 79d7818c8fe22a425f9524b884ddf59d25f0ef89
parent 7d7b454b4c9ed86569671993bd03ca868b676665
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 02:48:43 +0000

transport: support exact indexed tag selectors

- bound and canonicalize exact single-letter tag filters
- map tag filters into Nostr fetch and subscription requests
- bind continuation scopes to every selector dimension
- freeze error, serialization, documentation, and API contracts

Diffstat:
Mcontracts/api_baselines/radroots_transport.txt | 11+++++++++++
Mcrates/transport/README.md | 2+-
Mcrates/transport/src/error.rs | 10++++++++++
Mcrates/transport/src/source.rs | 164++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mcrates/transport/tests/error_contract.rs | 12++++++++++++
Mcrates/transport/tests/package_boundary.rs | 6++++++
Mcrates/transport/tests/source_contract.rs | 107++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mcrates/transport_nostr/README.md | 14+++++++-------
Mcrates/transport_nostr/src/source.rs | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Mcrates/transport_nostr/src/subscription.rs | 11+++++++++--
Mcrates/transport_nostr/tests/package_boundary.rs | 1+
11 files changed, 404 insertions(+), 25 deletions(-)

diff --git a/contracts/api_baselines/radroots_transport.txt b/contracts/api_baselines/radroots_transport.txt @@ -97,6 +97,7 @@ pub radroots_transport::error::Error::DeliveryTargetReceiptStatusMismatch pub radroots_transport::error::Error::DuplicateDeliveryTargetReceipt pub radroots_transport::error::Error::DuplicateFetchAuthor pub radroots_transport::error::Error::DuplicateFetchKind +pub radroots_transport::error::Error::DuplicateFetchTagValue pub radroots_transport::error::Error::DuplicateFetchTargetOutcome pub radroots_transport::error::Error::DuplicateRequiredTargetFingerprint pub radroots_transport::error::Error::DuplicateSubscriptionCheckpoint @@ -125,6 +126,8 @@ pub radroots_transport::error::Error::InvalidFetchCursor pub radroots_transport::error::Error::InvalidFetchDeadline pub radroots_transport::error::Error::InvalidFetchLimit pub radroots_transport::error::Error::InvalidFetchRequestId +pub radroots_transport::error::Error::InvalidFetchTagKey +pub radroots_transport::error::Error::InvalidFetchTagValue pub radroots_transport::error::Error::InvalidFetchTimeRange pub radroots_transport::error::Error::InvalidObservedAt pub radroots_transport::error::Error::InvalidPayloadBytes @@ -380,11 +383,13 @@ pub struct radroots_transport::source::FetchSelector impl radroots_transport::source::FetchSelector pub const fn radroots_transport::source::FetchSelector::all() -> Self pub fn radroots_transport::source::FetchSelector::authors(&self) -> &[radroots_identity::key::PublicKey] +pub fn radroots_transport::source::FetchSelector::exact_tag_filters(&self) -> impl core::iter::traits::iterator::Iterator<Item = (char, &[alloc::string::String])> + '_ pub fn radroots_transport::source::FetchSelector::kinds(&self) -> &[u32] pub fn radroots_transport::source::FetchSelector::matches(&self, &radroots_event::draft::SignedEvent) -> bool pub const fn radroots_transport::source::FetchSelector::since_unix_seconds(&self) -> core::option::Option<u64> pub const fn radroots_transport::source::FetchSelector::until_unix_seconds(&self) -> core::option::Option<u64> pub fn radroots_transport::source::FetchSelector::with_authors(self, alloc::vec::Vec<radroots_identity::key::PublicKey>) -> core::result::Result<Self, radroots_transport::error::Error> +pub fn radroots_transport::source::FetchSelector::with_exact_tag_value(self, char, impl core::convert::AsRef<str>) -> core::result::Result<Self, radroots_transport::error::Error> pub fn radroots_transport::source::FetchSelector::with_kinds(self, alloc::vec::Vec<u32>) -> core::result::Result<Self, radroots_transport::error::Error> pub fn radroots_transport::source::FetchSelector::with_since_unix_seconds(self, u64) -> core::result::Result<Self, radroots_transport::error::Error> pub fn radroots_transport::source::FetchSelector::with_until_unix_seconds(self, u64) -> core::result::Result<Self, radroots_transport::error::Error> @@ -464,6 +469,9 @@ pub const radroots_transport::source::FETCH_PAGE_MAX_EVENTS: u16 pub const radroots_transport::source::FETCH_REQUEST_ID_MAX_BYTES: usize pub const radroots_transport::source::FETCH_SELECTOR_MAX_AUTHORS: usize pub const radroots_transport::source::FETCH_SELECTOR_MAX_KINDS: usize +pub const radroots_transport::source::FETCH_SELECTOR_MAX_TAG_KEYS: usize +pub const radroots_transport::source::FETCH_SELECTOR_MAX_TAG_VALUES: usize +pub const radroots_transport::source::FETCH_SELECTOR_TAG_VALUE_MAX_BYTES: usize pub const radroots_transport::source::SUBSCRIPTION_MAX_EVENTS: u16 pub const radroots_transport::source::SUBSCRIPTION_REQUEST_ID_MAX_BYTES: usize pub trait radroots_transport::source::EventSource: core::marker::Send + core::marker::Sync @@ -576,6 +584,7 @@ pub radroots_transport::Error::DeliveryTargetReceiptStatusMismatch pub radroots_transport::Error::DuplicateDeliveryTargetReceipt pub radroots_transport::Error::DuplicateFetchAuthor pub radroots_transport::Error::DuplicateFetchKind +pub radroots_transport::Error::DuplicateFetchTagValue pub radroots_transport::Error::DuplicateFetchTargetOutcome pub radroots_transport::Error::DuplicateRequiredTargetFingerprint pub radroots_transport::Error::DuplicateSubscriptionCheckpoint @@ -604,6 +613,8 @@ pub radroots_transport::Error::InvalidFetchCursor pub radroots_transport::Error::InvalidFetchDeadline pub radroots_transport::Error::InvalidFetchLimit pub radroots_transport::Error::InvalidFetchRequestId +pub radroots_transport::Error::InvalidFetchTagKey +pub radroots_transport::Error::InvalidFetchTagValue pub radroots_transport::Error::InvalidFetchTimeRange pub radroots_transport::Error::InvalidObservedAt pub radroots_transport::Error::InvalidPayloadBytes diff --git a/crates/transport/README.md b/crates/transport/README.md @@ -22,7 +22,7 @@ The authoritative package charter is the 3. The caller creates a bounded [`FetchRequest`], [`SubscriptionRequest`], or [`DeliveryRequest`] with a request identity and absolute deadline. Inbound operations may carry a validated [`FetchSelector`] for exact kinds, - authors, and inclusive event-time bounds. + authors, indexed single-letter tag values, and inclusive event-time bounds. 4. A dyn-compatible [`EventSource`], [`EventSubscriber`], or [`EventSink`] implementation performs only the requested operation. 5. The caller validates [`FetchPage`] or [`DeliveryReceipt`] against the diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs @@ -25,6 +25,9 @@ pub enum Error { FetchSelectorTooLarge, DuplicateFetchKind, DuplicateFetchAuthor, + InvalidFetchTagKey, + InvalidFetchTagValue, + DuplicateFetchTagValue, InvalidFetchTimeRange, EmptyFetchCursor, InvalidFetchCursor, @@ -108,6 +111,13 @@ impl fmt::Display for Error { Self::DuplicateFetchAuthor => { f.write_str("transport fetch selector contains a duplicate author") } + Self::InvalidFetchTagKey => f.write_str("transport fetch selector tag key is invalid"), + Self::InvalidFetchTagValue => { + f.write_str("transport fetch selector tag value is invalid") + } + Self::DuplicateFetchTagValue => { + f.write_str("transport fetch selector contains a duplicate tag value") + } Self::InvalidFetchTimeRange => { f.write_str("transport fetch selector time range is invalid") } diff --git a/crates/transport/src/source.rs b/crates/transport/src/source.rs @@ -5,7 +5,12 @@ use crate::{ outcome::FetchTargetOutcome, target::{TargetFingerprint, TargetSet}, }; -use alloc::{boxed::Box, collections::BTreeSet, string::String, vec::Vec}; +use alloc::{ + boxed::Box, + collections::{BTreeMap, BTreeSet}, + string::String, + vec::Vec, +}; use core::{fmt, future::Future, pin::Pin}; use radroots_event::SignedEvent; use radroots_identity::PublicKey; @@ -25,6 +30,19 @@ pub const FETCH_PAGE_MAX_EVENTS: u16 = 1_000; pub const FETCH_SELECTOR_MAX_KINDS: usize = 64; /// Maximum distinct event authors in one source selector. pub const FETCH_SELECTOR_MAX_AUTHORS: usize = 256; +/// Maximum distinct exact single-letter tag keys in one source selector. +pub const FETCH_SELECTOR_MAX_TAG_KEYS: usize = 26; +/// Maximum exact tag values across one source selector. +pub const FETCH_SELECTOR_MAX_TAG_VALUES: usize = 256; +/// Maximum UTF-8 bytes in one exact tag value. +pub const FETCH_SELECTOR_TAG_VALUE_MAX_BYTES: usize = 4_096; + +// Deliberate representation indirection keeps every selector-bearing request +// and terminal value compact while allocating nothing for the common no-tag +// case. The map itself still owns its bounded tree nodes. +#[allow(clippy::box_collection)] +type ExactTagFilters = Box<BTreeMap<char, Vec<String>>>; + /// Maximum encoded live-subscription request identity length. pub const SUBSCRIPTION_REQUEST_ID_MAX_BYTES: usize = 256; /// Maximum number of events one live subscription may emit. @@ -199,15 +217,21 @@ impl SubscriptionBounds { /// Transport-neutral constraints applied before a source page is bounded. /// -/// An empty kind or author collection means "any" for that dimension. Time -/// bounds are inclusive Unix seconds. Adapters must apply every configured -/// dimension remotely when their protocol supports it and must defensively -/// exclude non-matching events before returning a page. +/// An empty kind, author, or tag collection means "any" for that dimension. +/// Values for one tag key are alternatives, while distinct tag keys are +/// conjunctive. Time bounds are inclusive Unix seconds. Adapters must apply +/// every configured dimension remotely when their protocol supports it and +/// must defensively exclude non-matching events before returning a page. #[cfg_attr(feature = "serde", derive(serde::Serialize))] #[derive(Clone, Debug, Default, Eq, PartialEq)] pub struct FetchSelector { kinds: Vec<u32>, authors: Vec<PublicKey>, + #[cfg_attr( + feature = "serde", + serde(serialize_with = "serde_impl::serialize_exact_tags") + )] + exact_tags: Option<ExactTagFilters>, since_unix_seconds: Option<u64>, until_unix_seconds: Option<u64>, } @@ -219,6 +243,7 @@ impl FetchSelector { Self { kinds: Vec::new(), authors: Vec::new(), + exact_tags: None, since_unix_seconds: None, until_unix_seconds: None, } @@ -250,6 +275,50 @@ impl FetchSelector { Ok(self) } + /// Requires one exact indexed single-letter tag value. + /// + /// Repeating a key adds an alternative value for that key. Different keys + /// are conjunctive. Keys are lowercase ASCII letters and values are + /// non-empty bounded UTF-8 strings. + pub fn with_exact_tag_value( + mut self, + key: char, + value: impl AsRef<str>, + ) -> Result<Self, Error> { + if !key.is_ascii_lowercase() { + return Err(Error::InvalidFetchTagKey); + } + let value = value.as_ref(); + if value.is_empty() + || value.len() > FETCH_SELECTOR_TAG_VALUE_MAX_BYTES + || value.chars().any(char::is_control) + { + return Err(Error::InvalidFetchTagValue); + } + let exact_tags = self.exact_tags.as_deref(); + let total_values = exact_tags + .into_iter() + .flat_map(BTreeMap::values) + .map(Vec::len) + .sum::<usize>(); + if (!exact_tags.is_some_and(|tags| tags.contains_key(&key)) + && exact_tags.is_some_and(|tags| tags.len() == FETCH_SELECTOR_MAX_TAG_KEYS)) + || total_values == FETCH_SELECTOR_MAX_TAG_VALUES + { + return Err(Error::FetchSelectorTooLarge); + } + let values = self + .exact_tags + .get_or_insert_with(|| Box::new(BTreeMap::new())) + .entry(key) + .or_default(); + match values.binary_search_by(|candidate| candidate.as_str().cmp(value)) { + Ok(_) => return Err(Error::DuplicateFetchTagValue), + Err(position) => values.insert(position, String::from(value)), + } + Ok(self) + } + /// Sets an inclusive lower event-time bound. pub fn with_since_unix_seconds(mut self, since: u64) -> Result<Self, Error> { if self.until_unix_seconds.is_some_and(|until| since > until) { @@ -278,6 +347,14 @@ impl FetchSelector { self.authors.as_slice() } + /// Returns exact tag filters in canonical key order. + pub fn exact_tag_filters(&self) -> impl Iterator<Item = (char, &[String])> + '_ { + self.exact_tags + .iter() + .flat_map(|tags| tags.iter()) + .map(|(key, values)| (*key, values.as_slice())) + } + /// Returns the inclusive lower event-time bound. pub const fn since_unix_seconds(&self) -> Option<u64> { self.since_unix_seconds @@ -293,6 +370,20 @@ impl FetchSelector { pub fn matches(&self, event: &SignedEvent) -> bool { (self.kinds.is_empty() || self.kinds.binary_search(&event.kind()).is_ok()) && (self.authors.is_empty() || self.authors.binary_search(event.pubkey()).is_ok()) + && self.exact_tags.as_deref().is_none_or(|tags| { + tags.iter().all(|(key, values)| { + event.envelope().tag_slices().iter().any(|tag| { + let elements = tag.as_slice(); + elements.first().is_some_and(|candidate| { + candidate.len() == 1 && candidate.starts_with(*key) + }) && elements.get(1).is_some_and(|value| { + values + .binary_search_by(|candidate| candidate.as_str().cmp(value)) + .is_ok() + }) + }) + }) + }) && self .since_unix_seconds .is_none_or(|since| event.created_at() >= since) @@ -934,6 +1025,19 @@ pub trait EventSubscriber: Send + Sync { mod serde_impl { use super::*; + pub(super) fn serialize_exact_tags<S>( + exact_tags: &Option<ExactTagFilters>, + serializer: S, + ) -> Result<S::Ok, S::Error> + where + S: serde::Serializer, + { + match exact_tags { + Some(exact_tags) => serde::Serialize::serialize(exact_tags, serializer), + None => serde::Serialize::serialize(&BTreeMap::<char, Vec<String>>::new(), serializer), + } + } + impl serde::Serialize for FetchRequestId { fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> where @@ -1025,6 +1129,44 @@ mod serde_impl { } } + #[derive(Default)] + struct ExactTagsWire(Vec<(char, Vec<String>)>); + + impl<'de> serde::Deserialize<'de> for ExactTagsWire { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: serde::Deserializer<'de>, + { + struct ExactTagsVisitor; + + impl<'de> serde::de::Visitor<'de> for ExactTagsVisitor { + type Value = ExactTagsWire; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("a map of unique exact single-letter tag filters") + } + + fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error> + where + A: serde::de::MapAccess<'de>, + { + let mut entries = Vec::<(char, Vec<String>)>::new(); + while let Some((key, values)) = map.next_entry::<char, Vec<String>>()? { + if entries.iter().any(|(candidate, _)| *candidate == key) { + return Err(serde::de::Error::custom( + "transport fetch selector contains a duplicate tag key", + )); + } + entries.push((key, values)); + } + Ok(ExactTagsWire(entries)) + } + } + + deserializer.deserialize_map(ExactTagsVisitor) + } + } + #[derive(serde::Deserialize)] #[serde(deny_unknown_fields)] struct FetchRequestWire { @@ -1090,6 +1232,8 @@ mod serde_impl { kinds: Vec<u32>, #[serde(default)] authors: Vec<PublicKey>, + #[serde(default)] + exact_tags: ExactTagsWire, since_unix_seconds: Option<u64>, until_unix_seconds: Option<u64>, } @@ -1098,6 +1242,16 @@ mod serde_impl { let selector = FetchSelector::all() .with_kinds(wire.kinds) .and_then(|selector| selector.with_authors(wire.authors)) + .and_then(|selector| { + wire.exact_tags + .0 + .into_iter() + .try_fold(selector, |selector, (key, values)| { + values.into_iter().try_fold(selector, |selector, value| { + selector.with_exact_tag_value(key, value) + }) + }) + }) .and_then(|selector| match wire.since_unix_seconds { Some(since) => selector.with_since_unix_seconds(since), None => Ok(selector), diff --git a/crates/transport/tests/error_contract.rs b/crates/transport/tests/error_contract.rs @@ -60,6 +60,18 @@ fn every_transport_error_has_stable_operator_facing_text() { "transport fetch selector contains a duplicate author", ), ( + Error::InvalidFetchTagKey, + "transport fetch selector tag key is invalid", + ), + ( + Error::InvalidFetchTagValue, + "transport fetch selector tag value is invalid", + ), + ( + Error::DuplicateFetchTagValue, + "transport fetch selector contains a duplicate tag value", + ), + ( Error::InvalidFetchTimeRange, "transport fetch selector time range is invalid", ), diff --git a/crates/transport/tests/package_boundary.rs b/crates/transport/tests/package_boundary.rs @@ -84,6 +84,7 @@ fn package_documentation_and_reviewed_api_baseline_are_complete() { "## Intended consumers", "radroots_crates_release_v1.toml", "examples/host_transport.rs", + "indexed single-letter tag values", ] { assert!(README.contains(required), "README is missing {required}"); } @@ -119,6 +120,11 @@ fn package_documentation_and_reviewed_api_baseline_are_complete() { "pub struct radroots_transport::SubscriptionRequest", "pub enum radroots_transport::SubscriptionEndReason", "pub const radroots_transport::source::SUBSCRIPTION_MAX_EVENTS: u16", + "pub const radroots_transport::source::FETCH_SELECTOR_MAX_TAG_VALUES: usize", + "pub const radroots_transport::source::FETCH_SELECTOR_MAX_TAG_KEYS: usize", + "pub const radroots_transport::source::FETCH_SELECTOR_TAG_VALUE_MAX_BYTES: usize", + "pub fn radroots_transport::source::FetchSelector::with_exact_tag_value", + "pub fn radroots_transport::source::FetchSelector::exact_tag_filters", "pub fn radroots_transport::target::Target::new(radroots_transport::TransportId", "pub fn radroots_transport::target::Target::kind(&self) -> &radroots_transport::TransportId", ] { diff --git a/crates/transport/tests/source_contract.rs b/crates/transport/tests/source_contract.rs @@ -14,8 +14,9 @@ use radroots_transport::{ outcome::{FetchTargetOutcome, FetchTargetState}, source::{ EventProvenance, FETCH_CURSOR_MAX_BYTES, FETCH_PAGE_MAX_EVENTS, FETCH_REQUEST_ID_MAX_BYTES, - FETCH_SELECTOR_MAX_AUTHORS, FETCH_SELECTOR_MAX_KINDS, FetchBounds, FetchCursor, - FetchSelector, NextPage, ObservedEvent, + FETCH_SELECTOR_MAX_AUTHORS, FETCH_SELECTOR_MAX_KINDS, FETCH_SELECTOR_MAX_TAG_VALUES, + FETCH_SELECTOR_TAG_VALUE_MAX_BYTES, FetchBounds, FetchCursor, FetchSelector, NextPage, + ObservedEvent, }, }; @@ -25,13 +26,15 @@ fn target(uri: &str) -> Target { #[test] fn fetch_selector_is_bounded_canonical_and_request_bound() { - let event = signed_event(); + let event = tagged_event(); let author = *event.pubkey(); let selector = FetchSelector::all() .with_kinds(vec![1, 0]) .expect("kind selector") .with_authors(vec![author]) .expect("author selector") + .with_exact_tag_value('d', "trade-1") + .expect("tag selector") .with_since_unix_seconds(1_700_000_000) .expect("since") .with_until_unix_seconds(1_700_000_100) @@ -39,7 +42,26 @@ fn fetch_selector_is_bounded_canonical_and_request_bound() { assert_eq!(selector.kinds(), &[0, 1]); assert_eq!(selector.authors(), &[author]); + let exact_tags = selector.exact_tag_filters().collect::<Vec<_>>(); + assert_eq!(exact_tags.len(), 1); + assert_eq!(exact_tags[0].0, 'd'); + assert_eq!(exact_tags[0].1, &[String::from("trade-1")]); assert!(selector.matches(&event)); + let encoded = serde_json::to_string(&selector).expect("selector JSON"); + assert_eq!( + serde_json::from_str::<FetchSelector>(encoded.as_str()).expect("selector round trip"), + selector + ); + assert!( + serde_json::from_value::<FetchSelector>(serde_json::json!({ + "kinds": [], + "authors": [], + "exact_tags": {"D": ["trade-1"]}, + "since_unix_seconds": null, + "until_unix_seconds": null + })) + .is_err() + ); assert_eq!( FetchSelector::all() .with_kinds(vec![1, 1]) @@ -64,6 +86,73 @@ fn fetch_selector_is_bounded_canonical_and_request_bound() { .expect_err("too many authors"), Error::FetchSelectorTooLarge ); + for invalid in ['D', '0', '#', 'é'] { + assert_eq!( + FetchSelector::all() + .with_exact_tag_value(invalid, "trade-1") + .expect_err("invalid tag key"), + Error::InvalidFetchTagKey + ); + } + for invalid in [String::new(), String::from("line\nbreak")] { + assert_eq!( + FetchSelector::all() + .with_exact_tag_value('d', invalid) + .expect_err("invalid tag value"), + Error::InvalidFetchTagValue + ); + } + assert_eq!( + FetchSelector::all() + .with_exact_tag_value('d', "x".repeat(FETCH_SELECTOR_TAG_VALUE_MAX_BYTES + 1)) + .expect_err("oversized tag value"), + Error::InvalidFetchTagValue + ); + assert!( + FetchSelector::all() + .with_exact_tag_value('d', "x".repeat(FETCH_SELECTOR_TAG_VALUE_MAX_BYTES)) + .is_ok() + ); + assert_eq!( + FetchSelector::all() + .with_exact_tag_value('d', "trade-1") + .and_then(|selector| selector.with_exact_tag_value('d', "trade-1")) + .expect_err("duplicate tag value"), + Error::DuplicateFetchTagValue + ); + let maximum = (0..FETCH_SELECTOR_MAX_TAG_VALUES) + .try_fold(FetchSelector::all(), |selector, index| { + selector.with_exact_tag_value('d', format!("trade-{index:03}")) + }); + assert!(maximum.is_ok()); + assert_eq!( + maximum + .and_then(|selector| selector.with_exact_tag_value('d', "trade-overflow")) + .expect_err("too many tag values"), + Error::FetchSelectorTooLarge + ); + let every_key = (0..radroots_transport::source::FETCH_SELECTOR_MAX_TAG_KEYS).try_fold( + FetchSelector::all(), + |selector, index| { + selector.with_exact_tag_value( + char::from(b'a' + u8::try_from(index).expect("bounded key index")), + "value", + ) + }, + ); + assert_eq!( + every_key + .expect("all lowercase keys") + .exact_tag_filters() + .count(), + radroots_transport::source::FETCH_SELECTOR_MAX_TAG_KEYS + ); + assert!( + serde_json::from_str::<FetchSelector>( + r#"{"kinds":[],"authors":[],"exact_tags":{"d":["one"],"d":["two"]},"since_unix_seconds":null,"until_unix_seconds":null}"#, + ) + .is_err() + ); assert_eq!( FetchSelector::all() .with_since_unix_seconds(2) @@ -89,6 +178,15 @@ fn signed_event() -> SignedEvent { SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") } +fn tagged_event() -> SignedEvent { + let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#; + let mut wire = Nip01EventWire::parse_json(raw).expect("wire event"); + wire.tags = vec![vec![String::from("d"), String::from("trade-1")]]; + wire.id = wire.computed_event_id().expect("event id").into_string(); + let raw = serde_json::to_string(&wire).expect("event JSON"); + SignedEvent::from_wire_verified_id(wire, raw.as_str()).expect("signed event") +} + fn request(targets: TargetSet, limit: u16) -> FetchRequest { FetchRequest::new( "fetch-request", @@ -184,6 +282,9 @@ fn selectors_expose_bounds_and_reject_each_nonmatching_dimension() { .with_authors(vec![other_author]) .expect("authors"), FetchSelector::all() + .with_exact_tag_value('d', "other-trade") + .expect("tag"), + FetchSelector::all() .with_since_unix_seconds(event.created_at() + 1) .expect("since"), FetchSelector::all() diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md @@ -130,13 +130,13 @@ check before handing control to another network boundary. ## 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. +transport-neutral kind, author, exact indexed single-letter tag, 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; diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs @@ -4,7 +4,7 @@ use crate::{NostrTransport, RelayCursor, RelayUrl, status}; use core::cmp::Ordering; use core::time::Duration; use futures::{StreamExt, stream}; -use nostr_sdk::prelude::{Filter, JsonUtil, Kind, Timestamp}; +use nostr_sdk::prelude::{Filter, JsonUtil, Kind, SingleLetterTag, Timestamp}; use radroots_transport::{ BoxFuture, EventSource, FetchPage, FetchRequest, outcome::{FetchTargetOutcome, FetchTargetState}, @@ -94,10 +94,6 @@ impl RelaySourceClient for LiveRelaySourceClient { if authors.len() != selector.authors().len() { return Ok(Vec::new()); } - self.client.add_relay(url.as_str()).await?; - self.client - .try_connect_relay(url.as_str(), connect_timeout) - .await?; let mut filter = Filter::new().limit(UPSTREAM_FETCH_LIMIT); if !kinds.is_empty() { filter = filter.kinds(kinds); @@ -105,6 +101,9 @@ impl RelaySourceClient for LiveRelaySourceClient { if !authors.is_empty() { filter = filter.authors(authors); } + filter = apply_exact_tag_filters(filter, &selector).map_err(|()| { + String::from("validated indexed tag cannot be encoded") + })?; if let Some(since) = selector.since_unix_seconds() { filter = filter.since(Timestamp::from_secs(since)); } @@ -112,12 +111,20 @@ impl RelaySourceClient for LiveRelaySourceClient { filter = filter.until(Timestamp::from_secs(until)); } self.client + .add_relay(url.as_str()) + .await + .map_err(|error| error.to_string())?; + self.client + .try_connect_relay(url.as_str(), connect_timeout) + .await + .map_err(|error| error.to_string())?; + self.client .fetch_events_from([url.as_str()], filter, timeout) .await .map(|events| events.iter().map(JsonUtil::as_json).collect()) + .map_err(|error| error.to_string()) } - .await - .map_err(|error: nostr_sdk::client::Error| error.to_string()); + .await; RelayFetchBatch { relay, result } } })) @@ -387,11 +394,51 @@ fn request_scope(request: &FetchRequest) -> String { hasher.update(author.as_bytes()); } hasher.update([0]); + hash_exact_tag_filters(&mut hasher, request.selector()); hash_optional_u64(&mut hasher, request.selector().since_unix_seconds()); hash_optional_u64(&mut hasher, request.selector().until_unix_seconds()); hex_encode(&hasher.finalize()) } +pub(crate) fn hash_exact_tag_filters( + hasher: &mut Sha256, + selector: &radroots_transport::source::FetchSelector, +) { + let mut filters = selector.exact_tag_filters().peekable(); + if filters.peek().is_none() { + return; + } + hasher.update([2]); + for (key, values) in filters { + hasher.update([u8::try_from(key).expect("validated ASCII tag key")]); + hasher.update( + u64::try_from(values.len()) + .expect("bounded tag-value count") + .to_be_bytes(), + ); + for value in values { + hasher.update( + u64::try_from(value.len()) + .expect("bounded tag-value length") + .to_be_bytes(), + ); + hasher.update(value.as_bytes()); + } + } + hasher.update([0]); +} + +pub(crate) fn apply_exact_tag_filters( + mut filter: Filter, + selector: &radroots_transport::source::FetchSelector, +) -> Result<Filter, ()> { + for (key, values) in selector.exact_tag_filters() { + let tag = SingleLetterTag::from_char(key).map_err(|_| ())?; + filter = filter.custom_tags(tag, values.iter().map(String::as_str)); + } + Ok(filter) +} + fn hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) { match value { Some(value) => { @@ -739,6 +786,19 @@ mod tests { .expect_err("scope mismatch"); assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); + let tagged_selector = FetchSelector::all() + .with_exact_tag_value('d', "trade-1") + .expect("tag selector"); + let error = futures::executor::block_on( + transport().fetch( + request(1) + .with_selector(tagged_selector) + .with_cursor(cursor.clone()), + ), + ) + .expect_err("tag scope mismatch"); + assert_eq!(error, radroots_transport::Error::InvalidFetchCursor); + let other_targets = TargetSet::new(vec![Target::nostr_relay("wss://one.example").expect("one")]) .expect("targets"); @@ -758,6 +818,23 @@ mod tests { } #[test] + fn exact_tag_selector_maps_to_nostr_filter_and_is_defensively_enforced() { + let selector = FetchSelector::all() + .with_exact_tag_value('d', "trade-2") + .and_then(|selector| selector.with_exact_tag_value('d', "trade-1")) + .expect("tag selector"); + let encoded = apply_exact_tag_filters(Filter::new(), &selector) + .expect("Nostr filter") + .as_json(); + assert!(encoded.contains("\"#d\":[\"trade-1\",\"trade-2\"]")); + + let selected = + futures::executor::block_on(transport().fetch(request(10).with_selector(selector))) + .expect("selected page"); + assert!(selected.events().is_empty()); + } + + #[test] fn live_source_short_circuits_selectors_that_cannot_be_encoded() { let client = LiveRelaySourceClient::isolated(); let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay"); diff --git a/crates/transport_nostr/src/subscription.rs b/crates/transport_nostr/src/subscription.rs @@ -90,12 +90,12 @@ impl RelaySubscriptionClient for LiveRelaySubscriptionClient { for target in query.targets { let url = target.relay.as_str().to_owned(); + let filter = subscription_filter(&query.selector, target.since_unix_seconds)?; 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)); } @@ -593,6 +593,7 @@ fn subscription_filter( if !authors.is_empty() { filter = filter.authors(authors); } + filter = crate::source::apply_exact_tag_filters(filter, selector)?; if let Some(since) = since { filter = filter.since(Timestamp::from_secs(since)); } @@ -670,6 +671,7 @@ fn cursor_scope( hasher.update(author.as_bytes()); } hasher.update([0]); + crate::source::hash_exact_tag_filters(&mut hasher, request.selector()); hash_optional_u64(&mut hasher, request.selector().since_unix_seconds()); hash_optional_u64(&mut hasher, request.selector().until_unix_seconds()); hex_encode(&hasher.finalize()) @@ -861,6 +863,8 @@ mod tests { let selector = FetchSelector::all() .with_kinds(vec![1]) .expect("kind") + .with_exact_tag_value('d', "trade-1") + .expect("tag") .with_since_unix_seconds(1_700_000_000) .expect("since") .with_until_unix_seconds(1_800_000_000) @@ -892,6 +896,7 @@ mod tests { .expect("filter") .as_json(); assert!(filter.contains("\"kinds\":[1]")); + assert!(filter.contains("\"#d\":[\"trade-1\"]")); assert!(filter.contains("\"since\":1700000000")); assert!(filter.contains("\"until\":1800000000")); } @@ -968,7 +973,9 @@ mod tests { Some(radroots_transport::Error::InvalidFetchCursor) ); - let other_selector = FetchSelector::all().with_kinds(vec![2]).expect("selector"); + let other_selector = FetchSelector::all() + .with_exact_tag_value('d', "trade-1") + .expect("selector"); let scoped = encode_cursor( &base, &target, diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs @@ -78,6 +78,7 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() { "examples/configure_transport.rs", "contracts/api_baselines/radroots_transport_nostr.txt", "Live subscriptions use the same explicit readable targets", + "exact indexed single-letter tag", "inclusive `since` timestamp", "event-ID tie breaker", "upstream auto-close deadline",