commit 93e2280834d72aceb3c02d9467d73abf864552c6
parent ef9b2138adebd0e620184e122d725517edd55209
Author: triesap <tyson@radroots.org>
Date: Mon, 6 Jul 2026 21:25:03 +0000
runtime: add transport registry workers
- add feature-gated runtime transport registry and adapter dispatch ports
- add bounded queue, delivery worker, lease recovery, and deferred-target handling
- add Reticulum preview adapter registration and verified inbound observation sink contract
- validate runtime, outbox, event_store, and contract lanes
Diffstat:
4 files changed, 843 insertions(+), 0 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -4605,10 +4605,13 @@ dependencies = [
"clap",
"config",
"getrandom 0.2.17",
+ "radroots_events",
"radroots_log",
"radroots_protected_store",
"radroots_runtime_paths",
"radroots_secret_vault",
+ "radroots_transport",
+ "radroots_transport_reticulum",
"serde",
"serde_json",
"tempfile",
diff --git a/crates/runtime/Cargo.toml b/crates/runtime/Cargo.toml
@@ -15,6 +15,10 @@ readme = "README"
[features]
default = []
cli = ["dep:clap"]
+transport = ["dep:radroots_events", "dep:radroots_transport"]
+transport-nostr = ["transport"]
+transport-reticulum = ["transport", "dep:radroots_transport_reticulum"]
+transport-workers = ["transport"]
[dependencies]
anyhow = { workspace = true }
@@ -23,9 +27,17 @@ clap = { workspace = true, features = ["derive", "env"], optional = true }
config = { workspace = true }
getrandom = { workspace = true }
radroots_log = { workspace = true, features = ["std"] }
+radroots_events = { workspace = true, optional = true, default-features = false, features = [
+ "std",
+ "serde",
+] }
radroots_protected_store = { workspace = true, features = ["std"] }
radroots_runtime_paths = { workspace = true }
radroots_secret_vault = { workspace = true, features = ["std"] }
+radroots_transport = { workspace = true, optional = true, default-features = false, features = [
+ "serde",
+] }
+radroots_transport_reticulum = { workspace = true, optional = true, default-features = false }
serde = { workspace = true }
serde_json = { workspace = true }
tempfile = { workspace = true }
diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs
@@ -8,6 +8,8 @@ pub mod secret_file;
pub mod service;
pub mod signals;
pub mod tracing;
+#[cfg(feature = "transport")]
+pub mod transport;
#[cfg(feature = "cli")]
pub use cli::{parse_and_load_path, parse_and_load_path_with_env_overrides};
@@ -40,3 +42,20 @@ pub use signals::shutdown_signal;
pub use tracing::{
default_shared_runtime_logs_dir, default_shared_runtime_logs_dir_for, init, init_with_logs_dir,
};
+#[cfg(feature = "transport-reticulum")]
+pub use transport::RadrootsRuntimeReticulumPreviewTransport;
+#[cfg(feature = "transport")]
+pub use transport::{
+ RadrootsRuntimeBoundedQueue, RadrootsRuntimeQueueStatus, RadrootsRuntimeQueueTask,
+ RadrootsRuntimeTransportAdapter, RadrootsRuntimeTransportDispatchRequest,
+ RadrootsRuntimeTransportError, RadrootsRuntimeTransportFuture, RadrootsRuntimeTransportPayload,
+ RadrootsRuntimeTransportRegistry,
+};
+#[cfg(feature = "transport-workers")]
+pub use transport::{
+ RadrootsRuntimeDeliveryJob, RadrootsRuntimeDeliveryJobReceipt, RadrootsRuntimeDeliveryPlan,
+ RadrootsRuntimeDeliveryTarget, RadrootsRuntimeDeliveryWorker,
+ RadrootsRuntimeDeliveryWorkerConfig, RadrootsRuntimeInboundObservation,
+ RadrootsRuntimeInboundObservationSink, RadrootsRuntimeLeaseRecord,
+ record_verified_inbound_observation, recover_expired_leases,
+};
diff --git a/crates/runtime/src/transport.rs b/crates/runtime/src/transport.rs
@@ -0,0 +1,809 @@
+use std::collections::{BTreeMap, VecDeque};
+use std::future::Future;
+use std::pin::Pin;
+use std::sync::Arc;
+
+use radroots_events::draft::RadrootsSignedNostrEvent;
+#[cfg(feature = "transport-workers")]
+use radroots_transport::RadrootsTransportTargetReceipt;
+use radroots_transport::{
+ RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
+ RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportKind,
+ RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetSet,
+};
+use thiserror::Error;
+
+pub type RadrootsRuntimeTransportFuture<'a, T> =
+ Pin<Box<dyn Future<Output = Result<T, RadrootsRuntimeTransportError>> + Send + 'a>>;
+
+#[derive(Debug, Error)]
+pub enum RadrootsRuntimeTransportError {
+ #[error("transport adapter for `{0}` is not registered")]
+ AdapterNotRegistered(String),
+
+ #[error("transport adapter for `{0}` is already registered")]
+ AdapterAlreadyRegistered(String),
+
+ #[error("transport dispatch received no targets")]
+ EmptyDispatchTargets,
+
+ #[error("transport target error: {0}")]
+ TransportTarget(String),
+
+ #[error("runtime queue capacity must be greater than zero")]
+ InvalidQueueCapacity,
+
+ #[error("inbound observation for event `{event_id}` is not verified")]
+ InboundObservationUnverified { event_id: String },
+
+ #[error(
+ "inbound observation event `{observation_event_id}` does not match signed event `{signed_event_id}`"
+ )]
+ InboundObservationEventMismatch {
+ observation_event_id: String,
+ signed_event_id: String,
+ },
+
+ #[error("transport adapter `{kind}` failed: {message}")]
+ Adapter { kind: String, message: String },
+}
+
+impl From<RadrootsTransportError> for RadrootsRuntimeTransportError {
+ fn from(value: RadrootsTransportError) -> Self {
+ Self::TransportTarget(value.to_string())
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub enum RadrootsRuntimeTransportPayload {
+ SignedNostrEvent(RadrootsSignedNostrEvent),
+ DigestOnly(String),
+}
+
+impl RadrootsRuntimeTransportPayload {
+ pub fn digest(&self) -> String {
+ match self {
+ Self::SignedNostrEvent(event) => event.id.clone(),
+ Self::DigestOnly(value) => value.clone(),
+ }
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeTransportDispatchRequest {
+ pub request_id: String,
+ pub payload: RadrootsRuntimeTransportPayload,
+ pub target_set: RadrootsTransportTargetSet,
+ pub satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+ pub now_ms: i64,
+}
+
+impl RadrootsRuntimeTransportDispatchRequest {
+ pub fn new(
+ request_id: impl Into<String>,
+ payload: RadrootsRuntimeTransportPayload,
+ targets: Vec<RadrootsTransportTarget>,
+ satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+ now_ms: i64,
+ ) -> Result<Self, RadrootsRuntimeTransportError> {
+ if targets.is_empty() {
+ return Err(RadrootsRuntimeTransportError::EmptyDispatchTargets);
+ }
+ let target_set = RadrootsTransportTargetSet::new(targets)?;
+ Ok(Self {
+ request_id: request_id.into(),
+ payload,
+ target_set,
+ satisfaction_policy,
+ now_ms,
+ })
+ }
+
+ pub fn transport_delivery_request(&self) -> RadrootsTransportDeliveryRequest {
+ RadrootsTransportDeliveryRequest::new(
+ self.request_id.clone(),
+ self.payload.digest(),
+ self.target_set.clone(),
+ self.satisfaction_policy.clone(),
+ )
+ }
+}
+
+pub trait RadrootsRuntimeTransportAdapter: Send + Sync {
+ fn transport_kind(&self) -> RadrootsTransportKind;
+
+ fn deliver<'a>(
+ &'a self,
+ request: RadrootsRuntimeTransportDispatchRequest,
+ ) -> RadrootsRuntimeTransportFuture<'a, RadrootsTransportDeliveryReceipt>;
+}
+
+#[derive(Clone, Default)]
+pub struct RadrootsRuntimeTransportRegistry {
+ adapters: BTreeMap<RadrootsTransportKind, Arc<dyn RadrootsRuntimeTransportAdapter>>,
+}
+
+impl RadrootsRuntimeTransportRegistry {
+ pub fn new() -> Self {
+ Self::default()
+ }
+
+ pub fn register<A>(&mut self, adapter: A) -> Result<(), RadrootsRuntimeTransportError>
+ where
+ A: RadrootsRuntimeTransportAdapter + 'static,
+ {
+ let kind = adapter.transport_kind();
+ if self.adapters.contains_key(&kind) {
+ return Err(RadrootsRuntimeTransportError::AdapterAlreadyRegistered(
+ kind.canonical_label(),
+ ));
+ }
+ self.adapters.insert(kind, Arc::new(adapter));
+ Ok(())
+ }
+
+ pub fn adapter(
+ &self,
+ kind: &RadrootsTransportKind,
+ ) -> Result<Arc<dyn RadrootsRuntimeTransportAdapter>, RadrootsRuntimeTransportError> {
+ self.adapters.get(kind).cloned().ok_or_else(|| {
+ RadrootsRuntimeTransportError::AdapterNotRegistered(kind.canonical_label())
+ })
+ }
+
+ pub fn registered_kinds(&self) -> Vec<RadrootsTransportKind> {
+ self.adapters.keys().cloned().collect()
+ }
+}
+
+#[cfg(feature = "transport-reticulum")]
+#[derive(Clone, Debug)]
+pub struct RadrootsRuntimeReticulumPreviewTransport {
+ transport: radroots_transport_reticulum::RadrootsReticulumPreviewTransport,
+}
+
+#[cfg(feature = "transport-reticulum")]
+impl RadrootsRuntimeReticulumPreviewTransport {
+ pub fn new(transport: radroots_transport_reticulum::RadrootsReticulumPreviewTransport) -> Self {
+ Self { transport }
+ }
+}
+
+#[cfg(feature = "transport-reticulum")]
+impl Default for RadrootsRuntimeReticulumPreviewTransport {
+ fn default() -> Self {
+ Self::new(radroots_transport_reticulum::RadrootsReticulumPreviewTransport::default())
+ }
+}
+
+#[cfg(feature = "transport-reticulum")]
+impl RadrootsRuntimeTransportAdapter for RadrootsRuntimeReticulumPreviewTransport {
+ fn transport_kind(&self) -> RadrootsTransportKind {
+ RadrootsTransportKind::Reticulum
+ }
+
+ fn deliver<'a>(
+ &'a self,
+ request: RadrootsRuntimeTransportDispatchRequest,
+ ) -> RadrootsRuntimeTransportFuture<'a, RadrootsTransportDeliveryReceipt> {
+ Box::pin(async move {
+ self.transport
+ .deliver(request.transport_delivery_request())
+ .map_err(|error| RadrootsRuntimeTransportError::Adapter {
+ kind: RadrootsTransportKind::Reticulum.canonical_label(),
+ message: error.to_string(),
+ })
+ })
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeQueueStatus {
+ pub capacity: usize,
+ pub queued: usize,
+ pub in_flight: usize,
+ pub shutdown: bool,
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeQueueTask<T> {
+ pub sequence: u64,
+ pub payload: T,
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeBoundedQueue<T> {
+ capacity: usize,
+ next_sequence: u64,
+ in_flight: usize,
+ shutdown: bool,
+ queue: VecDeque<RadrootsRuntimeQueueTask<T>>,
+}
+
+impl<T> RadrootsRuntimeBoundedQueue<T> {
+ pub fn new(capacity: usize) -> Result<Self, RadrootsRuntimeTransportError> {
+ if capacity == 0 {
+ return Err(RadrootsRuntimeTransportError::InvalidQueueCapacity);
+ }
+ Ok(Self {
+ capacity,
+ next_sequence: 1,
+ in_flight: 0,
+ shutdown: false,
+ queue: VecDeque::new(),
+ })
+ }
+
+ pub fn try_enqueue(&mut self, payload: T) -> Option<u64> {
+ if self.shutdown || self.queue.len() >= self.capacity {
+ return None;
+ }
+ let sequence = self.next_sequence;
+ self.next_sequence += 1;
+ self.queue
+ .push_back(RadrootsRuntimeQueueTask { sequence, payload });
+ Some(sequence)
+ }
+
+ pub fn pop(&mut self) -> Option<RadrootsRuntimeQueueTask<T>> {
+ let task = self.queue.pop_front()?;
+ self.in_flight += 1;
+ Some(task)
+ }
+
+ pub fn complete_task(&mut self) {
+ self.in_flight = self.in_flight.saturating_sub(1);
+ }
+
+ pub fn shutdown(&mut self) {
+ self.shutdown = true;
+ }
+
+ pub fn status(&self) -> RadrootsRuntimeQueueStatus {
+ RadrootsRuntimeQueueStatus {
+ capacity: self.capacity,
+ queued: self.queue.len(),
+ in_flight: self.in_flight,
+ shutdown: self.shutdown,
+ }
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeDeliveryTarget {
+ pub delivery_target_id: i64,
+ pub target: RadrootsTransportTarget,
+ pub status: RadrootsTransportDeliveryTargetStatus,
+}
+
+impl RadrootsRuntimeDeliveryTarget {
+ pub fn ready(delivery_target_id: i64, target: RadrootsTransportTarget) -> Self {
+ Self {
+ delivery_target_id,
+ target,
+ status: RadrootsTransportDeliveryTargetStatus::Pending,
+ }
+ }
+
+ pub fn deferred_until_implemented(
+ delivery_target_id: i64,
+ target: RadrootsTransportTarget,
+ ) -> Self {
+ Self {
+ delivery_target_id,
+ target,
+ status: RadrootsTransportDeliveryTargetStatus::Deferred,
+ }
+ }
+
+ pub fn is_ready_for_attempt(&self) -> bool {
+ self.status == RadrootsTransportDeliveryTargetStatus::Pending
+ || self.status == RadrootsTransportDeliveryTargetStatus::Failed
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeDeliveryPlan {
+ pub delivery_plan_id: i64,
+ pub satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+ pub targets: Vec<RadrootsRuntimeDeliveryTarget>,
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeDeliveryJob {
+ pub outbox_event_id: i64,
+ pub payload: RadrootsRuntimeTransportPayload,
+ pub plans: Vec<RadrootsRuntimeDeliveryPlan>,
+ pub now_ms: i64,
+}
+
+#[cfg(feature = "transport-workers")]
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeDeliveryWorkerConfig {
+ pub bounded_queue_capacity: usize,
+}
+
+#[cfg(feature = "transport-workers")]
+impl Default for RadrootsRuntimeDeliveryWorkerConfig {
+ fn default() -> Self {
+ Self {
+ bounded_queue_capacity: 64,
+ }
+ }
+}
+
+#[cfg(feature = "transport-workers")]
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeDeliveryJobReceipt {
+ pub outbox_event_id: i64,
+ pub dispatch_count: usize,
+ pub target_receipts: Vec<RadrootsTransportTargetReceipt>,
+}
+
+#[cfg(feature = "transport-workers")]
+pub struct RadrootsRuntimeDeliveryWorker<'a> {
+ registry: &'a RadrootsRuntimeTransportRegistry,
+ config: RadrootsRuntimeDeliveryWorkerConfig,
+}
+
+#[cfg(feature = "transport-workers")]
+impl<'a> RadrootsRuntimeDeliveryWorker<'a> {
+ pub fn new(
+ registry: &'a RadrootsRuntimeTransportRegistry,
+ config: RadrootsRuntimeDeliveryWorkerConfig,
+ ) -> Self {
+ Self { registry, config }
+ }
+
+ pub fn queue_capacity(&self) -> usize {
+ self.config.bounded_queue_capacity
+ }
+
+ pub async fn execute_job(
+ &self,
+ job: RadrootsRuntimeDeliveryJob,
+ ) -> Result<RadrootsRuntimeDeliveryJobReceipt, RadrootsRuntimeTransportError> {
+ let mut target_receipts = Vec::new();
+ let mut dispatch_count = 0usize;
+ for plan in job.plans {
+ let mut by_kind =
+ BTreeMap::<RadrootsTransportKind, Vec<RadrootsRuntimeDeliveryTarget>>::new();
+ for target in plan
+ .targets
+ .into_iter()
+ .filter(RadrootsRuntimeDeliveryTarget::is_ready_for_attempt)
+ {
+ by_kind
+ .entry(target.target.kind.clone())
+ .or_default()
+ .push(target);
+ }
+ for (kind, targets) in by_kind {
+ let adapter = self.registry.adapter(&kind)?;
+ let transport_targets = targets
+ .iter()
+ .map(|target| target.target.clone())
+ .collect::<Vec<_>>();
+ let request = RadrootsRuntimeTransportDispatchRequest::new(
+ format!(
+ "outbox-event-{}-plan-{}",
+ job.outbox_event_id, plan.delivery_plan_id
+ ),
+ job.payload.clone(),
+ transport_targets,
+ plan.satisfaction_policy.clone(),
+ job.now_ms,
+ )?;
+ let receipt = adapter.deliver(request).await?;
+ target_receipts.extend(receipt.target_receipts);
+ dispatch_count += 1;
+ }
+ }
+ Ok(RadrootsRuntimeDeliveryJobReceipt {
+ outbox_event_id: job.outbox_event_id,
+ dispatch_count,
+ target_receipts,
+ })
+ }
+}
+
+#[cfg(feature = "transport-workers")]
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeLeaseRecord {
+ pub record_id: i64,
+ pub claimed: bool,
+ pub claim_expires_at_ms: i64,
+ pub recovered: bool,
+}
+
+#[cfg(feature = "transport-workers")]
+pub fn recover_expired_leases(leases: &mut [RadrootsRuntimeLeaseRecord], now_ms: i64) -> usize {
+ let mut recovered = 0usize;
+ for lease in leases {
+ if lease.claimed && lease.claim_expires_at_ms <= now_ms {
+ lease.claimed = false;
+ lease.recovered = true;
+ recovered += 1;
+ }
+ }
+ recovered
+}
+
+#[cfg(feature = "transport-workers")]
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRuntimeInboundObservation {
+ pub event_id: String,
+ pub verified: bool,
+ pub transport_kind: RadrootsTransportKind,
+ pub endpoint_uri: String,
+ pub observed_at_ms: i64,
+}
+
+#[cfg(feature = "transport-workers")]
+impl RadrootsRuntimeInboundObservation {
+ pub fn verified_signed_event(
+ event: &RadrootsSignedNostrEvent,
+ transport_kind: RadrootsTransportKind,
+ endpoint_uri: impl Into<String>,
+ observed_at_ms: i64,
+ ) -> Self {
+ Self {
+ event_id: event.id.clone(),
+ verified: true,
+ transport_kind,
+ endpoint_uri: endpoint_uri.into(),
+ observed_at_ms,
+ }
+ }
+
+ pub fn require_verified_for_signed_event(
+ &self,
+ event: &RadrootsSignedNostrEvent,
+ ) -> Result<(), RadrootsRuntimeTransportError> {
+ if !self.verified {
+ return Err(
+ RadrootsRuntimeTransportError::InboundObservationUnverified {
+ event_id: self.event_id.clone(),
+ },
+ );
+ }
+ if self.event_id != event.id {
+ return Err(
+ RadrootsRuntimeTransportError::InboundObservationEventMismatch {
+ observation_event_id: self.event_id.clone(),
+ signed_event_id: event.id.clone(),
+ },
+ );
+ }
+ Ok(())
+ }
+}
+
+#[cfg(feature = "transport-workers")]
+pub trait RadrootsRuntimeInboundObservationSink: Send + Sync {
+ fn record_verified_observation<'a>(
+ &'a self,
+ event: RadrootsSignedNostrEvent,
+ observation: RadrootsRuntimeInboundObservation,
+ ) -> RadrootsRuntimeTransportFuture<'a, ()>;
+}
+
+#[cfg(feature = "transport-workers")]
+pub fn record_verified_inbound_observation<'a, S>(
+ sink: &'a S,
+ event: RadrootsSignedNostrEvent,
+ observation: RadrootsRuntimeInboundObservation,
+) -> RadrootsRuntimeTransportFuture<'a, ()>
+where
+ S: RadrootsRuntimeInboundObservationSink + ?Sized,
+{
+ Box::pin(async move {
+ observation.require_verified_for_signed_event(&event)?;
+ sink.record_verified_observation(event, observation).await
+ })
+}
+
+#[cfg(test)]
+mod tests {
+ use super::{
+ RadrootsRuntimeBoundedQueue, RadrootsRuntimeTransportAdapter,
+ RadrootsRuntimeTransportDispatchRequest, RadrootsRuntimeTransportError,
+ RadrootsRuntimeTransportFuture, RadrootsRuntimeTransportPayload,
+ RadrootsRuntimeTransportRegistry,
+ };
+ #[cfg(feature = "transport-workers")]
+ use super::{
+ RadrootsRuntimeDeliveryJob, RadrootsRuntimeDeliveryPlan, RadrootsRuntimeDeliveryTarget,
+ RadrootsRuntimeDeliveryWorker, RadrootsRuntimeDeliveryWorkerConfig,
+ RadrootsRuntimeInboundObservation, RadrootsRuntimeInboundObservationSink,
+ RadrootsRuntimeLeaseRecord, record_verified_inbound_observation, recover_expired_leases,
+ };
+ #[cfg(feature = "transport-workers")]
+ use radroots_events::draft::{RadrootsSignedNostrEvent, RadrootsSignedNostrEventParts};
+ use radroots_transport::{
+ RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryTargetStatus,
+ RadrootsTransportKind, RadrootsTransportOutcome, RadrootsTransportSatisfactionPolicy,
+ RadrootsTransportTarget, RadrootsTransportTargetReceipt,
+ };
+
+ struct StaticAdapter {
+ kind: RadrootsTransportKind,
+ status: RadrootsTransportDeliveryTargetStatus,
+ }
+
+ impl RadrootsRuntimeTransportAdapter for StaticAdapter {
+ fn transport_kind(&self) -> RadrootsTransportKind {
+ self.kind.clone()
+ }
+
+ fn deliver<'a>(
+ &'a self,
+ request: RadrootsRuntimeTransportDispatchRequest,
+ ) -> RadrootsRuntimeTransportFuture<'a, RadrootsTransportDeliveryReceipt> {
+ Box::pin(async move {
+ Ok(RadrootsTransportDeliveryReceipt {
+ request_id: request.request_id,
+ target_receipts: request
+ .target_set
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| {
+ RadrootsTransportTargetReceipt::new(
+ target,
+ RadrootsTransportOutcome::new(self.status),
+ )
+ })
+ .collect(),
+ })
+ })
+ }
+ }
+
+ fn target(kind: RadrootsTransportKind, uri: &str) -> RadrootsTransportTarget {
+ RadrootsTransportTarget::new(kind, uri).expect("target")
+ }
+
+ #[cfg(feature = "transport-workers")]
+ fn signed_event() -> RadrootsSignedNostrEvent {
+ RadrootsSignedNostrEvent::new(RadrootsSignedNostrEventParts {
+ id: "d".repeat(64),
+ pubkey: "e".repeat(64),
+ created_at: 10,
+ kind: 1,
+ tags: Vec::new(),
+ content: "hello".to_owned(),
+ sig: "f".repeat(128),
+ raw_json: "{\"id\":\"fixture\"}".to_owned(),
+ })
+ .expect("signed event")
+ }
+
+ #[cfg(feature = "transport-workers")]
+ struct RecordingInboundSink {
+ expected_event_id: String,
+ }
+
+ #[cfg(feature = "transport-workers")]
+ impl RadrootsRuntimeInboundObservationSink for RecordingInboundSink {
+ fn record_verified_observation<'a>(
+ &'a self,
+ event: RadrootsSignedNostrEvent,
+ observation: RadrootsRuntimeInboundObservation,
+ ) -> RadrootsRuntimeTransportFuture<'a, ()> {
+ Box::pin(async move {
+ assert_eq!(event.id, self.expected_event_id);
+ assert_eq!(observation.event_id, self.expected_event_id);
+ assert!(observation.verified);
+ Ok(())
+ })
+ }
+ }
+
+ #[tokio::test]
+ async fn registry_dispatches_nostr_adapter_by_transport_kind() {
+ let mut registry = RadrootsRuntimeTransportRegistry::new();
+ registry
+ .register(StaticAdapter {
+ kind: RadrootsTransportKind::Nostr,
+ status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ })
+ .expect("register");
+ assert_eq!(
+ registry.registered_kinds(),
+ vec![RadrootsTransportKind::Nostr]
+ );
+ assert!(matches!(
+ registry.register(StaticAdapter {
+ kind: RadrootsTransportKind::Nostr,
+ status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ }),
+ Err(RadrootsRuntimeTransportError::AdapterAlreadyRegistered(_))
+ ));
+
+ let adapter = registry
+ .adapter(&RadrootsTransportKind::Nostr)
+ .expect("nostr adapter");
+ let receipt = adapter
+ .deliver(
+ RadrootsRuntimeTransportDispatchRequest::new(
+ "nostr-delivery",
+ RadrootsRuntimeTransportPayload::DigestOnly("sha256:event".to_owned()),
+ vec![target(RadrootsTransportKind::Nostr, "wss://relay.example")],
+ RadrootsTransportSatisfactionPolicy::AnyTarget,
+ 1_000,
+ )
+ .expect("request"),
+ )
+ .await
+ .expect("receipt");
+
+ assert_eq!(receipt.satisfied_target_count(), 1);
+ assert_eq!(
+ receipt.target_receipts[0].status,
+ RadrootsTransportDeliveryTargetStatus::Accepted
+ );
+ }
+
+ #[cfg(feature = "transport-reticulum")]
+ #[tokio::test]
+ async fn registry_dispatches_reticulum_preview_without_success() {
+ let mut registry = RadrootsRuntimeTransportRegistry::new();
+ registry
+ .register(super::RadrootsRuntimeReticulumPreviewTransport::default())
+ .expect("register");
+ let adapter = registry
+ .adapter(&RadrootsTransportKind::Reticulum)
+ .expect("reticulum adapter");
+ let receipt = adapter
+ .deliver(
+ RadrootsRuntimeTransportDispatchRequest::new(
+ "reticulum-delivery",
+ RadrootsRuntimeTransportPayload::DigestOnly("sha256:event".to_owned()),
+ vec![target(
+ RadrootsTransportKind::Reticulum,
+ "reticulum:preview",
+ )],
+ RadrootsTransportSatisfactionPolicy::AnyTarget,
+ 1_000,
+ )
+ .expect("request"),
+ )
+ .await
+ .expect("receipt");
+
+ assert_eq!(receipt.satisfied_target_count(), 0);
+ assert_eq!(
+ receipt.target_receipts[0].status,
+ RadrootsTransportDeliveryTargetStatus::Unavailable
+ );
+ }
+
+ #[test]
+ fn bounded_queue_tracks_capacity_inflight_and_shutdown() {
+ let mut queue = RadrootsRuntimeBoundedQueue::new(2).expect("queue");
+ assert_eq!(queue.try_enqueue("a"), Some(1));
+ assert_eq!(queue.try_enqueue("b"), Some(2));
+ assert_eq!(queue.try_enqueue("c"), None);
+ let task = queue.pop().expect("task");
+ assert_eq!(task.sequence, 1);
+ assert_eq!(queue.status().in_flight, 1);
+ queue.complete_task();
+ queue.shutdown();
+ assert_eq!(queue.try_enqueue("d"), None);
+ assert!(queue.status().shutdown);
+ }
+
+ #[cfg(feature = "transport-workers")]
+ #[tokio::test]
+ async fn delivery_worker_dispatches_ready_targets_and_skips_deferred() {
+ let mut registry = RadrootsRuntimeTransportRegistry::new();
+ registry
+ .register(StaticAdapter {
+ kind: RadrootsTransportKind::Nostr,
+ status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ })
+ .expect("register");
+ let worker = RadrootsRuntimeDeliveryWorker::new(
+ ®istry,
+ RadrootsRuntimeDeliveryWorkerConfig {
+ bounded_queue_capacity: 8,
+ },
+ );
+ let ready = RadrootsRuntimeDeliveryTarget::ready(
+ 1,
+ target(RadrootsTransportKind::Nostr, "wss://relay.example"),
+ );
+ let deferred = RadrootsRuntimeDeliveryTarget::deferred_until_implemented(
+ 2,
+ target(RadrootsTransportKind::Reticulum, "reticulum:preview"),
+ );
+ let receipt = worker
+ .execute_job(RadrootsRuntimeDeliveryJob {
+ outbox_event_id: 42,
+ payload: RadrootsRuntimeTransportPayload::DigestOnly("sha256:event".to_owned()),
+ plans: vec![RadrootsRuntimeDeliveryPlan {
+ delivery_plan_id: 7,
+ satisfaction_policy: RadrootsTransportSatisfactionPolicy::AnyTarget,
+ targets: vec![ready, deferred],
+ }],
+ now_ms: 1_000,
+ })
+ .await
+ .expect("worker receipt");
+
+ assert_eq!(worker.queue_capacity(), 8);
+ assert_eq!(receipt.dispatch_count, 1);
+ assert_eq!(receipt.target_receipts.len(), 1);
+ assert_eq!(
+ receipt.target_receipts[0].status,
+ RadrootsTransportDeliveryTargetStatus::Accepted
+ );
+ }
+
+ #[cfg(feature = "transport-workers")]
+ #[test]
+ fn lease_recovery_marks_expired_claims() {
+ let mut leases = vec![
+ RadrootsRuntimeLeaseRecord {
+ record_id: 1,
+ claimed: true,
+ claim_expires_at_ms: 900,
+ recovered: false,
+ },
+ RadrootsRuntimeLeaseRecord {
+ record_id: 2,
+ claimed: true,
+ claim_expires_at_ms: 1_100,
+ recovered: false,
+ },
+ ];
+ assert_eq!(recover_expired_leases(&mut leases, 1_000), 1);
+ assert!(!leases[0].claimed);
+ assert!(leases[0].recovered);
+ assert!(leases[1].claimed);
+ }
+
+ #[cfg(feature = "transport-workers")]
+ #[tokio::test]
+ async fn inbound_observation_sink_requires_verified_signed_events() {
+ let event = signed_event();
+ let sink = RecordingInboundSink {
+ expected_event_id: event.id.clone(),
+ };
+
+ let observation = RadrootsRuntimeInboundObservation::verified_signed_event(
+ &event,
+ RadrootsTransportKind::Nostr,
+ "wss://relay.example",
+ 1_000,
+ );
+ record_verified_inbound_observation(&sink, event.clone(), observation)
+ .await
+ .expect("record observation");
+
+ let unverified = RadrootsRuntimeInboundObservation {
+ event_id: event.id.clone(),
+ verified: false,
+ transport_kind: RadrootsTransportKind::Nostr,
+ endpoint_uri: "wss://relay.example".to_owned(),
+ observed_at_ms: 1_000,
+ };
+ assert!(matches!(
+ record_verified_inbound_observation(&sink, event.clone(), unverified).await,
+ Err(RadrootsRuntimeTransportError::InboundObservationUnverified { .. })
+ ));
+
+ let mismatched = RadrootsRuntimeInboundObservation {
+ event_id: "a".repeat(64),
+ verified: true,
+ transport_kind: RadrootsTransportKind::Nostr,
+ endpoint_uri: "wss://relay.example".to_owned(),
+ observed_at_ms: 1_000,
+ };
+ assert!(matches!(
+ record_verified_inbound_observation(&sink, event, mismatched).await,
+ Err(RadrootsRuntimeTransportError::InboundObservationEventMismatch { .. })
+ ));
+ }
+}