commit 32d24adbeb0f874c1f242be47a6f3a49c53ccc4d
parent ff893f0207fef17a6a2a406503bba475d50e41a9
Author: triesap <tyson@radroots.org>
Date: Thu, 30 Jul 2026 16:00:35 +0000
transport: split event source and event sink traits
- add independent dyn-compatible EventSource and EventSink host SPIs
- document status deadline cancellation and remote publication behavior
- retain the monolithic contract only as a hidden migration bridge
- cover source-only sink-only and bidirectional dynamic dispatch
Diffstat:
7 files changed, 297 insertions(+), 1 deletion(-)
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -1,5 +1,8 @@
use core::fmt;
+/// Transport contract failure.
+pub type Error = RadrootsTransportError;
+
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RadrootsTransportError {
UnsupportedOperation,
diff --git a/crates/transport/src/lib.rs b/crates/transport/src/lib.rs
@@ -26,7 +26,7 @@ pub use delivery::{
RadrootsTransportDeliveryRequest, RadrootsTransportSatisfactionClass,
RadrootsTransportSatisfactionPolicy, RadrootsTransportTargetReceipt,
};
-pub use error::RadrootsTransportError;
+pub use error::{Error, RadrootsTransportError};
pub use id::{RadrootsTransportKind, TRANSPORT_ID_MAX_BYTES, TransportId};
pub use kind::{
RadrootsTransportCapabilityAvailability, RadrootsTransportCapabilityMaturity,
@@ -43,6 +43,8 @@ pub use reticulum::{
ReticulumGatewaySemanticsV1, ReticulumPayloadPolicyV1, ReticulumPrivacySemanticsV1,
ReticulumRoutingMetadataV1,
};
+pub use sink::{DeliveryReceipt, DeliveryRequest, EventSink, SinkStatus};
+pub use source::{BoxFuture, EventSource, FetchPage, FetchRequest, SourceStatus};
pub use status::{
RadrootsTransportCapabilities, RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome,
RadrootsTransportOutcomeKind, RadrootsTransportStatus,
diff --git a/crates/transport/src/sink.rs b/crates/transport/src/sink.rs
@@ -1 +1,42 @@
//! Outbound event delivery SPI and request models.
+
+use crate::{
+ Error, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
+ RadrootsTransportStatus, source::BoxFuture,
+};
+
+/// Current sink status.
+///
+/// This compatibility alias is replaced by the dedicated sink status model in
+/// the ordered capability/status checkpoint.
+pub type SinkStatus = RadrootsTransportStatus;
+
+/// Bounded delivery request.
+///
+/// The dedicated delivery checkpoint replaces this compatibility alias with
+/// the final request model.
+pub type DeliveryRequest = RadrootsTransportDeliveryRequest;
+
+/// Per-target delivery result.
+pub type DeliveryReceipt = RadrootsTransportDeliveryReceipt;
+
+/// Host SPI for outbound event delivery.
+///
+/// This trait supports external implementations and is dyn-compatible. Its
+/// futures are `Send`; implementations must not borrow request data after a
+/// future completes. `status` observes sink state and does not initiate
+/// delivery. `deliver` performs only the attempts authorized by its request,
+/// returns partial success per target, and owns no hidden retry loop.
+///
+/// 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
+/// explicit request deadline bounds work independently of future cancellation.
+pub trait EventSink: Send + Sync {
+ /// Returns the sink's current runtime status.
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>>;
+
+ /// Delivers an event according to the request's bounded target policy.
+ fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>>;
+}
diff --git a/crates/transport/src/source.rs b/crates/transport/src/source.rs
@@ -1 +1,46 @@
//! Inbound event source SPI and page models.
+
+use crate::{
+ Error, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest, RadrootsTransportStatus,
+};
+use alloc::boxed::Box;
+use core::{future::Future, pin::Pin};
+
+/// Heap-backed future returned by transport SPIs.
+pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
+
+/// Current source status.
+///
+/// This compatibility alias is replaced by the dedicated source status model
+/// in the ordered capability/status checkpoint.
+pub type SourceStatus = RadrootsTransportStatus;
+
+/// Bounded source request.
+///
+/// The dedicated bounded-page checkpoint replaces this compatibility alias
+/// with the final request model.
+pub type FetchRequest = RadrootsTransportFetchRequest;
+
+/// One bounded page returned by an event source.
+pub type FetchPage = RadrootsTransportFetchReceipt;
+
+/// Host SPI for inbound event retrieval.
+///
+/// This trait supports external implementations and is dyn-compatible. Its
+/// futures are `Send`; implementations must not borrow request data after a
+/// future completes. `status` observes source state and does not initiate a
+/// fetch. `fetch` performs at most the work bounded by its request and owns no
+/// hidden retry loop.
+///
+/// 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
+/// explicit request deadline bounds work independently of future cancellation.
+pub trait EventSource: Send + Sync {
+ /// Returns the source's current runtime status.
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>>;
+
+ /// Fetches one bounded page of transport-neutral events.
+ fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>>;
+}
diff --git a/crates/transport/src/transport.rs b/crates/transport/src/transport.rs
@@ -12,6 +12,8 @@ use core::pin::Pin;
pub type RadrootsTransportFuture<'a, T> =
Pin<Box<dyn Future<Output = Result<T, RadrootsTransportError>> + Send + 'a>>;
+/// Compatibility contract retained until the planned workspace consumer cutover.
+#[doc(hidden)]
pub trait RadrootsTransport: Send + Sync {
fn transport_kind(&self) -> RadrootsTransportKind;
diff --git a/crates/transport/tests/package_boundary.rs b/crates/transport/tests/package_boundary.rs
@@ -8,6 +8,9 @@ use radroots_transport::{
const MANIFEST: &str = include_str!("../Cargo.toml");
const ROOT: &str = include_str!("../src/lib.rs");
+const SOURCE: &str = include_str!("../src/source.rs");
+const SINK: &str = include_str!("../src/sink.rs");
+const LEGACY_TRANSPORT: &str = include_str!("../src/transport.rs");
#[test]
fn manifest_has_final_identity_features_and_required_radroots_dependencies() {
@@ -58,6 +61,36 @@ fn crate_root_declares_the_approved_public_module_skeleton() {
);
}
+#[test]
+fn source_and_sink_are_independent_dyn_compatible_host_spis() {
+ for required in [
+ "pub trait EventSource: Send + Sync",
+ "fn status(&self)",
+ "fn fetch(",
+ "Dropping a returned future requests cancellation.",
+ "explicit request deadline",
+ ] {
+ assert!(
+ SOURCE.contains(required),
+ "source SPI is missing {required}"
+ );
+ }
+ assert!(!SOURCE.contains("fn deliver("));
+
+ for required in [
+ "pub trait EventSink: Send + Sync",
+ "fn status(&self)",
+ "fn deliver(",
+ "Dropping a returned future requests cancellation.",
+ "explicit request deadline",
+ ] {
+ assert!(SINK.contains(required), "sink SPI is missing {required}");
+ }
+ assert!(!SINK.contains("fn fetch("));
+
+ assert!(LEGACY_TRANSPORT.contains("#[doc(hidden)]\npub trait RadrootsTransport: Send + Sync"));
+}
+
fn table_keys<'a>(source: &'a str, heading: &str) -> BTreeSet<&'a str> {
let mut in_table = false;
source
diff --git a/crates/transport/tests/spi.rs b/crates/transport/tests/spi.rs
@@ -0,0 +1,170 @@
+use futures::executor::block_on;
+use radroots_transport::{
+ BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, EventSource, FetchPage, FetchRequest,
+ RadrootsTransportCapabilities, RadrootsTransportImplementationState, RadrootsTransportKind,
+ RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportPayload,
+ RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget,
+ RadrootsTransportTargetReceipt, RadrootsTransportTargetSet, SinkStatus, SourceStatus,
+};
+
+struct SourceOnly;
+
+impl EventSource for SourceOnly {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
+ Box::pin(async { Ok(source_status()) })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
+ Box::pin(async move { Ok(FetchPage::new(request.request_id, Vec::new(), 0)) })
+ }
+}
+
+struct SinkOnly;
+
+impl EventSink for SinkOnly {
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
+ Box::pin(async { Ok(sink_status()) })
+ }
+
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ Box::pin(async move {
+ let receipts = request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| {
+ RadrootsTransportTargetReceipt::new(
+ target,
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Delivered),
+ )
+ })
+ .collect();
+ DeliveryReceipt::for_request(&request, receipts)
+ })
+ }
+}
+
+struct Bidirectional {
+ source: SourceOnly,
+ sink: SinkOnly,
+}
+
+impl EventSource for Bidirectional {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> {
+ EventSource::status(&self.source)
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
+ self.source.fetch(request)
+ }
+}
+
+impl EventSink for Bidirectional {
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
+ EventSink::status(&self.sink)
+ }
+
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ self.sink.deliver(request)
+ }
+}
+
+fn source_status() -> RadrootsTransportStatus {
+ RadrootsTransportStatus::new(
+ RadrootsTransportKind::Local,
+ true,
+ RadrootsTransportImplementationState::Real,
+ false,
+ "source ready",
+ )
+ .with_capabilities(RadrootsTransportCapabilities::fetch_only())
+}
+
+fn sink_status() -> RadrootsTransportStatus {
+ RadrootsTransportStatus::new(
+ RadrootsTransportKind::Local,
+ true,
+ RadrootsTransportImplementationState::Real,
+ true,
+ "sink ready",
+ )
+ .with_capabilities(RadrootsTransportCapabilities::deliver_only())
+}
+
+fn target_set() -> RadrootsTransportTargetSet {
+ RadrootsTransportTargetSet::new(vec![
+ RadrootsTransportTarget::local("local:spi").expect("local target"),
+ ])
+ .expect("target set")
+}
+
+fn assert_source_dyn_compatible(_: &dyn EventSource) {}
+fn assert_sink_dyn_compatible(_: &dyn EventSink) {}
+
+#[test]
+fn source_only_and_sink_only_implementations_are_independently_dispatchable() {
+ let source = SourceOnly;
+ let sink = SinkOnly;
+ assert_source_dyn_compatible(&source);
+ assert_sink_dyn_compatible(&sink);
+
+ let source_status = block_on(EventSource::status(&source)).expect("source status");
+ assert!(source_status.capabilities.fetch);
+ assert!(!source_status.capabilities.deliver);
+ let page =
+ block_on(source.fetch(FetchRequest::new("fetch-1", target_set()))).expect("fetch page");
+ assert_eq!(page.request_id, "fetch-1");
+
+ let sink_status = block_on(EventSink::status(&sink)).expect("sink status");
+ assert!(sink_status.capabilities.deliver);
+ assert!(!sink_status.capabilities.fetch);
+ let receipt = block_on(
+ sink.deliver(
+ DeliveryRequest::new(
+ "deliver-1",
+ RadrootsTransportPayload::opaque_bytes("spi", [1]).expect("payload"),
+ target_set(),
+ RadrootsTransportSatisfactionPolicy::all_delivered(),
+ )
+ .expect("delivery request"),
+ ),
+ )
+ .expect("delivery receipt");
+ assert_eq!(receipt.request_id(), "deliver-1");
+}
+
+#[test]
+fn a_bidirectional_adapter_exposes_both_dyn_contracts() {
+ let adapter = Bidirectional {
+ source: SourceOnly,
+ sink: SinkOnly,
+ };
+ assert_source_dyn_compatible(&adapter);
+ assert_sink_dyn_compatible(&adapter);
+
+ assert!(
+ block_on(EventSource::status(&adapter))
+ .expect("source status")
+ .capabilities
+ .fetch
+ );
+ assert!(
+ block_on(EventSink::status(&adapter))
+ .expect("sink status")
+ .capabilities
+ .deliver
+ );
+}