commit e02a52881616c1c548dd36e2127b51eac6c8414f
parent fffa131884a7cf0d83fcbb7bc235684ca36d05c8
Author: triesap <tyson@radroots.org>
Date: Thu, 9 Jul 2026 19:49:21 +0000
transport: add neutral outcome contracts
- add typed neutral transport outcome kinds and transport trait surface
- add required-target satisfaction policy with fingerprint-aware receipt evaluation
- map relay outcomes into neutral transport outcomes and harden publisher satisfaction
- align runtime, outbox, and Reticulum preview paths with typed outcomes
Diffstat:
13 files changed, 1000 insertions(+), 74 deletions(-)
diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs
@@ -2270,6 +2270,17 @@ fn satisfaction_policy_storage_value(policy: &RadrootsTransportSatisfactionPolic
satisfaction_class_storage_value(*class)
)
}
+ RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => {
+ let fingerprints = targets
+ .iter()
+ .map(RadrootsTransportTargetFingerprint::as_str)
+ .collect::<Vec<_>>()
+ .join(",");
+ format!(
+ "required_{}:{fingerprints}",
+ satisfaction_class_storage_value(*class)
+ )
+ }
}
}
@@ -2302,6 +2313,18 @@ fn parse_satisfaction_policy(
required_count_u16(required_success_count)?,
))
}
+ stored if stored.starts_with("required_accepted:") => parse_required_target_policy(
+ stored,
+ "required_accepted:",
+ RadrootsTransportSatisfactionClass::Accepted,
+ required_success_count,
+ ),
+ stored if stored.starts_with("required_delivered:") => parse_required_target_policy(
+ stored,
+ "required_delivered:",
+ RadrootsTransportSatisfactionClass::Delivered,
+ required_success_count,
+ ),
_ => Err(RadrootsOutboxError::InvalidStoredEnum {
field: "outbox_delivery_plan.satisfaction_policy",
value: value.to_owned(),
@@ -2309,6 +2332,34 @@ fn parse_satisfaction_policy(
}
}
+fn parse_required_target_policy(
+ stored: &str,
+ prefix: &str,
+ class: RadrootsTransportSatisfactionClass,
+ required_success_count: i64,
+) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsOutboxError> {
+ let targets = stored
+ .strip_prefix(prefix)
+ .expect("stored required target policy prefix")
+ .split(',')
+ .map(RadrootsTransportTargetFingerprint::parse)
+ .collect::<Result<Vec<_>, _>>()?;
+ if required_success_count
+ != i64::try_from(targets.len()).map_err(|_| RadrootsOutboxError::IntegerRange {
+ field: "required_success_count",
+ value: required_success_count,
+ })?
+ {
+ return Err(RadrootsOutboxError::InvalidStoredEnum {
+ field: "outbox_delivery_plan.satisfaction_policy",
+ value: stored.to_owned(),
+ });
+ }
+ Ok(RadrootsTransportSatisfactionPolicy::required_targets(
+ class, targets,
+ )?)
+}
+
fn required_count_u16(required_success_count: i64) -> Result<u16, RadrootsOutboxError> {
u16::try_from(required_success_count).map_err(|_| RadrootsOutboxError::IntegerRange {
field: "required_success_count",
@@ -2368,6 +2419,35 @@ mod tests {
RadrootsTransportTarget::new(RadrootsTransportKind::Proxy, uri).expect("proxy target")
}
+ #[test]
+ fn required_target_satisfaction_policy_storage_round_trips() {
+ let first = nostr_target("wss://required-one.example");
+ let second = nostr_target("wss://required-two.example");
+ let policy = RadrootsTransportSatisfactionPolicy::required_targets(
+ RadrootsTransportSatisfactionClass::Delivered,
+ vec![first.fingerprint.clone(), second.fingerprint.clone()],
+ )
+ .expect("required targets policy");
+
+ let stored = satisfaction_policy_storage_value(&policy);
+ assert_eq!(
+ stored,
+ format!(
+ "required_delivered:{},{}",
+ first.fingerprint.as_str(),
+ second.fingerprint.as_str()
+ )
+ );
+ assert_eq!(
+ parse_satisfaction_policy(stored.as_str(), 2).expect("parse required targets"),
+ policy
+ );
+ assert!(matches!(
+ parse_satisfaction_policy(stored.as_str(), 1),
+ Err(RadrootsOutboxError::InvalidStoredEnum { .. })
+ ));
+ }
+
fn malformed_reticulum_target(uri: &str) -> RadrootsTransportTarget {
let endpoint_uri = RadrootsTransportTargetUri::parse(uri).expect("target uri");
let endpoint_fingerprint = RadrootsTransportTargetFingerprint::from_target(
diff --git a/crates/runtime/src/transport.rs b/crates/runtime/src/transport.rs
@@ -333,6 +333,13 @@ impl Default for RadrootsRuntimeDeliveryWorkerConfig {
#[cfg(feature = "transport-workers")]
#[derive(Clone, Debug, PartialEq, Eq)]
+struct RadrootsRuntimeDeliveryTargetState {
+ target: RadrootsTransportTarget,
+ status: RadrootsTransportDeliveryTargetStatus,
+}
+
+#[cfg(feature = "transport-workers")]
+#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RadrootsRuntimeDeliveryPlanSatisfactionState {
Satisfied,
Unsatisfied,
@@ -353,34 +360,26 @@ pub struct RadrootsRuntimeDeliveryPlanReceipt {
#[cfg(feature = "transport-workers")]
impl RadrootsRuntimeDeliveryPlanReceipt {
- fn from_statuses(
+ fn from_target_states(
delivery_plan_id: i64,
satisfaction_policy: RadrootsTransportSatisfactionPolicy,
- target_statuses: &[RadrootsTransportDeliveryTargetStatus],
+ target_states: &[RadrootsRuntimeDeliveryTargetState],
target_receipts: Vec<RadrootsTransportTargetReceipt>,
) -> Result<Self, RadrootsRuntimeTransportError> {
let required_target_count =
- satisfaction_policy.required_target_count(target_statuses.len())?;
+ satisfaction_policy.required_target_count(target_states.len())?;
let satisfied_target_count =
- satisfaction_policy
- .target_satisfaction_class()
- .map_or(0, |satisfaction_class| {
- target_statuses
- .iter()
- .filter(|status| status.counts_as_satisfied(satisfaction_class))
- .count()
- });
- let satisfaction_state = if satisfaction_policy
- .is_satisfied_by(target_statuses.len(), satisfied_target_count)?
- {
- RadrootsRuntimeDeliveryPlanSatisfactionState::Satisfied
- } else {
- RadrootsRuntimeDeliveryPlanSatisfactionState::Unsatisfied
- };
+ satisfied_target_count_for_policy(&satisfaction_policy, target_states);
+ let satisfaction_state =
+ if target_states_satisfy_policy(&satisfaction_policy, target_states)? {
+ RadrootsRuntimeDeliveryPlanSatisfactionState::Satisfied
+ } else {
+ RadrootsRuntimeDeliveryPlanSatisfactionState::Unsatisfied
+ };
Ok(Self {
delivery_plan_id,
satisfaction_policy,
- target_count: target_statuses.len(),
+ target_count: target_states.len(),
attempted_target_count: target_receipts.len(),
required_target_count,
satisfied_target_count,
@@ -395,6 +394,51 @@ impl RadrootsRuntimeDeliveryPlanReceipt {
}
#[cfg(feature = "transport-workers")]
+fn satisfied_target_count_for_policy(
+ policy: &RadrootsTransportSatisfactionPolicy,
+ target_states: &[RadrootsRuntimeDeliveryTargetState],
+) -> usize {
+ match policy {
+ RadrootsTransportSatisfactionPolicy::NoWait => 0,
+ RadrootsTransportSatisfactionPolicy::Any { class }
+ | RadrootsTransportSatisfactionPolicy::All { class }
+ | RadrootsTransportSatisfactionPolicy::Quorum { class, .. } => target_states
+ .iter()
+ .filter(|state| state.status.counts_as_satisfied(*class))
+ .count(),
+ RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => targets
+ .iter()
+ .filter(|required| {
+ target_states.iter().any(|state| {
+ state.target.fingerprint == **required
+ && state.status.counts_as_satisfied(*class)
+ })
+ })
+ .count(),
+ }
+}
+
+#[cfg(feature = "transport-workers")]
+fn target_states_satisfy_policy(
+ policy: &RadrootsTransportSatisfactionPolicy,
+ target_states: &[RadrootsRuntimeDeliveryTargetState],
+) -> Result<bool, RadrootsRuntimeTransportError> {
+ match policy {
+ RadrootsTransportSatisfactionPolicy::NoWait => Ok(true),
+ RadrootsTransportSatisfactionPolicy::Any { .. }
+ | RadrootsTransportSatisfactionPolicy::All { .. }
+ | RadrootsTransportSatisfactionPolicy::Quorum { .. } => Ok(policy.is_satisfied_by(
+ target_states.len(),
+ satisfied_target_count_for_policy(policy, target_states),
+ )?),
+ RadrootsTransportSatisfactionPolicy::RequiredTargets { targets, .. } => {
+ policy.required_target_count(target_states.len())?;
+ Ok(satisfied_target_count_for_policy(policy, target_states) == targets.len())
+ }
+ }
+}
+
+#[cfg(feature = "transport-workers")]
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RadrootsRuntimeDeliveryJobReceipt {
pub outbox_event_id: i64,
@@ -433,10 +477,18 @@ impl<'a> RadrootsRuntimeDeliveryWorker<'a> {
for plan in job.plans {
let delivery_plan_id = plan.delivery_plan_id;
let satisfaction_policy = plan.satisfaction_policy.clone();
- let mut target_statuses = plan
+ let mut target_states = plan
.targets
.iter()
- .map(|target| (target.target.fingerprint.as_str().to_owned(), target.status))
+ .map(|target| {
+ (
+ target.target.fingerprint.as_str().to_owned(),
+ RadrootsRuntimeDeliveryTargetState {
+ target: target.target.clone(),
+ status: target.status,
+ },
+ )
+ })
.collect::<BTreeMap<_, _>>();
let mut by_kind =
BTreeMap::<RadrootsTransportKind, Vec<RadrootsRuntimeDeliveryTarget>>::new();
@@ -469,19 +521,22 @@ impl<'a> RadrootsRuntimeDeliveryWorker<'a> {
)?;
let receipt = adapter.deliver(request).await?;
for target_receipt in receipt.target_receipts {
- target_statuses.insert(
+ target_states.insert(
target_receipt.target.fingerprint.as_str().to_owned(),
- target_receipt.status,
+ RadrootsRuntimeDeliveryTargetState {
+ target: target_receipt.target.clone(),
+ status: target_receipt.status,
+ },
);
plan_target_receipts.push(target_receipt);
}
dispatch_count += 1;
}
- let final_statuses = target_statuses.into_values().collect::<Vec<_>>();
- let plan_receipt = RadrootsRuntimeDeliveryPlanReceipt::from_statuses(
+ let final_target_states = target_states.into_values().collect::<Vec<_>>();
+ let plan_receipt = RadrootsRuntimeDeliveryPlanReceipt::from_target_states(
delivery_plan_id,
satisfaction_policy,
- &final_statuses,
+ &final_target_states,
plan_target_receipts,
)?;
target_receipts.extend(plan_receipt.target_receipts.iter().cloned());
@@ -616,14 +671,14 @@ mod tests {
use radroots_events::draft::{RadrootsSignedNostrEvent, RadrootsSignedNostrEventParts};
use radroots_transport::{
RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryTargetStatus,
- RadrootsTransportKind, RadrootsTransportOutcome, RadrootsTransportSatisfactionClass,
- RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget,
- RadrootsTransportTargetReceipt,
+ RadrootsTransportKind, RadrootsTransportOutcome, RadrootsTransportOutcomeKind,
+ RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy,
+ RadrootsTransportTarget, RadrootsTransportTargetReceipt,
};
struct StaticAdapter {
kind: RadrootsTransportKind,
- status: RadrootsTransportDeliveryTargetStatus,
+ outcome_kind: RadrootsTransportOutcomeKind,
}
impl RadrootsRuntimeTransportAdapter for StaticAdapter {
@@ -646,7 +701,7 @@ mod tests {
.map(|target| {
RadrootsTransportTargetReceipt::new(
target,
- RadrootsTransportOutcome::new(self.status),
+ RadrootsTransportOutcome::new(self.outcome_kind),
)
})
.collect(),
@@ -701,7 +756,7 @@ mod tests {
registry
.register(StaticAdapter {
kind: RadrootsTransportKind::Nostr,
- status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
})
.expect("register");
assert_eq!(
@@ -711,7 +766,7 @@ mod tests {
assert!(matches!(
registry.register(StaticAdapter {
kind: RadrootsTransportKind::Nostr,
- status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
}),
Err(RadrootsRuntimeTransportError::AdapterAlreadyRegistered(_))
));
@@ -802,7 +857,7 @@ mod tests {
registry
.register(StaticAdapter {
kind: RadrootsTransportKind::Nostr,
- status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
})
.expect("register");
let worker = RadrootsRuntimeDeliveryWorker::new(
@@ -894,12 +949,72 @@ mod tests {
#[cfg(feature = "transport-workers")]
#[tokio::test]
+ async fn delivery_worker_required_targets_use_fingerprints() {
+ let mut registry = RadrootsRuntimeTransportRegistry::new();
+ registry
+ .register(StaticAdapter {
+ kind: RadrootsTransportKind::Nostr,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
+ })
+ .expect("register");
+ let worker = RadrootsRuntimeDeliveryWorker::new(
+ ®istry,
+ RadrootsRuntimeDeliveryWorkerConfig {
+ bounded_queue_capacity: 8,
+ },
+ );
+ let required_target = target(
+ RadrootsTransportKind::Reticulum,
+ "reticulum:preview-unavailable",
+ );
+ let required_fingerprint = required_target.fingerprint.clone();
+ 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::required_targets(
+ RadrootsTransportSatisfactionClass::Accepted,
+ vec![required_fingerprint],
+ )
+ .expect("required target policy"),
+ targets: vec![
+ RadrootsRuntimeDeliveryTarget::ready(
+ 1,
+ target(RadrootsTransportKind::Nostr, "wss://relay.example"),
+ ),
+ RadrootsRuntimeDeliveryTarget::deferred_until_implemented(
+ 2,
+ required_target,
+ ),
+ ],
+ }],
+ now_ms: 1_000,
+ })
+ .await
+ .expect("worker receipt");
+
+ assert_eq!(receipt.dispatch_count, 1);
+ assert!(!receipt.all_plans_satisfied);
+ assert_eq!(receipt.plan_receipts[0].target_count, 2);
+ assert_eq!(receipt.plan_receipts[0].attempted_target_count, 1);
+ assert_eq!(receipt.plan_receipts[0].required_target_count, 1);
+ assert_eq!(receipt.plan_receipts[0].satisfied_target_count, 0);
+ assert_eq!(
+ receipt.plan_receipts[0].satisfaction_state,
+ RadrootsRuntimeDeliveryPlanSatisfactionState::Unsatisfied
+ );
+ }
+
+ #[cfg(feature = "transport-workers")]
+ #[tokio::test]
async fn delivery_worker_keeps_accepted_and_delivered_satisfaction_distinct() {
let mut registry = RadrootsRuntimeTransportRegistry::new();
registry
.register(StaticAdapter {
kind: RadrootsTransportKind::Nostr,
- status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
})
.expect("register");
let worker = RadrootsRuntimeDeliveryWorker::new(
@@ -942,7 +1057,7 @@ mod tests {
registry
.register(StaticAdapter {
kind: RadrootsTransportKind::Nostr,
- status: RadrootsTransportDeliveryTargetStatus::Accepted,
+ outcome_kind: RadrootsTransportOutcomeKind::Accepted,
})
.expect("register");
let worker = RadrootsRuntimeDeliveryWorker::new(
diff --git a/crates/transport/src/delivery.rs b/crates/transport/src/delivery.rs
@@ -1,7 +1,8 @@
use crate::{
RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportOutcome,
- RadrootsTransportTarget, RadrootsTransportTargetSet,
+ RadrootsTransportTarget, RadrootsTransportTargetFingerprint, RadrootsTransportTargetSet,
};
+use alloc::collections::BTreeSet;
use alloc::string::String;
use alloc::vec::Vec;
@@ -26,6 +27,10 @@ pub enum RadrootsTransportSatisfactionPolicy {
class: RadrootsTransportSatisfactionClass,
threshold: u16,
},
+ RequiredTargets {
+ class: RadrootsTransportSatisfactionClass,
+ targets: Vec<RadrootsTransportTargetFingerprint>,
+ },
}
impl RadrootsTransportSatisfactionPolicy {
@@ -71,10 +76,28 @@ impl RadrootsTransportSatisfactionPolicy {
}
}
+ pub fn required_targets(
+ class: RadrootsTransportSatisfactionClass,
+ targets: Vec<RadrootsTransportTargetFingerprint>,
+ ) -> Result<Self, RadrootsTransportError> {
+ validate_required_targets(&targets)?;
+ Ok(Self::RequiredTargets { class, targets })
+ }
+
pub fn target_satisfaction_class(&self) -> Option<RadrootsTransportSatisfactionClass> {
match self {
Self::NoWait => None,
- Self::Any { class } | Self::All { class } | Self::Quorum { class, .. } => Some(*class),
+ Self::Any { class }
+ | Self::All { class }
+ | Self::Quorum { class, .. }
+ | Self::RequiredTargets { class, .. } => Some(*class),
+ }
+ }
+
+ pub fn required_target_fingerprints(&self) -> Option<&[RadrootsTransportTargetFingerprint]> {
+ match self {
+ Self::RequiredTargets { targets, .. } => Some(targets),
+ Self::NoWait | Self::Any { .. } | Self::All { .. } | Self::Quorum { .. } => None,
}
}
@@ -98,6 +121,13 @@ impl RadrootsTransportSatisfactionPolicy {
Ok(usize::from(*threshold))
}
Self::Quorum { .. } => Err(RadrootsTransportError::InvalidSatisfactionPolicy),
+ Self::RequiredTargets { targets, .. } => {
+ validate_required_targets(targets)?;
+ if targets.len() > total_targets {
+ return Err(RadrootsTransportError::InvalidSatisfactionPolicy);
+ }
+ Ok(targets.len())
+ }
}
}
@@ -106,11 +136,29 @@ impl RadrootsTransportSatisfactionPolicy {
total_targets: usize,
satisfied_targets: usize,
) -> Result<bool, RadrootsTransportError> {
+ if matches!(self, Self::RequiredTargets { .. }) {
+ return Err(RadrootsTransportError::InvalidSatisfactionPolicy);
+ }
let required = self.required_target_count(total_targets)?;
Ok(satisfied_targets >= required)
}
}
+fn validate_required_targets(
+ targets: &[RadrootsTransportTargetFingerprint],
+) -> Result<(), RadrootsTransportError> {
+ if targets.is_empty() {
+ return Err(RadrootsTransportError::EmptyRequiredTargetSet);
+ }
+ let mut fingerprints = BTreeSet::new();
+ for target in targets {
+ if !fingerprints.insert(target.as_str()) {
+ return Err(RadrootsTransportError::DuplicateRequiredTargetFingerprint);
+ }
+ }
+ Ok(())
+}
+
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RadrootsTransportDeliveryRequest {
@@ -171,4 +219,28 @@ impl RadrootsTransportDeliveryReceipt {
.filter(|receipt| receipt.status.counts_as_satisfied(satisfaction_class))
.count()
}
+
+ pub fn is_satisfied_by(
+ &self,
+ policy: &RadrootsTransportSatisfactionPolicy,
+ ) -> Result<bool, RadrootsTransportError> {
+ match policy {
+ RadrootsTransportSatisfactionPolicy::NoWait => Ok(true),
+ RadrootsTransportSatisfactionPolicy::Any { class }
+ | RadrootsTransportSatisfactionPolicy::All { class }
+ | RadrootsTransportSatisfactionPolicy::Quorum { class, .. } => policy.is_satisfied_by(
+ self.target_receipts.len(),
+ self.satisfied_target_count(*class),
+ ),
+ RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => {
+ validate_required_targets(targets)?;
+ Ok(targets.iter().all(|required| {
+ self.target_receipts.iter().any(|receipt| {
+ receipt.target.fingerprint == *required
+ && receipt.status.counts_as_satisfied(*class)
+ })
+ }))
+ }
+ }
+ }
}
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -14,6 +14,8 @@ pub enum RadrootsTransportError {
DuplicateTargetFingerprint,
InvalidTargetFingerprint,
InvalidSatisfactionPolicy,
+ EmptyRequiredTargetSet,
+ DuplicateRequiredTargetFingerprint,
}
impl fmt::Display for RadrootsTransportError {
@@ -37,6 +39,10 @@ impl fmt::Display for RadrootsTransportError {
Self::InvalidSatisfactionPolicy => {
f.write_str("transport satisfaction policy is invalid")
}
+ Self::EmptyRequiredTargetSet => f.write_str("transport required target set is empty"),
+ Self::DuplicateRequiredTargetFingerprint => {
+ f.write_str("transport required target set contains duplicate fingerprints")
+ }
}
}
}
diff --git a/crates/transport/src/lib.rs b/crates/transport/src/lib.rs
@@ -9,6 +9,7 @@ mod kind;
mod message;
mod status;
mod target;
+mod transport;
pub use delivery::{
RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
@@ -22,9 +23,13 @@ pub use message::{
RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE,
};
pub use status::{
- RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome, RadrootsTransportStatus,
+ RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome, RadrootsTransportOutcomeKind,
+ RadrootsTransportStatus,
};
pub use target::{
RadrootsTransportMeshScopeId, RadrootsTransportTarget, RadrootsTransportTargetFingerprint,
RadrootsTransportTargetLabel, RadrootsTransportTargetSet, RadrootsTransportTargetUri,
};
+pub use transport::{
+ RadrootsTransport, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest,
+};
diff --git a/crates/transport/src/status.rs b/crates/transport/src/status.rs
@@ -68,21 +68,107 @@ impl RadrootsTransportDeliveryTargetStatus {
}
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
+pub enum RadrootsTransportOutcomeKind {
+ Accepted,
+ DuplicateAccepted,
+ Delivered,
+ Forwarded,
+ StoredByGateway,
+ Seen,
+ DeferredUntilImplemented,
+ Rejected,
+ RouteUnavailable,
+ PayloadTooLarge,
+ PolicyDenied,
+ Timeout,
+ ConnectionFailed,
+ TransportUnavailable,
+}
+
+impl RadrootsTransportOutcomeKind {
+ pub fn as_str(self) -> &'static str {
+ match self {
+ Self::Accepted => "accepted",
+ Self::DuplicateAccepted => "duplicate_accepted",
+ Self::Delivered => "delivered",
+ Self::Forwarded => "forwarded",
+ Self::StoredByGateway => "stored_by_gateway",
+ Self::Seen => "seen",
+ Self::DeferredUntilImplemented => "deferred_until_implemented",
+ Self::Rejected => "rejected",
+ Self::RouteUnavailable => "route_unavailable",
+ Self::PayloadTooLarge => "payload_too_large",
+ Self::PolicyDenied => "policy_denied",
+ Self::Timeout => "timeout",
+ Self::ConnectionFailed => "connection_failed",
+ Self::TransportUnavailable => "transport_unavailable",
+ }
+ }
+
+ pub fn target_status(self) -> RadrootsTransportDeliveryTargetStatus {
+ match self {
+ Self::Accepted | Self::DuplicateAccepted => {
+ RadrootsTransportDeliveryTargetStatus::Accepted
+ }
+ Self::Delivered => RadrootsTransportDeliveryTargetStatus::Delivered,
+ Self::Forwarded => RadrootsTransportDeliveryTargetStatus::Forwarded,
+ Self::StoredByGateway => RadrootsTransportDeliveryTargetStatus::StoredByGateway,
+ Self::Seen => RadrootsTransportDeliveryTargetStatus::Seen,
+ Self::DeferredUntilImplemented => {
+ RadrootsTransportDeliveryTargetStatus::DeferredUntilImplemented
+ }
+ Self::PolicyDenied => RadrootsTransportDeliveryTargetStatus::SkippedPolicyDenied,
+ Self::Timeout | Self::ConnectionFailed | Self::TransportUnavailable => {
+ RadrootsTransportDeliveryTargetStatus::FailedRetryable
+ }
+ Self::Rejected | Self::RouteUnavailable | Self::PayloadTooLarge => {
+ RadrootsTransportDeliveryTargetStatus::FailedTerminal
+ }
+ }
+ }
+
+ pub fn counts_as_satisfied(
+ self,
+ satisfaction_class: RadrootsTransportSatisfactionClass,
+ ) -> bool {
+ self.target_status().counts_as_satisfied(satisfaction_class)
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RadrootsTransportOutcome {
+ pub kind: RadrootsTransportOutcomeKind,
pub status: RadrootsTransportDeliveryTargetStatus,
pub code: Option<String>,
pub message: Option<String>,
}
impl RadrootsTransportOutcome {
- pub fn new(status: RadrootsTransportDeliveryTargetStatus) -> Self {
+ pub fn new(kind: RadrootsTransportOutcomeKind) -> Self {
Self {
- status,
+ kind,
+ status: kind.target_status(),
code: None,
message: None,
}
}
+
+ pub fn with_message(mut self, message: impl Into<String>) -> Self {
+ self.message = Some(message.into());
+ self
+ }
+
+ pub fn with_code(mut self, code: impl Into<String>) -> Self {
+ self.code = Some(code.into());
+ self
+ }
+
+ pub fn with_target_status(mut self, status: RadrootsTransportDeliveryTargetStatus) -> Self {
+ self.status = status;
+ self
+ }
}
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
diff --git a/crates/transport/src/transport.rs b/crates/transport/src/transport.rs
@@ -0,0 +1,58 @@
+use crate::{
+ RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, RadrootsTransportError,
+ RadrootsTransportStatus, RadrootsTransportTargetReceipt, RadrootsTransportTargetSet,
+};
+use alloc::string::String;
+use alloc::vec::Vec;
+
+pub trait RadrootsTransport {
+ fn status(&self) -> Result<RadrootsTransportStatus, RadrootsTransportError>;
+
+ fn deliver(
+ &self,
+ request: RadrootsTransportDeliveryRequest,
+ ) -> Result<RadrootsTransportDeliveryReceipt, RadrootsTransportError>;
+
+ fn fetch(
+ &self,
+ request: RadrootsTransportFetchRequest,
+ ) -> Result<RadrootsTransportFetchReceipt, RadrootsTransportError>;
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsTransportFetchRequest {
+ pub request_id: String,
+ pub target_set: RadrootsTransportTargetSet,
+}
+
+impl RadrootsTransportFetchRequest {
+ pub fn new(request_id: impl Into<String>, target_set: RadrootsTransportTargetSet) -> Self {
+ Self {
+ request_id: request_id.into(),
+ target_set,
+ }
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsTransportFetchReceipt {
+ pub request_id: String,
+ pub target_receipts: Vec<RadrootsTransportTargetReceipt>,
+ pub fetched_count: usize,
+}
+
+impl RadrootsTransportFetchReceipt {
+ pub fn new(
+ request_id: impl Into<String>,
+ target_receipts: Vec<RadrootsTransportTargetReceipt>,
+ fetched_count: usize,
+ ) -> Self {
+ Self {
+ request_id: request_id.into(),
+ target_receipts,
+ fetched_count,
+ }
+ }
+}
diff --git a/crates/transport/tests/source_boundary.rs b/crates/transport/tests/source_boundary.rs
@@ -27,6 +27,8 @@ const GENERIC_TRANSPORT_STATUS_SOURCE_ROOTS: &[&str] = &[
const CORE_STATUS_CONTRACT_SOURCE_ROOTS: &[&str] = &["transport/src", "transport_reticulum/src"];
+const CORE_TRANSPORT_CONTRACT_SOURCE_ROOTS: &[&str] = &["transport/src"];
+
const FORBIDDEN_TRANSPORT_CONCEPTS: &[ForbiddenConcept] = &[
ForbiddenConcept {
pattern: "\"radrootsd_proxy\"",
@@ -160,6 +162,21 @@ const FORBIDDEN_GENERIC_TRANSPORT_STATUS_CONCEPTS: &[ForbiddenConcept] = &[
},
];
+const FORBIDDEN_CORE_TRANSPORT_CONCEPTS: &[ForbiddenConcept] = &[
+ ForbiddenConcept {
+ pattern: concat!("Radroots", "Relay"),
+ reason: "core transport contracts must not expose Nostr relay-shaped APIs",
+ },
+ ForbiddenConcept {
+ pattern: concat!("Relay", "Transport"),
+ reason: "core transport contracts must use transport-neutral names",
+ },
+ ForbiddenConcept {
+ pattern: concat!("relay", "_transport"),
+ reason: "core transport contracts must use transport-neutral names",
+ },
+];
+
#[test]
fn transport_hardening_sources_reject_removed_protocol_identifiers() {
let crates_root = Path::new(env!("CARGO_MANIFEST_DIR"))
@@ -260,6 +277,37 @@ fn generic_transport_status_sources_reject_retired_relay_shaped_names() {
}
#[test]
+fn core_transport_sources_reject_relay_shaped_public_contracts() {
+ let crates_root = Path::new(env!("CARGO_MANIFEST_DIR"))
+ .parent()
+ .expect("transport crate parent");
+ let mut findings = Vec::new();
+
+ for relative_root in CORE_TRANSPORT_CONTRACT_SOURCE_ROOTS {
+ for path in rust_source_files(crates_root.join(relative_root).as_path()) {
+ let source_raw = read_source(path.as_path());
+ let source = production_source(source_raw.as_str());
+ let relative_path = relative_path(crates_root, path.as_path());
+
+ for concept in FORBIDDEN_CORE_TRANSPORT_CONCEPTS {
+ if contains_forbidden_concept(source, concept.pattern) {
+ findings.push(format!(
+ "{} contains relay-shaped core transport concept `{}`: {}",
+ relative_path, concept.pattern, concept.reason
+ ));
+ }
+ }
+ }
+ }
+
+ assert!(
+ findings.is_empty(),
+ "core transport public contract source-boundary violations:\n{}",
+ findings.join("\n")
+ );
+}
+
+#[test]
fn transport_publish_capabilities_keep_canonical_status_fields() {
let source_raw = read_source(
Path::new(env!("CARGO_MANIFEST_DIR"))
diff --git a/crates/transport/tests/transport.rs b/crates/transport/tests/transport.rs
@@ -1,12 +1,13 @@
use radroots_transport::{
RADROOTS_RETICULUM_PREVIEW_ENDPOINT_URI, RADROOTS_RETICULUM_PREVIEW_SCOPE_ID,
- RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
- RadrootsTransportDeliveryTargetStatus, RadrootsTransportError,
- RadrootsTransportImplementationState, RadrootsTransportKind, RadrootsTransportMeshScopeId,
- RadrootsTransportOutcome, RadrootsTransportSatisfactionClass,
- RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget,
- RadrootsTransportTargetFingerprint, RadrootsTransportTargetLabel,
- RadrootsTransportTargetReceipt, RadrootsTransportTargetSet, RadrootsTransportTargetUri,
+ RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
+ RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportFetchReceipt,
+ RadrootsTransportFetchRequest, RadrootsTransportImplementationState, RadrootsTransportKind,
+ RadrootsTransportMeshScopeId, RadrootsTransportOutcome, RadrootsTransportOutcomeKind,
+ RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy,
+ RadrootsTransportStatus, RadrootsTransportTarget, RadrootsTransportTargetFingerprint,
+ RadrootsTransportTargetLabel, RadrootsTransportTargetReceipt, RadrootsTransportTargetSet,
+ RadrootsTransportTargetUri,
};
#[test]
@@ -231,9 +232,7 @@ fn deferred_transport_outcomes_are_terminal_but_not_satisfied() {
request_id: "reticulum-preview".to_owned(),
target_receipts: vec![RadrootsTransportTargetReceipt::new(
target,
- RadrootsTransportOutcome::new(
- RadrootsTransportDeliveryTargetStatus::DeferredUntilImplemented,
- ),
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::DeferredUntilImplemented),
)],
};
@@ -334,6 +333,14 @@ fn transport_errors_have_stable_display_strings() {
RadrootsTransportError::InvalidSatisfactionPolicy,
"transport satisfaction policy is invalid",
),
+ (
+ RadrootsTransportError::EmptyRequiredTargetSet,
+ "transport required target set is empty",
+ ),
+ (
+ RadrootsTransportError::DuplicateRequiredTargetFingerprint,
+ "transport required target set contains duplicate fingerprints",
+ ),
];
for (error, message) in cases {
@@ -571,3 +578,288 @@ fn satisfaction_and_target_status_cover_all_contract_states() {
assert!(RadrootsTransportDeliveryTargetStatus::FailedRetryable.is_retryable_failure());
assert!(RadrootsTransportDeliveryTargetStatus::FailedTerminal.is_terminal_failure());
}
+
+#[test]
+fn typed_outcome_kinds_drive_status_and_satisfaction_semantics() {
+ let cases = [
+ (
+ RadrootsTransportOutcomeKind::Accepted,
+ "accepted",
+ RadrootsTransportDeliveryTargetStatus::Accepted,
+ true,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::DuplicateAccepted,
+ "duplicate_accepted",
+ RadrootsTransportDeliveryTargetStatus::Accepted,
+ true,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Delivered,
+ "delivered",
+ RadrootsTransportDeliveryTargetStatus::Delivered,
+ true,
+ true,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Forwarded,
+ "forwarded",
+ RadrootsTransportDeliveryTargetStatus::Forwarded,
+ true,
+ true,
+ ),
+ (
+ RadrootsTransportOutcomeKind::StoredByGateway,
+ "stored_by_gateway",
+ RadrootsTransportDeliveryTargetStatus::StoredByGateway,
+ true,
+ true,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Seen,
+ "seen",
+ RadrootsTransportDeliveryTargetStatus::Seen,
+ true,
+ true,
+ ),
+ (
+ RadrootsTransportOutcomeKind::DeferredUntilImplemented,
+ "deferred_until_implemented",
+ RadrootsTransportDeliveryTargetStatus::DeferredUntilImplemented,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Rejected,
+ "rejected",
+ RadrootsTransportDeliveryTargetStatus::FailedTerminal,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::RouteUnavailable,
+ "route_unavailable",
+ RadrootsTransportDeliveryTargetStatus::FailedTerminal,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::PayloadTooLarge,
+ "payload_too_large",
+ RadrootsTransportDeliveryTargetStatus::FailedTerminal,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::PolicyDenied,
+ "policy_denied",
+ RadrootsTransportDeliveryTargetStatus::SkippedPolicyDenied,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Timeout,
+ "timeout",
+ RadrootsTransportDeliveryTargetStatus::FailedRetryable,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::ConnectionFailed,
+ "connection_failed",
+ RadrootsTransportDeliveryTargetStatus::FailedRetryable,
+ false,
+ false,
+ ),
+ (
+ RadrootsTransportOutcomeKind::TransportUnavailable,
+ "transport_unavailable",
+ RadrootsTransportDeliveryTargetStatus::FailedRetryable,
+ false,
+ false,
+ ),
+ ];
+
+ for (kind, label, status, accepted, delivered) in cases {
+ let outcome = RadrootsTransportOutcome::new(kind).with_message("transport detail");
+ assert_eq!(kind.as_str(), label);
+ assert_eq!(outcome.kind, kind);
+ assert_eq!(outcome.status, status);
+ assert_eq!(
+ kind.counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted),
+ accepted
+ );
+ assert_eq!(
+ kind.counts_as_satisfied(RadrootsTransportSatisfactionClass::Delivered),
+ delivered
+ );
+ assert_eq!(outcome.message.as_deref(), Some("transport detail"));
+ }
+
+ let preview_unavailable =
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::TransportUnavailable)
+ .with_target_status(RadrootsTransportDeliveryTargetStatus::PreviewUnavailable);
+ assert_eq!(
+ preview_unavailable.status,
+ RadrootsTransportDeliveryTargetStatus::PreviewUnavailable
+ );
+}
+
+#[test]
+fn required_target_satisfaction_uses_fingerprints_not_target_counts() {
+ let required = RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, "wss://one.example")
+ .expect("required target");
+ let optional = RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, "wss://two.example")
+ .expect("optional target");
+ let policy = RadrootsTransportSatisfactionPolicy::required_targets(
+ RadrootsTransportSatisfactionClass::Accepted,
+ vec![required.fingerprint.clone()],
+ )
+ .expect("required target policy");
+ assert_eq!(policy.required_target_count(2).expect("required count"), 1);
+ assert_eq!(
+ policy
+ .is_satisfied_by(2, 1)
+ .expect_err("count-only required targets are invalid"),
+ RadrootsTransportError::InvalidSatisfactionPolicy
+ );
+ assert_eq!(
+ policy
+ .required_target_count(0)
+ .expect_err("required target count exceeds total targets"),
+ RadrootsTransportError::InvalidSatisfactionPolicy
+ );
+ assert_eq!(
+ policy.target_satisfaction_class(),
+ Some(RadrootsTransportSatisfactionClass::Accepted)
+ );
+ assert_eq!(
+ policy
+ .required_target_fingerprints()
+ .expect("required targets"),
+ &[required.fingerprint.clone()]
+ );
+
+ let optional_only = RadrootsTransportDeliveryReceipt {
+ request_id: "required-target".to_owned(),
+ target_receipts: vec![RadrootsTransportTargetReceipt::new(
+ optional.clone(),
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted),
+ )],
+ };
+ assert!(
+ !optional_only
+ .is_satisfied_by(&policy)
+ .expect("missing required target")
+ );
+
+ let required_delivered = RadrootsTransportDeliveryReceipt {
+ request_id: "required-target".to_owned(),
+ target_receipts: vec![
+ RadrootsTransportTargetReceipt::new(
+ optional,
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Rejected),
+ ),
+ RadrootsTransportTargetReceipt::new(
+ required.clone(),
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted),
+ ),
+ ],
+ };
+ assert!(
+ required_delivered
+ .is_satisfied_by(&policy)
+ .expect("required target accepted")
+ );
+ assert_eq!(
+ RadrootsTransportSatisfactionPolicy::required_targets(
+ RadrootsTransportSatisfactionClass::Accepted,
+ Vec::new(),
+ )
+ .expect_err("empty required targets"),
+ RadrootsTransportError::EmptyRequiredTargetSet
+ );
+ assert_eq!(
+ RadrootsTransportSatisfactionPolicy::required_targets(
+ RadrootsTransportSatisfactionClass::Accepted,
+ vec![required.fingerprint.clone(), required.fingerprint],
+ )
+ .expect_err("duplicate required target"),
+ RadrootsTransportError::DuplicateRequiredTargetFingerprint
+ );
+}
+
+#[test]
+fn neutral_transport_trait_covers_status_delivery_and_fetch() {
+ struct MemoryTransport {
+ target: RadrootsTransportTarget,
+ }
+
+ impl RadrootsTransport for MemoryTransport {
+ fn status(&self) -> Result<RadrootsTransportStatus, RadrootsTransportError> {
+ Ok(RadrootsTransportStatus::new(
+ RadrootsTransportKind::Local,
+ true,
+ RadrootsTransportImplementationState::Real,
+ true,
+ "ready",
+ ))
+ }
+
+ fn deliver(
+ &self,
+ request: RadrootsTransportDeliveryRequest,
+ ) -> Result<RadrootsTransportDeliveryReceipt, RadrootsTransportError> {
+ Ok(RadrootsTransportDeliveryReceipt {
+ request_id: request.request_id,
+ target_receipts: vec![RadrootsTransportTargetReceipt::new(
+ self.target.clone(),
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Delivered),
+ )],
+ })
+ }
+
+ fn fetch(
+ &self,
+ request: RadrootsTransportFetchRequest,
+ ) -> Result<RadrootsTransportFetchReceipt, RadrootsTransportError> {
+ Ok(RadrootsTransportFetchReceipt::new(
+ request.request_id,
+ vec![RadrootsTransportTargetReceipt::new(
+ self.target.clone(),
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Seen),
+ )],
+ 1,
+ ))
+ }
+ }
+
+ let target = RadrootsTransportTarget::new(RadrootsTransportKind::Local, "local:memory")
+ .expect("local target");
+ let target_set = RadrootsTransportTargetSet::new(vec![target.clone()]).expect("target set");
+ let transport = MemoryTransport { target };
+ let status = transport.status().expect("status");
+ assert_eq!(status.kind, RadrootsTransportKind::Local);
+ let delivery = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "deliver-1",
+ "sha256:payload",
+ target_set.clone(),
+ RadrootsTransportSatisfactionPolicy::all_delivered(),
+ ))
+ .expect("deliver");
+ assert_eq!(
+ delivery.target_receipts[0].outcome.kind,
+ RadrootsTransportOutcomeKind::Delivered
+ );
+ let fetch = transport
+ .fetch(RadrootsTransportFetchRequest::new("fetch-1", target_set))
+ .expect("fetch");
+ assert_eq!(fetch.fetched_count, 1);
+ assert_eq!(
+ fetch.target_receipts[0].outcome.kind,
+ RadrootsTransportOutcomeKind::Seen
+ );
+}
diff --git a/crates/transport_nostr/src/outcome.rs b/crates/transport_nostr/src/outcome.rs
@@ -1,6 +1,6 @@
#![forbid(unsafe_code)]
-use radroots_transport::{RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome};
+use radroots_transport::{RadrootsTransportOutcome, RadrootsTransportOutcomeKind};
use serde::{Deserialize, Serialize};
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
@@ -79,6 +79,27 @@ impl RadrootsRelayOutcomeKind {
| Self::RelayUrlRejected
)
}
+
+ pub fn transport_outcome_kind(self) -> RadrootsTransportOutcomeKind {
+ match self {
+ Self::Accepted => RadrootsTransportOutcomeKind::Accepted,
+ Self::DuplicateAccepted | Self::SkippedAlreadyAccepted => {
+ RadrootsTransportOutcomeKind::DuplicateAccepted
+ }
+ Self::Blocked | Self::Invalid | Self::Restricted | Self::Muted | Self::Unsupported => {
+ RadrootsTransportOutcomeKind::Rejected
+ }
+ Self::RelayUrlRejected => RadrootsTransportOutcomeKind::RouteUnavailable,
+ Self::PaymentRequired | Self::PowRequired | Self::AuthRequired => {
+ RadrootsTransportOutcomeKind::PolicyDenied
+ }
+ Self::RateLimited | Self::Error | Self::Unknown => {
+ RadrootsTransportOutcomeKind::TransportUnavailable
+ }
+ Self::Timeout => RadrootsTransportOutcomeKind::Timeout,
+ Self::ConnectionFailed => RadrootsTransportOutcomeKind::ConnectionFailed,
+ }
+ }
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
@@ -179,16 +200,11 @@ impl RadrootsRelayOutcome {
}
pub fn to_transport_outcome(&self) -> RadrootsTransportOutcome {
- let status = if self.counts_toward_quorum() {
- RadrootsTransportDeliveryTargetStatus::Accepted
- } else if self.is_retryable() {
- RadrootsTransportDeliveryTargetStatus::FailedRetryable
- } else {
- RadrootsTransportDeliveryTargetStatus::FailedTerminal
- };
- let mut outcome = RadrootsTransportOutcome::new(status);
- outcome.code = Some(self.kind.as_str().to_owned());
- outcome.message = self.message.clone();
+ let mut outcome = RadrootsTransportOutcome::new(self.kind.transport_outcome_kind())
+ .with_code(self.kind.as_str());
+ if let Some(message) = &self.message {
+ outcome = outcome.with_message(message.clone());
+ }
outcome
}
}
diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs
@@ -5,7 +5,9 @@ use crate::{RadrootsRelayOutcome, RadrootsRelayTargetSet, RadrootsRelayTransport
use core::time::Duration;
use futures::future::BoxFuture;
use radroots_events::draft::RadrootsSignedNostrEvent;
-use radroots_transport::RadrootsTransportSatisfactionPolicy;
+use radroots_transport::{
+ RadrootsTransportKind, RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget,
+};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Arc, Mutex, PoisonError};
@@ -110,9 +112,9 @@ where
A: RadrootsRelayPublishAdapter,
{
let event_id = request.signed_event.id.clone();
- let quorum = request
- .satisfaction_policy
- .required_target_count(request.targets.len())?;
+ let satisfaction_policy = request.satisfaction_policy.clone();
+ let target_count = request.targets.len();
+ let quorum = satisfaction_policy.required_target_count(target_count)?;
let relays = adapter.publish(request).await?;
let attempted_count = relays.iter().filter(|receipt| receipt.attempted).count();
let accepted_count = relays
@@ -127,6 +129,7 @@ where
.iter()
.filter(|receipt| receipt.outcome.is_terminal_failure())
.count();
+ let quorum_met = relay_publish_satisfies_policy(&satisfaction_policy, target_count, &relays)?;
Ok(RadrootsRelayPublishReceipt {
event_id,
attempted_count,
@@ -134,11 +137,56 @@ where
retryable_count,
terminal_count,
quorum,
- quorum_met: accepted_count >= quorum,
+ quorum_met,
relays,
})
}
+fn relay_publish_satisfies_policy(
+ policy: &RadrootsTransportSatisfactionPolicy,
+ target_count: usize,
+ relays: &[RadrootsRelayPublishRelayReceipt],
+) -> Result<bool, RadrootsRelayTransportError> {
+ match policy {
+ RadrootsTransportSatisfactionPolicy::NoWait => Ok(true),
+ RadrootsTransportSatisfactionPolicy::Any { class }
+ | RadrootsTransportSatisfactionPolicy::All { class }
+ | RadrootsTransportSatisfactionPolicy::Quorum { class, .. } => {
+ let satisfied_count = relays
+ .iter()
+ .filter(|receipt| {
+ receipt
+ .outcome
+ .to_transport_outcome()
+ .status
+ .counts_as_satisfied(*class)
+ })
+ .count();
+ Ok(policy.is_satisfied_by(target_count, satisfied_count)?)
+ }
+ RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => {
+ policy.required_target_count(target_count)?;
+ let mut satisfied_required_targets = BTreeSet::new();
+ for receipt in relays {
+ let target =
+ RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, &receipt.relay_url)?;
+ if targets.contains(&target.fingerprint)
+ && receipt
+ .outcome
+ .to_transport_outcome()
+ .status
+ .counts_as_satisfied(*class)
+ {
+ satisfied_required_targets.insert(target.fingerprint);
+ }
+ }
+ Ok(targets
+ .iter()
+ .all(|target| satisfied_required_targets.contains(target)))
+ }
+ }
+}
+
#[derive(Clone, Default)]
pub struct RadrootsMockRelayPublishAdapter {
outcomes: BTreeMap<String, RadrootsRelayOutcome>,
diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs
@@ -17,7 +17,8 @@ use radroots_outbox::{
RadrootsOutboxOperationStatus,
};
use radroots_transport::{
- RadrootsTransportKind, RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget,
+ RadrootsTransportKind, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy,
+ RadrootsTransportTarget,
};
use radroots_transport_nostr::{
RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsOutboxPublishPolicy,
@@ -54,6 +55,23 @@ impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter {
}
}
+struct PartialPublishAdapter;
+
+impl RadrootsRelayPublishAdapter for PartialPublishAdapter {
+ fn publish<'a>(
+ &'a self,
+ _request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
+ {
+ Box::pin(async {
+ Ok(vec![RadrootsRelayPublishRelayReceipt::attempted(
+ RELAY_PRIMARY_WSS,
+ RadrootsRelayOutcome::accepted(),
+ )])
+ })
+ }
+}
+
struct NostrJsonFailurePublishAdapter;
impl RadrootsRelayPublishAdapter for NostrJsonFailurePublishAdapter {
@@ -573,6 +591,10 @@ fn outcome_prefix_classification_covers_required_kinds() {
assert!(RadrootsRelayOutcome::relay_url_rejected("unsafe relay").is_terminal_failure());
assert!(RadrootsRelayOutcome::classify("mute: pubkey muted").is_terminal_failure());
assert_eq!(
+ RadrootsRelayOutcome::accepted().to_transport_outcome().kind,
+ radroots_transport::RadrootsTransportOutcomeKind::Accepted
+ );
+ assert_eq!(
RadrootsRelayOutcome::accepted()
.to_transport_outcome()
.status,
@@ -581,16 +603,34 @@ fn outcome_prefix_classification_covers_required_kinds() {
assert_eq!(
RadrootsRelayOutcome::timeout("timeout: no OK")
.to_transport_outcome()
+ .kind,
+ radroots_transport::RadrootsTransportOutcomeKind::Timeout
+ );
+ assert_eq!(
+ RadrootsRelayOutcome::timeout("timeout: no OK")
+ .to_transport_outcome()
.status,
radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedRetryable
);
assert_eq!(
RadrootsRelayOutcome::classify("restricted: denied")
.to_transport_outcome()
+ .kind,
+ radroots_transport::RadrootsTransportOutcomeKind::Rejected
+ );
+ assert_eq!(
+ RadrootsRelayOutcome::classify("restricted: denied")
+ .to_transport_outcome()
.status,
radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedTerminal
);
assert_eq!(
+ RadrootsRelayOutcome::relay_url_rejected("unsafe")
+ .to_transport_outcome()
+ .kind,
+ radroots_transport::RadrootsTransportOutcomeKind::RouteUnavailable
+ );
+ assert_eq!(
RadrootsRelayOutcome::connection_failed("offline")
.kind
.as_str(),
@@ -701,6 +741,65 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() {
assert!(matches!(error, RadrootsRelayTransportError::Transport(_)));
}
+#[tokio::test]
+async fn publish_required_target_policy_uses_relay_fingerprints() {
+ let signed = signed_post("required relay");
+ let required_target =
+ RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, RELAY_PRIMARY_WSS)
+ .expect("required target");
+ let targets = RadrootsRelayTargetSet::new(
+ vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("targets");
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(
+ RELAY_PRIMARY_WSS,
+ RadrootsRelayOutcome::classify("restricted: required relay rejected"),
+ )
+ .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted());
+
+ let receipt = publish_signed_event(
+ &adapter,
+ RadrootsRelayPublishRequest::new(signed, targets, 1_070).with_satisfaction_policy(
+ RadrootsTransportSatisfactionPolicy::required_targets(
+ RadrootsTransportSatisfactionClass::Accepted,
+ vec![required_target.fingerprint],
+ )
+ .expect("required relay policy"),
+ ),
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(receipt.accepted_count, 1);
+ assert_eq!(receipt.quorum, 1);
+ assert!(!receipt.quorum_met);
+}
+
+#[tokio::test]
+async fn publish_all_policy_uses_requested_target_count() {
+ let signed = signed_post("partial adapter");
+ let targets = RadrootsRelayTargetSet::new(
+ vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("targets");
+
+ let receipt = publish_signed_event(
+ &PartialPublishAdapter,
+ RadrootsRelayPublishRequest::new(signed, targets, 1_080)
+ .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::all_accepted()),
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(receipt.attempted_count, 1);
+ assert_eq!(receipt.accepted_count, 1);
+ assert_eq!(receipt.quorum, 2);
+ assert!(!receipt.quorum_met);
+}
+
#[test]
fn fetch_requests_reject_empty_filter_sets() {
assert!(matches!(
diff --git a/crates/transport_reticulum/src/lib.rs b/crates/transport_reticulum/src/lib.rs
@@ -12,8 +12,8 @@ use radroots_transport::{
RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, RadrootsTransportDeliveryReceipt,
RadrootsTransportDeliveryRequest, RadrootsTransportDeliveryTargetStatus,
RadrootsTransportImplementationState, RadrootsTransportKind, RadrootsTransportMeshScopeId,
- RadrootsTransportOutcome, RadrootsTransportStatus, RadrootsTransportTarget,
- RadrootsTransportTargetReceipt,
+ RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportStatus,
+ RadrootsTransportTarget, RadrootsTransportTargetReceipt,
};
const DEFAULT_PROFILE_ID: &str = "transport.reticulum.preview";
@@ -359,11 +359,12 @@ fn ensure_reticulum_targets(
fn preview_outcome(behavior: RadrootsReticulumPreviewBehavior) -> RadrootsTransportOutcome {
let mut outcome = match behavior {
RadrootsReticulumPreviewBehavior::RejectDeliveryAttempts => {
- RadrootsTransportOutcome::new(RadrootsTransportDeliveryTargetStatus::PreviewUnavailable)
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::TransportUnavailable)
+ .with_target_status(RadrootsTransportDeliveryTargetStatus::PreviewUnavailable)
+ }
+ RadrootsReticulumPreviewBehavior::DeferDeliveryPlans => {
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::DeferredUntilImplemented)
}
- RadrootsReticulumPreviewBehavior::DeferDeliveryPlans => RadrootsTransportOutcome::new(
- RadrootsTransportDeliveryTargetStatus::DeferredUntilImplemented,
- ),
};
outcome.code = Some(
match behavior {