commit 89f5bf05516d24cd57007e676e5bd8c8889e5325
parent 0e5f7595dd2dc56e77416bea5b970eda0fa3022d
Author: triesap <tyson@radroots.org>
Date: Mon, 6 Jul 2026 20:46:21 +0000
transport_nostr: isolate Nostr transport
- rename relay transport package to radroots_transport_nostr without old-name reexports
- publish outbox delivery targets instead of relay status rows
- drive Nostr publish quorum from transport satisfaction policy
Diffstat:
22 files changed, 2952 insertions(+), 2918 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -4555,24 +4555,6 @@ dependencies = [
]
[[package]]
-name = "radroots_relay_transport"
-version = "0.1.0-alpha.2"
-dependencies = [
- "futures",
- "nostr",
- "radroots_event_store",
- "radroots_events",
- "radroots_nostr",
- "radroots_outbox",
- "radroots_transport",
- "serde",
- "serde_json",
- "thiserror 1.0.69",
- "tokio",
- "url",
-]
-
-[[package]]
name = "radroots_replica_db"
version = "0.1.0-alpha.2"
dependencies = [
@@ -4876,6 +4858,24 @@ dependencies = [
]
[[package]]
+name = "radroots_transport_nostr"
+version = "0.1.0-alpha.2"
+dependencies = [
+ "futures",
+ "nostr",
+ "radroots_event_store",
+ "radroots_events",
+ "radroots_nostr",
+ "radroots_outbox",
+ "radroots_transport",
+ "serde",
+ "serde_json",
+ "thiserror 1.0.69",
+ "tokio",
+ "url",
+]
+
+[[package]]
name = "radroots_types"
version = "0.1.0-alpha.2"
dependencies = [
diff --git a/Cargo.toml b/Cargo.toml
@@ -19,7 +19,7 @@ members = [
"crates/nostr_runtime",
"crates/outbox",
"crates/publish_proxy_protocol",
- "crates/relay_transport",
+ "crates/transport_nostr",
"crates/runtime",
"crates/secret_vault",
"crates/simplex_app_store",
@@ -86,7 +86,7 @@ radroots_net = { path = "crates/net", version = "0.1.0-alpha.2", default-feature
radroots_nostr_runtime = { path = "crates/nostr_runtime", version = "0.1.0-alpha.2", default-features = false }
radroots_outbox = { path = "crates/outbox", version = "0.1.0-alpha.2", default-features = false }
radroots_publish_proxy_protocol = { path = "crates/publish_proxy_protocol", version = "0.1.0-alpha.2", default-features = false }
-radroots_relay_transport = { path = "crates/relay_transport", version = "0.1.0-alpha.2", default-features = false }
+radroots_transport_nostr = { path = "crates/transport_nostr", version = "0.1.0-alpha.2", default-features = false }
radroots_simplex_agent_proto = { path = "crates/simplex_agent_proto", version = "0.1.0-alpha.2", default-features = false }
radroots_simplex_agent_runtime = { path = "crates/simplex_agent_runtime", version = "0.1.0-alpha.2", default-features = false }
radroots_simplex_agent_store = { path = "crates/simplex_agent_store", version = "0.1.0-alpha.2", default-features = false }
diff --git a/contracts/coverage.toml b/contracts/coverage.toml
@@ -53,7 +53,7 @@ crates = [
"radroots_outbox",
"radroots_protected_store",
"radroots_publish_proxy_protocol",
- "radroots_relay_transport",
+ "radroots_transport_nostr",
"radroots_replica_db",
"radroots_replica_db_schema",
"radroots_replica_sync",
diff --git a/contracts/manifest.toml b/contracts/manifest.toml
@@ -45,7 +45,7 @@ deferred_publication = [
"radroots_event_store",
"radroots_outbox",
"radroots_publish_proxy_protocol",
- "radroots_relay_transport",
+ "radroots_transport_nostr",
"radroots_transport",
"radroots_mesh",
"radroots_mesh_agent_proto",
diff --git a/crates/relay_transport/Cargo.toml b/crates/relay_transport/Cargo.toml
@@ -1,62 +0,0 @@
-[package]
-name = "radroots_relay_transport"
-publish = false
-version = "0.1.0-alpha.2"
-edition.workspace = true
-authors = ["Tyson Lupul <tyson@radroots.org>"]
-rust-version.workspace = true
-license.workspace = true
-description = "Relay transport layer for Radroots"
-repository.workspace = true
-homepage.workspace = true
-readme = "README"
-
-[features]
-default = ["std", "client", "storage", "runtime-tokio"]
-std = []
-client = [
- "dep:radroots_nostr",
- "radroots_nostr/std",
- "radroots_nostr/client",
- "radroots_nostr/events",
-]
-storage = ["dep:radroots_event_store", "dep:radroots_outbox", "client"]
-runtime-tokio = [
- "dep:tokio",
- "storage",
- "radroots_event_store/runtime-tokio",
- "radroots_outbox/runtime-tokio",
-]
-
-[dependencies]
-radroots_events = { workspace = true, default-features = false, features = [
- "std",
- "serde",
-] }
-radroots_event_store = { workspace = true, optional = true, default-features = false, features = [
- "sqlite",
- "runtime-tokio",
-] }
-radroots_nostr = { workspace = true, optional = true, default-features = false, features = [
- "std",
- "client",
- "events",
-] }
-radroots_outbox = { workspace = true, optional = true, default-features = false, features = [
- "sqlite",
- "runtime-tokio",
-] }
-radroots_transport = { workspace = true, default-features = false }
-futures = { workspace = true }
-nostr = { workspace = true }
-serde = { workspace = true, features = ["derive", "std"] }
-serde_json = { workspace = true, features = ["std"] }
-thiserror = { workspace = true }
-tokio = { workspace = true, optional = true, features = ["rt"] }
-url = { workspace = true }
-
-[dev-dependencies]
-tokio = { workspace = true, features = ["macros", "rt"] }
-
-[lints.rust]
-unexpected_cfgs = { level = "warn", check-cfg = ['cfg(coverage_nightly)'] }
diff --git a/crates/relay_transport/README b/crates/relay_transport/README
@@ -1,3 +0,0 @@
-# radroots_relay_transport
-
-Deterministic Nostr relay transport substrate for exact signed-event publish, fetch ingest, and outbox relay status coordination.
diff --git a/crates/relay_transport/src/error.rs b/crates/relay_transport/src/error.rs
@@ -1,64 +0,0 @@
-#![forbid(unsafe_code)]
-
-use thiserror::Error;
-
-#[derive(Debug, Error)]
-pub enum RadrootsRelayTransportError {
- #[error("Relay URL parse failed for `{url}`: {reason}")]
- RelayUrlParse { url: String, reason: String },
-
- #[error("Relay URL `{url}` uses ws outside localhost relay policy")]
- WsRequiresLocalhostPolicy { url: String },
-
- #[error("Relay URL `{url}` has unsupported scheme `{scheme}`")]
- UnsupportedRelayScheme { url: String, scheme: String },
-
- #[error("Relay URL `{url}` must include a host")]
- EmptyRelayHost { url: String },
-
- #[error("Relay URL `{url}` must not include userinfo")]
- RelayUrlUserinfo { url: String },
-
- #[error("Relay URL `{url}` must not include query or fragment")]
- RelayUrlQueryOrFragment { url: String },
-
- #[error("Relay URL `{url}` targets forbidden destination: {reason}")]
- RelayUrlForbiddenDestination { url: String, reason: String },
-
- #[error("Relay URL `{url}` resolved to forbidden address `{address}`: {reason}")]
- RelayUrlResolvedForbiddenDestination {
- url: String,
- address: String,
- reason: String,
- },
-
- #[error("Relay target set must not be empty")]
- EmptyTargetSet,
-
- #[error("Relay fetch filters must not be empty")]
- EmptyFetchFilters,
-
- #[error("Relay fetch {field} must be greater than zero")]
- InvalidFetchLimit { field: &'static str },
-
- #[error("JSON error: {0}")]
- Json(#[from] serde_json::Error),
-
- #[error("Nostr event JSON error: {0}")]
- NostrEventJson(String),
-
- #[cfg(feature = "storage")]
- #[error("Event store error: {0}")]
- EventStore(#[from] radroots_event_store::RadrootsEventStoreError),
-
- #[cfg(feature = "storage")]
- #[error("Outbox error: {0}")]
- Outbox(#[from] radroots_outbox::RadrootsOutboxError),
-
- #[cfg(feature = "storage")]
- #[error("Outbox claim {0} does not contain a signed event")]
- MissingSignedOutboxEvent(i64),
-
- #[error("Relay transport error: {0}")]
- Transport(String),
-}
diff --git a/crates/relay_transport/src/outbox.rs b/crates/relay_transport/src/outbox.rs
@@ -1,305 +0,0 @@
-#![forbid(unsafe_code)]
-
-use crate::{
- RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter,
- RadrootsRelayPublishReceipt, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest,
- RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrlPolicy,
- publish_signed_event,
-};
-use radroots_event_store::{
- RadrootsEventIngest, RadrootsEventStore, RadrootsTransportObservation,
- RadrootsTransportObservationType,
-};
-use radroots_events::RadrootsNostrEvent;
-use radroots_events::draft::RadrootsSignedNostrEvent;
-use radroots_outbox::{
- RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxEventStoreIngestReceipt,
- RadrootsOutboxRelayStatus,
-};
-use radroots_transport::RadrootsTransportKind;
-
-#[derive(Clone, Debug, PartialEq, Eq)]
-pub struct RadrootsOutboxPublishPolicy {
- pub accepted_quorum: Option<usize>,
- pub next_attempt_after_ms: i64,
- pub republish_accepted_relays: bool,
- pub relay_url_policy: RadrootsRelayUrlPolicy,
-}
-
-impl RadrootsOutboxPublishPolicy {
- pub fn new(next_attempt_after_ms: i64) -> Self {
- Self {
- accepted_quorum: None,
- next_attempt_after_ms,
- republish_accepted_relays: false,
- relay_url_policy: RadrootsRelayUrlPolicy::Public,
- }
- }
-
- pub fn with_accepted_quorum(mut self, accepted_quorum: usize) -> Self {
- self.accepted_quorum = Some(accepted_quorum);
- self
- }
-
- pub fn republish_accepted_relays(mut self, enabled: bool) -> Self {
- self.republish_accepted_relays = enabled;
- self
- }
-
- pub fn relay_url_policy(mut self, policy: RadrootsRelayUrlPolicy) -> Self {
- self.relay_url_policy = policy;
- self
- }
-}
-
-#[derive(Clone, Debug, PartialEq, Eq)]
-pub struct RadrootsOutboxPublishReceipt {
- pub local_ingest: RadrootsOutboxEventStoreIngestReceipt,
- pub publish: RadrootsRelayPublishReceipt,
-}
-
-pub async fn publish_claimed_outbox_event<A>(
- outbox: &RadrootsOutbox,
- event_store: &RadrootsEventStore,
- adapter: &A,
- claimed: &RadrootsOutboxClaimedEvent,
- policy: RadrootsOutboxPublishPolicy,
- now_ms: i64,
-) -> Result<RadrootsOutboxPublishReceipt, RadrootsRelayTransportError>
-where
- A: RadrootsRelayPublishAdapter,
-{
- let signed_event = claimed.signed_event.clone().ok_or(
- RadrootsRelayTransportError::MissingSignedOutboxEvent(claimed.outbox_event_id),
- )?;
- let local_ingest = outbox
- .ingest_signed_event_local(
- event_store,
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- now_ms,
- )
- .await?;
- let publishable = publishable_relays(outbox, claimed, policy.republish_accepted_relays).await?;
- let overall_quorum = policy
- .accepted_quorum
- .unwrap_or(publishable.total_target_count);
- outbox
- .set_publish_quorum(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- overall_quorum as i64,
- now_ms,
- )
- .await?;
- if publishable.accepted_count >= overall_quorum {
- outbox
- .complete_publish_attempt(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- "relay publish incomplete",
- "relay publish terminal",
- policy.next_attempt_after_ms,
- now_ms,
- )
- .await?;
- let publish = RadrootsRelayPublishReceipt {
- event_id: signed_event.id,
- attempted_count: 0,
- accepted_count: publishable.accepted_count,
- retryable_count: 0,
- terminal_count: 0,
- quorum: overall_quorum,
- quorum_met: true,
- relays: Vec::new(),
- };
- return Ok(RadrootsOutboxPublishReceipt {
- local_ingest,
- publish,
- });
- }
- let targets = RadrootsRelayTargetSet::new(publishable.relays, policy.relay_url_policy)?;
- let target_strings = targets.relay_strings();
- let quorum = overall_quorum.saturating_sub(publishable.accepted_count);
- let request = RadrootsRelayPublishRequest::new(signed_event.clone(), targets, now_ms)
- .with_accepted_quorum(quorum);
- let publish = match publish_signed_event(adapter, request).await {
- Ok(receipt) => receipt,
- Err(RadrootsRelayTransportError::Transport(message)) => adapter_transport_failure_receipt(
- signed_event.id.clone(),
- target_strings,
- quorum,
- message,
- ),
- Err(error) => return Err(error),
- };
-
- for relay in &publish.relays {
- match relay.outcome.kind {
- RadrootsRelayOutcomeKind::Accepted | RadrootsRelayOutcomeKind::DuplicateAccepted => {
- outbox
- .mark_relay_accepted(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- relay.relay_url.as_str(),
- now_ms,
- )
- .await?;
- ingest_publish_observation(
- event_store,
- &signed_event,
- relay.relay_url.as_str(),
- relay.outcome.message.as_deref(),
- now_ms,
- )
- .await?;
- }
- _ if relay.outcome.is_retryable() => {
- outbox
- .mark_relay_failed_retryable(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- relay.relay_url.as_str(),
- relay
- .outcome
- .message
- .as_deref()
- .unwrap_or("relay publish retryable"),
- now_ms,
- )
- .await?;
- }
- _ => {
- outbox
- .mark_relay_failed_terminal(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- relay.relay_url.as_str(),
- relay
- .outcome
- .message
- .as_deref()
- .unwrap_or("relay publish terminal"),
- now_ms,
- )
- .await?;
- }
- }
- }
-
- outbox
- .complete_publish_attempt(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- "relay publish incomplete",
- "relay publish terminal",
- policy.next_attempt_after_ms,
- now_ms,
- )
- .await?;
-
- Ok(RadrootsOutboxPublishReceipt {
- local_ingest,
- publish,
- })
-}
-
-fn adapter_transport_failure_receipt(
- event_id: String,
- relay_urls: Vec<String>,
- quorum: usize,
- message: String,
-) -> RadrootsRelayPublishReceipt {
- let relays = relay_urls
- .into_iter()
- .map(|relay_url| {
- RadrootsRelayPublishRelayReceipt::attempted(
- relay_url,
- RadrootsRelayOutcome::connection_failed(message.clone()),
- )
- })
- .collect::<Vec<_>>();
- RadrootsRelayPublishReceipt {
- event_id,
- attempted_count: relays.len(),
- accepted_count: 0,
- retryable_count: relays.len(),
- terminal_count: 0,
- quorum,
- quorum_met: false,
- relays,
- }
-}
-
-struct PublishableRelays {
- relays: Vec<String>,
- total_target_count: usize,
- accepted_count: usize,
-}
-
-async fn publishable_relays(
- outbox: &RadrootsOutbox,
- claimed: &RadrootsOutboxClaimedEvent,
- republish_accepted_relays: bool,
-) -> Result<PublishableRelays, RadrootsRelayTransportError> {
- let statuses = outbox.relay_statuses(claimed.outbox_event_id).await?;
- let mut relays = Vec::new();
- let mut total_target_count = 0usize;
- let mut accepted_count = 0usize;
- for status in statuses {
- if !claimed
- .target_relays
- .iter()
- .any(|relay_url| relay_url == &status.relay_url)
- {
- continue;
- }
- total_target_count += 1;
- if status.status == RadrootsOutboxRelayStatus::Accepted {
- accepted_count += 1;
- }
- if republish_accepted_relays || status.status != RadrootsOutboxRelayStatus::Accepted {
- relays.push(status.relay_url);
- }
- }
- Ok(PublishableRelays {
- relays,
- total_target_count,
- accepted_count,
- })
-}
-
-async fn ingest_publish_observation(
- event_store: &RadrootsEventStore,
- signed_event: &RadrootsSignedNostrEvent,
- relay_url: &str,
- message: Option<&str>,
- observed_at_ms: i64,
-) -> Result<(), RadrootsRelayTransportError> {
- let mut observation = RadrootsTransportObservation::new(
- RadrootsTransportKind::Nostr,
- relay_url,
- RadrootsTransportObservationType::NostrPublishAck,
- observed_at_ms,
- )?;
- if let Some(message) = message {
- observation = observation.with_redacted_message(message);
- }
- let ingest = RadrootsEventIngest::new(event_from_signed(signed_event), observed_at_ms)
- .with_raw_json(signed_event.raw_json.clone())
- .with_observation(observation);
- event_store.ingest_event(ingest).await?;
- Ok(())
-}
-
-fn event_from_signed(signed_event: &RadrootsSignedNostrEvent) -> RadrootsNostrEvent {
- RadrootsNostrEvent {
- id: signed_event.id.clone(),
- author: signed_event.pubkey.clone(),
- created_at: signed_event.created_at,
- kind: signed_event.kind,
- tags: signed_event.tags.clone(),
- content: signed_event.content.clone(),
- sig: signed_event.sig.clone(),
- }
-}
diff --git a/crates/relay_transport/src/outcome.rs b/crates/relay_transport/src/outcome.rs
@@ -1,157 +0,0 @@
-#![forbid(unsafe_code)]
-
-use serde::{Deserialize, Serialize};
-
-#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
-pub enum RadrootsRelayOutcomeKind {
- Accepted,
- DuplicateAccepted,
- Blocked,
- RateLimited,
- Invalid,
- PowRequired,
- Restricted,
- AuthRequired,
- Muted,
- Unsupported,
- PaymentRequired,
- Error,
- Timeout,
- ConnectionFailed,
- RelayUrlRejected,
- SkippedAlreadyAccepted,
- Unknown,
-}
-
-impl RadrootsRelayOutcomeKind {
- pub fn counts_toward_quorum(self) -> bool {
- matches!(
- self,
- Self::Accepted | Self::DuplicateAccepted | Self::SkippedAlreadyAccepted
- )
- }
-
- pub fn is_retryable(self) -> bool {
- matches!(
- self,
- Self::RateLimited
- | Self::PowRequired
- | Self::AuthRequired
- | Self::Error
- | Self::Timeout
- | Self::ConnectionFailed
- | Self::Unknown
- )
- }
-
- pub fn is_terminal_failure(self) -> bool {
- matches!(
- self,
- Self::Blocked
- | Self::Invalid
- | Self::Restricted
- | Self::Muted
- | Self::Unsupported
- | Self::PaymentRequired
- | Self::RelayUrlRejected
- )
- }
-}
-
-#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
-pub struct RadrootsRelayOutcome {
- pub kind: RadrootsRelayOutcomeKind,
- pub message: Option<String>,
-}
-
-impl RadrootsRelayOutcome {
- pub fn accepted() -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::Accepted,
- message: None,
- }
- }
-
- pub fn duplicate_accepted(message: impl Into<String>) -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::DuplicateAccepted,
- message: Some(message.into()),
- }
- }
-
- pub fn connection_failed(message: impl Into<String>) -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::ConnectionFailed,
- message: Some(message.into()),
- }
- }
-
- pub fn timeout(message: impl Into<String>) -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::Timeout,
- message: Some(message.into()),
- }
- }
-
- pub fn relay_url_rejected(message: impl Into<String>) -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::RelayUrlRejected,
- message: Some(message.into()),
- }
- }
-
- pub fn skipped_already_accepted(message: impl Into<String>) -> Self {
- Self {
- kind: RadrootsRelayOutcomeKind::SkippedAlreadyAccepted,
- message: Some(message.into()),
- }
- }
-
- pub fn classify(message: impl AsRef<str>) -> Self {
- let message = message.as_ref().trim();
- let lower = message.to_ascii_lowercase();
- let kind = if lower.starts_with("duplicate:") {
- RadrootsRelayOutcomeKind::DuplicateAccepted
- } else if lower.starts_with("blocked:") {
- RadrootsRelayOutcomeKind::Blocked
- } else if lower.starts_with("rate-limited:") {
- RadrootsRelayOutcomeKind::RateLimited
- } else if lower.starts_with("invalid:") {
- RadrootsRelayOutcomeKind::Invalid
- } else if lower.starts_with("pow:") {
- RadrootsRelayOutcomeKind::PowRequired
- } else if lower.starts_with("restricted:") {
- RadrootsRelayOutcomeKind::Restricted
- } else if lower.starts_with("auth-required:") {
- RadrootsRelayOutcomeKind::AuthRequired
- } else if lower.starts_with("mute:") {
- RadrootsRelayOutcomeKind::Muted
- } else if lower.starts_with("unsupported:") {
- RadrootsRelayOutcomeKind::Unsupported
- } else if lower.starts_with("payment-required:") {
- RadrootsRelayOutcomeKind::PaymentRequired
- } else if lower.starts_with("error:") {
- RadrootsRelayOutcomeKind::Error
- } else if lower.starts_with("timeout:") {
- RadrootsRelayOutcomeKind::Timeout
- } else {
- RadrootsRelayOutcomeKind::Unknown
- };
- Self {
- kind,
- message: Some(message.to_owned()),
- }
- }
-
- pub fn counts_toward_quorum(&self) -> bool {
- self.kind.counts_toward_quorum()
- }
-
- pub fn is_retryable(&self) -> bool {
- self.kind.is_retryable()
- }
-
- pub fn is_terminal_failure(&self) -> bool {
- self.kind.is_terminal_failure()
- }
-}
diff --git a/crates/relay_transport/src/publish.rs b/crates/relay_transport/src/publish.rs
@@ -1,444 +0,0 @@
-#![forbid(unsafe_code)]
-
-use crate::{RadrootsRelayOutcome, RadrootsRelayTargetSet, RadrootsRelayTransportError};
-#[cfg(feature = "client")]
-use core::time::Duration;
-use futures::future::BoxFuture;
-use radroots_events::draft::RadrootsSignedNostrEvent;
-use serde::{Deserialize, Serialize};
-use std::collections::{BTreeMap, BTreeSet};
-use std::sync::{Arc, Mutex, PoisonError};
-
-#[cfg(feature = "client")]
-use crate::RadrootsRelayOutcomeKind;
-#[cfg(feature = "client")]
-use nostr::JsonUtil;
-#[cfg(feature = "client")]
-use radroots_nostr::prelude::{RadrootsNostrClient, RadrootsNostrEvent};
-
-#[cfg(feature = "client")]
-const RELAY_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
-
-#[derive(Clone, Debug, PartialEq, Eq)]
-pub struct RadrootsRelayPublishRequest {
- pub signed_event: RadrootsSignedNostrEvent,
- pub targets: RadrootsRelayTargetSet,
- pub accepted_quorum: usize,
- pub now_ms: i64,
-}
-
-impl RadrootsRelayPublishRequest {
- pub fn new(
- signed_event: RadrootsSignedNostrEvent,
- targets: RadrootsRelayTargetSet,
- now_ms: i64,
- ) -> Self {
- let accepted_quorum = targets.len();
- Self {
- signed_event,
- targets,
- accepted_quorum,
- now_ms,
- }
- }
-
- pub fn with_accepted_quorum(mut self, accepted_quorum: usize) -> Self {
- self.accepted_quorum = accepted_quorum;
- self
- }
-}
-
-#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
-pub struct RadrootsRelayPublishRelayReceipt {
- pub relay_url: String,
- pub outcome: RadrootsRelayOutcome,
- pub attempted: bool,
-}
-
-impl RadrootsRelayPublishRelayReceipt {
- pub fn attempted(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self {
- Self {
- relay_url: relay_url.into(),
- outcome,
- attempted: true,
- }
- }
-
- pub fn skipped(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self {
- Self {
- relay_url: relay_url.into(),
- outcome,
- attempted: false,
- }
- }
-}
-
-#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
-pub struct RadrootsRelayPublishReceipt {
- pub event_id: String,
- pub attempted_count: usize,
- pub accepted_count: usize,
- pub retryable_count: usize,
- pub terminal_count: usize,
- pub quorum: usize,
- pub quorum_met: bool,
- pub relays: Vec<RadrootsRelayPublishRelayReceipt>,
-}
-
-pub trait RadrootsRelayPublishAdapter: Send + Sync {
- fn publish<'a>(
- &'a self,
- request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>;
-}
-
-pub async fn publish_signed_event<A>(
- adapter: &A,
- request: RadrootsRelayPublishRequest,
-) -> Result<RadrootsRelayPublishReceipt, RadrootsRelayTransportError>
-where
- A: RadrootsRelayPublishAdapter,
-{
- let event_id = request.signed_event.id.clone();
- let quorum = request.accepted_quorum;
- let relays = adapter.publish(request).await?;
- let attempted_count = relays.iter().filter(|receipt| receipt.attempted).count();
- let accepted_count = relays
- .iter()
- .filter(|receipt| receipt.outcome.counts_toward_quorum())
- .count();
- let retryable_count = relays
- .iter()
- .filter(|receipt| receipt.outcome.is_retryable())
- .count();
- let terminal_count = relays
- .iter()
- .filter(|receipt| receipt.outcome.is_terminal_failure())
- .count();
- Ok(RadrootsRelayPublishReceipt {
- event_id,
- attempted_count,
- accepted_count,
- retryable_count,
- terminal_count,
- quorum,
- quorum_met: accepted_count >= quorum,
- relays,
- })
-}
-
-#[derive(Clone, Default)]
-pub struct RadrootsMockRelayPublishAdapter {
- outcomes: BTreeMap<String, RadrootsRelayOutcome>,
- captured_raw_events: Arc<Mutex<Vec<String>>>,
-}
-
-impl RadrootsMockRelayPublishAdapter {
- pub fn new() -> Self {
- Self::default()
- }
-
- pub fn with_outcome(
- mut self,
- relay_url: impl Into<String>,
- outcome: RadrootsRelayOutcome,
- ) -> Self {
- self.outcomes.insert(relay_url.into(), outcome);
- self
- }
-
- pub fn captured_raw_events(&self) -> Vec<String> {
- self.captured_raw_events
- .lock()
- .expect("captured raw event lock")
- .clone()
- }
-}
-
-impl RadrootsRelayPublishAdapter for RadrootsMockRelayPublishAdapter {
- fn publish<'a>(
- &'a self,
- request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
- {
- Box::pin(async move {
- self.captured_raw_events
- .lock()
- .map_err(captured_raw_event_lock_error)?
- .push(request.signed_event.raw_json.clone());
- Ok(request
- .targets
- .relays()
- .iter()
- .map(|relay| {
- let outcome = self
- .outcomes
- .get(relay.as_str())
- .cloned()
- .unwrap_or_else(RadrootsRelayOutcome::accepted);
- RadrootsRelayPublishRelayReceipt::attempted(relay.as_str(), outcome)
- })
- .collect())
- })
- }
-}
-
-#[cfg_attr(coverage_nightly, coverage(off))]
-fn captured_raw_event_lock_error<T>(_error: PoisonError<T>) -> RadrootsRelayTransportError {
- RadrootsRelayTransportError::Transport("captured raw event lock poisoned".to_owned())
-}
-
-#[cfg(feature = "client")]
-#[derive(Clone)]
-pub struct RadrootsNostrClientPublishAdapter {
- client: RadrootsNostrClient,
-}
-
-#[cfg(feature = "client")]
-impl RadrootsNostrClientPublishAdapter {
- #[cfg_attr(coverage_nightly, coverage(off))]
- pub fn new(client: RadrootsNostrClient) -> Self {
- Self { client }
- }
-}
-
-#[cfg(feature = "client")]
-impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter {
- #[cfg_attr(coverage_nightly, coverage(off))]
- fn publish<'a>(
- &'a self,
- request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
- {
- Box::pin(async move {
- let event = RadrootsNostrEvent::from_json(request.signed_event.raw_json.as_str())
- .map_err(|error| RadrootsRelayTransportError::NostrEventJson(error.to_string()))?;
- ensure_raw_event_matches_signed_event(&event, &request.signed_event)?;
- let target_strings = request.targets.relay_strings();
- for relay_url in &target_strings {
- self.client
- .add_write_relay(relay_url.as_str())
- .await
- .map_err(|error| RadrootsRelayTransportError::Transport(error.to_string()))?;
- }
- let connection_output = self.client.try_connect(RELAY_CONNECT_TIMEOUT).await;
- let target_url_set = target_strings
- .iter()
- .map(|relay_url| relay_url.trim_end_matches('/').to_owned())
- .collect::<BTreeSet<_>>();
- let connected_strings = self
- .client
- .relays()
- .await
- .into_values()
- .filter(|relay| relay.is_connected())
- .map(|relay| relay.url().to_string())
- .filter(|relay_url| target_url_set.contains(relay_url.trim_end_matches('/')))
- .collect::<Vec<_>>();
- let connection_failures = connection_output
- .failed
- .iter()
- .map(|(relay, reason)| {
- (
- relay.to_string().trim_end_matches('/').to_owned(),
- reason.clone(),
- )
- })
- .collect::<BTreeMap<_, _>>();
- if connected_strings.is_empty() {
- return Ok(target_strings
- .into_iter()
- .map(|relay_url| {
- let target_url = relay_url.trim_end_matches('/');
- let reason = connection_failures
- .get(target_url)
- .cloned()
- .unwrap_or_else(|| "relay did not connect".to_owned());
- RadrootsRelayPublishRelayReceipt::attempted(
- relay_url,
- RadrootsRelayOutcome::connection_failed(reason),
- )
- })
- .collect());
- }
- let output = match self.client.send_event_to(connected_strings, &event).await {
- Ok(output) => output,
- Err(error) => {
- let message = error.to_string();
- return Ok(target_strings
- .into_iter()
- .map(|relay_url| {
- RadrootsRelayPublishRelayReceipt::attempted(
- relay_url,
- RadrootsRelayOutcome::connection_failed(message.clone()),
- )
- })
- .collect());
- }
- };
- let mut receipts = Vec::new();
- for relay_url in &target_strings {
- let target_url = relay_url.trim_end_matches('/');
- let success = output
- .success
- .iter()
- .any(|success_url| success_url.to_string().trim_end_matches('/') == target_url);
- if success {
- receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
- relay_url,
- RadrootsRelayOutcome {
- kind: RadrootsRelayOutcomeKind::Accepted,
- message: Some(
- "nostr-relay-pool-success-ok-message-unavailable".to_owned(),
- ),
- },
- ));
- continue;
- }
- if let Some(reason) = connection_failures.get(target_url) {
- receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
- relay_url,
- RadrootsRelayOutcome::connection_failed(reason.clone()),
- ));
- continue;
- }
- let failed = output.failed.iter().find_map(|(failed_url, message)| {
- if failed_url.to_string().trim_end_matches('/') == target_url {
- Some(message.clone())
- } else {
- None
- }
- });
- let outcome = failed
- .map(RadrootsRelayOutcome::classify)
- .unwrap_or_else(|| {
- RadrootsRelayOutcome::classify("error: relay output omitted target")
- });
- receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
- relay_url, outcome,
- ));
- }
- Ok(receipts)
- })
- }
-}
-
-#[cfg(feature = "client")]
-fn ensure_raw_event_matches_signed_event(
- event: &RadrootsNostrEvent,
- signed_event: &RadrootsSignedNostrEvent,
-) -> Result<(), RadrootsRelayTransportError> {
- let mismatches = [
- ("id", event.id.to_hex(), signed_event.id.clone()),
- ("pubkey", event.pubkey.to_hex(), signed_event.pubkey.clone()),
- (
- "created_at",
- event.created_at.as_secs().to_string(),
- signed_event.created_at.to_string(),
- ),
- (
- "kind",
- (event.kind.as_u16() as u32).to_string(),
- signed_event.kind.to_string(),
- ),
- (
- "content",
- event.content.clone(),
- signed_event.content.clone(),
- ),
- ("sig", event.sig.to_string(), signed_event.sig.clone()),
- ];
- for (field, raw, wrapped) in mismatches {
- if raw != wrapped {
- return Err(RadrootsRelayTransportError::NostrEventJson(format!(
- "raw event JSON {field} does not match signed event {field}"
- )));
- }
- }
- let raw_tags = event
- .tags
- .iter()
- .map(|tag| tag.as_slice().to_vec())
- .collect::<Vec<_>>();
- if raw_tags != signed_event.tags {
- return Err(RadrootsRelayTransportError::NostrEventJson(
- "raw event JSON tags do not match signed event tags".to_owned(),
- ));
- }
- Ok(())
-}
-
-#[cfg(all(test, feature = "client"))]
-mod tests {
- use super::{RadrootsNostrEvent, ensure_raw_event_matches_signed_event};
- use nostr::JsonUtil;
- use radroots_events::draft::{RadrootsFrozenEventDraft, RadrootsSignedNostrEvent};
- use radroots_events::kinds::KIND_POST;
- use radroots_nostr::prelude::{
- RadrootsNostrKeys, RadrootsNostrSecretKey, radroots_nostr_sign_frozen_draft,
- };
-
- const FIXTURE_ALICE_SECRET_KEY_HEX: &str =
- "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5";
- const FIXTURE_ALICE_PUBLIC_KEY_HEX: &str =
- "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
-
- fn signed_post(content: &str) -> (RadrootsNostrEvent, RadrootsSignedNostrEvent) {
- let secret_key =
- RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key");
- let keys = RadrootsNostrKeys::new(secret_key);
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- 1_700_000_000,
- vec![vec!["t".to_owned(), "soil".to_owned()]],
- content,
- FIXTURE_ALICE_PUBLIC_KEY_HEX,
- )
- .expect("draft");
- let signed_event = radroots_nostr_sign_frozen_draft(&keys, &draft).expect("signed event");
- let raw_event =
- RadrootsNostrEvent::from_json(signed_event.raw_json.as_str()).expect("raw event");
- (raw_event, signed_event)
- }
-
- fn assert_mismatch(raw_event: &RadrootsNostrEvent, signed_event: RadrootsSignedNostrEvent) {
- assert!(ensure_raw_event_matches_signed_event(raw_event, &signed_event).is_err());
- }
-
- #[test]
- fn raw_event_match_guard_accepts_exact_event_and_rejects_field_mismatches() {
- let (raw_event, signed_event) = signed_post("matched");
- ensure_raw_event_matches_signed_event(&raw_event, &signed_event).expect("matching event");
-
- let mut mismatched = signed_event.clone();
- mismatched.id = "00".repeat(32);
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event.clone();
- mismatched.pubkey = "11".repeat(32);
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event.clone();
- mismatched.created_at += 1;
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event.clone();
- mismatched.kind += 1;
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event.clone();
- mismatched.content.push_str(" changed");
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event.clone();
- mismatched.sig = "22".repeat(64);
- assert_mismatch(&raw_event, mismatched);
-
- let mut mismatched = signed_event;
- mismatched
- .tags
- .push(vec!["t".to_owned(), "compost".to_owned()]);
- assert_mismatch(&raw_event, mismatched);
- }
-}
diff --git a/crates/relay_transport/tests/transport.rs b/crates/relay_transport/tests/transport.rs
@@ -1,1857 +0,0 @@
-use futures::future::BoxFuture;
-use nostr::JsonUtil;
-use radroots_event_store::{
- RadrootsEventStore, RadrootsEventVerificationStatus, RadrootsTransportObservationType,
-};
-use radroots_events::draft::{RadrootsFrozenEventDraft, RadrootsSignedNostrEvent};
-use radroots_events::kinds::KIND_POST;
-use radroots_nostr::prelude::{
- RadrootsNostrFilter, RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrSecretKey,
- RadrootsNostrTimestamp, radroots_nostr_build_event, radroots_nostr_filter_tag,
- radroots_nostr_sign_frozen_draft,
-};
-use radroots_outbox::{
- RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxEventState,
- RadrootsOutboxOperationInput, RadrootsOutboxOperationStatus, RadrootsOutboxRelayStatus,
-};
-use radroots_relay_transport::{
- RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsOutboxPublishPolicy,
- RadrootsRelayFetchFilters, RadrootsRelayFetchItem, RadrootsRelayFetchMode,
- RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchRequest, RadrootsRelayOutcome,
- RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt,
- RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError,
- RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events,
- fetch_relay_events_blocking, publish_claimed_outbox_event, publish_signed_event,
-};
-use radroots_transport::RadrootsTransportKind;
-use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
-
-const FIXTURE_ALICE_SECRET_KEY_HEX: &str =
- "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5";
-const FIXTURE_ALICE_PUBLIC_KEY_HEX: &str =
- "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
-const RELAY_PRIMARY_WSS: &str = "wss://relay.example.com";
-const RELAY_SECONDARY_WSS: &str = "wss://relay-2.example.com";
-const RELAY_TERTIARY_WSS: &str = "wss://relay-3.example.com";
-
-struct TransportFailurePublishAdapter;
-
-impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter {
- fn publish<'a>(
- &'a self,
- _request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
- {
- Box::pin(async {
- Err(RadrootsRelayTransportError::Transport(
- "adapter boundary unavailable".to_owned(),
- ))
- })
- }
-}
-
-struct NostrJsonFailurePublishAdapter;
-
-impl RadrootsRelayPublishAdapter for NostrJsonFailurePublishAdapter {
- fn publish<'a>(
- &'a self,
- _request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
- {
- Box::pin(async {
- Err(RadrootsRelayTransportError::NostrEventJson(
- "adapter rejected raw event".to_owned(),
- ))
- })
- }
-}
-
-fn fixture_keys() -> RadrootsNostrKeys {
- let secret_key =
- RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key");
- RadrootsNostrKeys::new(secret_key)
-}
-
-fn signed_post(content: &str) -> RadrootsSignedNostrEvent {
- signed_event_with_kind_and_hashtag(content, KIND_POST, "soil")
-}
-
-fn signed_event_with_kind_and_hashtag(
- content: &str,
- kind: u32,
- hashtag: &str,
-) -> RadrootsSignedNostrEvent {
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- kind,
- 1_700_000_000,
- vec![vec!["t".to_owned(), hashtag.to_owned()]],
- content,
- FIXTURE_ALICE_PUBLIC_KEY_HEX,
- )
- .expect("draft");
- radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event")
-}
-
-fn signed_raw_event_with_kind_and_hashtag(content: &str, kind: u32, hashtag: &str) -> nostr::Event {
- radroots_nostr_build_event(
- kind,
- content,
- vec![vec!["t".to_owned(), hashtag.to_owned()]],
- )
- .expect("event builder")
- .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000))
- .sign_with_keys(&fixture_keys())
- .expect("signed event")
-}
-
-async fn complete_claimed_signing(
- outbox: &RadrootsOutbox,
- claimed: &RadrootsOutboxClaimedEvent,
- now_ms: i64,
-) -> RadrootsSignedNostrEvent {
- if let Some(signed_event) = claimed.signed_event.clone() {
- return signed_event;
- }
- let signed_event =
- radroots_nostr_sign_frozen_draft(&fixture_keys(), &claimed.draft).expect("signed event");
- outbox
- .complete_signing(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- signed_event,
- now_ms,
- )
- .await
- .expect("complete signing")
-}
-
-fn unsupported_raw_event() -> String {
- let event = radroots_nostr_build_event(999, "unsupported", Vec::new())
- .expect("event builder")
- .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001))
- .sign_with_keys(&fixture_keys())
- .expect("signed unsupported event");
- event.as_json()
-}
-
-fn post_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter {
- radroots_nostr_filter_tag(
- RadrootsNostrFilter::new()
- .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
- .limit(limit),
- "t",
- vec!["soil".to_owned()],
- )
- .expect("post relay fetch filter")
-}
-
-fn unsupported_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter {
- RadrootsNostrFilter::new()
- .kind(RadrootsNostrKind::Custom(999))
- .limit(limit)
-}
-
-fn fixture_relay_fetch_request(
- observed_at_ms: i64,
- max_events: usize,
-) -> RadrootsRelayFetchRequest {
- RadrootsRelayFetchRequest::fetch(
- observed_at_ms,
- max_events,
- [
- post_relay_fetch_filter(max_events),
- unsupported_relay_fetch_filter(max_events),
- ],
- )
- .expect("fixture relay fetch request")
-}
-
-fn post_relay_fetch_request(observed_at_ms: i64, max_events: usize) -> RadrootsRelayFetchRequest {
- RadrootsRelayFetchRequest::fetch(
- observed_at_ms,
- max_events,
- [post_relay_fetch_filter(max_events)],
- )
- .expect("post relay fetch request")
-}
-
-fn tampered_raw_event() -> String {
- let signed = signed_post("trusted");
- let mut value =
- serde_json::from_str::<serde_json::Value>(signed.raw_json.as_str()).expect("raw json");
- value["content"] = serde_json::Value::String("tampered".to_owned());
- serde_json::to_string(&value).expect("tampered json")
-}
-
-#[test]
-fn relay_url_validation_and_target_normalization() {
- let relay = RadrootsRelayUrl::parse("wss://Relay.Example.com", RadrootsRelayUrlPolicy::Public)
- .expect("relay");
- assert_eq!(relay.as_str(), RELAY_PRIMARY_WSS);
- assert_eq!(relay.clone().into_string(), RELAY_PRIMARY_WSS);
- let relay_path = RadrootsRelayUrl::parse(
- "wss://Relay.Example.com/nostr",
- RadrootsRelayUrlPolicy::Public,
- )
- .expect("relay path");
- assert_eq!(relay_path.as_str(), "wss://relay.example.com/nostr");
-
- assert!(
- RadrootsRelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Public).is_err()
- );
- let local = RadrootsRelayUrl::parse("ws://localhost:7777", RadrootsRelayUrlPolicy::Localhost)
- .expect("local relay");
- assert_eq!(local.as_str(), "ws://localhost:7777");
- let local_ipv4 =
- RadrootsRelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Localhost)
- .expect("local ipv4 relay");
- assert_eq!(local_ipv4.as_str(), "ws://127.0.0.1:7777");
- let local_ipv6 = RadrootsRelayUrl::parse("ws://[::1]:7777", RadrootsRelayUrlPolicy::Localhost)
- .expect("local ipv6 relay");
- assert_eq!(local_ipv6.as_str(), "ws://[::1]:7777");
- assert!(
- RadrootsRelayUrl::parse("ws://example.com", RadrootsRelayUrlPolicy::Localhost).is_err()
- );
- assert!(
- RadrootsRelayUrl::parse("ws://192.168.1.10:7777", RadrootsRelayUrlPolicy::Localhost)
- .is_err()
- );
- assert!(matches!(
- RadrootsRelayUrl::parse("wss://127.0.0.1", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
- ));
- assert!(matches!(
- RadrootsRelayUrl::parse("wss://10.1.2.3", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
- ));
- assert!(matches!(
- RadrootsRelayUrl::parse("wss://[::1]", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
- ));
- assert!(matches!(
- RadrootsRelayUrl::parse("wss://[fd00::1]", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
- ));
- for relay_url in [
- "wss://0.0.0.0",
- "wss://169.254.1.2",
- "wss://224.0.0.1",
- "wss://255.255.255.255",
- "wss://100.64.0.1",
- "wss://192.0.0.8",
- "wss://198.18.0.1",
- "wss://240.0.0.1",
- "wss://[::]",
- "wss://[ff02::1]",
- "wss://[fe80::1]",
- "wss://[2001:db8::1]",
- "wss://[2001:1::1]",
- "wss://[::ffff:192.168.1.10]",
- ] {
- assert!(matches!(
- RadrootsRelayUrl::parse(relay_url, RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
- ));
- }
- let public_relay =
- RadrootsRelayUrl::parse("wss://relay.example.com", RadrootsRelayUrlPolicy::Public)
- .expect("public relay");
- public_relay
- .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))])
- .expect("public resolved ip");
- assert!(matches!(
- public_relay
- .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10))]),
- Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. })
- ));
- assert!(matches!(
- public_relay.validate_public_resolved_ip_addrs([IpAddr::V6(
- "::ffff:192.168.1.10"
- .parse::<Ipv6Addr>()
- .expect("mapped ipv6")
- )]),
- Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. })
- ));
- public_relay
- .validate_public_resolved_ip_addrs([IpAddr::V6(
- "2001:4860:4860::8888"
- .parse::<Ipv6Addr>()
- .expect("public ipv6"),
- )])
- .expect("public resolved ipv6");
- public_relay
- .validate_public_resolved_ip_addrs(Vec::<IpAddr>::new())
- .expect("empty resolved set");
-
- assert!(
- RadrootsRelayUrl::parse("https://relay.example.com", RadrootsRelayUrlPolicy::Public)
- .is_err()
- );
- assert!(
- RadrootsRelayUrl::parse(
- "wss://user@relay.example.com",
- RadrootsRelayUrlPolicy::Public
- )
- .is_err()
- );
- assert!(matches!(
- RadrootsRelayUrl::parse(
- "wss://user:password@relay.example.com",
- RadrootsRelayUrlPolicy::Public
- ),
- Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. })
- ));
- assert!(matches!(
- RadrootsRelayUrl::parse(
- "wss://:password@relay.example.com",
- RadrootsRelayUrlPolicy::Public
- ),
- Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. })
- ));
- assert!(
- RadrootsRelayUrl::parse(
- "wss://relay.example.com:bad",
- RadrootsRelayUrlPolicy::Public
- )
- .is_err()
- );
- assert!(RadrootsRelayUrl::parse("wss://", RadrootsRelayUrlPolicy::Public).is_err());
- assert!(matches!(
- RadrootsRelayUrl::parse("radroots:relay", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::EmptyRelayHost { .. })
- ));
- assert!(matches!(
- RadrootsRelayUrl::parse("relay.example.com", RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::RelayUrlParse { .. })
- ));
- assert!(
- RadrootsRelayUrl::parse(
- "wss://relay.example.com?subscription=1",
- RadrootsRelayUrlPolicy::Public
- )
- .is_err()
- );
- assert!(
- RadrootsRelayUrl::parse(
- "wss://relay.example.com#fragment",
- RadrootsRelayUrlPolicy::Public
- )
- .is_err()
- );
-
- let targets = RadrootsRelayTargetSet::new(
- vec![
- RELAY_TERTIARY_WSS,
- RELAY_PRIMARY_WSS,
- RELAY_PRIMARY_WSS,
- RELAY_SECONDARY_WSS,
- ],
- RadrootsRelayUrlPolicy::Public,
- )
- .expect("targets");
- assert_eq!(
- targets.relay_strings(),
- vec![
- RELAY_TERTIARY_WSS.to_owned(),
- RELAY_PRIMARY_WSS.to_owned(),
- RELAY_SECONDARY_WSS.to_owned()
- ]
- );
-
- let from_urls = RadrootsRelayTargetSet::from_urls(vec![
- relay_path.clone(),
- relay_path.clone(),
- RadrootsRelayUrl::parse(RELAY_SECONDARY_WSS, RadrootsRelayUrlPolicy::Public)
- .expect("secondary"),
- ])
- .expect("from urls");
- assert_eq!(from_urls.len(), 2);
- assert!(!from_urls.is_empty());
- assert_eq!(from_urls.relays()[0], relay_path);
- assert_eq!(
- from_urls.relays()[0].to_string(),
- "wss://relay.example.com/nostr"
- );
- assert!(matches!(
- RadrootsRelayTargetSet::new(Vec::<&str>::new(), RadrootsRelayUrlPolicy::Public),
- Err(RadrootsRelayTransportError::EmptyTargetSet)
- ));
- assert!(matches!(
- RadrootsRelayTargetSet::from_urls(Vec::new()),
- Err(RadrootsRelayTransportError::EmptyTargetSet)
- ));
-}
-
-#[test]
-fn outcome_prefix_classification_covers_required_kinds() {
- let cases = [
- ("blocked: policy", RadrootsRelayOutcomeKind::Blocked),
- (
- "rate-limited: slow down",
- RadrootsRelayOutcomeKind::RateLimited,
- ),
- ("invalid: bad event", RadrootsRelayOutcomeKind::Invalid),
- ("pow: difficulty 24", RadrootsRelayOutcomeKind::PowRequired),
- (
- "restricted: group write denied",
- RadrootsRelayOutcomeKind::Restricted,
- ),
- (
- "auth-required: challenge",
- RadrootsRelayOutcomeKind::AuthRequired,
- ),
- ("mute: pubkey muted", RadrootsRelayOutcomeKind::Muted),
- (
- "unsupported: event kind",
- RadrootsRelayOutcomeKind::Unsupported,
- ),
- (
- "payment-required: paid relay",
- RadrootsRelayOutcomeKind::PaymentRequired,
- ),
- (
- "duplicate: already have it",
- RadrootsRelayOutcomeKind::DuplicateAccepted,
- ),
- ("error: relay failed", RadrootsRelayOutcomeKind::Error),
- ("timeout: no OK", RadrootsRelayOutcomeKind::Timeout),
- ("strange relay text", RadrootsRelayOutcomeKind::Unknown),
- ];
-
- for (message, kind) in cases {
- let outcome = RadrootsRelayOutcome::classify(message);
- assert_eq!(outcome.kind, kind);
- }
-
- assert!(RadrootsRelayOutcome::classify("duplicate: already have it").counts_toward_quorum());
- assert!(
- RadrootsRelayOutcome::skipped_already_accepted("already accepted").counts_toward_quorum()
- );
- assert!(RadrootsRelayOutcome::classify("auth-required: challenge").is_retryable());
- assert!(RadrootsRelayOutcome::classify("restricted: denied").is_terminal_failure());
- assert!(RadrootsRelayOutcome::relay_url_rejected("unsafe relay").is_terminal_failure());
- assert!(RadrootsRelayOutcome::classify("mute: pubkey muted").is_terminal_failure());
-}
-
-#[tokio::test]
-async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() {
- let signed = signed_post("hello");
- let targets = RadrootsRelayTargetSet::new(
- vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS],
- RadrootsRelayUrlPolicy::Public,
- )
- .expect("targets");
- let adapter = RadrootsMockRelayPublishAdapter::new()
- .with_outcome(
- RELAY_SECONDARY_WSS,
- RadrootsRelayOutcome::classify("duplicate: already have it"),
- )
- .with_outcome(
- RELAY_TERTIARY_WSS,
- RadrootsRelayOutcome::classify("auth-required: challenge"),
- );
-
- let receipt = publish_signed_event(
- &adapter,
- radroots_relay_transport::RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_000)
- .with_accepted_quorum(2),
- )
- .await
- .expect("publish");
-
- assert_eq!(adapter.captured_raw_events(), vec![signed.raw_json]);
- assert_eq!(receipt.attempted_count, 3);
- assert_eq!(receipt.accepted_count, 2);
- assert_eq!(receipt.retryable_count, 1);
- assert!(receipt.quorum_met);
- serde_json::to_string(&receipt).expect("receipt json");
-}
-
-#[tokio::test]
-async fn publish_receipts_track_terminal_skipped_and_adapter_errors() {
- let signed = signed_post("terminal");
- let targets = RadrootsRelayTargetSet::new(
- vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
- RadrootsRelayUrlPolicy::Public,
- )
- .expect("targets");
- let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome(
- RELAY_SECONDARY_WSS,
- RadrootsRelayOutcome::classify("restricted: group write denied"),
- );
-
- let receipt = publish_signed_event(
- &adapter,
- RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_050).with_accepted_quorum(2),
- )
- .await
- .expect("publish");
-
- assert_eq!(receipt.event_id, signed.id);
- assert_eq!(receipt.attempted_count, 2);
- assert_eq!(receipt.accepted_count, 1);
- assert_eq!(receipt.retryable_count, 0);
- assert_eq!(receipt.terminal_count, 1);
- assert_eq!(receipt.quorum, 2);
- assert!(!receipt.quorum_met);
-
- let skipped = RadrootsRelayPublishRelayReceipt::skipped(
- RELAY_TERTIARY_WSS,
- RadrootsRelayOutcome::timeout("timeout: no OK"),
- );
- assert_eq!(skipped.relay_url, RELAY_TERTIARY_WSS);
- assert!(!skipped.attempted);
- assert_eq!(skipped.outcome.kind, RadrootsRelayOutcomeKind::Timeout);
-
- let error = publish_signed_event(
- &TransportFailurePublishAdapter,
- RadrootsRelayPublishRequest::new(
- signed,
- RadrootsRelayTargetSet::new(vec![RELAY_PRIMARY_WSS], RadrootsRelayUrlPolicy::Public)
- .expect("targets"),
- 1_060,
- ),
- )
- .await
- .expect_err("transport failure");
- assert!(matches!(error, RadrootsRelayTransportError::Transport(_)));
-}
-
-#[test]
-fn fetch_requests_reject_empty_filter_sets() {
- assert!(matches!(
- RadrootsRelayFetchRequest::fetch(1_000, 10, Vec::<RadrootsNostrFilter>::new()),
- Err(RadrootsRelayTransportError::EmptyFetchFilters)
- ));
- assert!(matches!(
- RadrootsRelayFetchRequest::subscription(1_000, 10, Vec::<RadrootsNostrFilter>::new()),
- Err(RadrootsRelayTransportError::EmptyFetchFilters)
- ));
-}
-
-#[test]
-fn fetch_requests_reject_zero_limits_and_timeouts() {
- let filter = post_relay_fetch_filter(1);
- let filters = RadrootsRelayFetchFilters::new([filter.clone()]).expect("filters");
- let as_ref_filters: &[RadrootsNostrFilter] = filters.as_ref();
- assert_eq!(as_ref_filters.len(), 1);
-
- assert!(matches!(
- RadrootsRelayFetchRequest::fetch(1_000, 0, [filter.clone()]),
- Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events"
- ));
- assert!(matches!(
- RadrootsRelayFetchRequest::subscription(1_000, 0, [filter.clone()]),
- Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events"
- ));
-
- let request =
- RadrootsRelayFetchRequest::fetch(1_000, 1, [filter]).expect("valid fetch request");
- assert!(matches!(
- request.clone().with_timeout_ms(0),
- Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "timeout_ms"
- ));
- assert!(matches!(
- request.clone().with_raw_event_scan_limit(0),
- Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_raw_events"
- ));
-
- let request = request
- .with_timeout_ms(1)
- .expect("minimum timeout")
- .with_raw_event_scan_limit(1)
- .expect("minimum raw scan limit");
- assert_eq!(request.timeout_ms(), 1);
- assert_eq!(request.max_raw_events(), 1);
-
- let request = RadrootsRelayFetchRequest::subscription(1_005, 2, [post_relay_fetch_filter(2)])
- .expect("subscription request")
- .with_relay_urls([RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS])
- .with_timeout_ms(25)
- .expect("timeout")
- .with_raw_event_scan_limit(3)
- .expect("raw limit");
- assert_eq!(request.mode(), RadrootsRelayFetchMode::Subscription);
- assert_eq!(request.observed_at_ms(), 1_005);
- assert_eq!(request.max_events(), 2);
- assert_eq!(request.max_raw_events(), 3);
- assert_eq!(
- request.relay_urls(),
- &[RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()]
- );
- assert_eq!(request.filters().len(), 1);
- assert_eq!(request.timeout_ms(), 25);
-}
-
-#[test]
-fn fetch_blocking_facade_runs_mock_adapter() {
- let signed = signed_post("blocking fetch");
- let accepted_id = signed.id.clone();
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: signed.raw_json,
- observed_at_ms: 1_090,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- ]);
-
- let receipt = fetch_relay_events_blocking(&adapter, post_relay_fetch_request(1_090, 10))
- .expect("blocking fetch");
-
- assert_eq!(receipt.events.len(), 1);
- assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id);
- assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]);
-}
-
-#[tokio::test]
-async fn fetch_ingests_events_and_records_transport_observations() {
- let signed = signed_post("hello");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: signed.raw_json.clone(),
- observed_at_ms: 1_000,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: signed.raw_json.clone(),
- observed_at_ms: 1_001,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- raw_json: unsupported_raw_event(),
- observed_at_ms: 1_002,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- raw_json: tampered_raw_event(),
- observed_at_ms: 1_003,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_TERTIARY_WSS.to_owned(),
- raw_json: "{not json".to_owned(),
- observed_at_ms: 1_004,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- RadrootsRelayFetchItem::Closed {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- message: "auth-required: challenge".to_owned(),
- },
- RadrootsRelayFetchItem::Closed {
- relay_url: RELAY_TERTIARY_WSS.to_owned(),
- message: "restricted: group write denied".to_owned(),
- },
- RadrootsRelayFetchItem::Notice {
- relay_url: RELAY_TERTIARY_WSS.to_owned(),
- message: "notice: test".to_owned(),
- },
- ]);
-
- let receipt =
- fetch_and_ingest_relay_events(&adapter, &store, fixture_relay_fetch_request(1_000, 10))
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 3);
- assert_eq!(receipt.duplicate_count, 1);
- assert_eq!(receipt.unsupported_count, 1);
- assert_eq!(receipt.malformed_count, 1);
- assert_eq!(receipt.eose_count, 1);
- assert_eq!(receipt.closed_count, 2);
- assert_eq!(receipt.notice_count, 1);
- assert_eq!(receipt.relay_outcomes.len(), 4);
- assert_eq!(receipt.relay_outcomes[0].relay_url, RELAY_PRIMARY_WSS);
- assert_eq!(
- receipt.relay_outcomes[0].kind,
- RadrootsRelayFetchOutcomeKind::Eose
- );
- assert!(receipt.relay_outcomes[0].relay_outcome.is_none());
- assert_eq!(receipt.relay_outcomes[1].relay_url, RELAY_SECONDARY_WSS);
- assert_eq!(
- receipt.relay_outcomes[1]
- .relay_outcome
- .as_ref()
- .expect("auth outcome")
- .kind,
- RadrootsRelayOutcomeKind::AuthRequired
- );
- assert_eq!(receipt.relay_outcomes[2].relay_url, RELAY_TERTIARY_WSS);
- assert_eq!(
- receipt.relay_outcomes[2]
- .relay_outcome
- .as_ref()
- .expect("restricted outcome")
- .kind,
- RadrootsRelayOutcomeKind::Restricted
- );
- assert_eq!(
- receipt.relay_outcomes[3].kind,
- RadrootsRelayFetchOutcomeKind::Notice
- );
- assert!(receipt.relay_outcomes[3].relay_outcome.is_none());
- assert_eq!(
- receipt.events[0].verification_status.as_deref(),
- Some(RadrootsEventVerificationStatus::Verified.as_str())
- );
- assert!(receipt.events[0].projection_eligible);
- assert_eq!(
- receipt.events[1].verification_status.as_deref(),
- Some(RadrootsEventVerificationStatus::Verified.as_str())
- );
- assert!(!receipt.events[1].projection_eligible);
- assert_eq!(
- receipt.events[2].verification_status.as_deref(),
- Some(RadrootsEventVerificationStatus::Verified.as_str())
- );
- assert!(!receipt.events[2].projection_eligible);
- assert_eq!(
- receipt.events[3].verification_status.as_deref(),
- Some(RadrootsEventVerificationStatus::IdMismatch.as_str())
- );
- assert!(!receipt.events[3].projection_eligible);
- assert_eq!(receipt.events[4].verification_status, None);
- assert!(!receipt.events[4].projection_eligible);
-
- let observations = store
- .observations_for_event(signed.id.as_str())
- .await
- .expect("observations");
- assert_eq!(observations.len(), 1);
- assert_eq!(observations[0].transport_kind, RadrootsTransportKind::Nostr);
- assert_eq!(observations[0].endpoint_uri.as_str(), RELAY_PRIMARY_WSS);
- assert_eq!(
- observations[0].observation_type,
- RadrootsTransportObservationType::NostrFetch
- );
- assert_eq!(observations[0].observation_count, 2);
-}
-
-#[tokio::test]
-async fn fetch_rejects_out_of_filter_events_before_store_mutation() {
- let accepted = signed_post("filter match");
- let wrong_tag = signed_event_with_kind_and_hashtag("filter wrong tag", KIND_POST, "compost");
- let wrong_kind = signed_raw_event_with_kind_and_hashtag("filter wrong kind", 999, "soil");
- let wrong_kind_event_id = wrong_kind.id.to_hex();
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: wrong_tag.raw_json.clone(),
- observed_at_ms: 1_005,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: accepted.raw_json.clone(),
- observed_at_ms: 1_006,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- raw_json: wrong_kind.as_json(),
- observed_at_ms: 1_007,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- ]);
- let filter = radroots_nostr_filter_tag(
- RadrootsNostrFilter::new()
- .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
- .limit(10),
- "t",
- vec!["soil".to_owned()],
- )
- .expect("filter");
-
- let receipt = fetch_and_ingest_relay_events(
- &adapter,
- &store,
- RadrootsRelayFetchRequest::fetch(1_005, 10, [filter]).expect("fetch request"),
- )
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 1);
- assert_eq!(receipt.out_of_filter_count, 2);
- assert_eq!(receipt.malformed_count, 0);
- assert_eq!(receipt.unsupported_count, 0);
- assert_eq!(receipt.events.len(), 3);
- assert!(receipt.events[0].out_of_filter);
- assert!(!receipt.events[1].out_of_filter);
- assert!(receipt.events[2].out_of_filter);
- assert!(
- store
- .get_event(accepted.id.as_str())
- .await
- .expect("accepted lookup")
- .is_some()
- );
- assert!(
- store
- .get_event(wrong_tag.id.as_str())
- .await
- .expect("wrong tag lookup")
- .is_none()
- );
- assert!(
- store
- .get_event(wrong_kind_event_id.as_str())
- .await
- .expect("wrong kind lookup")
- .is_none()
- );
-}
-
-#[tokio::test]
-async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_control_outcomes() {
- let accepted = signed_post("accepted capped event");
- let skipped = signed_post("skipped capped event");
- let wrong_tag = signed_event_with_kind_and_hashtag("wrong capped tag", KIND_POST, "compost");
- let accepted_id = accepted.id.clone();
- let skipped_id = skipped.id.clone();
- let wrong_tag_id = wrong_tag.id.clone();
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: "{not json".to_owned(),
- observed_at_ms: 1_099,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: wrong_tag.raw_json,
- observed_at_ms: 1_100,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: accepted.raw_json.clone(),
- observed_at_ms: 1_101,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: skipped.raw_json,
- observed_at_ms: 1_102,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- RadrootsRelayFetchItem::Closed {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- message: "auth-required: challenge".to_owned(),
- },
- RadrootsRelayFetchItem::Notice {
- relay_url: RELAY_TERTIARY_WSS.to_owned(),
- message: "notice: still visible".to_owned(),
- },
- ]);
-
- let receipt =
- fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_100, 1))
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 1);
- assert_eq!(receipt.duplicate_count, 0);
- assert_eq!(receipt.unsupported_count, 0);
- assert_eq!(receipt.malformed_count, 1);
- assert_eq!(receipt.out_of_filter_count, 1);
- assert_eq!(receipt.skipped_over_limit_count, 1);
- assert_eq!(receipt.events.len(), 4);
- assert!(receipt.events[0].malformed);
- assert!(receipt.events[1].out_of_filter);
- assert!(receipt.events[2].inserted);
- assert!(receipt.events[3].skipped_over_limit);
- assert_eq!(receipt.eose_count, 1);
- assert_eq!(receipt.closed_count, 1);
- assert_eq!(receipt.notice_count, 1);
- assert_eq!(receipt.relay_outcomes.len(), 3);
- assert_eq!(
- receipt.relay_outcomes[0].kind,
- RadrootsRelayFetchOutcomeKind::Eose
- );
- assert_eq!(
- receipt.relay_outcomes[1]
- .relay_outcome
- .as_ref()
- .expect("closed outcome")
- .kind,
- RadrootsRelayOutcomeKind::AuthRequired
- );
- assert_eq!(
- receipt.relay_outcomes[2].kind,
- RadrootsRelayFetchOutcomeKind::Notice
- );
- assert!(
- store
- .get_event(accepted_id.as_str())
- .await
- .expect("accepted lookup")
- .is_some()
- );
- assert!(
- store
- .get_event(skipped_id.as_str())
- .await
- .expect("skipped lookup")
- .is_none()
- );
- assert!(
- store
- .get_event(wrong_tag_id.as_str())
- .await
- .expect("wrong tag lookup")
- .is_none()
- );
-}
-
-#[tokio::test]
-async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() {
- let accepted = signed_event_with_kind_and_hashtag("shared fetch accepted", KIND_POST, "soil");
- let skipped = signed_event_with_kind_and_hashtag("shared fetch skipped", KIND_POST, "soil");
- let wrong_tag =
- signed_event_with_kind_and_hashtag("shared fetch wrong tag", KIND_POST, "compost");
- let filter = radroots_nostr_filter_tag(
- RadrootsNostrFilter::new()
- .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
- .limit(10),
- "t",
- vec!["soil".to_owned()],
- )
- .expect("filter");
- let accepted_id = accepted.id.clone();
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: "{not json".to_owned(),
- observed_at_ms: 2_100,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: wrong_tag.raw_json,
- observed_at_ms: 2_101,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: accepted.raw_json.clone(),
- observed_at_ms: 2_102,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: skipped.raw_json,
- observed_at_ms: 2_103,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- RadrootsRelayFetchItem::Closed {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- message: "auth-required: challenge".to_owned(),
- },
- RadrootsRelayFetchItem::Notice {
- relay_url: RELAY_TERTIARY_WSS.to_owned(),
- message: "notice: still visible".to_owned(),
- },
- ]);
-
- let receipt = fetch_relay_events(
- &adapter,
- RadrootsRelayFetchRequest::fetch(2_100, 1, [filter])
- .expect("fetch request")
- .with_relay_urls([RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS]),
- )
- .await
- .expect("fetch events");
-
- assert_eq!(
- receipt.target_relays,
- vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS]
- );
- assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]);
- assert_eq!(receipt.failed_relays.len(), 1);
- assert_eq!(receipt.failed_relays[0].relay_url, RELAY_SECONDARY_WSS);
- assert_eq!(receipt.events.len(), 1);
- assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id);
- assert_eq!(receipt.malformed_count, 1);
- assert_eq!(receipt.out_of_filter_count, 1);
- assert_eq!(receipt.skipped_over_limit_count, 1);
- assert_eq!(receipt.eose_count, 1);
- assert_eq!(receipt.closed_count, 1);
- assert_eq!(receipt.notice_count, 1);
- assert_eq!(receipt.event_receipts.len(), 4);
- assert!(receipt.event_receipts[0].malformed);
- assert!(receipt.event_receipts[1].out_of_filter);
- assert!(!receipt.event_receipts[2].malformed);
- assert!(receipt.event_receipts[3].skipped_over_limit);
-}
-
-#[tokio::test]
-async fn fetch_raw_scan_limit_bounds_noisy_adapter_output() {
- let accepted = signed_post("raw scan accepted event");
- let wrong_tag = signed_event_with_kind_and_hashtag("raw scan wrong tag", KIND_POST, "compost");
- let accepted_id = accepted.id.clone();
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: "{not json".to_owned(),
- observed_at_ms: 1_130,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: wrong_tag.raw_json,
- observed_at_ms: 1_131,
- },
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: accepted.raw_json,
- observed_at_ms: 1_132,
- },
- RadrootsRelayFetchItem::Eose {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- },
- ]);
-
- let receipt = fetch_and_ingest_relay_events(
- &adapter,
- &store,
- post_relay_fetch_request(1_130, 1)
- .with_raw_event_scan_limit(2)
- .expect("raw scan limit"),
- )
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 0);
- assert_eq!(receipt.malformed_count, 1);
- assert_eq!(receipt.out_of_filter_count, 1);
- assert_eq!(receipt.skipped_over_limit_count, 1);
- assert_eq!(receipt.events.len(), 2);
- assert_eq!(receipt.eose_count, 1);
- assert!(
- store
- .get_event(accepted_id.as_str())
- .await
- .expect("accepted lookup")
- .is_none()
- );
-}
-
-#[tokio::test]
-async fn fetch_subscription_mode_and_store_errors_are_reported() {
- let signed = signed_post("subscription");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: signed.raw_json.clone(),
- observed_at_ms: 1_200,
- }]);
-
- let receipt = fetch_and_ingest_relay_events(
- &adapter,
- &store,
- RadrootsRelayFetchRequest::subscription(1_200, 10, [post_relay_fetch_filter(10)])
- .expect("subscription request"),
- )
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 1);
- let observations = store
- .observations_for_event(signed.id.as_str())
- .await
- .expect("observations");
- assert_eq!(observations.len(), 1);
- assert_eq!(
- observations[0].observation_type,
- RadrootsTransportObservationType::NostrSubscription
- );
-
- let closed_store = RadrootsEventStore::open_memory().await.expect("store");
- closed_store.pool().close().await;
- let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event {
- relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: signed.raw_json,
- observed_at_ms: 1_210,
- }]);
- let receipt =
- fetch_and_ingest_relay_events(&adapter, &closed_store, post_relay_fetch_request(1_210, 10))
- .await
- .expect("fetch ingest");
-
- assert_eq!(receipt.inserted_count, 0);
- assert_eq!(receipt.malformed_count, 1);
- assert!(receipt.events[0].malformed);
- assert!(receipt.events[0].message.is_some());
-}
-
-#[tokio::test]
-async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() {
- let signed = signed_post("hello");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![
- RELAY_PRIMARY_WSS.to_owned(),
- RELAY_SECONDARY_WSS.to_owned(),
- RELAY_TERTIARY_WSS.to_owned(),
- ],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
- assert_eq!(publish_claim.state, RadrootsOutboxEventState::Publishing);
-
- let adapter = RadrootsMockRelayPublishAdapter::new()
- .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
- .with_outcome(
- RELAY_SECONDARY_WSS,
- RadrootsRelayOutcome::timeout("timeout: no OK"),
- )
- .with_outcome(
- RELAY_TERTIARY_WSS,
- RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"),
- );
- let first = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(first.publish.attempted_count, 3);
- assert_eq!(first.publish.accepted_count, 2);
- assert!(!first.publish.quorum_met);
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
- assert_eq!(event.accepted_quorum, 3);
-
- let statuses = outbox
- .relay_statuses(receipt.outbox_event_id)
- .await
- .expect("statuses");
- assert_eq!(
- statuses
- .iter()
- .find(|status| status.relay_url == RELAY_PRIMARY_WSS)
- .expect("primary")
- .status,
- RadrootsOutboxRelayStatus::Accepted
- );
- assert_eq!(
- statuses
- .iter()
- .find(|status| status.relay_url == RELAY_SECONDARY_WSS)
- .expect("secondary")
- .status,
- RadrootsOutboxRelayStatus::FailedRetryable
- );
- assert_eq!(
- statuses
- .iter()
- .find(|status| status.relay_url == RELAY_TERTIARY_WSS)
- .expect("tertiary")
- .status,
- RadrootsOutboxRelayStatus::Accepted
- );
-
- let retry_claim = outbox
- .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500)
- .await
- .expect("claim")
- .expect("retry claim");
- let retry_adapter = RadrootsMockRelayPublishAdapter::new()
- .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted());
- let second = publish_claimed_outbox_event(
- &outbox,
- &store,
- &retry_adapter,
- &retry_claim,
- RadrootsOutboxPublishPolicy::new(3_000),
- 2_600,
- )
- .await
- .expect("retry publish");
-
- assert_eq!(second.local_ingest.event_id, signed.id);
- assert_eq!(second.publish.attempted_count, 1);
- assert_eq!(retry_adapter.captured_raw_events().len(), 1);
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::Published);
- assert_eq!(event.accepted_quorum, 3);
- let operation = outbox
- .get_operation(receipt.operation_id)
- .await
- .expect("operation")
- .expect("operation");
- assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
-
- let observations = store
- .observations_for_event(signed.id.as_str())
- .await
- .expect("observations");
- assert_eq!(observations.len(), 3);
-}
-
-#[tokio::test]
-async fn outbox_publish_transport_failure_releases_retryable_claim() {
- let signed = signed_post("adapter transport failure");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
-
- let published = publish_claimed_outbox_event(
- &outbox,
- &store,
- &TransportFailurePublishAdapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(published.publish.attempted_count, 2);
- assert_eq!(published.publish.accepted_count, 0);
- assert_eq!(published.publish.retryable_count, 2);
- assert_eq!(published.publish.terminal_count, 0);
- assert!(!published.publish.quorum_met);
- assert!(
- published
- .publish
- .relays
- .iter()
- .all(|relay| relay.outcome.kind == RadrootsRelayOutcomeKind::ConnectionFailed)
- );
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
- assert!(event.claim_token.is_none());
- assert_eq!(event.next_attempt_after_ms, 2_500);
-
- let statuses = outbox
- .relay_statuses(receipt.outbox_event_id)
- .await
- .expect("statuses");
- assert_eq!(statuses.len(), 2);
- assert!(
- statuses
- .iter()
- .all(|status| status.status == RadrootsOutboxRelayStatus::FailedRetryable)
- );
- assert!(
- outbox
- .claim_next_ready_event("publisher", "publish-b", 4_000, 2_499)
- .await
- .expect("early claim")
- .is_none()
- );
- let retry_claim = outbox
- .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500)
- .await
- .expect("retry claim")
- .expect("retry claim");
- assert_eq!(retry_claim.outbox_event_id, receipt.outbox_event_id);
- assert_eq!(retry_claim.state, RadrootsOutboxEventState::Publishing);
-}
-
-#[tokio::test]
-async fn outbox_publish_marks_published_without_adapter_when_all_relays_already_accepted() {
- let signed = signed_post("already accepted");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
- outbox
- .mark_relay_accepted(
- publish_claim.outbox_event_id,
- publish_claim.claim_token.as_str(),
- RELAY_PRIMARY_WSS,
- 2_150,
- )
- .await
- .expect("primary accepted");
- outbox
- .mark_relay_accepted(
- publish_claim.outbox_event_id,
- publish_claim.claim_token.as_str(),
- RELAY_SECONDARY_WSS,
- 2_151,
- )
- .await
- .expect("secondary accepted");
-
- let adapter = RadrootsMockRelayPublishAdapter::new();
- let published = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(published.local_ingest.event_id, signed.id);
- assert_eq!(published.publish.event_id, signed.id);
- assert_eq!(published.publish.attempted_count, 0);
- assert_eq!(published.publish.accepted_count, 2);
- assert_eq!(published.publish.quorum, 2);
- assert!(published.publish.quorum_met);
- assert!(published.publish.relays.is_empty());
- assert!(adapter.captured_raw_events().is_empty());
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::Published);
- assert_eq!(event.accepted_quorum, 2);
- assert!(event.claim_token.is_none());
- let operation = outbox
- .get_operation(receipt.operation_id)
- .await
- .expect("operation")
- .expect("operation");
- assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
-}
-
-#[tokio::test]
-async fn outbox_publish_uses_persisted_accepted_count_for_explicit_quorum() {
- let signed = signed_post("explicit quorum already accepted");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![
- RELAY_PRIMARY_WSS.to_owned(),
- RELAY_SECONDARY_WSS.to_owned(),
- RELAY_TERTIARY_WSS.to_owned(),
- ],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
- outbox
- .mark_relay_accepted(
- publish_claim.outbox_event_id,
- publish_claim.claim_token.as_str(),
- RELAY_PRIMARY_WSS,
- 2_150,
- )
- .await
- .expect("primary accepted");
- outbox
- .mark_relay_accepted(
- publish_claim.outbox_event_id,
- publish_claim.claim_token.as_str(),
- RELAY_SECONDARY_WSS,
- 2_151,
- )
- .await
- .expect("secondary accepted");
-
- let adapter = RadrootsMockRelayPublishAdapter::new();
- let published = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500).with_accepted_quorum(2),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(published.publish.attempted_count, 0);
- assert_eq!(published.publish.accepted_count, 2);
- assert_eq!(published.publish.quorum, 2);
- assert!(published.publish.quorum_met);
- assert!(adapter.captured_raw_events().is_empty());
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::Published);
- assert_eq!(event.accepted_quorum, 2);
- let statuses = outbox
- .relay_statuses(receipt.outbox_event_id)
- .await
- .expect("statuses");
- assert_eq!(
- statuses
- .iter()
- .find(|status| status.relay_url == RELAY_TERTIARY_WSS)
- .expect("tertiary")
- .status,
- RadrootsOutboxRelayStatus::Pending
- );
-}
-
-#[tokio::test]
-async fn outbox_publish_marks_published_when_policy_quorum_is_met_with_failure_diagnostics() {
- let signed = signed_post("quorum");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![
- RELAY_PRIMARY_WSS.to_owned(),
- RELAY_SECONDARY_WSS.to_owned(),
- RELAY_TERTIARY_WSS.to_owned(),
- ],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
-
- let adapter = RadrootsMockRelayPublishAdapter::new()
- .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
- .with_outcome(
- RELAY_SECONDARY_WSS,
- RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"),
- )
- .with_outcome(
- RELAY_TERTIARY_WSS,
- RadrootsRelayOutcome::classify("restricted: group write denied"),
- );
- let published = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500).with_accepted_quorum(2),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(published.publish.quorum, 2);
- assert_eq!(published.publish.accepted_count, 2);
- assert_eq!(published.publish.terminal_count, 1);
- assert!(published.publish.quorum_met);
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::Published);
- assert_eq!(event.accepted_quorum, 2);
- assert!(event.claim_token.is_none());
- let operation = outbox
- .get_operation(receipt.operation_id)
- .await
- .expect("operation")
- .expect("operation");
- assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
-
- let statuses = outbox
- .relay_statuses(receipt.outbox_event_id)
- .await
- .expect("statuses");
- assert_eq!(
- statuses
- .iter()
- .find(|status| status.relay_url == RELAY_TERTIARY_WSS)
- .expect("tertiary")
- .status,
- RadrootsOutboxRelayStatus::FailedTerminal
- );
- assert!(
- outbox
- .claim_next_ready_event("publisher", "publish-b", 4_000, 2_300)
- .await
- .expect("claim")
- .is_none()
- );
-
- let observations = store
- .observations_for_event(signed.id.as_str())
- .await
- .expect("observations");
- assert_eq!(observations.len(), 2);
-}
-
-#[tokio::test]
-async fn outbox_publish_republishes_accepted_relays_when_policy_requests_it() {
- let signed = signed_post("republish accepted");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags.clone(),
- signed.content.clone(),
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
- outbox
- .mark_relay_accepted(
- publish_claim.outbox_event_id,
- publish_claim.claim_token.as_str(),
- RELAY_PRIMARY_WSS,
- 2_150,
- )
- .await
- .expect("primary accepted");
-
- let adapter = RadrootsMockRelayPublishAdapter::new()
- .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
- .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted());
- let published = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500)
- .republish_accepted_relays(true)
- .relay_url_policy(RadrootsRelayUrlPolicy::Public),
- 2_200,
- )
- .await
- .expect("publish");
-
- assert_eq!(published.local_ingest.event_id, signed.id);
- assert_eq!(published.publish.attempted_count, 2);
- assert_eq!(published.publish.accepted_count, 2);
- assert_eq!(published.publish.quorum, 1);
- assert!(published.publish.quorum_met);
- assert_eq!(adapter.captured_raw_events().len(), 1);
-
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.state, RadrootsOutboxEventState::Published);
- let statuses = outbox
- .relay_statuses(receipt.outbox_event_id)
- .await
- .expect("statuses");
- assert!(
- statuses
- .iter()
- .all(|status| status.status == RadrootsOutboxRelayStatus::Accepted)
- );
-}
-
-#[tokio::test]
-async fn outbox_publish_requires_claimed_signed_event() {
- let signed = signed_post("missing signature");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags,
- signed.content,
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![RELAY_PRIMARY_WSS.to_owned()],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- let adapter = RadrootsMockRelayPublishAdapter::new();
-
- let error = publish_claimed_outbox_event(
- &outbox,
- &store,
- &adapter,
- &claimed,
- RadrootsOutboxPublishPolicy::new(2_500),
- 1_100,
- )
- .await
- .expect_err("missing signed event");
-
- assert!(matches!(
- error,
- RadrootsRelayTransportError::MissingSignedOutboxEvent(event_id)
- if event_id == receipt.outbox_event_id
- ));
- assert!(adapter.captured_raw_events().is_empty());
-}
-
-#[tokio::test]
-async fn outbox_publish_propagates_non_transport_adapter_errors_after_target_filtering() {
- let signed = signed_post("adapter non transport failure");
- let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let draft = RadrootsFrozenEventDraft::new(
- "radroots.social.post.v1",
- KIND_POST,
- signed.created_at,
- signed.tags,
- signed.content,
- signed.pubkey.as_str(),
- )
- .expect("draft");
- let receipt = outbox
- .enqueue_operation(RadrootsOutboxOperationInput::new(
- "publish_post",
- draft,
- vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
- 1_000,
- ))
- .await
- .expect("enqueue");
- let claimed = outbox
- .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
- .await
- .expect("claim")
- .expect("claim");
- complete_claimed_signing(&outbox, &claimed, 1_100).await;
- outbox.recover_expired_claims(2_001).await.expect("recover");
- let mut publish_claim = outbox
- .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
- .await
- .expect("claim")
- .expect("publish claim");
- publish_claim.target_relays = vec![RELAY_PRIMARY_WSS.to_owned()];
-
- let error = publish_claimed_outbox_event(
- &outbox,
- &store,
- &NostrJsonFailurePublishAdapter,
- &publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500),
- 2_200,
- )
- .await
- .expect_err("adapter error");
-
- assert!(matches!(
- error,
- RadrootsRelayTransportError::NostrEventJson(_)
- ));
- let event = outbox
- .get_event(receipt.outbox_event_id)
- .await
- .expect("event")
- .expect("event");
- assert_eq!(event.accepted_quorum, 1);
-}
-
-#[tokio::test]
-async fn smoke_relay_fetch_processes_one_thousand_event_receipts() {
- let store = RadrootsEventStore::open_memory().await.expect("store");
- let mut items = Vec::new();
- for index in 0..1_000 {
- let signed = signed_post(format!("fetch-smoke-{index}").as_str());
- let relay_url = match index % 3 {
- 0 => RELAY_PRIMARY_WSS,
- 1 => RELAY_SECONDARY_WSS,
- _ => RELAY_TERTIARY_WSS,
- };
- items.push(RadrootsRelayFetchItem::Event {
- relay_url: relay_url.to_owned(),
- raw_json: signed.raw_json,
- observed_at_ms: 10_000 + index,
- });
- }
- let adapter = RadrootsMockRelayFetchAdapter::new(items);
- let receipt =
- fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(10_000, 1_000))
- .await
- .expect("fetch");
-
- assert_eq!(receipt.inserted_count, 1_000);
- assert_eq!(receipt.duplicate_count, 0);
- assert_eq!(receipt.malformed_count, 0);
- assert_eq!(receipt.unsupported_count, 0);
- assert_eq!(receipt.events.len(), 1_000);
- assert!(receipt.events.iter().all(|event| event.projection_eligible));
- let replay = store
- .events_since_cursor("fetch-smoke", 1_000)
- .await
- .expect("replay");
- assert_eq!(replay.len(), 1_000);
-}
diff --git a/crates/transport_nostr/Cargo.toml b/crates/transport_nostr/Cargo.toml
@@ -0,0 +1,62 @@
+[package]
+name = "radroots_transport_nostr"
+publish = false
+version = "0.1.0-alpha.2"
+edition.workspace = true
+authors = ["Tyson Lupul <tyson@radroots.org>"]
+rust-version.workspace = true
+license.workspace = true
+description = "Nostr transport layer for Radroots"
+repository.workspace = true
+homepage.workspace = true
+readme = "README"
+
+[features]
+default = ["std", "client", "storage", "runtime-tokio"]
+std = []
+client = [
+ "dep:radroots_nostr",
+ "radroots_nostr/std",
+ "radroots_nostr/client",
+ "radroots_nostr/events",
+]
+storage = ["dep:radroots_event_store", "dep:radroots_outbox", "client"]
+runtime-tokio = [
+ "dep:tokio",
+ "storage",
+ "radroots_event_store/runtime-tokio",
+ "radroots_outbox/runtime-tokio",
+]
+
+[dependencies]
+radroots_events = { workspace = true, default-features = false, features = [
+ "std",
+ "serde",
+] }
+radroots_event_store = { workspace = true, optional = true, default-features = false, features = [
+ "sqlite",
+ "runtime-tokio",
+] }
+radroots_nostr = { workspace = true, optional = true, default-features = false, features = [
+ "std",
+ "client",
+ "events",
+] }
+radroots_outbox = { workspace = true, optional = true, default-features = false, features = [
+ "sqlite",
+ "runtime-tokio",
+] }
+radroots_transport = { workspace = true, default-features = false }
+futures = { workspace = true }
+nostr = { workspace = true }
+serde = { workspace = true, features = ["derive", "std"] }
+serde_json = { workspace = true, features = ["std"] }
+thiserror = { workspace = true }
+tokio = { workspace = true, optional = true, features = ["rt"] }
+url = { workspace = true }
+
+[dev-dependencies]
+tokio = { workspace = true, features = ["macros", "rt"] }
+
+[lints.rust]
+unexpected_cfgs = { level = "warn", check-cfg = ['cfg(coverage_nightly)'] }
diff --git a/crates/transport_nostr/README b/crates/transport_nostr/README
@@ -0,0 +1,3 @@
+# radroots_transport_nostr
+
+Deterministic Nostr relay transport substrate for exact signed-event publish, fetch ingest, and outbox delivery target coordination.
diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs
@@ -0,0 +1,73 @@
+#![forbid(unsafe_code)]
+
+use thiserror::Error;
+
+#[derive(Debug, Error)]
+pub enum RadrootsRelayTransportError {
+ #[error("Relay URL parse failed for `{url}`: {reason}")]
+ RelayUrlParse { url: String, reason: String },
+
+ #[error("Relay URL `{url}` uses ws outside localhost relay policy")]
+ WsRequiresLocalhostPolicy { url: String },
+
+ #[error("Relay URL `{url}` has unsupported scheme `{scheme}`")]
+ UnsupportedRelayScheme { url: String, scheme: String },
+
+ #[error("Relay URL `{url}` must include a host")]
+ EmptyRelayHost { url: String },
+
+ #[error("Relay URL `{url}` must not include userinfo")]
+ RelayUrlUserinfo { url: String },
+
+ #[error("Relay URL `{url}` must not include query or fragment")]
+ RelayUrlQueryOrFragment { url: String },
+
+ #[error("Relay URL `{url}` targets forbidden destination: {reason}")]
+ RelayUrlForbiddenDestination { url: String, reason: String },
+
+ #[error("Relay URL `{url}` resolved to forbidden address `{address}`: {reason}")]
+ RelayUrlResolvedForbiddenDestination {
+ url: String,
+ address: String,
+ reason: String,
+ },
+
+ #[error("Relay target set must not be empty")]
+ EmptyTargetSet,
+
+ #[error("Relay fetch filters must not be empty")]
+ EmptyFetchFilters,
+
+ #[error("Relay fetch {field} must be greater than zero")]
+ InvalidFetchLimit { field: &'static str },
+
+ #[error("Transport contract error: {0}")]
+ TransportContract(String),
+
+ #[error("JSON error: {0}")]
+ Json(#[from] serde_json::Error),
+
+ #[error("Nostr event JSON error: {0}")]
+ NostrEventJson(String),
+
+ #[cfg(feature = "storage")]
+ #[error("Event store error: {0}")]
+ EventStore(#[from] radroots_event_store::RadrootsEventStoreError),
+
+ #[cfg(feature = "storage")]
+ #[error("Outbox error: {0}")]
+ Outbox(#[from] radroots_outbox::RadrootsOutboxError),
+
+ #[cfg(feature = "storage")]
+ #[error("Outbox claim {0} does not contain a signed event")]
+ MissingSignedOutboxEvent(i64),
+
+ #[error("Relay transport error: {0}")]
+ Transport(String),
+}
+
+impl From<radroots_transport::RadrootsTransportError> for RadrootsRelayTransportError {
+ fn from(value: radroots_transport::RadrootsTransportError) -> Self {
+ Self::TransportContract(value.to_string())
+ }
+}
diff --git a/crates/relay_transport/src/fetch.rs b/crates/transport_nostr/src/fetch.rs
diff --git a/crates/relay_transport/src/lib.rs b/crates/transport_nostr/src/lib.rs
diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs
@@ -0,0 +1,344 @@
+#![forbid(unsafe_code)]
+
+use crate::{
+ RadrootsRelayOutcome, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt,
+ RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTargetSet,
+ RadrootsRelayTransportError, RadrootsRelayUrlPolicy, publish_signed_event,
+};
+use radroots_event_store::{
+ RadrootsEventIngest, RadrootsEventStore, RadrootsTransportObservation,
+ RadrootsTransportObservationType,
+};
+use radroots_events::RadrootsNostrEvent;
+use radroots_events::draft::RadrootsSignedNostrEvent;
+use radroots_outbox::{
+ RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryTargetRecord,
+ RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxEventStoreIngestReceipt,
+};
+use radroots_transport::{RadrootsTransportKind, RadrootsTransportSatisfactionPolicy};
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsOutboxPublishPolicy {
+ pub next_attempt_after_ms: i64,
+ pub republish_accepted_relays: bool,
+ pub relay_url_policy: RadrootsRelayUrlPolicy,
+}
+
+impl RadrootsOutboxPublishPolicy {
+ pub fn new(next_attempt_after_ms: i64) -> Self {
+ Self {
+ next_attempt_after_ms,
+ republish_accepted_relays: false,
+ relay_url_policy: RadrootsRelayUrlPolicy::Public,
+ }
+ }
+
+ pub fn republish_accepted_relays(mut self, enabled: bool) -> Self {
+ self.republish_accepted_relays = enabled;
+ self
+ }
+
+ pub fn relay_url_policy(mut self, policy: RadrootsRelayUrlPolicy) -> Self {
+ self.relay_url_policy = policy;
+ self
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsOutboxPublishReceipt {
+ pub local_ingest: RadrootsOutboxEventStoreIngestReceipt,
+ pub publish: RadrootsRelayPublishReceipt,
+}
+
+pub async fn publish_claimed_outbox_event<A>(
+ outbox: &RadrootsOutbox,
+ event_store: &RadrootsEventStore,
+ adapter: &A,
+ claimed: &RadrootsOutboxClaimedEvent,
+ policy: RadrootsOutboxPublishPolicy,
+ now_ms: i64,
+) -> Result<RadrootsOutboxPublishReceipt, RadrootsRelayTransportError>
+where
+ A: RadrootsRelayPublishAdapter,
+{
+ let signed_event = claimed.signed_event.clone().ok_or(
+ RadrootsRelayTransportError::MissingSignedOutboxEvent(claimed.outbox_event_id),
+ )?;
+ let local_ingest = outbox
+ .ingest_signed_event_local(
+ event_store,
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ now_ms,
+ )
+ .await?;
+ let publishable = publishable_relays(outbox, claimed, policy.republish_accepted_relays).await?;
+ if publishable.relays.is_empty() {
+ outbox
+ .complete_publish_attempt(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ "relay publish incomplete",
+ "relay publish terminal",
+ policy.next_attempt_after_ms,
+ now_ms,
+ )
+ .await?;
+ let publish = RadrootsRelayPublishReceipt {
+ event_id: signed_event.id,
+ attempted_count: 0,
+ accepted_count: publishable.accepted_count,
+ retryable_count: 0,
+ terminal_count: 0,
+ quorum: publishable.satisfaction_required_count,
+ quorum_met: publishable.accepted_count >= publishable.satisfaction_required_count,
+ relays: Vec::new(),
+ };
+ return Ok(RadrootsOutboxPublishReceipt {
+ local_ingest,
+ publish,
+ });
+ }
+ let targets = RadrootsRelayTargetSet::new(
+ publishable
+ .relays
+ .iter()
+ .map(|target| target.relay_url.as_str()),
+ policy.relay_url_policy,
+ )?;
+ let target_strings = targets.relay_strings();
+ let request = RadrootsRelayPublishRequest::new(signed_event.clone(), targets, now_ms)
+ .with_satisfaction_policy(satisfaction_policy_for_required_accept_count(
+ publishable.required_accept_count,
+ publishable.relays.len(),
+ )?);
+ let publish = match publish_signed_event(adapter, request).await {
+ Ok(receipt) => receipt,
+ Err(RadrootsRelayTransportError::Transport(message)) => adapter_transport_failure_receipt(
+ signed_event.id.clone(),
+ target_strings,
+ publishable.required_accept_count,
+ message,
+ ),
+ Err(error) => return Err(error),
+ };
+
+ for relay in &publish.relays {
+ if let Some(target) = publishable.target_for_relay(relay.relay_url.as_str()) {
+ if relay.outcome.counts_toward_quorum() {
+ outbox
+ .mark_delivery_target_accepted(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ target.delivery_target_id,
+ now_ms,
+ )
+ .await?;
+ ingest_publish_observation(
+ event_store,
+ &signed_event,
+ relay.relay_url.as_str(),
+ relay.outcome.message.as_deref(),
+ now_ms,
+ )
+ .await?;
+ } else if relay.outcome.is_retryable() {
+ outbox
+ .mark_delivery_target_failed_retryable(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ target.delivery_target_id,
+ relay
+ .outcome
+ .message
+ .as_deref()
+ .unwrap_or("relay publish retryable"),
+ now_ms,
+ )
+ .await?;
+ } else {
+ outbox
+ .mark_delivery_target_failed_terminal(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ target.delivery_target_id,
+ relay
+ .outcome
+ .message
+ .as_deref()
+ .unwrap_or("relay publish terminal"),
+ now_ms,
+ )
+ .await?;
+ }
+ }
+ }
+
+ outbox
+ .complete_publish_attempt(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ "relay publish incomplete",
+ "relay publish terminal",
+ policy.next_attempt_after_ms,
+ now_ms,
+ )
+ .await?;
+
+ Ok(RadrootsOutboxPublishReceipt {
+ local_ingest,
+ publish,
+ })
+}
+
+fn adapter_transport_failure_receipt(
+ event_id: String,
+ relay_urls: Vec<String>,
+ quorum: usize,
+ message: String,
+) -> RadrootsRelayPublishReceipt {
+ let relays = relay_urls
+ .into_iter()
+ .map(|relay_url| {
+ RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url,
+ RadrootsRelayOutcome::connection_failed(message.clone()),
+ )
+ })
+ .collect::<Vec<_>>();
+ RadrootsRelayPublishReceipt {
+ event_id,
+ attempted_count: relays.len(),
+ accepted_count: 0,
+ retryable_count: relays.len(),
+ terminal_count: 0,
+ quorum,
+ quorum_met: false,
+ relays,
+ }
+}
+
+struct PublishableRelays {
+ relays: Vec<PublishableRelay>,
+ accepted_count: usize,
+ satisfaction_required_count: usize,
+ required_accept_count: usize,
+}
+
+impl PublishableRelays {
+ fn target_for_relay(&self, relay_url: &str) -> Option<&PublishableRelay> {
+ self.relays
+ .iter()
+ .find(|target| target.relay_url == relay_url)
+ }
+}
+
+struct PublishableRelay {
+ delivery_target_id: i64,
+ relay_url: String,
+}
+
+async fn publishable_relays(
+ outbox: &RadrootsOutbox,
+ claimed: &RadrootsOutboxClaimedEvent,
+ republish_accepted_relays: bool,
+) -> Result<PublishableRelays, RadrootsRelayTransportError> {
+ let targets = outbox.delivery_targets(claimed.outbox_event_id).await?;
+ let plans = outbox.delivery_plans(claimed.outbox_event_id).await?;
+ let satisfaction_required_count = plans
+ .iter()
+ .map(|plan| plan.required_success_count as usize)
+ .sum::<usize>();
+ let required_accept_count = plans
+ .iter()
+ .map(|plan| {
+ let satisfied_count = targets
+ .iter()
+ .filter(|target| target.delivery_plan_id == plan.delivery_plan_id)
+ .filter(|target| target.status.counts_as_satisfied())
+ .count();
+ (plan.required_success_count as usize).saturating_sub(satisfied_count)
+ })
+ .sum::<usize>();
+ let mut relays = Vec::new();
+ let mut accepted_count = 0usize;
+ for target in targets {
+ if !is_nostr_target(&target) {
+ continue;
+ }
+ if target.status.counts_as_satisfied() {
+ accepted_count += 1;
+ }
+ if required_accept_count > 0
+ && (target.status.is_ready_for_attempt()
+ || (republish_accepted_relays
+ && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted))
+ {
+ relays.push(PublishableRelay {
+ delivery_target_id: target.delivery_target_id,
+ relay_url: target.endpoint_uri.as_str().to_owned(),
+ });
+ }
+ }
+ Ok(PublishableRelays {
+ relays,
+ accepted_count,
+ satisfaction_required_count,
+ required_accept_count,
+ })
+}
+
+fn is_nostr_target(target: &RadrootsOutboxDeliveryTargetRecord) -> bool {
+ target.transport_kind == RadrootsTransportKind::Nostr
+}
+
+fn satisfaction_policy_for_required_accept_count(
+ required_accept_count: usize,
+ target_count: usize,
+) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> {
+ if required_accept_count >= target_count {
+ return Ok(RadrootsTransportSatisfactionPolicy::AllTargets);
+ }
+ let count = u16::try_from(required_accept_count).map_err(|_| {
+ RadrootsRelayTransportError::Transport(
+ "required Nostr relay acceptance count exceeds supported transport policy range"
+ .to_owned(),
+ )
+ })?;
+ Ok(RadrootsTransportSatisfactionPolicy::AtLeast(count))
+}
+
+async fn ingest_publish_observation(
+ event_store: &RadrootsEventStore,
+ signed_event: &RadrootsSignedNostrEvent,
+ relay_url: &str,
+ message: Option<&str>,
+ observed_at_ms: i64,
+) -> Result<(), RadrootsRelayTransportError> {
+ let mut observation = RadrootsTransportObservation::new(
+ RadrootsTransportKind::Nostr,
+ relay_url,
+ RadrootsTransportObservationType::NostrPublishAck,
+ observed_at_ms,
+ )?;
+ if let Some(message) = message {
+ observation = observation.with_redacted_message(message);
+ }
+ let ingest = RadrootsEventIngest::new(event_from_signed(signed_event), observed_at_ms)
+ .with_raw_json(signed_event.raw_json.clone())
+ .with_observation(observation);
+ event_store.ingest_event(ingest).await?;
+ Ok(())
+}
+
+fn event_from_signed(signed_event: &RadrootsSignedNostrEvent) -> RadrootsNostrEvent {
+ RadrootsNostrEvent {
+ id: signed_event.id.clone(),
+ author: signed_event.pubkey.clone(),
+ created_at: signed_event.created_at,
+ kind: signed_event.kind,
+ tags: signed_event.tags.clone(),
+ content: signed_event.content.clone(),
+ sig: signed_event.sig.clone(),
+ }
+}
diff --git a/crates/transport_nostr/src/outcome.rs b/crates/transport_nostr/src/outcome.rs
@@ -0,0 +1,194 @@
+#![forbid(unsafe_code)]
+
+use radroots_transport::{RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome};
+use serde::{Deserialize, Serialize};
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
+pub enum RadrootsRelayOutcomeKind {
+ Accepted,
+ DuplicateAccepted,
+ Blocked,
+ RateLimited,
+ Invalid,
+ PowRequired,
+ Restricted,
+ AuthRequired,
+ Muted,
+ Unsupported,
+ PaymentRequired,
+ Error,
+ Timeout,
+ ConnectionFailed,
+ RelayUrlRejected,
+ SkippedAlreadyAccepted,
+ Unknown,
+}
+
+impl RadrootsRelayOutcomeKind {
+ pub fn as_str(self) -> &'static str {
+ match self {
+ Self::Accepted => "accepted",
+ Self::DuplicateAccepted => "duplicate_accepted",
+ Self::Blocked => "blocked",
+ Self::RateLimited => "rate_limited",
+ Self::Invalid => "invalid",
+ Self::PowRequired => "pow_required",
+ Self::Restricted => "restricted",
+ Self::AuthRequired => "auth_required",
+ Self::Muted => "muted",
+ Self::Unsupported => "unsupported",
+ Self::PaymentRequired => "payment_required",
+ Self::Error => "error",
+ Self::Timeout => "timeout",
+ Self::ConnectionFailed => "connection_failed",
+ Self::RelayUrlRejected => "relay_url_rejected",
+ Self::SkippedAlreadyAccepted => "skipped_already_accepted",
+ Self::Unknown => "unknown",
+ }
+ }
+
+ pub fn counts_toward_quorum(self) -> bool {
+ matches!(
+ self,
+ Self::Accepted | Self::DuplicateAccepted | Self::SkippedAlreadyAccepted
+ )
+ }
+
+ pub fn is_retryable(self) -> bool {
+ matches!(
+ self,
+ Self::RateLimited
+ | Self::PowRequired
+ | Self::AuthRequired
+ | Self::Error
+ | Self::Timeout
+ | Self::ConnectionFailed
+ | Self::Unknown
+ )
+ }
+
+ pub fn is_terminal_failure(self) -> bool {
+ matches!(
+ self,
+ Self::Blocked
+ | Self::Invalid
+ | Self::Restricted
+ | Self::Muted
+ | Self::Unsupported
+ | Self::PaymentRequired
+ | Self::RelayUrlRejected
+ )
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
+pub struct RadrootsRelayOutcome {
+ pub kind: RadrootsRelayOutcomeKind,
+ pub message: Option<String>,
+}
+
+impl RadrootsRelayOutcome {
+ pub fn accepted() -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::Accepted,
+ message: None,
+ }
+ }
+
+ pub fn duplicate_accepted(message: impl Into<String>) -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::DuplicateAccepted,
+ message: Some(message.into()),
+ }
+ }
+
+ pub fn connection_failed(message: impl Into<String>) -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::ConnectionFailed,
+ message: Some(message.into()),
+ }
+ }
+
+ pub fn timeout(message: impl Into<String>) -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::Timeout,
+ message: Some(message.into()),
+ }
+ }
+
+ pub fn relay_url_rejected(message: impl Into<String>) -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::RelayUrlRejected,
+ message: Some(message.into()),
+ }
+ }
+
+ pub fn skipped_already_accepted(message: impl Into<String>) -> Self {
+ Self {
+ kind: RadrootsRelayOutcomeKind::SkippedAlreadyAccepted,
+ message: Some(message.into()),
+ }
+ }
+
+ pub fn classify(message: impl AsRef<str>) -> Self {
+ let message = message.as_ref().trim();
+ let lower = message.to_ascii_lowercase();
+ let kind = if lower.starts_with("duplicate:") {
+ RadrootsRelayOutcomeKind::DuplicateAccepted
+ } else if lower.starts_with("blocked:") {
+ RadrootsRelayOutcomeKind::Blocked
+ } else if lower.starts_with("rate-limited:") {
+ RadrootsRelayOutcomeKind::RateLimited
+ } else if lower.starts_with("invalid:") {
+ RadrootsRelayOutcomeKind::Invalid
+ } else if lower.starts_with("pow:") {
+ RadrootsRelayOutcomeKind::PowRequired
+ } else if lower.starts_with("restricted:") {
+ RadrootsRelayOutcomeKind::Restricted
+ } else if lower.starts_with("auth-required:") {
+ RadrootsRelayOutcomeKind::AuthRequired
+ } else if lower.starts_with("mute:") {
+ RadrootsRelayOutcomeKind::Muted
+ } else if lower.starts_with("unsupported:") {
+ RadrootsRelayOutcomeKind::Unsupported
+ } else if lower.starts_with("payment-required:") {
+ RadrootsRelayOutcomeKind::PaymentRequired
+ } else if lower.starts_with("error:") {
+ RadrootsRelayOutcomeKind::Error
+ } else if lower.starts_with("timeout:") {
+ RadrootsRelayOutcomeKind::Timeout
+ } else {
+ RadrootsRelayOutcomeKind::Unknown
+ };
+ Self {
+ kind,
+ message: Some(message.to_owned()),
+ }
+ }
+
+ pub fn counts_toward_quorum(&self) -> bool {
+ self.kind.counts_toward_quorum()
+ }
+
+ pub fn is_retryable(&self) -> bool {
+ self.kind.is_retryable()
+ }
+
+ pub fn is_terminal_failure(&self) -> bool {
+ self.kind.is_terminal_failure()
+ }
+
+ pub fn to_transport_outcome(&self) -> RadrootsTransportOutcome {
+ let status = if self.counts_toward_quorum() {
+ RadrootsTransportDeliveryTargetStatus::Accepted
+ } else if self.is_retryable() {
+ RadrootsTransportDeliveryTargetStatus::Failed
+ } else {
+ RadrootsTransportDeliveryTargetStatus::Rejected
+ };
+ let mut outcome = RadrootsTransportOutcome::new(status);
+ outcome.code = Some(self.kind.as_str().to_owned());
+ outcome.message = self.message.clone();
+ outcome
+ }
+}
diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs
@@ -0,0 +1,449 @@
+#![forbid(unsafe_code)]
+
+use crate::{RadrootsRelayOutcome, RadrootsRelayTargetSet, RadrootsRelayTransportError};
+#[cfg(feature = "client")]
+use core::time::Duration;
+use futures::future::BoxFuture;
+use radroots_events::draft::RadrootsSignedNostrEvent;
+use radroots_transport::RadrootsTransportSatisfactionPolicy;
+use serde::{Deserialize, Serialize};
+use std::collections::{BTreeMap, BTreeSet};
+use std::sync::{Arc, Mutex, PoisonError};
+
+#[cfg(feature = "client")]
+use crate::RadrootsRelayOutcomeKind;
+#[cfg(feature = "client")]
+use nostr::JsonUtil;
+#[cfg(feature = "client")]
+use radroots_nostr::prelude::{RadrootsNostrClient, RadrootsNostrEvent};
+
+#[cfg(feature = "client")]
+const RELAY_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
+
+#[derive(Clone, Debug, PartialEq, Eq)]
+pub struct RadrootsRelayPublishRequest {
+ pub signed_event: RadrootsSignedNostrEvent,
+ pub targets: RadrootsRelayTargetSet,
+ pub satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+ pub now_ms: i64,
+}
+
+impl RadrootsRelayPublishRequest {
+ pub fn new(
+ signed_event: RadrootsSignedNostrEvent,
+ targets: RadrootsRelayTargetSet,
+ now_ms: i64,
+ ) -> Self {
+ Self {
+ signed_event,
+ targets,
+ satisfaction_policy: RadrootsTransportSatisfactionPolicy::AllTargets,
+ now_ms,
+ }
+ }
+
+ pub fn with_satisfaction_policy(
+ mut self,
+ satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+ ) -> Self {
+ self.satisfaction_policy = satisfaction_policy;
+ self
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
+pub struct RadrootsRelayPublishRelayReceipt {
+ pub relay_url: String,
+ pub outcome: RadrootsRelayOutcome,
+ pub attempted: bool,
+}
+
+impl RadrootsRelayPublishRelayReceipt {
+ pub fn attempted(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self {
+ Self {
+ relay_url: relay_url.into(),
+ outcome,
+ attempted: true,
+ }
+ }
+
+ pub fn skipped(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self {
+ Self {
+ relay_url: relay_url.into(),
+ outcome,
+ attempted: false,
+ }
+ }
+}
+
+#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
+pub struct RadrootsRelayPublishReceipt {
+ pub event_id: String,
+ pub attempted_count: usize,
+ pub accepted_count: usize,
+ pub retryable_count: usize,
+ pub terminal_count: usize,
+ pub quorum: usize,
+ pub quorum_met: bool,
+ pub relays: Vec<RadrootsRelayPublishRelayReceipt>,
+}
+
+pub trait RadrootsRelayPublishAdapter: Send + Sync {
+ fn publish<'a>(
+ &'a self,
+ request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>;
+}
+
+pub async fn publish_signed_event<A>(
+ adapter: &A,
+ request: RadrootsRelayPublishRequest,
+) -> Result<RadrootsRelayPublishReceipt, RadrootsRelayTransportError>
+where
+ A: RadrootsRelayPublishAdapter,
+{
+ let event_id = request.signed_event.id.clone();
+ let quorum = request
+ .satisfaction_policy
+ .required_target_count(request.targets.len())?;
+ let relays = adapter.publish(request).await?;
+ let attempted_count = relays.iter().filter(|receipt| receipt.attempted).count();
+ let accepted_count = relays
+ .iter()
+ .filter(|receipt| receipt.outcome.counts_toward_quorum())
+ .count();
+ let retryable_count = relays
+ .iter()
+ .filter(|receipt| receipt.outcome.is_retryable())
+ .count();
+ let terminal_count = relays
+ .iter()
+ .filter(|receipt| receipt.outcome.is_terminal_failure())
+ .count();
+ Ok(RadrootsRelayPublishReceipt {
+ event_id,
+ attempted_count,
+ accepted_count,
+ retryable_count,
+ terminal_count,
+ quorum,
+ quorum_met: accepted_count >= quorum,
+ relays,
+ })
+}
+
+#[derive(Clone, Default)]
+pub struct RadrootsMockRelayPublishAdapter {
+ outcomes: BTreeMap<String, RadrootsRelayOutcome>,
+ captured_raw_events: Arc<Mutex<Vec<String>>>,
+}
+
+impl RadrootsMockRelayPublishAdapter {
+ pub fn new() -> Self {
+ Self::default()
+ }
+
+ pub fn with_outcome(
+ mut self,
+ relay_url: impl Into<String>,
+ outcome: RadrootsRelayOutcome,
+ ) -> Self {
+ self.outcomes.insert(relay_url.into(), outcome);
+ self
+ }
+
+ pub fn captured_raw_events(&self) -> Vec<String> {
+ self.captured_raw_events
+ .lock()
+ .expect("captured raw event lock")
+ .clone()
+ }
+}
+
+impl RadrootsRelayPublishAdapter for RadrootsMockRelayPublishAdapter {
+ fn publish<'a>(
+ &'a self,
+ request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
+ {
+ Box::pin(async move {
+ self.captured_raw_events
+ .lock()
+ .map_err(captured_raw_event_lock_error)?
+ .push(request.signed_event.raw_json.clone());
+ Ok(request
+ .targets
+ .relays()
+ .iter()
+ .map(|relay| {
+ let outcome = self
+ .outcomes
+ .get(relay.as_str())
+ .cloned()
+ .unwrap_or_else(RadrootsRelayOutcome::accepted);
+ RadrootsRelayPublishRelayReceipt::attempted(relay.as_str(), outcome)
+ })
+ .collect())
+ })
+ }
+}
+
+#[cfg_attr(coverage_nightly, coverage(off))]
+fn captured_raw_event_lock_error<T>(_error: PoisonError<T>) -> RadrootsRelayTransportError {
+ RadrootsRelayTransportError::Transport("captured raw event lock poisoned".to_owned())
+}
+
+#[cfg(feature = "client")]
+#[derive(Clone)]
+pub struct RadrootsNostrClientPublishAdapter {
+ client: RadrootsNostrClient,
+}
+
+#[cfg(feature = "client")]
+impl RadrootsNostrClientPublishAdapter {
+ #[cfg_attr(coverage_nightly, coverage(off))]
+ pub fn new(client: RadrootsNostrClient) -> Self {
+ Self { client }
+ }
+}
+
+#[cfg(feature = "client")]
+impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter {
+ #[cfg_attr(coverage_nightly, coverage(off))]
+ fn publish<'a>(
+ &'a self,
+ request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
+ {
+ Box::pin(async move {
+ let event = RadrootsNostrEvent::from_json(request.signed_event.raw_json.as_str())
+ .map_err(|error| RadrootsRelayTransportError::NostrEventJson(error.to_string()))?;
+ ensure_raw_event_matches_signed_event(&event, &request.signed_event)?;
+ let target_strings = request.targets.relay_strings();
+ for relay_url in &target_strings {
+ self.client
+ .add_write_relay(relay_url.as_str())
+ .await
+ .map_err(|error| RadrootsRelayTransportError::Transport(error.to_string()))?;
+ }
+ let connection_output = self.client.try_connect(RELAY_CONNECT_TIMEOUT).await;
+ let target_url_set = target_strings
+ .iter()
+ .map(|relay_url| relay_url.trim_end_matches('/').to_owned())
+ .collect::<BTreeSet<_>>();
+ let connected_strings = self
+ .client
+ .relays()
+ .await
+ .into_values()
+ .filter(|relay| relay.is_connected())
+ .map(|relay| relay.url().to_string())
+ .filter(|relay_url| target_url_set.contains(relay_url.trim_end_matches('/')))
+ .collect::<Vec<_>>();
+ let connection_failures = connection_output
+ .failed
+ .iter()
+ .map(|(relay, reason)| {
+ (
+ relay.to_string().trim_end_matches('/').to_owned(),
+ reason.clone(),
+ )
+ })
+ .collect::<BTreeMap<_, _>>();
+ if connected_strings.is_empty() {
+ return Ok(target_strings
+ .into_iter()
+ .map(|relay_url| {
+ let target_url = relay_url.trim_end_matches('/');
+ let reason = connection_failures
+ .get(target_url)
+ .cloned()
+ .unwrap_or_else(|| "relay did not connect".to_owned());
+ RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url,
+ RadrootsRelayOutcome::connection_failed(reason),
+ )
+ })
+ .collect());
+ }
+ let output = match self.client.send_event_to(connected_strings, &event).await {
+ Ok(output) => output,
+ Err(error) => {
+ let message = error.to_string();
+ return Ok(target_strings
+ .into_iter()
+ .map(|relay_url| {
+ RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url,
+ RadrootsRelayOutcome::connection_failed(message.clone()),
+ )
+ })
+ .collect());
+ }
+ };
+ let mut receipts = Vec::new();
+ for relay_url in &target_strings {
+ let target_url = relay_url.trim_end_matches('/');
+ let success = output
+ .success
+ .iter()
+ .any(|success_url| success_url.to_string().trim_end_matches('/') == target_url);
+ if success {
+ receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url,
+ RadrootsRelayOutcome {
+ kind: RadrootsRelayOutcomeKind::Accepted,
+ message: Some(
+ "nostr-relay-pool-success-ok-message-unavailable".to_owned(),
+ ),
+ },
+ ));
+ continue;
+ }
+ if let Some(reason) = connection_failures.get(target_url) {
+ receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url,
+ RadrootsRelayOutcome::connection_failed(reason.clone()),
+ ));
+ continue;
+ }
+ let failed = output.failed.iter().find_map(|(failed_url, message)| {
+ if failed_url.to_string().trim_end_matches('/') == target_url {
+ Some(message.clone())
+ } else {
+ None
+ }
+ });
+ let outcome = failed
+ .map(RadrootsRelayOutcome::classify)
+ .unwrap_or_else(|| {
+ RadrootsRelayOutcome::classify("error: relay output omitted target")
+ });
+ receipts.push(RadrootsRelayPublishRelayReceipt::attempted(
+ relay_url, outcome,
+ ));
+ }
+ Ok(receipts)
+ })
+ }
+}
+
+#[cfg(feature = "client")]
+fn ensure_raw_event_matches_signed_event(
+ event: &RadrootsNostrEvent,
+ signed_event: &RadrootsSignedNostrEvent,
+) -> Result<(), RadrootsRelayTransportError> {
+ let mismatches = [
+ ("id", event.id.to_hex(), signed_event.id.clone()),
+ ("pubkey", event.pubkey.to_hex(), signed_event.pubkey.clone()),
+ (
+ "created_at",
+ event.created_at.as_secs().to_string(),
+ signed_event.created_at.to_string(),
+ ),
+ (
+ "kind",
+ (event.kind.as_u16() as u32).to_string(),
+ signed_event.kind.to_string(),
+ ),
+ (
+ "content",
+ event.content.clone(),
+ signed_event.content.clone(),
+ ),
+ ("sig", event.sig.to_string(), signed_event.sig.clone()),
+ ];
+ for (field, raw, wrapped) in mismatches {
+ if raw != wrapped {
+ return Err(RadrootsRelayTransportError::NostrEventJson(format!(
+ "raw event JSON {field} does not match signed event {field}"
+ )));
+ }
+ }
+ let raw_tags = event
+ .tags
+ .iter()
+ .map(|tag| tag.as_slice().to_vec())
+ .collect::<Vec<_>>();
+ if raw_tags != signed_event.tags {
+ return Err(RadrootsRelayTransportError::NostrEventJson(
+ "raw event JSON tags do not match signed event tags".to_owned(),
+ ));
+ }
+ Ok(())
+}
+
+#[cfg(all(test, feature = "client"))]
+mod tests {
+ use super::{RadrootsNostrEvent, ensure_raw_event_matches_signed_event};
+ use nostr::JsonUtil;
+ use radroots_events::draft::{RadrootsFrozenEventDraft, RadrootsSignedNostrEvent};
+ use radroots_events::kinds::KIND_POST;
+ use radroots_nostr::prelude::{
+ RadrootsNostrKeys, RadrootsNostrSecretKey, radroots_nostr_sign_frozen_draft,
+ };
+
+ const FIXTURE_ALICE_SECRET_KEY_HEX: &str =
+ "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5";
+ const FIXTURE_ALICE_PUBLIC_KEY_HEX: &str =
+ "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+
+ fn signed_post(content: &str) -> (RadrootsNostrEvent, RadrootsSignedNostrEvent) {
+ let secret_key =
+ RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key");
+ let keys = RadrootsNostrKeys::new(secret_key);
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ 1_700_000_000,
+ vec![vec!["t".to_owned(), "soil".to_owned()]],
+ content,
+ FIXTURE_ALICE_PUBLIC_KEY_HEX,
+ )
+ .expect("draft");
+ let signed_event = radroots_nostr_sign_frozen_draft(&keys, &draft).expect("signed event");
+ let raw_event =
+ RadrootsNostrEvent::from_json(signed_event.raw_json.as_str()).expect("raw event");
+ (raw_event, signed_event)
+ }
+
+ fn assert_mismatch(raw_event: &RadrootsNostrEvent, signed_event: RadrootsSignedNostrEvent) {
+ assert!(ensure_raw_event_matches_signed_event(raw_event, &signed_event).is_err());
+ }
+
+ #[test]
+ fn raw_event_match_guard_accepts_exact_event_and_rejects_field_mismatches() {
+ let (raw_event, signed_event) = signed_post("matched");
+ ensure_raw_event_matches_signed_event(&raw_event, &signed_event).expect("matching event");
+
+ let mut mismatched = signed_event.clone();
+ mismatched.id = "00".repeat(32);
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event.clone();
+ mismatched.pubkey = "11".repeat(32);
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event.clone();
+ mismatched.created_at += 1;
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event.clone();
+ mismatched.kind += 1;
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event.clone();
+ mismatched.content.push_str(" changed");
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event.clone();
+ mismatched.sig = "22".repeat(64);
+ assert_mismatch(&raw_event, mismatched);
+
+ let mut mismatched = signed_event;
+ mismatched
+ .tags
+ .push(vec!["t".to_owned(), "compost".to_owned()]);
+ assert_mismatch(&raw_event, mismatched);
+ }
+}
diff --git a/crates/relay_transport/src/relay.rs b/crates/transport_nostr/src/relay.rs
diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs
@@ -0,0 +1,1801 @@
+use futures::future::BoxFuture;
+use nostr::JsonUtil;
+use radroots_event_store::{
+ RadrootsEventStore, RadrootsEventVerificationStatus, RadrootsTransportObservationType,
+};
+use radroots_events::draft::{RadrootsFrozenEventDraft, RadrootsSignedNostrEvent};
+use radroots_events::kinds::KIND_POST;
+use radroots_nostr::prelude::{
+ RadrootsNostrFilter, RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrSecretKey,
+ RadrootsNostrTimestamp, radroots_nostr_build_event, radroots_nostr_filter_tag,
+ radroots_nostr_sign_frozen_draft,
+};
+use radroots_outbox::{
+ RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryPlanInput,
+ RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxEventState, RadrootsOutboxOperationInput,
+ RadrootsOutboxOperationStatus,
+};
+use radroots_transport::{
+ RadrootsTransportKind, RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget,
+};
+use radroots_transport_nostr::{
+ RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsOutboxPublishPolicy,
+ RadrootsRelayFetchFilters, RadrootsRelayFetchItem, RadrootsRelayFetchMode,
+ RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchRequest, RadrootsRelayOutcome,
+ RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt,
+ RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError,
+ RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events,
+ fetch_relay_events_blocking, publish_claimed_outbox_event, publish_signed_event,
+};
+use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
+
+const FIXTURE_ALICE_SECRET_KEY_HEX: &str =
+ "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5";
+const FIXTURE_ALICE_PUBLIC_KEY_HEX: &str =
+ "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+const RELAY_PRIMARY_WSS: &str = "wss://relay.example.com";
+const RELAY_SECONDARY_WSS: &str = "wss://relay-2.example.com";
+const RELAY_TERTIARY_WSS: &str = "wss://relay-3.example.com";
+
+struct TransportFailurePublishAdapter;
+
+impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter {
+ fn publish<'a>(
+ &'a self,
+ _request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
+ {
+ Box::pin(async {
+ Err(RadrootsRelayTransportError::Transport(
+ "adapter boundary unavailable".to_owned(),
+ ))
+ })
+ }
+}
+
+struct NostrJsonFailurePublishAdapter;
+
+impl RadrootsRelayPublishAdapter for NostrJsonFailurePublishAdapter {
+ fn publish<'a>(
+ &'a self,
+ _request: RadrootsRelayPublishRequest,
+ ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
+ {
+ Box::pin(async {
+ Err(RadrootsRelayTransportError::NostrEventJson(
+ "adapter rejected raw event".to_owned(),
+ ))
+ })
+ }
+}
+
+fn fixture_keys() -> RadrootsNostrKeys {
+ let secret_key =
+ RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key");
+ RadrootsNostrKeys::new(secret_key)
+}
+
+fn signed_post(content: &str) -> RadrootsSignedNostrEvent {
+ signed_event_with_kind_and_hashtag(content, KIND_POST, "soil")
+}
+
+fn signed_event_with_kind_and_hashtag(
+ content: &str,
+ kind: u32,
+ hashtag: &str,
+) -> RadrootsSignedNostrEvent {
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ kind,
+ 1_700_000_000,
+ vec![vec!["t".to_owned(), hashtag.to_owned()]],
+ content,
+ FIXTURE_ALICE_PUBLIC_KEY_HEX,
+ )
+ .expect("draft");
+ radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event")
+}
+
+fn signed_raw_event_with_kind_and_hashtag(content: &str, kind: u32, hashtag: &str) -> nostr::Event {
+ radroots_nostr_build_event(
+ kind,
+ content,
+ vec![vec!["t".to_owned(), hashtag.to_owned()]],
+ )
+ .expect("event builder")
+ .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000))
+ .sign_with_keys(&fixture_keys())
+ .expect("signed event")
+}
+
+async fn complete_claimed_signing(
+ outbox: &RadrootsOutbox,
+ claimed: &RadrootsOutboxClaimedEvent,
+ now_ms: i64,
+) -> RadrootsSignedNostrEvent {
+ if let Some(signed_event) = claimed.signed_event.clone() {
+ return signed_event;
+ }
+ let signed_event =
+ radroots_nostr_sign_frozen_draft(&fixture_keys(), &claimed.draft).expect("signed event");
+ outbox
+ .complete_signing(
+ claimed.outbox_event_id,
+ claimed.claim_token.as_str(),
+ signed_event,
+ now_ms,
+ )
+ .await
+ .expect("complete signing")
+}
+
+fn nostr_target(relay_url: &str) -> RadrootsTransportTarget {
+ RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, relay_url).expect("nostr target")
+}
+
+fn outbox_operation_input<I, S>(
+ draft: RadrootsFrozenEventDraft,
+ relays: I,
+ satisfaction_policy: RadrootsTransportSatisfactionPolicy,
+) -> RadrootsOutboxOperationInput
+where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+{
+ let targets = relays
+ .into_iter()
+ .map(|relay| nostr_target(relay.as_ref()))
+ .collect::<Vec<_>>();
+ RadrootsOutboxOperationInput::new(
+ "publish_post",
+ draft,
+ RadrootsOutboxDeliveryPlanInput::new(
+ "transport.nostr.local",
+ 1,
+ satisfaction_policy,
+ targets,
+ ),
+ 1_000,
+ )
+}
+
+fn all_targets_outbox_operation_input<I, S>(
+ draft: RadrootsFrozenEventDraft,
+ relays: I,
+) -> RadrootsOutboxOperationInput
+where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+{
+ outbox_operation_input(
+ draft,
+ relays,
+ RadrootsTransportSatisfactionPolicy::AllTargets,
+ )
+}
+
+fn unsupported_raw_event() -> String {
+ let event = radroots_nostr_build_event(999, "unsupported", Vec::new())
+ .expect("event builder")
+ .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001))
+ .sign_with_keys(&fixture_keys())
+ .expect("signed unsupported event");
+ event.as_json()
+}
+
+fn post_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter {
+ radroots_nostr_filter_tag(
+ RadrootsNostrFilter::new()
+ .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
+ .limit(limit),
+ "t",
+ vec!["soil".to_owned()],
+ )
+ .expect("post relay fetch filter")
+}
+
+fn unsupported_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter {
+ RadrootsNostrFilter::new()
+ .kind(RadrootsNostrKind::Custom(999))
+ .limit(limit)
+}
+
+fn fixture_relay_fetch_request(
+ observed_at_ms: i64,
+ max_events: usize,
+) -> RadrootsRelayFetchRequest {
+ RadrootsRelayFetchRequest::fetch(
+ observed_at_ms,
+ max_events,
+ [
+ post_relay_fetch_filter(max_events),
+ unsupported_relay_fetch_filter(max_events),
+ ],
+ )
+ .expect("fixture relay fetch request")
+}
+
+fn post_relay_fetch_request(observed_at_ms: i64, max_events: usize) -> RadrootsRelayFetchRequest {
+ RadrootsRelayFetchRequest::fetch(
+ observed_at_ms,
+ max_events,
+ [post_relay_fetch_filter(max_events)],
+ )
+ .expect("post relay fetch request")
+}
+
+fn tampered_raw_event() -> String {
+ let signed = signed_post("trusted");
+ let mut value =
+ serde_json::from_str::<serde_json::Value>(signed.raw_json.as_str()).expect("raw json");
+ value["content"] = serde_json::Value::String("tampered".to_owned());
+ serde_json::to_string(&value).expect("tampered json")
+}
+
+#[test]
+fn relay_url_validation_and_target_normalization() {
+ let relay = RadrootsRelayUrl::parse("wss://Relay.Example.com", RadrootsRelayUrlPolicy::Public)
+ .expect("relay");
+ assert_eq!(relay.as_str(), RELAY_PRIMARY_WSS);
+ assert_eq!(relay.clone().into_string(), RELAY_PRIMARY_WSS);
+ let relay_path = RadrootsRelayUrl::parse(
+ "wss://Relay.Example.com/nostr",
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("relay path");
+ assert_eq!(relay_path.as_str(), "wss://relay.example.com/nostr");
+
+ assert!(
+ RadrootsRelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Public).is_err()
+ );
+ let local = RadrootsRelayUrl::parse("ws://localhost:7777", RadrootsRelayUrlPolicy::Localhost)
+ .expect("local relay");
+ assert_eq!(local.as_str(), "ws://localhost:7777");
+ let local_ipv4 =
+ RadrootsRelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Localhost)
+ .expect("local ipv4 relay");
+ assert_eq!(local_ipv4.as_str(), "ws://127.0.0.1:7777");
+ let local_ipv6 = RadrootsRelayUrl::parse("ws://[::1]:7777", RadrootsRelayUrlPolicy::Localhost)
+ .expect("local ipv6 relay");
+ assert_eq!(local_ipv6.as_str(), "ws://[::1]:7777");
+ assert!(
+ RadrootsRelayUrl::parse("ws://example.com", RadrootsRelayUrlPolicy::Localhost).is_err()
+ );
+ assert!(
+ RadrootsRelayUrl::parse("ws://192.168.1.10:7777", RadrootsRelayUrlPolicy::Localhost)
+ .is_err()
+ );
+ assert!(matches!(
+ RadrootsRelayUrl::parse("wss://127.0.0.1", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
+ ));
+ assert!(matches!(
+ RadrootsRelayUrl::parse("wss://10.1.2.3", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
+ ));
+ assert!(matches!(
+ RadrootsRelayUrl::parse("wss://[::1]", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
+ ));
+ assert!(matches!(
+ RadrootsRelayUrl::parse("wss://[fd00::1]", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
+ ));
+ for relay_url in [
+ "wss://0.0.0.0",
+ "wss://169.254.1.2",
+ "wss://224.0.0.1",
+ "wss://255.255.255.255",
+ "wss://100.64.0.1",
+ "wss://192.0.0.8",
+ "wss://198.18.0.1",
+ "wss://240.0.0.1",
+ "wss://[::]",
+ "wss://[ff02::1]",
+ "wss://[fe80::1]",
+ "wss://[2001:db8::1]",
+ "wss://[2001:1::1]",
+ "wss://[::ffff:192.168.1.10]",
+ ] {
+ assert!(matches!(
+ RadrootsRelayUrl::parse(relay_url, RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. })
+ ));
+ }
+ let public_relay =
+ RadrootsRelayUrl::parse("wss://relay.example.com", RadrootsRelayUrlPolicy::Public)
+ .expect("public relay");
+ public_relay
+ .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))])
+ .expect("public resolved ip");
+ assert!(matches!(
+ public_relay
+ .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10))]),
+ Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. })
+ ));
+ assert!(matches!(
+ public_relay.validate_public_resolved_ip_addrs([IpAddr::V6(
+ "::ffff:192.168.1.10"
+ .parse::<Ipv6Addr>()
+ .expect("mapped ipv6")
+ )]),
+ Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. })
+ ));
+ public_relay
+ .validate_public_resolved_ip_addrs([IpAddr::V6(
+ "2001:4860:4860::8888"
+ .parse::<Ipv6Addr>()
+ .expect("public ipv6"),
+ )])
+ .expect("public resolved ipv6");
+ public_relay
+ .validate_public_resolved_ip_addrs(Vec::<IpAddr>::new())
+ .expect("empty resolved set");
+
+ assert!(
+ RadrootsRelayUrl::parse("https://relay.example.com", RadrootsRelayUrlPolicy::Public)
+ .is_err()
+ );
+ assert!(
+ RadrootsRelayUrl::parse(
+ "wss://user@relay.example.com",
+ RadrootsRelayUrlPolicy::Public
+ )
+ .is_err()
+ );
+ assert!(matches!(
+ RadrootsRelayUrl::parse(
+ "wss://user:password@relay.example.com",
+ RadrootsRelayUrlPolicy::Public
+ ),
+ Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. })
+ ));
+ assert!(matches!(
+ RadrootsRelayUrl::parse(
+ "wss://:password@relay.example.com",
+ RadrootsRelayUrlPolicy::Public
+ ),
+ Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. })
+ ));
+ assert!(
+ RadrootsRelayUrl::parse(
+ "wss://relay.example.com:bad",
+ RadrootsRelayUrlPolicy::Public
+ )
+ .is_err()
+ );
+ assert!(RadrootsRelayUrl::parse("wss://", RadrootsRelayUrlPolicy::Public).is_err());
+ assert!(matches!(
+ RadrootsRelayUrl::parse("radroots:relay", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::EmptyRelayHost { .. })
+ ));
+ assert!(matches!(
+ RadrootsRelayUrl::parse("relay.example.com", RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::RelayUrlParse { .. })
+ ));
+ assert!(
+ RadrootsRelayUrl::parse(
+ "wss://relay.example.com?subscription=1",
+ RadrootsRelayUrlPolicy::Public
+ )
+ .is_err()
+ );
+ assert!(
+ RadrootsRelayUrl::parse(
+ "wss://relay.example.com#fragment",
+ RadrootsRelayUrlPolicy::Public
+ )
+ .is_err()
+ );
+
+ let targets = RadrootsRelayTargetSet::new(
+ vec![
+ RELAY_TERTIARY_WSS,
+ RELAY_PRIMARY_WSS,
+ RELAY_PRIMARY_WSS,
+ RELAY_SECONDARY_WSS,
+ ],
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("targets");
+ assert_eq!(
+ targets.relay_strings(),
+ vec![
+ RELAY_TERTIARY_WSS.to_owned(),
+ RELAY_PRIMARY_WSS.to_owned(),
+ RELAY_SECONDARY_WSS.to_owned()
+ ]
+ );
+
+ let from_urls = RadrootsRelayTargetSet::from_urls(vec![
+ relay_path.clone(),
+ relay_path.clone(),
+ RadrootsRelayUrl::parse(RELAY_SECONDARY_WSS, RadrootsRelayUrlPolicy::Public)
+ .expect("secondary"),
+ ])
+ .expect("from urls");
+ assert_eq!(from_urls.len(), 2);
+ assert!(!from_urls.is_empty());
+ assert_eq!(from_urls.relays()[0], relay_path);
+ assert_eq!(
+ from_urls.relays()[0].to_string(),
+ "wss://relay.example.com/nostr"
+ );
+ assert!(matches!(
+ RadrootsRelayTargetSet::new(Vec::<&str>::new(), RadrootsRelayUrlPolicy::Public),
+ Err(RadrootsRelayTransportError::EmptyTargetSet)
+ ));
+ assert!(matches!(
+ RadrootsRelayTargetSet::from_urls(Vec::new()),
+ Err(RadrootsRelayTransportError::EmptyTargetSet)
+ ));
+}
+
+#[test]
+fn outcome_prefix_classification_covers_required_kinds() {
+ let cases = [
+ ("blocked: policy", RadrootsRelayOutcomeKind::Blocked),
+ (
+ "rate-limited: slow down",
+ RadrootsRelayOutcomeKind::RateLimited,
+ ),
+ ("invalid: bad event", RadrootsRelayOutcomeKind::Invalid),
+ ("pow: difficulty 24", RadrootsRelayOutcomeKind::PowRequired),
+ (
+ "restricted: group write denied",
+ RadrootsRelayOutcomeKind::Restricted,
+ ),
+ (
+ "auth-required: challenge",
+ RadrootsRelayOutcomeKind::AuthRequired,
+ ),
+ ("mute: pubkey muted", RadrootsRelayOutcomeKind::Muted),
+ (
+ "unsupported: event kind",
+ RadrootsRelayOutcomeKind::Unsupported,
+ ),
+ (
+ "payment-required: paid relay",
+ RadrootsRelayOutcomeKind::PaymentRequired,
+ ),
+ (
+ "duplicate: already have it",
+ RadrootsRelayOutcomeKind::DuplicateAccepted,
+ ),
+ ("error: relay failed", RadrootsRelayOutcomeKind::Error),
+ ("timeout: no OK", RadrootsRelayOutcomeKind::Timeout),
+ ("strange relay text", RadrootsRelayOutcomeKind::Unknown),
+ ];
+
+ for (message, kind) in cases {
+ let outcome = RadrootsRelayOutcome::classify(message);
+ assert_eq!(outcome.kind, kind);
+ }
+
+ assert!(RadrootsRelayOutcome::classify("duplicate: already have it").counts_toward_quorum());
+ assert!(
+ RadrootsRelayOutcome::skipped_already_accepted("already accepted").counts_toward_quorum()
+ );
+ assert!(RadrootsRelayOutcome::classify("auth-required: challenge").is_retryable());
+ assert!(RadrootsRelayOutcome::classify("restricted: denied").is_terminal_failure());
+ 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()
+ .status,
+ radroots_transport::RadrootsTransportDeliveryTargetStatus::Accepted
+ );
+}
+
+#[tokio::test]
+async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() {
+ let signed = signed_post("hello");
+ let targets = RadrootsRelayTargetSet::new(
+ vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS],
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("targets");
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(
+ RELAY_SECONDARY_WSS,
+ RadrootsRelayOutcome::classify("duplicate: already have it"),
+ )
+ .with_outcome(
+ RELAY_TERTIARY_WSS,
+ RadrootsRelayOutcome::classify("auth-required: challenge"),
+ );
+
+ let receipt = publish_signed_event(
+ &adapter,
+ radroots_transport_nostr::RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_000)
+ .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::AtLeast(2)),
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(adapter.captured_raw_events(), vec![signed.raw_json]);
+ assert_eq!(receipt.attempted_count, 3);
+ assert_eq!(receipt.accepted_count, 2);
+ assert_eq!(receipt.retryable_count, 1);
+ assert!(receipt.quorum_met);
+ serde_json::to_string(&receipt).expect("receipt json");
+}
+
+#[tokio::test]
+async fn publish_receipts_track_terminal_skipped_and_adapter_errors() {
+ let signed = signed_post("terminal");
+ let targets = RadrootsRelayTargetSet::new(
+ vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
+ RadrootsRelayUrlPolicy::Public,
+ )
+ .expect("targets");
+ let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome(
+ RELAY_SECONDARY_WSS,
+ RadrootsRelayOutcome::classify("restricted: group write denied"),
+ );
+
+ let receipt = publish_signed_event(
+ &adapter,
+ RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_050)
+ .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::AllTargets),
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(receipt.event_id, signed.id);
+ assert_eq!(receipt.attempted_count, 2);
+ assert_eq!(receipt.accepted_count, 1);
+ assert_eq!(receipt.retryable_count, 0);
+ assert_eq!(receipt.terminal_count, 1);
+ assert_eq!(receipt.quorum, 2);
+ assert!(!receipt.quorum_met);
+
+ let skipped = RadrootsRelayPublishRelayReceipt::skipped(
+ RELAY_TERTIARY_WSS,
+ RadrootsRelayOutcome::timeout("timeout: no OK"),
+ );
+ assert_eq!(skipped.relay_url, RELAY_TERTIARY_WSS);
+ assert!(!skipped.attempted);
+ assert_eq!(skipped.outcome.kind, RadrootsRelayOutcomeKind::Timeout);
+
+ let error = publish_signed_event(
+ &TransportFailurePublishAdapter,
+ RadrootsRelayPublishRequest::new(
+ signed,
+ RadrootsRelayTargetSet::new(vec![RELAY_PRIMARY_WSS], RadrootsRelayUrlPolicy::Public)
+ .expect("targets"),
+ 1_060,
+ ),
+ )
+ .await
+ .expect_err("transport failure");
+ assert!(matches!(error, RadrootsRelayTransportError::Transport(_)));
+}
+
+#[test]
+fn fetch_requests_reject_empty_filter_sets() {
+ assert!(matches!(
+ RadrootsRelayFetchRequest::fetch(1_000, 10, Vec::<RadrootsNostrFilter>::new()),
+ Err(RadrootsRelayTransportError::EmptyFetchFilters)
+ ));
+ assert!(matches!(
+ RadrootsRelayFetchRequest::subscription(1_000, 10, Vec::<RadrootsNostrFilter>::new()),
+ Err(RadrootsRelayTransportError::EmptyFetchFilters)
+ ));
+}
+
+#[test]
+fn fetch_requests_reject_zero_limits_and_timeouts() {
+ let filter = post_relay_fetch_filter(1);
+ let filters = RadrootsRelayFetchFilters::new([filter.clone()]).expect("filters");
+ let as_ref_filters: &[RadrootsNostrFilter] = filters.as_ref();
+ assert_eq!(as_ref_filters.len(), 1);
+
+ assert!(matches!(
+ RadrootsRelayFetchRequest::fetch(1_000, 0, [filter.clone()]),
+ Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events"
+ ));
+ assert!(matches!(
+ RadrootsRelayFetchRequest::subscription(1_000, 0, [filter.clone()]),
+ Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events"
+ ));
+
+ let request =
+ RadrootsRelayFetchRequest::fetch(1_000, 1, [filter]).expect("valid fetch request");
+ assert!(matches!(
+ request.clone().with_timeout_ms(0),
+ Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "timeout_ms"
+ ));
+ assert!(matches!(
+ request.clone().with_raw_event_scan_limit(0),
+ Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_raw_events"
+ ));
+
+ let request = request
+ .with_timeout_ms(1)
+ .expect("minimum timeout")
+ .with_raw_event_scan_limit(1)
+ .expect("minimum raw scan limit");
+ assert_eq!(request.timeout_ms(), 1);
+ assert_eq!(request.max_raw_events(), 1);
+
+ let request = RadrootsRelayFetchRequest::subscription(1_005, 2, [post_relay_fetch_filter(2)])
+ .expect("subscription request")
+ .with_relay_urls([RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS])
+ .with_timeout_ms(25)
+ .expect("timeout")
+ .with_raw_event_scan_limit(3)
+ .expect("raw limit");
+ assert_eq!(request.mode(), RadrootsRelayFetchMode::Subscription);
+ assert_eq!(request.observed_at_ms(), 1_005);
+ assert_eq!(request.max_events(), 2);
+ assert_eq!(request.max_raw_events(), 3);
+ assert_eq!(
+ request.relay_urls(),
+ &[RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()]
+ );
+ assert_eq!(request.filters().len(), 1);
+ assert_eq!(request.timeout_ms(), 25);
+}
+
+#[test]
+fn fetch_blocking_facade_runs_mock_adapter() {
+ let signed = signed_post("blocking fetch");
+ let accepted_id = signed.id.clone();
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: signed.raw_json,
+ observed_at_ms: 1_090,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ ]);
+
+ let receipt = fetch_relay_events_blocking(&adapter, post_relay_fetch_request(1_090, 10))
+ .expect("blocking fetch");
+
+ assert_eq!(receipt.events.len(), 1);
+ assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id);
+ assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]);
+}
+
+#[tokio::test]
+async fn fetch_ingests_events_and_records_transport_observations() {
+ let signed = signed_post("hello");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: signed.raw_json.clone(),
+ observed_at_ms: 1_000,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: signed.raw_json.clone(),
+ observed_at_ms: 1_001,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ raw_json: unsupported_raw_event(),
+ observed_at_ms: 1_002,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ raw_json: tampered_raw_event(),
+ observed_at_ms: 1_003,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_TERTIARY_WSS.to_owned(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 1_004,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ RadrootsRelayFetchItem::Closed {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ message: "auth-required: challenge".to_owned(),
+ },
+ RadrootsRelayFetchItem::Closed {
+ relay_url: RELAY_TERTIARY_WSS.to_owned(),
+ message: "restricted: group write denied".to_owned(),
+ },
+ RadrootsRelayFetchItem::Notice {
+ relay_url: RELAY_TERTIARY_WSS.to_owned(),
+ message: "notice: test".to_owned(),
+ },
+ ]);
+
+ let receipt =
+ fetch_and_ingest_relay_events(&adapter, &store, fixture_relay_fetch_request(1_000, 10))
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 3);
+ assert_eq!(receipt.duplicate_count, 1);
+ assert_eq!(receipt.unsupported_count, 1);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.eose_count, 1);
+ assert_eq!(receipt.closed_count, 2);
+ assert_eq!(receipt.notice_count, 1);
+ assert_eq!(receipt.relay_outcomes.len(), 4);
+ assert_eq!(receipt.relay_outcomes[0].relay_url, RELAY_PRIMARY_WSS);
+ assert_eq!(
+ receipt.relay_outcomes[0].kind,
+ RadrootsRelayFetchOutcomeKind::Eose
+ );
+ assert!(receipt.relay_outcomes[0].relay_outcome.is_none());
+ assert_eq!(receipt.relay_outcomes[1].relay_url, RELAY_SECONDARY_WSS);
+ assert_eq!(
+ receipt.relay_outcomes[1]
+ .relay_outcome
+ .as_ref()
+ .expect("auth outcome")
+ .kind,
+ RadrootsRelayOutcomeKind::AuthRequired
+ );
+ assert_eq!(receipt.relay_outcomes[2].relay_url, RELAY_TERTIARY_WSS);
+ assert_eq!(
+ receipt.relay_outcomes[2]
+ .relay_outcome
+ .as_ref()
+ .expect("restricted outcome")
+ .kind,
+ RadrootsRelayOutcomeKind::Restricted
+ );
+ assert_eq!(
+ receipt.relay_outcomes[3].kind,
+ RadrootsRelayFetchOutcomeKind::Notice
+ );
+ assert!(receipt.relay_outcomes[3].relay_outcome.is_none());
+ assert_eq!(
+ receipt.events[0].verification_status.as_deref(),
+ Some(RadrootsEventVerificationStatus::Verified.as_str())
+ );
+ assert!(receipt.events[0].projection_eligible);
+ assert_eq!(
+ receipt.events[1].verification_status.as_deref(),
+ Some(RadrootsEventVerificationStatus::Verified.as_str())
+ );
+ assert!(!receipt.events[1].projection_eligible);
+ assert_eq!(
+ receipt.events[2].verification_status.as_deref(),
+ Some(RadrootsEventVerificationStatus::Verified.as_str())
+ );
+ assert!(!receipt.events[2].projection_eligible);
+ assert_eq!(
+ receipt.events[3].verification_status.as_deref(),
+ Some(RadrootsEventVerificationStatus::IdMismatch.as_str())
+ );
+ assert!(!receipt.events[3].projection_eligible);
+ assert_eq!(receipt.events[4].verification_status, None);
+ assert!(!receipt.events[4].projection_eligible);
+
+ let observations = store
+ .observations_for_event(signed.id.as_str())
+ .await
+ .expect("observations");
+ assert_eq!(observations.len(), 1);
+ assert_eq!(observations[0].transport_kind, RadrootsTransportKind::Nostr);
+ assert_eq!(observations[0].endpoint_uri.as_str(), RELAY_PRIMARY_WSS);
+ assert_eq!(
+ observations[0].observation_type,
+ RadrootsTransportObservationType::NostrFetch
+ );
+ assert_eq!(observations[0].observation_count, 2);
+}
+
+#[tokio::test]
+async fn fetch_rejects_out_of_filter_events_before_store_mutation() {
+ let accepted = signed_post("filter match");
+ let wrong_tag = signed_event_with_kind_and_hashtag("filter wrong tag", KIND_POST, "compost");
+ let wrong_kind = signed_raw_event_with_kind_and_hashtag("filter wrong kind", 999, "soil");
+ let wrong_kind_event_id = wrong_kind.id.to_hex();
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json.clone(),
+ observed_at_ms: 1_005,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: accepted.raw_json.clone(),
+ observed_at_ms: 1_006,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ raw_json: wrong_kind.as_json(),
+ observed_at_ms: 1_007,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ ]);
+ let filter = radroots_nostr_filter_tag(
+ RadrootsNostrFilter::new()
+ .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
+ .limit(10),
+ "t",
+ vec!["soil".to_owned()],
+ )
+ .expect("filter");
+
+ let receipt = fetch_and_ingest_relay_events(
+ &adapter,
+ &store,
+ RadrootsRelayFetchRequest::fetch(1_005, 10, [filter]).expect("fetch request"),
+ )
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 2);
+ assert_eq!(receipt.malformed_count, 0);
+ assert_eq!(receipt.unsupported_count, 0);
+ assert_eq!(receipt.events.len(), 3);
+ assert!(receipt.events[0].out_of_filter);
+ assert!(!receipt.events[1].out_of_filter);
+ assert!(receipt.events[2].out_of_filter);
+ assert!(
+ store
+ .get_event(accepted.id.as_str())
+ .await
+ .expect("accepted lookup")
+ .is_some()
+ );
+ assert!(
+ store
+ .get_event(wrong_tag.id.as_str())
+ .await
+ .expect("wrong tag lookup")
+ .is_none()
+ );
+ assert!(
+ store
+ .get_event(wrong_kind_event_id.as_str())
+ .await
+ .expect("wrong kind lookup")
+ .is_none()
+ );
+}
+
+#[tokio::test]
+async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_control_outcomes() {
+ let accepted = signed_post("accepted capped event");
+ let skipped = signed_post("skipped capped event");
+ let wrong_tag = signed_event_with_kind_and_hashtag("wrong capped tag", KIND_POST, "compost");
+ let accepted_id = accepted.id.clone();
+ let skipped_id = skipped.id.clone();
+ let wrong_tag_id = wrong_tag.id.clone();
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 1_099,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json,
+ observed_at_ms: 1_100,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: accepted.raw_json.clone(),
+ observed_at_ms: 1_101,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: skipped.raw_json,
+ observed_at_ms: 1_102,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ RadrootsRelayFetchItem::Closed {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ message: "auth-required: challenge".to_owned(),
+ },
+ RadrootsRelayFetchItem::Notice {
+ relay_url: RELAY_TERTIARY_WSS.to_owned(),
+ message: "notice: still visible".to_owned(),
+ },
+ ]);
+
+ let receipt =
+ fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_100, 1))
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 1);
+ assert_eq!(receipt.duplicate_count, 0);
+ assert_eq!(receipt.unsupported_count, 0);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 1);
+ assert_eq!(receipt.skipped_over_limit_count, 1);
+ assert_eq!(receipt.events.len(), 4);
+ assert!(receipt.events[0].malformed);
+ assert!(receipt.events[1].out_of_filter);
+ assert!(receipt.events[2].inserted);
+ assert!(receipt.events[3].skipped_over_limit);
+ assert_eq!(receipt.eose_count, 1);
+ assert_eq!(receipt.closed_count, 1);
+ assert_eq!(receipt.notice_count, 1);
+ assert_eq!(receipt.relay_outcomes.len(), 3);
+ assert_eq!(
+ receipt.relay_outcomes[0].kind,
+ RadrootsRelayFetchOutcomeKind::Eose
+ );
+ assert_eq!(
+ receipt.relay_outcomes[1]
+ .relay_outcome
+ .as_ref()
+ .expect("closed outcome")
+ .kind,
+ RadrootsRelayOutcomeKind::AuthRequired
+ );
+ assert_eq!(
+ receipt.relay_outcomes[2].kind,
+ RadrootsRelayFetchOutcomeKind::Notice
+ );
+ assert!(
+ store
+ .get_event(accepted_id.as_str())
+ .await
+ .expect("accepted lookup")
+ .is_some()
+ );
+ assert!(
+ store
+ .get_event(skipped_id.as_str())
+ .await
+ .expect("skipped lookup")
+ .is_none()
+ );
+ assert!(
+ store
+ .get_event(wrong_tag_id.as_str())
+ .await
+ .expect("wrong tag lookup")
+ .is_none()
+ );
+}
+
+#[tokio::test]
+async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() {
+ let accepted = signed_event_with_kind_and_hashtag("shared fetch accepted", KIND_POST, "soil");
+ let skipped = signed_event_with_kind_and_hashtag("shared fetch skipped", KIND_POST, "soil");
+ let wrong_tag =
+ signed_event_with_kind_and_hashtag("shared fetch wrong tag", KIND_POST, "compost");
+ let filter = radroots_nostr_filter_tag(
+ RadrootsNostrFilter::new()
+ .kind(RadrootsNostrKind::Custom(KIND_POST as u16))
+ .limit(10),
+ "t",
+ vec!["soil".to_owned()],
+ )
+ .expect("filter");
+ let accepted_id = accepted.id.clone();
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 2_100,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json,
+ observed_at_ms: 2_101,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: accepted.raw_json.clone(),
+ observed_at_ms: 2_102,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: skipped.raw_json,
+ observed_at_ms: 2_103,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ RadrootsRelayFetchItem::Closed {
+ relay_url: RELAY_SECONDARY_WSS.to_owned(),
+ message: "auth-required: challenge".to_owned(),
+ },
+ RadrootsRelayFetchItem::Notice {
+ relay_url: RELAY_TERTIARY_WSS.to_owned(),
+ message: "notice: still visible".to_owned(),
+ },
+ ]);
+
+ let receipt = fetch_relay_events(
+ &adapter,
+ RadrootsRelayFetchRequest::fetch(2_100, 1, [filter])
+ .expect("fetch request")
+ .with_relay_urls([RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS]),
+ )
+ .await
+ .expect("fetch events");
+
+ assert_eq!(
+ receipt.target_relays,
+ vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS]
+ );
+ assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]);
+ assert_eq!(receipt.failed_relays.len(), 1);
+ assert_eq!(receipt.failed_relays[0].relay_url, RELAY_SECONDARY_WSS);
+ assert_eq!(receipt.events.len(), 1);
+ assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 1);
+ assert_eq!(receipt.skipped_over_limit_count, 1);
+ assert_eq!(receipt.eose_count, 1);
+ assert_eq!(receipt.closed_count, 1);
+ assert_eq!(receipt.notice_count, 1);
+ assert_eq!(receipt.event_receipts.len(), 4);
+ assert!(receipt.event_receipts[0].malformed);
+ assert!(receipt.event_receipts[1].out_of_filter);
+ assert!(!receipt.event_receipts[2].malformed);
+ assert!(receipt.event_receipts[3].skipped_over_limit);
+}
+
+#[tokio::test]
+async fn fetch_raw_scan_limit_bounds_noisy_adapter_output() {
+ let accepted = signed_post("raw scan accepted event");
+ let wrong_tag = signed_event_with_kind_and_hashtag("raw scan wrong tag", KIND_POST, "compost");
+ let accepted_id = accepted.id.clone();
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 1_130,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json,
+ observed_at_ms: 1_131,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: accepted.raw_json,
+ observed_at_ms: 1_132,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ ]);
+
+ let receipt = fetch_and_ingest_relay_events(
+ &adapter,
+ &store,
+ post_relay_fetch_request(1_130, 1)
+ .with_raw_event_scan_limit(2)
+ .expect("raw scan limit"),
+ )
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 0);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 1);
+ assert_eq!(receipt.skipped_over_limit_count, 1);
+ assert_eq!(receipt.events.len(), 2);
+ assert_eq!(receipt.eose_count, 1);
+ assert!(
+ store
+ .get_event(accepted_id.as_str())
+ .await
+ .expect("accepted lookup")
+ .is_none()
+ );
+}
+
+#[tokio::test]
+async fn fetch_subscription_mode_and_store_errors_are_reported() {
+ let signed = signed_post("subscription");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: signed.raw_json.clone(),
+ observed_at_ms: 1_200,
+ }]);
+
+ let receipt = fetch_and_ingest_relay_events(
+ &adapter,
+ &store,
+ RadrootsRelayFetchRequest::subscription(1_200, 10, [post_relay_fetch_filter(10)])
+ .expect("subscription request"),
+ )
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 1);
+ let observations = store
+ .observations_for_event(signed.id.as_str())
+ .await
+ .expect("observations");
+ assert_eq!(observations.len(), 1);
+ assert_eq!(
+ observations[0].observation_type,
+ RadrootsTransportObservationType::NostrSubscription
+ );
+
+ let closed_store = RadrootsEventStore::open_memory().await.expect("store");
+ closed_store.pool().close().await;
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: signed.raw_json,
+ observed_at_ms: 1_210,
+ }]);
+ let receipt =
+ fetch_and_ingest_relay_events(&adapter, &closed_store, post_relay_fetch_request(1_210, 10))
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 0);
+ assert_eq!(receipt.malformed_count, 1);
+ assert!(receipt.events[0].malformed);
+ assert!(receipt.events[0].message.is_some());
+}
+
+#[tokio::test]
+async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() {
+ let signed = signed_post("hello");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags.clone(),
+ signed.content.clone(),
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![
+ RELAY_PRIMARY_WSS.to_owned(),
+ RELAY_SECONDARY_WSS.to_owned(),
+ RELAY_TERTIARY_WSS.to_owned(),
+ ],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+ assert_eq!(publish_claim.state, RadrootsOutboxEventState::Publishing);
+
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
+ .with_outcome(
+ RELAY_SECONDARY_WSS,
+ RadrootsRelayOutcome::timeout("timeout: no OK"),
+ )
+ .with_outcome(
+ RELAY_TERTIARY_WSS,
+ RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"),
+ );
+ let first = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &adapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(first.publish.attempted_count, 3);
+ assert_eq!(first.publish.accepted_count, 2);
+ assert!(!first.publish.quorum_met);
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
+
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS)
+ .expect("primary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::Accepted
+ );
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_SECONDARY_WSS)
+ .expect("secondary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable
+ );
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_TERTIARY_WSS)
+ .expect("tertiary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::Accepted
+ );
+
+ let retry_claim = outbox
+ .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500)
+ .await
+ .expect("claim")
+ .expect("retry claim");
+ let retry_adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted());
+ let second = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &retry_adapter,
+ &retry_claim,
+ RadrootsOutboxPublishPolicy::new(3_000),
+ 2_600,
+ )
+ .await
+ .expect("retry publish");
+
+ assert_eq!(second.local_ingest.event_id, signed.id);
+ assert_eq!(second.publish.attempted_count, 1);
+ assert_eq!(retry_adapter.captured_raw_events().len(), 1);
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Published);
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
+
+ let observations = store
+ .observations_for_event(signed.id.as_str())
+ .await
+ .expect("observations");
+ assert_eq!(observations.len(), 3);
+}
+
+#[tokio::test]
+async fn outbox_publish_transport_failure_releases_retryable_claim() {
+ let signed = signed_post("adapter transport failure");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags.clone(),
+ signed.content.clone(),
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+
+ let published = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &TransportFailurePublishAdapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(published.publish.attempted_count, 2);
+ assert_eq!(published.publish.accepted_count, 0);
+ assert_eq!(published.publish.retryable_count, 2);
+ assert_eq!(published.publish.terminal_count, 0);
+ assert!(!published.publish.quorum_met);
+ assert!(
+ published
+ .publish
+ .relays
+ .iter()
+ .all(|relay| relay.outcome.kind == RadrootsRelayOutcomeKind::ConnectionFailed)
+ );
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
+ assert!(event.claim_token.is_none());
+ assert_eq!(event.next_attempt_after_ms, 2_500);
+
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ assert_eq!(targets.len(), 2);
+ assert!(
+ targets
+ .iter()
+ .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable)
+ );
+ assert!(
+ outbox
+ .claim_next_ready_event("publisher", "publish-b", 4_000, 2_499)
+ .await
+ .expect("early claim")
+ .is_none()
+ );
+ let retry_claim = outbox
+ .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500)
+ .await
+ .expect("retry claim")
+ .expect("retry claim");
+ assert_eq!(retry_claim.outbox_event_id, receipt.outbox_event_id);
+ assert_eq!(retry_claim.state, RadrootsOutboxEventState::Publishing);
+}
+
+#[tokio::test]
+async fn outbox_publish_marks_published_without_adapter_when_all_relays_already_accepted() {
+ let signed = signed_post("already accepted");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags.clone(),
+ signed.content.clone(),
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+ let initial_targets = publish_claim.delivery_targets.clone();
+ outbox
+ .mark_delivery_target_accepted(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ initial_targets[0].delivery_target_id,
+ 2_150,
+ )
+ .await
+ .expect("primary accepted");
+ outbox
+ .mark_delivery_target_accepted(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ initial_targets[1].delivery_target_id,
+ 2_151,
+ )
+ .await
+ .expect("secondary accepted");
+
+ let adapter = RadrootsMockRelayPublishAdapter::new();
+ let published = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &adapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(published.local_ingest.event_id, signed.id);
+ assert_eq!(published.publish.event_id, signed.id);
+ assert_eq!(published.publish.attempted_count, 0);
+ assert_eq!(published.publish.accepted_count, 2);
+ assert_eq!(published.publish.quorum, 2);
+ assert!(published.publish.quorum_met);
+ assert!(published.publish.relays.is_empty());
+ assert!(adapter.captured_raw_events().is_empty());
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Published);
+ assert!(event.claim_token.is_none());
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
+}
+
+#[tokio::test]
+async fn outbox_publish_marks_published_when_delivery_plan_satisfaction_is_met_with_failure_diagnostics()
+ {
+ let signed = signed_post("quorum");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags.clone(),
+ signed.content.clone(),
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(outbox_operation_input(
+ draft,
+ vec![
+ RELAY_PRIMARY_WSS.to_owned(),
+ RELAY_SECONDARY_WSS.to_owned(),
+ RELAY_TERTIARY_WSS.to_owned(),
+ ],
+ RadrootsTransportSatisfactionPolicy::AtLeast(2),
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
+ .with_outcome(
+ RELAY_SECONDARY_WSS,
+ RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"),
+ )
+ .with_outcome(
+ RELAY_TERTIARY_WSS,
+ RadrootsRelayOutcome::classify("restricted: group write denied"),
+ );
+ let published = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &adapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(published.publish.quorum, 2);
+ assert_eq!(published.publish.accepted_count, 2);
+ assert_eq!(published.publish.terminal_count, 1);
+ assert!(published.publish.quorum_met);
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Published);
+ assert!(event.claim_token.is_none());
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete);
+
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_TERTIARY_WSS)
+ .expect("tertiary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::FailedTerminal
+ );
+ assert!(
+ outbox
+ .claim_next_ready_event("publisher", "publish-b", 4_000, 2_300)
+ .await
+ .expect("claim")
+ .is_none()
+ );
+
+ let observations = store
+ .observations_for_event(signed.id.as_str())
+ .await
+ .expect("observations");
+ assert_eq!(observations.len(), 2);
+}
+
+#[tokio::test]
+async fn outbox_publish_republishes_accepted_relays_when_policy_requests_it() {
+ let signed = signed_post("republish accepted");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags.clone(),
+ signed.content.clone(),
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+ let initial_targets = publish_claim.delivery_targets.clone();
+ outbox
+ .mark_delivery_target_accepted(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ initial_targets[0].delivery_target_id,
+ 2_150,
+ )
+ .await
+ .expect("primary accepted");
+
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
+ .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted());
+ let published = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &adapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500)
+ .republish_accepted_relays(true)
+ .relay_url_policy(RadrootsRelayUrlPolicy::Public),
+ 2_200,
+ )
+ .await
+ .expect("publish");
+
+ assert_eq!(published.local_ingest.event_id, signed.id);
+ assert_eq!(published.publish.attempted_count, 2);
+ assert_eq!(published.publish.accepted_count, 2);
+ assert_eq!(published.publish.quorum, 1);
+ assert!(published.publish.quorum_met);
+ assert_eq!(adapter.captured_raw_events().len(), 1);
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Published);
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ assert!(
+ targets
+ .iter()
+ .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::Accepted)
+ );
+}
+
+#[tokio::test]
+async fn outbox_publish_requires_claimed_signed_event() {
+ let signed = signed_post("missing signature");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags,
+ signed.content,
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![RELAY_PRIMARY_WSS.to_owned()],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ let adapter = RadrootsMockRelayPublishAdapter::new();
+
+ let error = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &adapter,
+ &claimed,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 1_100,
+ )
+ .await
+ .expect_err("missing signed event");
+
+ assert!(matches!(
+ error,
+ RadrootsRelayTransportError::MissingSignedOutboxEvent(event_id)
+ if event_id == receipt.outbox_event_id
+ ));
+ assert!(adapter.captured_raw_events().is_empty());
+}
+
+#[tokio::test]
+async fn outbox_publish_propagates_non_transport_adapter_errors_after_target_filtering() {
+ let signed = signed_post("adapter non transport failure");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsFrozenEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at,
+ signed.tags,
+ signed.content,
+ signed.pubkey.as_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_targets_outbox_operation_input(
+ draft,
+ vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "sign-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+ complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ outbox.recover_expired_claims(2_001).await.expect("recover");
+ let mut publish_claim = outbox
+ .claim_next_ready_event("publisher", "publish-a", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .expect("publish claim");
+ publish_claim.delivery_targets.truncate(1);
+
+ let error = publish_claimed_outbox_event(
+ &outbox,
+ &store,
+ &NostrJsonFailurePublishAdapter,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect_err("adapter error");
+
+ assert!(matches!(
+ error,
+ RadrootsRelayTransportError::NostrEventJson(_)
+ ));
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Publishing);
+}
+
+#[tokio::test]
+async fn smoke_relay_fetch_processes_one_thousand_event_receipts() {
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let mut items = Vec::new();
+ for index in 0..1_000 {
+ let signed = signed_post(format!("fetch-smoke-{index}").as_str());
+ let relay_url = match index % 3 {
+ 0 => RELAY_PRIMARY_WSS,
+ 1 => RELAY_SECONDARY_WSS,
+ _ => RELAY_TERTIARY_WSS,
+ };
+ items.push(RadrootsRelayFetchItem::Event {
+ relay_url: relay_url.to_owned(),
+ raw_json: signed.raw_json,
+ observed_at_ms: 10_000 + index,
+ });
+ }
+ let adapter = RadrootsMockRelayFetchAdapter::new(items);
+ let receipt =
+ fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(10_000, 1_000))
+ .await
+ .expect("fetch");
+
+ assert_eq!(receipt.inserted_count, 1_000);
+ assert_eq!(receipt.duplicate_count, 0);
+ assert_eq!(receipt.malformed_count, 0);
+ assert_eq!(receipt.unsupported_count, 0);
+ assert_eq!(receipt.events.len(), 1_000);
+ assert!(receipt.events.iter().all(|event| event.projection_eligible));
+ let replay = store
+ .events_since_cursor("fetch-smoke", 1_000)
+ .await
+ .expect("replay");
+ assert_eq!(replay.len(), 1_000);
+}
diff --git a/tools/xtask/src/hygiene.rs b/tools/xtask/src/hygiene.rs
@@ -24,7 +24,7 @@ pub fn validate_forbidden_identifiers(root: &Path) -> Result<(), String> {
let mut failures = Vec::new();
reject_substrings(
root,
- &[PathBuf::from("crates/relay_transport/src")],
+ &[PathBuf::from("crates/transport_nostr/src")],
&["RadrootsEventIngest::verified"],
"relay fetch must not bypass event-store verification",
&[],
@@ -408,7 +408,7 @@ mod tests {
let root = unique_temp_dir("clean");
write_file(
&root,
- "crates/relay_transport/src/fetch.rs",
+ "crates/transport_nostr/src/fetch.rs",
"fn fetch() { let _ = RadrootsEventIngest::new; }\n",
);
write_file(
@@ -430,7 +430,7 @@ mod tests {
let root = unique_temp_dir("dirty");
write_file(
&root,
- "crates/relay_transport/src/fetch.rs",
+ "crates/transport_nostr/src/fetch.rs",
"fn fetch() { let _ = RadrootsEventIngest::verified; }\n",
);
write_file(
@@ -476,7 +476,7 @@ mod tests {
let root = unique_temp_dir("run");
write_file(
&root,
- "crates/relay_transport/src/fetch.rs",
+ "crates/transport_nostr/src/fetch.rs",
"fn fetch() { let _ = RadrootsEventIngest::new; }\n",
);
run(&["forbidden-identifiers".to_string()], &root).expect("hygiene run");