commit be39b6fbcace4970f50cebbe6b9f816c3c7288da
parent 9e4ab3c50a4bcf5ba5f954017f5e6731a32e14fb
Author: triesap <tyson@radroots.org>
Date: Thu, 30 Jul 2026 17:12:13 +0000
transport: define bounded fetch pages and provenance
- validate request identities page limits deadlines and opaque cursors
- bind signed event observations to exact transport target provenance
- represent partial target outcomes and explicit continuation or cancellation
- verify serde invariants cancellation boundaries no-std wasm and workspace gates
Diffstat:
7 files changed, 951 insertions(+), 15 deletions(-)
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -18,6 +18,18 @@ pub enum RadrootsTransportError {
TargetSetTooLarge,
DuplicateTargetFingerprint,
InvalidTargetFingerprint,
+ EmptyFetchRequestId,
+ InvalidFetchRequestId,
+ InvalidFetchLimit,
+ InvalidFetchDeadline,
+ EmptyFetchCursor,
+ InvalidFetchCursor,
+ InvalidObservedAt,
+ UnexpectedFetchProvenance,
+ UnexpectedFetchTargetOutcome,
+ DuplicateFetchTargetOutcome,
+ FetchPageLimitExceeded,
+ FetchPageRequestMismatch,
InvalidSatisfactionPolicy,
EmptyRequiredTargetSet,
DuplicateRequiredTargetFingerprint,
@@ -63,6 +75,28 @@ impl fmt::Display for RadrootsTransportError {
Self::InvalidTargetFingerprint => {
f.write_str("transport target fingerprint is invalid")
}
+ Self::EmptyFetchRequestId => f.write_str("transport fetch request id is empty"),
+ Self::InvalidFetchRequestId => f.write_str("transport fetch request id is invalid"),
+ Self::InvalidFetchLimit => f.write_str("transport fetch limit is invalid"),
+ Self::InvalidFetchDeadline => f.write_str("transport fetch deadline is invalid"),
+ Self::EmptyFetchCursor => f.write_str("transport fetch cursor is empty"),
+ Self::InvalidFetchCursor => f.write_str("transport fetch cursor is invalid"),
+ Self::InvalidObservedAt => f.write_str("transport event observation time is invalid"),
+ Self::UnexpectedFetchProvenance => {
+ f.write_str("transport fetch page contains unexpected event provenance")
+ }
+ Self::UnexpectedFetchTargetOutcome => {
+ f.write_str("transport fetch page contains an unexpected target outcome")
+ }
+ Self::DuplicateFetchTargetOutcome => {
+ f.write_str("transport fetch page contains a duplicate target outcome")
+ }
+ Self::FetchPageLimitExceeded => {
+ f.write_str("transport fetch page exceeds its requested event limit")
+ }
+ Self::FetchPageRequestMismatch => {
+ f.write_str("transport fetch page does not match its request")
+ }
Self::InvalidSatisfactionPolicy => {
f.write_str("transport satisfaction policy is invalid")
}
diff --git a/crates/transport/src/outcome.rs b/crates/transport/src/outcome.rs
@@ -1 +1,81 @@
//! Normalized transport operation outcomes.
+
+use crate::target::TargetFingerprint;
+use alloc::string::String;
+
+/// Target-local result of one bounded fetch attempt.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum FetchTargetState {
+ /// The target reached its current end without error.
+ Complete,
+ /// The target produced some results but did not reach its current end.
+ Partial,
+ /// The target was not available for this operation.
+ Unavailable,
+ /// The attempt failed and a caller may choose to retry.
+ FailedRetryable,
+ /// The attempt failed and retrying the same request is not useful.
+ FailedTerminal,
+ /// Work for this target stopped because the operation was cancelled.
+ Cancelled,
+}
+
+impl FetchTargetState {
+ /// Whether a caller may choose to retry this target.
+ pub const fn is_retryable(self) -> bool {
+ matches!(
+ self,
+ Self::Partial | Self::Unavailable | Self::FailedRetryable
+ )
+ }
+
+ /// Whether this target reached a terminal state for the current request.
+ pub const fn is_terminal(self) -> bool {
+ matches!(self, Self::Complete | Self::FailedTerminal)
+ }
+}
+
+/// Explicit result for one requested source target.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct FetchTargetOutcome {
+ target: TargetFingerprint,
+ state: FetchTargetState,
+ message: Option<String>,
+}
+
+impl FetchTargetOutcome {
+ /// Creates a target-specific normalized outcome.
+ pub const fn new(target: TargetFingerprint, state: FetchTargetState) -> Self {
+ Self {
+ target,
+ state,
+ message: None,
+ }
+ }
+
+ /// Attaches caller-safe diagnostic detail.
+ #[must_use]
+ pub fn with_message(mut self, message: impl Into<String>) -> Self {
+ self.message = Some(message.into());
+ self
+ }
+
+ /// Returns the exact requested target fingerprint.
+ pub const fn target(&self) -> &TargetFingerprint {
+ &self.target
+ }
+
+ /// Returns normalized state.
+ pub const fn state(&self) -> FetchTargetState {
+ self.state
+ }
+
+ /// Returns caller-safe diagnostic detail.
+ pub fn message(&self) -> Option<&str> {
+ self.message.as_deref()
+ }
+}
diff --git a/crates/transport/src/source.rs b/crates/transport/src/source.rs
@@ -1,22 +1,382 @@
-//! Inbound event source SPI and page models.
+//! Inbound event source SPI and bounded page models.
-use crate::{Error, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest};
-use alloc::boxed::Box;
-use core::{future::Future, pin::Pin};
+use crate::{
+ Error, TransportId,
+ outcome::FetchTargetOutcome,
+ target::{TargetFingerprint, TargetSet},
+};
+use alloc::{boxed::Box, collections::BTreeSet, string::String, vec::Vec};
+use core::{fmt, future::Future, pin::Pin};
+use radroots_event::SignedEvent;
pub use crate::status::SourceStatus;
+/// Maximum encoded request identity length.
+pub const FETCH_REQUEST_ID_MAX_BYTES: usize = 256;
+/// Maximum opaque cursor length.
+pub const FETCH_CURSOR_MAX_BYTES: usize = 2_048;
+/// Maximum number of events one page may request.
+pub const FETCH_PAGE_MAX_EVENTS: u16 = 1_000;
+
/// Heap-backed future returned by transport SPIs.
pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
-/// Bounded source request.
+/// Validated caller identity for one fetch operation.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct FetchRequestId(String);
+
+impl FetchRequestId {
+ /// Parses a non-empty, bounded, printable request identity.
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if value.is_empty() {
+ return Err(Error::EmptyFetchRequestId);
+ }
+ if value.len() > FETCH_REQUEST_ID_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control)
+ {
+ return Err(Error::InvalidFetchRequestId);
+ }
+ Ok(Self(value))
+ }
+
+ /// Returns the validated request identity.
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+impl fmt::Display for FetchRequestId {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(self.as_str())
+ }
+}
+
+/// Opaque adapter-owned continuation token.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct FetchCursor(String);
+
+impl FetchCursor {
+ /// Parses a bounded printable cursor without interpreting its contents.
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if value.is_empty() {
+ return Err(Error::EmptyFetchCursor);
+ }
+ if value.len() > FETCH_CURSOR_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control)
+ {
+ return Err(Error::InvalidFetchCursor);
+ }
+ Ok(Self(value))
+ }
+
+ /// Returns the opaque cursor exactly as supplied by its adapter.
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+impl fmt::Display for FetchCursor {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(self.as_str())
+ }
+}
+
+/// Hard bounds for one source operation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct FetchBounds {
+ limit: u16,
+ deadline_unix_ms: u64,
+}
+
+impl FetchBounds {
+ /// Creates bounds with a non-zero page limit and absolute deadline.
+ pub const fn new(limit: u16, deadline_unix_ms: u64) -> Result<Self, Error> {
+ if limit == 0 || limit > FETCH_PAGE_MAX_EVENTS {
+ return Err(Error::InvalidFetchLimit);
+ }
+ if deadline_unix_ms == 0 {
+ return Err(Error::InvalidFetchDeadline);
+ }
+ Ok(Self {
+ limit,
+ deadline_unix_ms,
+ })
+ }
+
+ /// Maximum number of events the adapter may return.
+ pub const fn limit(self) -> u16 {
+ self.limit
+ }
+
+ /// Absolute Unix deadline in milliseconds.
+ pub const fn deadline_unix_ms(self) -> u64 {
+ self.deadline_unix_ms
+ }
+}
+
+/// Bounded request for one page from one or more transport targets.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct FetchRequest {
+ request_id: FetchRequestId,
+ target_set: TargetSet,
+ bounds: FetchBounds,
+ cursor: Option<FetchCursor>,
+}
+
+impl FetchRequest {
+ /// Creates a first-page request.
+ pub fn new(
+ request_id: impl Into<String>,
+ target_set: TargetSet,
+ bounds: FetchBounds,
+ ) -> Result<Self, Error> {
+ Ok(Self {
+ request_id: FetchRequestId::parse(request_id)?,
+ target_set,
+ bounds,
+ cursor: None,
+ })
+ }
+
+ /// Sets the adapter-owned cursor for a continuation request.
+ #[must_use]
+ pub fn with_cursor(mut self, cursor: FetchCursor) -> Self {
+ self.cursor = Some(cursor);
+ self
+ }
+
+ /// Returns the request identity.
+ pub fn request_id(&self) -> &FetchRequestId {
+ &self.request_id
+ }
+
+ /// Returns the exact requested target set.
+ pub fn target_set(&self) -> &TargetSet {
+ &self.target_set
+ }
+
+ /// Returns the hard operation bounds.
+ pub const fn bounds(&self) -> FetchBounds {
+ self.bounds
+ }
+
+ /// Returns the continuation cursor, when this is not a first-page request.
+ pub fn cursor(&self) -> Option<&FetchCursor> {
+ self.cursor.as_ref()
+ }
+}
+
+/// Transport observation attached to one inbound event.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventProvenance {
+ transport_id: TransportId,
+ target: TargetFingerprint,
+ observed_at_unix_ms: u64,
+ cursor: Option<FetchCursor>,
+}
+
+impl EventProvenance {
+ /// Creates provenance for an event observed from one exact target.
+ pub fn new(
+ transport_id: TransportId,
+ target: TargetFingerprint,
+ observed_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if observed_at_unix_ms == 0 {
+ return Err(Error::InvalidObservedAt);
+ }
+ Ok(Self {
+ transport_id,
+ target,
+ observed_at_unix_ms,
+ cursor: None,
+ })
+ }
+
+ /// Attaches the adapter cursor that located this event.
+ #[must_use]
+ pub fn with_cursor(mut self, cursor: FetchCursor) -> Self {
+ self.cursor = Some(cursor);
+ self
+ }
+
+ /// Returns the transport that produced this observation.
+ pub const fn transport_id(&self) -> TransportId {
+ self.transport_id
+ }
+
+ /// Returns the exact target fingerprint that produced this observation.
+ pub const fn target(&self) -> &TargetFingerprint {
+ &self.target
+ }
+
+ /// Returns the host-recorded observation time.
+ pub const fn observed_at_unix_ms(&self) -> u64 {
+ self.observed_at_unix_ms
+ }
+
+ /// Returns the optional adapter cursor at the observation point.
+ pub const fn cursor(&self) -> Option<&FetchCursor> {
+ self.cursor.as_ref()
+ }
+}
+
+/// ID-checked signed event plus transport provenance.
///
-/// The dedicated bounded-page checkpoint replaces this compatibility alias
-/// with the final request model.
-pub type FetchRequest = RadrootsTransportFetchRequest;
+/// Signature verification, contract validation, canonical admission, storage,
+/// and projection results intentionally remain outside this transport model.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ObservedEvent {
+ event: SignedEvent,
+ provenance: EventProvenance,
+}
+
+impl ObservedEvent {
+ /// Attaches provenance to an inbound signed event.
+ pub const fn new(event: SignedEvent, provenance: EventProvenance) -> Self {
+ Self { event, provenance }
+ }
+
+ /// Returns the unverified signed event payload.
+ pub const fn event(&self) -> &SignedEvent {
+ &self.event
+ }
+
+ /// Returns the transport observation.
+ pub const fn provenance(&self) -> &EventProvenance {
+ &self.provenance
+ }
+}
+
+/// State required to continue or conclude a bounded fetch.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum NextPage {
+ /// Every requested target reached its current end.
+ Complete,
+ /// More results are available from this exact cursor.
+ Cursor(FetchCursor),
+ /// The operation was cancelled and may optionally be resumed.
+ Cancelled { resume_from: Option<FetchCursor> },
+}
+
+/// One validated, request-bound page of inbound observations.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct FetchPage {
+ request_id: FetchRequestId,
+ target_set: TargetSet,
+ limit: u16,
+ events: Vec<ObservedEvent>,
+ target_outcomes: Vec<FetchTargetOutcome>,
+ next_page: NextPage,
+}
+
+impl FetchPage {
+ /// Creates and validates one page against its originating request.
+ pub fn for_request(
+ request: &FetchRequest,
+ events: Vec<ObservedEvent>,
+ target_outcomes: Vec<FetchTargetOutcome>,
+ next_page: NextPage,
+ ) -> Result<Self, Error> {
+ let page = Self {
+ request_id: request.request_id.clone(),
+ target_set: request.target_set.clone(),
+ limit: request.bounds.limit,
+ events,
+ target_outcomes,
+ next_page,
+ };
+ page.validate()?;
+ Ok(page)
+ }
+
+ /// Validates internal cardinality, target provenance, and outcome identity.
+ pub fn validate(&self) -> Result<(), Error> {
+ if self.limit == 0 || self.limit > FETCH_PAGE_MAX_EVENTS {
+ return Err(Error::InvalidFetchLimit);
+ }
+ if self.events.len() > usize::from(self.limit) {
+ return Err(Error::FetchPageLimitExceeded);
+ }
+
+ for observed in &self.events {
+ let provenance = observed.provenance();
+ let Some(target) = self
+ .target_set
+ .targets()
+ .iter()
+ .find(|target| target.fingerprint() == provenance.target())
+ else {
+ return Err(Error::UnexpectedFetchProvenance);
+ };
+ if *target.kind() != provenance.transport_id() {
+ return Err(Error::UnexpectedFetchProvenance);
+ }
+ }
+
+ let requested: BTreeSet<&str> = self
+ .target_set
+ .targets()
+ .iter()
+ .map(|target| target.fingerprint().as_str())
+ .collect();
+ let mut outcomes = BTreeSet::new();
+ for outcome in &self.target_outcomes {
+ if !requested.contains(outcome.target().as_str()) {
+ return Err(Error::UnexpectedFetchTargetOutcome);
+ }
+ if !outcomes.insert(outcome.target().as_str()) {
+ return Err(Error::DuplicateFetchTargetOutcome);
+ }
+ }
+ Ok(())
+ }
+
+ /// Validates that this page is bound to the exact originating request.
+ pub fn validate_for_request(&self, request: &FetchRequest) -> Result<(), Error> {
+ self.validate()?;
+ if &self.request_id != request.request_id()
+ || self.target_set != *request.target_set()
+ || self.limit != request.bounds().limit()
+ {
+ return Err(Error::FetchPageRequestMismatch);
+ }
+ Ok(())
+ }
+
+ /// Returns the request identity.
+ pub const fn request_id(&self) -> &FetchRequestId {
+ &self.request_id
+ }
-/// One bounded page returned by an event source.
-pub type FetchPage = RadrootsTransportFetchReceipt;
+ /// Returns the observations in adapter order.
+ pub fn events(&self) -> &[ObservedEvent] {
+ self.events.as_slice()
+ }
+
+ /// Returns zero or more target-specific outcomes; omitted targets remain unreported.
+ pub fn target_outcomes(&self) -> &[FetchTargetOutcome] {
+ self.target_outcomes.as_slice()
+ }
+
+ /// Returns continuation, completion, or cancellation state.
+ pub const fn next_page(&self) -> &NextPage {
+ &self.next_page
+ }
+}
/// Host SPI for inbound event retrieval.
///
@@ -29,7 +389,7 @@ pub type FetchPage = RadrootsTransportFetchReceipt;
/// Dropping a returned future requests cancellation. If it is dropped before
/// a remote request is published, the implementation must leave no remote
/// operation behind. Once publication may have occurred, cancellation cannot
-/// claim rollback; a later observation may report the remote outcome. An
+/// claim rollback; a later observation may report the remote outcome. The
/// explicit request deadline bounds work independently of future cancellation.
pub trait EventSource: Send + Sync {
/// Returns the source's current runtime status.
@@ -38,3 +398,141 @@ pub trait EventSource: Send + Sync {
/// Fetches one bounded page of transport-neutral events.
fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>>;
}
+
+#[cfg(feature = "serde")]
+mod serde_impl {
+ use super::*;
+
+ impl serde::Serialize for FetchRequestId {
+ fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
+ where
+ S: serde::Serializer,
+ {
+ serializer.serialize_str(self.as_str())
+ }
+ }
+
+ impl<'de> serde::Deserialize<'de> for FetchRequestId {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let value = <String as serde::Deserialize>::deserialize(deserializer)?;
+ Self::parse(value).map_err(serde::de::Error::custom)
+ }
+ }
+
+ impl serde::Serialize for FetchCursor {
+ fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
+ where
+ S: serde::Serializer,
+ {
+ serializer.serialize_str(self.as_str())
+ }
+ }
+
+ impl<'de> serde::Deserialize<'de> for FetchCursor {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let value = <String as serde::Deserialize>::deserialize(deserializer)?;
+ Self::parse(value).map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct FetchBoundsWire {
+ limit: u16,
+ deadline_unix_ms: u64,
+ }
+
+ impl<'de> serde::Deserialize<'de> for FetchBounds {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = FetchBoundsWire::deserialize(deserializer)?;
+ Self::new(wire.limit, wire.deadline_unix_ms).map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct FetchRequestWire {
+ request_id: String,
+ target_set: TargetSet,
+ bounds: FetchBounds,
+ cursor: Option<FetchCursor>,
+ }
+
+ impl<'de> serde::Deserialize<'de> for FetchRequest {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = FetchRequestWire::deserialize(deserializer)?;
+ Self::new(wire.request_id, wire.target_set, wire.bounds)
+ .map(|request| match wire.cursor {
+ Some(cursor) => request.with_cursor(cursor),
+ None => request,
+ })
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct EventProvenanceWire {
+ transport_id: TransportId,
+ target: TargetFingerprint,
+ observed_at_unix_ms: u64,
+ cursor: Option<FetchCursor>,
+ }
+
+ impl<'de> serde::Deserialize<'de> for EventProvenance {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = EventProvenanceWire::deserialize(deserializer)?;
+ Self::new(wire.transport_id, wire.target, wire.observed_at_unix_ms)
+ .map(|provenance| match wire.cursor {
+ Some(cursor) => provenance.with_cursor(cursor),
+ None => provenance,
+ })
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct FetchPageWire {
+ request_id: FetchRequestId,
+ target_set: TargetSet,
+ limit: u16,
+ events: Vec<ObservedEvent>,
+ target_outcomes: Vec<FetchTargetOutcome>,
+ next_page: NextPage,
+ }
+
+ impl<'de> serde::Deserialize<'de> for FetchPage {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = FetchPageWire::deserialize(deserializer)?;
+ let page = Self {
+ request_id: wire.request_id,
+ target_set: wire.target_set,
+ limit: wire.limit,
+ events: wire.events,
+ target_outcomes: wire.target_outcomes,
+ next_page: wire.next_page,
+ };
+ page.validate().map_err(serde::de::Error::custom)?;
+ Ok(page)
+ }
+ }
+}
diff --git a/crates/transport/tests/source_contract.rs b/crates/transport/tests/source_contract.rs
@@ -0,0 +1,292 @@
+use core::{future::Future, pin::Pin, task::Context};
+use std::sync::{
+ Arc,
+ atomic::{AtomicBool, Ordering},
+};
+
+use futures::{future, task::noop_waker_ref};
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
+use radroots_transport::{
+ BoxFuture, Error, EventSource, FetchPage, FetchRequest, SourceStatus, Target, TargetSet,
+ TransportId,
+ capability::{Availability, Maturity, SourceCapabilities},
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ source::{
+ EventProvenance, FETCH_CURSOR_MAX_BYTES, FETCH_PAGE_MAX_EVENTS, FETCH_REQUEST_ID_MAX_BYTES,
+ FetchBounds, FetchCursor, NextPage, ObservedEvent,
+ },
+};
+
+fn target(uri: &str) -> Target {
+ Target::nostr_relay(uri).expect("nostr target")
+}
+
+fn signed_event() -> SignedEvent {
+ let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
+ let wire = Nip01EventWire::parse_json(raw).expect("wire event");
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
+}
+
+fn request(targets: TargetSet, limit: u16) -> FetchRequest {
+ FetchRequest::new(
+ "fetch-request",
+ targets,
+ FetchBounds::new(limit, 1_700_000_100_000).expect("bounds"),
+ )
+ .expect("request")
+}
+
+#[test]
+fn fetch_bounds_request_ids_and_cursors_fail_closed() {
+ assert_eq!(
+ FetchBounds::new(0, 1).expect_err("zero limit"),
+ Error::InvalidFetchLimit
+ );
+ assert_eq!(
+ FetchBounds::new(FETCH_PAGE_MAX_EVENTS + 1, 1).expect_err("oversized limit"),
+ Error::InvalidFetchLimit
+ );
+ assert_eq!(
+ FetchBounds::new(1, 0).expect_err("zero deadline"),
+ Error::InvalidFetchDeadline
+ );
+ assert_eq!(
+ FetchRequest::new(
+ "",
+ TargetSet::new(vec![target("wss://one.example")]).expect("targets"),
+ FetchBounds::new(1, 1).expect("bounds"),
+ )
+ .expect_err("empty request id"),
+ Error::EmptyFetchRequestId
+ );
+ assert_eq!(
+ FetchRequest::new(
+ "x".repeat(FETCH_REQUEST_ID_MAX_BYTES + 1),
+ TargetSet::new(vec![target("wss://one.example")]).expect("targets"),
+ FetchBounds::new(1, 1).expect("bounds"),
+ )
+ .expect_err("oversized request id"),
+ Error::InvalidFetchRequestId
+ );
+ assert_eq!(
+ FetchCursor::parse("").expect_err("empty cursor"),
+ Error::EmptyFetchCursor
+ );
+ assert_eq!(
+ FetchCursor::parse("x".repeat(FETCH_CURSOR_MAX_BYTES + 1)).expect_err("oversized cursor"),
+ Error::InvalidFetchCursor
+ );
+}
+
+#[test]
+fn page_preserves_cursor_provenance_and_partial_target_outcomes() {
+ let first = target("wss://one.example");
+ let second = target("wss://two.example");
+ let targets = TargetSet::new(vec![first.clone(), second.clone()]).expect("targets");
+ let request = request(targets, 2).with_cursor(FetchCursor::parse("page-1").expect("cursor"));
+ let provenance = EventProvenance::new(
+ TransportId::NOSTR,
+ first.fingerprint().clone(),
+ 1_700_000_000_001,
+ )
+ .expect("provenance")
+ .with_cursor(FetchCursor::parse("event-1").expect("event cursor"));
+ let observed = ObservedEvent::new(signed_event(), provenance);
+ let outcomes = vec![
+ FetchTargetOutcome::new(first.fingerprint().clone(), FetchTargetState::Complete),
+ FetchTargetOutcome::new(
+ second.fingerprint().clone(),
+ FetchTargetState::FailedRetryable,
+ )
+ .with_message("relay unavailable"),
+ ];
+ let page = FetchPage::for_request(
+ &request,
+ vec![observed],
+ outcomes,
+ NextPage::Cursor(FetchCursor::parse("page-2").expect("next cursor")),
+ )
+ .expect("page");
+
+ page.validate_for_request(&request)
+ .expect("request binding");
+ assert_eq!(page.events()[0].event().id_str(), signed_event().id_str());
+ assert_eq!(page.events()[0].provenance().target(), first.fingerprint());
+ assert_eq!(page.target_outcomes().len(), 2);
+ assert!(page.target_outcomes()[1].state().is_retryable());
+ assert_eq!(
+ page.target_outcomes()[1].message(),
+ Some("relay unavailable")
+ );
+ assert!(matches!(page.next_page(), NextPage::Cursor(cursor) if cursor.as_str() == "page-2"));
+
+ let encoded = serde_json::to_string(&page).expect("serialize page");
+ assert!(!encoded.contains("admission"));
+ assert!(!encoded.contains("storage"));
+ assert_eq!(
+ serde_json::from_str::<FetchPage>(&encoded).expect("deserialize page"),
+ page
+ );
+ let mut invalid_time = serde_json::to_value(&page).expect("page value");
+ invalid_time["events"][0]["provenance"]["observed_at_unix_ms"] = 0.into();
+ assert!(serde_json::from_value::<FetchPage>(invalid_time).is_err());
+}
+
+#[test]
+fn pages_reject_oversize_unrequested_and_duplicate_evidence() {
+ let requested = target("wss://one.example");
+ let foreign = target("wss://foreign.example");
+ let request = request(TargetSet::new(vec![requested.clone()]).expect("targets"), 1);
+ let observed = ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(TransportId::NOSTR, requested.fingerprint().clone(), 1)
+ .expect("provenance"),
+ );
+ assert_eq!(
+ FetchPage::for_request(
+ &request,
+ vec![observed.clone(), observed],
+ Vec::new(),
+ NextPage::Complete,
+ )
+ .expect_err("oversized page"),
+ Error::FetchPageLimitExceeded
+ );
+
+ let foreign_observation = ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(TransportId::NOSTR, foreign.fingerprint().clone(), 1)
+ .expect("foreign provenance"),
+ );
+ assert_eq!(
+ FetchPage::for_request(
+ &request,
+ vec![foreign_observation],
+ Vec::new(),
+ NextPage::Complete,
+ )
+ .expect_err("foreign provenance"),
+ Error::UnexpectedFetchProvenance
+ );
+
+ let wrong_transport = ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(
+ TransportId::parse("future-mesh").expect("custom transport"),
+ requested.fingerprint().clone(),
+ 1,
+ )
+ .expect("wrong transport provenance"),
+ );
+ assert_eq!(
+ FetchPage::for_request(
+ &request,
+ vec![wrong_transport],
+ Vec::new(),
+ NextPage::Complete,
+ )
+ .expect_err("transport mismatch"),
+ Error::UnexpectedFetchProvenance
+ );
+
+ let duplicate =
+ FetchTargetOutcome::new(requested.fingerprint().clone(), FetchTargetState::Partial);
+ assert_eq!(
+ FetchPage::for_request(
+ &request,
+ Vec::new(),
+ vec![duplicate.clone(), duplicate],
+ NextPage::Cancelled {
+ resume_from: Some(FetchCursor::parse("resume").expect("resume cursor")),
+ },
+ )
+ .expect_err("duplicate outcome"),
+ Error::DuplicateFetchTargetOutcome
+ );
+ assert_eq!(
+ FetchPage::for_request(
+ &request,
+ Vec::new(),
+ vec![FetchTargetOutcome::new(
+ foreign.fingerprint().clone(),
+ FetchTargetState::Unavailable,
+ )],
+ NextPage::Complete,
+ )
+ .expect_err("foreign outcome"),
+ Error::UnexpectedFetchTargetOutcome
+ );
+
+ let page = FetchPage::for_request(&request, Vec::new(), Vec::new(), NextPage::Complete)
+ .expect("empty page");
+ let other_request = FetchRequest::new(
+ "other-request",
+ request.target_set().clone(),
+ request.bounds(),
+ )
+ .expect("other request");
+ assert_eq!(
+ page.validate_for_request(&other_request)
+ .expect_err("request mismatch"),
+ Error::FetchPageRequestMismatch
+ );
+}
+
+struct CancellationSource {
+ published: Arc<AtomicBool>,
+ cancelled_after_publish: Arc<AtomicBool>,
+}
+
+struct PublicationGuard(Arc<AtomicBool>);
+
+impl Drop for PublicationGuard {
+ fn drop(&mut self) {
+ self.0.store(true, Ordering::SeqCst);
+ }
+}
+
+impl EventSource for CancellationSource {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> {
+ Box::pin(async {
+ Ok(SourceStatus::new(
+ TransportId::NOSTR,
+ true,
+ Maturity::Stable,
+ Availability::Available,
+ SourceCapabilities::FETCH,
+ "ready",
+ ))
+ })
+ }
+
+ fn fetch(&self, _request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> {
+ let published = Arc::clone(&self.published);
+ let cancelled = Arc::clone(&self.cancelled_after_publish);
+ Box::pin(async move {
+ published.store(true, Ordering::SeqCst);
+ let _guard = PublicationGuard(cancelled);
+ future::pending::<Result<FetchPage, Error>>().await
+ })
+ }
+}
+
+#[test]
+fn dropping_fetch_futures_respects_before_and_after_publication_boundaries() {
+ let source = CancellationSource {
+ published: Arc::new(AtomicBool::new(false)),
+ cancelled_after_publish: Arc::new(AtomicBool::new(false)),
+ };
+ let targets = TargetSet::new(vec![target("wss://one.example")]).expect("targets");
+
+ let unpolled = source.fetch(request(targets.clone(), 1));
+ drop(unpolled);
+ assert!(!source.published.load(Ordering::SeqCst));
+ assert!(!source.cancelled_after_publish.load(Ordering::SeqCst));
+
+ let mut published = source.fetch(request(targets, 1));
+ let mut context = Context::from_waker(noop_waker_ref());
+ assert!(Pin::new(&mut published).poll(&mut context).is_pending());
+ assert!(source.published.load(Ordering::SeqCst));
+ drop(published);
+ assert!(source.cancelled_after_publish.load(Ordering::SeqCst));
+}
diff --git a/crates/transport/tests/spi.rs b/crates/transport/tests/spi.rs
@@ -5,6 +5,7 @@ use radroots_transport::{
RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetReceipt,
RadrootsTransportTargetSet, SinkStatus, SourceStatus, TransportId,
capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
+ source::{FetchBounds, NextPage},
};
struct SourceOnly;
@@ -18,7 +19,9 @@ impl EventSource for SourceOnly {
&self,
request: FetchRequest,
) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
- Box::pin(async move { Ok(FetchPage::new(request.request_id, Vec::new(), 0)) })
+ Box::pin(async move {
+ FetchPage::for_request(&request, Vec::new(), Vec::new(), NextPage::Complete)
+ })
}
}
@@ -123,9 +126,14 @@ fn source_only_and_sink_only_implementations_are_independently_dispatchable() {
let source_status = block_on(EventSource::status(&source)).expect("source status");
assert!(source_status.capabilities().can_fetch());
- let page =
- block_on(source.fetch(FetchRequest::new("fetch-1", target_set()))).expect("fetch page");
- assert_eq!(page.request_id, "fetch-1");
+ let request = FetchRequest::new(
+ "fetch-1",
+ target_set(),
+ FetchBounds::new(10, 1_700_000_000_000).expect("fetch bounds"),
+ )
+ .expect("fetch request");
+ let page = block_on(source.fetch(request)).expect("fetch page");
+ assert_eq!(page.request_id().as_str(), "fetch-1");
let sink_status = block_on(EventSink::status(&sink)).expect("sink status");
assert!(sink_status.capabilities().can_deliver());
diff --git a/crates/transport/tests/transport.rs b/crates/transport/tests/transport.rs
@@ -2086,6 +2086,18 @@ fn target_contract_covers_parser_and_authority_boundaries() {
fn every_transport_error_has_a_stable_display_message() {
let remaining = [
RadrootsTransportError::RequiredTargetNotRequested,
+ RadrootsTransportError::EmptyFetchRequestId,
+ RadrootsTransportError::InvalidFetchRequestId,
+ RadrootsTransportError::InvalidFetchLimit,
+ RadrootsTransportError::InvalidFetchDeadline,
+ RadrootsTransportError::EmptyFetchCursor,
+ RadrootsTransportError::InvalidFetchCursor,
+ RadrootsTransportError::InvalidObservedAt,
+ RadrootsTransportError::UnexpectedFetchProvenance,
+ RadrootsTransportError::UnexpectedFetchTargetOutcome,
+ RadrootsTransportError::DuplicateFetchTargetOutcome,
+ RadrootsTransportError::FetchPageLimitExceeded,
+ RadrootsTransportError::FetchPageRequestMismatch,
RadrootsTransportError::EmptyDeliveryRequestId,
RadrootsTransportError::InvalidDeliveryRequestId,
RadrootsTransportError::InvalidDeliveryTimestamp,
diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs
@@ -790,6 +790,18 @@ fn transport_error_to_relay_error(error: RadrootsTransportError) -> RadrootsRela
| RadrootsTransportError::TargetSetTooLarge
| RadrootsTransportError::DuplicateTargetFingerprint
| RadrootsTransportError::InvalidTargetFingerprint
+ | RadrootsTransportError::EmptyFetchRequestId
+ | RadrootsTransportError::InvalidFetchRequestId
+ | RadrootsTransportError::InvalidFetchLimit
+ | RadrootsTransportError::InvalidFetchDeadline
+ | RadrootsTransportError::EmptyFetchCursor
+ | RadrootsTransportError::InvalidFetchCursor
+ | RadrootsTransportError::InvalidObservedAt
+ | RadrootsTransportError::UnexpectedFetchProvenance
+ | RadrootsTransportError::UnexpectedFetchTargetOutcome
+ | RadrootsTransportError::DuplicateFetchTargetOutcome
+ | RadrootsTransportError::FetchPageLimitExceeded
+ | RadrootsTransportError::FetchPageRequestMismatch
| RadrootsTransportError::UnexpectedDeliveryTargetReceipt
| RadrootsTransportError::DuplicateDeliveryTargetReceipt
| RadrootsTransportError::MissingDeliveryTargetReceipt