commit 5b0ad5e9d96770fcda89d21d86a50453076b7e3c
parent de48b1a1168deb1fa722cb8fe2020e2032b21473
Author: triesap <tyson@radroots.org>
Date: Wed, 1 Jul 2026 21:11:35 +0000
cli: migrate relay reads to shared transport
Remove the CLI-local direct relay fetch module and route trade event list, order rebind checks, sync pull, and market refresh through the shared relay transport fetch receipt boundary with fail-closed filters and relay evidence.
Validation: cargo extbuild run -- cargo fmt --all; cargo extbuild run -- cargo check --workspace --all-targets; cargo extbuild run -- cargo test --workspace cli_production_sources_reject_direct_relay_fetch_helpers; cargo extbuild run -- cargo test --workspace migrated_cli_paths_are_guarded_against_workflow_bypasses; cargo extbuild run -- cargo test --workspace sync_pull; cargo extbuild run -- cargo test --workspace local_order_event_list_attempts_configured_shared_relay_transport; cargo extbuild run -- cargo test --workspace direct_relay
Diffstat:
9 files changed, 242 insertions(+), 356 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -3630,6 +3630,7 @@ dependencies = [
"radroots_nostr_connect",
"radroots_nostr_signer",
"radroots_protected_store",
+ "radroots_relay_transport",
"radroots_replica_db",
"radroots_replica_db_schema",
"radroots_replica_sync",
@@ -3851,6 +3852,7 @@ dependencies = [
"serde",
"serde_json",
"thiserror 1.0.69",
+ "tokio",
"url",
]
diff --git a/Cargo.toml b/Cargo.toml
@@ -42,6 +42,7 @@ radroots_protected_store = { path = "../lib/crates/protected_store", features =
radroots_replica_db = { path = "../lib/crates/replica_db" }
radroots_replica_db_schema = { path = "../lib/crates/replica_db_schema" }
radroots_replica_sync = { path = "../lib/crates/replica_sync" }
+radroots_relay_transport = { path = "../lib/crates/relay_transport", default-features = false, features = ["runtime-tokio"] }
radroots_runtime = { path = "../lib/crates/runtime" }
radroots_runtime_paths = { path = "../lib/crates/runtime_paths" }
radroots_sdk = { path = "../sdk/crates/sdk", features = ["local-runtime-radrootsd-proxy", "relay-runtime"] }
diff --git a/src/runtime/direct_relay.rs b/src/runtime/direct_relay.rs
@@ -1,203 +0,0 @@
-use std::time::Duration;
-
-use radroots_nostr::prelude::{
- RadrootsNostrClient, RadrootsNostrError, RadrootsNostrEvent, RadrootsNostrFilter,
- RadrootsNostrOutput,
-};
-
-const RELAY_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
-const RELAY_FETCH_TIMEOUT: Duration = Duration::from_secs(10);
-
-#[derive(Debug, Clone, PartialEq, Eq)]
-pub struct DirectRelayFailure {
- pub relay: String,
- pub reason: String,
-}
-
-#[derive(Debug, Clone)]
-pub struct DirectRelayFetchReceipt {
- pub target_relays: Vec<String>,
- pub connected_relays: Vec<String>,
- pub failed_relays: Vec<DirectRelayFailure>,
- pub events: Vec<RadrootsNostrEvent>,
-}
-
-#[derive(Debug, thiserror::Error)]
-pub enum DirectRelayFetchError {
- #[error("direct relay fetch requires at least one configured relay")]
- MissingRelays,
- #[error("failed to build async runtime for direct relay fetch: {0}")]
- Runtime(String),
- #[error("failed to configure relay `{relay}` for direct relay fetch: {source}")]
- RelayConfig {
- relay: String,
- #[source]
- source: RadrootsNostrError,
- },
- #[error("direct relay connection failed: {reason}")]
- Connect {
- reason: String,
- target_relays: Vec<String>,
- failed_relays: Vec<DirectRelayFailure>,
- },
- #[error("direct relay fetch failed: {0}")]
- Fetch(#[source] RadrootsNostrError),
-}
-
-pub fn fetch_events_from_relays(
- relay_urls: &[String],
- filter: RadrootsNostrFilter,
-) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError> {
- fetch_events_from_relays_with_timeout(relay_urls, filter, RELAY_FETCH_TIMEOUT)
-}
-
-pub fn fetch_events_from_relays_with_timeout(
- relay_urls: &[String],
- filter: RadrootsNostrFilter,
- fetch_timeout: Duration,
-) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError> {
- if relay_urls.is_empty() {
- return Err(DirectRelayFetchError::MissingRelays);
- }
-
- let runtime = tokio::runtime::Builder::new_multi_thread()
- .enable_all()
- .build()
- .map_err(|error| DirectRelayFetchError::Runtime(error.to_string()))?;
-
- runtime.block_on(fetch_events_from_relays_async(
- relay_urls,
- filter,
- fetch_timeout,
- RELAY_CONNECT_TIMEOUT,
- ))
-}
-
-async fn fetch_events_from_relays_async(
- relay_urls: &[String],
- filter: RadrootsNostrFilter,
- fetch_timeout: Duration,
- connect_timeout: Duration,
-) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError> {
- let client = RadrootsNostrClient::new_signerless();
-
- for relay_url in relay_urls {
- client.add_read_relay(relay_url).await.map_err(|source| {
- DirectRelayFetchError::RelayConfig {
- relay: relay_url.clone(),
- source,
- }
- })?;
- }
-
- let connection_output = client.try_connect(connect_timeout).await;
- let failed_relays = relay_failures_from_output(&connection_output);
- if connection_output.success.is_empty() {
- return Err(DirectRelayFetchError::Connect {
- reason: summarize_failures(&failed_relays),
- target_relays: relay_urls.to_vec(),
- failed_relays,
- });
- }
-
- let events = client
- .fetch_events(filter, fetch_timeout)
- .await
- .map_err(DirectRelayFetchError::Fetch)?;
-
- Ok(DirectRelayFetchReceipt {
- target_relays: relay_urls.to_vec(),
- connected_relays: connection_output
- .success
- .iter()
- .map(ToString::to_string)
- .collect(),
- failed_relays,
- events,
- })
-}
-
-fn relay_failures_from_output<T: std::fmt::Debug>(
- output: &RadrootsNostrOutput<T>,
-) -> Vec<DirectRelayFailure> {
- output
- .failed
- .iter()
- .map(|(relay, reason)| DirectRelayFailure {
- relay: relay.to_string(),
- reason: reason.to_string(),
- })
- .collect()
-}
-
-fn summarize_failures(failed_relays: &[DirectRelayFailure]) -> String {
- if failed_relays.is_empty() {
- return "no relay acknowledged the operation".to_owned();
- }
-
- failed_relays
- .iter()
- .map(|failure| format!("{}: {}", failure.relay, failure.reason))
- .collect::<Vec<_>>()
- .join("; ")
-}
-
-#[cfg(test)]
-mod tests {
- use std::time::Duration;
-
- use radroots_nostr::prelude::RadrootsNostrFilter;
-
- use super::{
- DirectRelayFetchError, fetch_events_from_relays_async,
- fetch_events_from_relays_with_timeout,
- };
-
- #[test]
- fn fetch_events_requires_relays_before_runtime_work() {
- let err = fetch_events_from_relays_with_timeout(
- &[],
- RadrootsNostrFilter::new(),
- Duration::from_millis(1),
- )
- .expect_err("missing relay error");
-
- assert!(matches!(err, DirectRelayFetchError::MissingRelays));
- }
-
- #[test]
- fn fetch_events_rejects_invalid_relay_urls() {
- let err = fetch_events_from_relays_with_timeout(
- &["not-a-relay".to_owned()],
- RadrootsNostrFilter::new(),
- Duration::from_millis(1),
- )
- .expect_err("relay config error");
-
- assert!(matches!(err, DirectRelayFetchError::RelayConfig { .. }));
- }
-
- #[tokio::test]
- async fn fetch_events_reports_connection_failure() {
- let err = fetch_events_from_relays_async(
- &["ws://127.0.0.1:9".to_owned()],
- RadrootsNostrFilter::new(),
- Duration::from_millis(1),
- Duration::from_millis(50),
- )
- .await
- .expect_err("connection failure");
-
- match err {
- DirectRelayFetchError::Connect {
- target_relays,
- failed_relays,
- ..
- } => {
- assert_eq!(target_relays, vec!["ws://127.0.0.1:9"]);
- assert_eq!(failed_relays.len(), 1);
- }
- _ => panic!("expected connection failure"),
- }
- }
-}
diff --git a/src/runtime/mod.rs b/src/runtime/mod.rs
@@ -1,6 +1,5 @@
pub mod account;
pub mod config;
-pub mod direct_relay;
pub mod farm;
pub mod farm_config;
pub mod find;
diff --git a/src/runtime/order.rs b/src/runtime/order.rs
@@ -39,6 +39,9 @@ use radroots_nostr::prelude::{
RadrootsNostrEvent, RadrootsNostrFilter, radroots_event_from_nostr, radroots_nostr_filter_tag,
radroots_nostr_kind,
};
+use radroots_relay_transport::{
+ RadrootsRelayFetchFailure, RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError,
+};
use radroots_replica_db::{
ReplicaSql, ReplicaTradeProductSummaryRow, nostr_event_head, trade_product,
};
@@ -72,14 +75,13 @@ use crate::cli::global::{
use crate::runtime::RuntimeError;
use crate::runtime::account;
use crate::runtime::config::RuntimeConfig;
-use crate::runtime::direct_relay::{
- DirectRelayFailure, DirectRelayFetchError, DirectRelayFetchReceipt, fetch_events_from_relays,
-};
use crate::runtime::local_events::{
get_shared_record, list_shared_records_before, list_shared_records_latest,
shared_local_events_db_path,
};
-use crate::runtime::sdk::{CliSdkAdapterError, CliSdkSession};
+use crate::runtime::sdk::{
+ CliSdkAdapterError, CliSdkSession, fetch_relay_events_via_shared_transport,
+};
use crate::runtime::sync::{RelayIngestScope, relay_provenance_relays_for_scope};
use crate::view::runtime::{
OrderAppRecordExportView, OrderAppRecordListView, OrderAppRecordSummaryView,
@@ -99,7 +101,7 @@ const ORDER_DECISION_SOURCE: &str = "SDK trade decision · local key";
const ORDER_REVISION_PROPOSAL_SOURCE: &str = "SDK trade revision proposal · local key";
const ORDER_REVISION_DECISION_SOURCE: &str = "SDK trade revision decision · local key";
const ORDER_CANCELLATION_SOURCE: &str = "SDK trade cancellation · local key";
-const ORDER_EVENT_LIST_SOURCE: &str = "direct Nostr relay fetch · selected seller identity";
+const ORDER_EVENT_LIST_SOURCE: &str = "shared relay transport fetch · selected seller identity";
const ORDER_STATUS_SDK_SOURCE: &str = "SDK local trade projection";
const ORDER_EVENT_LIST_RELAY_ACTION: &str =
"radroots --relay wss://relay.example.com trade event list";
@@ -1184,23 +1186,18 @@ pub fn event_list(
};
let seller_pubkey = actor_context.seller_pubkey;
let filter = order_request_filter(seller_pubkey.as_str(), order_id)?;
- let receipt = match fetch_events_from_relays(&config.relay.urls, filter) {
- Ok(receipt) => receipt,
- Err(DirectRelayFetchError::Connect {
- reason,
- target_relays,
- failed_relays,
- }) => {
- return Ok(order_event_list_unavailable(
- seller_pubkey,
- actor_context.source,
- reason,
- target_relays,
- failed_relays,
- ));
- }
- Err(error) => return Err(RuntimeError::Network(error.to_string())),
- };
+ let receipt =
+ fetch_relay_events_via_shared_transport(&config.relay.urls, now_unix_ms(), 1_000, filter)
+ .map_err(order_relay_fetch_error)?;
+ if receipt.connected_relays.is_empty() && !receipt.failed_relays.is_empty() {
+ return Ok(order_event_list_unavailable(
+ seller_pubkey,
+ actor_context.source,
+ relay_fetch_failure_reason(&receipt.failed_relays),
+ receipt.target_relays,
+ receipt.failed_relays,
+ ));
+ }
Ok(order_event_list_from_receipt(
seller_pubkey,
@@ -1804,7 +1801,7 @@ fn order_event_list_unavailable(
actor_context_source: &'static str,
reason: String,
target_relays: Vec<String>,
- failed_relays: Vec<DirectRelayFailure>,
+ failed_relays: Vec<RadrootsRelayFetchFailure>,
) -> OrderEventListView {
OrderEventListView {
state: "unavailable".to_owned(),
@@ -1818,7 +1815,7 @@ fn order_event_list_unavailable(
decoded_count: 0,
skipped_count: 0,
count: 0,
- reason: Some(format!("direct relay connection failed: {reason}")),
+ reason: Some(format!("relay transport fetch failed: {reason}")),
orders: Vec::new(),
actions: Vec::new(),
}
@@ -1828,13 +1825,14 @@ fn order_event_list_from_receipt(
seller_pubkey: String,
order_id: Option<&str>,
actor_context_source: &'static str,
- receipt: DirectRelayFetchReceipt,
+ receipt: RadrootsRelayFetchedEventsReceipt,
) -> OrderEventListView {
- let DirectRelayFetchReceipt {
+ let RadrootsRelayFetchedEventsReceipt {
target_relays,
connected_relays,
failed_relays,
events,
+ ..
} = receipt;
let fetched_count = events.len();
let mut skipped_count = 0usize;
@@ -1842,7 +1840,7 @@ fn order_event_list_from_receipt(
let mut orders = Vec::new();
for event in events {
- match order_event_list_entry_from_event(&event, seller_pubkey.as_str()) {
+ match order_event_list_entry_from_event(&event.event, seller_pubkey.as_str()) {
Ok(entry) => {
decoded_count += 1;
if order_id.is_none_or(|order_id| entry.id == order_id) {
@@ -4523,13 +4521,14 @@ fn order_rebind_existing_request_check(
loaded.document.order.seller_pubkey.as_str(),
Some(loaded.document.order.order_id.as_str()),
)?;
- let receipt = fetch_events_from_relays(&config.relay.urls, filter)
- .map_err(|error| RuntimeError::Network(error.to_string()))?;
+ let receipt =
+ fetch_relay_events_via_shared_transport(&config.relay.urls, now_unix_ms(), 1_000, filter)
+ .map_err(order_relay_fetch_error)?;
let mut event_ids = receipt
.events
.iter()
- .filter_map(|event| {
- order_submit_request_from_event(event, loaded)
+ .filter_map(|fetched| {
+ order_submit_request_from_event(&fetched.event, loaded)
.ok()
.map(|request| request.request_event_id)
})
@@ -5238,11 +5237,26 @@ fn order_buyer_failure_detail(
detail
}
-fn relay_failures(failures: Vec<DirectRelayFailure>) -> Vec<RelayFailureView> {
+fn order_relay_fetch_error(error: RadrootsRelayTransportError) -> RuntimeError {
+ RuntimeError::Network(error.to_string())
+}
+
+fn relay_fetch_failure_reason(failed_relays: &[RadrootsRelayFetchFailure]) -> String {
+ if failed_relays.is_empty() {
+ return "no relay acknowledged the fetch".to_owned();
+ }
+ failed_relays
+ .iter()
+ .map(|failure| format!("{}: {}", failure.relay_url, failure.reason))
+ .collect::<Vec<_>>()
+ .join("; ")
+}
+
+fn relay_failures(failures: Vec<RadrootsRelayFetchFailure>) -> Vec<RelayFailureView> {
failures
.into_iter()
.map(|failure| RelayFailureView {
- relay: failure.relay,
+ relay: failure.relay_url,
reason: failure.reason,
})
.collect()
@@ -5431,6 +5445,13 @@ fn now_unix() -> u64 {
.unwrap_or_default()
}
+fn now_unix_ms() -> i64 {
+ SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .map(|value| i64::try_from(value.as_millis()).unwrap_or(i64::MAX))
+ .unwrap_or_default()
+}
+
fn next_order_id() -> String {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
diff --git a/src/runtime/sdk.rs b/src/runtime/sdk.rs
@@ -15,6 +15,10 @@ use radroots_nostr_connect::prelude::{
RADROOTS_NOSTR_CONNECT_RPC_KIND, RadrootsNostrConnectBunkerUri,
RadrootsNostrConnectClientTarget, RadrootsNostrConnectError, RadrootsNostrConnectUri,
};
+use radroots_relay_transport::{
+ RadrootsNostrClientFetchAdapter, RadrootsRelayFetchRequest, RadrootsRelayFetchedEventsReceipt,
+ RadrootsRelayTransportError, fetch_relay_events_blocking,
+};
use radroots_sdk::{
RadrootsClient, RadrootsClientBuilder, RadrootsSdkError, RadrootsSdkLocalKeySigner,
RadrootsSdkMycNip46RequestPolicy, RadrootsSdkMycNip46Signer, RadrootsSdkNip46Transport,
@@ -37,6 +41,7 @@ use crate::runtime::config::{
const SDK_STORAGE_DIR_NAME: &str = "sdk";
const RADROOTSD_PROXY_SECRET_SERVICE: &str = "org.radroots.cli.radrootsd-proxy";
+const CLI_RELAY_FETCH_TIMEOUT_MS: u64 = 10_000;
pub(crate) const MYC_NIP46_SESSION_SECRET_SERVICE: &str = "org.radroots.cli.myc-nip46-session";
#[derive(Debug, thiserror::Error)]
@@ -588,6 +593,18 @@ pub(crate) fn sdk_runtime() -> Result<Runtime, RuntimeError> {
})
}
+pub(crate) fn fetch_relay_events_via_shared_transport(
+ relay_urls: &[String],
+ observed_at_ms: i64,
+ max_events: usize,
+ filter: RadrootsNostrFilter,
+) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> {
+ let request = RadrootsRelayFetchRequest::fetch(observed_at_ms, max_events, [filter])?
+ .with_relay_urls(relay_urls.iter().cloned())
+ .with_timeout_ms(CLI_RELAY_FETCH_TIMEOUT_MS);
+ fetch_relay_events_blocking(&RadrootsNostrClientFetchAdapter, request)
+}
+
fn memory_builder(config: &CliSdkConfig) -> RadrootsClientBuilder {
config.relay_urls.iter().fold(
RadrootsClient::builder()
@@ -699,14 +716,6 @@ mod tests {
lifecycle: &'static str,
}
- struct DirectRelayConsumerException {
- path: &'static str,
- required_tokens: &'static [&'static str],
- owner: &'static str,
- reason: &'static str,
- lifecycle: &'static str,
- }
-
struct MigratedCliPathGuard {
label: &'static str,
path: &'static str,
@@ -768,9 +777,16 @@ mod tests {
DirectRrRsDependency {
section: "dependencies",
name: "radroots_nostr",
- owner: "non-migrated-direct-relay-workflows",
- reason: "direct relay fetch/publish and event conversion for active non-migrated commands",
- lifecycle: "retain until direct relay command families migrate or are retired",
+ owner: "cli-signer-and-event-runtime",
+ reason: "remote signer relay transport, account event conversion, and direct publish command transport",
+ lifecycle: "retain while CLI owns signer transport and direct publish selection",
+ },
+ DirectRrRsDependency {
+ section: "dependencies",
+ name: "radroots_relay_transport",
+ owner: "cli-shared-relay-read-boundary",
+ reason: "shared fail-closed relay fetch receipts for trade event list, sync pull, and market refresh",
+ lifecycle: "retain until those read surfaces are fully SDK-owned",
},
DirectRrRsDependency {
section: "dependencies",
@@ -865,21 +881,12 @@ mod tests {
},
];
- const DIRECT_RELAY_CONSUMER_EXCEPTIONS: &[DirectRelayConsumerException] = &[
- DirectRelayConsumerException {
- path: "src/runtime/order.rs",
- required_tokens: &["pub fn event_list(", "fetch_events_from_relays"],
- owner: "non-migrated-trade-event-and-derived-projection-reads",
- reason: "event listing, draft checks, and derived projection maintenance still use direct relay reads outside migrated trade mutation paths",
- lifecycle: "retain only for non-migrated read/projection surfaces",
- },
- DirectRelayConsumerException {
- path: "src/runtime/sync.rs",
- required_tokens: &["fetch_events_from_relays", "pull_with_fetcher"],
- owner: "sync.pull-and-market-refresh",
- reason: "relay ingest into the derived projection cache",
- lifecycle: "retain until relay ingest and derived projection repair migrate to SDK APIs",
- },
+ const DIRECT_RELAY_FETCH_DISALLOWED_TOKENS: &[&str] = &[
+ "pub mod direct_relay",
+ "use crate::runtime::direct_relay",
+ "fetch_events_from_relays",
+ "fetch_events_from_relays_with_timeout",
+ ".fetch_events(",
];
const MIGRATED_CLI_PATH_GUARDS: &[MigratedCliPathGuard] = &[
@@ -1262,27 +1269,31 @@ mod tests {
}
#[test]
- fn direct_relay_consumer_exceptions_are_explicit() {
- let actual = direct_relay_consumer_exception_paths();
- let expected = DIRECT_RELAY_CONSUMER_EXCEPTIONS
+ fn cli_production_sources_reject_direct_relay_fetch_helpers() {
+ let manifest_dir = Path::new(env!("CARGO_MANIFEST_DIR"));
+ let mut files = Vec::new();
+ collect_rs_files(manifest_dir.join("src").as_path(), &mut files);
+ files.sort();
+
+ let findings = files
.iter()
- .map(|consumer| consumer.path.to_owned())
- .collect::<BTreeSet<_>>();
+ .flat_map(|file| {
+ let source = fs::read_to_string(file).expect("read cli source");
+ let relative_path = relative_source_path(manifest_dir, file.as_path());
+ match production_source_without_tests(&relative_path, &source) {
+ Ok(production_source) => {
+ direct_relay_fetch_findings(&relative_path, production_source.as_str())
+ }
+ Err(error) => vec![error],
+ }
+ })
+ .collect::<Vec<_>>();
- assert_eq!(actual, expected);
- for consumer in DIRECT_RELAY_CONSUMER_EXCEPTIONS {
- let source = crate_source(consumer.path);
- for token in consumer.required_tokens {
- assert!(
- source.contains(token),
- "{} does not contain direct-relay exception token `{token}`",
- consumer.path
- );
- }
- assert!(!consumer.owner.trim().is_empty());
- assert!(!consumer.reason.trim().is_empty());
- assert!(!consumer.lifecycle.trim().is_empty());
- }
+ assert!(
+ findings.is_empty(),
+ "CLI production sources contain direct relay fetch helpers:\n{}",
+ findings.join("\n")
+ );
}
#[test]
@@ -1516,27 +1527,6 @@ mod tests {
format!("{}:{}", dependency.section, dependency.name)
}
- fn direct_relay_consumer_exception_paths() -> BTreeSet<String> {
- let manifest_dir = Path::new(env!("CARGO_MANIFEST_DIR"));
- let mut files = Vec::new();
- collect_rs_files(manifest_dir.join("src/runtime").as_path(), &mut files);
- files
- .into_iter()
- .filter(|file| {
- !matches!(
- file.file_name().and_then(|name| name.to_str()),
- Some("direct_relay.rs" | "sdk.rs")
- )
- })
- .filter_map(|file| {
- let source = fs::read_to_string(&file).expect("read runtime source");
- source
- .contains("use crate::runtime::direct_relay")
- .then(|| relative_source_path(manifest_dir, file.as_path()))
- })
- .collect()
- }
-
fn relative_source_path(root: &Path, path: &Path) -> String {
path.strip_prefix(root)
.expect("source path under manifest root")
@@ -1640,6 +1630,20 @@ mod tests {
.collect()
}
+ fn direct_relay_fetch_findings(label: &str, source: &str) -> Vec<String> {
+ DIRECT_RELAY_FETCH_DISALLOWED_TOKENS
+ .iter()
+ .flat_map(|token| {
+ source.match_indices(token).map(move |(index, _)| {
+ format!(
+ "{label}:{} uses direct relay fetch token `{token}`",
+ line_number(source, index)
+ )
+ })
+ })
+ .collect()
+ }
+
fn production_source_without_tests(path: &str, source: &str) -> Result<String, String> {
let code_source = rust_code_without_non_code(path, source)?;
let mut production_source = String::with_capacity(code_source.len());
diff --git a/src/runtime/sync.rs b/src/runtime/sync.rs
@@ -11,6 +11,9 @@ use radroots_events::kinds::{
use radroots_nostr::prelude::{
RadrootsNostrFilter, RadrootsNostrTimestamp, radroots_event_from_nostr, radroots_nostr_kind,
};
+use radroots_relay_transport::{
+ RadrootsRelayFetchFailure, RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError,
+};
use radroots_replica_db::{ReplicaSql, migrations};
use radroots_replica_sync::{
RadrootsReplicaEventsError, RadrootsReplicaIngestOutcome, radroots_replica_ingest_event,
@@ -27,10 +30,10 @@ use serde_json::json;
use crate::cli::global::SyncWatchArgs;
use crate::runtime::RuntimeError;
use crate::runtime::config::RuntimeConfig;
-use crate::runtime::direct_relay::{
- DirectRelayFailure, DirectRelayFetchError, DirectRelayFetchReceipt, fetch_events_from_relays,
+use crate::runtime::sdk::{
+ CliSdkAdapterError, CliSdkSession, fetch_relay_events_via_shared_transport,
+ sdk_relay_url_policy,
};
-use crate::runtime::sdk::{CliSdkAdapterError, CliSdkSession, sdk_relay_url_policy};
use crate::view::runtime::{
RelayFailureView, SyncActionView, SyncFreshnessView, SyncQueueView, SyncRunFreshnessView,
SyncStatusView, SyncWatchFrameView, SyncWatchView,
@@ -44,7 +47,7 @@ const SYNC_PULL_ACTION: &str = "radroots sync pull";
const SYNC_PUSH_ACTION: &str = "radroots sync push";
const SYNC_READY_ACTION: &str = "radroots market product search eggs";
const MARKET_READY_ACTION: &str = "radroots market product search eggs";
-const INGEST_SOURCE: &str = "direct Nostr relay fetch · local replica ingest";
+const INGEST_SOURCE: &str = "shared relay transport fetch · local replica ingest";
const RELAY_FETCH_LIMIT: usize = 1_000;
const RELAY_FETCH_MAX_PAGES: usize = 5;
const MARKET_FRESHNESS_STALE_AFTER_SECONDS: u64 = 15 * 60;
@@ -130,11 +133,11 @@ pub fn status(config: &RuntimeConfig) -> Result<SyncStatusView, CliSdkAdapterErr
}
pub fn pull(config: &RuntimeConfig) -> Result<SyncActionView, RuntimeError> {
- pull_with_fetcher(config, fetch_events_from_relays_windowed)
+ pull_with_fetcher(config, shared_relay_transport_fetch_windowed)
}
pub fn market_refresh(config: &RuntimeConfig) -> Result<SyncActionView, RuntimeError> {
- market_refresh_with_fetcher(config, fetch_events_from_relays_windowed)
+ market_refresh_with_fetcher(config, shared_relay_transport_fetch_windowed)
}
fn pull_with_fetcher<F>(config: &RuntimeConfig, fetcher: F) -> Result<SyncActionView, RuntimeError>
@@ -142,7 +145,7 @@ where
F: FnOnce(
&[String],
RadrootsNostrFilter,
- ) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError>,
+ ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>,
{
relay_ingest(config, RelayIngestScope::SyncPull, fetcher)
}
@@ -155,25 +158,30 @@ where
F: FnOnce(
&[String],
RadrootsNostrFilter,
- ) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError>,
+ ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>,
{
relay_ingest(config, RelayIngestScope::MarketRefresh, fetcher)
}
-fn fetch_events_from_relays_windowed(
+fn shared_relay_transport_fetch_windowed(
relay_urls: &[String],
base_filter: RadrootsNostrFilter,
-) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError> {
+) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> {
let mut next_filter = base_filter.clone();
- let mut merged: Option<DirectRelayFetchReceipt> = None;
+ let mut merged: Option<RadrootsRelayFetchedEventsReceipt> = None;
for _ in 0..RELAY_FETCH_MAX_PAGES {
- let receipt = fetch_events_from_relays(relay_urls, next_filter)?;
+ let receipt = fetch_relay_events_via_shared_transport(
+ relay_urls,
+ unix_now_ms(),
+ RELAY_FETCH_LIMIT,
+ next_filter,
+ )?;
let page_len = receipt.events.len();
let oldest_created_at = receipt
.events
.iter()
- .map(|event| event.created_at.as_secs())
+ .map(|event| event.event.created_at.as_secs())
.min();
merge_fetch_receipt(&mut merged, receipt);
if page_len < RELAY_FETCH_LIMIT {
@@ -191,7 +199,7 @@ fn fetch_events_from_relays_windowed(
.limit(RELAY_FETCH_LIMIT);
}
- merged.ok_or(DirectRelayFetchError::MissingRelays)
+ merged.ok_or(RadrootsRelayTransportError::EmptyTargetSet)
}
fn relay_ingest<F>(
@@ -203,7 +211,7 @@ where
F: FnOnce(
&[String],
RadrootsNostrFilter,
- ) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError>,
+ ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>,
{
let snapshot = inspect_sync(config)?;
if snapshot.state == "unconfigured" {
@@ -229,14 +237,11 @@ where
let started_at = unix_now();
let receipt = match fetcher(&config.relay.urls, scope.filter()) {
- Ok(receipt) => receipt,
- Err(DirectRelayFetchError::Connect {
- reason,
- target_relays,
- failed_relays,
- }) => {
- let failed_relays = relay_failures(failed_relays);
- let failure_reason = format!("direct relay connection failed: {reason}");
+ Ok(receipt) if receipt.connected_relays.is_empty() && !receipt.failed_relays.is_empty() => {
+ let target_relays = receipt.target_relays;
+ let failed_relays = relay_failures(receipt.failed_relays);
+ let reason = relay_failure_reason(&failed_relays);
+ let failure_reason = format!("relay transport fetch failed: {reason}");
let executor = SqliteExecutor::open(&config.local.replica_db_path)?;
migrations::run_all_up(&executor)?;
record_sync_run(
@@ -259,6 +264,7 @@ where
view.freshness = freshness_for_scope_from_executor(config, &executor, scope)?;
return Ok(view);
}
+ Ok(receipt) => receipt,
Err(error) => {
let failure_reason = error.to_string();
let executor = SqliteExecutor::open(&config.local.replica_db_path)?;
@@ -1100,7 +1106,7 @@ fn sync_record_from_failure(
fn sync_record_from_ingest(
scope: RelayIngestScope,
relays: &[String],
- receipt: &DirectRelayFetchReceipt,
+ receipt: &RadrootsRelayFetchedEventsReceipt,
ingest: &RelayIngestCounts,
started_at: u64,
) -> Result<SyncRunRecord, RuntimeError> {
@@ -1305,7 +1311,7 @@ impl RelayIngestScope {
fn ingest_events(
executor: &SqliteExecutor,
- receipt: &DirectRelayFetchReceipt,
+ receipt: &RadrootsRelayFetchedEventsReceipt,
scope: RelayIngestScope,
) -> Result<RelayIngestCounts, RuntimeError> {
let mut counts = RelayIngestCounts {
@@ -1314,11 +1320,11 @@ fn ingest_events(
};
for event in &receipt.events {
- if !scope.supports_kind(event_kind(event)) {
+ if !scope.supports_kind(event_kind(&event.event)) {
counts.unsupported_count += 1;
continue;
}
- let event = radroots_event_from_nostr(event);
+ let event = radroots_event_from_nostr(&event.event);
match radroots_replica_ingest_event(executor, &event) {
Ok(RadrootsReplicaIngestOutcome::Applied) => counts.ingested_count += 1,
Ok(RadrootsReplicaIngestOutcome::Skipped) => counts.skipped_count += 1,
@@ -1339,19 +1345,19 @@ fn event_kind(event: &radroots_nostr::prelude::RadrootsNostrEvent) -> u32 {
u32::from(event.kind.as_u16())
}
-fn relay_failures(failures: Vec<DirectRelayFailure>) -> Vec<RelayFailureView> {
+fn relay_failures(failures: Vec<RadrootsRelayFetchFailure>) -> Vec<RelayFailureView> {
failures
.into_iter()
.map(|failure| RelayFailureView {
- relay: failure.relay,
+ relay: failure.relay_url,
reason: failure.reason,
})
.collect()
}
fn merge_fetch_receipt(
- target: &mut Option<DirectRelayFetchReceipt>,
- receipt: DirectRelayFetchReceipt,
+ target: &mut Option<RadrootsRelayFetchedEventsReceipt>,
+ receipt: RadrootsRelayFetchedEventsReceipt,
) {
match target {
Some(target) => {
@@ -1364,12 +1370,20 @@ fn merge_fetch_receipt(
if !target
.failed_relays
.iter()
- .any(|existing| existing.relay == failure.relay)
+ .any(|existing| existing.relay_url == failure.relay_url)
{
target.failed_relays.push(failure);
}
}
target.events.extend(receipt.events);
+ target.event_receipts.extend(receipt.event_receipts);
+ target.malformed_count += receipt.malformed_count;
+ target.out_of_filter_count += receipt.out_of_filter_count;
+ target.skipped_over_limit_count += receipt.skipped_over_limit_count;
+ target.eose_count += receipt.eose_count;
+ target.closed_count += receipt.closed_count;
+ target.notice_count += receipt.notice_count;
+ target.relay_outcomes.extend(receipt.relay_outcomes);
}
None => *target = Some(receipt),
}
@@ -1390,6 +1404,13 @@ fn unix_now() -> u64 {
.unwrap_or(0)
}
+fn unix_now_ms() -> i64 {
+ SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .map(|duration| i64::try_from(duration.as_millis()).unwrap_or(i64::MAX))
+ .unwrap_or(0)
+}
+
fn relative_age(age_seconds: u64) -> String {
match age_seconds {
0 => "now".to_owned(),
@@ -1420,6 +1441,10 @@ mod tests {
use radroots_nostr::prelude::{
RadrootsNostrEvent, RadrootsNostrFilter, RadrootsNostrTimestamp, radroots_nostr_build_event,
};
+ use radroots_relay_transport::{
+ RadrootsRelayFetchFailure, RadrootsRelayFetchedEvent, RadrootsRelayFetchedEventsReceipt,
+ RadrootsRelayTransportError,
+ };
use radroots_sdk::{
PushOutboxEventReceipt, PushOutboxEventState, PushOutboxReceipt,
PushOutboxRelayOutcomeKind, PushOutboxRelayReceipt, SyncEventStoreStatus, SyncOutboxStatus,
@@ -1429,8 +1454,7 @@ mod tests {
use tempfile::tempdir;
use super::{
- DirectRelayFailure, DirectRelayFetchError, DirectRelayFetchReceipt, RelayIngestScope,
- freshness_for_scope, market_refresh_with_fetcher, pull_with_fetcher,
+ RelayIngestScope, freshness_for_scope, market_refresh_with_fetcher, pull_with_fetcher,
relay_provenance_relays_for_scope, sdk_push_dry_run_view, sdk_push_view,
sdk_sync_status_view,
};
@@ -1926,15 +1950,14 @@ mod tests {
let seller = identity(13);
let view = pull_with_fetcher(&config, |relays, _| {
- Ok(DirectRelayFetchReceipt {
- target_relays: relays.to_vec(),
- connected_relays: vec![relays[0].clone()],
- failed_relays: vec![DirectRelayFailure {
- relay: relays[1].clone(),
+ Ok(relay_fetch_receipt_with_failed(
+ relays,
+ vec![listing_event(&seller)],
+ vec![RadrootsRelayFetchFailure {
+ relay_url: relays[1].clone(),
reason: "connection refused".to_owned(),
}],
- events: vec![listing_event(&seller)],
- })
+ ))
})
.expect("sync pull partial relay fetch");
@@ -2054,14 +2077,53 @@ mod tests {
) -> impl FnOnce(
&[String],
RadrootsNostrFilter,
- ) -> Result<DirectRelayFetchReceipt, DirectRelayFetchError> {
- move |relays, _| {
- Ok(DirectRelayFetchReceipt {
- target_relays: relays.to_vec(),
- connected_relays: relays.to_vec(),
- failed_relays: Vec::new(),
- events,
+ ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> {
+ move |relays, _| Ok(relay_fetch_receipt(relays, events))
+ }
+
+ fn relay_fetch_receipt(
+ relays: &[String],
+ events: Vec<RadrootsNostrEvent>,
+ ) -> RadrootsRelayFetchedEventsReceipt {
+ relay_fetch_receipt_with_failed(relays, events, Vec::new())
+ }
+
+ fn relay_fetch_receipt_with_failed(
+ relays: &[String],
+ events: Vec<RadrootsNostrEvent>,
+ failed_relays: Vec<RadrootsRelayFetchFailure>,
+ ) -> RadrootsRelayFetchedEventsReceipt {
+ let connected_relays = relays
+ .iter()
+ .filter(|relay| {
+ failed_relays
+ .iter()
+ .all(|failure| failure.relay_url.as_str() != relay.as_str())
})
+ .cloned()
+ .collect::<Vec<_>>();
+ let closed_count = failed_relays.len();
+ RadrootsRelayFetchedEventsReceipt {
+ target_relays: relays.to_vec(),
+ connected_relays: connected_relays.clone(),
+ failed_relays,
+ events: events
+ .into_iter()
+ .map(|event| RadrootsRelayFetchedEvent {
+ relay_url: connected_relays.first().cloned().unwrap_or_default(),
+ event,
+ raw_json: String::new(),
+ observed_at_ms: 0,
+ })
+ .collect(),
+ event_receipts: Vec::new(),
+ malformed_count: 0,
+ out_of_filter_count: 0,
+ skipped_over_limit_count: 0,
+ eose_count: connected_relays.len(),
+ closed_count,
+ notice_count: 0,
+ relay_outcomes: Vec::new(),
}
}
diff --git a/tests/signer_runtime_modes.rs b/tests/signer_runtime_modes.rs
@@ -2362,7 +2362,7 @@ fn local_seller_publish_commands_attempt_configured_relay() {
}
#[test]
-fn local_order_event_list_attempts_configured_direct_relay() {
+fn local_order_event_list_attempts_configured_shared_relay_transport() {
let sandbox = RadrootsCliSandbox::new();
sandbox.json_success(&["--format", "json", "account", "create"]);
let relay = "ws://127.0.0.1:9";
@@ -2372,7 +2372,7 @@ fn local_order_event_list_attempts_configured_direct_relay() {
]);
assert!(!output.status.success());
- assert_direct_relay_connection_failure(&value, "trade.event.list", &["trade", "event", "list"]);
+ assert_relay_transport_fetch_failure(&value, "trade.event.list", &["trade", "event", "list"]);
assert_eq!(value["errors"][0]["detail"]["state"], "unavailable");
assert_eq!(value["errors"][0]["detail"]["target_relays"][0], relay);
assert_eq!(
@@ -2715,7 +2715,7 @@ fn configure_myc_mode(sandbox: &RadrootsCliSandbox, executable: &Path) {
));
}
-fn assert_direct_relay_connection_failure(
+fn assert_relay_transport_fetch_failure(
value: &serde_json::Value,
operation_id: &str,
args: &[&str],
@@ -2727,7 +2727,7 @@ fn assert_direct_relay_connection_failure(
assert_eq!(value["errors"][0]["detail"]["class"], "network");
assert_contains(
&value["errors"][0]["message"],
- "direct relay connection failed",
+ "relay transport fetch failed",
);
assert_no_removed_command_reference(value, args);
assert_no_daemon_runtime_reference(value, args);
diff --git a/tests/target_cli.rs b/tests/target_cli.rs
@@ -113,7 +113,7 @@ impl RadrootsdProxyJsonRpcServer {
.expect("radrootsd proxy nonblocking");
let endpoint = format!("http://{}", listener.local_addr().expect("proxy addr"));
let handle = thread::spawn(move || {
- let deadline = Instant::now() + Duration::from_secs(10);
+ let deadline = Instant::now() + Duration::from_secs(30);
loop {
match listener.accept() {
Ok((stream, _)) => {