commit ed3016d3620d98cfde1b53901e12eebdd73f6b85
parent ec20c8477848b0d111f30b286c572cfc6b83ef7f
Author: triesap <tyson@radroots.org>
Date: Sun, 16 Aug 2026 21:01:26 +0000
transport: define bounded live subscriptions
- Intent: add a runtime-neutral live-subscription SPI with sealed request, event, checkpoint, and terminal models.
- Why: services require explicit bounded cancellation, deadline, provenance, and reconnect semantics without changing existing fetch producers.
- Verification: extbuild package tests, no-default checks, Clippy, Rustdoc, workspace check, doctest, and API baseline comparison passed.
- Risk: public alpha API expands, while concrete adapter execution and consumer adoption remain deferred to later checkpoints.
Diffstat:
9 files changed, 1259 insertions(+), 30 deletions(-)
diff --git a/contracts/api_baselines/radroots_transport.txt b/contracts/api_baselines/radroots_transport.txt
@@ -99,6 +99,7 @@ pub radroots_transport::error::Error::DuplicateFetchAuthor
pub radroots_transport::error::Error::DuplicateFetchKind
pub radroots_transport::error::Error::DuplicateFetchTargetOutcome
pub radroots_transport::error::Error::DuplicateRequiredTargetFingerprint
+pub radroots_transport::error::Error::DuplicateSubscriptionCheckpoint
pub radroots_transport::error::Error::DuplicateTargetFingerprint
pub radroots_transport::error::Error::EmptyDeliveryRequestId
pub radroots_transport::error::Error::EmptyFetchCursor
@@ -107,6 +108,7 @@ pub radroots_transport::error::Error::EmptyPayloadBytes
pub radroots_transport::error::Error::EmptyPayloadId
pub radroots_transport::error::Error::EmptyPayloadLabel
pub radroots_transport::error::Error::EmptyRequiredTargetSet
+pub radroots_transport::error::Error::EmptySubscriptionRequestId
pub radroots_transport::error::Error::EmptyTargetLabel
pub radroots_transport::error::Error::EmptyTargetScope
pub radroots_transport::error::Error::EmptyTargetSet
@@ -130,6 +132,10 @@ pub radroots_transport::error::Error::InvalidPayloadDigest
pub radroots_transport::error::Error::InvalidPayloadId
pub radroots_transport::error::Error::InvalidPayloadLabel
pub radroots_transport::error::Error::InvalidSatisfactionPolicy
+pub radroots_transport::error::Error::InvalidSubscriptionDeadline
+pub radroots_transport::error::Error::InvalidSubscriptionEnd
+pub radroots_transport::error::Error::InvalidSubscriptionLimit
+pub radroots_transport::error::Error::InvalidSubscriptionRequestId
pub radroots_transport::error::Error::InvalidTargetFingerprint
pub radroots_transport::error::Error::InvalidTargetLabel
pub radroots_transport::error::Error::InvalidTargetScope
@@ -138,12 +144,18 @@ pub radroots_transport::error::Error::InvalidTransportKind
pub radroots_transport::error::Error::MissingDeliveryTargetReceipt
pub radroots_transport::error::Error::PayloadDigestMismatch
pub radroots_transport::error::Error::RequiredTargetNotRequested
+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::TargetSetTooLarge
pub radroots_transport::error::Error::TransportOutcomeStatusMismatch
pub radroots_transport::error::Error::UnexpectedDeliveryTargetReceipt
pub radroots_transport::error::Error::UnexpectedFetchEvent
pub radroots_transport::error::Error::UnexpectedFetchProvenance
pub radroots_transport::error::Error::UnexpectedFetchTargetOutcome
+pub radroots_transport::error::Error::UnexpectedSubscriptionCheckpoint
+pub radroots_transport::error::Error::UnexpectedSubscriptionEvent
pub radroots_transport::error::Error::UnsupportedOperation
impl core::fmt::Display for radroots_transport::error::Error
pub fn radroots_transport::error::Error::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
@@ -295,6 +307,14 @@ pub radroots_transport::source::NextPage::Cancelled
pub radroots_transport::source::NextPage::Cancelled::resume_from: core::option::Option<radroots_transport::source::FetchCursor>
pub radroots_transport::source::NextPage::Complete
pub radroots_transport::source::NextPage::Cursor(radroots_transport::source::FetchCursor)
+pub enum radroots_transport::source::SubscriptionEndReason
+pub radroots_transport::source::SubscriptionEndReason::Cancelled
+pub radroots_transport::source::SubscriptionEndReason::Deadline
+pub radroots_transport::source::SubscriptionEndReason::EventLimit
+pub radroots_transport::source::SubscriptionEndReason::SourceClosed
+pub enum radroots_transport::source::SubscriptionNext
+pub radroots_transport::source::SubscriptionNext::End(radroots_transport::source::SubscriptionEnd)
+pub radroots_transport::source::SubscriptionNext::Event(alloc::boxed::Box<radroots_transport::source::SubscriptionEvent>)
pub struct radroots_transport::source::EventProvenance
impl radroots_transport::source::EventProvenance
pub const fn radroots_transport::source::EventProvenance::cursor(&self) -> core::option::Option<&radroots_transport::source::FetchCursor>
@@ -384,15 +404,78 @@ pub const fn radroots_transport::SourceStatus::maturity(&self) -> radroots_trans
pub fn radroots_transport::SourceStatus::message(&self) -> &str
pub fn radroots_transport::SourceStatus::new(radroots_transport::TransportId, bool, radroots_transport::capability::Maturity, radroots_transport::capability::Availability, radroots_transport::capability::SourceCapabilities, impl core::convert::Into<alloc::string::String>) -> Self
pub const fn radroots_transport::SourceStatus::transport_id(&self) -> radroots_transport::TransportId
+pub struct radroots_transport::source::SubscriptionBounds
+impl radroots_transport::source::SubscriptionBounds
+pub const fn radroots_transport::source::SubscriptionBounds::deadline_unix_ms(self) -> u64
+pub const fn radroots_transport::source::SubscriptionBounds::event_limit(self) -> u16
+pub const fn radroots_transport::source::SubscriptionBounds::new(u16, u64) -> core::result::Result<Self, radroots_transport::error::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionBounds
+pub fn radroots_transport::source::SubscriptionBounds::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::source::SubscriptionCheckpoint
+impl radroots_transport::source::SubscriptionCheckpoint
+pub const fn radroots_transport::source::SubscriptionCheckpoint::cursor(&self) -> &radroots_transport::source::FetchCursor
+pub const fn radroots_transport::source::SubscriptionCheckpoint::new(radroots_transport::target::TargetFingerprint, radroots_transport::source::FetchCursor) -> Self
+pub const fn radroots_transport::source::SubscriptionCheckpoint::target(&self) -> &radroots_transport::target::TargetFingerprint
+pub struct radroots_transport::source::SubscriptionEnd
+impl radroots_transport::source::SubscriptionEnd
+pub fn radroots_transport::source::SubscriptionEnd::checkpoints(&self) -> &[radroots_transport::source::SubscriptionCheckpoint]
+pub const fn radroots_transport::source::SubscriptionEnd::event_count(&self) -> u16
+pub fn radroots_transport::source::SubscriptionEnd::for_request<I>(&radroots_transport::source::SubscriptionRequest, u16, I, radroots_transport::source::SubscriptionEndReason) -> core::result::Result<Self, radroots_transport::error::Error> where I: core::iter::traits::collect::IntoIterator<Item = radroots_transport::source::SubscriptionCheckpoint>
+pub const fn radroots_transport::source::SubscriptionEnd::reason(&self) -> radroots_transport::source::SubscriptionEndReason
+pub const fn radroots_transport::source::SubscriptionEnd::request(&self) -> &radroots_transport::source::SubscriptionRequest
+pub fn radroots_transport::source::SubscriptionEnd::validate_for_request(&self, &radroots_transport::source::SubscriptionRequest) -> core::result::Result<(), radroots_transport::error::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionEnd
+pub fn radroots_transport::source::SubscriptionEnd::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::source::SubscriptionEvent
+impl radroots_transport::source::SubscriptionEvent
+pub const fn radroots_transport::source::SubscriptionEvent::checkpoint(&self) -> &radroots_transport::source::SubscriptionCheckpoint
+pub fn radroots_transport::source::SubscriptionEvent::for_request(&radroots_transport::source::SubscriptionRequest, radroots_transport::source::ObservedEvent, radroots_transport::source::SubscriptionCheckpoint) -> core::result::Result<Self, radroots_transport::error::Error>
+pub const fn radroots_transport::source::SubscriptionEvent::observed(&self) -> &radroots_transport::source::ObservedEvent
+pub const fn radroots_transport::source::SubscriptionEvent::request(&self) -> &radroots_transport::source::SubscriptionRequest
+pub const fn radroots_transport::source::SubscriptionEvent::request_id(&self) -> &radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionEvent::validate_for_request(&self, &radroots_transport::source::SubscriptionRequest) -> core::result::Result<(), radroots_transport::error::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionEvent
+pub fn radroots_transport::source::SubscriptionEvent::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::source::SubscriptionRequest
+impl radroots_transport::source::SubscriptionRequest
+pub const fn radroots_transport::source::SubscriptionRequest::bounds(&self) -> radroots_transport::source::SubscriptionBounds
+pub fn radroots_transport::source::SubscriptionRequest::checkpoints(&self) -> &[radroots_transport::source::SubscriptionCheckpoint]
+pub fn radroots_transport::source::SubscriptionRequest::new(impl core::convert::AsRef<str>, radroots_transport::target::TargetSet, radroots_transport::source::SubscriptionBounds) -> core::result::Result<Self, radroots_transport::error::Error>
+pub const fn radroots_transport::source::SubscriptionRequest::request_id(&self) -> &radroots_transport::source::SubscriptionRequestId
+pub const fn radroots_transport::source::SubscriptionRequest::selector(&self) -> &radroots_transport::source::FetchSelector
+pub const fn radroots_transport::source::SubscriptionRequest::target_set(&self) -> &radroots_transport::target::TargetSet
+pub fn radroots_transport::source::SubscriptionRequest::with_checkpoints<I>(self, I) -> core::result::Result<Self, radroots_transport::error::Error> where I: core::iter::traits::collect::IntoIterator<Item = radroots_transport::source::SubscriptionCheckpoint>
+pub fn radroots_transport::source::SubscriptionRequest::with_selector(self, radroots_transport::source::FetchSelector) -> Self
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionRequest
+pub fn radroots_transport::source::SubscriptionRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::source::SubscriptionRequestId(_)
+impl radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionRequestId::as_str(&self) -> &str
+pub fn radroots_transport::source::SubscriptionRequestId::parse(impl core::convert::AsRef<str>) -> core::result::Result<Self, radroots_transport::error::Error>
+impl core::fmt::Display for radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionRequestId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl serde_core::ser::Serialize for radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionRequestId::serialize<S>(&self, S) -> core::result::Result<<S as serde_core::ser::Serializer>::Ok, <S as serde_core::ser::Serializer>::Error> where S: serde_core::ser::Serializer
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionRequestId::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
pub const radroots_transport::source::FETCH_CURSOR_MAX_BYTES: usize
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::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
pub fn radroots_transport::source::EventSource::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::source::EventSource::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::SourceStatus, radroots_transport::error::Error>>
+pub trait radroots_transport::source::EventSubscriber: core::marker::Send + core::marker::Sync
+pub fn radroots_transport::source::EventSubscriber::subscribe(&self, radroots_transport::source::SubscriptionRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::BoxSubscription, radroots_transport::error::Error>>
+pub trait radroots_transport::source::EventSubscription: core::marker::Send
+pub fn radroots_transport::source::EventSubscription::cancel(&mut self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::SubscriptionEnd, radroots_transport::error::Error>>
+pub fn radroots_transport::source::EventSubscription::next(&mut self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::SubscriptionNext, radroots_transport::error::Error>>
+pub fn radroots_transport::source::EventSubscription::request(&self) -> &radroots_transport::source::SubscriptionRequest
pub type radroots_transport::source::BoxFuture<'a, T> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = T> + core::marker::Send + 'a)>>
+pub type radroots_transport::source::BoxSubscription = alloc::boxed::Box<dyn radroots_transport::source::EventSubscription>
pub mod radroots_transport::target
pub struct radroots_transport::target::EndpointUri(_)
impl radroots_transport::target::EndpointUri
@@ -494,6 +577,7 @@ pub radroots_transport::Error::DuplicateFetchAuthor
pub radroots_transport::Error::DuplicateFetchKind
pub radroots_transport::Error::DuplicateFetchTargetOutcome
pub radroots_transport::Error::DuplicateRequiredTargetFingerprint
+pub radroots_transport::Error::DuplicateSubscriptionCheckpoint
pub radroots_transport::Error::DuplicateTargetFingerprint
pub radroots_transport::Error::EmptyDeliveryRequestId
pub radroots_transport::Error::EmptyFetchCursor
@@ -502,6 +586,7 @@ pub radroots_transport::Error::EmptyPayloadBytes
pub radroots_transport::Error::EmptyPayloadId
pub radroots_transport::Error::EmptyPayloadLabel
pub radroots_transport::Error::EmptyRequiredTargetSet
+pub radroots_transport::Error::EmptySubscriptionRequestId
pub radroots_transport::Error::EmptyTargetLabel
pub radroots_transport::Error::EmptyTargetScope
pub radroots_transport::Error::EmptyTargetSet
@@ -525,6 +610,10 @@ pub radroots_transport::Error::InvalidPayloadDigest
pub radroots_transport::Error::InvalidPayloadId
pub radroots_transport::Error::InvalidPayloadLabel
pub radroots_transport::Error::InvalidSatisfactionPolicy
+pub radroots_transport::Error::InvalidSubscriptionDeadline
+pub radroots_transport::Error::InvalidSubscriptionEnd
+pub radroots_transport::Error::InvalidSubscriptionLimit
+pub radroots_transport::Error::InvalidSubscriptionRequestId
pub radroots_transport::Error::InvalidTargetFingerprint
pub radroots_transport::Error::InvalidTargetLabel
pub radroots_transport::Error::InvalidTargetScope
@@ -533,15 +622,29 @@ pub radroots_transport::Error::InvalidTransportKind
pub radroots_transport::Error::MissingDeliveryTargetReceipt
pub radroots_transport::Error::PayloadDigestMismatch
pub radroots_transport::Error::RequiredTargetNotRequested
+pub radroots_transport::Error::SubscriptionCheckpointSetTooLarge
+pub radroots_transport::Error::SubscriptionEndLimitExceeded
+pub radroots_transport::Error::SubscriptionEndRequestMismatch
+pub radroots_transport::Error::SubscriptionEventCheckpointMismatch
pub radroots_transport::Error::TargetSetTooLarge
pub radroots_transport::Error::TransportOutcomeStatusMismatch
pub radroots_transport::Error::UnexpectedDeliveryTargetReceipt
pub radroots_transport::Error::UnexpectedFetchEvent
pub radroots_transport::Error::UnexpectedFetchProvenance
pub radroots_transport::Error::UnexpectedFetchTargetOutcome
+pub radroots_transport::Error::UnexpectedSubscriptionCheckpoint
+pub radroots_transport::Error::UnexpectedSubscriptionEvent
pub radroots_transport::Error::UnsupportedOperation
impl core::fmt::Display for radroots_transport::error::Error
pub fn radroots_transport::error::Error::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub enum radroots_transport::SubscriptionEndReason
+pub radroots_transport::SubscriptionEndReason::Cancelled
+pub radroots_transport::SubscriptionEndReason::Deadline
+pub radroots_transport::SubscriptionEndReason::EventLimit
+pub radroots_transport::SubscriptionEndReason::SourceClosed
+pub enum radroots_transport::SubscriptionNext
+pub radroots_transport::SubscriptionNext::End(radroots_transport::source::SubscriptionEnd)
+pub radroots_transport::SubscriptionNext::Event(alloc::boxed::Box<radroots_transport::source::SubscriptionEvent>)
pub struct radroots_transport::DeliveryReceipt
impl radroots_transport::sink::DeliveryReceipt
pub fn radroots_transport::sink::DeliveryReceipt::for_request(&radroots_transport::sink::DeliveryRequest, alloc::vec::Vec<radroots_transport::sink::DeliveryTargetReceipt>) -> core::result::Result<Self, radroots_transport::error::Error>
@@ -617,6 +720,38 @@ pub const fn radroots_transport::SourceStatus::maturity(&self) -> radroots_trans
pub fn radroots_transport::SourceStatus::message(&self) -> &str
pub fn radroots_transport::SourceStatus::new(radroots_transport::TransportId, bool, radroots_transport::capability::Maturity, radroots_transport::capability::Availability, radroots_transport::capability::SourceCapabilities, impl core::convert::Into<alloc::string::String>) -> Self
pub const fn radroots_transport::SourceStatus::transport_id(&self) -> radroots_transport::TransportId
+pub struct radroots_transport::SubscriptionEnd
+impl radroots_transport::source::SubscriptionEnd
+pub fn radroots_transport::source::SubscriptionEnd::checkpoints(&self) -> &[radroots_transport::source::SubscriptionCheckpoint]
+pub const fn radroots_transport::source::SubscriptionEnd::event_count(&self) -> u16
+pub fn radroots_transport::source::SubscriptionEnd::for_request<I>(&radroots_transport::source::SubscriptionRequest, u16, I, radroots_transport::source::SubscriptionEndReason) -> core::result::Result<Self, radroots_transport::error::Error> where I: core::iter::traits::collect::IntoIterator<Item = radroots_transport::source::SubscriptionCheckpoint>
+pub const fn radroots_transport::source::SubscriptionEnd::reason(&self) -> radroots_transport::source::SubscriptionEndReason
+pub const fn radroots_transport::source::SubscriptionEnd::request(&self) -> &radroots_transport::source::SubscriptionRequest
+pub fn radroots_transport::source::SubscriptionEnd::validate_for_request(&self, &radroots_transport::source::SubscriptionRequest) -> core::result::Result<(), radroots_transport::error::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionEnd
+pub fn radroots_transport::source::SubscriptionEnd::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::SubscriptionEvent
+impl radroots_transport::source::SubscriptionEvent
+pub const fn radroots_transport::source::SubscriptionEvent::checkpoint(&self) -> &radroots_transport::source::SubscriptionCheckpoint
+pub fn radroots_transport::source::SubscriptionEvent::for_request(&radroots_transport::source::SubscriptionRequest, radroots_transport::source::ObservedEvent, radroots_transport::source::SubscriptionCheckpoint) -> core::result::Result<Self, radroots_transport::error::Error>
+pub const fn radroots_transport::source::SubscriptionEvent::observed(&self) -> &radroots_transport::source::ObservedEvent
+pub const fn radroots_transport::source::SubscriptionEvent::request(&self) -> &radroots_transport::source::SubscriptionRequest
+pub const fn radroots_transport::source::SubscriptionEvent::request_id(&self) -> &radroots_transport::source::SubscriptionRequestId
+pub fn radroots_transport::source::SubscriptionEvent::validate_for_request(&self, &radroots_transport::source::SubscriptionRequest) -> core::result::Result<(), radroots_transport::error::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionEvent
+pub fn radroots_transport::source::SubscriptionEvent::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_transport::SubscriptionRequest
+impl radroots_transport::source::SubscriptionRequest
+pub const fn radroots_transport::source::SubscriptionRequest::bounds(&self) -> radroots_transport::source::SubscriptionBounds
+pub fn radroots_transport::source::SubscriptionRequest::checkpoints(&self) -> &[radroots_transport::source::SubscriptionCheckpoint]
+pub fn radroots_transport::source::SubscriptionRequest::new(impl core::convert::AsRef<str>, radroots_transport::target::TargetSet, radroots_transport::source::SubscriptionBounds) -> core::result::Result<Self, radroots_transport::error::Error>
+pub const fn radroots_transport::source::SubscriptionRequest::request_id(&self) -> &radroots_transport::source::SubscriptionRequestId
+pub const fn radroots_transport::source::SubscriptionRequest::selector(&self) -> &radroots_transport::source::FetchSelector
+pub const fn radroots_transport::source::SubscriptionRequest::target_set(&self) -> &radroots_transport::target::TargetSet
+pub fn radroots_transport::source::SubscriptionRequest::with_checkpoints<I>(self, I) -> core::result::Result<Self, radroots_transport::error::Error> where I: core::iter::traits::collect::IntoIterator<Item = radroots_transport::source::SubscriptionCheckpoint>
+pub fn radroots_transport::source::SubscriptionRequest::with_selector(self, radroots_transport::source::FetchSelector) -> Self
+impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::source::SubscriptionRequest
+pub fn radroots_transport::source::SubscriptionRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
pub struct radroots_transport::Target
impl radroots_transport::target::Target
pub fn radroots_transport::target::Target::fingerprint(&self) -> &radroots_transport::target::TargetFingerprint
@@ -680,4 +815,11 @@ pub fn radroots_transport::EventSink::status(&self) -> radroots_transport::sourc
pub trait radroots_transport::EventSource: core::marker::Send + core::marker::Sync
pub fn radroots_transport::EventSource::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::EventSource::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::SourceStatus, radroots_transport::error::Error>>
+pub trait radroots_transport::EventSubscriber: core::marker::Send + core::marker::Sync
+pub fn radroots_transport::EventSubscriber::subscribe(&self, radroots_transport::source::SubscriptionRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::BoxSubscription, radroots_transport::error::Error>>
+pub trait radroots_transport::EventSubscription: core::marker::Send
+pub fn radroots_transport::EventSubscription::cancel(&mut self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::SubscriptionEnd, radroots_transport::error::Error>>
+pub fn radroots_transport::EventSubscription::next(&mut self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::SubscriptionNext, radroots_transport::error::Error>>
+pub fn radroots_transport::EventSubscription::request(&self) -> &radroots_transport::source::SubscriptionRequest
pub type radroots_transport::BoxFuture<'a, T> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = T> + core::marker::Send + 'a)>>
+pub type radroots_transport::BoxSubscription = alloc::boxed::Box<dyn radroots_transport::source::EventSubscription>
diff --git a/crates/transport/README.md b/crates/transport/README.md
@@ -3,7 +3,8 @@
`radroots_transport` is the transport-neutral host SPI for moving verified
Radroots events between explicit targets. It owns extensible transport and
target identities, separate source and sink capabilities, bounded fetch and
-delivery requests, provenance, partial outcomes, and normalized errors.
+delivery requests, bounded live subscriptions, provenance, partial outcomes,
+and normalized errors.
The crate is `no_std` with `alloc`. It does not own network clients, storage,
outbox claiming, retry loops, schedulers, timers, or fallback policy. Concrete
@@ -18,11 +19,12 @@ The authoritative package charter is the
1. A host selects a validated [`TransportId`] and constructs one or more
canonical [`Target`] values.
2. [`TargetSet`] rejects an empty, oversized, or duplicate-fingerprint set.
-3. The caller creates a bounded [`FetchRequest`] or [`DeliveryRequest`] with a
- request identity and absolute deadline. A fetch may carry a validated
- [`FetchSelector`] for exact kinds, authors, and inclusive event-time bounds.
-4. A dyn-compatible [`EventSource`] or [`EventSink`] implementation performs
- only the requested operation.
+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.
+4. A dyn-compatible [`EventSource`], [`EventSubscriber`], or [`EventSink`]
+ implementation performs only the requested operation.
5. The caller validates [`FetchPage`] or [`DeliveryReceipt`] against the
originating request and decides whether normalized retryable outcomes merit
another explicit operation.
@@ -31,9 +33,11 @@ The authoritative package charter is the
[`Target`]: crate::Target
[`TargetSet`]: crate::TargetSet
[`FetchRequest`]: crate::FetchRequest
+[`SubscriptionRequest`]: crate::SubscriptionRequest
[`FetchSelector`]: crate::source::FetchSelector
[`DeliveryRequest`]: crate::DeliveryRequest
[`EventSource`]: crate::EventSource
+[`EventSubscriber`]: crate::EventSubscriber
[`EventSink`]: crate::EventSink
[`FetchPage`]: crate::FetchPage
[`DeliveryReceipt`]: crate::DeliveryReceipt
@@ -54,15 +58,18 @@ The complete externally implementable SPI example is
## Host SPI contract
-`EventSource` and `EventSink` are independent, externally implementable,
-dyn-compatible `Send + Sync` traits. An adapter may implement either or both.
-Their methods return boxed `Future + Send` values so this package does not
-select an async runtime or require an async-trait macro.
+`EventSource`, `EventSubscriber`, and `EventSink` are independent, externally
+implementable, dyn-compatible `Send + Sync` traits. An adapter may implement
+any subset. Their methods return boxed `Future + Send` values so this package
+does not select an async runtime or require an async-trait macro.
- `status` is observational and must not begin fetch or delivery work.
- `fetch` returns one bounded page plus per-target outcomes and explicit
continuation state. Returned events must satisfy the request selector;
request-bound page validation rejects adapter drift.
+- `subscribe` returns a sealed live-operation capability. Its `next` method
+ returns request-bound events with target checkpoints or one stable terminal
+ result; its `cancel` method returns that same terminal result once ended.
- `deliver` returns one receipt for the exact request and every requested
target; it never performs an implicit retry.
- Implementations must not install an executor or spawn hidden workers. The
@@ -89,10 +96,10 @@ private-network exceptions.
## Bounds, deadlines, cancellation, and commit points
-Request IDs, endpoint values, cursors, target sets, outcome details, and page
-sizes are bounded by public constants. Fetch and delivery requests carry an
-absolute Unix-millisecond deadline; constructors validate the deadline but do
-not read a clock.
+Request IDs, endpoint values, cursors, target sets, outcome details, page
+sizes, and live event counts are bounded by public constants. Fetch,
+subscription, and delivery requests carry an absolute Unix-millisecond
+deadline; constructors validate the deadline but do not read a clock.
Dropping a returned future requests cancellation. Before an adapter publishes
or commits a remote operation, cancellation must leave no durable effect.
@@ -100,6 +107,12 @@ After that boundary, cancellation does not imply rollback: the adapter must
preserve and report the final remote state when it can be observed. Each
concrete adapter documents its exact publication or commit point.
+Live subscriptions additionally retain canonical per-target checkpoints in
+request target order. Reconnect callers pass only an explicit checkpoint
+subset; adapters must not invent hidden replay state. Once a subscription
+returns its terminal result, subsequent `next` and `cancel` calls return the
+same result.
+
## Outcomes, partial success, and retry
Fetch pages and delivery receipts preserve target-local results. Delivery
@@ -111,14 +124,16 @@ state must not be rewritten as success or silent fallback.
## Serialization and provenance
-With `serde`, passive identities, requests, pages, receipts, status values, and
-outcomes use validated representations. Deserialization rechecks canonical
-identity, bounds, request/receipt cardinality, and provenance rather than
-trusting serialized fingerprints or counts.
+With `serde`, passive identities, requests, subscriptions, pages, receipts,
+status values, and outcomes use validated representations. Deserialization
+rechecks canonical identity, bounds, request/result cardinality, and provenance
+rather than trusting serialized fingerprints or counts.
Inbound events carry the transport ID, exact target fingerprint, observation
time, and optional continuation cursor that produced them. A page cannot claim
-events or outcomes for targets outside its request.
+events or outcomes for targets outside its request. A live event additionally
+binds its exact originating request and requires its resulting target
+checkpoint to match the event's provenance cursor.
## Security and side effects
@@ -126,8 +141,8 @@ events or outcomes for targets outside its request.
adapters remain responsible for scheme, DNS/IP, TLS, and private-network
policy before connection.
- Delivery accepts a validated signed event, not arbitrary bytes.
-- Receipt and page constructors reject missing, duplicate, unexpected, or
- forged target evidence.
+- Receipt, page, and subscription constructors reject missing, duplicate,
+ unexpected, or forged target evidence.
- This crate forbids unsafe code and contains no socket, filesystem, database,
global state, timer, executor, or retry implementation.
- No transport may silently route through another transport.
@@ -137,7 +152,7 @@ events or outcomes for targets outside its request.
| Feature | Default | Contract |
| --- | --- | --- |
| `std` | yes | Enables standard-library support in canonical dependency values; it adds no I/O, runtime, or global initialization. |
-| `serde` | yes | Adds validated serialization for passive identities, requests, status, provenance, pages, receipts, policies, and outcomes. |
+| `serde` | yes | Adds validated serialization for passive identities, requests, subscriptions, status, provenance, pages, receipts, policies, and outcomes. |
Features are additive. `--no-default-features` provides the `no_std + alloc`
core, and `serde` is supported independently of `std`.
@@ -148,7 +163,7 @@ core, and `serde` is supported independently of `std`.
- `radroots_storage` persists transport-neutral outbox and evidence values.
- `radroots_sync` plans bounded operations and applies explicit retry policy.
- `radroots_sdk` composes storage, signing, source, and sink implementations.
-- Service and future adapter hosts implement one or both SPIs and own runtime,
+- Service and future adapter hosts implement any of the SPIs and own runtime,
network, cancellation, and lifecycle behavior.
Applications that only need ordinary Radroots operations should normally use
diff --git a/crates/transport/examples/host_transport.rs b/crates/transport/examples/host_transport.rs
@@ -1,6 +1,6 @@
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource, FetchPage, FetchRequest,
- TransportId,
+ BoxSubscription, DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource,
+ EventSubscriber, FetchPage, FetchRequest, SubscriptionRequest, TransportId,
capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
sink::SinkStatus,
source::{BoxFuture, SourceStatus},
@@ -27,6 +27,15 @@ impl EventSource for HostTransport {
}
}
+impl EventSubscriber for HostTransport {
+ fn subscribe(
+ &self,
+ _request: SubscriptionRequest,
+ ) -> BoxFuture<'_, Result<BoxSubscription, Error>> {
+ Box::pin(async { Err(Error::UnsupportedOperation) })
+ }
+}
+
impl EventSink for HostTransport {
fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
Box::pin(async {
@@ -52,10 +61,12 @@ impl EventSink for HostTransport {
fn main() {
let transport = HostTransport;
let source: &dyn EventSource = &transport;
+ let subscriber: &dyn EventSubscriber = &transport;
let sink: &dyn EventSink = &transport;
let future = source.status();
drop(future); // The composing host chooses and drives its async executor.
+ let _ = subscriber;
let future = sink.status();
drop(future);
}
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -35,6 +35,18 @@ pub enum Error {
DuplicateFetchTargetOutcome,
FetchPageLimitExceeded,
FetchPageRequestMismatch,
+ EmptySubscriptionRequestId,
+ InvalidSubscriptionRequestId,
+ InvalidSubscriptionLimit,
+ InvalidSubscriptionDeadline,
+ SubscriptionCheckpointSetTooLarge,
+ UnexpectedSubscriptionCheckpoint,
+ DuplicateSubscriptionCheckpoint,
+ UnexpectedSubscriptionEvent,
+ SubscriptionEventCheckpointMismatch,
+ SubscriptionEndLimitExceeded,
+ InvalidSubscriptionEnd,
+ SubscriptionEndRequestMismatch,
InvalidSatisfactionPolicy,
EmptyRequiredTargetSet,
DuplicateRequiredTargetFingerprint,
@@ -119,6 +131,42 @@ impl fmt::Display for Error {
Self::FetchPageRequestMismatch => {
f.write_str("transport fetch page does not match its request")
}
+ Self::EmptySubscriptionRequestId => {
+ f.write_str("transport subscription request id is empty")
+ }
+ Self::InvalidSubscriptionRequestId => {
+ f.write_str("transport subscription request id is invalid")
+ }
+ Self::InvalidSubscriptionLimit => {
+ f.write_str("transport subscription event limit is invalid")
+ }
+ Self::InvalidSubscriptionDeadline => {
+ f.write_str("transport subscription deadline is invalid")
+ }
+ Self::SubscriptionCheckpointSetTooLarge => {
+ f.write_str("transport subscription checkpoint set exceeds its target limit")
+ }
+ Self::UnexpectedSubscriptionCheckpoint => {
+ f.write_str("transport subscription contains an unexpected checkpoint")
+ }
+ Self::DuplicateSubscriptionCheckpoint => {
+ f.write_str("transport subscription contains a duplicate checkpoint")
+ }
+ Self::UnexpectedSubscriptionEvent => {
+ f.write_str("transport subscription contains an unexpected event")
+ }
+ Self::SubscriptionEventCheckpointMismatch => {
+ f.write_str("transport subscription event checkpoint does not match provenance")
+ }
+ Self::SubscriptionEndLimitExceeded => {
+ f.write_str("transport subscription exceeded its requested event limit")
+ }
+ Self::InvalidSubscriptionEnd => {
+ f.write_str("transport subscription terminal result is invalid")
+ }
+ Self::SubscriptionEndRequestMismatch => {
+ f.write_str("transport subscription result does not match its request")
+ }
Self::InvalidSatisfactionPolicy => {
f.write_str("transport satisfaction policy is invalid")
}
diff --git a/crates/transport/src/lib.rs b/crates/transport/src/lib.rs
@@ -19,7 +19,11 @@ pub mod target;
pub use error::Error;
pub use id::{TRANSPORT_ID_MAX_BYTES, TransportId};
pub use sink::{DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, SinkStatus};
-pub use source::{BoxFuture, EventSource, FetchPage, FetchRequest, SourceStatus};
+pub use source::{
+ BoxFuture, BoxSubscription, EventSource, EventSubscriber, EventSubscription, FetchPage,
+ FetchRequest, SourceStatus, SubscriptionEnd, SubscriptionEndReason, SubscriptionEvent,
+ SubscriptionNext, SubscriptionRequest,
+};
pub use target::{TARGET_SET_MAX_ITEMS, Target, TargetSet};
#[cfg(test)]
diff --git a/crates/transport/src/source.rs b/crates/transport/src/source.rs
@@ -10,6 +10,9 @@ use core::{fmt, future::Future, pin::Pin};
use radroots_event::SignedEvent;
use radroots_identity::PublicKey;
+#[cfg(feature = "serde")]
+use crate::target::TARGET_SET_MAX_ITEMS;
+
pub use crate::status::SourceStatus;
/// Maximum encoded request identity length.
@@ -22,6 +25,10 @@ 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 encoded live-subscription request identity length.
+pub const SUBSCRIPTION_REQUEST_ID_MAX_BYTES: usize = 256;
+/// Maximum number of events one live subscription may emit.
+pub const SUBSCRIPTION_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>>;
@@ -90,6 +97,38 @@ impl fmt::Display for FetchCursor {
}
}
+/// Validated caller identity for one bounded live subscription.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct SubscriptionRequestId(String);
+
+impl SubscriptionRequestId {
+ /// Parses a non-empty, bounded, printable request identity.
+ pub fn parse(value: impl AsRef<str>) -> Result<Self, Error> {
+ let value = value.as_ref();
+ if value.is_empty() {
+ return Err(Error::EmptySubscriptionRequestId);
+ }
+ if value.len() > SUBSCRIPTION_REQUEST_ID_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control)
+ {
+ return Err(Error::InvalidSubscriptionRequestId);
+ }
+ Ok(Self(String::from(value)))
+ }
+
+ /// Returns the validated request identity.
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+impl fmt::Display for SubscriptionRequestId {
+ 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)]
@@ -124,6 +163,40 @@ impl FetchBounds {
}
}
+/// Hard bounds for one live-subscription operation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct SubscriptionBounds {
+ event_limit: u16,
+ deadline_unix_ms: u64,
+}
+
+impl SubscriptionBounds {
+ /// Creates bounds with a non-zero event limit and absolute deadline.
+ pub const fn new(event_limit: u16, deadline_unix_ms: u64) -> Result<Self, Error> {
+ if event_limit == 0 || event_limit > SUBSCRIPTION_MAX_EVENTS {
+ return Err(Error::InvalidSubscriptionLimit);
+ }
+ if deadline_unix_ms == 0 {
+ return Err(Error::InvalidSubscriptionDeadline);
+ }
+ Ok(Self {
+ event_limit,
+ deadline_unix_ms,
+ })
+ }
+
+ /// Maximum number of events the adapter may emit.
+ pub const fn event_limit(self) -> u16 {
+ self.event_limit
+ }
+
+ /// Absolute Unix deadline in milliseconds.
+ pub const fn deadline_unix_ms(self) -> u64 {
+ self.deadline_unix_ms
+ }
+}
+
/// Transport-neutral constraints applied before a source page is bounded.
///
/// An empty kind or author collection means "any" for that dimension. Time
@@ -508,6 +581,301 @@ impl FetchPage {
}
}
+/// Per-target continuation point for a live subscription.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SubscriptionCheckpoint {
+ target: TargetFingerprint,
+ cursor: FetchCursor,
+}
+
+impl SubscriptionCheckpoint {
+ /// Binds an opaque adapter cursor to one exact target.
+ pub const fn new(target: TargetFingerprint, cursor: FetchCursor) -> Self {
+ Self { target, cursor }
+ }
+
+ /// Returns the exact target fingerprint.
+ pub const fn target(&self) -> &TargetFingerprint {
+ &self.target
+ }
+
+ /// Returns the opaque adapter cursor.
+ pub const fn cursor(&self) -> &FetchCursor {
+ &self.cursor
+ }
+}
+
+/// Bounded request for a live stream from one or more transport targets.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SubscriptionRequest {
+ request_id: SubscriptionRequestId,
+ target_set: TargetSet,
+ bounds: SubscriptionBounds,
+ selector: FetchSelector,
+ checkpoints: Vec<SubscriptionCheckpoint>,
+}
+
+impl SubscriptionRequest {
+ /// Creates a live request with no prior target checkpoints.
+ pub fn new(
+ request_id: impl AsRef<str>,
+ target_set: TargetSet,
+ bounds: SubscriptionBounds,
+ ) -> Result<Self, Error> {
+ Ok(Self {
+ request_id: SubscriptionRequestId::parse(request_id)?,
+ target_set,
+ bounds,
+ selector: FetchSelector::all(),
+ checkpoints: Vec::new(),
+ })
+ }
+
+ /// Applies explicit transport-neutral event constraints.
+ #[must_use]
+ pub fn with_selector(mut self, selector: FetchSelector) -> Self {
+ self.selector = selector;
+ self
+ }
+
+ /// Applies a bounded, unique checkpoint subset in canonical target order.
+ pub fn with_checkpoints<I>(mut self, checkpoints: I) -> Result<Self, Error>
+ where
+ I: IntoIterator<Item = SubscriptionCheckpoint>,
+ {
+ self.checkpoints = normalize_subscription_checkpoints(&self.target_set, checkpoints)?;
+ Ok(self)
+ }
+
+ /// Returns the request identity.
+ pub const fn request_id(&self) -> &SubscriptionRequestId {
+ &self.request_id
+ }
+
+ /// Returns the exact requested target set.
+ pub const fn target_set(&self) -> &TargetSet {
+ &self.target_set
+ }
+
+ /// Returns the hard operation bounds.
+ pub const fn bounds(&self) -> SubscriptionBounds {
+ self.bounds
+ }
+
+ /// Returns the exact event constraints for this request.
+ pub const fn selector(&self) -> &FetchSelector {
+ &self.selector
+ }
+
+ /// Returns prior per-target checkpoints in canonical target order.
+ pub fn checkpoints(&self) -> &[SubscriptionCheckpoint] {
+ self.checkpoints.as_slice()
+ }
+}
+
+/// One request-bound live event and its resulting per-target checkpoint.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SubscriptionEvent {
+ request: SubscriptionRequest,
+ observed: ObservedEvent,
+ checkpoint: SubscriptionCheckpoint,
+}
+
+impl SubscriptionEvent {
+ /// Creates and validates one live event against its originating request.
+ pub fn for_request(
+ request: &SubscriptionRequest,
+ observed: ObservedEvent,
+ checkpoint: SubscriptionCheckpoint,
+ ) -> Result<Self, Error> {
+ let event = Self {
+ request: request.clone(),
+ observed,
+ checkpoint,
+ };
+ event.validate_for_request(request)?;
+ Ok(event)
+ }
+
+ /// Validates target, transport, selector, cursor, and request identity.
+ pub fn validate_for_request(&self, request: &SubscriptionRequest) -> Result<(), Error> {
+ if self.request != *request {
+ return Err(Error::UnexpectedSubscriptionEvent);
+ }
+ let provenance = self.observed.provenance();
+ let Some(target) = request
+ .target_set
+ .targets()
+ .iter()
+ .find(|target| target.fingerprint() == provenance.target())
+ else {
+ return Err(Error::UnexpectedSubscriptionEvent);
+ };
+ if *target.kind() != provenance.transport_id()
+ || !request.selector.matches(self.observed.event())
+ {
+ return Err(Error::UnexpectedSubscriptionEvent);
+ }
+ if self.checkpoint.target() != provenance.target()
+ || provenance.cursor() != Some(self.checkpoint.cursor())
+ {
+ return Err(Error::SubscriptionEventCheckpointMismatch);
+ }
+ Ok(())
+ }
+
+ /// Returns the request identity.
+ pub const fn request_id(&self) -> &SubscriptionRequestId {
+ self.request.request_id()
+ }
+
+ /// Returns the exact originating request.
+ pub const fn request(&self) -> &SubscriptionRequest {
+ &self.request
+ }
+
+ /// Returns the observed event.
+ pub const fn observed(&self) -> &ObservedEvent {
+ &self.observed
+ }
+
+ /// Returns the checkpoint established by this event.
+ pub const fn checkpoint(&self) -> &SubscriptionCheckpoint {
+ &self.checkpoint
+ }
+}
+
+/// Stable terminal reason for a bounded live subscription.
+#[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 SubscriptionEndReason {
+ /// The requested maximum event count was emitted.
+ EventLimit,
+ /// The absolute request deadline was reached.
+ Deadline,
+ /// Explicit or future-drop cancellation was observed.
+ Cancelled,
+ /// The underlying source closed before another event was available.
+ SourceClosed,
+}
+
+/// Request-bound terminal result for a live subscription.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SubscriptionEnd {
+ request: SubscriptionRequest,
+ event_count: u16,
+ checkpoints: Vec<SubscriptionCheckpoint>,
+ reason: SubscriptionEndReason,
+}
+
+impl SubscriptionEnd {
+ /// Creates a terminal result with canonical final checkpoints.
+ pub fn for_request<I>(
+ request: &SubscriptionRequest,
+ event_count: u16,
+ checkpoints: I,
+ reason: SubscriptionEndReason,
+ ) -> Result<Self, Error>
+ where
+ I: IntoIterator<Item = SubscriptionCheckpoint>,
+ {
+ if event_count > request.bounds.event_limit {
+ return Err(Error::SubscriptionEndLimitExceeded);
+ }
+ if reason == SubscriptionEndReason::EventLimit && event_count != request.bounds.event_limit
+ {
+ return Err(Error::InvalidSubscriptionEnd);
+ }
+ Ok(Self {
+ request: request.clone(),
+ event_count,
+ checkpoints: normalize_subscription_checkpoints(&request.target_set, checkpoints)?,
+ reason,
+ })
+ }
+
+ /// Validates that this result belongs to the exact originating request.
+ pub fn validate_for_request(&self, request: &SubscriptionRequest) -> Result<(), Error> {
+ if self.request != *request || self.event_count > request.bounds.event_limit {
+ return Err(Error::SubscriptionEndRequestMismatch);
+ }
+ normalize_subscription_checkpoints(&request.target_set, self.checkpoints.clone())?;
+ Ok(())
+ }
+
+ /// Returns the exact originating request.
+ pub const fn request(&self) -> &SubscriptionRequest {
+ &self.request
+ }
+
+ /// Returns the number of events emitted before termination.
+ pub const fn event_count(&self) -> u16 {
+ self.event_count
+ }
+
+ /// Returns final per-target checkpoints in canonical target order.
+ pub fn checkpoints(&self) -> &[SubscriptionCheckpoint] {
+ self.checkpoints.as_slice()
+ }
+
+ /// Returns why the bounded operation terminated.
+ pub const fn reason(&self) -> SubscriptionEndReason {
+ self.reason
+ }
+}
+
+/// One event or the stable terminal result from a live subscription.
+#[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 SubscriptionNext {
+ /// One request-bound event was observed.
+ Event(Box<SubscriptionEvent>),
+ /// The subscription reached its stable terminal state.
+ End(SubscriptionEnd),
+}
+
+fn normalize_subscription_checkpoints<I>(
+ target_set: &TargetSet,
+ checkpoints: I,
+) -> Result<Vec<SubscriptionCheckpoint>, Error>
+where
+ I: IntoIterator<Item = SubscriptionCheckpoint>,
+{
+ let maximum = target_set.len();
+ let mut checkpoints: Vec<_> = checkpoints.into_iter().take(maximum + 1).collect();
+ if checkpoints.len() > maximum {
+ return Err(Error::SubscriptionCheckpointSetTooLarge);
+ }
+
+ let mut seen = BTreeSet::new();
+ for checkpoint in &checkpoints {
+ if target_position(target_set, checkpoint.target()).is_none() {
+ return Err(Error::UnexpectedSubscriptionCheckpoint);
+ }
+ if !seen.insert(checkpoint.target().as_str()) {
+ return Err(Error::DuplicateSubscriptionCheckpoint);
+ }
+ }
+ checkpoints.sort_by_key(|checkpoint| {
+ target_position(target_set, checkpoint.target()).expect("checkpoint target validated")
+ });
+ Ok(checkpoints)
+}
+
+fn target_position(target_set: &TargetSet, fingerprint: &TargetFingerprint) -> Option<usize> {
+ target_set
+ .targets()
+ .iter()
+ .position(|target| target.fingerprint() == fingerprint)
+}
+
/// Host SPI for inbound event retrieval.
///
/// This trait supports external implementations and is dyn-compatible. Its
@@ -529,6 +897,39 @@ pub trait EventSource: Send + Sync {
fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>>;
}
+/// One active, bounded live-subscription operation.
+///
+/// Implementations must enforce the request's event limit and absolute
+/// deadline without hidden retries. Once [`SubscriptionNext::End`] has been
+/// returned, every later `next` or `cancel` call must return the exact same
+/// terminal result. Dropping a pending future requests cancellation but does
+/// not claim that an already-observed remote event was rolled back.
+pub trait EventSubscription: Send {
+ /// Returns the exact request governing this operation.
+ fn request(&self) -> &SubscriptionRequest;
+
+ /// Returns the next request-bound event or stable terminal result.
+ fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, Error>>;
+
+ /// Requests cancellation and returns the stable terminal result.
+ fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, Error>>;
+}
+
+/// Heap-owned live-subscription capability returned by adapters.
+pub type BoxSubscription = Box<dyn EventSubscription>;
+
+/// Host SPI for beginning bounded live subscriptions.
+///
+/// This is separate from [`EventSource`] so existing bounded-fetch producers
+/// remain source-compatible until they explicitly adopt live delivery.
+pub trait EventSubscriber: Send + Sync {
+ /// Begins one exact bounded live-subscription request.
+ fn subscribe(
+ &self,
+ request: SubscriptionRequest,
+ ) -> BoxFuture<'_, Result<BoxSubscription, Error>>;
+}
+
#[cfg(feature = "serde")]
mod serde_impl {
use super::*;
@@ -571,6 +972,25 @@ mod serde_impl {
}
}
+ impl serde::Serialize for SubscriptionRequestId {
+ 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 SubscriptionRequestId {
+ 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.as_str()).map_err(serde::de::Error::custom)
+ }
+ }
+
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct FetchBoundsWire {
@@ -590,6 +1010,23 @@ mod serde_impl {
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
+ struct SubscriptionBoundsWire {
+ event_limit: u16,
+ deadline_unix_ms: u64,
+ }
+
+ impl<'de> serde::Deserialize<'de> for SubscriptionBounds {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = SubscriptionBoundsWire::deserialize(deserializer)?;
+ Self::new(wire.event_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,
@@ -615,6 +1052,32 @@ mod serde_impl {
}
}
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct SubscriptionRequestWire {
+ request_id: String,
+ target_set: TargetSet,
+ bounds: SubscriptionBounds,
+ #[serde(default)]
+ selector: FetchSelector,
+ #[serde(default)]
+ #[serde(deserialize_with = "deserialize_subscription_checkpoints")]
+ checkpoints: Vec<SubscriptionCheckpoint>,
+ }
+
+ impl<'de> serde::Deserialize<'de> for SubscriptionRequest {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = SubscriptionRequestWire::deserialize(deserializer)?;
+ Self::new(wire.request_id.as_str(), wire.target_set, wire.bounds)
+ .map(|request| request.with_selector(wire.selector))
+ .and_then(|request| request.with_checkpoints(wire.checkpoints))
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
impl<'de> serde::Deserialize<'de> for FetchSelector {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
@@ -703,4 +1166,85 @@ mod serde_impl {
Ok(page)
}
}
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct SubscriptionEventWire {
+ request: SubscriptionRequest,
+ observed: ObservedEvent,
+ checkpoint: SubscriptionCheckpoint,
+ }
+
+ impl<'de> serde::Deserialize<'de> for SubscriptionEvent {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = SubscriptionEventWire::deserialize(deserializer)?;
+ Self::for_request(&wire.request, wire.observed, wire.checkpoint)
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct SubscriptionEndWire {
+ request: SubscriptionRequest,
+ event_count: u16,
+ #[serde(deserialize_with = "deserialize_subscription_checkpoints")]
+ checkpoints: Vec<SubscriptionCheckpoint>,
+ reason: SubscriptionEndReason,
+ }
+
+ impl<'de> serde::Deserialize<'de> for SubscriptionEnd {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = SubscriptionEndWire::deserialize(deserializer)?;
+ Self::for_request(
+ &wire.request,
+ wire.event_count,
+ wire.checkpoints,
+ wire.reason,
+ )
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
+ fn deserialize_subscription_checkpoints<'de, D>(
+ deserializer: D,
+ ) -> Result<Vec<SubscriptionCheckpoint>, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ struct CheckpointVisitor;
+
+ impl<'de> serde::de::Visitor<'de> for CheckpointVisitor {
+ type Value = Vec<SubscriptionCheckpoint>;
+
+ fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("a bounded transport subscription checkpoint sequence")
+ }
+
+ fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
+ where
+ A: serde::de::SeqAccess<'de>,
+ {
+ let capacity = sequence.size_hint().unwrap_or(0).min(TARGET_SET_MAX_ITEMS);
+ let mut checkpoints = Vec::with_capacity(capacity);
+ while let Some(checkpoint) = sequence.next_element()? {
+ if checkpoints.len() == TARGET_SET_MAX_ITEMS {
+ return Err(serde::de::Error::custom(
+ Error::SubscriptionCheckpointSetTooLarge,
+ ));
+ }
+ checkpoints.push(checkpoint);
+ }
+ Ok(checkpoints)
+ }
+ }
+
+ deserializer.deserialize_seq(CheckpointVisitor)
+ }
}
diff --git a/crates/transport/tests/error_contract.rs b/crates/transport/tests/error_contract.rs
@@ -97,6 +97,54 @@ fn every_transport_error_has_stable_operator_facing_text() {
"transport fetch page does not match its request",
),
(
+ Error::EmptySubscriptionRequestId,
+ "transport subscription request id is empty",
+ ),
+ (
+ Error::InvalidSubscriptionRequestId,
+ "transport subscription request id is invalid",
+ ),
+ (
+ Error::InvalidSubscriptionLimit,
+ "transport subscription event limit is invalid",
+ ),
+ (
+ Error::InvalidSubscriptionDeadline,
+ "transport subscription deadline is invalid",
+ ),
+ (
+ Error::SubscriptionCheckpointSetTooLarge,
+ "transport subscription checkpoint set exceeds its target limit",
+ ),
+ (
+ Error::UnexpectedSubscriptionCheckpoint,
+ "transport subscription contains an unexpected checkpoint",
+ ),
+ (
+ Error::DuplicateSubscriptionCheckpoint,
+ "transport subscription contains a duplicate checkpoint",
+ ),
+ (
+ Error::UnexpectedSubscriptionEvent,
+ "transport subscription contains an unexpected event",
+ ),
+ (
+ Error::SubscriptionEventCheckpointMismatch,
+ "transport subscription event checkpoint does not match provenance",
+ ),
+ (
+ Error::SubscriptionEndLimitExceeded,
+ "transport subscription exceeded its requested event limit",
+ ),
+ (
+ Error::InvalidSubscriptionEnd,
+ "transport subscription terminal result is invalid",
+ ),
+ (
+ Error::SubscriptionEndRequestMismatch,
+ "transport subscription result does not match its request",
+ ),
+ (
Error::InvalidSatisfactionPolicy,
"transport satisfaction policy is invalid",
),
diff --git a/crates/transport/tests/package_boundary.rs b/crates/transport/tests/package_boundary.rs
@@ -3,9 +3,10 @@ use std::{collections::BTreeSet, fs, path::Path};
#[allow(unused_imports)]
use radroots_transport::{
DeliveryReceipt as _, DeliveryRequest as _, Error as _, EventSink as _, EventSource as _,
- FetchPage as _, FetchRequest as _, Target as _, TargetSet as _, TransportId as _,
- capability as _, endpoint as _, error as _, outcome as _, policy as _, sink as _, source as _,
- target as _,
+ EventSubscriber as _, EventSubscription as _, FetchPage as _, FetchRequest as _,
+ SubscriptionEnd as _, SubscriptionEvent as _, SubscriptionRequest as _, Target as _,
+ TargetSet as _, TransportId as _, capability as _, endpoint as _, error as _, outcome as _,
+ policy as _, sink as _, source as _, target as _,
};
const MANIFEST: &str = include_str!("../Cargo.toml");
@@ -88,10 +89,13 @@ fn package_documentation_and_reviewed_api_baseline_are_complete() {
}
for required in [
"impl EventSource for HostTransport",
+ "impl EventSubscriber for HostTransport",
"impl EventSink for HostTransport",
"fn fetch(&self, _request: FetchRequest) -> BoxFuture",
+ "_request: SubscriptionRequest",
"_request: DeliveryRequest",
"let source: &dyn EventSource",
+ "let subscriber: &dyn EventSubscriber",
"let sink: &dyn EventSink",
"drop(future)",
] {
@@ -107,7 +111,14 @@ fn package_documentation_and_reviewed_api_baseline_are_complete() {
"pub mod radroots_transport::source",
"pub mod radroots_transport::target",
"pub trait radroots_transport::EventSource",
+ "pub trait radroots_transport::EventSubscriber",
+ "pub trait radroots_transport::EventSubscription",
"pub trait radroots_transport::EventSink",
+ "pub struct radroots_transport::source::SubscriptionBounds",
+ "pub struct radroots_transport::source::SubscriptionCheckpoint",
+ "pub struct radroots_transport::SubscriptionRequest",
+ "pub enum radroots_transport::SubscriptionEndReason",
+ "pub const radroots_transport::source::SUBSCRIPTION_MAX_EVENTS: u16",
"pub fn radroots_transport::target::Target::new(radroots_transport::TransportId",
"pub fn radroots_transport::target::Target::kind(&self) -> &radroots_transport::TransportId",
] {
@@ -154,7 +165,7 @@ fn every_public_module_has_crate_level_documentation() {
}
#[test]
-fn source_and_sink_are_independent_dyn_compatible_host_spis() {
+fn source_subscription_and_sink_are_independent_dyn_compatible_host_spis() {
for required in [
"pub trait EventSource: Send + Sync",
"fn status(&self)",
@@ -170,6 +181,21 @@ fn source_and_sink_are_independent_dyn_compatible_host_spis() {
assert!(!SOURCE.contains("fn deliver("));
for required in [
+ "pub trait EventSubscriber: Send + Sync",
+ "pub trait EventSubscription: Send",
+ "fn subscribe(",
+ "fn next(&mut self)",
+ "fn cancel(&mut self)",
+ "Once [`SubscriptionNext::End`] has been",
+ "exact same",
+ ] {
+ assert!(
+ SOURCE.contains(required),
+ "subscription SPI is missing {required}"
+ );
+ }
+
+ for required in [
"pub trait EventSink: Send + Sync",
"fn status(&self)",
"fn deliver(",
diff --git a/crates/transport/tests/subscription_contract.rs b/crates/transport/tests/subscription_contract.rs
@@ -0,0 +1,391 @@
+use core::{future::Future, pin::Pin, task::Context};
+use std::sync::{
+ Arc,
+ atomic::{AtomicBool, Ordering},
+};
+
+use futures::{executor::block_on, future, task::noop_waker_ref};
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
+use radroots_transport::{
+ BoxFuture, BoxSubscription, Error, EventSubscriber, EventSubscription, SubscriptionEnd,
+ SubscriptionEndReason, SubscriptionEvent, SubscriptionNext, SubscriptionRequest, Target,
+ TargetSet, TransportId,
+ source::{
+ EventProvenance, FetchCursor, FetchSelector, ObservedEvent, SUBSCRIPTION_MAX_EVENTS,
+ SUBSCRIPTION_REQUEST_ID_MAX_BYTES, SubscriptionBounds, SubscriptionCheckpoint,
+ SubscriptionRequestId,
+ },
+};
+
+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) -> SubscriptionRequest {
+ SubscriptionRequest::new(
+ "subscription-request",
+ targets,
+ SubscriptionBounds::new(limit, 1_700_000_100_000).expect("bounds"),
+ )
+ .expect("request")
+}
+
+#[test]
+fn subscription_identity_and_bounds_are_exact_and_bounded() {
+ assert_eq!(
+ SubscriptionRequestId::parse("").expect_err("empty id"),
+ Error::EmptySubscriptionRequestId
+ );
+ for invalid in [" request", "request ", "request\nid"] {
+ assert_eq!(
+ SubscriptionRequestId::parse(invalid).expect_err("invalid id"),
+ Error::InvalidSubscriptionRequestId
+ );
+ }
+ assert_eq!(
+ SubscriptionRequestId::parse("x".repeat(SUBSCRIPTION_REQUEST_ID_MAX_BYTES + 1))
+ .expect_err("oversized id"),
+ Error::InvalidSubscriptionRequestId
+ );
+ let maximum = SubscriptionRequestId::parse("x".repeat(SUBSCRIPTION_REQUEST_ID_MAX_BYTES))
+ .expect("maximum id");
+ assert_eq!(maximum.as_str().len(), SUBSCRIPTION_REQUEST_ID_MAX_BYTES);
+ assert_eq!(maximum.to_string(), maximum.as_str());
+
+ assert_eq!(
+ SubscriptionBounds::new(0, 1).expect_err("zero limit"),
+ Error::InvalidSubscriptionLimit
+ );
+ assert_eq!(
+ SubscriptionBounds::new(SUBSCRIPTION_MAX_EVENTS + 1, 1).expect_err("oversized limit"),
+ Error::InvalidSubscriptionLimit
+ );
+ assert_eq!(
+ SubscriptionBounds::new(1, 0).expect_err("zero deadline"),
+ Error::InvalidSubscriptionDeadline
+ );
+ let maximum =
+ SubscriptionBounds::new(SUBSCRIPTION_MAX_EVENTS, u64::MAX).expect("maximum bounds");
+ assert_eq!(maximum.event_limit(), SUBSCRIPTION_MAX_EVENTS);
+ assert_eq!(maximum.deadline_unix_ms(), u64::MAX);
+}
+
+#[test]
+fn checkpoints_are_bounded_unique_and_canonical_for_the_target_set() {
+ let first = target("wss://one.example");
+ let second = target("wss://two.example");
+ let targets = TargetSet::new(vec![first.clone(), second.clone()]).expect("targets");
+ let first_checkpoint = SubscriptionCheckpoint::new(
+ first.fingerprint().clone(),
+ FetchCursor::parse("first").expect("cursor"),
+ );
+ let second_checkpoint = SubscriptionCheckpoint::new(
+ second.fingerprint().clone(),
+ FetchCursor::parse("second").expect("cursor"),
+ );
+ let configured = request(targets.clone(), 2)
+ .with_checkpoints([second_checkpoint.clone(), first_checkpoint.clone()])
+ .expect("checkpoints");
+ assert_eq!(
+ configured.checkpoints(),
+ &[first_checkpoint.clone(), second_checkpoint]
+ );
+
+ assert_eq!(
+ request(targets.clone(), 2)
+ .with_checkpoints([first_checkpoint.clone(), first_checkpoint.clone()])
+ .expect_err("duplicate"),
+ Error::DuplicateSubscriptionCheckpoint
+ );
+ let foreign = target("wss://foreign.example");
+ assert_eq!(
+ request(targets.clone(), 2)
+ .with_checkpoints([SubscriptionCheckpoint::new(
+ foreign.fingerprint().clone(),
+ FetchCursor::parse("foreign").expect("cursor"),
+ )])
+ .expect_err("foreign"),
+ Error::UnexpectedSubscriptionCheckpoint
+ );
+ assert_eq!(
+ request(TargetSet::new(vec![first]).expect("targets"), 1)
+ .with_checkpoints(core::iter::repeat(first_checkpoint))
+ .expect_err("infinite iterator is bounded"),
+ Error::SubscriptionCheckpointSetTooLarge
+ );
+}
+
+#[test]
+fn live_events_bind_selector_target_transport_and_checkpoint() {
+ let requested = target("wss://one.example");
+ let targets = TargetSet::new(vec![requested.clone()]).expect("targets");
+ let request = request(targets, 2)
+ .with_selector(FetchSelector::all().with_kinds(vec![0]).expect("selector"));
+ let cursor = FetchCursor::parse("event-1").expect("cursor");
+ let observed = ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(
+ TransportId::NOSTR,
+ requested.fingerprint().clone(),
+ 1_700_000_000_001,
+ )
+ .expect("provenance")
+ .with_cursor(cursor.clone()),
+ );
+ let event = SubscriptionEvent::for_request(
+ &request,
+ observed.clone(),
+ SubscriptionCheckpoint::new(requested.fingerprint().clone(), cursor),
+ )
+ .expect("event");
+ event
+ .validate_for_request(&request)
+ .expect("request binding");
+ assert_eq!(event.request_id(), request.request_id());
+ assert_eq!(event.observed().event().id_str(), signed_event().id_str());
+
+ assert_eq!(
+ SubscriptionEvent::for_request(
+ &request,
+ observed,
+ SubscriptionCheckpoint::new(
+ requested.fingerprint().clone(),
+ FetchCursor::parse("different").expect("cursor"),
+ ),
+ )
+ .expect_err("cursor mismatch"),
+ Error::SubscriptionEventCheckpointMismatch
+ );
+
+ let filtered =
+ request.with_selector(FetchSelector::all().with_kinds(vec![1]).expect("selector"));
+ let observed = ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(TransportId::NOSTR, requested.fingerprint().clone(), 1)
+ .expect("provenance")
+ .with_cursor(FetchCursor::parse("event-2").expect("cursor")),
+ );
+ assert_eq!(
+ SubscriptionEvent::for_request(
+ &filtered,
+ observed,
+ SubscriptionCheckpoint::new(
+ requested.fingerprint().clone(),
+ FetchCursor::parse("event-2").expect("cursor"),
+ ),
+ )
+ .expect_err("selector mismatch"),
+ Error::UnexpectedSubscriptionEvent
+ );
+}
+
+struct StableSubscription {
+ request: SubscriptionRequest,
+ terminal: SubscriptionEnd,
+}
+
+impl EventSubscription for StableSubscription {
+ fn request(&self) -> &SubscriptionRequest {
+ &self.request
+ }
+
+ fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, Error>> {
+ let terminal = self.terminal.clone();
+ Box::pin(async move { Ok(SubscriptionNext::End(terminal)) })
+ }
+
+ fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, Error>> {
+ let terminal = self.terminal.clone();
+ Box::pin(async move { Ok(terminal) })
+ }
+}
+
+struct StableSubscriber;
+
+impl EventSubscriber for StableSubscriber {
+ fn subscribe(
+ &self,
+ request: SubscriptionRequest,
+ ) -> BoxFuture<'_, Result<BoxSubscription, Error>> {
+ Box::pin(async move {
+ let terminal =
+ SubscriptionEnd::for_request(&request, 0, [], SubscriptionEndReason::SourceClosed)?;
+ Ok(Box::new(StableSubscription { request, terminal }) as BoxSubscription)
+ })
+ }
+}
+
+#[test]
+fn subscription_spi_is_dyn_compatible_and_terminal_results_are_idempotent() {
+ let subscriber: &dyn EventSubscriber = &StableSubscriber;
+ let targets = TargetSet::new(vec![target("wss://one.example")]).expect("targets");
+ let request = request(targets, 1);
+ let mut subscription = block_on(subscriber.subscribe(request.clone())).expect("subscribe");
+ assert_eq!(subscription.request(), &request);
+
+ let first = block_on(subscription.next()).expect("next");
+ let second = block_on(subscription.next()).expect("next again");
+ let cancelled = block_on(subscription.cancel()).expect("cancel after end");
+ assert_eq!(first, second);
+ assert_eq!(first, SubscriptionNext::End(cancelled.clone()));
+ assert_eq!(cancelled.reason(), SubscriptionEndReason::SourceClosed);
+ assert_eq!(cancelled.event_count(), 0);
+ assert!(cancelled.checkpoints().is_empty());
+ cancelled
+ .validate_for_request(&request)
+ .expect("request-bound end");
+ assert_eq!(cancelled.request(), &request);
+
+ assert_eq!(
+ SubscriptionEnd::for_request(&request, 2, [], SubscriptionEndReason::EventLimit,)
+ .expect_err("event limit exceeded"),
+ Error::SubscriptionEndLimitExceeded
+ );
+ assert_eq!(
+ SubscriptionEnd::for_request(&request, 0, [], SubscriptionEndReason::EventLimit,)
+ .expect_err("event limit reason requires the exact limit"),
+ Error::InvalidSubscriptionEnd
+ );
+
+ let other = SubscriptionRequest::new(
+ "other-request",
+ request.target_set().clone(),
+ request.bounds(),
+ )
+ .expect("other request");
+ assert_eq!(
+ cancelled
+ .validate_for_request(&other)
+ .expect_err("request mismatch"),
+ Error::SubscriptionEndRequestMismatch
+ );
+
+ for (reason, event_count) in [
+ (SubscriptionEndReason::EventLimit, 1),
+ (SubscriptionEndReason::Deadline, 0),
+ (SubscriptionEndReason::Cancelled, 0),
+ (SubscriptionEndReason::SourceClosed, 0),
+ ] {
+ assert_eq!(
+ SubscriptionEnd::for_request(&request, event_count, [], reason)
+ .expect("terminal reason")
+ .reason(),
+ reason
+ );
+ }
+}
+
+struct CancellationGuard(Arc<AtomicBool>);
+
+impl Drop for CancellationGuard {
+ fn drop(&mut self) {
+ self.0.store(true, Ordering::SeqCst);
+ }
+}
+
+struct PendingSubscription {
+ request: SubscriptionRequest,
+ terminal: SubscriptionEnd,
+ cancellation_observed: Arc<AtomicBool>,
+}
+
+impl EventSubscription for PendingSubscription {
+ fn request(&self) -> &SubscriptionRequest {
+ &self.request
+ }
+
+ fn next(&mut self) -> BoxFuture<'_, Result<SubscriptionNext, Error>> {
+ let cancellation_observed = Arc::clone(&self.cancellation_observed);
+ Box::pin(async move {
+ let _guard = CancellationGuard(cancellation_observed);
+ future::pending::<Result<SubscriptionNext, Error>>().await
+ })
+ }
+
+ fn cancel(&mut self) -> BoxFuture<'_, Result<SubscriptionEnd, Error>> {
+ let terminal = self.terminal.clone();
+ Box::pin(async move { Ok(terminal) })
+ }
+}
+
+#[test]
+fn dropping_a_polled_subscription_future_requests_cancellation() {
+ let targets = TargetSet::new(vec![target("wss://one.example")]).expect("targets");
+ let request = request(targets, 1);
+ let terminal = SubscriptionEnd::for_request(&request, 0, [], SubscriptionEndReason::Cancelled)
+ .expect("terminal");
+ let cancellation_observed = Arc::new(AtomicBool::new(false));
+ let mut subscription = PendingSubscription {
+ request,
+ terminal,
+ cancellation_observed: Arc::clone(&cancellation_observed),
+ };
+
+ let unpolled = subscription.next();
+ drop(unpolled);
+ assert!(!cancellation_observed.load(Ordering::SeqCst));
+
+ let mut pending = subscription.next();
+ let mut context = Context::from_waker(noop_waker_ref());
+ assert!(Pin::new(&mut pending).poll(&mut context).is_pending());
+ drop(pending);
+ assert!(cancellation_observed.load(Ordering::SeqCst));
+ assert_eq!(
+ block_on(subscription.cancel()).expect("cancel").reason(),
+ SubscriptionEndReason::Cancelled
+ );
+}
+
+#[cfg(feature = "serde")]
+#[test]
+fn subscription_wire_models_revalidate_bounds_and_request_binding() {
+ let requested = target("wss://one.example");
+ let request = request(TargetSet::new(vec![requested.clone()]).expect("targets"), 1);
+ let cursor = FetchCursor::parse("event-1").expect("cursor");
+ let event = SubscriptionEvent::for_request(
+ &request,
+ ObservedEvent::new(
+ signed_event(),
+ EventProvenance::new(TransportId::NOSTR, requested.fingerprint().clone(), 1)
+ .expect("provenance")
+ .with_cursor(cursor.clone()),
+ ),
+ SubscriptionCheckpoint::new(requested.fingerprint().clone(), cursor),
+ )
+ .expect("event");
+ let encoded = serde_json::to_string(&SubscriptionNext::Event(Box::new(event.clone())))
+ .expect("serialize event");
+ assert_eq!(
+ serde_json::from_str::<SubscriptionNext>(&encoded).expect("deserialize event"),
+ SubscriptionNext::Event(Box::new(event))
+ );
+
+ let mut invalid = serde_json::to_value(&request).expect("request value");
+ invalid["bounds"]["event_limit"] = 0.into();
+ assert!(serde_json::from_value::<SubscriptionRequest>(invalid).is_err());
+
+ let checkpoint = serde_json::to_value(SubscriptionCheckpoint::new(
+ requested.fingerprint().clone(),
+ FetchCursor::parse("checkpoint").expect("cursor"),
+ ))
+ .expect("checkpoint value");
+ let mut oversized = serde_json::to_value(&request).expect("request value");
+ oversized["checkpoints"] = serde_json::Value::Array(vec![
+ checkpoint;
+ radroots_transport::TARGET_SET_MAX_ITEMS
+ + 1
+ ]);
+ assert!(serde_json::from_value::<SubscriptionRequest>(oversized).is_err());
+
+ let terminal = SubscriptionEnd::for_request(&request, 1, [], SubscriptionEndReason::Deadline)
+ .expect("terminal");
+ let mut unknown = serde_json::to_value(&terminal).expect("terminal value");
+ unknown["unknown"] = true.into();
+ assert!(serde_json::from_value::<SubscriptionEnd>(unknown).is_err());
+}