cli

Command-line interface for Radroots
git clone https://radroots.dev/git/cli.git
Log | Files | Refs | README | LICENSE

commit 220f6ac80ccb63ca467827367fe8ea6e9cba131e
parent 63e3e1bedf1482f5b4dd3ed256e803c933b70be8
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 15:12:15 +0000

cli: use shared sync engine

- Compose canonical storage, transport, signer, and sync capabilities in the CLI host.
- Replace CLI pagination, ingest, projection, and outbox orchestration with shared SDK operations.
- Map native sync status, pull, delivery, partial failure, and cancellation receipts to terminal views.
- Verify formatting, all targets, and the complete CLI test suite through extbuild.

Diffstat:
MCargo.lock | 3+++
MCargo.toml | 2++
Msrc/ops/error.rs | 46+++++++++++++---------------------------------
Msrc/ops/exec/core.rs | 2+-
Msrc/ops/exec/market.rs | 13+++++++++----
Msrc/ops/exec/runtime.rs | 17+++++++----------
Msrc/runtime/config.rs | 32++++++++++++++++++++++----------
Msrc/runtime/listing.rs | 16+++++-----------
Msrc/runtime/provider.rs | 2+-
Msrc/runtime/sdk.rs | 2271++++++++-----------------------------------------------------------------------
Msrc/runtime/sync.rs | 3606+++++++++++--------------------------------------------------------------------
Msrc/runtime/trade.rs | 2+-
Msrc/runtime/transport.rs | 200++++++++++++++++++++++++-------------------------------------------------------
Msrc/view/runtime.rs | 2+-
14 files changed, 827 insertions(+), 5387 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -2143,6 +2143,8 @@ dependencies = [ "radroots_signing", "radroots_sql_core", "radroots_storage", + "radroots_storage_sqlite", + "radroots_sync", "radroots_trade", "radroots_transport", "radroots_transport_nostr", @@ -2541,6 +2543,7 @@ dependencies = [ "radroots_storage", "radroots_trade", "radroots_transport", + "serde", "sha2", ] diff --git a/Cargo.toml b/Cargo.toml @@ -54,6 +54,8 @@ radroots_secret_vault = { version = "=0.1.0-alpha", features = ["std", "os-keyri radroots_signing = { version = "=0.1.0-alpha", features = ["std"] } radroots_sql_core = { version = "=0.1.0-alpha", features = ["native"] } radroots_storage = "=0.1.0-alpha" +radroots_storage_sqlite = "=0.1.0-alpha" +radroots_sync = "=0.1.0-alpha" radroots_trade = "=0.1.0-alpha" radroots_protocol = "=0.1.0-alpha" serde = { version = "1.0", features = ["derive"] } diff --git a/src/ops/error.rs b/src/ops/error.rs @@ -353,6 +353,13 @@ impl OperationAdapterError { match error { CliSdkAdapterError::Runtime(error) => Self::runtime_failure(operation_id, error), CliSdkAdapterError::Sdk(error) => Self::sdk_failure(operation_id, error), + CliSdkAdapterError::Sync(error) => Self::Runtime(error.to_string()), + CliSdkAdapterError::Storage(error) => Self::Runtime(error.to_string()), + CliSdkAdapterError::Transport(error) => Self::NetworkUnavailable { + operation_id: operation_id.to_owned(), + message: error.to_string(), + }, + CliSdkAdapterError::Io(error) => Self::Runtime(error.to_string()), } } @@ -957,42 +964,15 @@ mod tests { use super::*; #[test] - fn sdk_storage_error_maps_to_typed_output_without_string_classification() { - let error = OperationAdapterError::sdk_failure( - "store.inspect", - RadrootsSdkError::EventStore { - message: "database is locked".to_owned(), - }, - ); + fn final_sdk_error_maps_from_its_protocol_report() { + let sdk_error = radroots_sdk::ClientBuilder::new() + .build() + .expect_err("missing storage must fail"); + let error = OperationAdapterError::sdk_failure("store.inspect", sdk_error); let output = error.to_output_error(); - - assert_eq!(output.code, "event_store"); - assert_eq!(output.exit_code, CliExitCode::RuntimeUnavailable.code()); let detail = output.detail.expect("detail"); assert_eq!(detail["operation_id"], "store.inspect"); - assert_eq!(detail["class"], "storage"); - assert_eq!(detail["retryable"], true); - assert_eq!(detail["detail"]["message"], "database is locked"); - assert_eq!(detail["actions"], json!(["radroots store inspect"])); - } - - #[test] - fn sdk_request_error_maps_recovery_to_operation_retry_action() { - let error = OperationAdapterError::sdk_failure( - "listing.publish", - RadrootsSdkError::InvalidRequest { - message: "idempotency key must not contain boundary whitespace".to_owned(), - }, - ); - - let output = error.to_output_error(); - - assert_eq!(output.code, "invalid_request"); - assert_eq!(output.exit_code, CliExitCode::InvalidInput.code()); - let detail = output.detail.expect("detail"); - assert_eq!(detail["class"], "request"); - assert_eq!(detail["retryable"], false); - assert_eq!(detail["actions"], json!(["radroots listing publish"])); + assert!(detail["class"].is_string()); } } diff --git a/src/ops/exec/core.rs b/src/ops/exec/core.rs @@ -1,6 +1,6 @@ use std::path::PathBuf; -use radroots_transport::RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE; +use crate::runtime::config::RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE; use serde::Serialize; use serde_json::{Value, json}; diff --git a/src/ops/exec/market.rs b/src/ops/exec/market.rs @@ -276,7 +276,10 @@ mod tests { assert_eq!(envelope.operation_id, "market.pull"); assert_eq!(envelope.result["state"], "unconfigured"); assert_eq!(envelope.result["direction"], "pull"); - assert_eq!(envelope.result["actions"][0], "radroots store inspect"); + assert_eq!( + envelope.result["actions"][0], + "radroots transport config update --kind nostr --nostr-relay wss://relay.example.com" + ); } #[test] @@ -299,7 +302,7 @@ mod tests { assert_eq!(envelope.operation_id, "market.pull"); assert!(envelope.dry_run); assert_eq!(envelope.result["state"], "unconfigured"); - assert_eq!(envelope.result["replica_store"], "missing"); + assert_eq!(envelope.result["replica_store"], "canonical"); assert_eq!(envelope.result["direction"], "pull"); } @@ -308,7 +311,9 @@ mod tests { let dir = tempdir().expect("tempdir"); let mut config = sample_config(dir.path()); config.output.dry_run = true; - config.transport.nostr_relay_urls = vec!["wss://relay.example.com".to_owned()]; + config.transport = crate::runtime::config::TransportConfig::from_nostr_relay_urls(vec![ + "wss://relay.example.com".to_owned(), + ]); crate::runtime::store::init(&config).expect("store init"); let service = OperationAdapter::new(MarketOperationService::new(&config)); @@ -325,7 +330,7 @@ mod tests { .expect("market refresh envelope"); assert_eq!(envelope.operation_id, "market.pull"); - assert_eq!(envelope.result["state"], "ready"); + assert_eq!(envelope.result["state"], "dry_run"); assert_eq!( envelope.result["target_transport_endpoints"][0], "wss://relay.example.com" diff --git a/src/ops/exec/runtime.rs b/src/ops/exec/runtime.rs @@ -341,18 +341,15 @@ mod tests { .expect("sync status envelope"); assert_eq!(envelope.operation_id, "sync.status"); - assert_eq!(envelope.result["state"], "ready"); - assert_eq!( - envelope.result["source"], - "SDK canonical event store and outbox" - ); - assert_eq!( - envelope.result["replica_store"], - "derived_projection_not_checked" - ); + assert_eq!(envelope.result["state"], "degraded"); + assert_eq!(envelope.result["source"], "canonical SDK sync engine"); + assert_eq!(envelope.result["replica_store"], "canonical"); assert_eq!(envelope.result["queue"]["pending_count"], 0); assert_eq!(envelope.result["queue"]["total_count"], 0); - assert_eq!(envelope.result["actions"][0], "radroots sync pull"); + assert_eq!( + envelope.result["actions"][0], + "radroots transport config update --kind nostr --nostr-relay wss://relay.example.com" + ); } fn sample_config(root: &Path, relays: Vec<String>) -> RuntimeConfig { diff --git a/src/runtime/config.rs b/src/runtime/config.rs @@ -7,10 +7,8 @@ use std::path::PathBuf; use radroots_runtime::{parse_bool_value, parse_strict_env_file, parse_u64_value}; use radroots_runtime_paths::RadrootsPathResolver; use radroots_secret_vault::{RadrootsHostVaultPolicy, RadrootsSecretBackend}; -use radroots_transport::{RADROOTS_RETICULUM_SCOPE_ID, RadrootsTransportMeshScopeId}; -use radroots_transport_nostr::{ - RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, -}; +use radroots_transport::RadrootsTransportMeshScopeId; +use radroots_transport_nostr::{Error as RelayTransportError, RelayUrl, RelayUrlPolicy}; use serde::Deserialize; use url::Url; @@ -20,6 +18,9 @@ pub use crate::runtime::paths::PathsConfig; use crate::runtime::paths::{ENV_CLI_PATHS_PROFILE, ENV_CLI_PATHS_REPO_LOCAL_ROOT, resolve_paths}; const DEFAULT_LOG_FILTER: &str = "info"; +pub(crate) const RADROOTS_RETICULUM_SCOPE_ID: &str = "local"; +pub(crate) const RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE: &str = + "Reticulum transport is not available in this release"; const DEFAULT_ENV_PATH: &str = ".env"; const DEFAULT_LOCAL_STATE_DIR: &str = "replica"; const DEFAULT_LOCAL_DB_FILE: &str = "replica.sqlite"; @@ -1956,9 +1957,20 @@ fn validate_relay_url(value: &str, source: &str) -> Result<String, RuntimeError> "{source} contains an empty relay url" ))); } - RadrootsRelayUrl::parse(trimmed, nostr_relay_url_policy_for_url(trimmed)).map(|relay| relay.into_string()).map_err(|error| match error { - RadrootsRelayTransportError::UnsupportedRelayScheme { .. } - | RadrootsRelayTransportError::WsRequiresLocalhostPolicy { .. } => { + let relay = if trimmed.starts_with("ws://") { + radroots_transport::Target::nostr_relay(trimmed) + .map(|target| target.uri().as_str().to_owned()) + .map_err(|error| RelayTransportError::InvalidRelayUrl { + url: trimmed.to_owned(), + reason: error.to_string(), + }) + } else { + RelayUrl::parse(trimmed, nostr_relay_url_policy_for_url(trimmed)) + .map(|relay| relay.as_str().to_owned()) + }; + relay.map_err(|error| match error { + RelayTransportError::RelaySchemeDenied { .. } + | RelayTransportError::RelayDestinationDenied { .. } => { RuntimeError::Config(format!("{source} must use websocket relay urls allowed by the active Nostr policy, got `{trimmed}`")) } other => RuntimeError::Config(format!( @@ -1967,11 +1979,11 @@ fn validate_relay_url(value: &str, source: &str) -> Result<String, RuntimeError> }) } -pub(crate) fn nostr_relay_url_policy_for_url(value: &str) -> RadrootsRelayUrlPolicy { +pub(crate) fn nostr_relay_url_policy_for_url(value: &str) -> RelayUrlPolicy { if value.trim_start().starts_with("ws://") { - RadrootsRelayUrlPolicy::Localhost + RelayUrlPolicy::Local } else { - RadrootsRelayUrlPolicy::Public + RelayUrlPolicy::Public } } diff --git a/src/runtime/listing.rs b/src/runtime/listing.rs @@ -2382,18 +2382,12 @@ fn build_listing_discounts( )); } }; - let discount = Discount { - scope: DiscountScope::Bin, - threshold: DiscountThreshold::BinCount { bin_id, min }, + let discount = Discount::try_new( + DiscountScope::Bin, + DiscountThreshold::BinCount { bin_id, min }, value, - }; - if !discount.is_non_negative() { - return Err(issue_for_field( - contents, - field_prefix.as_str(), - "discount value must not be negative", - )); - } + ) + .map_err(|error| issue_for_field(contents, field_prefix.as_str(), error.to_string()))?; discounts.push(discount); } Ok((!discounts.is_empty()).then_some(discounts)) diff --git a/src/runtime/provider.rs b/src/runtime/provider.rs @@ -1,3 +1,4 @@ +use crate::runtime::config::RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE; #[cfg(test)] use crate::runtime::config::{ CapabilityBindingInspection, CapabilityBindingInspectionState, INFERENCE_HYF_STDIO_CAPABILITY, @@ -6,7 +7,6 @@ use crate::runtime::config::{RuntimeConfig, TransportProfileKind}; #[cfg(test)] use crate::runtime::hyf; use crate::view::runtime::PublishRuntimeView; -use radroots_transport::RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE; #[cfg(test)] const WRITE_PLANE_TARGET_DETAIL: &str = diff --git a/src/runtime/sdk.rs b/src/runtime/sdk.rs @@ -1,31 +1,39 @@ -use std::fs; -use std::future::Future; -use std::path::PathBuf; -use std::time::{SystemTime, UNIX_EPOCH}; +//! CLI-owned composition of the final SDK capability graph. + +use std::{ + fs, + future::Future, + path::PathBuf, + sync::Arc, + time::{SystemTime, UNIX_EPOCH}, +}; -use radroots_sdk::{ - Client, ClientBuilder, Error as SdkError, MeshScopeId, MultiTargetProfile, NostrProfile, - NostrRelayUrlPolicy, PushOutboxTargetOutcomeKind, PushOutboxTransportOutcomeKind, - RadrootsClient, RadrootsClientBuilder, RadrootsSdkStorageConfig, RadrootsdExecutionProfile, - ReticulumAgentEndpoint, ReticulumBehavior as SdkReticulumBehavior, ReticulumProfile, - TargetPolicy, TransportProfile, +use radroots_sdk::{Client, ClientBuilder, Error as SdkError}; +use radroots_signing::{ + SignReceipt, SignRequest, Signer, SignerStatus, signer::BoxFuture as SigningFuture, }; -use radroots_transport_nostr::{ - RadrootsNostrClientFetchAdapter, RadrootsRelayFetchRequest, RadrootsRelayFetchedEventsReceipt, - RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, - fetch_relay_events_blocking, +use radroots_storage::{event::SourceGeneration, memory::MemoryStorage}; +use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage}; +use radroots_sync::policy::{Clock, DeadlinePolicy, IdSource, OperationKind, SyncId, SyncStorage}; +use radroots_transport::{ + BoxFuture, DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, + FetchPage, FetchRequest, SinkStatus, SourceStatus, TransportId, + capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, }; +use radroots_transport_nostr::{Config as NostrConfig, NostrTransport, RelayUrlPolicy}; use tokio::runtime::{Builder as TokioRuntimeBuilder, Runtime}; -use crate::runtime::RuntimeError; -use crate::runtime::account; -use crate::runtime::config::{ - ReticulumBehavior, RuntimeConfig, TransportProfileKind, nostr_relay_url_policy_for_url, +use crate::runtime::{ + RuntimeError, + config::{RuntimeConfig, TransportProfileKind}, + signing, }; -use crate::runtime::signing; const SDK_STORAGE_DIR_NAME: &str = "sdk"; -const CLI_RELAY_FETCH_TIMEOUT_MS: u64 = 10_000; +const SYNC_TIMEOUT_MS: u64 = 30_000; +const LOCAL_TRANSPORT_MESSAGE: &str = + "local-only transport is configured without network fetch or delivery"; + pub(crate) use signing::{MYC_NIP46_SESSION_SECRET_SERVICE, myc_managed_account_ref_matches}; #[derive(Debug, thiserror::Error)] @@ -34,76 +42,34 @@ pub enum CliSdkAdapterError { Runtime(#[from] RuntimeError), #[error("{0}")] Sdk(#[from] SdkError), -} - -pub fn sdk_transport_outcome_kind_label(kind: PushOutboxTransportOutcomeKind) -> String { - kind.as_str().to_owned() -} - -pub fn sdk_target_outcome_kind_label(kind: PushOutboxTargetOutcomeKind) -> String { - kind.as_str().to_owned() + #[error("{0}")] + Sync(#[from] radroots_sync::Error), + #[error("{0}")] + Storage(#[from] radroots_storage_sqlite::Error), + #[error("{0}")] + Transport(#[from] radroots_transport_nostr::Error), + #[error("{0}")] + Io(#[from] std::io::Error), } #[derive(Debug, Clone, PartialEq, Eq)] pub struct CliSdkConfig { pub storage_root: PathBuf, - pub geonames_cache_root: PathBuf, - pub nostr_relay_url_policy: NostrRelayUrlPolicy, - pub transport_profile: TransportProfile, - pub radrootsd_execution_profile: Option<RadrootsdExecutionProfile>, } impl CliSdkConfig { - pub fn from_runtime_config(config: &RuntimeConfig) -> Result<Self, RuntimeError> { - Ok(Self { - storage_root: sdk_storage_root(config), - geonames_cache_root: config.paths.shared_cache_root.clone(), - nostr_relay_url_policy: sdk_nostr_relay_url_policy(config), - transport_profile: sdk_transport_profile(config)?, - radrootsd_execution_profile: sdk_radrootsd_execution_profile(config)?, - }) - } - - pub fn from_runtime_config_for_storage_status(config: &RuntimeConfig) -> Self { + pub fn from_runtime_config(config: &RuntimeConfig) -> Self { Self { storage_root: sdk_storage_root(config), - geonames_cache_root: config.paths.shared_cache_root.clone(), - nostr_relay_url_policy: sdk_nostr_relay_url_policy(config), - transport_profile: TransportProfile::local_only(), - radrootsd_execution_profile: None, } } - fn sqlite_options(&self) -> Result<radroots_sdk::storage::SqliteOptions, RuntimeError> { + fn sqlite_options(&self) -> Result<OpenOptions, CliSdkAdapterError> { fs::create_dir_all(&self.storage_root)?; - let paths = radroots_sdk::storage::SqlitePaths::from_directory(&self.storage_root) - .map_err(|error| RuntimeError::Config(format!("invalid SDK storage paths: {error}")))?; - let mut options = radroots_sdk::storage::SqliteOptions::new( - paths, - radroots_sdk::storage::SqliteOpenMode::Create, - ); + let paths = Paths::from_directory(&self.storage_root)?; + let mut options = OpenOptions::new(paths, OpenMode::Create); if !self.storage_root.join("runtime.sqlite").exists() { - let mut bytes = [0_u8; 32]; - getrandom::getrandom(&mut bytes).map_err(|error| { - RuntimeError::Config(format!("failed to generate SDK source identity: {error}")) - })?; - let generation = - radroots_storage::event::SourceGeneration::new(bytes).map_err(|error| { - RuntimeError::Config(format!("invalid SDK source identity: {error}")) - })?; - let created_at_unix_ms = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map_err(|error| RuntimeError::Config(format!("system clock error: {error}")))? - .as_millis() - .try_into() - .map_err(|_| { - RuntimeError::Config("system clock is outside SDK range".to_owned()) - })?; - options = options - .with_source_generation(generation, created_at_unix_ms) - .map_err(|error| { - RuntimeError::Config(format!("invalid SDK source identity: {error}")) - })?; + options = options.with_source_generation(new_source_generation()?, now_unix_ms()?)?; } Ok(options) } @@ -117,38 +83,15 @@ pub struct CliSdkSession { impl CliSdkSession { pub fn connect(config: &RuntimeConfig) -> Result<Self, CliSdkAdapterError> { - let sdk_config = CliSdkConfig::from_runtime_config(config)?; - let runtime = sdk_runtime()?; - let options = sdk_config.sqlite_options()?; - let sdk = runtime.block_on(ClientBuilder::sqlite(options))?.build()?; - Ok(Self { - runtime, - sdk, - config: sdk_config, - }) + Self::connect_inner(config, None, false) } pub fn connect_storage_status(config: &RuntimeConfig) -> Result<Self, CliSdkAdapterError> { - let sdk_config = CliSdkConfig::from_runtime_config_for_storage_status(config); - let runtime = sdk_runtime()?; - let options = sdk_config.sqlite_options()?; - let sdk = runtime.block_on(ClientBuilder::sqlite(options))?.build()?; - Ok(Self { - runtime, - sdk, - config: sdk_config, - }) + Self::connect(config) } pub fn connect_memory(config: &RuntimeConfig) -> Result<Self, CliSdkAdapterError> { - let sdk_config = CliSdkConfig::from_runtime_config(config)?; - let runtime = sdk_runtime()?; - let sdk = ClientBuilder::memory_default().build()?; - Ok(Self { - runtime, - sdk, - config: sdk_config, - }) + Self::connect_inner(config, None, true) } pub fn connect_for_actor( @@ -157,25 +100,14 @@ impl CliSdkSession { actor_pubkey: &str, actor_label: &str, ) -> Result<Self, CliSdkAdapterError> { - let sdk_config = CliSdkConfig::from_runtime_config(config)?; let runtime = sdk_runtime()?; - let signer_provider = runtime.block_on(signing::provider_for_actor( + let provider = runtime.block_on(signing::provider_for_actor( config, actor_account_id, actor_pubkey, actor_label, ))?; - let sdk = runtime.block_on( - sdk_config - .builder() - .signer_provider(signer_provider) - .build(), - )?; - Ok(Self { - runtime, - sdk, - config: sdk_config, - }) + Self::compose(config, runtime, Some(provider), false) } pub fn connect_memory_for_actor( @@ -184,19 +116,77 @@ impl CliSdkSession { actor_pubkey: &str, actor_label: &str, ) -> Result<Self, CliSdkAdapterError> { - let sdk_config = CliSdkConfig::from_runtime_config(config)?; let runtime = sdk_runtime()?; - let signer_provider = runtime.block_on(signing::provider_for_actor( + let provider = runtime.block_on(signing::provider_for_actor( config, actor_account_id, actor_pubkey, actor_label, ))?; - let sdk = runtime.block_on( - memory_builder(&sdk_config) - .signer_provider(signer_provider) - .build(), - )?; + Self::compose(config, runtime, Some(provider), true) + } + + fn connect_inner( + config: &RuntimeConfig, + signer: Option<radroots_sdk::signing::Provider>, + memory: bool, + ) -> Result<Self, CliSdkAdapterError> { + Self::compose(config, sdk_runtime()?, signer, memory) + } + + fn compose( + config: &RuntimeConfig, + runtime: Runtime, + signer: Option<radroots_sdk::signing::Provider>, + memory: bool, + ) -> Result<Self, CliSdkAdapterError> { + let sdk_config = CliSdkConfig::from_runtime_config(config); + if memory { + let storage = Arc::new(MemoryStorage::new(new_source_generation()?)); + Self::compose_with_storage(config, runtime, sdk_config, signer, storage) + } else { + let storage = + Arc::new(runtime.block_on(SqliteStorage::open(sdk_config.sqlite_options()?))?); + Self::compose_with_storage(config, runtime, sdk_config, signer, storage) + } + } + + fn compose_with_storage<T>( + config: &RuntimeConfig, + runtime: Runtime, + sdk_config: CliSdkConfig, + signer: Option<radroots_sdk::signing::Provider>, + storage: Arc<T>, + ) -> Result<Self, CliSdkAdapterError> + where + T: radroots_storage::Storage + SyncStorage + 'static, + { + let (source, sink) = transport_capabilities(config)?; + let signer_capability = signer + .as_ref() + .map(|provider| Arc::new(SharedProvider(provider.clone())) as Arc<dyn Signer>); + let mut engine = radroots_sync::Engine::builder( + storage.clone(), + Arc::new(SystemClock), + Arc::new(RandomIds), + DeadlinePolicy::new(SYNC_TIMEOUT_MS, SYNC_TIMEOUT_MS, SYNC_TIMEOUT_MS)?, + ) + .source(Arc::clone(&source)) + .sink(Arc::clone(&sink)); + if let Some(capability) = signer_capability.as_ref() { + engine = engine.signer(Arc::clone(capability)); + } + let engine = engine.build()?; + + let mut builder = ClientBuilder::new() + .storage(storage) + .source(source) + .sink(sink) + .sync_engine(engine); + if let Some(provider) = signer { + builder = builder.signing(provider); + } + let sdk = builder.build()?; Ok(Self { runtime, sdk, @@ -212,10 +202,7 @@ impl CliSdkSession { &self.config } - pub fn block_on<F>(&self, future: F) -> F::Output - where - F: Future, - { + pub fn block_on<F: Future>(&self, future: F) -> F::Output { self.runtime.block_on(future) } } @@ -242,1968 +229,162 @@ 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 relay_targets = RadrootsRelayTargetSet::from_urls( - relay_urls - .iter() - .map(|url| RadrootsRelayUrl::parse(url, nostr_relay_url_policy_for_url(url))) - .collect::<Result<Vec<_>, _>>()?, - )?; - let request = - RadrootsRelayFetchRequest::fetch(observed_at_ms, max_events, relay_targets, [filter])? - .with_timeout_ms(CLI_RELAY_FETCH_TIMEOUT_MS)?; - fetch_relay_events_blocking(&RadrootsNostrClientFetchAdapter, request) -} - -fn memory_builder(config: &CliSdkConfig) -> RadrootsClientBuilder { - RadrootsClient::builder() - .geonames_cache_root(config.geonames_cache_root.clone()) - .transport_profile(config.transport_profile.clone()) -} - -pub fn sdk_nostr_relay_url_policy(config: &RuntimeConfig) -> NostrRelayUrlPolicy { +pub fn sdk_nostr_relay_url_policy(config: &RuntimeConfig) -> RelayUrlPolicy { if config .transport .nostr_relay_urls .iter() .any(|relay_url| relay_url.starts_with("ws://")) { - NostrRelayUrlPolicy::Localhost + RelayUrlPolicy::Local } else { - NostrRelayUrlPolicy::Public + RelayUrlPolicy::Public } } -pub fn sdk_target_policy(_config: &RuntimeConfig) -> TargetPolicy { - TargetPolicy::default_profile() -} - -fn sdk_transport_profile(config: &RuntimeConfig) -> Result<TransportProfile, RuntimeError> { - match config.transport.profile { - TransportProfileKind::LocalOnly => Ok(TransportProfile::local_only()), - TransportProfileKind::Nostr => { - let profile = NostrProfile::new( - config.transport.nostr_relay_urls.iter().map(String::as_str), - sdk_nostr_relay_url_policy(config), - ) - .map_err(|error| RuntimeError::Config(error.to_string()))?; - Ok(TransportProfile::nostr(profile)) - } - TransportProfileKind::Reticulum => { - Ok(TransportProfile::reticulum(sdk_reticulum_profile(config)?)) - } - TransportProfileKind::MultiTarget => { - let nostr = NostrProfile::new( - config.transport.nostr_relay_urls.iter().map(String::as_str), - sdk_nostr_relay_url_policy(config), - ) - .map_err(|error| RuntimeError::Config(error.to_string()))?; - Ok(TransportProfile::multi_target(MultiTargetProfile::new( - nostr, - sdk_reticulum_profile(config)?, - ))) - } - } -} - -fn sdk_reticulum_profile(config: &RuntimeConfig) -> Result<ReticulumProfile, RuntimeError> { - let behavior = match config.transport.reticulum_behavior { - ReticulumBehavior::RejectDeliveryAttempts => SdkReticulumBehavior::RejectDeliveryAttempts, - ReticulumBehavior::DeferDeliveryPlans => SdkReticulumBehavior::DeferDeliveryPlans, - }; - let scope = MeshScopeId::parse(config.transport.reticulum_scope.as_str()) - .map_err(|error| RuntimeError::Config(error.to_string()))?; - let mut profile = ReticulumProfile::deferred_until_implemented() - .with_behavior(behavior) - .with_scope(scope); - if let Some(agent_endpoint) = config.transport.reticulum_agent_endpoint.as_ref() { - profile = profile.with_agent_endpoint( - ReticulumAgentEndpoint::parse(agent_endpoint.as_str()) - .map_err(|error| RuntimeError::Config(error.to_string()))?, - ); - } - Ok(profile) +pub(crate) fn sync_targets( + config: &RuntimeConfig, +) -> Result<radroots_transport::TargetSet, TransportError> { + let targets = config + .transport + .nostr_relay_urls + .iter() + .map(|relay| radroots_transport::Target::nostr_relay(relay)) + .collect::<Result<Vec<_>, _>>()?; + radroots_transport::TargetSet::new(targets) } -fn sdk_radrootsd_execution_profile( +fn transport_capabilities( config: &RuntimeConfig, -) -> Result<Option<RadrootsdExecutionProfile>, RuntimeError> { - if config.transport.radrootsd_execution.token_file.is_none() - && config - .transport - .radrootsd_execution - .token_secret_id - .is_none() +) -> Result<(Arc<dyn EventSource>, Arc<dyn EventSink>), CliSdkAdapterError> { + if matches!( + config.transport.profile, + TransportProfileKind::Nostr | TransportProfileKind::MultiTarget + ) && !config.transport.nostr_relay_urls.is_empty() { - return Ok(None); + let transport = Arc::new(NostrTransport::new(NostrConfig::new( + sdk_nostr_relay_url_policy(config), + &config.transport.nostr_relay_urls, + )?)); + let source: Arc<dyn EventSource> = transport.clone(); + let sink: Arc<dyn EventSink> = transport; + Ok((source, sink)) + } else { + let transport = Arc::new(UnavailableTransport); + let source: Arc<dyn EventSource> = transport.clone(); + let sink: Arc<dyn EventSink> = transport; + Ok((source, sink)) } - Ok(Some( - RadrootsdExecutionProfile::new(config.transport.radrootsd_execution.url.clone()) - .with_bearer_token(radrootsd_execution_bearer_token(config)?), - )) } -fn radrootsd_execution_bearer_token(config: &RuntimeConfig) -> Result<String, RuntimeError> { - if let Some(path) = config.transport.radrootsd_execution.token_file.as_ref() { - let token = fs::read_to_string(path).map_err(|error| { - RuntimeError::Config(format!( - "failed to read radrootsd execution token file {}: {error}", - path.display() - )) - })?; - return normalize_radrootsd_execution_bearer_token( - token.as_str(), - format!("radrootsd execution token file {}", path.display()).as_str(), - ); - } - - if let Some(secret_id) = config - .transport - .radrootsd_execution - .token_secret_id - .as_ref() - { - let vault = account::account_secret_vault(config)?; - let token = vault.load_secret(secret_id).map_err(|error| { - RuntimeError::Config(format!( - "failed to load radrootsd execution token secret `{secret_id}`: {error}" - )) - })?; - let token = token.ok_or_else(|| { - RuntimeError::Config(format!( - "radrootsd execution token secret `{secret_id}` was not found" - )) - })?; - return normalize_radrootsd_execution_bearer_token( - token.as_str(), - format!("radrootsd execution token secret `{secret_id}`").as_str(), - ); - } - - Err(RuntimeError::Config( - "radrootsd execution requires a configured token file or token secret id".to_owned(), - )) +fn now_unix_ms() -> Result<u64, RuntimeError> { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|error| RuntimeError::Config(format!("system clock error: {error}")))? + .as_millis() + .try_into() + .map_err(|_| RuntimeError::Config("system clock is outside SDK range".to_owned())) } -fn normalize_radrootsd_execution_bearer_token( - raw: &str, - source: &str, -) -> Result<String, RuntimeError> { - let token = raw.trim(); - if token.is_empty() { - return Err(RuntimeError::Config(format!("{source} is empty"))); - } - if token.bytes().any(|byte| byte.is_ascii_control()) { - return Err(RuntimeError::Config(format!( - "{source} contains unsupported control characters" - ))); - } - Ok(token.to_owned()) +fn new_source_generation() -> Result<SourceGeneration, RuntimeError> { + let mut bytes = [0_u8; 32]; + getrandom::getrandom(&mut bytes).map_err(|error| { + RuntimeError::Config(format!("failed to generate SDK source identity: {error}")) + })?; + SourceGeneration::new(bytes) + .map_err(|error| RuntimeError::Config(format!("invalid SDK source identity: {error}"))) } -#[cfg(test)] -mod tests { - use std::collections::BTreeSet; - use std::fs; - use std::path::{Path, PathBuf}; - use std::time::Duration; - - use radroots_authority::RadrootsEventSigner; - use radroots_sdk::{ - PushOutboxTargetOutcomeKind, PushOutboxTransportOutcomeKind, RadrootsdExecutionAuth, - SdkStorageKind, StorageStatusRequest, - }; - use radroots_secret_vault::RadrootsSecretBackend; - use tempfile::tempdir; - - use super::*; - use crate::runtime::config::{ - AccountConfig, AccountSecretContractConfig, HyfConfig, IdentityConfig, InteractionConfig, - LocalConfig, LoggingConfig, MycConfig, OutputConfig, OutputFormat, PathsConfig, - ReticulumBehavior, RhiConfig, RpcConfig, SignerBackend, SignerConfig, Verbosity, - }; - - struct DirectRrRsDependency { - section: &'static str, - name: &'static str, - owner: &'static str, - reason: &'static str, - lifecycle: &'static str, - } - - struct MigratedCliPathGuard { - label: &'static str, - path: &'static str, - start: &'static str, - end: &'static str, - required_tokens: &'static [&'static str], - } - - struct SdkOutcomeLabelHelperGuard { - label: &'static str, - start: &'static str, - end: &'static str, - } - - const DIRECT_RR_RS_DEPENDENCIES: &[DirectRrRsDependency] = &[ - DirectRrRsDependency { - section: "dependencies", - name: "radroots_authority", - owner: "cli-sdk-adapter", - reason: "local account signer materialization for SDK and remaining CLI-authored signing", - lifecycle: "retain until all signed mutation construction moves behind SDK signer requests", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_core", - owner: "cli-drafts-and-rendering", - reason: "CLI draft parsing, numeric validation, and display DTOs", - lifecycle: "retain while CLI owns TOML draft UX and command rendering", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_event", - owner: "cli-drafts-and-non-migrated-workflows", - reason: "event DTOs for local drafts, views, relay reads, and release-product mutation inspection", - lifecycle: "retain until the remaining event-authoring and inspection surfaces migrate", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_event_codec", - owner: "cli-drafts-and-non-migrated-workflows", - reason: "event encoding and decoding for farm, listing draft, order, sync pull, and validation inspection", - lifecycle: "retain until those command families are SDK-backed", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_identity", - owner: "cli-account-and-signer-ux", - reason: "account identity views, local signer materialization, and direct-relay workflows outside the migrated paths", - lifecycle: "retain while CLI owns account selection and local identity custody UX", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_runtime_store", - owner: "cli-app-interop", - reason: "shared local work and signed-event interop with the desktop app", - lifecycle: "retain until a shared runtime-store SDK boundary replaces direct CLI access", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_log", - owner: "cli-runtime-shell", - reason: "CLI logging initialization and file layout", - lifecycle: "permanent CLI runtime ownership", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_mesh", - owner: "cli-mesh-reticulum-admission-policy", - reason: "canonical Reticulum admission policy and structured delivery denial reporting", - lifecycle: "retain while CLI exposes Reticulum availability status and policy inspection", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_nostr", - 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_transport", - owner: "cli-transport-config", - reason: "canonical transport scope identifiers and fail-closed transport profile configuration", - lifecycle: "retain while CLI owns runtime transport config parsing", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_transport_nostr", - owner: "cli-nostr-transport-read-boundary", - reason: "shared fail-closed Nostr 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", - name: "radroots_nostr_connect", - owner: "sdk-myc-nip46-transport", - reason: "CLI Myc signer target parsing and NIP-46 relay transport wiring for SDK signing", - lifecycle: "retain while CLI owns signer backend wiring", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_nostr_accounts", - owner: "cli-account-store", - reason: "CLI account selection, import, local signer status, and account persistence", - lifecycle: "retain while CLI owns local account UX and storage", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_nostr_signer", - owner: "cli-signer-readiness", - reason: "signer readiness reporting for active mutation command surfaces", - lifecycle: "retain until signer readiness is fully SDK-owned", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_replica_store", - owner: "derived-projection-and-market-reads", - reason: "derived projection status, export, market reads, sync pull, basket lookup, and trade draft preflight", - lifecycle: "retain until those derived projection surfaces move behind SDK APIs", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_replica_schema", - owner: "derived-projection-and-market-reads", - reason: "typed query filters for market, basket, and order lookup projections", - lifecycle: "retain until those derived projection surfaces move behind SDK APIs", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_replica_sync", - owner: "sync-pull-and-derived-projection", - reason: "relay ingest, sync pull, market refresh, and derived projection state reporting", - lifecycle: "retain until relay ingest and projection repair move behind SDK APIs", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_runtime", - owner: "cli-config", - reason: "strict environment and config value parsing", - lifecycle: "permanent CLI configuration ownership unless a shared runtime config crate replaces it", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_runtime_paths", - owner: "cli-runtime-paths", - reason: "profile-aware CLI config, data, logs, and secrets path resolution", - lifecycle: "permanent CLI runtime ownership", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_secret_vault", - owner: "cli-account-store", - reason: "local account secret backend selection and readiness", - lifecycle: "retain while CLI owns local account custody UX", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_protected_store", - owner: "cli-account-store", - reason: "protected file secret vault selection for local account and Myc session material", - lifecycle: "retain while CLI owns account and signer session custody UX", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_sql_core", - owner: "derived-projection-and-runtime-store", - reason: "SQLite executor for derived projection and shared runtime-store storage", - lifecycle: "transitional until those storage surfaces move behind SDK or shared runtime APIs", - }, - DirectRrRsDependency { - section: "dependencies", - name: "radroots_trade", - owner: "cli-drafts-and-validation", - reason: "listing draft validation, trade economics, reducer helpers, and release-product mutation parsing", - lifecycle: "retain until remaining trade validation and draft behavior migrates", - }, - DirectRrRsDependency { - section: "dev-dependencies", - name: "radroots_outbox", - owner: "cli-test-fixtures", - reason: "test-only outbox fixture assertions for CLI transport and SDK workflow coverage", - lifecycle: "retain while CLI integration tests assert local outbox side effects directly", - }, - ]; - - const NOSTR_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 SDK_OUTCOME_LABEL_SOURCE_DISALLOWED_TOKENS: &[(&str, &str)] = &[ - ( - concat!("sdk_enum", "_label"), - "serde-backed SDK enum label extraction", - ), - ( - concat!("serde_json::to_value", "(kind)"), - "serde-backed SDK enum label extraction", - ), - ( - concat!("panic!", "(\"SDK enum"), - "production SDK enum label panic", - ), - ( - concat!("fn sdk_target", "_outcome_kind("), - "local SDK target outcome label helper", - ), - ]; - - const SDK_OUTCOME_LABEL_HELPER_GUARDS: &[SdkOutcomeLabelHelperGuard] = &[ - SdkOutcomeLabelHelperGuard { - label: "push outbox transport outcome label helper", - start: "pub fn sdk_transport_outcome_kind_label(", - end: "pub fn sdk_target_outcome_kind_label(", - }, - SdkOutcomeLabelHelperGuard { - label: "push outbox target outcome label helper", - start: "pub fn sdk_target_outcome_kind_label(", - end: "#[derive(Debug, Clone, PartialEq, Eq)]", - }, - ]; - - const MIGRATED_CLI_PATH_GUARDS: &[MigratedCliPathGuard] = &[ - MigratedCliPathGuard { - label: "listing publish", - path: "src/runtime/listing.rs", - start: "pub fn publish_via_sdk(", - end: "fn sdk_listing_publish_input(", - required_tokens: &[ - "session.sdk().listings().prepare_publish", - "session.sdk().listings().enqueue_publish", - "session.sdk().sync().push_outbox", - ], - }, - MigratedCliPathGuard { - label: "farm publish", - path: "src/runtime/farm.rs", - start: "fn publish_via_sdk(", - end: "#[derive(Debug, Clone)]\nstruct SdkFarmPublishInput", - required_tokens: &[ - "prepare_publish(FarmPreparePublishRequest::new", - "enqueue_publish(request)", - "session.sdk().sync().push_outbox", - ], - }, - MigratedCliPathGuard { - label: "sync status", - path: "src/runtime/sync.rs", - start: "pub fn status(config: &RuntimeConfig) -> Result<SyncStatusView, CliSdkAdapterError>", - end: "pub fn pull(", - required_tokens: &["session.sdk().sync().status"], - }, - MigratedCliPathGuard { - label: "sync push", - path: "src/runtime/sync.rs", - start: "pub fn push(config: &RuntimeConfig) -> Result<SyncActionView, CliSdkAdapterError>", - end: "pub fn watch(", - required_tokens: &["session.sdk().sync().push_outbox", "PushOutboxRequest::new"], - }, - MigratedCliPathGuard { - label: "trade release product commands", - path: "src/runtime/trade.rs", - start: "pub fn submit_proposal(", - end: "pub fn get_trade(", - required_tokens: &[ - "SubmitProposalRequest::new", - "ProposeRevisionRequest::new", - "DecideCandidateRequest::new", - "CancelTradeRequest::new", - "ResumeOperationRequest::new", - "session.sdk().trades().commands().submit_proposal", - "session.sdk().trades().commands().propose_revision", - "session.sdk().trades().commands().decide_candidate", - "session.sdk().trades().commands().cancel_trade", - "session.sdk().trades().commands().resume_operation", - ], - }, - MigratedCliPathGuard { - label: "trade release product queries", - path: "src/runtime/trade.rs", - start: "pub fn get_trade(", - end: "pub fn seal_private_artifact(", - required_tokens: &[ - "GetTradeRequest::new", - "ListTradesRequest::new", - "RefreshTradeEvidenceRequest::new", - "InspectEvidenceRequest::new", - "session.sdk().trades().queries().get_trade", - "session.sdk().trades().queries().list_trades", - "session.sdk().trades().queries().refresh_evidence", - "session.sdk().trades().queries().inspect_evidence", - ], - }, - MigratedCliPathGuard { - label: "trade private artifact SDK", - path: "src/runtime/trade.rs", - start: "pub fn seal_private_artifact(", - end: "pub fn scaffold_proposal_draft(", - required_tokens: &[ - "TradePrivateArtifactSealRequest::binding_terms", - "TradePrivateArtifactOpenRequest::new", - "TradePrivateArtifactDeleteRequest::new", - "session.sdk().trades().seal_private_artifact", - "session.sdk().trades().open_private_artifact", - "session.sdk().trades().delete_private_artifact", - ], - }, - MigratedCliPathGuard { - label: "store status", - path: "src/runtime/store.rs", - start: "pub fn status(config: &RuntimeConfig) -> Result<LocalStatusView, CliSdkAdapterError>", - end: "fn derived_projection_status(", - required_tokens: &[ - "session.sdk()", - "storage_status(StorageStatusRequest::new())", - "integrity(IntegrityRequest::new())", - ], - }, - MigratedCliPathGuard { - label: "store backup", - path: "src/runtime/store.rs", - start: "pub fn backup(\n config: &RuntimeConfig", - end: "pub fn backup_preflight(", - required_tokens: &["session.sdk().backup", "BackupRequest"], - }, - MigratedCliPathGuard { - label: "store backup preflight", - path: "src/runtime/store.rs", - start: "pub fn backup_preflight(", - end: "pub fn restore(", - required_tokens: &[ - "storage_status(StorageStatusRequest::new())", - "integrity(IntegrityRequest::new())", - ], - }, - MigratedCliPathGuard { - label: "store restore", - path: "src/runtime/store.rs", - start: "pub fn restore(", - end: "pub fn export(", - required_tokens: &[ - "RestoreRequest::new", - "sdk_runtime()", - "RadrootsClient::restore", - ], - }, - ]; +struct SystemClock; - const MIGRATED_PATH_DISALLOWED_TOKENS: &[&str] = &[ - "fetch_events_from_relays", - "fetch_relay_events_via_shared_transport", - "publish_parts_with_identity", - "publish_via_direct_relay", - "mutate_via_direct_relay", - "radroots_replica_pending_publish", - "radroots_replica_pending_publish_batch", - "radroots_replica_sync_status", - "ReplicaSql::new", - "SqliteExecutor::open(&config.local.replica_store_path)", - "outbox_idempotency_digest", - "canonical_target_transport_endpoints", - "radroots_sdk::protocol::order", - "build_order_request_draft", - "build_order_decision_draft", - "build_order_cancellation_draft", - "parse_order_root_tag", - "parse_order_prev_tag", - "build_transition_proof_request_tags", - "build_transition_proof_result_tags", - "build_job_feedback_tags", - "KIND_TRADE_TRANSITION_PROOF", - "KIND_JOB_FEEDBACK", - "status_client(", - "TradeStatusClient", - "TradeValidationClient", - ]; - - const REMOVED_SDK_ROOT_TRADE_ALIAS_NAMES: &[&str] = &[ - "trade_buyer", - "trade_seller", - "trade_status", - "trade_resync", - "trade_validation", - ]; - - const REMOVED_SDK_STATUS_SURFACE_TOKENS: &[&str] = &[ - "status_client(", - "TradeStatusClient", - "TradeValidationClient", - ]; - - mod removed_surface_fixtures { - pub const ACCOUNT_CREATE_MODE: &str = "AccountCreateMode"; - pub const CONFIG_IDENTITY_PATH_EXISTS: &str = "config.identity.path.exists()"; - pub const CREATE_OR_MIGRATE_DEFAULT_ACCOUNT: &str = "create_or_migrate_default_account"; - pub const MIGRATED_OUTPUT_JSON: &str = "\"migrated\""; - pub const MIGRATE_LEGACY_IDENTITY_FILE: &str = "migrate_legacy_identity_file"; - } - - #[test] - fn maps_runtime_config_to_sdk_builder_inputs() { - let root = tempdir().expect("tempdir"); - let config = sample_config( - root.path(), - vec!["wss://relay.one".to_owned(), "wss://relay.two".to_owned()], - ); - - let sdk_config = CliSdkConfig::from_runtime_config(&config).expect("sdk config"); - - assert_eq!(sdk_config.storage_root, config.local.root.join("sdk")); - assert_eq!( - sdk_config.nostr_relay_url_policy, - NostrRelayUrlPolicy::Public - ); - let TransportProfile::Nostr { profile } = sdk_config.transport_profile else { - panic!("expected Nostr transport profile"); - }; - assert_eq!( - profile.relay_urls(), - vec!["wss://relay.one".to_owned(), "wss://relay.two".to_owned()] - ); - } - - #[test] - fn maps_multi_target_runtime_config_to_sdk_multi_target_profile() { - let root = tempdir().expect("tempdir"); - let mut config = sample_config( - root.path(), - vec!["wss://relay.one".to_owned(), "wss://relay.two".to_owned()], - ); - config.transport.profile = TransportProfileKind::MultiTarget; - config.transport.reticulum_behavior = ReticulumBehavior::DeferDeliveryPlans; - config.transport.reticulum_scope = "farmers_market".to_owned(); - config.transport.reticulum_agent_endpoint = Some("reticulum-agent:local".to_owned()); - - let sdk_config = CliSdkConfig::from_runtime_config(&config).expect("sdk config"); - - let TransportProfile::MultiTarget { profile } = sdk_config.transport_profile else { - panic!("expected multi-target transport profile"); - }; - assert_eq!( - profile.nostr().relay_urls(), - vec!["wss://relay.one".to_owned(), "wss://relay.two".to_owned()] - ); - assert_eq!( - profile.reticulum().behavior().as_str(), - "defer_delivery_plans" - ); - assert_eq!(profile.reticulum().scope().as_str(), "farmers_market"); - assert_eq!( - profile - .reticulum() - .agent_endpoint() - .expect("agent endpoint") - .as_str(), - "reticulum-agent:local" - ); - } - - #[test] - fn maps_radrootsd_execution_token_file_to_sdk_profile_auth() { - let root = tempdir().expect("tempdir"); - let mut config = sample_config(root.path(), Vec::new()); - let token_file = root.path().join("radrootsd_execution.token"); - fs::write(&token_file, "radrootsd-execution-file-token\n").expect("write token file"); - config.transport.radrootsd_execution.url = "http://127.0.0.1:7070".to_owned(); - config.transport.radrootsd_execution.token_file = Some(token_file); - - let sdk_config = CliSdkConfig::from_runtime_config(&config).expect("sdk config"); - - assert!(matches!( - sdk_config.transport_profile, - TransportProfile::LocalOnly - )); - let profile = sdk_config - .radrootsd_execution_profile - .expect("radrootsd execution profile"); - assert_eq!(profile.endpoint_url(), "http://127.0.0.1:7070"); - assert_eq!( - profile.auth(), - &RadrootsdExecutionAuth::BearerToken("radrootsd-execution-file-token".to_owned()) - ); - } - - #[test] - fn maps_radrootsd_execution_token_secret_id_to_sdk_profile_auth() { - let root = tempdir().expect("tempdir"); - let mut config = sample_config(root.path(), Vec::new()); - config.transport.radrootsd_execution.url = "http://127.0.0.1:7070".to_owned(); - config.transport.radrootsd_execution.token_secret_id = - Some("radrootsd_execution_token".to_owned()); - let vault = account::account_secret_vault(&config).expect("account vault"); - vault - .store_secret( - "radrootsd_execution_token", - "radrootsd-execution-secret-token", - ) - .expect("store radrootsd execution token"); - - let sdk_config = CliSdkConfig::from_runtime_config(&config).expect("sdk config"); - - let profile = sdk_config - .radrootsd_execution_profile - .expect("radrootsd execution profile"); - assert_eq!( - profile.auth(), - &RadrootsdExecutionAuth::BearerToken("radrootsd-execution-secret-token".to_owned()) - ); - } - - #[test] - fn radrootsd_execution_profile_requires_materialized_bearer_token() { - let root = tempdir().expect("tempdir"); - let mut config = sample_config(root.path(), Vec::new()); - config.transport.radrootsd_execution.url = "http://127.0.0.1:7070".to_owned(); - - let sdk_config = CliSdkConfig::from_runtime_config(&config).expect("sdk config"); - assert!(sdk_config.radrootsd_execution_profile.is_none()); - } - - #[test] - fn radrootsd_execution_token_resolution_rejects_empty_or_header_unsafe_tokens() { - assert!(matches!( - normalize_radrootsd_execution_bearer_token(" \n", "radrootsd execution token"), - Err(RuntimeError::Config(message)) if message.contains("empty") - )); - assert!(matches!( - normalize_radrootsd_execution_bearer_token( - "radrootsd\nexecution", - "radrootsd execution token" - ), - Err(RuntimeError::Config(message)) if message.contains("control characters") - )); - } - - #[test] - fn maps_localhost_ws_relays_to_localhost_sdk_policy() { - let root = tempdir().expect("tempdir"); - let config = sample_config(root.path(), vec!["ws://127.0.0.1:8080".to_owned()]); - - assert_eq!( - sdk_nostr_relay_url_policy(&config), - NostrRelayUrlPolicy::Localhost - ); - } - - #[test] - fn sdk_session_builds_once_and_runs_async_storage_smoke() { - let root = tempdir().expect("tempdir"); - let config = sample_config(root.path(), Vec::new()); - let session = CliSdkSession::connect(&config).expect("sdk session"); - - let status = session - .block_on(session.sdk().storage_status(StorageStatusRequest::new())) - .expect("storage status"); - - assert_eq!(session.config().storage_root, config.local.root.join("sdk")); - assert_eq!(status.storage, SdkStorageKind::Directory); - assert_eq!(status.event_store.total_events, 0); - assert_eq!(status.event_store.valid_stream_events, 0); - assert_eq!(status.outbox.total_events, 0); - } - - #[test] - fn sdk_sources_do_not_import_cli_types() { - let sdk_src = Path::new(env!("CARGO_MANIFEST_DIR")).join("../sdk/crates/sdk/src"); - let mut files = Vec::new(); - collect_rs_files(sdk_src.as_path(), &mut files); - let forbidden = vec![ - ("radroots_cli".to_owned(), "CLI crate identity"), - ("domains/radroots/cli".to_owned(), "CLI mount path"), - (["approval", "token"].join("_"), "retired approval string"), - ("OutputEnvelope".to_owned(), "CLI output envelope"), - ("next_actions".to_owned(), "CLI next-action rendering"), - ("exit_code".to_owned(), "CLI exit-code contract"), - ("docs/".to_owned(), "repository docs path"), - ("radroots store".to_owned(), "CLI command string"), - ("radroots sync".to_owned(), "CLI command string"), - ("radroots listing".to_owned(), "CLI command string"), - ("radroots trade".to_owned(), "CLI command string"), - ]; - - for file in files { - let source = fs::read_to_string(&file).expect("read sdk source"); - for (needle, description) in forbidden.iter() { - assert!( - !source.contains(needle.as_str()), - "SDK source contains {description} `{needle}` in {}", - file.display() - ); - } - } - } - - #[test] - fn sync_order_transport_label_helpers_cover_public_outcomes() { - for (kind, label) in [ - (PushOutboxTargetOutcomeKind::Accepted, "accepted"), - ( - PushOutboxTargetOutcomeKind::DuplicateAccepted, - "duplicate_accepted", - ), - (PushOutboxTargetOutcomeKind::Blocked, "blocked"), - (PushOutboxTargetOutcomeKind::RateLimited, "rate_limited"), - (PushOutboxTargetOutcomeKind::Invalid, "invalid"), - (PushOutboxTargetOutcomeKind::PowRequired, "pow_required"), - (PushOutboxTargetOutcomeKind::Restricted, "restricted"), - (PushOutboxTargetOutcomeKind::AuthRequired, "auth_required"), - (PushOutboxTargetOutcomeKind::Muted, "muted"), - (PushOutboxTargetOutcomeKind::Unsupported, "unsupported"), - ( - PushOutboxTargetOutcomeKind::PaymentRequired, - "payment_required", - ), - (PushOutboxTargetOutcomeKind::Error, "error"), - (PushOutboxTargetOutcomeKind::Timeout, "timeout"), - ( - PushOutboxTargetOutcomeKind::ConnectionFailed, - "connection_failed", - ), - ( - PushOutboxTargetOutcomeKind::TargetUriRejected, - "target_uri_rejected", - ), - ( - PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted, - "skipped_already_accepted", - ), - ( - PushOutboxTargetOutcomeKind::DeferredUntilImplemented, - "deferred_until_implemented", - ), - (PushOutboxTargetOutcomeKind::Unknown, "unknown"), - ] { - assert_eq!(sdk_target_outcome_kind_label(kind), label); - } - - for (kind, label) in [ - (PushOutboxTransportOutcomeKind::Accepted, "accepted"), - ( - PushOutboxTransportOutcomeKind::DuplicateAccepted, - "duplicate_accepted", - ), - (PushOutboxTransportOutcomeKind::Delivered, "delivered"), - (PushOutboxTransportOutcomeKind::Forwarded, "forwarded"), - ( - PushOutboxTransportOutcomeKind::StoredByGateway, - "stored_by_gateway", - ), - (PushOutboxTransportOutcomeKind::Seen, "seen"), - ( - PushOutboxTransportOutcomeKind::DeferredUntilImplemented, - "deferred_until_implemented", - ), - (PushOutboxTransportOutcomeKind::Rejected, "rejected"), - ( - PushOutboxTransportOutcomeKind::RouteUnavailable, - "route_unavailable", - ), - ( - PushOutboxTransportOutcomeKind::PayloadTooLarge, - "payload_too_large", - ), - ( - PushOutboxTransportOutcomeKind::PolicyDenied, - "policy_denied", - ), - (PushOutboxTransportOutcomeKind::Timeout, "timeout"), - ( - PushOutboxTransportOutcomeKind::ConnectionFailed, - "connection_failed", - ), - ( - PushOutboxTransportOutcomeKind::TransportUnavailable, - "transport_unavailable", - ), - ] { - assert_eq!(sdk_transport_outcome_kind_label(kind), label); - } - } - - #[test] - fn cli_direct_rr_rs_dependencies_are_classified() { - let manifest_path = Path::new(env!("CARGO_MANIFEST_DIR")).join("Cargo.toml"); - let manifest = fs::read_to_string(&manifest_path).expect("read manifest"); - let manifest = manifest.parse::<toml::Value>().expect("parse manifest"); - let actual = direct_rr_rs_dependency_keys(&manifest); - let expected = DIRECT_RR_RS_DEPENDENCIES - .iter() - .map(direct_rr_rs_dependency_key) - .collect::<BTreeSet<_>>(); - - assert_eq!(actual, expected); - for dependency in DIRECT_RR_RS_DEPENDENCIES { - assert!(!dependency.owner.trim().is_empty()); - assert!(!dependency.reason.trim().is_empty()); - assert!(!dependency.lifecycle.trim().is_empty()); - } - } - - #[test] - fn cli_production_sources_reject_unowned_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() - .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) => { - nostr_relay_fetch_findings(&relative_path, production_source.as_str()) - } - Err(error) => vec![error], - } - }) - .collect::<Vec<_>>(); - - assert!( - findings.is_empty(), - "CLI production sources contain unowned Nostr relay fetch helpers:\n{}", - findings.join("\n") - ); - } - - #[test] - fn cli_production_sources_reject_dead_code_suppressions() { - 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() - .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) => { - if production_source.contains("allow(dead_code)") { - vec![format!("{relative_path}: production dead-code suppression")] - } else { - Vec::new() - } - } - Err(error) => vec![error], - } - }) - .collect::<Vec<_>>(); - - assert!( - findings.is_empty(), - "CLI production sources contain dead-code suppressions:\n{}", - findings.join("\n") - ); - } - - #[test] - fn sync_order_transport_cli_production_sources_use_sdk_outcome_label_contract() { - 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() - .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) => sdk_outcome_label_contract_findings( - &relative_path, - source.as_str(), - production_source.as_str(), - ), - Err(error) => vec![error], - } - }) - .collect::<Vec<_>>(); - - assert!( - findings.is_empty(), - "CLI production sources violate the SDK outcome label contract:\n{}", - findings.join("\n") - ); - } - - #[test] - fn sdk_outcome_label_guard_rejects_wildcard_unknown_inside_helper_span() { - let source = sdk_outcome_label_helper_source( - concat!( - " match kind {\n", - " PushOutboxTransportOutcomeKind::Accepted => \"accepted\".to_owned(),\n", - " _ => \"unknown\".to_owned(),\n", - " }\n" - ), - "", - ); - let production_source = - production_source_without_tests("fixture.rs", source.as_str()).expect("source"); - let findings = - sdk_outcome_label_contract_findings("fixture.rs", source.as_str(), &production_source); - - assert!( - findings - .iter() - .any(|finding| finding.contains("does not delegate to SDK enum as_str labels")), - "SDK outcome label guard must reject local helper rendering:\n{}", - findings.join("\n") - ); - assert!( - findings - .iter() - .any(|finding| finding.contains("uses a wildcard outcome label arm")), - "SDK outcome label guard must reject wildcard unknown helper rendering:\n{}", - findings.join("\n") - ); - } - - #[test] - fn sdk_outcome_label_guard_allows_unrelated_unknown_fallbacks_outside_helper_spans() { - let source = sdk_outcome_label_helper_source( - " kind.as_str().to_owned()\n", - concat!( - "fn unrelated_status_label(value: Option<&str>) -> &str {\n", - " match value {\n", - " Some(value) => value,\n", - " _ => \"unknown\",\n", - " }\n", - "}\n" - ), - ); - let production_source = - production_source_without_tests("fixture.rs", source.as_str()).expect("source"); - let findings = - sdk_outcome_label_contract_findings("fixture.rs", source.as_str(), &production_source); - - assert!( - findings.is_empty(), - "SDK outcome label guard must ignore unrelated unknown fallbacks:\n{}", - findings.join("\n") - ); +impl Clock for SystemClock { + fn now_unix_ms(&self) -> Result<u64, radroots_sync::Error> { + now_unix_ms().map_err(|_| radroots_sync::Error::ClockUnavailable) } +} - #[test] - fn cli_account_create_sources_reject_implicit_identity_ingestion() { - let account_source = rust_code_without_non_code( - "src/runtime/account.rs", - crate_source("src/runtime/account.rs").as_str(), - ) - .expect("account source classification"); - let core_source = rust_code_without_non_code( - "src/ops/exec/core.rs", - crate_source("src/ops/exec/core.rs").as_str(), - ) - .expect("core source classification"); +struct RandomIds; - for token in [ - removed_surface_fixtures::CREATE_OR_MIGRATE_DEFAULT_ACCOUNT, - removed_surface_fixtures::ACCOUNT_CREATE_MODE, - removed_surface_fixtures::MIGRATE_LEGACY_IDENTITY_FILE, - removed_surface_fixtures::CONFIG_IDENTITY_PATH_EXISTS, - ] { - assert!( - !account_source.contains(token), - "CLI account runtime must not contain implicit account-create identity-ingestion token `{token}`" - ); - } - for token in [ - removed_surface_fixtures::ACCOUNT_CREATE_MODE, - removed_surface_fixtures::MIGRATED_OUTPUT_JSON, - ] { - assert!( - !core_source.contains(token), - "CLI account create output must not contain removed account-create token `{token}`" - ); - } +impl IdSource for RandomIds { + fn next_id(&self, _operation: OperationKind) -> Result<SyncId, radroots_sync::Error> { + let mut bytes = [0_u8; 16]; + getrandom::getrandom(&mut bytes).map_err(|_| radroots_sync::Error::InvalidSyncId)?; + SyncId::new(bytes) } +} - #[test] - fn migrated_cli_paths_are_guarded_against_workflow_bypasses() { - for guard in MIGRATED_CLI_PATH_GUARDS { - let source = crate_source(guard.path); - assert_migrated_path( - guard.label, - source_segment(&source, guard.start, guard.end), - guard.required_tokens, - ); - } - } - - #[test] - fn migrated_path_root_alias_scanner_preserves_status_helpers() { - assert!( - root_trade_alias_findings( - "allowed", - "fn trade_status_for_locator() { RadrootsSdkError::trade_status_limit_invalid(0, 1, 100); }", - ) - .is_empty() - ); - - let findings = root_trade_alias_findings( - "forbidden", - "sdk.trade_status (request); RadrootsClient::trade_resync(&sdk);", - ); +#[derive(Clone)] +struct SharedProvider(radroots_sdk::signing::Provider); - for alias in ["trade_status", "trade_resync"] { - assert!( - findings.iter().any(|finding| finding.contains(alias)), - "migrated path root alias scanner must reject `{alias}`" - ); - } +impl Signer for SharedProvider { + fn status(&self) -> SigningFuture<'_, Result<SignerStatus, radroots_signing::Error>> { + Box::pin(self.0.status()) } - #[test] - fn cli_production_sources_reject_removed_sdk_status_surfaces_repo_wide() { - 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() - .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) => { - let mut findings = removed_sdk_status_surface_findings( - &relative_path, - production_source.as_str(), - ); - findings.extend(root_trade_alias_findings( - &relative_path, - production_source.as_str(), - )); - findings - } - Err(error) => vec![error], - } - }) - .collect::<Vec<_>>(); - - assert!( - findings.is_empty(), - "CLI production sources contain removed SDK surfaces:\n{}", - findings.join("\n") - ); + fn sign( + &self, + request: SignRequest, + ) -> SigningFuture<'_, Result<SignReceipt, radroots_signing::Error>> { + Box::pin(self.0.sign(request)) } +} - #[test] - fn cli_production_source_scanner_strips_test_modules() { - let inline = concat!( - "fn production() {}\n", - "#[cfg(test)] mod tests { fn test_only() { sdk.status_client(); } }\n", - ); - let multiline = concat!( - "fn production() {}\n", - "#[cfg(test)]\n", - "#[doc(hidden)]\n", - "mod tests {\n", - " fn test_only() { let _ = TradeValidationClient; }\n", - " const BRACE: &str = \"}\";\n", - "}\n", - ); - - for source in [inline, multiline] { - let production_source = - production_source_without_tests("fixture.rs", source).expect("production source"); - assert!( - removed_sdk_status_surface_findings("fixture.rs", production_source.as_str()) - .is_empty() - ); - assert!(root_trade_alias_findings("fixture.rs", production_source.as_str()).is_empty()); - } +struct UnavailableTransport; + +impl EventSource for UnavailableTransport { + fn status(&self) -> BoxFuture<'_, Result<SourceStatus, TransportError>> { + Box::pin(async { + Ok(SourceStatus::new( + TransportId::LOCAL, + false, + Maturity::Stable, + Availability::Unavailable, + SourceCapabilities::NONE, + LOCAL_TRANSPORT_MESSAGE, + )) + }) } - #[test] - fn cli_production_source_scanner_strips_cfg_test_functions() { - let source = concat!( - "fn production() {}\n", - "#[cfg(test)]\n", - "fn test_only() { let _ = TradeValidationClient; }\n", - "fn after_tests() {}\n", - ); - let production_source = - production_source_without_tests("fixture.rs", source).expect("production source"); - - assert!(!production_source.contains("TradeValidationClient")); - assert!(production_source.contains("fn production()")); - assert!(production_source.contains("fn after_tests()")); - assert!( - removed_sdk_status_surface_findings("fixture.rs", production_source.as_str()) - .is_empty() - ); + fn fetch(&self, _request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, TransportError>> { + Box::pin(async { Err(TransportError::UnsupportedOperation) }) } +} - #[test] - fn cli_production_source_scanner_strips_cfg_test_fragments() { - let source = concat!( - "enum Provider {", - "#[cfg(test)] TestOnly,", - "Production", - "}\n", - "fn production(provider: Provider) { match provider {", - "#[cfg(test)] Provider::TestOnly => sdk.status_client(),", - "Provider::Production => {}", - "} }\n", - ); - let production_source = - production_source_without_tests("fixture.rs", source).expect("production source"); - - assert!(!production_source.contains("TestOnly")); - assert!( - removed_sdk_status_surface_findings("fixture.rs", production_source.as_str()) - .is_empty() - ); +impl EventSink for UnavailableTransport { + fn status(&self) -> BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { + Ok(SinkStatus::new( + TransportId::LOCAL, + false, + Maturity::Stable, + Availability::Unavailable, + SinkCapabilities::NONE, + LOCAL_TRANSPORT_MESSAGE, + )) + }) } - #[test] - fn cli_production_source_scanner_reports_malformed_cfg_test_items() { - let source = concat!( - "fn production() {}\n", - "#[cfg(test)]\n", - "mod tests { fn hidden() { let _ = TradeValidationClient; }\n", - "fn after_tests() { sdk.status_client(); }\n", - ); - let error = - production_source_without_tests("fixture.rs", source).expect_err("classification"); - - assert!(error.contains("fixture.rs:2")); - assert!(error.contains("cfg(test) item is not closed")); + fn deliver( + &self, + _request: DeliveryRequest, + ) -> BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + Box::pin(async { Err(TransportError::UnsupportedOperation) }) } +} - #[test] - fn removed_surface_scanner_ignores_comments_and_literals() { - let source = concat!( - "fn production() {\n", - " let literal = \"status_client(\";\n", - " let raw = r#\"RadrootsClient::trade_resync(&sdk)\"#;\n", - " let character = 'x';\n", - "}\n", - "// TradeValidationClient::new(root)\n", - "/* sdk.trade_status(request) */\n", - ); - let production_source = - production_source_without_tests("fixture.rs", source).expect("production source"); - - assert!( - removed_sdk_status_surface_findings("fixture.rs", production_source.as_str()) - .is_empty() - ); - assert!(root_trade_alias_findings("fixture.rs", production_source.as_str()).is_empty()); - } +#[cfg(test)] +mod tests { + use super::*; #[test] - fn repo_wide_removed_surface_scanner_reports_production_violations() { - let source = "fn production() { sdk.status_client(); RadrootsClient::trade_resync(&sdk); }"; - let status_findings = removed_sdk_status_surface_findings("fixture.rs", source); - let alias_findings = root_trade_alias_findings("fixture.rs", source); - - assert!( - status_findings - .iter() - .any(|finding| finding.contains("status_client(")) - ); - assert!( - alias_findings - .iter() - .any(|finding| finding.contains("trade_resync")) - ); - } - - fn collect_rs_files(dir: &Path, files: &mut Vec<PathBuf>) { - for entry in fs::read_dir(dir).expect("read dir") { - let path = entry.expect("entry").path(); - if path.is_dir() { - collect_rs_files(path.as_path(), files); - } else if path.extension().and_then(|extension| extension.to_str()) == Some("rs") { - files.push(path); - } - } - } - - fn direct_rr_rs_dependency_keys(manifest: &toml::Value) -> BTreeSet<String> { - ["dependencies", "dev-dependencies"] - .into_iter() - .flat_map(|section| { - manifest - .get(section) - .and_then(toml::Value::as_table) - .into_iter() - .flat_map(move |dependencies| { - dependencies.iter().filter_map(move |(name, value)| { - dependency_path(value) - .filter(|path| { - path.contains("../lib/crates") - || path.contains("domains/radroots/lib/crates") - }) - .map(|_| format!("{section}:{name}")) - }) - }) - }) - .collect() - } - - fn direct_rr_rs_dependency_key(dependency: &DirectRrRsDependency) -> String { - format!("{}:{}", dependency.section, dependency.name) - } - - fn relative_source_path(root: &Path, path: &Path) -> String { - path.strip_prefix(root) - .expect("source path under manifest root") - .to_string_lossy() - .replace('\\', "/") - } - - fn dependency_path(value: &toml::Value) -> Option<&str> { - value - .as_table() - .and_then(|table| table.get("path")) - .and_then(toml::Value::as_str) - } - - fn crate_source(path: &str) -> String { - fs::read_to_string(Path::new(env!("CARGO_MANIFEST_DIR")).join(path)).expect("read source") - } - - fn source_segment<'a>(source: &'a str, start: &str, end: &str) -> &'a str { - let start_index = source.find(start).expect("source segment start"); - let end_index = source[start_index..] - .find(end) - .map(|index| start_index + index) - .expect("source segment end"); - &source[start_index..end_index] - } - - fn assert_migrated_path(label: &str, source: &str, required_tokens: &[&str]) { - let source = - rust_code_without_non_code(label, source).expect("migrated path source classification"); - - for token in required_tokens { - assert!( - source.as_str().contains(token), - "{label} does not contain required SDK token `{token}`" - ); - } - - for token in MIGRATED_PATH_DISALLOWED_TOKENS { - assert!( - !source.as_str().contains(token), - "{label} contains disallowed migrated-path token `{token}`" - ); - } - - let findings = root_trade_alias_findings(label, source.as_str()); - assert!( - findings.is_empty(), - "{label} contains removed SDK root trade aliases:\n{}", - findings.join("\n") + fn random_host_policies_produce_valid_values() { + assert!(SystemClock.now_unix_ms().expect("clock") > 0); + assert_ne!( + RandomIds + .next_id(OperationKind::Pull) + .expect("sync id") + .as_bytes(), + &[0; 16] ); } - - fn root_trade_alias_findings(label: &str, source: &str) -> Vec<String> { - let mut findings = Vec::new(); - - for alias in REMOVED_SDK_ROOT_TRADE_ALIAS_NAMES { - for (index, _) in source.match_indices(alias) { - let before = source[..index].chars().next_back(); - let after_index = index + alias.len(); - let after = source[after_index..].chars().next(); - - if before.is_some_and(is_rust_identifier_character) - || after.is_some_and(is_rust_identifier_character) - { - continue; - } - - if source[after_index..] - .chars() - .find(|character| !character.is_whitespace()) - != Some('(') - { - continue; - } - - let prefix = source[..index].trim_end(); - if prefix.ends_with('.') || prefix.ends_with("::") { - findings.push(format!( - "{label}:{} uses removed SDK root trade alias `{alias}`", - line_number(source, index) - )); - } - } - } - - findings - } - - fn removed_sdk_status_surface_findings(label: &str, source: &str) -> Vec<String> { - REMOVED_SDK_STATUS_SURFACE_TOKENS - .iter() - .flat_map(|token| { - source.match_indices(token).map(move |(index, _)| { - format!( - "{label}:{} uses removed SDK surface `{token}`", - line_number(source, index) - ) - }) - }) - .collect() - } - - fn nostr_relay_fetch_findings(label: &str, source: &str) -> Vec<String> { - NOSTR_RELAY_FETCH_DISALLOWED_TOKENS - .iter() - .flat_map(|token| { - source.match_indices(token).map(move |(index, _)| { - format!( - "{label}:{} uses unowned Nostr relay fetch token `{token}`", - line_number(source, index) - ) - }) - }) - .collect() - } - - fn sdk_outcome_label_contract_findings( - label: &str, - raw_source: &str, - production_source: &str, - ) -> Vec<String> { - let mut findings = SDK_OUTCOME_LABEL_SOURCE_DISALLOWED_TOKENS - .iter() - .flat_map(|(token, reason)| { - production_source.match_indices(token).map(move |(index, _)| { - format!( - "{label}:{} uses forbidden SDK outcome label token `{token}` for {reason}", - line_number(production_source, index) - ) - }) - }) - .collect::<Vec<_>>(); - - findings.extend(sdk_outcome_label_helper_findings(label, raw_source)); - - findings - } - - fn sdk_outcome_label_helper_findings(label: &str, source: &str) -> Vec<String> { - let mut findings = Vec::new(); - for guard in SDK_OUTCOME_LABEL_HELPER_GUARDS { - let (start_index, end_index) = - match optional_source_segment_bounds(source, guard.start, guard.end) { - Ok(Some(bounds)) => bounds, - Ok(None) => continue, - Err(error) => { - findings.push(format!("{label}: {error}")); - continue; - } - }; - let raw_segment = &source[start_index..end_index]; - let code_segment = - rust_code_without_non_code(label, raw_segment).expect("SDK label helper source"); - if !code_segment.contains("kind.as_str().to_owned()") { - findings.push(format!( - "{label}:{} {helper} does not delegate to SDK enum as_str labels", - line_number(source, start_index), - helper = guard.label - )); - } - if let Some(index) = code_segment.find("match kind") { - findings.push(format!( - "{label}:{} {helper} matches SDK outcome variants locally", - line_number(source, start_index + index), - helper = guard.label - )); - } - if let Some(index) = code_segment.find("_ =>") { - findings.push(format!( - "{label}:{} {helper} uses a wildcard outcome label arm", - line_number(source, start_index + index), - helper = guard.label - )); - } - } - findings - } - - fn sdk_outcome_label_helper_source(transport_body: &str, extra_source: &str) -> String { - format!( - "pub fn sdk_transport_outcome_kind_label(kind: PushOutboxTransportOutcomeKind) -> String {{\n\ -{transport_body}\ -}}\n\ -pub fn sdk_target_outcome_kind_label(kind: PushOutboxTargetOutcomeKind) -> String {{\n\ - kind.as_str().to_owned()\n\ -}}\n\ -#[derive(Debug, Clone, PartialEq, Eq)]\n\ -struct FixtureConfig;\n\ -{extra_source}" - ) - } - - fn optional_source_segment_bounds( - source: &str, - start: &str, - end: &str, - ) -> Result<Option<(usize, usize)>, String> { - let Some(start_index) = source.find(start) else { - return Ok(None); - }; - let Some(end_index) = source[start_index..] - .find(end) - .map(|index| start_index + index) - else { - return Err(format!( - "SDK outcome label helper source segment starting `{start}` is missing end marker `{end}`" - )); - }; - Ok(Some((start_index, end_index))) - } - - 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()); - let mut cursor = 0; - - while let Some((attribute_start, attribute_end)) = - find_cfg_test_attribute(code_source.as_str(), cursor) - { - let item_end = - find_cfg_test_item_end(path, code_source.as_str(), attribute_start, attribute_end)?; - - production_source.push_str(&code_source[cursor..attribute_start]); - push_masked_source( - &mut production_source, - &code_source[attribute_start..item_end], - ); - cursor = item_end; - } - - production_source.push_str(&code_source[cursor..]); - Ok(production_source) - } - - fn rust_code_without_non_code(path: &str, source: &str) -> Result<String, String> { - let bytes = source.as_bytes(); - let mut code = String::with_capacity(source.len()); - let mut cursor = 0; - - while cursor < bytes.len() { - match bytes[cursor] { - b'"' => { - let end = skip_quoted_rust_literal(source, cursor, b'"').ok_or_else(|| { - classification_error(path, source, cursor, "unterminated string literal") - })?; - push_masked_source(&mut code, &source[cursor..end]); - cursor = end; - } - b'\'' => { - if let Some(end) = skip_rust_char_literal(source, cursor) { - push_masked_source(&mut code, &source[cursor..end]); - cursor = end; - } else { - let character = source[cursor..].chars().next().expect("quote"); - code.push(character); - cursor += character.len_utf8(); - } - } - b'/' if bytes.get(cursor + 1) == Some(&b'/') => { - let end = source[cursor..] - .find('\n') - .map_or(source.len(), |newline| cursor + newline); - push_masked_source(&mut code, &source[cursor..end]); - cursor = end; - } - b'/' if bytes.get(cursor + 1) == Some(&b'*') => { - let end = skip_rust_block_comment(source, cursor).ok_or_else(|| { - classification_error(path, source, cursor, "unterminated block comment") - })?; - push_masked_source(&mut code, &source[cursor..end]); - cursor = end; - } - b'r' => { - if let Some(end) = skip_raw_rust_string(source, cursor) { - push_masked_source(&mut code, &source[cursor..end]); - cursor = end; - } else { - code.push('r'); - cursor += 1; - } - } - _ => { - let character = source[cursor..].chars().next().expect("source character"); - code.push(character); - cursor += character.len_utf8(); - } - } - } - - Ok(code) - } - - fn push_masked_source(output: &mut String, source: &str) { - for character in source.chars() { - if character == '\n' { - output.push('\n'); - } else { - output.push(' '); - } - } - } - - fn find_cfg_test_attribute(source: &str, start: usize) -> Option<(usize, usize)> { - let mut cursor = start; - while let Some(relative_start) = source[cursor..].find("#[") { - let attribute_start = cursor + relative_start; - let content_start = attribute_start + 2; - let content_end = source[content_start..] - .find(']') - .map(|relative_end| content_start + relative_end)?; - let normalized = source[content_start..content_end] - .chars() - .filter(|character| !character.is_whitespace()) - .collect::<String>(); - let attribute_end = content_end + 1; - if normalized == "cfg(test)" { - return Some((attribute_start, attribute_end)); - } - cursor = attribute_end; - } - None - } - - fn find_cfg_test_item_end( - path: &str, - source: &str, - attribute_start: usize, - attribute_end: usize, - ) -> Result<usize, String> { - let mut cursor = skip_rust_whitespace(source, attribute_end); - while source[cursor..].starts_with("#[") { - let attribute_end = source[cursor..] - .find(']') - .map(|end| cursor + end + 1) - .ok_or_else(|| { - classification_error(path, source, cursor, "unterminated attribute") - })?; - cursor = skip_rust_whitespace(source, attribute_end); - } - - cursor = skip_optional_visibility(path, source, cursor)?; - find_rust_item_end(path, source, attribute_start, cursor) - } - - fn skip_optional_visibility(path: &str, source: &str, cursor: usize) -> Result<usize, String> { - if !starts_with_rust_keyword(source, cursor, "pub") { - return Ok(cursor); - } - - let mut cursor = skip_rust_whitespace(source, cursor + "pub".len()); - if source[cursor..].starts_with('(') { - cursor = skip_balanced_parentheses(source, cursor).ok_or_else(|| { - classification_error(path, source, cursor, "malformed visibility") - })?; - cursor = skip_rust_whitespace(source, cursor); - } - Ok(cursor) - } - - fn find_rust_item_end( - path: &str, - source: &str, - attribute_start: usize, - item_start: usize, - ) -> Result<usize, String> { - let bytes = source.as_bytes(); - let mut cursor = item_start; - let mut brace_depth = 0usize; - let mut bracket_depth = 0usize; - let mut paren_depth = 0usize; - let mut saw_brace = false; - - while cursor < bytes.len() { - match bytes[cursor] { - b'(' => paren_depth += 1, - b')' => { - paren_depth = paren_depth.checked_sub(1).ok_or_else(|| { - classification_error(path, source, cursor, "unbalanced closing parenthesis") - })?; - } - b'[' => bracket_depth += 1, - b']' => { - bracket_depth = bracket_depth.checked_sub(1).ok_or_else(|| { - classification_error(path, source, cursor, "unbalanced closing bracket") - })?; - } - b'{' if paren_depth == 0 && bracket_depth == 0 => { - brace_depth += 1; - saw_brace = true; - } - b'}' if paren_depth == 0 && bracket_depth == 0 => { - brace_depth = brace_depth.checked_sub(1).ok_or_else(|| { - classification_error(path, source, cursor, "unbalanced closing brace") - })?; - cursor += 1; - if saw_brace && brace_depth == 0 { - return Ok(cursor); - } - continue; - } - b';' if paren_depth == 0 && bracket_depth == 0 && brace_depth == 0 => { - return Ok(cursor + 1); - } - b',' if paren_depth == 0 && bracket_depth == 0 && brace_depth == 0 => { - return Ok(cursor + 1); - } - _ => {} - } - cursor += 1; - } - - Err(classification_error( - path, - source, - attribute_start, - "cfg(test) item is not closed", - )) - } - - fn skip_balanced_parentheses(source: &str, open_index: usize) -> Option<usize> { - let mut depth = 0usize; - for (relative_index, character) in source[open_index..].char_indices() { - match character { - '(' => depth += 1, - ')' => { - depth = depth.checked_sub(1)?; - if depth == 0 { - return Some(open_index + relative_index + character.len_utf8()); - } - } - _ => {} - } - } - None - } - - fn skip_rust_whitespace(source: &str, mut cursor: usize) -> usize { - while cursor < source.len() { - let Some(character) = source[cursor..].chars().next() else { - return cursor; - }; - if !character.is_whitespace() { - return cursor; - } - cursor += character.len_utf8(); - } - cursor - } - - fn starts_with_rust_keyword(source: &str, cursor: usize, keyword: &str) -> bool { - source[cursor..].starts_with(keyword) - && source[cursor + keyword.len()..] - .chars() - .next() - .is_none_or(|character| !is_rust_identifier_character(character)) - } - - fn skip_rust_block_comment(source: &str, start: usize) -> Option<usize> { - let bytes = source.as_bytes(); - let mut cursor = start; - let mut depth = 0usize; - - while cursor + 1 < bytes.len() { - if bytes[cursor] == b'/' && bytes[cursor + 1] == b'*' { - depth += 1; - cursor += 2; - continue; - } - if bytes[cursor] == b'*' && bytes[cursor + 1] == b'/' { - depth = depth.checked_sub(1)?; - cursor += 2; - if depth == 0 { - return Some(cursor); - } - continue; - } - cursor += 1; - } - - None - } - - fn skip_quoted_rust_literal(source: &str, start: usize, delimiter: u8) -> Option<usize> { - let bytes = source.as_bytes(); - let mut cursor = start + 1; - let mut escaped = false; - while cursor < bytes.len() { - let byte = bytes[cursor]; - if escaped { - escaped = false; - } else if byte == b'\\' { - escaped = true; - } else if byte == delimiter { - return Some(cursor + 1); - } - cursor += 1; - } - None - } - - fn skip_rust_char_literal(source: &str, start: usize) -> Option<usize> { - let end = skip_quoted_rust_literal(source, start, b'\'')?; - let literal = &source[start + 1..end - 1]; - if literal.starts_with('\\') || literal.chars().count() == 1 { - Some(end) - } else { - None - } - } - - fn skip_raw_rust_string(source: &str, start: usize) -> Option<usize> { - let bytes = source.as_bytes(); - let mut cursor = start + 1; - let mut hashes = 0usize; - - while bytes.get(cursor) == Some(&b'#') { - hashes += 1; - cursor += 1; - } - - if bytes.get(cursor) != Some(&b'"') { - return None; - } - cursor += 1; - - while cursor < bytes.len() { - if bytes[cursor] == b'"' { - let mut matched = true; - for offset in 0..hashes { - if bytes.get(cursor + 1 + offset) != Some(&b'#') { - matched = false; - break; - } - } - if matched { - return Some(cursor + 1 + hashes); - } - } - cursor += 1; - } - - None - } - - fn is_rust_identifier_character(character: char) -> bool { - character == '_' || character.is_ascii_alphanumeric() - } - - fn line_number(source: &str, index: usize) -> usize { - source[..index] - .bytes() - .filter(|byte| *byte == b'\n') - .count() - + 1 - } - - fn classification_error(path: &str, source: &str, index: usize, reason: &str) -> String { - format!( - "{path}:{} source classification failed: {reason}", - line_number(source, index) - ) - } - - fn sample_config(root: &Path, relays: Vec<String>) -> RuntimeConfig { - let data = root.join("data"); - let cache = root.join("cache"); - let logs = root.join("logs"); - let secrets = root.join("secrets"); - RuntimeConfig { - output: OutputConfig { - format: OutputFormat::Json, - verbosity: Verbosity::Normal, - dry_run: false, - }, - interaction: InteractionConfig { - input_enabled: false, - assume_yes: false, - stdin_tty: false, - stdout_tty: false, - prompts_allowed: false, - confirmations_allowed: false, - }, - paths: PathsConfig { - profile: "interactive_user".to_owned(), - profile_source: "test".to_owned(), - allowed_profiles: vec!["interactive_user".to_owned(), "repo_local".to_owned()], - root_source: "test".to_owned(), - repo_local_root: None, - repo_local_root_source: None, - subordinate_path_override_source: "runtime_config".to_owned(), - app_namespace: "apps/cli".to_owned(), - shared_accounts_namespace: "shared/accounts".to_owned(), - shared_identities_namespace: "shared/identities".to_owned(), - app_config_path: root.join("config/apps/cli/config.toml"), - workspace_config_path: None, - app_data_root: data.join("apps/cli"), - shared_cache_root: cache.clone(), - app_logs_root: logs.join("apps/cli"), - shared_accounts_data_root: data.join("shared/accounts"), - shared_accounts_secrets_root: secrets.join("shared/accounts"), - default_identity_path: secrets.join("shared/identities/default.json"), - }, - logging: LoggingConfig { - filter: "info".to_owned(), - directory: None, - stdout: false, - }, - account: AccountConfig { - selector: None, - store_path: data.join("shared/accounts/store.json"), - secrets_dir: secrets.join("shared/accounts"), - secret_backend: RadrootsSecretBackend::EncryptedFile, - }, - account_secret_contract: AccountSecretContractConfig { - default_backend: "host_vault".to_owned(), - allowed_backends: vec!["host_vault".to_owned(), "encrypted_file".to_owned()], - host_vault_policy: Some("desktop".to_owned()), - uses_protected_store: true, - }, - identity: IdentityConfig { - path: secrets.join("shared/identities/default.json"), - }, - signer: SignerConfig { - backend: SignerBackend::Local, - }, - transport: crate::runtime::config::TransportConfig::from_nostr_relay_urls( - relays.clone(), - ), - local: LocalConfig { - root: data.join("apps/cli/replica"), - replica_store_path: data.join("apps/cli/replica/replica.sqlite"), - backups_dir: data.join("apps/cli/replica/backups"), - exports_dir: data.join("apps/cli/replica/exports"), - }, - myc: MycConfig { - executable: PathBuf::from("myc"), - status_timeout_ms: 2_000, - }, - hyf: HyfConfig { - enabled: false, - executable: PathBuf::from("hyfd"), - }, - mesh: crate::runtime::config::MeshConfig::disabled(), - rpc: RpcConfig { - url: "http://127.0.0.1:7070".to_owned(), - }, - rhi: RhiConfig { - validator_set: None, - require_cryptographic_proof: false, - }, - capability_bindings: Vec::new(), - } - } } diff --git a/src/runtime/sync.rs b/src/runtime/sync.rs @@ -1,375 +1,202 @@ -use std::thread; -use std::time::{Duration, SystemTime, UNIX_EPOCH}; +//! CLI scheduling and presentation over the canonical shared sync engine. -use radroots_event::envelope::kind::{ - KIND_CLASSIFIED_LISTING, KIND_FARM, KIND_LIST_SET_APP_CURATION, KIND_LIST_SET_BOOKMARK, - KIND_LIST_SET_CALENDAR, KIND_LIST_SET_CURATION, KIND_LIST_SET_EMOJI, KIND_LIST_SET_FOLLOW, - KIND_LIST_SET_GENERIC, KIND_LIST_SET_INTEREST, KIND_LIST_SET_KIND_MUTE, - KIND_LIST_SET_MEDIA_STARTER_PACK, KIND_LIST_SET_PICTURE, KIND_LIST_SET_RELAY, - KIND_LIST_SET_RELEASE_ARTIFACT, KIND_LIST_SET_STARTER_PACK, KIND_LIST_SET_VIDEO, KIND_PLOT, - KIND_PROFILE, +use std::{ + thread, + time::{Duration, SystemTime, UNIX_EPOCH}, }; -use radroots_nostr::prelude::{ - RadrootsNostrFilter, RadrootsNostrTimestamp, radroots_event_from_nostr, radroots_nostr_kind, -}; -use radroots_replica_store::{ReplicaSql, migrations}; -use radroots_replica_sync::{ - RadrootsReplicaEventsError, RadrootsReplicaIngestOutcome, radroots_replica_ingest_event, - radroots_replica_sync_status, -}; -use radroots_sdk::{ - PushOutboxEventReceipt, PushOutboxEventState, PushOutboxReceipt, PushOutboxRequest, - PushOutboxTargetOutcomeKind, PushOutboxTargetReceipt, SyncStatusReceipt, SyncStatusRequest, + +use crate::{ + cli::global::SyncWatchArgs, + runtime::{ + RuntimeError, + config::{RuntimeConfig, TransportProfileKind}, + sdk::{CliSdkAdapterError, CliSdkSession, sync_targets}, + }, + view::runtime::{ + SyncActionView, SyncFreshnessView, SyncQueueView, SyncStatusView, + SyncTransportStatusInspectView, SyncTransportTargetView, SyncWatchFrameView, SyncWatchView, + TransportOperationCapabilitiesView, TransportTargetFailureView, + }, }; +use radroots_replica_store::ReplicaSql; +#[cfg(test)] +use radroots_replica_store::migrations; use radroots_sql_core::{SqlExecutor, SqlxSqliteExecutor}; -use radroots_transport::{ - RADROOTS_RETICULUM_ENDPOINT_URI, RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, - RadrootsTransportCapabilityAvailability, RadrootsTransportCapabilityMaturity, - RadrootsTransportImplementationState, RadrootsTransportKind, RadrootsTransportStatus, - RadrootsTransportTarget, -}; -use radroots_transport_nostr::{ - RadrootsRelayFetchFailure, RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError, -}; -use serde::Deserialize; -use serde_json::json; - -use crate::cli::global::SyncWatchArgs; -use crate::runtime::RuntimeError; -use crate::runtime::config::{RuntimeConfig, TransportProfileKind}; -use crate::runtime::sdk::{ - CliSdkAdapterError, CliSdkSession, fetch_relay_events_via_shared_transport, - sdk_nostr_relay_url_policy, sdk_target_outcome_kind_label, sdk_transport_outcome_kind_label, +use radroots_storage::outbox::LeaseOwner; +use radroots_sync::{ + PullRequest, SyncStatus, + ingest::RegistryPolicy, + policy::SyncId, + pull::PullTermination, + push::{DeliveryRunReceipt, DeliveryRunRequest}, }; -use crate::view::runtime::{ - SyncActionView, SyncFreshnessView, SyncQueueView, SyncRunFreshnessView, SyncStatusView, - SyncTransportStatusInspectView, SyncTransportTargetView, SyncWatchFrameView, SyncWatchView, - TransportOperationCapabilitiesView, TransportTargetFailureView, +use radroots_transport::{ + SinkStatus, SourceStatus, Target, + capability::{Availability, Maturity}, + outcome::FetchTargetState, }; -const SYNC_SOURCE: &str = "local replica · local first"; -const SDK_SYNC_SOURCE: &str = "SDK canonical event store and outbox"; -const SDK_PUSH_SOURCE: &str = "SDK outbox push"; -const RELAY_PULL_SETUP_ACTION: &str = - "radroots transport config update --kind nostr --nostr-relay wss://relay.example.com"; -const SYNC_PULL_ACTION: &str = "radroots sync pull"; -const SYNC_PUSH_ACTION: &str = "radroots sync push"; -const SYNC_READY_ACTION: &str = "radroots market search eggs"; -const MARKET_READY_ACTION: &str = "radroots market search eggs"; -const INGEST_SOURCE: &str = "shared Nostr transport fetch · local replica ingest"; -const RELAY_FETCH_LIMIT: usize = 1_000; -const RELAY_FETCH_MAX_PAGES: usize = 5; +const SDK_SYNC_SOURCE: &str = "canonical SDK sync engine"; const MARKET_FRESHNESS_STALE_AFTER_SECONDS: u64 = 15 * 60; const SYNC_PULL_FRESHNESS_STALE_AFTER_SECONDS: u64 = 30 * 60; -const SYNC_RUN_TABLE: &str = "radroots_cli_sync_run"; -const MARKET_REFRESH_KINDS: &[u32] = &[KIND_PROFILE, KIND_FARM, KIND_CLASSIFIED_LISTING]; -const SYNC_PULL_KINDS: &[u32] = &[ - KIND_PROFILE, - KIND_FARM, - KIND_PLOT, - KIND_CLASSIFIED_LISTING, - KIND_LIST_SET_FOLLOW, - KIND_LIST_SET_GENERIC, - KIND_LIST_SET_RELAY, - KIND_LIST_SET_BOOKMARK, - KIND_LIST_SET_CURATION, - KIND_LIST_SET_VIDEO, - KIND_LIST_SET_PICTURE, - KIND_LIST_SET_KIND_MUTE, - KIND_LIST_SET_INTEREST, - KIND_LIST_SET_EMOJI, - KIND_LIST_SET_RELEASE_ARTIFACT, - KIND_LIST_SET_APP_CURATION, - KIND_LIST_SET_CALENDAR, - KIND_LIST_SET_STARTER_PACK, - KIND_LIST_SET_MEDIA_STARTER_PACK, -]; - -#[derive(Debug, Clone)] -struct SyncSnapshot { - state: String, - source: String, - local_root: String, - replica_store: String, - configured_transport_target_count: usize, - configured_transport_targets: Vec<SyncTransportTargetView>, - transport_statuses: Vec<SyncTransportStatusInspectView>, - publish_policy: String, - freshness: SyncFreshnessView, - queue: SyncQueueView, - reason: Option<String>, - actions: Vec<String>, -} - -struct SyncTransportMetadata { - configured_transport_target_count: usize, - configured_transport_targets: Vec<SyncTransportTargetView>, - transport_statuses: Vec<SyncTransportStatusInspectView>, -} - -#[derive(Debug, Clone)] -struct SyncRunRecord { - scope: String, - relay_set_fingerprint: String, - target_transport_endpoints_json: String, - attempted_transport_endpoints_json: String, - failed_transport_targets_json: String, - started_at: u64, - completed_at: Option<u64>, - state: String, - fetched_count: usize, - ingested_count: usize, - skipped_count: usize, - unsupported_count: usize, - failed_count: usize, - failure_reason: Option<String>, -} - -#[derive(Debug, Deserialize)] -struct SyncRunRow { - scope: String, - relay_set_fingerprint: String, - target_transport_endpoints_json: String, - attempted_transport_endpoints_json: String, - failed_transport_targets_json: String, - started_at: i64, - completed_at: Option<i64>, - state: String, - fetched_count: i64, - ingested_count: i64, - skipped_count: i64, - unsupported_count: i64, - failed_count: i64, - failure_reason: Option<String>, -} +const PULL_PAGE_LIMIT: u16 = 1_000; +const PULL_MAX_PAGES: u16 = 5; +const DELIVERY_LIMIT: u16 = 1_000; +const DELIVERY_LEASE_MS: u64 = 30_000; pub fn status(config: &RuntimeConfig) -> Result<SyncStatusView, CliSdkAdapterError> { let session = CliSdkSession::connect(config)?; - let receipt = session.block_on(session.sdk().sync().status(SyncStatusRequest::new()))?; - Ok(sdk_sync_status_view(config, receipt)) + let receipt = canonical_status(&session)?; + Ok(status_view(config, &receipt)) } pub fn pull(config: &RuntimeConfig) -> Result<SyncActionView, RuntimeError> { - pull_with_fetcher(config, shared_relay_transport_fetch_windowed) + pull_for_scope(config, RelayIngestScope::SyncPull) } pub fn market_refresh(config: &RuntimeConfig) -> Result<SyncActionView, RuntimeError> { - market_refresh_with_fetcher(config, shared_relay_transport_fetch_windowed) -} - -fn pull_with_fetcher<F>(config: &RuntimeConfig, fetcher: F) -> Result<SyncActionView, RuntimeError> -where - F: FnOnce( - &[String], - RadrootsNostrFilter, - ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>, -{ - relay_ingest(config, RelayIngestScope::SyncPull, fetcher) -} - -fn market_refresh_with_fetcher<F>( - config: &RuntimeConfig, - fetcher: F, -) -> Result<SyncActionView, RuntimeError> -where - F: FnOnce( - &[String], - RadrootsNostrFilter, - ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>, -{ - relay_ingest(config, RelayIngestScope::MarketPull, fetcher) -} - -fn shared_relay_transport_fetch_windowed( - relay_urls: &[String], - base_filter: RadrootsNostrFilter, -) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> { - let mut next_filter = base_filter.clone(); - let mut merged: Option<RadrootsRelayFetchedEventsReceipt> = None; - - for _ in 0..RELAY_FETCH_MAX_PAGES { - 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.event.created_at.as_secs()) - .min(); - merge_fetch_receipt(&mut merged, receipt); - if page_len < RELAY_FETCH_LIMIT { - break; - } - let Some(oldest_created_at) = oldest_created_at else { - break; - }; - if oldest_created_at == 0 { - break; - } - next_filter = base_filter - .clone() - .until(RadrootsNostrTimestamp::from(oldest_created_at - 1)) - .limit(RELAY_FETCH_LIMIT); - } - - merged.ok_or(RadrootsRelayTransportError::EmptyTargetSet) + pull_for_scope(config, RelayIngestScope::MarketPull) } -fn relay_ingest<F>( +fn pull_for_scope( config: &RuntimeConfig, scope: RelayIngestScope, - fetcher: F, -) -> Result<SyncActionView, RuntimeError> -where - F: FnOnce( - &[String], - RadrootsNostrFilter, - ) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError>, -{ - let snapshot = inspect_sync(config)?; - if snapshot.state != "ready" { - return Ok(empty_action_from_snapshot(snapshot, "pull")); +) -> Result<SyncActionView, RuntimeError> { + let session = CliSdkSession::connect(config).map_err(adapter_runtime_error)?; + let before = canonical_status(&session).map_err(adapter_runtime_error)?; + if !nostr_targets_configured(config) { + return Ok(empty_action( + config, + &before, + "pull", + "unconfigured", + Some("sync pull requires at least one configured Nostr relay".to_owned()), + )); } - if config.output.dry_run { - let mut view = empty_action_from_snapshot(snapshot, "pull"); - view.state = "ready".to_owned(); - view.reason = Some("dry run requested; Nostr fetch skipped".to_owned()); + let mut view = empty_action( + config, + &before, + "pull", + "dry_run", + Some("dry run requested; canonical sync pull skipped".to_owned()), + ); view.target_transport_endpoints = config.transport.nostr_relay_urls.clone(); view.fetched_count = Some(0); view.ingested_count = Some(0); - view.publishable_count = None; - view.published_count = None; view.skipped_count = Some(0); view.unsupported_count = Some(0); view.failed_count = Some(0); view.reason_code = Some("dry_run".to_owned()); - view.actions = vec![scope.ready_action().to_owned()]; return Ok(view); } - let started_at = unix_now(); - let receipt = match fetcher(&config.transport.nostr_relay_urls, scope.filter()) { - Ok(receipt) if receipt.connected_relays.is_empty() && !receipt.failed_relays.is_empty() => { - let target_transport_endpoints = receipt.target_relays; - let failed_transport_targets = relay_failures(receipt.failed_relays); - let reason = relay_failure_reason(&failed_transport_targets); - let failure_reason = format!("Nostr transport fetch failed: {reason}"); - let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - record_sync_run( - &executor, - &sync_record_from_failure( - scope, - &config.transport.nostr_relay_urls, - target_transport_endpoints.clone(), - failed_transport_targets.clone(), - started_at, - failure_reason.clone(), - )?, - )?; - let mut view = empty_action_from_snapshot(snapshot, "pull"); - view.state = "unavailable".to_owned(); - view.reason = Some(failure_reason); - view.reason_code = Some("nostr_fetch_failed".to_owned()); - view.target_transport_endpoints = target_transport_endpoints; - view.failed_transport_targets = failed_transport_targets; - 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 = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - record_sync_run( - &executor, - &sync_record_from_failure( - scope, - &config.transport.nostr_relay_urls, - config.transport.nostr_relay_urls.clone(), - Vec::new(), - started_at, - failure_reason.clone(), - )?, - )?; - let mut view = empty_action_from_snapshot(snapshot, "pull"); - view.state = "unavailable".to_owned(); - view.reason = Some(failure_reason); - view.reason_code = Some("nostr_fetch_failed".to_owned()); - view.target_transport_endpoints = config.transport.nostr_relay_urls.clone(); - view.freshness = freshness_for_scope_from_executor(config, &executor, scope)?; - return Ok(view); + let targets = sync_targets(config).map_err(|error| RuntimeError::Config(error.to_string()))?; + let request = PullRequest::new(targets.clone(), PULL_PAGE_LIMIT, PULL_MAX_PAGES) + .map_err(|error| RuntimeError::Config(error.to_string()))?; + let operations = session + .sdk() + .sync() + .map_err(CliSdkAdapterError::from) + .map_err(adapter_runtime_error)? + .ok_or_else(|| { + RuntimeError::Config("canonical sync engine is not configured".to_owned()) + })?; + let receipt = session + .block_on(operations.pull(request, &RegistryPolicy::verified())) + .map_err(|error| adapter_runtime_error(error.into()))?; + let after = canonical_status(&session).map_err(adapter_runtime_error)?; + let ingested = receipt + .ingest_outcomes() + .iter() + .filter(|item| item.is_ok()) + .count(); + let ingest_failed = receipt.ingest_outcomes().len().saturating_sub(ingested); + let mut failed_targets = Vec::new(); + for outcome in receipt.target_outcomes() { + if matches!(outcome.state(), FetchTargetState::Complete) { + continue; } + let target = targets + .targets() + .iter() + .find(|target| target.fingerprint() == outcome.target()); + failed_targets.push(TransportTargetFailureView { + transport_kind: target + .map_or("nostr", |target| target.kind().as_str()) + .to_owned(), + endpoint_uri: target.map_or_else( + || outcome.target().as_str().to_owned(), + |target| target.uri().as_str().to_owned(), + ), + target_scope: target + .and_then(|target| target.scope()) + .map(|value| value.as_str().to_owned()), + target_label: target + .and_then(|target| target.label()) + .map(|value| value.as_str().to_owned()), + transport_outcome_kind: Some(fetch_state_label(outcome.state()).to_owned()), + reason: outcome + .message() + .unwrap_or(fetch_state_label(outcome.state())) + .to_owned(), + }); + } + let failed_count = ingest_failed + failed_targets.len(); + let state = if receipt.termination() == PullTermination::SourceFailed { + "unavailable" + } else if failed_count > 0 || receipt.termination() != PullTermination::Complete { + "partial" + } else { + "ready" }; - - let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - let ingest = ingest_events(&executor, &receipt, scope)?; - record_sync_run( - &executor, - &sync_record_from_ingest( - scope, - &config.transport.nostr_relay_urls, - &receipt, - &ingest, - started_at, - )?, - )?; - let failed_transport_targets = relay_failures(receipt.failed_relays); - let failed_count = ingest.failed_count + failed_transport_targets.len(); - let reason_code = - relay_ingest_reason_code(&ingest, &failed_transport_targets).map(str::to_owned); - let reason = relay_ingest_reason(&ingest, &failed_transport_targets); - let freshness = freshness_for_scope_from_executor(config, &executor, scope)?; - let queue = radroots_replica_sync_status(&executor)?; - - Ok(SyncActionView { - direction: "pull".to_owned(), - state: "ready".to_owned(), - source: INGEST_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "ready".to_owned(), - configured_transport_target_count: snapshot.configured_transport_target_count, - configured_transport_targets: snapshot.configured_transport_targets, - transport_statuses: snapshot.transport_statuses, - publish_policy: "any".to_owned(), - freshness, - queue: derived_projection_sync_queue(queue.expected_count, queue.pending_count), - target_transport_endpoints: receipt.target_relays, - attempted_transport_endpoints: receipt.connected_relays, - accepted_transport_endpoints: Vec::new(), - failed_transport_targets, - fetched_count: Some(ingest.fetched_count), - ingested_count: Some(ingest.ingested_count), - publishable_count: None, - published_count: None, - skipped_count: Some(ingest.skipped_count), - unsupported_count: Some(ingest.unsupported_count), - failed_count: Some(failed_count), - publish_plan: None, - reason_code, - reason, - actions: vec![scope.ready_action().to_owned()], - }) + let reason = (state != "ready").then(|| { + format!( + "canonical pull terminated as {} with {failed_count} failed outcome(s)", + pull_termination_label(receipt.termination()) + ) + }); + let mut view = empty_action(config, &after, "pull", state, reason); + view.target_transport_endpoints = config.transport.nostr_relay_urls.clone(); + view.attempted_transport_endpoints = config.transport.nostr_relay_urls.clone(); + view.accepted_transport_endpoints = config + .transport + .nostr_relay_urls + .iter() + .filter(|endpoint| { + !failed_targets + .iter() + .any(|failure| &failure.endpoint_uri == *endpoint) + }) + .cloned() + .collect(); + view.failed_transport_targets = failed_targets; + view.fetched_count = Some(receipt.events_observed()); + view.ingested_count = Some(ingested); + view.skipped_count = Some(0); + view.unsupported_count = Some(0); + view.failed_count = Some(failed_count); + view.reason_code = + (state != "ready").then(|| pull_termination_label(receipt.termination()).to_owned()); + view.actions = vec![scope.ready_action().to_owned()]; + Ok(view) } pub fn push(config: &RuntimeConfig) -> Result<SyncActionView, CliSdkAdapterError> { let session = CliSdkSession::connect(config)?; + let before = canonical_status(&session)?; if config.output.dry_run { - let status = session.block_on(session.sdk().sync().status(SyncStatusRequest::new()))?; - return Ok(sdk_push_dry_run_view(config, status)); + return Ok(empty_action( + config, + &before, + "push", + "dry_run", + Some("dry run requested; canonical outbox delivery skipped".to_owned()), + )); } - - let receipt = session.block_on(session.sdk().sync().push_outbox( - PushOutboxRequest::new().with_nostr_relay_url_policy(sdk_nostr_relay_url_policy(config)), - ))?; - let status = session.block_on(session.sdk().sync().status(SyncStatusRequest::new()))?; - Ok(sdk_push_view(config, receipt, status)) + let receipt = deliver_pending(&session)?; + let after = canonical_status(&session)?; + Ok(delivery_action(config, &after, &receipt)) } pub fn watch(config: &RuntimeConfig, args: &SyncWatchArgs) -> Result<SyncWatchView, RuntimeError> { @@ -378,43 +205,100 @@ pub fn watch(config: &RuntimeConfig, args: &SyncWatchArgs) -> Result<SyncWatchVi "`sync watch --frames` must be greater than 0".to_owned(), )); } - let mut frames = Vec::with_capacity(args.frames); - let mut last_snapshot = None; - + let mut final_view = None; for index in 0..args.frames { - let snapshot = inspect_sync(config)?; + let view = status(config).map_err(adapter_runtime_error)?; frames.push(SyncWatchFrameView { sequence: index + 1, observed_at: unix_now(), - state: snapshot.state.clone(), - configured_transport_target_count: snapshot.configured_transport_target_count, - freshness: snapshot.freshness.clone(), - queue: snapshot.queue.clone(), + state: view.state.clone(), + configured_transport_target_count: view.configured_transport_target_count, + freshness: view.freshness.clone(), + queue: view.queue.clone(), }); - last_snapshot = Some(snapshot); - + final_view = Some(view); if index + 1 < args.frames { thread::sleep(Duration::from_millis(args.interval_ms)); } } - - let snapshot = last_snapshot.expect("watch frames are non-empty"); + let view = final_view + .ok_or_else(|| RuntimeError::Config("sync watch produced no frames".to_owned()))?; Ok(SyncWatchView { - state: snapshot.state, - source: snapshot.source, + state: view.state, + source: view.source, interval_ms: args.interval_ms, frames, - reason: snapshot.reason, - actions: snapshot.actions, + reason: view.reason, + actions: view.actions, }) } -fn empty_action_from_snapshot(snapshot: SyncSnapshot, direction: &str) -> SyncActionView { +pub(crate) fn canonical_status(session: &CliSdkSession) -> Result<SyncStatus, CliSdkAdapterError> { + let operations = session.sdk().sync()?.ok_or_else(|| { + RuntimeError::Config("canonical sync engine is not configured".to_owned()) + })?; + Ok(session.block_on(operations.status(&[]))?) +} + +pub(crate) fn deliver_pending( + session: &CliSdkSession, +) -> Result<DeliveryRunReceipt, CliSdkAdapterError> { + let mut seed = [0_u8; 16]; + getrandom::getrandom(&mut seed).map_err(|error| { + RuntimeError::Config(format!( + "failed to generate delivery lease identity: {error}" + )) + })?; + let owner = LeaseOwner::parse(format!("radroots-cli-{}", std::process::id())) + .map_err(|error| RuntimeError::Config(error.to_string()))?; + let request = + DeliveryRunRequest::new(owner, SyncId::new(seed)?, DELIVERY_LEASE_MS, DELIVERY_LIMIT)?; + let operations = session.sdk().sync()?.ok_or_else(|| { + RuntimeError::Config("canonical sync engine is not configured".to_owned()) + })?; + Ok(session.block_on(operations.deliver_pending(request))?) +} + +fn status_view(config: &RuntimeConfig, status: &SyncStatus) -> SyncStatusView { + let configured_targets = configured_targets(config); + let state = sync_health_label(status); + SyncStatusView { + state: state.to_owned(), + source: SDK_SYNC_SOURCE.to_owned(), + local_root: config.local.root.display().to_string(), + replica_store: "canonical".to_owned(), + configured_transport_target_count: configured_targets.len(), + configured_transport_targets: configured_targets, + transport_statuses: transport_status_views(status), + publish_policy: "canonical satisfaction policy".to_owned(), + freshness: missing_freshness(), + queue: queue_view(status), + reason: (!nostr_targets_configured(config)) + .then(|| "no Nostr relay targets are configured".to_owned()), + actions: if !nostr_targets_configured(config) { + vec!["radroots transport config update --kind nostr --nostr-relay wss://relay.example.com".to_owned()] + } else { + vec![ + "radroots sync pull".to_owned(), + "radroots sync push".to_owned(), + ] + }, + } +} + +fn empty_action( + config: &RuntimeConfig, + status: &SyncStatus, + direction: &str, + state: &str, + reason: Option<String>, +) -> SyncActionView { + let snapshot = status_view(config, status); SyncActionView { direction: direction.to_owned(), - state: snapshot.state, - source: snapshot.source, + state: state.to_owned(), + source: SDK_SYNC_SOURCE.to_owned(), local_root: snapshot.local_root, replica_store: snapshot.replica_store, configured_transport_target_count: snapshot.configured_transport_target_count, @@ -436,2876 +320,336 @@ fn empty_action_from_snapshot(snapshot: SyncSnapshot, direction: &str) -> SyncAc failed_count: None, publish_plan: None, reason_code: None, - reason: snapshot.reason, + reason, actions: snapshot.actions, } } -fn sdk_sync_status_view(config: &RuntimeConfig, receipt: SyncStatusReceipt) -> SyncStatusView { - let actions = sdk_sync_status_actions(&receipt); - let configured_transport_target_count = - receipt.transport_profile.configured_transport_target_count; - let configured_transport_targets = sdk_transport_targets(&receipt); - let transport_statuses = sdk_transport_statuses(&receipt); - SyncStatusView { - state: "ready".to_owned(), - source: SDK_SYNC_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "derived_projection_not_checked".to_owned(), - configured_transport_target_count, - configured_transport_targets, - transport_statuses, - publish_policy: "any".to_owned(), - freshness: sdk_sync_freshness(&receipt), - queue: sdk_sync_queue(&receipt), - reason: None, - actions, - } -} - -fn sdk_push_dry_run_view(config: &RuntimeConfig, status: SyncStatusReceipt) -> SyncActionView { - let publishable_count = usize_from_i64(status.outbox.ready_signed_events); - let state = if publishable_count > 0 { - "dry_run" - } else { +fn delivery_action( + config: &RuntimeConfig, + status: &SyncStatus, + receipt: &DeliveryRunReceipt, +) -> SyncActionView { + let attempted = receipt.outcomes().len(); + let succeeded = receipt.succeeded(); + let failed = receipt.failed(); + let state = if attempted == 0 { "ready" - }; - let reason = if publishable_count > 0 { - Some("dry run requested; SDK outbox push skipped".to_owned()) - } else if status.outbox.total_events > 0 { - Some("SDK outbox has no ready signed events to push".to_owned()) + } else if failed == 0 { + "published" + } else if succeeded > 0 { + "partial" } else { - None + "unavailable" }; - sdk_push_action_view( - config, - state, - sdk_sync_queue(&status), - sdk_sync_freshness(&status), - status.transport_profile.configured_transport_target_count, - sdk_transport_targets(&status), - sdk_transport_statuses(&status), - status - .transport_profile - .configured_transport_targets - .iter() - .map(|target| target.endpoint_uri.clone()) - .collect(), - Vec::new(), - Vec::new(), - Vec::new(), - publishable_count, - 0, - 0, - Some(0), - reason, - sdk_sync_push_actions(state, publishable_count > 0), - ) -} - -fn sdk_push_view( - config: &RuntimeConfig, - receipt: PushOutboxReceipt, - status: SyncStatusReceipt, -) -> SyncActionView { - let failed_count = receipt.retryable_events + receipt.terminal_events; - let state = sdk_push_state(&receipt, failed_count); - let reason = sdk_push_reason(&receipt, failed_count); - sdk_push_action_view( + let mut view = empty_action( config, + status, + "push", state, - sdk_sync_queue(&status), - sdk_sync_freshness(&status), - status.transport_profile.configured_transport_target_count, - sdk_transport_targets(&status), - sdk_transport_statuses(&status), - sdk_push_target_transport_endpoints(&receipt, &status), - sdk_push_attempted_transport_endpoints(&receipt), - sdk_push_accepted_transport_endpoints(&receipt), - sdk_push_failed_transport_targets(&receipt), - receipt.attempted_events, - receipt.published_events, - failed_count, - Some(0), - reason, - sdk_sync_push_actions(state, failed_count > 0), - ) + (failed > 0).then(|| format!("{failed} canonical delivery outcome(s) failed")), + ); + view.target_transport_endpoints = config.transport.nostr_relay_urls.clone(); + view.attempted_transport_endpoints = if attempted > 0 { + config.transport.nostr_relay_urls.clone() + } else { + Vec::new() + }; + view.accepted_transport_endpoints = if succeeded > 0 { + config.transport.nostr_relay_urls.clone() + } else { + Vec::new() + }; + view.publishable_count = Some(attempted); + view.published_count = Some(succeeded); + view.failed_count = Some(failed); + view.reason_code = (failed > 0).then(|| "delivery_partial_failure".to_owned()); + view.actions = vec!["radroots sync status".to_owned()]; + view } -#[expect( - clippy::too_many_arguments, - reason = "sync push view construction mirrors the V1 output contract fields" -)] -fn sdk_push_action_view( - config: &RuntimeConfig, - state: &str, - queue: SyncQueueView, - freshness: SyncFreshnessView, - configured_transport_target_count: usize, - configured_transport_targets: Vec<SyncTransportTargetView>, - transport_statuses: Vec<SyncTransportStatusInspectView>, - target_transport_endpoints: Vec<String>, - attempted_transport_endpoints: Vec<String>, - accepted_transport_endpoints: Vec<String>, - failed_transport_targets: Vec<TransportTargetFailureView>, - publishable_count: usize, - published_count: usize, - failed_count: usize, - skipped_count: Option<usize>, - reason: Option<String>, - actions: Vec<String>, -) -> SyncActionView { - SyncActionView { - direction: "push".to_owned(), - state: state.to_owned(), - source: SDK_PUSH_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "derived_projection_not_checked".to_owned(), - configured_transport_target_count, - configured_transport_targets, - transport_statuses, - publish_policy: "any".to_owned(), - freshness, - queue, - target_transport_endpoints, - attempted_transport_endpoints, - accepted_transport_endpoints, - failed_transport_targets, - fetched_count: None, - ingested_count: None, - publishable_count: Some(publishable_count), - published_count: Some(published_count), - skipped_count, - unsupported_count: Some(0), - failed_count: Some(failed_count), - publish_plan: None, - reason_code: sdk_sync_push_reason_code(state).map(str::to_owned), - reason, - actions, +fn queue_view(status: &SyncStatus) -> SyncQueueView { + let outbox = status.outbox(); + let total = outbox.total().and_then(|value| usize::try_from(value).ok()); + SyncQueueView { + expected_count: total.unwrap_or_default(), + pending_count: usize::try_from(outbox.pending + outbox.leased + outbox.retryable) + .unwrap_or(usize::MAX), + total_count: total, + retryable_count: usize::try_from(outbox.retryable).ok(), + terminal_count: usize::try_from(outbox.satisfied + outbox.exhausted).ok(), + failed_terminal_count: usize::try_from(outbox.exhausted).ok(), + deferred_until_implemented_count: Some(0), + ready_signed_count: usize::try_from(outbox.pending + outbox.retryable).ok(), + publishing_count: usize::try_from(outbox.leased).ok(), + last_attempt_at_ms: None, + last_error: None, } } -fn sdk_transport_targets(receipt: &SyncStatusReceipt) -> Vec<SyncTransportTargetView> { - receipt - .transport_profile - .configured_transport_targets +fn configured_targets(config: &RuntimeConfig) -> Vec<SyncTransportTargetView> { + config + .transport + .nostr_relay_urls .iter() + .filter_map(|endpoint| Target::nostr_relay(endpoint).ok()) .map(|target| SyncTransportTargetView { - transport_kind: target.transport_kind.clone(), - endpoint_uri: target.endpoint_uri.clone(), - endpoint_fingerprint: target.endpoint_fingerprint.clone(), - target_scope: target.target_scope.clone(), - target_label: target.target_label.clone(), + transport_kind: target.kind().as_str().to_owned(), + endpoint_uri: target.uri().as_str().to_owned(), + endpoint_fingerprint: target.fingerprint().as_str().to_owned(), + target_scope: target.scope().map(|scope| scope.as_str().to_owned()), + target_label: target.label().map(|label| label.as_str().to_owned()), }) .collect() } -fn sdk_transport_statuses(receipt: &SyncStatusReceipt) -> Vec<SyncTransportStatusInspectView> { - receipt - .transport_profile - .transport_statuses - .iter() - .map(sdk_transport_status_view) - .collect() +fn transport_status_views(status: &SyncStatus) -> Vec<SyncTransportStatusInspectView> { + let mut views = Vec::new(); + if let Some(source) = status.source().status() { + views.push(source_status_view(source)); + } + if let Some(sink) = status.sink().status() + && !views + .iter() + .any(|view| view.transport == sink.transport_id().as_str()) + { + views.push(sink_status_view(sink)); + } + views } -fn sdk_transport_status_view( - status: &radroots_sdk::SyncTransportStatusSummary, -) -> SyncTransportStatusInspectView { +fn source_status_view(status: &SourceStatus) -> SyncTransportStatusInspectView { SyncTransportStatusInspectView { - transport: status.transport.clone(), - profile_id: status.profile_id.clone(), - endpoint_uri: status.endpoint_uri.clone(), - configured: status.configured, - implementation: status.implementation.clone(), - maturity: status.maturity.clone(), - availability: status.availability.clone(), - usable_for_delivery: status.usable_for_delivery, + transport: status.transport_id().as_str().to_owned(), + profile_id: None, + endpoint_uri: None, + configured: status.is_configured(), + implementation: "real".to_owned(), + maturity: maturity_label(status.maturity()).to_owned(), + availability: availability_label(status.availability()).to_owned(), + usable_for_delivery: false, capabilities: TransportOperationCapabilitiesView { - deliver: status.capabilities.deliver, - fetch: status.capabilities.fetch, + deliver: false, + fetch: status.capabilities().can_fetch(), }, - message: status.message.clone(), + message: status.message().to_owned(), } } -fn sdk_sync_status_actions(receipt: &SyncStatusReceipt) -> Vec<String> { - let mut actions = Vec::new(); - if receipt.outbox.ready_signed_events > 0 { - actions.push(SYNC_PUSH_ACTION.to_owned()); - } - if receipt.event_store.total_events == 0 { - actions.push(SYNC_PULL_ACTION.to_owned()); +fn sink_status_view(status: &SinkStatus) -> SyncTransportStatusInspectView { + SyncTransportStatusInspectView { + transport: status.transport_id().as_str().to_owned(), + profile_id: None, + endpoint_uri: None, + configured: status.is_configured(), + implementation: "real".to_owned(), + maturity: maturity_label(status.maturity()).to_owned(), + availability: availability_label(status.availability()).to_owned(), + usable_for_delivery: status.capabilities().can_deliver() + && status.availability() != Availability::Unavailable, + capabilities: TransportOperationCapabilitiesView { + deliver: status.capabilities().can_deliver(), + fetch: false, + }, + message: status.message().to_owned(), } - actions } -fn sdk_sync_push_actions(state: &str, retryable: bool) -> Vec<String> { - match state { - "published" | "ready" | "deferred_until_implemented" => { - vec!["radroots sync status".to_owned()] - } - "dry_run" | "partial" | "unavailable" if retryable => { - vec![ - SYNC_PUSH_ACTION.to_owned(), - "radroots sync status".to_owned(), - ] - } - _ => vec!["radroots sync status".to_owned()], +fn sync_health_label(status: &SyncStatus) -> &'static str { + match status.health() { + radroots_protocol::runtime::v1::SyncHealth::Healthy => "ready", + radroots_protocol::runtime::v1::SyncHealth::Degraded => "degraded", + radroots_protocol::runtime::v1::SyncHealth::Unavailable => "unavailable", } } -fn sdk_sync_push_reason_code(state: &str) -> Option<&'static str> { - match state { - "dry_run" => Some("dry_run"), - "partial" => Some("sdk_outbox_push_partial"), - "unavailable" => Some("sdk_outbox_push_failed"), - "deferred_until_implemented" => Some("sdk_outbox_push_deferred_until_implemented"), - _ => None, +fn maturity_label(value: Maturity) -> &'static str { + match value { + Maturity::Experimental => "experimental", + Maturity::Preview => "preview", + Maturity::Stable => "stable", } } -fn sdk_push_state(receipt: &PushOutboxReceipt, failed_count: usize) -> &'static str { - if receipt.attempted_events == 0 { - return sdk_push_reported_deferred_state(receipt).unwrap_or("ready"); - } - if receipt.published_events > 0 && failed_count > 0 { - "partial" - } else if failed_count > 0 { - "unavailable" - } else if receipt.published_events > 0 { - "published" - } else { - "ready" +fn availability_label(value: Availability) -> &'static str { + match value { + Availability::Available => "available", + Availability::Degraded => "degraded", + Availability::Unavailable => "unavailable", } } -fn sdk_push_reported_deferred_state(receipt: &PushOutboxReceipt) -> Option<&'static str> { - let mut deferred = false; - for event in &receipt.events { - if event.final_state == PushOutboxEventState::DeferredUntilImplemented { - deferred = true; - } - for target in &event.targets { - if target.outcome_kind == PushOutboxTargetOutcomeKind::DeferredUntilImplemented { - deferred = true; - } - } +fn fetch_state_label(value: FetchTargetState) -> &'static str { + match value { + FetchTargetState::Complete => "complete", + FetchTargetState::Partial => "partial", + FetchTargetState::Unavailable => "unavailable", + FetchTargetState::FailedRetryable => "failed_retryable", + FetchTargetState::FailedTerminal => "failed_terminal", + FetchTargetState::Cancelled => "cancelled", } - deferred.then_some("deferred_until_implemented") } -fn sdk_push_reason(receipt: &PushOutboxReceipt, failed_count: usize) -> Option<String> { - if receipt.attempted_events == 0 { - if let Some(state) = sdk_push_reported_deferred_state(receipt) { - return Some(match state { - "deferred_until_implemented" => { - "SDK outbox push reported Reticulum work as deferred until implemented without network delivery" - } - _ => "SDK outbox push reported Reticulum work without network delivery", - } - .to_owned()); - } - return Some("SDK outbox had no ready signed events to push".to_owned()); - } - if failed_count > 0 && receipt.published_events > 0 { - return Some(format!( - "SDK outbox push published {} event(s); {failed_count} event(s) remain retryable or terminal", - receipt.published_events - )); +fn pull_termination_label(value: PullTermination) -> &'static str { + match value { + PullTermination::Complete => "complete", + PullTermination::PageLimit => "page_limit", + PullTermination::Deadline => "deadline", + PullTermination::Cancelled => "cancelled", + PullTermination::SourceFailed => "source_failed", } - if failed_count > 0 { - return Some( - "SDK outbox push did not reach accepted quorum for any ready event".to_owned(), - ); - } - None } -fn sdk_sync_queue(receipt: &SyncStatusReceipt) -> SyncQueueView { - let pending_count = usize_from_i64( - receipt - .outbox - .pending_events - .saturating_add(receipt.outbox.retryable_events), - ); - SyncQueueView { - expected_count: usize_from_i64(receipt.outbox.total_events), - pending_count, - total_count: Some(usize_from_i64(receipt.outbox.total_events)), - retryable_count: Some(usize_from_i64(receipt.outbox.retryable_events)), - terminal_count: Some(usize_from_i64(receipt.outbox.terminal_events)), - failed_terminal_count: Some(usize_from_i64(receipt.outbox.failed_terminal_events)), - deferred_until_implemented_count: Some(usize_from_i64( - receipt.outbox.deferred_until_implemented_events, - )), - ready_signed_count: Some(usize_from_i64(receipt.outbox.ready_signed_events)), - publishing_count: Some(usize_from_i64(receipt.outbox.publishing_events)), - last_attempt_at_ms: receipt.outbox.last_attempt_at_ms, - last_error: receipt.outbox.last_error.clone(), - } +fn nostr_targets_configured(config: &RuntimeConfig) -> bool { + matches!( + config.transport.profile, + TransportProfileKind::Nostr | TransportProfileKind::MultiTarget + ) && !config.transport.nostr_relay_urls.is_empty() } -fn derived_projection_sync_queue(expected_count: usize, pending_count: usize) -> SyncQueueView { - SyncQueueView { - expected_count, - pending_count, - total_count: None, - retryable_count: None, - terminal_count: None, - failed_terminal_count: None, - deferred_until_implemented_count: None, - ready_signed_count: None, - publishing_count: None, - last_attempt_at_ms: None, - last_error: None, - } +fn adapter_runtime_error(error: CliSdkAdapterError) -> RuntimeError { + RuntimeError::Network(error.to_string()) } -fn sdk_sync_freshness(receipt: &SyncStatusReceipt) -> SyncFreshnessView { - let Some(last_event_updated_at_ms) = receipt.event_store.last_event_updated_at_ms else { - return missing_freshness(); - }; - let last_event_at = u64::try_from(last_event_updated_at_ms / 1_000).unwrap_or(0); - let observed_at = u64::try_from(receipt.observed_at_ms / 1_000).unwrap_or_else(|_| unix_now()); - let age_seconds = observed_at.saturating_sub(last_event_at); +pub(crate) fn missing_freshness() -> SyncFreshnessView { SyncFreshnessView { - state: "synced".to_owned(), - display: format!("SDK event store updated {}", relative_age(age_seconds)), - age_seconds: Some(age_seconds), - last_event_at: Some(last_event_at), + state: "never".to_owned(), + display: "never synced".to_owned(), + age_seconds: None, + last_event_at: None, run: None, } } -fn sdk_push_target_transport_endpoints( - receipt: &PushOutboxReceipt, - status: &SyncStatusReceipt, -) -> Vec<String> { - let mut targets = Vec::new(); - for target in receipt.events.iter().flat_map(|event| event.targets.iter()) { - if !targets.contains(&target.endpoint_uri) { - targets.push(target.endpoint_uri.clone()); - } - } - if targets.is_empty() { - targets.extend( - status - .transport_profile - .configured_transport_targets - .iter() - .map(|target| target.endpoint_uri.clone()), - ); - } - targets +#[cfg(test)] +pub(crate) fn freshness_for_scope( + config: &RuntimeConfig, + scope: RelayIngestScope, +) -> Result<SyncFreshnessView, RuntimeError> { + let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; + migrations::run_all_up(&executor)?; + freshness_for_scope_from_executor(config, &executor, scope) } -fn sdk_push_attempted_transport_endpoints(receipt: &PushOutboxReceipt) -> Vec<String> { - sdk_push_targets_matching(receipt, |_, target| target.attempted) +#[cfg(test)] +pub(crate) fn relay_provenance_relays_for_scope( + _config: &RuntimeConfig, + _scope: RelayIngestScope, +) -> Result<Vec<String>, RuntimeError> { + Ok(Vec::new()) } -fn sdk_push_accepted_transport_endpoints(receipt: &PushOutboxReceipt) -> Vec<String> { - sdk_push_targets_matching(receipt, |_, target| { - sdk_target_accepted(target.outcome_kind) +pub(crate) fn freshness_for_scope_from_executor( + _config: &RuntimeConfig, + executor: &SqlxSqliteExecutor, + scope: RelayIngestScope, +) -> Result<SyncFreshnessView, RuntimeError> { + let last_event_at = ReplicaSql::new(executor).nostr_event_last_created_at()?; + let age_seconds = last_event_at.map(|last| unix_now().saturating_sub(last)); + let state = match age_seconds { + None => "never", + Some(age) if age > scope.stale_after_seconds() => "stale", + Some(_) => "fresh", + }; + Ok(SyncFreshnessView { + state: state.to_owned(), + display: match age_seconds { + Some(age) => format!("{} {state} {}s ago", scope.display(), age), + None => format!("{} never synced", scope.display()), + }, + age_seconds, + last_event_at, + run: None, }) } -fn sdk_push_targets_matching( - receipt: &PushOutboxReceipt, - predicate: impl Fn(&PushOutboxEventReceipt, &PushOutboxTargetReceipt) -> bool, -) -> Vec<String> { - let mut targets = Vec::new(); - for event in &receipt.events { - for target in &event.targets { - if predicate(event, target) && !targets.contains(&target.endpoint_uri) { - targets.push(target.endpoint_uri.clone()); - } - } - } - targets -} - -fn sdk_push_failed_transport_targets( - receipt: &PushOutboxReceipt, -) -> Vec<TransportTargetFailureView> { - receipt - .events - .iter() - .flat_map(|event| event.targets.iter()) - .filter(|target| !sdk_target_accepted(target.outcome_kind)) - .map(|target| TransportTargetFailureView { - transport_kind: target.transport_kind.clone(), - endpoint_uri: target.endpoint_uri.clone(), - target_scope: target.target_scope.clone(), - target_label: target.target_label.clone(), - transport_outcome_kind: target - .transport_outcome_kind - .map(sdk_transport_outcome_kind_label), - reason: target - .message - .clone() - .unwrap_or_else(|| sdk_target_outcome_kind_label(target.outcome_kind)), - }) - .collect() -} - -fn sdk_target_accepted(kind: PushOutboxTargetOutcomeKind) -> bool { +pub(crate) fn freshness_requires_refresh(freshness: &SyncFreshnessView) -> bool { matches!( - kind, - PushOutboxTargetOutcomeKind::Accepted | PushOutboxTargetOutcomeKind::DuplicateAccepted + freshness.state.as_str(), + "never" | "stale" | "relay_set_changed" | "refresh_failed" ) } -fn usize_from_i64(value: i64) -> usize { - usize::try_from(value.max(0)).unwrap_or(usize::MAX) +pub(crate) fn ensure_sync_run_table(executor: &SqlxSqliteExecutor) -> Result<(), RuntimeError> { + executor.exec( + "CREATE TABLE IF NOT EXISTS radroots_cli_sync_run ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + scope TEXT NOT NULL, + relay_set_fingerprint TEXT NOT NULL, + target_transport_endpoints_json TEXT NOT NULL, + attempted_transport_endpoints_json TEXT NOT NULL, + failed_transport_targets_json TEXT NOT NULL, + started_at INTEGER NOT NULL, + completed_at INTEGER, + state TEXT NOT NULL, + fetched_count INTEGER NOT NULL, + ingested_count INTEGER NOT NULL, + skipped_count INTEGER NOT NULL, + unsupported_count INTEGER NOT NULL, + failed_count INTEGER NOT NULL, + failure_reason TEXT + ); + CREATE INDEX IF NOT EXISTS idx_radroots_cli_sync_run_scope_started + ON radroots_cli_sync_run(scope, started_at DESC);", + "[]", + )?; + Ok(()) } -fn sync_transport_metadata(config: &RuntimeConfig) -> Result<SyncTransportMetadata, RuntimeError> { - let configured_transport_targets = sync_configured_transport_targets(config)?; - Ok(SyncTransportMetadata { - configured_transport_target_count: configured_transport_targets.len(), - configured_transport_targets, - transport_statuses: sync_transport_statuses(config), - }) +#[derive(Debug, Clone, Copy)] +pub(crate) enum RelayIngestScope { + SyncPull, + MarketPull, } -fn sync_configured_transport_targets( - config: &RuntimeConfig, -) -> Result<Vec<SyncTransportTargetView>, RuntimeError> { - let mut targets = Vec::new(); - match config.transport.profile { - TransportProfileKind::LocalOnly => {} - TransportProfileKind::Nostr => { - for relay_url in &config.transport.nostr_relay_urls { - targets.push(sync_transport_target_view( - RadrootsTransportKind::Nostr, - relay_url, - )?); - } - } - TransportProfileKind::Reticulum => { - targets.push(sync_transport_target_view( - RadrootsTransportKind::Reticulum, - RADROOTS_RETICULUM_ENDPOINT_URI, - )?); - } - TransportProfileKind::MultiTarget => { - for relay_url in &config.transport.nostr_relay_urls { - targets.push(sync_transport_target_view( - RadrootsTransportKind::Nostr, - relay_url, - )?); - } - targets.push(sync_transport_target_view( - RadrootsTransportKind::Reticulum, - RADROOTS_RETICULUM_ENDPOINT_URI, - )?); +impl RelayIngestScope { + fn display(self) -> &'static str { + match self { + Self::SyncPull => "sync pull", + Self::MarketPull => "market refresh", } } - Ok(targets) -} -fn sync_transport_target_view( - kind: RadrootsTransportKind, - endpoint_uri: &str, -) -> Result<SyncTransportTargetView, RuntimeError> { - let target = match kind { - RadrootsTransportKind::Nostr => RadrootsTransportTarget::nostr_relay(endpoint_uri), - RadrootsTransportKind::Reticulum => { - RadrootsTransportTarget::reticulum_with_metadata(endpoint_uri, None, None) + fn stale_after_seconds(self) -> u64 { + match self { + Self::SyncPull => SYNC_PULL_FRESHNESS_STALE_AFTER_SECONDS, + Self::MarketPull => MARKET_FRESHNESS_STALE_AFTER_SECONDS, } - RadrootsTransportKind::Local => RadrootsTransportTarget::local(endpoint_uri), } - .map_err(|error| RuntimeError::Config(format!("invalid transport target: {error}")))?; - Ok(SyncTransportTargetView { - transport_kind: target.kind().canonical_label(), - endpoint_uri: target.uri().as_str().to_owned(), - endpoint_fingerprint: target.fingerprint().as_str().to_owned(), - target_scope: target.scope().map(|scope| scope.as_str().to_owned()), - target_label: target.label().map(|label| label.as_str().to_owned()), - }) -} -fn sync_transport_statuses(config: &RuntimeConfig) -> Vec<SyncTransportStatusInspectView> { - let profile_id = config.transport.profile.as_str(); - let nostr_targets_configured = !config.transport.nostr_relay_urls.is_empty(); - match config.transport.profile { - TransportProfileKind::LocalOnly => vec![sync_transport_status_view( - RadrootsTransportStatus::new( - RadrootsTransportKind::Local, - true, - RadrootsTransportImplementationState::Real, - false, - "Local-only profile writes only to local state", - ) - .with_profile_id(profile_id), - )], - TransportProfileKind::Nostr => { - vec![sync_transport_status_view(nostr_transport_status( - profile_id, - nostr_targets_configured, - ))] - } - TransportProfileKind::Reticulum => { - vec![sync_transport_status_view(reticulum_transport_status( - profile_id, - ))] - } - TransportProfileKind::MultiTarget => vec![ - sync_transport_status_view(nostr_transport_status( - profile_id, - nostr_targets_configured, - )), - sync_transport_status_view(reticulum_transport_status(profile_id)), - ], - } -} - -fn nostr_transport_status(profile_id: &str, targets_configured: bool) -> RadrootsTransportStatus { - RadrootsTransportStatus::new( - RadrootsTransportKind::Nostr, - targets_configured, - RadrootsTransportImplementationState::Real, - targets_configured, - if targets_configured { - "Nostr relay transport is configured for delivery" - } else { - "Nostr transport requires configured Nostr relay targets" - }, - ) - .with_profile_id(profile_id) -} - -fn reticulum_transport_status(profile_id: &str) -> RadrootsTransportStatus { - RadrootsTransportStatus::new( - RadrootsTransportKind::Reticulum, - true, - RadrootsTransportImplementationState::Real, - false, - RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, - ) - .with_profile_id(profile_id) - .with_endpoint_uri(RADROOTS_RETICULUM_ENDPOINT_URI) - .with_maturity(RadrootsTransportCapabilityMaturity::Preview) - .with_availability(RadrootsTransportCapabilityAvailability::Unavailable) -} - -fn sync_transport_status_view(status: RadrootsTransportStatus) -> SyncTransportStatusInspectView { - SyncTransportStatusInspectView { - transport: status.kind.canonical_label(), - profile_id: status.profile_id, - endpoint_uri: status.endpoint_uri, - configured: status.configured, - implementation: transport_implementation_label(status.implementation).to_owned(), - maturity: transport_maturity_label(status.maturity).to_owned(), - availability: transport_availability_label(status.availability).to_owned(), - usable_for_delivery: status.usable_for_delivery, - capabilities: TransportOperationCapabilitiesView { - deliver: status.capabilities.deliver, - fetch: status.capabilities.fetch, - }, - message: status.message, - } -} - -fn transport_implementation_label(state: RadrootsTransportImplementationState) -> &'static str { - match state { - RadrootsTransportImplementationState::Real => "real", - RadrootsTransportImplementationState::Mock => "mock", - } -} - -fn transport_maturity_label(maturity: RadrootsTransportCapabilityMaturity) -> &'static str { - match maturity { - RadrootsTransportCapabilityMaturity::Preview => "preview", - RadrootsTransportCapabilityMaturity::Stable => "stable", - } -} - -fn transport_availability_label( - availability: RadrootsTransportCapabilityAvailability, -) -> &'static str { - match availability { - RadrootsTransportCapabilityAvailability::Available => "available", - RadrootsTransportCapabilityAvailability::Degraded => "degraded", - RadrootsTransportCapabilityAvailability::Unavailable => "unavailable", - } -} - -fn inspect_sync(config: &RuntimeConfig) -> Result<SyncSnapshot, RuntimeError> { - let transport_metadata = sync_transport_metadata(config)?; - if !config.local.replica_store_path.exists() { - return Ok(SyncSnapshot { - state: "unconfigured".to_owned(), - source: SYNC_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "missing".to_owned(), - configured_transport_target_count: transport_metadata.configured_transport_target_count, - configured_transport_targets: transport_metadata.configured_transport_targets, - transport_statuses: transport_metadata.transport_statuses, - publish_policy: "any".to_owned(), - freshness: missing_freshness(), - queue: derived_projection_sync_queue(0, 0), - reason: Some("local replica database is not initialized".to_owned()), - actions: vec!["radroots store inspect".to_owned()], - }); - } - - let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - let queue = radroots_replica_sync_status(&executor)?; - let freshness = - freshness_for_scope_from_executor(config, &executor, RelayIngestScope::SyncPull)?; - let publish_policy = "any".to_owned(); - let mut actions = Vec::new(); - - if config.transport.nostr_relay_urls.is_empty() { - let (state, reason) = if transport_metadata.configured_transport_target_count == 0 { - ( - "unconfigured", - "no transport targets are configured for this operator session", - ) - } else { - ( - "unavailable", - "active transport profile does not expose Nostr relay fetch targets", - ) - }; - actions.push( - if transport_metadata.configured_transport_target_count == 0 { - RELAY_PULL_SETUP_ACTION.to_owned() - } else { - "radroots transport config inspect".to_owned() - }, - ); - return Ok(SyncSnapshot { - state: state.to_owned(), - source: SYNC_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "ready".to_owned(), - configured_transport_target_count: transport_metadata.configured_transport_target_count, - configured_transport_targets: transport_metadata.configured_transport_targets, - transport_statuses: transport_metadata.transport_statuses, - publish_policy, - freshness, - queue: derived_projection_sync_queue(queue.expected_count, queue.pending_count), - reason: Some(reason.to_owned()), - actions, - }); - } - - actions.push(SYNC_PULL_ACTION.to_owned()); - if queue.pending_count > 0 { - actions.push(SYNC_PUSH_ACTION.to_owned()); - } - - Ok(SyncSnapshot { - state: "ready".to_owned(), - source: SYNC_SOURCE.to_owned(), - local_root: config.local.root.display().to_string(), - replica_store: "ready".to_owned(), - configured_transport_target_count: transport_metadata.configured_transport_target_count, - configured_transport_targets: transport_metadata.configured_transport_targets, - transport_statuses: transport_metadata.transport_statuses, - publish_policy, - freshness, - queue: derived_projection_sync_queue(queue.expected_count, queue.pending_count), - reason: None, - actions, - }) -} - -pub(crate) fn missing_freshness() -> SyncFreshnessView { - SyncFreshnessView { - state: "never".to_owned(), - display: "never synced".to_owned(), - age_seconds: None, - last_event_at: None, - run: None, - } -} - -#[cfg(test)] -pub(crate) fn freshness_for_scope( - config: &RuntimeConfig, - scope: RelayIngestScope, -) -> Result<SyncFreshnessView, RuntimeError> { - let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - freshness_for_scope_from_executor(config, &executor, scope) -} - -#[cfg(test)] -pub(crate) fn relay_provenance_relays_for_scope( - config: &RuntimeConfig, - scope: RelayIngestScope, -) -> Result<Vec<String>, RuntimeError> { - if !config.local.replica_store_path.exists() { - return Ok(Vec::new()); - } - let executor = SqlxSqliteExecutor::open(&config.local.replica_store_path)?; - migrations::run_all_up(&executor)?; - ensure_sync_run_table(&executor)?; - let current_fingerprint = relay_set_fingerprint(&config.transport.nostr_relay_urls); - let Some(run) = latest_sync_run(&executor, scope)? else { - return Ok(Vec::new()); - }; - if run.relay_set_fingerprint != current_fingerprint || !sync_run_successful(&run) { - return Ok(Vec::new()); - } - let mut relays: Vec<String> = - serde_json::from_str(run.attempted_transport_endpoints_json.as_str())?; - relays.sort(); - relays.dedup(); - Ok(relays) -} - -pub(crate) fn freshness_for_scope_from_executor( - config: &RuntimeConfig, - executor: &SqlxSqliteExecutor, - scope: RelayIngestScope, -) -> Result<SyncFreshnessView, RuntimeError> { - let last_event_at = ReplicaSql::new(executor).nostr_event_last_created_at()?; - let now = unix_now(); - let age_seconds = last_event_at.map(|last_event_at| now.saturating_sub(last_event_at)); - ensure_sync_run_table(executor)?; - let current_fingerprint = relay_set_fingerprint(&config.transport.nostr_relay_urls); - let latest = latest_sync_run(executor, scope)?; - let current = latest - .as_ref() - .filter(|run| run.relay_set_fingerprint == current_fingerprint); - let last_success = current.filter(|run| sync_run_successful(run)); - let state = freshness_state(scope, latest.as_ref(), current, last_success, age_seconds); - let display = freshness_display(scope, state.as_str(), age_seconds, current); - - Ok(SyncFreshnessView { - state, - display, - age_seconds, - last_event_at, - run: latest.map(|run| sync_run_freshness_view(scope, run, current_fingerprint)), - }) -} - -pub(crate) fn freshness_requires_refresh(freshness: &SyncFreshnessView) -> bool { - matches!( - freshness.state.as_str(), - "never" | "stale" | "relay_set_changed" | "refresh_failed" - ) -} - -fn freshness_state( - scope: RelayIngestScope, - latest: Option<&SyncRunRecord>, - current: Option<&SyncRunRecord>, - last_success: Option<&SyncRunRecord>, - age_seconds: Option<u64>, -) -> String { - let Some(latest) = latest else { - return "never".to_owned(); - }; - let Some(current) = current else { - return "relay_set_changed".to_owned(); - }; - if !sync_run_successful(current) { - return "refresh_failed".to_owned(); - } - if last_success.is_none() { - return "refresh_failed".to_owned(); - } - if age_seconds.is_none() { - return "fresh".to_owned(); - } - if age_seconds.unwrap_or_default() > scope.stale_after_seconds() { - return "stale".to_owned(); - } - if latest.state == "partial" { - return "partial".to_owned(); - } - "fresh".to_owned() -} - -fn freshness_display( - scope: RelayIngestScope, - state: &str, - age_seconds: Option<u64>, - run: Option<&SyncRunRecord>, -) -> String { - match state { - "fresh" => match age_seconds { - Some(age_seconds) => format!("{} fresh {}", scope.display(), relative_age(age_seconds)), - None => format!("{} fresh; no market events yet", scope.display()), - }, - "partial" => match age_seconds { - Some(age_seconds) => format!( - "{} partially refreshed {}", - scope.display(), - relative_age(age_seconds) - ), - None => format!( - "{} partially refreshed; no market events yet", - scope.display() - ), - }, - "stale" => match age_seconds { - Some(age_seconds) => format!("{} stale {}", scope.display(), relative_age(age_seconds)), - None => format!("{} stale", scope.display()), - }, - "relay_set_changed" => format!("{} relay set changed; refresh required", scope.display()), - "refresh_failed" => run - .and_then(|run| run.failure_reason.clone()) - .unwrap_or_else(|| format!("{} refresh failed", scope.display())), - _ => format!("{} never synced", scope.display()), - } -} - -fn sync_run_successful(run: &SyncRunRecord) -> bool { - matches!(run.state.as_str(), "success" | "partial") -} - -fn sync_run_freshness_view( - scope: RelayIngestScope, - run: SyncRunRecord, - current_fingerprint: String, -) -> SyncRunFreshnessView { - let relay_set_current = run.relay_set_fingerprint == current_fingerprint; - let successful = sync_run_successful(&run); - let last_successful_at = successful.then_some(run.completed_at.unwrap_or(run.started_at)); - SyncRunFreshnessView { - scope: run.scope, - relay_set_fingerprint: run.relay_set_fingerprint, - relay_set_current, - last_state: run.state, - last_attempted_at: Some(run.started_at), - last_successful_at, - last_completed_at: run.completed_at, - stale_after_seconds: Some(scope.stale_after_seconds()), - fetched_count: Some(run.fetched_count), - ingested_count: Some(run.ingested_count), - skipped_count: Some(run.skipped_count), - unsupported_count: Some(run.unsupported_count), - failed_count: Some(run.failed_count), - failure_reason: run.failure_reason, - } -} - -pub(crate) fn ensure_sync_run_table(executor: &SqlxSqliteExecutor) -> Result<(), RuntimeError> { - executor.exec( - "CREATE TABLE IF NOT EXISTS radroots_cli_sync_run ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - scope TEXT NOT NULL, - relay_set_fingerprint TEXT NOT NULL, - target_transport_endpoints_json TEXT NOT NULL, - attempted_transport_endpoints_json TEXT NOT NULL, - failed_transport_targets_json TEXT NOT NULL, - started_at INTEGER NOT NULL, - completed_at INTEGER, - state TEXT NOT NULL, - fetched_count INTEGER NOT NULL, - ingested_count INTEGER NOT NULL, - skipped_count INTEGER NOT NULL, - unsupported_count INTEGER NOT NULL, - failed_count INTEGER NOT NULL, - failure_reason TEXT - ); - CREATE INDEX IF NOT EXISTS idx_radroots_cli_sync_run_scope_started - ON radroots_cli_sync_run(scope, started_at DESC);", - "[]", - )?; - Ok(()) -} - -fn latest_sync_run( - executor: &SqlxSqliteExecutor, - scope: RelayIngestScope, -) -> Result<Option<SyncRunRecord>, RuntimeError> { - let rows = executor.query_raw( - &format!( - "SELECT scope, - relay_set_fingerprint, - target_transport_endpoints_json, - attempted_transport_endpoints_json, - failed_transport_targets_json, - started_at, - completed_at, - state, - fetched_count, - ingested_count, - skipped_count, - unsupported_count, - failed_count, - failure_reason - FROM {SYNC_RUN_TABLE} - WHERE scope = ?1 - ORDER BY started_at DESC, id DESC - LIMIT 1" - ), - json!([scope.id()]).to_string().as_str(), - )?; - let mut rows: Vec<SyncRunRow> = serde_json::from_str(rows.as_str())?; - Ok(rows.pop().map(sync_run_record_from_row)) -} - -fn sync_run_record_from_row(row: SyncRunRow) -> SyncRunRecord { - SyncRunRecord { - scope: row.scope, - relay_set_fingerprint: row.relay_set_fingerprint, - target_transport_endpoints_json: row.target_transport_endpoints_json, - attempted_transport_endpoints_json: row.attempted_transport_endpoints_json, - failed_transport_targets_json: row.failed_transport_targets_json, - started_at: u64_from_db(row.started_at), - completed_at: row.completed_at.map(u64_from_db), - state: row.state, - fetched_count: usize_from_db(row.fetched_count), - ingested_count: usize_from_db(row.ingested_count), - skipped_count: usize_from_db(row.skipped_count), - unsupported_count: usize_from_db(row.unsupported_count), - failed_count: usize_from_db(row.failed_count), - failure_reason: row.failure_reason, - } -} - -fn record_sync_run( - executor: &SqlxSqliteExecutor, - record: &SyncRunRecord, -) -> Result<(), RuntimeError> { - ensure_sync_run_table(executor)?; - executor.exec( - &format!( - "INSERT INTO {SYNC_RUN_TABLE} ( - scope, - relay_set_fingerprint, - target_transport_endpoints_json, - attempted_transport_endpoints_json, - failed_transport_targets_json, - started_at, - completed_at, - state, - fetched_count, - ingested_count, - skipped_count, - unsupported_count, - failed_count, - failure_reason - ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)" - ), - json!([ - record.scope.as_str(), - record.relay_set_fingerprint.as_str(), - record.target_transport_endpoints_json.as_str(), - record.attempted_transport_endpoints_json.as_str(), - record.failed_transport_targets_json.as_str(), - i64_from_u64(record.started_at), - record.completed_at.map(i64_from_u64), - record.state.as_str(), - i64_from_usize(record.fetched_count), - i64_from_usize(record.ingested_count), - i64_from_usize(record.skipped_count), - i64_from_usize(record.unsupported_count), - i64_from_usize(record.failed_count), - record.failure_reason.as_deref(), - ]) - .to_string() - .as_str(), - )?; - Ok(()) -} - -fn sync_record_from_failure( - scope: RelayIngestScope, - relays: &[String], - target_transport_endpoints: Vec<String>, - failed_transport_targets: Vec<TransportTargetFailureView>, - started_at: u64, - reason: String, -) -> Result<SyncRunRecord, RuntimeError> { - Ok(SyncRunRecord { - scope: scope.id().to_owned(), - relay_set_fingerprint: relay_set_fingerprint(relays), - target_transport_endpoints_json: serde_json::to_string(&target_transport_endpoints)?, - attempted_transport_endpoints_json: serde_json::to_string(&Vec::<String>::new())?, - failed_transport_targets_json: serde_json::to_string(&failed_transport_targets)?, - started_at, - completed_at: Some(unix_now()), - state: "failed".to_owned(), - fetched_count: 0, - ingested_count: 0, - skipped_count: 0, - unsupported_count: 0, - failed_count: 1, - failure_reason: Some(reason), - }) -} - -fn sync_record_from_ingest( - scope: RelayIngestScope, - relays: &[String], - receipt: &RadrootsRelayFetchedEventsReceipt, - ingest: &RelayIngestCounts, - started_at: u64, -) -> Result<SyncRunRecord, RuntimeError> { - let failed_transport_targets = relay_failures(receipt.failed_relays.clone()); - let state = if ingest.failed_count > 0 || !failed_transport_targets.is_empty() { - "partial" - } else { - "success" - }; - Ok(SyncRunRecord { - scope: scope.id().to_owned(), - relay_set_fingerprint: relay_set_fingerprint(relays), - target_transport_endpoints_json: serde_json::to_string(&receipt.target_relays)?, - attempted_transport_endpoints_json: serde_json::to_string(&receipt.connected_relays)?, - failed_transport_targets_json: serde_json::to_string(&failed_transport_targets)?, - started_at, - completed_at: Some(unix_now()), - state: state.to_owned(), - fetched_count: ingest.fetched_count, - ingested_count: ingest.ingested_count, - skipped_count: ingest.skipped_count, - unsupported_count: ingest.unsupported_count, - failed_count: ingest.failed_count + failed_transport_targets.len(), - failure_reason: ingest.reason(), - }) -} - -fn relay_set_fingerprint(relays: &[String]) -> String { - let mut normalized = relays - .iter() - .map(|relay| relay.trim().to_ascii_lowercase()) - .filter(|relay| !relay.is_empty()) - .collect::<Vec<_>>(); - normalized.sort(); - normalized.dedup(); - let mut hash = 0xcbf29ce484222325_u64; - for relay in normalized { - for byte in relay.as_bytes() { - hash ^= u64::from(*byte); - hash = hash.wrapping_mul(0x100000001b3); - } - hash ^= 0xff; - hash = hash.wrapping_mul(0x100000001b3); - } - format!("relayset_{hash:016x}") -} - -fn u64_from_db(value: i64) -> u64 { - u64::try_from(value).unwrap_or_default() -} - -fn usize_from_db(value: i64) -> usize { - usize::try_from(value).unwrap_or_default() -} - -fn i64_from_u64(value: u64) -> i64 { - i64::try_from(value).unwrap_or(i64::MAX) -} - -fn i64_from_usize(value: usize) -> i64 { - i64::try_from(value).unwrap_or(i64::MAX) -} - -#[derive(Debug, Clone, Default)] -struct RelayIngestCounts { - fetched_count: usize, - ingested_count: usize, - skipped_count: usize, - unsupported_count: usize, - failed_count: usize, - first_failure_reason: Option<String>, -} - -impl RelayIngestCounts { - fn reason_code(&self) -> Option<&'static str> { - if self.failed_count > 0 { - Some("sync_ingest_failed") - } else if self.skipped_count > 0 { - Some("sync_no_overwrite") - } else { - None - } - } - - fn reason(&self) -> Option<String> { - if self.failed_count > 0 { - return Some(match &self.first_failure_reason { - Some(reason) => format!( - "{} fetched event(s) failed ingest: {}", - self.failed_count, reason - ), - None => format!("{} fetched event(s) failed ingest", self.failed_count), - }); - } - if self.skipped_count > 0 { - return Some(format!( - "{} fetched event(s) skipped because the local replica already has current or newer state", - self.skipped_count - )); - } - None - } -} - -fn relay_ingest_reason_code( - ingest: &RelayIngestCounts, - failed_transport_targets: &[TransportTargetFailureView], -) -> Option<&'static str> { - ingest - .reason_code() - .or_else(|| (!failed_transport_targets.is_empty()).then_some("nostr_fetch_partial")) -} - -fn relay_ingest_reason( - ingest: &RelayIngestCounts, - failed_transport_targets: &[TransportTargetFailureView], -) -> Option<String> { - let mut parts = Vec::new(); - if let Some(reason) = ingest.reason() { - parts.push(reason); - } - if !failed_transport_targets.is_empty() { - parts.push(format!( - "{} relay(s) failed during fetch: {}", - failed_transport_targets.len(), - relay_failure_reason(failed_transport_targets) - )); - } - - if parts.is_empty() { - None - } else { - Some(parts.join("; ")) - } -} - -fn relay_failure_reason(failed_transport_targets: &[TransportTargetFailureView]) -> String { - failed_transport_targets - .iter() - .map(|failure| format!("{}: {}", failure.endpoint_uri, failure.reason)) - .collect::<Vec<_>>() - .join("; ") -} - -#[derive(Debug, Clone, Copy)] -pub(crate) enum RelayIngestScope { - SyncPull, - MarketPull, -} - -impl RelayIngestScope { - fn id(self) -> &'static str { - match self { - Self::SyncPull => "sync_pull", - Self::MarketPull => "market_refresh", - } - } - - fn display(self) -> &'static str { - match self { - Self::SyncPull => "sync pull", - Self::MarketPull => "market refresh", - } - } - - fn stale_after_seconds(self) -> u64 { - match self { - Self::SyncPull => SYNC_PULL_FRESHNESS_STALE_AFTER_SECONDS, - Self::MarketPull => MARKET_FRESHNESS_STALE_AFTER_SECONDS, - } - } - - fn kinds(self) -> &'static [u32] { - match self { - Self::SyncPull => SYNC_PULL_KINDS, - Self::MarketPull => MARKET_REFRESH_KINDS, - } - } - - fn filter(self) -> RadrootsNostrFilter { - RadrootsNostrFilter::new() - .kinds( - self.kinds() - .iter() - .copied() - .map(|kind| radroots_nostr_kind(kind as u16)), - ) - .limit(RELAY_FETCH_LIMIT) - } - - fn ready_action(self) -> &'static str { - match self { - Self::SyncPull => SYNC_READY_ACTION, - Self::MarketPull => MARKET_READY_ACTION, - } - } - - fn supports_kind(self, kind: u32) -> bool { - self.kinds().contains(&kind) - } -} - -fn ingest_events( - executor: &SqlxSqliteExecutor, - receipt: &RadrootsRelayFetchedEventsReceipt, - scope: RelayIngestScope, -) -> Result<RelayIngestCounts, RuntimeError> { - let mut counts = RelayIngestCounts { - fetched_count: receipt.events.len(), - ..RelayIngestCounts::default() - }; - - for event in &receipt.events { - if !scope.supports_kind(event_kind(&event.event)) { - counts.unsupported_count += 1; - continue; - } - let event = match radroots_event_from_nostr(&event.event) { - Ok(event) => event, - Err(error) => { - counts.failed_count += 1; - if counts.first_failure_reason.is_none() { - counts.first_failure_reason = Some(error.to_string()); - } - continue; - } - }; - match radroots_replica_ingest_event(executor, &event) { - Ok(RadrootsReplicaIngestOutcome::Applied) => counts.ingested_count += 1, - Ok(RadrootsReplicaIngestOutcome::Excluded) => counts.unsupported_count += 1, - Ok(RadrootsReplicaIngestOutcome::Rejected) => { - counts.failed_count += 1; - if counts.first_failure_reason.is_none() { - counts.first_failure_reason = - Some("event was rejected by the local replica projection".to_owned()); - } - } - Ok(RadrootsReplicaIngestOutcome::Skipped) => counts.skipped_count += 1, - Err(error @ RadrootsReplicaEventsError::Sql(_)) => return Err(error.into()), - Err(error) => { - counts.failed_count += 1; - if counts.first_failure_reason.is_none() { - counts.first_failure_reason = Some(error.to_string()); - } - } - } - } - - Ok(counts) -} - -fn event_kind(event: &radroots_nostr::prelude::RadrootsNostrEvent) -> u32 { - u32::from(event.kind.as_u16()) -} - -fn relay_failures(failures: Vec<RadrootsRelayFetchFailure>) -> Vec<TransportTargetFailureView> { - failures - .into_iter() - .map(|failure| TransportTargetFailureView { - transport_kind: "nostr".to_owned(), - endpoint_uri: failure.relay_url, - target_scope: None, - target_label: None, - transport_outcome_kind: None, - reason: failure.reason, - }) - .collect() -} - -fn merge_fetch_receipt( - target: &mut Option<RadrootsRelayFetchedEventsReceipt>, - receipt: RadrootsRelayFetchedEventsReceipt, -) { - match target { - Some(target) => { - push_unique_many(&mut target.target_relays, receipt.target_relays.iter()); - push_unique_many( - &mut target.connected_relays, - receipt.connected_relays.iter(), - ); - for failure in receipt.failed_relays { - if !target - .failed_relays - .iter() - .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), - } -} - -fn push_unique_many<'a>(target: &mut Vec<String>, values: impl Iterator<Item = &'a String>) { - for value in values { - if !target.contains(value) { - target.push(value.clone()); - } - } -} - -fn unix_now() -> u64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|duration| duration.as_secs()) - .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(), - 1..=59 => format!("{age_seconds}s ago"), - 60..=3_599 => format!("{}m ago", age_seconds / 60), - 3_600..=86_399 => format!("{}h ago", age_seconds / 3_600), - _ => format!("{}d ago", age_seconds / 86_400), - } -} - -#[cfg(test)] -mod tests { - use std::path::{Path, PathBuf}; - - use radroots_event::envelope::kind::{ - KIND_CLASSIFIED_LISTING, KIND_FARM, KIND_LIST_SET_GENERIC, KIND_POST, - }; - use radroots_event::farm::{Farm, FarmRef}; - use radroots_event::id::RadrootsEventId; - use radroots_event::list::RadrootsListEntry; - use radroots_event::list_set::RadrootsListSet; - use radroots_event::plot::RadrootsPlot; - use radroots_event::profile::AuthoredProfile; - use radroots_event::wire::{DEFAULT_CONTENT_MAX_BYTES, Nip01EventWireParts}; - use radroots_event_codec::farm::encode as farm_encode; - use radroots_event_codec::list_set::encode as list_set_encode; - use radroots_event_codec::plot::encode as plot_encode; - use radroots_event_codec::profile::authored::authored_profile_to_wire_parts; - use radroots_identity::RadrootsIdentity; - use radroots_nostr::prelude::{ - RadrootsNostrEvent, RadrootsNostrFilter, RadrootsNostrTimestamp, - }; - use radroots_sdk::{ - PushOutboxEventReceipt, PushOutboxEventState, PushOutboxReceipt, - PushOutboxTargetOutcomeKind, PushOutboxTargetReceipt, PushOutboxTransportOutcomeKind, - SyncEventStoreStatus, SyncOutboxStatus, SyncStatusReceipt, SyncStatusSource, - SyncTransportOperationCapabilitiesSummary, SyncTransportProfileSummary, - SyncTransportStatusSummary, SyncTransportTargetSummary, - }; - use radroots_secret_vault::RadrootsSecretBackend; - use radroots_transport::{ - RADROOTS_RETICULUM_ENDPOINT_URI, RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, - }; - use radroots_transport_nostr::{ - RadrootsRelayFetchFailure, RadrootsRelayFetchedEvent, RadrootsRelayFetchedEventsReceipt, - RadrootsRelayTransportError, - }; - use tempfile::tempdir; - - use super::{ - RelayIngestScope, freshness_for_scope, inspect_sync, market_refresh_with_fetcher, - pull_with_fetcher, relay_provenance_relays_for_scope, sdk_push_dry_run_view, sdk_push_view, - sdk_sync_status_view, - }; - use crate::cli::global::{FindQueryArgs, RecordLookupArgs}; - use crate::runtime::config::{ - AccountConfig, AccountSecretContractConfig, HyfConfig, IdentityConfig, InteractionConfig, - LocalConfig, LoggingConfig, MycConfig, OutputConfig, OutputFormat, PathsConfig, - ReticulumBehavior, RpcConfig, RuntimeConfig, SignerBackend, SignerConfig, - TransportProfileKind, Verbosity, - }; - - const FARM_D_TAG: &str = "AAAAAAAAAAAAAAAAAAAAAA"; - const PLOT_D_TAG: &str = "AAAAAAAAAAAAAAAAAAAAAQ"; - const LISTING_D_TAG: &str = "AAAAAAAAAAAAAAAAAAAAAg"; - - #[test] - fn sync_pull_dry_run_skips_relay_fetch() { - let dir = tempdir().expect("tempdir"); - let mut config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - config.output.dry_run = true; - crate::runtime::store::init(&config).expect("store init"); - - let view = pull_with_fetcher(&config, |_, _| panic!("dry run must not fetch")) - .expect("sync pull dry run"); - - assert_eq!(view.state, "ready"); - assert_eq!( - view.target_transport_endpoints, - vec!["wss://relay.example.com"] - ); - assert_eq!(view.fetched_count, Some(0)); - assert_eq!(view.ingested_count, Some(0)); - assert_eq!(view.skipped_count, Some(0)); - assert_eq!(view.unsupported_count, Some(0)); - assert_eq!(view.failed_count, Some(0)); - } - - #[test] - fn sync_pull_no_relay_action_is_actionable() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), Vec::new()); - crate::runtime::store::init(&config).expect("store init"); - - let view = pull_with_fetcher(&config, |_, _| { - panic!("unconfigured sync pull must not fetch") - }) - .expect("sync pull unconfigured"); - - assert_eq!(view.state, "unconfigured"); - assert_eq!( - view.actions, - vec![ - "radroots transport config update --kind nostr --nostr-relay wss://relay.example.com" - ] - ); - } - - #[test] - fn sync_pull_empty_nostr_profile_reports_nostr_misconfigured() { - let dir = tempdir().expect("tempdir"); - let config = nostr_config(dir.path(), Vec::new()); - crate::runtime::store::init(&config).expect("store init"); - - let view = pull_with_fetcher(&config, |_, _| { - panic!("empty Nostr profile sync pull must not fetch") - }) - .expect("sync pull empty Nostr profile"); - - assert_eq!(view.state, "unconfigured"); - assert_eq!(view.configured_transport_target_count, 0); - assert!(view.configured_transport_targets.is_empty()); - assert_eq!(view.transport_statuses.len(), 1); - assert_eq!(view.transport_statuses[0].transport, "nostr"); - assert_eq!( - view.transport_statuses[0].profile_id.as_deref(), - Some("nostr") - ); - assert!(!view.transport_statuses[0].configured); - assert_eq!(view.transport_statuses[0].implementation, "real"); - assert!(!view.transport_statuses[0].usable_for_delivery); - assert!(!view.transport_statuses[0].capabilities.deliver); - assert!(!view.transport_statuses[0].capabilities.fetch); - assert!(view.target_transport_endpoints.is_empty()); - assert_eq!( - view.actions, - vec![ - "radroots transport config update --kind nostr --nostr-relay wss://relay.example.com" - ] - ); - } - - #[test] - fn sync_pull_multi_target_reports_profile_targets_and_nostr_fetch_detail() { - let dir = tempdir().expect("tempdir"); - let config = multi_target_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - - let view = - pull_with_fetcher(&config, fake_fetcher(Vec::new())).expect("sync pull multi_target"); - - assert_eq!(view.state, "ready"); - assert_eq!(view.configured_transport_target_count, 2); - assert_eq!(view.configured_transport_targets.len(), 2); - assert_eq!(view.configured_transport_targets[0].transport_kind, "nostr"); - assert_eq!( - view.configured_transport_targets[0].endpoint_uri, - "wss://relay.example.com" - ); - assert_eq!( - view.configured_transport_targets[1].transport_kind, - "reticulum" - ); - assert_eq!(view.transport_statuses.len(), 2); - assert_eq!(view.transport_statuses[0].transport, "nostr"); - assert!(view.transport_statuses[0].usable_for_delivery); - assert!(view.transport_statuses[0].capabilities.deliver); - assert!(!view.transport_statuses[0].capabilities.fetch); - assert_eq!(view.transport_statuses[1].transport, "reticulum"); - assert!(view.transport_statuses[1].configured); - assert_eq!(view.transport_statuses[1].implementation, "real"); - assert_eq!(view.transport_statuses[1].maturity, "preview"); - assert_eq!(view.transport_statuses[1].availability, "unavailable"); - assert!(!view.transport_statuses[1].usable_for_delivery); - assert!(!view.transport_statuses[1].capabilities.deliver); - assert!(!view.transport_statuses[1].capabilities.fetch); - assert_eq!( - view.target_transport_endpoints, - vec!["wss://relay.example.com"] - ); - assert_eq!( - view.attempted_transport_endpoints, - vec!["wss://relay.example.com"] - ); - } - - #[test] - fn sync_pull_empty_multi_target_profile_reports_nostr_misconfigured() { - let dir = tempdir().expect("tempdir"); - let config = multi_target_config(dir.path(), Vec::new()); - crate::runtime::store::init(&config).expect("store init"); - - let view = pull_with_fetcher(&config, |_, _| { - panic!("empty MultiTarget profile sync pull must not fetch") - }) - .expect("sync pull empty MultiTarget profile"); - - assert_eq!(view.state, "unavailable"); - assert_eq!(view.configured_transport_target_count, 1); - assert_eq!(view.configured_transport_targets.len(), 1); - assert_eq!( - view.configured_transport_targets[0].transport_kind, - "reticulum" - ); - assert_eq!(view.transport_statuses.len(), 2); - assert_eq!(view.transport_statuses[0].transport, "nostr"); - assert_eq!( - view.transport_statuses[0].profile_id.as_deref(), - Some("multi_target") - ); - assert!(!view.transport_statuses[0].configured); - assert_eq!(view.transport_statuses[0].implementation, "real"); - assert!(!view.transport_statuses[0].usable_for_delivery); - assert!(!view.transport_statuses[0].capabilities.deliver); - assert!(!view.transport_statuses[0].capabilities.fetch); - assert_eq!(view.transport_statuses[1].transport, "reticulum"); - assert!(view.transport_statuses[1].configured); - assert_eq!(view.transport_statuses[1].implementation, "real"); - assert_eq!(view.transport_statuses[1].maturity, "preview"); - assert_eq!(view.transport_statuses[1].availability, "unavailable"); - assert!(!view.transport_statuses[1].usable_for_delivery); - assert!(!view.transport_statuses[1].capabilities.deliver); - assert!(!view.transport_statuses[1].capabilities.fetch); - assert!(view.target_transport_endpoints.is_empty()); - assert_eq!(view.actions, vec!["radroots transport config inspect"]); - } - - #[test] - fn sync_inspect_multi_target_uses_profile_target_count() { - let dir = tempdir().expect("tempdir"); - let config = multi_target_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - - let snapshot = inspect_sync(&config).expect("inspect sync multi_target"); - - assert_eq!(snapshot.state, "ready"); - assert_eq!(snapshot.configured_transport_target_count, 2); - assert_eq!(snapshot.configured_transport_targets.len(), 2); - assert_eq!(snapshot.transport_statuses.len(), 2); - assert_eq!(snapshot.transport_statuses[1].transport, "reticulum"); - assert_eq!(snapshot.transport_statuses[1].implementation, "real"); - assert_eq!(snapshot.transport_statuses[1].maturity, "preview"); - assert_eq!(snapshot.transport_statuses[1].availability, "unavailable"); - } - - #[test] - fn sync_pull_reticulum_does_not_report_fetch_success() { - let dir = tempdir().expect("tempdir"); - let config = reticulum_config(dir.path()); - crate::runtime::store::init(&config).expect("store init"); - - let view = pull_with_fetcher(&config, |_, _| { - panic!("reticulum sync pull must not run Nostr fetch") - }) - .expect("sync pull reticulum"); - - assert_eq!(view.state, "unavailable"); - assert_eq!(view.configured_transport_target_count, 1); - assert_eq!( - view.configured_transport_targets[0].transport_kind, - "reticulum" - ); - assert!(view.transport_statuses[0].configured); - assert_eq!(view.transport_statuses[0].implementation, "real"); - assert_eq!(view.transport_statuses[0].maturity, "preview"); - assert_eq!(view.transport_statuses[0].availability, "unavailable"); - assert!(!view.transport_statuses[0].usable_for_delivery); - assert!(!view.transport_statuses[0].capabilities.deliver); - assert!(!view.transport_statuses[0].capabilities.fetch); - assert_eq!(view.fetched_count, None); - assert!(view.target_transport_endpoints.is_empty()); - assert!( - view.reason - .as_deref() - .expect("reticulum unavailable reason") - .contains("does not expose Nostr relay fetch targets") - ); - } - - #[test] - fn sync_status_empty_sdk_store_reports_canonical_source() { - let dir = tempdir().expect("tempdir"); - let config = sample_config( - dir.path(), - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned(), - ], - ); - - let view = sdk_sync_status_view( - &config, - sdk_status_receipt( - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - None, - None, - &["wss://relay-a.example.com", "wss://relay-b.example.com"], - ), - ); - - assert_eq!(view.state, "ready"); - assert_eq!(view.source, "SDK canonical event store and outbox"); - assert_eq!(view.replica_store, "derived_projection_not_checked"); - assert_eq!(view.configured_transport_target_count, 2); - assert_eq!(view.queue.total_count, Some(0)); - assert_eq!(view.queue.pending_count, 0); - assert_eq!(view.queue.retryable_count, Some(0)); - assert_eq!(view.queue.terminal_count, Some(0)); - assert_eq!(view.queue.deferred_until_implemented_count, Some(0)); - assert_eq!(view.queue.ready_signed_count, Some(0)); - assert_eq!(view.actions, vec!["radroots sync pull"]); - } - - #[test] - fn sync_status_multi_target_reports_profile_targets_and_statuses() { - let dir = tempdir().expect("tempdir"); - let config = multi_target_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - let mut receipt = sdk_status_receipt( - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - None, - None, - &["wss://relay.example.com"], - ); - receipt.transport_profile = SyncTransportProfileSummary { - transport_profile_id: "multi_target".to_owned(), - configured_transport_target_count: 2, - configured_transport_targets: vec![ - SyncTransportTargetSummary { - transport_kind: "nostr".to_owned(), - endpoint_uri: "wss://relay.example.com".to_owned(), - endpoint_fingerprint: "0".repeat(64), - target_scope: None, - target_label: None, - }, - SyncTransportTargetSummary { - transport_kind: "reticulum".to_owned(), - endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), - endpoint_fingerprint: "1".repeat(64), - target_scope: Some("local".to_owned()), - target_label: None, - }, - ], - transport_statuses: vec![ - SyncTransportStatusSummary { - transport: "nostr".to_owned(), - profile_id: Some("multi_target".to_owned()), - endpoint_uri: None, - configured: true, - implementation: "real".to_owned(), - maturity: "stable".to_owned(), - availability: "available".to_owned(), - usable_for_delivery: true, - capabilities: SyncTransportOperationCapabilitiesSummary { - deliver: true, - fetch: false, - }, - message: "Nostr relay transport is configured for delivery".to_owned(), - }, - SyncTransportStatusSummary { - transport: "reticulum".to_owned(), - profile_id: Some("multi_target".to_owned()), - endpoint_uri: Some(RADROOTS_RETICULUM_ENDPOINT_URI.to_owned()), - configured: true, - implementation: "real".to_owned(), - maturity: "preview".to_owned(), - availability: "unavailable".to_owned(), - usable_for_delivery: false, - capabilities: SyncTransportOperationCapabilitiesSummary { - deliver: false, - fetch: false, - }, - message: RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned(), - }, - ], - }; - - let view = sdk_sync_status_view(&config, receipt); - - assert_eq!(view.configured_transport_target_count, 2); - assert_eq!(view.configured_transport_targets.len(), 2); - assert_eq!(view.configured_transport_targets[0].transport_kind, "nostr"); - assert_eq!( - view.configured_transport_targets[1].transport_kind, - "reticulum" - ); - assert_eq!(view.transport_statuses.len(), 2); - assert_eq!(view.transport_statuses[1].implementation, "real"); - assert_eq!(view.transport_statuses[1].maturity, "preview"); - assert_eq!(view.transport_statuses[1].availability, "unavailable"); - assert!(!view.transport_statuses[1].usable_for_delivery); - assert!(!view.transport_statuses[1].capabilities.deliver); - assert!(!view.transport_statuses[1].capabilities.fetch); - } - - #[test] - fn sync_status_reports_distinct_raw_valid_and_outbox_counts() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - - let receipt = sdk_status_receipt( - 3, - 2, - 4, - 1, - 1, - 2, - 1, - 0, - 1, - 0, - Some(1_700_000_010_000), - Some("auth-required: login".to_owned()), - &["wss://relay.example.com"], - ); - assert_eq!(receipt.event_store.total_events, 3); - assert_eq!(receipt.event_store.valid_stream_events, 2); - - let view = sdk_sync_status_view(&config, receipt); - - assert_eq!(view.state, "ready"); - assert_eq!(view.queue.expected_count, 4); - assert_eq!(view.queue.pending_count, 2); - assert_eq!(view.queue.retryable_count, Some(1)); - assert_eq!(view.queue.terminal_count, Some(2)); - assert_eq!(view.queue.failed_terminal_count, Some(1)); - assert_eq!(view.queue.deferred_until_implemented_count, Some(0)); - assert_eq!(view.queue.ready_signed_count, Some(1)); - assert_eq!(view.queue.last_attempt_at_ms, Some(1_700_000_010_000)); - assert_eq!( - view.queue.last_error.as_deref(), - Some("auth-required: login") - ); - assert_eq!(view.actions, vec!["radroots sync push"]); - } - - #[test] - fn sync_push_dry_run_reports_sdk_ready_outbox_plan() { - let dir = tempdir().expect("tempdir"); - let mut config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - config.output.dry_run = true; - - let view = sdk_push_dry_run_view( - &config, - sdk_status_receipt( - 1, - 1, - 1, - 1, - 0, - 0, - 0, - 0, - 1, - 0, - None, - None, - &["wss://relay.example.com"], - ), - ); - - assert_eq!(view.state, "dry_run"); - assert_eq!(view.source, "SDK outbox push"); - assert_eq!(view.replica_store, "derived_projection_not_checked"); - assert_eq!( - view.target_transport_endpoints, - vec!["wss://relay.example.com"] - ); - assert_eq!(view.publishable_count, Some(1)); - assert_eq!(view.published_count, Some(0)); - assert_eq!(view.failed_count, Some(0)); - assert_eq!(view.reason_code.as_deref(), Some("dry_run")); - assert_eq!( - view.reason.as_deref(), - Some("dry run requested; SDK outbox push skipped") - ); - assert_eq!( - view.actions, - vec!["radroots sync push", "radroots sync status"] - ); - assert!(view.publish_plan.is_none()); - } - - #[test] - fn sync_push_empty_queue_reports_ready_sdk_state() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - - let view = sdk_push_view( - &config, - PushOutboxReceipt::default(), - sdk_status_receipt( - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - 0, - None, - None, - &["wss://relay.example.com"], - ), - ); - - assert_eq!(view.state, "ready"); - assert_eq!(view.publishable_count, Some(0)); - assert_eq!(view.published_count, Some(0)); - assert_eq!(view.failed_count, Some(0)); - assert_eq!( - view.reason.as_deref(), - Some("SDK outbox had no ready signed events to push") - ); - assert_eq!(view.actions, vec!["radroots sync status"]); - } - - #[test] - fn sync_push_reticulum_reports_non_attempted_deferred_work() { - let cases = [ - ( - PushOutboxEventState::DeferredUntilImplemented, - PushOutboxTargetOutcomeKind::DeferredUntilImplemented, - "deferred_until_implemented", - "sdk_outbox_push_deferred_until_implemented", - "SDK outbox push reported Reticulum work as deferred until implemented without network delivery", - 0, - ), - ( - PushOutboxEventState::DeferredUntilImplemented, - PushOutboxTargetOutcomeKind::DeferredUntilImplemented, - "deferred_until_implemented", - "sdk_outbox_push_deferred_until_implemented", - "SDK outbox push reported Reticulum work as deferred until implemented without network delivery", - 1, - ), - ]; - - for ( - final_state, - outcome_kind, - expected_state, - expected_reason_code, - expected_reason, - deferred_count, - ) in cases - { - let dir = tempdir().expect("tempdir"); - let config = reticulum_config(dir.path()); - let receipt = PushOutboxReceipt { - attempted_events: 0, - published_events: 0, - retryable_events: 0, - terminal_events: 0, - events: vec![sdk_reticulum_push_event(final_state, outcome_kind)], - }; - - let view = sdk_push_view( - &config, - receipt, - sdk_reticulum_status_receipt(deferred_count), - ); - - assert_eq!(view.state, expected_state); - assert_eq!(view.publishable_count, Some(0)); - assert_eq!(view.published_count, Some(0)); - assert_eq!(view.failed_count, Some(0)); - assert_eq!(view.reason_code.as_deref(), Some(expected_reason_code)); - assert_eq!(view.reason.as_deref(), Some(expected_reason)); - assert_eq!( - view.target_transport_endpoints, - vec![RADROOTS_RETICULUM_ENDPOINT_URI.to_owned()] - ); - assert!(view.attempted_transport_endpoints.is_empty()); - assert!(view.accepted_transport_endpoints.is_empty()); - assert_eq!(view.failed_transport_targets.len(), 1); - assert_eq!(view.failed_transport_targets[0].transport_kind, "reticulum"); - assert_eq!( - view.failed_transport_targets[0].endpoint_uri, - RADROOTS_RETICULUM_ENDPOINT_URI - ); - assert_eq!( - view.failed_transport_targets[0].reason, - RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE - ); - assert_eq!(view.actions, vec!["radroots sync status"]); - } - } - - #[test] - fn sync_push_maps_published_and_auth_required_sdk_receipts() { - let dir = tempdir().expect("tempdir"); - let config = sample_config( - dir.path(), - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned(), - ], - ); - let receipt = PushOutboxReceipt { - attempted_events: 2, - published_events: 1, - retryable_events: 1, - terminal_events: 0, - events: vec![ - sdk_push_event( - "a", - PushOutboxEventState::Published, - PushOutboxTargetOutcomeKind::Accepted, - "wss://relay-a.example.com", - Some("accepted".to_owned()), - ), - sdk_push_event( - "b", - PushOutboxEventState::PublishRetryable, - PushOutboxTargetOutcomeKind::AuthRequired, - "wss://relay-b.example.com", - Some("auth-required: login".to_owned()), - ), - ], - }; - - let view = sdk_push_view( - &config, - receipt, - sdk_status_receipt( - 2, - 2, - 2, - 0, - 1, - 1, - 0, - 0, - 0, - 0, - Some(1_700_000_020_000), - Some("auth-required: login".to_owned()), - &["wss://relay-a.example.com", "wss://relay-b.example.com"], - ), - ); - - assert_eq!(view.state, "partial"); - assert_eq!(view.publishable_count, Some(2)); - assert_eq!(view.published_count, Some(1)); - assert_eq!(view.failed_count, Some(1)); - assert_eq!(view.reason_code.as_deref(), Some("sdk_outbox_push_partial")); - assert_eq!( - view.target_transport_endpoints, - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned() - ] - ); - assert_eq!( - view.attempted_transport_endpoints, - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned() - ] - ); - assert_eq!( - view.accepted_transport_endpoints, - vec!["wss://relay-a.example.com".to_owned()] - ); - assert_eq!(view.failed_transport_targets.len(), 1); - assert_eq!( - view.failed_transport_targets[0].endpoint_uri, - "wss://relay-b.example.com" - ); - assert_eq!(view.failed_transport_targets[0].transport_kind, "nostr"); - assert_eq!( - view.failed_transport_targets[0].reason, - "auth-required: login" - ); - assert_eq!( - view.actions, - vec!["radroots sync push", "radroots sync status"] - ); - } - - #[expect( - clippy::too_many_arguments, - reason = "test fixture construction sets every storage status field explicitly" - )] - fn sdk_status_receipt( - raw_total_events: i64, - valid_stream_events: i64, - outbox_total_events: i64, - pending_events: i64, - retryable_events: i64, - terminal_events: i64, - failed_terminal_events: i64, - deferred_until_implemented_events: i64, - ready_signed_events: i64, - publishing_events: i64, - last_attempt_at_ms: Option<i64>, - last_error: Option<String>, - relays: &[&str], - ) -> SyncStatusReceipt { - assert!( - !relays.is_empty(), - "sdk_status_receipt must not build ready Nostr status without transport targets" - ); - SyncStatusReceipt { - source: SyncStatusSource::SdkCanonicalStores, - observed_at_ms: 1_700_000_030_000, - event_store: SyncEventStoreStatus { - total_events: raw_total_events, - valid_stream_events, - transport_observations: 0, - last_event_seq: (raw_total_events > 0).then_some(raw_total_events), - last_event_updated_at_ms: (raw_total_events > 0).then_some(1_700_000_000_000), - }, - outbox: SyncOutboxStatus { - total_events: outbox_total_events, - pending_events, - retryable_events, - terminal_events, - failed_terminal_events, - deferred_until_implemented_events, - ready_signed_events, - publishing_events, - last_attempt_at_ms, - last_error, - }, - transport_profile: SyncTransportProfileSummary { - transport_profile_id: "nostr".to_owned(), - configured_transport_target_count: relays.len(), - configured_transport_targets: relays - .iter() - .enumerate() - .map(|(index, relay)| SyncTransportTargetSummary { - transport_kind: "nostr".to_owned(), - endpoint_uri: (*relay).to_owned(), - endpoint_fingerprint: format!("test-fingerprint-{index}"), - target_scope: None, - target_label: None, - }) - .collect(), - transport_statuses: vec![SyncTransportStatusSummary { - transport: "nostr".to_owned(), - profile_id: Some("nostr".to_owned()), - endpoint_uri: None, - configured: true, - implementation: "real".to_owned(), - maturity: "stable".to_owned(), - availability: "available".to_owned(), - usable_for_delivery: true, - capabilities: SyncTransportOperationCapabilitiesSummary { - deliver: true, - fetch: false, - }, - message: "Nostr relay transport is configured for delivery".to_owned(), - }], - }, - } - } - - fn sdk_reticulum_status_receipt(deferred_until_implemented_events: i64) -> SyncStatusReceipt { - SyncStatusReceipt { - source: SyncStatusSource::SdkCanonicalStores, - observed_at_ms: 1_700_000_030_000, - event_store: SyncEventStoreStatus { - total_events: 1, - valid_stream_events: 1, - transport_observations: 0, - last_event_seq: Some(1), - last_event_updated_at_ms: Some(1_700_000_000_000), - }, - outbox: SyncOutboxStatus { - total_events: 1, - pending_events: 0, - retryable_events: 0, - terminal_events: 0, - failed_terminal_events: 0, - deferred_until_implemented_events, - ready_signed_events: 0, - publishing_events: 0, - last_attempt_at_ms: None, - last_error: None, - }, - transport_profile: SyncTransportProfileSummary { - transport_profile_id: "reticulum".to_owned(), - configured_transport_target_count: 1, - configured_transport_targets: vec![SyncTransportTargetSummary { - transport_kind: "reticulum".to_owned(), - endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), - endpoint_fingerprint: "1".repeat(64), - target_scope: Some("local".to_owned()), - target_label: None, - }], - transport_statuses: vec![SyncTransportStatusSummary { - transport: "reticulum".to_owned(), - profile_id: Some("reticulum".to_owned()), - endpoint_uri: Some(RADROOTS_RETICULUM_ENDPOINT_URI.to_owned()), - configured: true, - implementation: "real".to_owned(), - maturity: "preview".to_owned(), - availability: "unavailable".to_owned(), - usable_for_delivery: false, - capabilities: SyncTransportOperationCapabilitiesSummary { - deliver: false, - fetch: false, - }, - message: RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned(), - }], - }, - } - } - - fn sdk_push_event( - event_id_prefix: &str, - final_state: PushOutboxEventState, - outcome_kind: PushOutboxTargetOutcomeKind, - relay_url: &str, - message: Option<String>, - ) -> PushOutboxEventReceipt { - PushOutboxEventReceipt { - event_id: RadrootsEventId::parse(event_id_prefix.repeat(64).as_str()) - .expect("event id"), - outbox_event_id: 7, - final_state, - attempted_count: 1, - accepted_count: usize::from(matches!( - outcome_kind, - PushOutboxTargetOutcomeKind::Accepted - | PushOutboxTargetOutcomeKind::DuplicateAccepted - )), - retryable_count: usize::from(matches!( - outcome_kind, - PushOutboxTargetOutcomeKind::AuthRequired - | PushOutboxTargetOutcomeKind::Timeout - | PushOutboxTargetOutcomeKind::ConnectionFailed - )), - terminal_count: usize::from(matches!( - outcome_kind, - PushOutboxTargetOutcomeKind::Blocked - | PushOutboxTargetOutcomeKind::RateLimited - | PushOutboxTargetOutcomeKind::Invalid - | PushOutboxTargetOutcomeKind::PowRequired - | PushOutboxTargetOutcomeKind::Restricted - | PushOutboxTargetOutcomeKind::Error - | PushOutboxTargetOutcomeKind::Unknown - )), - quorum: 1, - quorum_met: matches!( - outcome_kind, - PushOutboxTargetOutcomeKind::Accepted - | PushOutboxTargetOutcomeKind::DuplicateAccepted - ), - targets: vec![PushOutboxTargetReceipt { - transport_kind: "nostr".to_owned(), - endpoint_uri: relay_url.to_owned(), - target_scope: None, - target_label: None, - outcome_kind, - transport_outcome_kind: test_transport_outcome_kind(outcome_kind), - attempted: true, - message, - }], - } - } - - fn sdk_reticulum_push_event( - final_state: PushOutboxEventState, - outcome_kind: PushOutboxTargetOutcomeKind, - ) -> PushOutboxEventReceipt { - PushOutboxEventReceipt { - event_id: RadrootsEventId::parse("c".repeat(64).as_str()).expect("event id"), - outbox_event_id: 9, - final_state, - attempted_count: 0, - accepted_count: 0, - retryable_count: 0, - terminal_count: 0, - quorum: 1, - quorum_met: false, - targets: vec![PushOutboxTargetReceipt { - transport_kind: "reticulum".to_owned(), - endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), - target_scope: Some("local".to_owned()), - target_label: None, - outcome_kind, - transport_outcome_kind: test_transport_outcome_kind(outcome_kind), - attempted: false, - message: Some(RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned()), - }], + fn ready_action(self) -> &'static str { + match self { + Self::SyncPull => "radroots market search eggs", + Self::MarketPull => "radroots market search eggs", } } +} - fn test_transport_outcome_kind( - kind: PushOutboxTargetOutcomeKind, - ) -> Option<PushOutboxTransportOutcomeKind> { - Some(match kind { - PushOutboxTargetOutcomeKind::Accepted => PushOutboxTransportOutcomeKind::Accepted, - PushOutboxTargetOutcomeKind::DuplicateAccepted => { - PushOutboxTransportOutcomeKind::DuplicateAccepted - } - PushOutboxTargetOutcomeKind::DeferredUntilImplemented => { - PushOutboxTransportOutcomeKind::DeferredUntilImplemented - } - PushOutboxTargetOutcomeKind::Timeout => PushOutboxTransportOutcomeKind::Timeout, - PushOutboxTargetOutcomeKind::ConnectionFailed => { - PushOutboxTransportOutcomeKind::ConnectionFailed - } - PushOutboxTargetOutcomeKind::AuthRequired - | PushOutboxTargetOutcomeKind::Blocked - | PushOutboxTargetOutcomeKind::RateLimited - | PushOutboxTargetOutcomeKind::Invalid - | PushOutboxTargetOutcomeKind::PowRequired - | PushOutboxTargetOutcomeKind::Restricted - | PushOutboxTargetOutcomeKind::Muted - | PushOutboxTargetOutcomeKind::Unsupported - | PushOutboxTargetOutcomeKind::PaymentRequired - | PushOutboxTargetOutcomeKind::Error - | PushOutboxTargetOutcomeKind::TargetUriRejected - | PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted - | PushOutboxTargetOutcomeKind::Unknown => PushOutboxTransportOutcomeKind::Rejected, - _ => PushOutboxTransportOutcomeKind::Rejected, - }) - } - - #[test] - fn sync_pull_ingests_relay_events_and_market_reads_without_daemon() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(7); - let seller_pubkey = seller.public_key_hex(); - let listing_addr = format!("{KIND_CLASSIFIED_LISTING}:{seller_pubkey}:{LISTING_D_TAG}"); - let events = vec![ - farm_event(&seller), - plot_event(&seller), - listing_event(&seller), - list_set_event(&seller), - ]; - - let view = pull_with_fetcher(&config, fake_fetcher(events)).expect("sync pull ingest"); - - assert_eq!(view.state, "ready"); - assert_eq!(view.fetched_count, Some(4)); - assert_eq!(view.ingested_count, Some(4)); - assert_eq!(view.skipped_count, Some(0)); - assert_eq!(view.unsupported_count, Some(0)); - assert_eq!(view.failed_count, Some(0)); - assert_eq!(view.reason, None); - - let search = crate::runtime::find::search( - &config, - &FindQueryArgs { - query: vec!["eggs".to_owned()], - }, - ) - .expect("market search"); - assert_eq!(search.state, "ready"); - assert_eq!(search.count, 1); - assert_eq!( - search.results[0].listing_addr.as_deref(), - Some(listing_addr.as_str()) - ); - - let listing = crate::runtime::listing::get( - &config, - &RecordLookupArgs { - key: "pasture-eggs".to_owned(), - }, - ) - .expect("listing get"); - assert_eq!(listing.state, "ready"); - assert_eq!(listing.listing_addr.as_deref(), Some(listing_addr.as_str())); - } - - #[test] - fn sync_pull_keeps_tagless_profile_out_of_legacy_replica_projection() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(17); - - let view = pull_with_fetcher(&config, fake_fetcher(vec![profile_event(&seller)])) - .expect("sync pull Profile"); - - assert_eq!(view.state, "ready"); - assert_eq!(view.fetched_count, Some(1)); - assert_eq!(view.ingested_count, Some(0)); - assert_eq!(view.skipped_count, Some(0)); - assert_eq!(view.unsupported_count, Some(0)); - assert_eq!(view.failed_count, Some(1)); - assert_eq!(view.reason_code.as_deref(), Some("sync_ingest_failed")); - assert!( - view.reason - .as_deref() - .is_some_and(|reason| reason.contains("profile_type required")) - ); - } - - #[test] - fn market_refresh_uses_market_scope_for_ingest() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(8); - let events = vec![listing_event(&seller), plot_event(&seller)]; - - let view = - market_refresh_with_fetcher(&config, fake_fetcher(events)).expect("market refresh"); +fn unix_now() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |duration| duration.as_secs()) +} - assert_eq!(view.state, "ready"); - assert_eq!(view.fetched_count, Some(2)); - assert_eq!(view.ingested_count, Some(1)); - assert_eq!(view.unsupported_count, Some(1)); - assert_eq!(view.failed_count, Some(0)); - } +#[cfg(test)] +mod tests { + use super::*; #[test] - fn market_refresh_records_relay_provenance_relays_for_order_drafts() { - let dir = tempdir().expect("tempdir"); - let config = sample_config( - dir.path(), - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned(), - ], - ); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(9); - - let _ = market_refresh_with_fetcher(&config, fake_fetcher(vec![listing_event(&seller)])) - .expect("market refresh"); - let relays = relay_provenance_relays_for_scope(&config, RelayIngestScope::MarketPull) - .expect("relay provenance"); - + fn receipt_labels_cover_partial_failure_and_cancellation() { assert_eq!( - relays, - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned() - ] - ); - } - - #[test] - fn relay_refresh_records_current_run_freshness() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(10); - - let view = market_refresh_with_fetcher(&config, fake_fetcher(vec![listing_event(&seller)])) - .expect("market refresh"); - - assert_eq!(view.freshness.state, "fresh"); - let run = view.freshness.run.as_ref().expect("run freshness"); - assert_eq!(run.scope, "market_refresh"); - assert_eq!(run.last_state, "success"); - assert!(run.relay_set_current); - assert_eq!(run.fetched_count, Some(1)); - assert_eq!(run.ingested_count, Some(1)); - } - - #[test] - fn sync_pull_reports_partial_relay_fetch_reason_code() { - let dir = tempdir().expect("tempdir"); - let config = sample_config( - dir.path(), - vec![ - "wss://relay-a.example.com".to_owned(), - "wss://relay-b.example.com".to_owned(), - ], + pull_termination_label(PullTermination::Cancelled), + "cancelled" ); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(13); - - let view = pull_with_fetcher(&config, |relays, _| { - Ok(relay_fetch_receipt_with_failed( - relays, - vec![listing_event(&seller)], - vec![RadrootsRelayFetchFailure { - relay_url: relays[1].clone(), - reason: "connection refused".to_owned(), - }], - )) - }) - .expect("sync pull partial relay fetch"); - - assert_eq!(view.state, "ready"); assert_eq!( - view.attempted_transport_endpoints, - vec!["wss://relay-a.example.com"] - ); - assert_eq!(view.failed_transport_targets.len(), 1); - assert_eq!(view.failed_count, Some(1)); - assert_eq!(view.reason_code.as_deref(), Some("nostr_fetch_partial")); - assert!( - view.reason - .as_deref() - .expect("partial relay reason") - .contains("relay(s) failed during fetch") - ); - let run = view.freshness.run.as_ref().expect("run freshness"); - assert_eq!(run.last_state, "partial"); - assert_eq!(run.failed_count, Some(1)); - } - - #[test] - fn sync_pull_reports_no_overwrite_skips_without_replacing_projection() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(12); - - let first = listing_event_with_title_at(&seller, "Pasture Eggs", 200); - let stale = listing_event_with_title_at(&seller, "Older Eggs", 199); - pull_with_fetcher(&config, fake_fetcher(vec![first])).expect("initial sync pull"); - let view = pull_with_fetcher(&config, fake_fetcher(vec![stale])).expect("stale sync pull"); - - assert_eq!(view.state, "ready"); - assert_eq!(view.fetched_count, Some(1)); - assert_eq!(view.ingested_count, Some(0)); - assert_eq!(view.skipped_count, Some(1)); - assert_eq!(view.reason_code.as_deref(), Some("sync_no_overwrite")); - assert!( - view.reason - .as_deref() - .expect("skip reason") - .contains("current or newer state") + fetch_state_label(FetchTargetState::FailedRetryable), + "failed_retryable" ); - let run = view.freshness.run.as_ref().expect("run freshness"); - assert_eq!(run.last_state, "success"); - assert_eq!(run.skipped_count, Some(1)); - - let search = crate::runtime::find::search( - &config, - &FindQueryArgs { - query: vec!["eggs".to_owned()], - }, - ) - .expect("market search"); - assert_eq!(search.results[0].title, "Pasture Eggs"); - } - - #[test] - fn sync_pull_freshness_reports_relay_set_changed() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay-a.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(11); - pull_with_fetcher(&config, fake_fetcher(vec![listing_event(&seller)])).expect("sync pull"); - let changed = sample_config(dir.path(), vec!["wss://relay-b.example.com".to_owned()]); - - let freshness = - freshness_for_scope(&changed, RelayIngestScope::SyncPull).expect("sync freshness"); - - assert_eq!(freshness.state, "relay_set_changed"); - let run = freshness.run.as_ref().expect("run freshness"); - assert_eq!(run.scope, "sync_pull"); - assert!(!run.relay_set_current); + assert_eq!(fetch_state_label(FetchTargetState::Partial), "partial"); } #[test] - fn relay_ingest_splits_unsupported_and_failed_events() { - let dir = tempdir().expect("tempdir"); - let config = sample_config(dir.path(), vec!["wss://relay.example.com".to_owned()]); - crate::runtime::store::init(&config).expect("store init"); - let seller = identity(9); - let events = vec![ - signed_event( - &seller, - Nip01EventWireParts { - kind: KIND_CLASSIFIED_LISTING, - content: "x".repeat(DEFAULT_CONTENT_MAX_BYTES + 1), - tags: Vec::new(), - }, - ), - signed_event( - &seller, - Nip01EventWireParts { - kind: KIND_POST, - content: "hello".to_owned(), - tags: Vec::new(), - }, - ), - signed_event( - &seller, - Nip01EventWireParts { - kind: KIND_CLASSIFIED_LISTING, - content: "not a listing".to_owned(), - tags: Vec::new(), - }, - ), - ]; - - let view = pull_with_fetcher(&config, fake_fetcher(events)).expect("sync pull ingest"); - - assert_eq!(view.state, "ready"); - assert_eq!(view.fetched_count, Some(3)); - assert_eq!(view.ingested_count, Some(0)); - assert_eq!(view.unsupported_count, Some(1)); - assert_eq!(view.failed_count, Some(2)); - assert!( - view.reason - .as_deref() - .expect("failure reason") - .contains("event envelope content size") - ); - } - - fn fake_fetcher( - events: Vec<RadrootsNostrEvent>, - ) -> impl FnOnce( - &[String], - RadrootsNostrFilter, - ) -> 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(), - duplicate_count: 0, - verification_failed_count: 0, - malformed_count: 0, - out_of_filter_count: 0, - skipped_over_limit_count: 0, - eose_count: connected_relays.len(), - truncated_count: 0, - closed_count, - notice_count: 0, - relay_outcomes: Vec::new(), - } - } - - fn profile_event(identity: &RadrootsIdentity) -> RadrootsNostrEvent { - let profile = AuthoredProfile::new("seller") - .expect("profile") - .with_display_name("Seller") - .with_about("market seller"); - signed_event( - identity, - authored_profile_to_wire_parts(&profile).expect("profile parts"), - ) - } - - fn farm_event(identity: &RadrootsIdentity) -> RadrootsNostrEvent { - let farm = Farm { - d_tag: FARM_D_TAG.to_owned(), - name: "Relay Farm".to_owned(), - about: Some("relay farm".to_owned()), - website: Some("https://farm.example.com".to_owned()), - picture: None, - banner: None, - location: None, - tags: None, - }; - signed_event( - identity, - farm_encode::to_wire_parts(&farm).expect("farm parts"), - ) - } - - fn plot_event(identity: &RadrootsIdentity) -> RadrootsNostrEvent { - let plot = RadrootsPlot { - d_tag: PLOT_D_TAG.to_owned(), - farm: FarmRef { - pubkey: identity.public_key_hex(), - d_tag: FARM_D_TAG.to_owned(), - }, - name: "Relay Plot".to_owned(), - about: Some("relay plot".to_owned()), - location: None, - tags: None, - }; - signed_event( - identity, - plot_encode::to_wire_parts(&plot).expect("plot parts"), - ) - } - - fn list_set_event(identity: &RadrootsIdentity) -> RadrootsNostrEvent { - let list_set = RadrootsListSet { - d_tag: "member_of.farms".to_owned(), - content: String::new(), - entries: vec![RadrootsListEntry { - tag: "p".to_owned(), - values: vec![identity.public_key_hex()], - }], - title: None, - description: None, - image: None, - }; - signed_event( - identity, - list_set_encode::to_wire_parts_with_kind(&list_set, KIND_LIST_SET_GENERIC) - .expect("list set parts"), - ) - } - - fn listing_event(identity: &RadrootsIdentity) -> RadrootsNostrEvent { - listing_event_with_title_at(identity, "Pasture Eggs", 0) - } - - fn listing_event_with_title_at( - identity: &RadrootsIdentity, - title: &str, - created_at: u64, - ) -> RadrootsNostrEvent { - let mut builder = wire_fixture_builder(Nip01EventWireParts { - kind: KIND_CLASSIFIED_LISTING, - content: "# Pasture Eggs".to_owned(), - tags: vec![ - vec!["d".to_owned(), LISTING_D_TAG.to_owned()], - vec![ - "a".to_owned(), - format!("{}:{}:{}", KIND_FARM, identity.public_key_hex(), FARM_D_TAG), - ], - vec!["p".to_owned(), identity.public_key_hex()], - vec!["key".to_owned(), "pasture-eggs".to_owned()], - vec!["title".to_owned(), title.to_owned()], - vec!["category".to_owned(), "eggs".to_owned()], - vec!["summary".to_owned(), "Pasture-raised eggs".to_owned()], - vec!["process".to_owned(), "washed".to_owned()], - vec!["lot".to_owned(), "lot-a".to_owned()], - vec!["profile".to_owned(), "dozen".to_owned()], - vec!["year".to_owned(), "2026".to_owned()], - vec!["radroots:primary_bin".to_owned(), "bin-a".to_owned()], - vec![ - "radroots:bin".to_owned(), - "bin-a".to_owned(), - "12".to_owned(), - "each".to_owned(), - "12".to_owned(), - "each".to_owned(), - "dozen".to_owned(), - ], - vec![ - "radroots:price".to_owned(), - "bin-a".to_owned(), - "6".to_owned(), - "USD".to_owned(), - "1".to_owned(), - "each".to_owned(), - "6".to_owned(), - "each".to_owned(), - ], - vec!["inventory".to_owned(), "5".to_owned()], - vec!["status".to_owned(), "active".to_owned()], - ], - }); - if created_at > 0 { - builder = builder.custom_created_at(RadrootsNostrTimestamp::from(created_at)); - } - builder - .sign_with_keys(identity.keys()) - .expect("signed event") - } - - fn signed_event(identity: &RadrootsIdentity, parts: Nip01EventWireParts) -> RadrootsNostrEvent { - wire_fixture_builder(parts) - .sign_with_keys(identity.keys()) - .expect("signed event") - } - - fn wire_fixture_builder(parts: Nip01EventWireParts) -> nostr::EventBuilder { - let kind = u16::try_from(parts.kind).expect("fixture kind must fit NIP-01"); - let tags = parts - .tags - .into_iter() - .filter(|tag| !tag.is_empty()) - .map(|tag| nostr::Tag::parse(tag).expect("fixture tag must parse")) - .collect::<Vec<_>>(); - nostr::EventBuilder::new(nostr::Kind::Custom(kind), parts.content) - .tags(tags) - .allow_self_tagging() - } - - fn identity(seed: u8) -> RadrootsIdentity { - RadrootsIdentity::from_secret_key_bytes(&[seed; 32]).expect("identity") - } - - fn sample_config(root: &Path, relays: Vec<String>) -> RuntimeConfig { - let data = root.join("data"); - let cache = root.join("cache"); - let logs = root.join("logs"); - let secrets = root.join("secrets"); - RuntimeConfig { - output: OutputConfig { - format: OutputFormat::Terminal, - verbosity: Verbosity::Normal, - dry_run: false, - }, - interaction: InteractionConfig { - input_enabled: true, - assume_yes: false, - stdin_tty: false, - stdout_tty: false, - prompts_allowed: false, - confirmations_allowed: false, - }, - paths: PathsConfig { - profile: "interactive_user".into(), - profile_source: "test".into(), - allowed_profiles: vec!["interactive_user".into(), "repo_local".into()], - root_source: "test".into(), - repo_local_root: None, - repo_local_root_source: None, - subordinate_path_override_source: "runtime_config".into(), - app_namespace: "apps/cli".into(), - shared_accounts_namespace: "shared/accounts".into(), - shared_identities_namespace: "shared/identities".into(), - app_config_path: root.join("config/apps/cli/config.toml"), - workspace_config_path: None, - app_data_root: data.join("apps/cli"), - shared_cache_root: cache.clone(), - app_logs_root: logs.join("apps/cli"), - shared_accounts_data_root: data.join("shared/accounts"), - shared_accounts_secrets_root: secrets.join("shared/accounts"), - default_identity_path: secrets.join("shared/identities/default.json"), - }, - logging: LoggingConfig { - filter: "info".into(), - directory: None, - stdout: false, - }, - account: AccountConfig { - selector: None, - store_path: data.join("shared/accounts/store.json"), - secrets_dir: secrets.join("shared/accounts"), - secret_backend: RadrootsSecretBackend::EncryptedFile, - }, - account_secret_contract: AccountSecretContractConfig { - default_backend: "host_vault".into(), - allowed_backends: vec!["host_vault".into(), "encrypted_file".into()], - host_vault_policy: Some("desktop".into()), - uses_protected_store: true, - }, - identity: IdentityConfig { - path: secrets.join("shared/identities/default.json"), - }, - signer: SignerConfig { - backend: SignerBackend::Local, - }, - transport: crate::runtime::config::TransportConfig::from_nostr_relay_urls( - relays.clone(), - ), - local: LocalConfig { - root: data.join("apps/cli/replica"), - replica_store_path: data.join("apps/cli/replica/replica.sqlite"), - backups_dir: data.join("apps/cli/replica/backups"), - exports_dir: data.join("apps/cli/replica/exports"), - }, - myc: MycConfig { - executable: PathBuf::from("myc"), - status_timeout_ms: 2_000, - }, - hyf: HyfConfig { - enabled: false, - executable: PathBuf::from("hyfd"), - }, - mesh: crate::runtime::config::MeshConfig::disabled(), - rpc: RpcConfig { - url: "http://127.0.0.1:7070".into(), - }, - rhi: crate::runtime::config::RhiConfig { - validator_set: None, - require_cryptographic_proof: false, - }, - capability_bindings: Vec::new(), - } - } - - fn multi_target_config(root: &Path, relays: Vec<String>) -> RuntimeConfig { - let mut config = sample_config(root, relays); - config.transport.profile = TransportProfileKind::MultiTarget; - config.transport.reticulum_behavior = ReticulumBehavior::DeferDeliveryPlans; - config - } - - fn nostr_config(root: &Path, relays: Vec<String>) -> RuntimeConfig { - let mut config = sample_config(root, relays); - config.transport.profile = TransportProfileKind::Nostr; - config - } - - fn reticulum_config(root: &Path) -> RuntimeConfig { - let mut config = sample_config(root, Vec::new()); - config.transport.profile = TransportProfileKind::Reticulum; - config.transport.reticulum_behavior = ReticulumBehavior::DeferDeliveryPlans; - config + fn missing_freshness_requires_a_pull() { + assert!(freshness_requires_refresh(&missing_freshness())); } } diff --git a/src/runtime/trade.rs b/src/runtime/trade.rs @@ -793,7 +793,7 @@ fn resolve_active_listing_state( })?; state .last_event_id - .parse::<radroots_event::id::RadrootsEventId>() + .parse::<radroots_event::EventId>() .map_err(|error| { RuntimeError::Config(format!("listing latest event id is invalid: {error}")) })?; diff --git a/src/runtime/transport.rs b/src/runtime/transport.rs @@ -1,12 +1,6 @@ -use radroots_sdk::{ - PushOutboxEventState, PushOutboxReceipt, PushOutboxRequest, PushOutboxTargetOutcomeKind, - SyncStatusRequest, -}; use radroots_transport::{ - RADROOTS_RETICULUM_ENDPOINT_URI, RADROOTS_RETICULUM_SCOPE_ID, - RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, RadrootsTransportCapabilityAvailability, - RadrootsTransportCapabilityMaturity, RadrootsTransportImplementationState, - RadrootsTransportKind, RadrootsTransportStatus, + RadrootsTransportCapabilityAvailability, RadrootsTransportCapabilityMaturity, + RadrootsTransportImplementationState, RadrootsTransportKind, RadrootsTransportStatus, }; use serde_json::Value as JsonValue; use std::fs; @@ -15,7 +9,7 @@ use toml::{Value, map::Map}; use crate::ops::OperationData; use crate::runtime::RuntimeError; use crate::runtime::config::{RuntimeConfig, TransportProfileKind}; -use crate::runtime::sdk::{CliSdkAdapterError, CliSdkSession, sdk_nostr_relay_url_policy}; +use crate::runtime::sdk::{CliSdkAdapterError, CliSdkSession}; use crate::view::runtime::{ TransportDeliveryInspectView, TransportDeliveryRetryView, TransportOperationCapabilitiesView, TransportProfileSummaryView, TransportProfileView, TransportRuntimeStatusView, @@ -23,6 +17,10 @@ use crate::view::runtime::{ }; const TRANSPORT_SOURCE: &str = "transport profile config"; +const RADROOTS_RETICULUM_ENDPOINT_URI: &str = "reticulum:local"; +const RADROOTS_RETICULUM_SCOPE_ID: &str = "local"; +const RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE: &str = + "Reticulum transport is not available in this release"; pub fn profile(config: &RuntimeConfig) -> TransportProfileView { active_profile_view(config) } @@ -132,7 +130,8 @@ pub fn outbox_status( } else { CliSdkSession::connect_storage_status(config)? }; - let receipt = session.block_on(session.sdk().sync().status(SyncStatusRequest::new()))?; + let receipt = crate::runtime::sync::canonical_status(&session)?; + let outbox = receipt.outbox(); let state = if profile.configured_state == "configured" { "ready".to_owned() } else { @@ -142,15 +141,15 @@ pub fn outbox_status( state, source: "SDK transport outbox".to_owned(), transport_profile: profile.profile_id, - total_count: receipt.outbox.total_events, - pending_count: receipt.outbox.pending_events, - retryable_count: receipt.outbox.retryable_events, - terminal_count: receipt.outbox.terminal_events, - deferred_until_implemented_count: receipt.outbox.deferred_until_implemented_events, - ready_signed_count: receipt.outbox.ready_signed_events, - publishing_count: receipt.outbox.publishing_events, - last_attempt_at_ms: receipt.outbox.last_attempt_at_ms, - last_error: receipt.outbox.last_error, + total_count: i64::try_from(outbox.total().unwrap_or_default()).unwrap_or(i64::MAX), + pending_count: i64::try_from(outbox.pending).unwrap_or(i64::MAX), + retryable_count: i64::try_from(outbox.retryable).unwrap_or(i64::MAX), + terminal_count: i64::try_from(outbox.satisfied + outbox.exhausted).unwrap_or(i64::MAX), + deferred_until_implemented_count: 0, + ready_signed_count: i64::try_from(outbox.pending + outbox.retryable).unwrap_or(i64::MAX), + publishing_count: i64::try_from(outbox.leased).unwrap_or(i64::MAX), + last_attempt_at_ms: None, + last_error: None, actions: vec!["radroots transport outbox push".to_owned()], }) } @@ -173,75 +172,40 @@ pub fn outbox_push( }); } let session = CliSdkSession::connect(config)?; - let receipt = session.block_on(session.sdk().sync().push_outbox( - PushOutboxRequest::new().with_nostr_relay_url_policy(sdk_nostr_relay_url_policy(config)), - ))?; - let target_count = receipt - .events - .iter() - .flat_map(|event| event.targets.iter()) - .count(); - let failed_count = receipt.retryable_events + receipt.terminal_events; - let state = transport_outbox_push_state(&receipt, failed_count).to_owned(); + let receipt = crate::runtime::sync::deliver_pending(&session)?; + let attempted_events = receipt.outcomes().len(); + let published_events = receipt.succeeded(); + let failed_count = receipt.failed(); + let state = if attempted_events == 0 { + "ready" + } else if failed_count == 0 { + "published" + } else if published_events > 0 { + "partial" + } else { + "unavailable" + }; Ok(TransportDeliveryRetryView { - state, + state: state.to_owned(), source: "SDK transport outbox".to_owned(), - attempted_events: receipt.attempted_events, - published_events: receipt.published_events, - retryable_events: receipt.retryable_events, - terminal_events: receipt.terminal_events, - target_count, - reason: transport_outbox_push_reason(&receipt), + attempted_events, + published_events, + retryable_events: failed_count, + terminal_events: 0, + target_count: config.transport.nostr_relay_urls.len(), + reason: if attempted_events == 0 { + Some("canonical outbox had no ready delivery plans".to_owned()) + } else if failed_count > 0 { + Some(format!( + "{failed_count} canonical delivery outcome(s) failed" + )) + } else { + None + }, actions: vec!["radroots transport outbox status".to_owned()], }) } -fn transport_outbox_push_state(receipt: &PushOutboxReceipt, failed_count: usize) -> &'static str { - if receipt.attempted_events == 0 { - return transport_outbox_reported_deferred_state(receipt).unwrap_or("ready"); - } - if receipt.published_events > 0 && failed_count > 0 { - "partial" - } else if failed_count > 0 { - "unavailable" - } else if receipt.published_events > 0 { - "published" - } else { - "ready" - } -} - -fn transport_outbox_reported_deferred_state(receipt: &PushOutboxReceipt) -> Option<&'static str> { - let mut deferred = false; - for event in &receipt.events { - if event.final_state == PushOutboxEventState::DeferredUntilImplemented { - deferred = true; - } - for target in &event.targets { - if target.outcome_kind == PushOutboxTargetOutcomeKind::DeferredUntilImplemented { - deferred = true; - } - } - } - deferred.then_some("deferred_until_implemented") -} - -fn transport_outbox_push_reason(receipt: &PushOutboxReceipt) -> Option<String> { - if receipt.attempted_events == 0 { - if let Some(state) = transport_outbox_reported_deferred_state(receipt) { - return Some(match state { - "deferred_until_implemented" => { - "SDK outbox push reported Reticulum work as deferred until implemented without network delivery" - } - _ => "SDK outbox push reported Reticulum work without network delivery", - } - .to_owned()); - } - return Some("SDK outbox had no ready signed events to push".to_owned()); - } - None -} - fn active_profile_view(config: &RuntimeConfig) -> TransportProfileView { match config.transport.profile { TransportProfileKind::LocalOnly => profile_view_from_parts( @@ -424,6 +388,7 @@ fn transport_implementation_label(state: RadrootsTransportImplementationState) - fn transport_maturity_label(maturity: RadrootsTransportCapabilityMaturity) -> &'static str { match maturity { + RadrootsTransportCapabilityMaturity::Experimental => "experimental", RadrootsTransportCapabilityMaturity::Preview => "preview", RadrootsTransportCapabilityMaturity::Stable => "stable", } @@ -545,66 +510,23 @@ fn string_array_input(input: &OperationData, key: &str) -> Vec<String> { #[cfg(test)] mod tests { use super::*; - use radroots_event::id::RadrootsEventId; - use radroots_sdk::{ - PushOutboxEventReceipt, PushOutboxTargetReceipt, PushOutboxTransportOutcomeKind, - }; #[test] - fn transport_outbox_push_reports_reticulum_deferred_without_attempts() { - let cases = [( - PushOutboxEventState::DeferredUntilImplemented, - PushOutboxTargetOutcomeKind::DeferredUntilImplemented, - "deferred_until_implemented", - "SDK outbox push reported Reticulum work as deferred until implemented without network delivery", - )]; - - for (final_state, outcome_kind, expected_state, expected_reason) in cases { - let receipt = reticulum_receipt(final_state, outcome_kind); - - assert_eq!(transport_outbox_push_state(&receipt, 0), expected_state); - assert_eq!( - transport_outbox_push_reason(&receipt).as_deref(), - Some(expected_reason) - ); - } + fn reticulum_profile_remains_explicitly_unavailable() { + let status = reticulum_transport_status("reticulum"); + assert!(!status.usable_for_delivery); + assert_eq!( + status.endpoint_uri.as_deref(), + Some(RADROOTS_RETICULUM_ENDPOINT_URI) + ); + assert_eq!(status.message, RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE); } - fn reticulum_receipt( - final_state: PushOutboxEventState, - outcome_kind: PushOutboxTargetOutcomeKind, - ) -> PushOutboxReceipt { - PushOutboxReceipt { - attempted_events: 0, - published_events: 0, - retryable_events: 0, - terminal_events: 0, - events: vec![PushOutboxEventReceipt { - event_id: RadrootsEventId::parse("d".repeat(64).as_str()).expect("event id"), - outbox_event_id: 11, - final_state, - attempted_count: 0, - accepted_count: 0, - retryable_count: 0, - terminal_count: 0, - quorum: 1, - quorum_met: false, - targets: vec![PushOutboxTargetReceipt { - transport_kind: "reticulum".to_owned(), - endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), - target_scope: Some(RADROOTS_RETICULUM_SCOPE_ID.to_owned()), - target_label: None, - outcome_kind, - transport_outcome_kind: Some(match outcome_kind { - PushOutboxTargetOutcomeKind::DeferredUntilImplemented => { - PushOutboxTransportOutcomeKind::DeferredUntilImplemented - } - _ => PushOutboxTransportOutcomeKind::TransportUnavailable, - }), - attempted: false, - message: Some(RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned()), - }], - }], - } + #[test] + fn experimental_maturity_has_a_stable_terminal_label() { + assert_eq!( + transport_maturity_label(RadrootsTransportCapabilityMaturity::Experimental), + "experimental" + ); } } diff --git a/src/view/runtime.rs b/src/view/runtime.rs @@ -1230,7 +1230,7 @@ impl MarketReadinessView { mod market_readiness_tests { use super::MarketReadinessView; - const LISTING_ADDR: &str = "30402:1111111111111111111111111111111111111111111111111111111111111111:AAAAAAAAAAAAAAAAAAAAAg"; + const LISTING_ADDR: &str = "30402:585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df:AAAAAAAAAAAAAAAAAAAAAg"; #[test] fn market_readiness_separates_protocol_marketplace_and_order_request_state() {