commit aeb160c6db51ccaf3cd4c532b044c9c4e98c5a89
parent cabaac3e0bc778364ec8905dfcffbffcaee198b0
Author: triesap <tyson@radroots.org>
Date: Thu, 30 Jul 2026 18:18:11 +0000
transport: create a transport conformance suite
- add reusable source and sink conformance harness contracts
- cover status, bounds, identities, targets, partial outcomes, and errors
- prove deadline and cancellation behavior at publication boundaries
- run the suite against source-only, sink-only, and combined private mocks
Diffstat:
3 files changed, 714 insertions(+), 0 deletions(-)
diff --git a/crates/transport/tests/conformance.rs b/crates/transport/tests/conformance.rs
@@ -0,0 +1,53 @@
+#[path = "conformance/suite.rs"]
+mod suite;
+#[path = "conformance/support.rs"]
+mod support;
+
+use radroots_transport::Error;
+use suite::{
+ assert_request_boundaries, assert_sink_cancellation, assert_sink_conformance,
+ assert_sink_error, assert_source_cancellation, assert_source_conformance, assert_source_error,
+};
+use support::{CombinedAdapter, MockSink, MockSource};
+
+#[test]
+fn source_only_adapter_satisfies_the_reusable_contract() {
+ let source = MockSource::successful();
+ assert_source_conformance(&source);
+}
+
+#[test]
+fn sink_only_adapter_satisfies_the_reusable_contract() {
+ let sink = MockSink::successful();
+ assert_sink_conformance(&sink);
+}
+
+#[test]
+fn combined_adapter_satisfies_both_reusable_contracts() {
+ let adapter = CombinedAdapter::successful();
+ assert_source_conformance(&adapter);
+ assert_sink_conformance(&adapter);
+}
+
+#[test]
+fn request_identity_and_operation_bounds_fail_closed() {
+ assert_request_boundaries();
+}
+
+#[test]
+fn normalized_adapter_errors_are_not_retried_or_rewritten() {
+ let source = MockSource::failing(Error::UnsupportedOperation);
+ assert_source_error(&source, Error::UnsupportedOperation);
+
+ let sink = MockSink::failing(Error::UnsupportedOperation);
+ assert_sink_error(&sink, Error::UnsupportedOperation);
+}
+
+#[test]
+fn source_and_sink_cancellation_observe_publication_boundaries() {
+ let source = MockSource::pending();
+ assert_source_cancellation(&source);
+
+ let sink = MockSink::pending();
+ assert_sink_cancellation(&sink);
+}
diff --git a/crates/transport/tests/conformance/suite.rs b/crates/transport/tests/conformance/suite.rs
@@ -0,0 +1,256 @@
+use core::{future::Future, pin::Pin, task::Context};
+use futures::{executor::block_on, task::noop_waker_ref};
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
+use radroots_transport::{
+ DeliveryRequest, Error, EventSink, EventSource, FetchRequest, SinkStatus, SourceStatus, Target,
+ TargetSet,
+ outcome::FetchTargetState,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::{DELIVERY_REQUEST_ID_MAX_BYTES, DeliveryPayload},
+ source::{FETCH_PAGE_MAX_EVENTS, FETCH_REQUEST_ID_MAX_BYTES, FetchBounds, FetchCursor},
+};
+
+pub(crate) const NOW_UNIX_MS: u64 = 1_700_000_000_000;
+
+pub(crate) trait SourceConformanceHarness {
+ fn source(&self) -> &dyn EventSource;
+ fn expected_status(&self) -> SourceStatus;
+ fn target_set(&self) -> TargetSet;
+ fn now_unix_ms(&self) -> u64;
+ fn captured_request(&self) -> Option<FetchRequest>;
+ fn published(&self) -> bool;
+ fn cancelled_after_publish(&self) -> bool;
+}
+
+pub(crate) trait SinkConformanceHarness {
+ fn sink(&self) -> &dyn EventSink;
+ fn expected_status(&self) -> SinkStatus;
+ fn target_set(&self) -> TargetSet;
+ fn now_unix_ms(&self) -> u64;
+ fn captured_request(&self) -> Option<DeliveryRequest>;
+ fn published(&self) -> bool;
+ fn cancelled_after_publish(&self) -> bool;
+}
+
+fn fetch_request(id: &str, targets: TargetSet, deadline: u64) -> FetchRequest {
+ FetchRequest::new(
+ id,
+ targets,
+ FetchBounds::new(2, deadline).expect("fetch bounds"),
+ )
+ .expect("fetch request")
+ .with_cursor(FetchCursor::parse("opaque-cursor").expect("cursor"))
+}
+
+fn delivery_request(id: &str, targets: TargetSet, deadline: u64) -> DeliveryRequest {
+ DeliveryRequest::new(
+ id,
+ DeliveryPayload::new(signed_event()),
+ targets,
+ SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::any()),
+ deadline,
+ )
+ .expect("delivery request")
+}
+
+pub(crate) fn assert_source_conformance(harness: &impl SourceConformanceHarness) {
+ let source = harness.source();
+ let status = block_on(source.status()).expect("source status");
+ assert_eq!(status, harness.expected_status());
+ assert!(status.is_configured());
+ assert!(status.capabilities().can_fetch());
+ assert!(!status.message().is_empty());
+
+ let request = fetch_request(
+ "source-conformance",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ );
+ let page = block_on(source.fetch(request.clone())).expect("fetch page");
+ page.validate_for_request(&request)
+ .expect("request binding");
+ assert_eq!(page.request_id().as_str(), request.request_id().as_str());
+ assert!(page.events().len() <= usize::from(request.bounds().limit()));
+ assert_eq!(page.target_outcomes().len(), request.target_set().len());
+ assert_eq!(
+ page.target_outcomes()[0].target(),
+ request.target_set().targets()[0].fingerprint()
+ );
+ assert_eq!(
+ page.target_outcomes()[0].state(),
+ FetchTargetState::Complete
+ );
+ assert_eq!(
+ page.target_outcomes()[1].target(),
+ request.target_set().targets()[1].fingerprint()
+ );
+ assert!(page.target_outcomes()[1].state().is_retryable());
+ assert_eq!(harness.captured_request().as_ref(), Some(&request));
+
+ let expired = fetch_request(
+ "source-expired",
+ harness.target_set(),
+ harness.now_unix_ms(),
+ );
+ assert_eq!(
+ block_on(source.fetch(expired)).expect_err("expired fetch"),
+ Error::InvalidFetchDeadline
+ );
+}
+
+pub(crate) fn assert_sink_conformance(harness: &impl SinkConformanceHarness) {
+ let sink = harness.sink();
+ let status = block_on(sink.status()).expect("sink status");
+ assert_eq!(status, harness.expected_status());
+ assert!(status.is_configured());
+ assert!(status.capabilities().can_deliver());
+ assert!(!status.message().is_empty());
+
+ let request = delivery_request(
+ "sink-conformance",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ );
+ let receipt = block_on(sink.deliver(request.clone())).expect("delivery receipt");
+ receipt
+ .validate_for_request(&request)
+ .expect("request binding");
+ assert_eq!(receipt.request_id().as_str(), request.request_id().as_str());
+ assert_eq!(receipt.target_receipts().len(), request.target_set().len());
+ for (target_receipt, requested_target) in receipt
+ .target_receipts()
+ .iter()
+ .zip(request.target_set().targets())
+ {
+ assert_eq!(target_receipt.target(), requested_target);
+ }
+ assert!(
+ receipt.target_receipts()[0]
+ .outcome()
+ .satisfies(SatisfactionClass::Delivered)
+ );
+ assert!(receipt.target_receipts()[1].outcome().is_retryable());
+ assert!(receipt.is_satisfied(&request).expect("satisfaction"));
+ assert_eq!(harness.captured_request().as_ref(), Some(&request));
+
+ let expired = delivery_request("sink-expired", harness.target_set(), harness.now_unix_ms());
+ assert_eq!(
+ block_on(sink.deliver(expired)).expect_err("expired delivery"),
+ Error::InvalidDeliveryDeadline
+ );
+}
+
+pub(crate) fn assert_request_boundaries() {
+ assert_eq!(
+ FetchBounds::new(0, 1).expect_err("zero fetch bound"),
+ Error::InvalidFetchLimit
+ );
+ assert_eq!(
+ FetchBounds::new(FETCH_PAGE_MAX_EVENTS + 1, 1).expect_err("oversized fetch bound"),
+ Error::InvalidFetchLimit
+ );
+ assert_eq!(
+ FetchRequest::new(
+ "x".repeat(FETCH_REQUEST_ID_MAX_BYTES + 1),
+ target_set(),
+ FetchBounds::new(1, 1).expect("fetch bounds"),
+ )
+ .expect_err("oversized fetch request id"),
+ Error::InvalidFetchRequestId
+ );
+ assert_eq!(
+ DeliveryRequest::new(
+ "x".repeat(DELIVERY_REQUEST_ID_MAX_BYTES + 1),
+ DeliveryPayload::new(signed_event()),
+ target_set(),
+ SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all()),
+ 1,
+ )
+ .expect_err("oversized delivery request id"),
+ Error::InvalidDeliveryRequestId
+ );
+}
+
+pub(crate) fn assert_source_error(harness: &impl SourceConformanceHarness, expected: Error) {
+ let request = fetch_request(
+ "source-error",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ );
+ assert_eq!(
+ block_on(harness.source().fetch(request)).expect_err("source error"),
+ expected
+ );
+}
+
+pub(crate) fn assert_sink_error(harness: &impl SinkConformanceHarness, expected: Error) {
+ let request = delivery_request(
+ "sink-error",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ );
+ assert_eq!(
+ block_on(harness.sink().deliver(request)).expect_err("sink error"),
+ expected
+ );
+}
+
+pub(crate) fn assert_source_cancellation(harness: &impl SourceConformanceHarness) {
+ let source = harness.source();
+ let unpolled = source.fetch(fetch_request(
+ "source-unpolled",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ ));
+ drop(unpolled);
+ assert!(!harness.published());
+ assert!(!harness.cancelled_after_publish());
+
+ let mut published = source.fetch(fetch_request(
+ "source-published",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ ));
+ let mut context = Context::from_waker(noop_waker_ref());
+ assert!(Pin::new(&mut published).poll(&mut context).is_pending());
+ assert!(harness.published());
+ drop(published);
+ assert!(harness.cancelled_after_publish());
+}
+
+pub(crate) fn assert_sink_cancellation(harness: &impl SinkConformanceHarness) {
+ let sink = harness.sink();
+ let unpolled = sink.deliver(delivery_request(
+ "sink-unpolled",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ ));
+ drop(unpolled);
+ assert!(!harness.published());
+ assert!(!harness.cancelled_after_publish());
+
+ let mut published = sink.deliver(delivery_request(
+ "sink-published",
+ harness.target_set(),
+ harness.now_unix_ms() + 100,
+ ));
+ let mut context = Context::from_waker(noop_waker_ref());
+ assert!(Pin::new(&mut published).poll(&mut context).is_pending());
+ assert!(harness.published());
+ drop(published);
+ assert!(harness.cancelled_after_publish());
+}
+
+fn target_set() -> TargetSet {
+ TargetSet::new(vec![
+ Target::local("local:conformance-a").expect("first target"),
+ Target::local("local:conformance-b").expect("second target"),
+ ])
+ .expect("target set")
+}
+
+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")
+}
diff --git a/crates/transport/tests/conformance/support.rs b/crates/transport/tests/conformance/support.rs
@@ -0,0 +1,405 @@
+use std::sync::{
+ Arc, Mutex,
+ atomic::{AtomicBool, Ordering},
+};
+
+use futures::future;
+use radroots_transport::{
+ BoxFuture, DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource, FetchPage,
+ FetchRequest, SinkStatus, SourceStatus, TransportId,
+ capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
+ outcome::{DeliveryOutcome, FetchTargetOutcome, FetchTargetState},
+ sink::DeliveryTargetReceipt,
+ source::NextPage,
+};
+
+use crate::suite::{NOW_UNIX_MS, SinkConformanceHarness, SourceConformanceHarness};
+
+#[derive(Clone)]
+enum Mode {
+ Success,
+ Fail(Error),
+ Pending,
+}
+
+#[derive(Default)]
+pub(crate) struct SourceState {
+ request: Mutex<Option<FetchRequest>>,
+ pub(crate) published: AtomicBool,
+ pub(crate) cancelled_after_publish: AtomicBool,
+}
+
+impl SourceState {
+ pub(crate) fn request(&self) -> Option<FetchRequest> {
+ self.request.lock().expect("source request lock").clone()
+ }
+}
+
+#[derive(Default)]
+pub(crate) struct SinkState {
+ request: Mutex<Option<DeliveryRequest>>,
+ pub(crate) published: AtomicBool,
+ pub(crate) cancelled_after_publish: AtomicBool,
+}
+
+impl SinkState {
+ pub(crate) fn request(&self) -> Option<DeliveryRequest> {
+ self.request.lock().expect("sink request lock").clone()
+ }
+}
+
+pub(crate) struct MockSource {
+ mode: Mode,
+ state: Arc<SourceState>,
+}
+
+impl MockSource {
+ pub(crate) fn successful() -> Self {
+ Self::new(Mode::Success)
+ }
+
+ pub(crate) fn failing(error: Error) -> Self {
+ Self::new(Mode::Fail(error))
+ }
+
+ pub(crate) fn pending() -> Self {
+ Self::new(Mode::Pending)
+ }
+
+ fn new(mode: Mode) -> Self {
+ Self {
+ mode,
+ state: Arc::new(SourceState::default()),
+ }
+ }
+}
+
+impl EventSource for MockSource {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> {
+ Box::pin(async { Ok(source_status()) })
+ }
+
+ fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> {
+ let state = Arc::clone(&self.state);
+ let mode = self.mode.clone();
+ Box::pin(async move {
+ *state.request.lock().expect("source request lock") = Some(request.clone());
+ if request.bounds().deadline_unix_ms() <= NOW_UNIX_MS {
+ return Err(Error::InvalidFetchDeadline);
+ }
+ match mode {
+ Mode::Fail(error) => Err(error),
+ Mode::Pending => {
+ state.published.store(true, Ordering::SeqCst);
+ let state_for_drop = Arc::clone(&state);
+ let _publish_guard = ScopeGuard::new(move || {
+ state_for_drop
+ .cancelled_after_publish
+ .store(true, Ordering::SeqCst);
+ });
+ future::pending().await
+ }
+ Mode::Success => {
+ let outcomes = request
+ .target_set()
+ .targets()
+ .iter()
+ .enumerate()
+ .map(|(index, target)| {
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ if index == 0 {
+ FetchTargetState::Complete
+ } else {
+ FetchTargetState::FailedRetryable
+ },
+ )
+ })
+ .collect();
+ FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete)
+ }
+ }
+ })
+ }
+}
+
+impl SourceConformanceHarness for MockSource {
+ fn source(&self) -> &dyn EventSource {
+ self
+ }
+
+ fn expected_status(&self) -> SourceStatus {
+ source_status()
+ }
+
+ fn target_set(&self) -> radroots_transport::TargetSet {
+ target_set()
+ }
+
+ fn now_unix_ms(&self) -> u64 {
+ NOW_UNIX_MS
+ }
+
+ fn captured_request(&self) -> Option<FetchRequest> {
+ self.state.request()
+ }
+
+ fn published(&self) -> bool {
+ self.state.published.load(Ordering::SeqCst)
+ }
+
+ fn cancelled_after_publish(&self) -> bool {
+ self.state.cancelled_after_publish.load(Ordering::SeqCst)
+ }
+}
+
+pub(crate) struct MockSink {
+ mode: Mode,
+ state: Arc<SinkState>,
+}
+
+impl MockSink {
+ pub(crate) fn successful() -> Self {
+ Self::new(Mode::Success)
+ }
+
+ pub(crate) fn failing(error: Error) -> Self {
+ Self::new(Mode::Fail(error))
+ }
+
+ pub(crate) fn pending() -> Self {
+ Self::new(Mode::Pending)
+ }
+
+ fn new(mode: Mode) -> Self {
+ Self {
+ mode,
+ state: Arc::new(SinkState::default()),
+ }
+ }
+}
+
+impl EventSink for MockSink {
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
+ Box::pin(async { Ok(sink_status()) })
+ }
+
+ fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
+ let state = Arc::clone(&self.state);
+ let mode = self.mode.clone();
+ Box::pin(async move {
+ *state.request.lock().expect("sink request lock") = Some(request.clone());
+ if request.deadline_unix_ms() <= NOW_UNIX_MS {
+ return Err(Error::InvalidDeliveryDeadline);
+ }
+ match mode {
+ Mode::Fail(error) => Err(error),
+ Mode::Pending => {
+ state.published.store(true, Ordering::SeqCst);
+ let state_for_drop = Arc::clone(&state);
+ let _publish_guard = ScopeGuard::new(move || {
+ state_for_drop
+ .cancelled_after_publish
+ .store(true, Ordering::SeqCst);
+ });
+ future::pending().await
+ }
+ Mode::Success => {
+ let receipts = request
+ .target_set()
+ .targets()
+ .iter()
+ .enumerate()
+ .map(|(index, target)| {
+ DeliveryTargetReceipt::attempted(
+ target.clone(),
+ if index == 0 {
+ DeliveryOutcome::delivered()
+ } else {
+ DeliveryOutcome::unavailable()
+ },
+ )
+ })
+ .collect();
+ DeliveryReceipt::for_request(&request, receipts)
+ }
+ }
+ })
+ }
+}
+
+impl SinkConformanceHarness for MockSink {
+ fn sink(&self) -> &dyn EventSink {
+ self
+ }
+
+ fn expected_status(&self) -> SinkStatus {
+ sink_status()
+ }
+
+ fn target_set(&self) -> radroots_transport::TargetSet {
+ target_set()
+ }
+
+ fn now_unix_ms(&self) -> u64 {
+ NOW_UNIX_MS
+ }
+
+ fn captured_request(&self) -> Option<DeliveryRequest> {
+ self.state.request()
+ }
+
+ fn published(&self) -> bool {
+ self.state.published.load(Ordering::SeqCst)
+ }
+
+ fn cancelled_after_publish(&self) -> bool {
+ self.state.cancelled_after_publish.load(Ordering::SeqCst)
+ }
+}
+
+struct ScopeGuard<F: FnOnce()>(Option<F>);
+
+impl<F: FnOnce()> ScopeGuard<F> {
+ fn new(callback: F) -> Self {
+ Self(Some(callback))
+ }
+}
+
+impl<F: FnOnce()> Drop for ScopeGuard<F> {
+ fn drop(&mut self) {
+ if let Some(callback) = self.0.take() {
+ callback();
+ }
+ }
+}
+
+pub(crate) struct CombinedAdapter {
+ source: MockSource,
+ sink: MockSink,
+}
+
+impl CombinedAdapter {
+ pub(crate) fn successful() -> Self {
+ Self {
+ source: MockSource::successful(),
+ sink: MockSink::successful(),
+ }
+ }
+}
+
+impl EventSource for CombinedAdapter {
+ fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> {
+ EventSource::status(&self.source)
+ }
+
+ fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> {
+ self.source.fetch(request)
+ }
+}
+
+impl EventSink for CombinedAdapter {
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
+ EventSink::status(&self.sink)
+ }
+
+ fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
+ self.sink.deliver(request)
+ }
+}
+
+impl SourceConformanceHarness for CombinedAdapter {
+ fn source(&self) -> &dyn EventSource {
+ self
+ }
+
+ fn expected_status(&self) -> SourceStatus {
+ source_status()
+ }
+
+ fn target_set(&self) -> radroots_transport::TargetSet {
+ target_set()
+ }
+
+ fn now_unix_ms(&self) -> u64 {
+ NOW_UNIX_MS
+ }
+
+ fn captured_request(&self) -> Option<FetchRequest> {
+ self.source.state.request()
+ }
+
+ fn published(&self) -> bool {
+ self.source.state.published.load(Ordering::SeqCst)
+ }
+
+ fn cancelled_after_publish(&self) -> bool {
+ self.source
+ .state
+ .cancelled_after_publish
+ .load(Ordering::SeqCst)
+ }
+}
+
+impl SinkConformanceHarness for CombinedAdapter {
+ fn sink(&self) -> &dyn EventSink {
+ self
+ }
+
+ fn expected_status(&self) -> SinkStatus {
+ sink_status()
+ }
+
+ fn target_set(&self) -> radroots_transport::TargetSet {
+ target_set()
+ }
+
+ fn now_unix_ms(&self) -> u64 {
+ NOW_UNIX_MS
+ }
+
+ fn captured_request(&self) -> Option<DeliveryRequest> {
+ self.sink.state.request()
+ }
+
+ fn published(&self) -> bool {
+ self.sink.state.published.load(Ordering::SeqCst)
+ }
+
+ fn cancelled_after_publish(&self) -> bool {
+ self.sink
+ .state
+ .cancelled_after_publish
+ .load(Ordering::SeqCst)
+ }
+}
+
+fn target_set() -> radroots_transport::TargetSet {
+ radroots_transport::TargetSet::new(vec![
+ radroots_transport::Target::local("local:conformance-a").expect("first target"),
+ radroots_transport::Target::local("local:conformance-b").expect("second target"),
+ ])
+ .expect("target set")
+}
+
+fn source_status() -> SourceStatus {
+ SourceStatus::new(
+ TransportId::LOCAL,
+ true,
+ Maturity::Stable,
+ Availability::Available,
+ SourceCapabilities::FETCH,
+ "mock source ready",
+ )
+}
+
+fn sink_status() -> SinkStatus {
+ SinkStatus::new(
+ TransportId::LOCAL,
+ true,
+ Maturity::Stable,
+ Availability::Available,
+ SinkCapabilities::DELIVER,
+ "mock sink ready",
+ )
+}