radrootsd

JSON-RPC bridge for Radroots event publishing
git clone https://radroots.dev/git/radrootsd.git
Log | Files | Refs | README | LICENSE

commit 80286dcef72e3267276c0281f19ebdc80d2a9282
parent 2b6f6fcc2a59e5fdedbc1f61049486b7d32be92e
Author: triesap <tyson@radroots.org>
Date:   Tue,  7 Jul 2026 03:20:35 +0000

transport: replace daemon publish proxy surface

- Move radrootsd publish handling to the transport.publish JSON-RPC method family.
- Store target-policy jobs and target outcomes in transport_publish tables.
- Scope Nostr relay routing under the transport_publish.nostr config object.
- Validate with cargo fmt, cargo check, cargo test, and no-legacy source scans.

Diffstat:
MCargo.lock | 55++++++++++++++++++++++++++++++++-----------------------
MCargo.toml | 8+++++---
Mconfig.toml | 14++++++++------
Msrc/app/cli.rs | 14++++++++------
Msrc/app/config.rs | 252+++++++++++++++++++++++++++++++++++++++++--------------------------------------
Msrc/app/paths.rs | 22++++++++++++----------
Msrc/app/runtime.rs | 105+++++++++++++++++++++++++++++++++++++++++++++++--------------------------------
Msrc/core/mod.rs | 2+-
Dsrc/core/publish_proxy/mod.rs | 3018-------------------------------------------------------------------------------
Msrc/core/state.rs | 24++++++++++++------------
Asrc/core/transport_publish.rs | 3270+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/transport/jsonrpc/auth.rs | 81++++++++++++++++++++++++++++++++++++++++++-------------------------------------
Msrc/transport/jsonrpc/methods/mod.rs | 56+++++++++++++++++++++++++++++---------------------------
Dsrc/transport/jsonrpc/methods/publish_proxy.rs | 526-------------------------------------------------------------------------------
Asrc/transport/jsonrpc/methods/transport_publish.rs | 437+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/transport/jsonrpc/mod.rs | 6+++---
Msrc/transport/jsonrpc/server.rs | 112++++++++++++++++++++++++++++++++++++++++++-------------------------------------
17 files changed, 4114 insertions(+), 3888 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -1858,27 +1858,6 @@ dependencies = [ ] [[package]] -name = "radroots_publish_proxy_protocol" -version = "0.1.0-alpha.2" -dependencies = [ - "serde", -] - -[[package]] -name = "radroots_relay_transport" -version = "0.1.0-alpha.2" -dependencies = [ - "futures", - "nostr", - "radroots_events", - "radroots_nostr", - "serde", - "serde_json", - "thiserror 1.0.69", - "url", -] - -[[package]] name = "radroots_runtime" version = "0.1.0-alpha.2" dependencies = [ @@ -1914,6 +1893,35 @@ name = "radroots_secret_vault" version = "0.1.0-alpha.2" [[package]] +name = "radroots_transport" +version = "0.1.0-alpha.2" +dependencies = [ + "sha2", +] + +[[package]] +name = "radroots_transport_nostr" +version = "0.1.0-alpha.2" +dependencies = [ + "futures", + "nostr", + "radroots_events", + "radroots_nostr", + "radroots_transport", + "serde", + "serde_json", + "thiserror 1.0.69", + "url", +] + +[[package]] +name = "radroots_transport_publish_protocol" +version = "0.1.0-alpha.2" +dependencies = [ + "serde", +] + +[[package]] name = "radrootsd" version = "0.1.0" dependencies = [ @@ -1924,10 +1932,11 @@ dependencies = [ "radroots_events", "radroots_identity", "radroots_nostr", - "radroots_publish_proxy_protocol", - "radroots_relay_transport", "radroots_runtime", "radroots_runtime_paths", + "radroots_transport", + "radroots_transport_nostr", + "radroots_transport_publish_protocol", "rand 0.9.2", "rusqlite", "serde", diff --git a/Cargo.toml b/Cargo.toml @@ -14,8 +14,9 @@ resolver = "2" radroots_events = { path = "../lib/crates/events" } radroots_identity = { path = "../lib/crates/identity" } radroots_nostr = { path = "../lib/crates/nostr" } -radroots_publish_proxy_protocol = { path = "../lib/crates/publish_proxy_protocol" } -radroots_relay_transport = { path = "../lib/crates/relay_transport", default-features = false } +radroots_transport_publish_protocol = { path = "../lib/crates/transport_publish_protocol" } +radroots_relay_transport = { package = "radroots_transport_nostr", path = "../lib/crates/transport_nostr", default-features = false } +radroots_transport = { path = "../lib/crates/transport", default-features = false } radroots_runtime = { path = "../lib/crates/runtime" } [lints.rust] @@ -25,8 +26,9 @@ unexpected_cfgs = { level = "warn", check-cfg = ['cfg(coverage_nightly)'] } radroots_events = { workspace = true, features = ["serde"] } radroots_identity = { workspace = true } radroots_nostr = { workspace = true, features = ["client", "codec", "events", "http"] } -radroots_publish_proxy_protocol = { workspace = true, features = ["std", "serde"] } +radroots_transport_publish_protocol = { workspace = true, features = ["std", "serde"] } radroots_relay_transport = { workspace = true, features = ["std", "client"] } +radroots_transport = { workspace = true } radroots_runtime = { workspace = true, features = ["cli"] } radroots_runtime_paths = { path = "../lib/crates/runtime_paths" } nostr = { version = "0.44.2", features = ["nip46"] } diff --git a/config.toml b/config.toml @@ -9,13 +9,13 @@ # interactive_user: # logs_dir = ~/.radroots/logs/services/radrootsd # service identity = ~/.radroots/secrets/services/radrootsd/identity.secret.json -# publish proxy database = ~/.radroots/data/services/radrootsd/publish_proxy.sqlite +# transport publish database = ~/.radroots/data/services/radrootsd/transport_publish.sqlite # service_host: # logs_dir = /var/log/radroots/services/radrootsd # service identity = /etc/radroots/secrets/services/radrootsd/identity.secret.json -# publish proxy database = /var/lib/radroots/services/radrootsd/publish_proxy.sqlite +# transport publish database = /var/lib/radroots/services/radrootsd/transport_publish.sqlite # the canonical live service identity is always an encrypted local envelope -# only override logs_dir or config.publish_proxy.database_path intentionally +# only override logs_dir or config.transport_publish.database_path intentionally [metadata] name = "radrootsd" @@ -36,15 +36,17 @@ relays = [ [config.rpc] addr = "127.0.0.1:7070" -[config.publish_proxy] +[config.transport_publish] enabled = true max_event_bytes = 131072 -max_relays_per_request = 20 +max_targets_per_request = 20 job_list_limit = 100 max_concurrent_publish_jobs = 8 + +[config.transport_publish.nostr] relay_url_policy = "localhost" author_relay_discovery_relays = [] -daemon_default_publish_relays = ["ws://127.0.0.1:8080"] +daemon_default_relays = ["ws://127.0.0.1:8080"] [config.nip46] public_jsonrpc_enabled = false diff --git a/src/app/cli.rs b/src/app/cli.rs @@ -18,17 +18,17 @@ pub struct Args { #[derive(Subcommand, Debug, Clone)] pub enum Command { - PublishProxy(PublishProxyCommand), + TransportPublish(TransportPublishCommand), } #[derive(ClapArgs, Debug, Clone)] -pub struct PublishProxyCommand { +pub struct TransportPublishCommand { #[command(subcommand)] - pub command: PublishProxySubcommand, + pub command: TransportPublishSubcommand, } #[derive(Subcommand, Debug, Clone)] -pub enum PublishProxySubcommand { +pub enum TransportPublishSubcommand { Principal(PrincipalCommand), } @@ -54,9 +54,11 @@ pub struct PrincipalInitArgs { #[arg(long)] pub allowed_kind: Vec<u32>, #[arg(long)] - pub allowed_relay_policy: Vec<String>, + pub allowed_target_policy: Vec<String>, + #[arg(long)] + pub allowed_nostr_source_policy: Vec<String>, #[arg(long)] pub job_visibility: String, #[arg(long)] - pub allow_request_relays: bool, + pub allow_request_targets: bool, } diff --git a/src/app/config.rs b/src/app/config.rs @@ -5,7 +5,7 @@ use serde::{Deserialize, Serialize}; use std::path::{Path, PathBuf}; use super::paths::{ - RadrootsdRuntimePaths, default_publish_proxy_database_path, process_path_selection, + RadrootsdRuntimePaths, default_transport_publish_database_path, process_path_selection, resolve_runtime_paths_with_resolver, }; @@ -49,32 +49,32 @@ fn default_nip46_public_jsonrpc_enabled() -> bool { false } -fn default_publish_proxy_enabled() -> bool { +fn default_transport_publish_enabled() -> bool { true } -fn default_publish_proxy_connect_timeout_secs() -> u64 { +fn default_transport_publish_connect_timeout_secs() -> u64 { 10 } -fn default_publish_proxy_max_event_bytes() -> usize { +fn default_transport_publish_max_event_bytes() -> usize { 128 * 1024 } -fn default_publish_proxy_max_relays_per_request() -> usize { +fn default_transport_publish_max_targets_per_request() -> usize { 20 } -fn default_publish_proxy_job_list_limit() -> usize { +fn default_transport_publish_job_list_limit() -> usize { 100 } -fn default_publish_proxy_max_concurrent_publish_jobs() -> usize { +fn default_transport_publish_max_concurrent_publish_jobs() -> usize { 8 } -fn default_publish_proxy_relay_url_policy() -> PublishProxyRelayUrlPolicy { - PublishProxyRelayUrlPolicy::Public +fn default_nostr_relay_url_policy() -> NostrRelayUrlPolicy { + NostrRelayUrlPolicy::Public } #[derive(Debug, Deserialize, Clone, Default)] @@ -103,66 +103,63 @@ impl RawServiceConfig { } #[derive(Debug, Deserialize, Clone)] -struct RawPublishProxyConfig { - #[serde(default = "default_publish_proxy_enabled")] +#[serde(deny_unknown_fields)] +struct RawTransportPublishConfig { + #[serde(default = "default_transport_publish_enabled")] pub enabled: bool, - #[serde(default = "default_publish_proxy_connect_timeout_secs")] + #[serde(default = "default_transport_publish_connect_timeout_secs")] pub connect_timeout_secs: u64, - #[serde(default = "default_publish_proxy_max_event_bytes")] + #[serde(default = "default_transport_publish_max_event_bytes")] pub max_event_bytes: usize, - #[serde(default = "default_publish_proxy_max_relays_per_request")] - pub max_relays_per_request: usize, - #[serde(default = "default_publish_proxy_job_list_limit")] + #[serde(default = "default_transport_publish_max_targets_per_request")] + pub max_targets_per_request: usize, + #[serde(default = "default_transport_publish_job_list_limit")] pub job_list_limit: usize, - #[serde(default = "default_publish_proxy_max_concurrent_publish_jobs")] + #[serde(default = "default_transport_publish_max_concurrent_publish_jobs")] pub max_concurrent_publish_jobs: usize, #[serde(default)] pub database_path: Option<PathBuf>, - #[serde(default = "default_publish_proxy_relay_url_policy")] - pub relay_url_policy: PublishProxyRelayUrlPolicy, #[serde(default)] - pub author_relay_discovery_relays: Vec<String>, - #[serde(default)] - pub daemon_default_publish_relays: Vec<String>, + pub nostr: TransportPublishNostrConfig, } -impl Default for RawPublishProxyConfig { +impl Default for RawTransportPublishConfig { fn default() -> Self { Self { - enabled: default_publish_proxy_enabled(), - connect_timeout_secs: default_publish_proxy_connect_timeout_secs(), - max_event_bytes: default_publish_proxy_max_event_bytes(), - max_relays_per_request: default_publish_proxy_max_relays_per_request(), - job_list_limit: default_publish_proxy_job_list_limit(), - max_concurrent_publish_jobs: default_publish_proxy_max_concurrent_publish_jobs(), + enabled: default_transport_publish_enabled(), + connect_timeout_secs: default_transport_publish_connect_timeout_secs(), + max_event_bytes: default_transport_publish_max_event_bytes(), + max_targets_per_request: default_transport_publish_max_targets_per_request(), + job_list_limit: default_transport_publish_job_list_limit(), + max_concurrent_publish_jobs: default_transport_publish_max_concurrent_publish_jobs(), database_path: None, - relay_url_policy: default_publish_proxy_relay_url_policy(), - author_relay_discovery_relays: Vec::new(), - daemon_default_publish_relays: Vec::new(), + nostr: TransportPublishNostrConfig::default(), } } } -impl RawPublishProxyConfig { - fn into_publish_proxy_config(self, paths: &RadrootsdRuntimePaths) -> PublishProxyConfig { - PublishProxyConfig { +impl RawTransportPublishConfig { + fn into_transport_publish_config( + self, + paths: &RadrootsdRuntimePaths, + ) -> TransportPublishConfig { + TransportPublishConfig { enabled: self.enabled, connect_timeout_secs: self.connect_timeout_secs, max_event_bytes: self.max_event_bytes, - max_relays_per_request: self.max_relays_per_request, + max_targets_per_request: self.max_targets_per_request, job_list_limit: self.job_list_limit, max_concurrent_publish_jobs: self.max_concurrent_publish_jobs, database_path: self .database_path - .unwrap_or_else(|| paths.publish_proxy_database_path.clone()), - relay_url_policy: self.relay_url_policy, - author_relay_discovery_relays: self.author_relay_discovery_relays, - daemon_default_publish_relays: self.daemon_default_publish_relays, + .unwrap_or_else(|| paths.transport_publish_database_path.clone()), + nostr: self.nostr, } } } #[derive(Debug, Deserialize, Clone)] +#[serde(deny_unknown_fields)] struct RawConfiguration { #[serde(flatten)] pub service: RawServiceConfig, @@ -173,9 +170,7 @@ struct RawConfiguration { #[serde(default)] pub nip46: Nip46Config, #[serde(default)] - pub publish_proxy: RawPublishProxyConfig, - #[serde(default, rename = "bridge")] - pub obsolete_publish_bridge_config: Option<serde::de::IgnoredAny>, + pub transport_publish: RawTransportPublishConfig, } #[derive(Debug, Deserialize, Clone)] @@ -193,11 +188,10 @@ impl RawSettings { rpc: self.config.rpc, rpc_addr: self.config.rpc_addr, nip46: self.config.nip46, - publish_proxy: self.config.publish_proxy.into_publish_proxy_config(paths), - obsolete_bridge_config_present: self + transport_publish: self .config - .obsolete_publish_bridge_config - .is_some(), + .transport_publish + .into_transport_publish_config(paths), }, } } @@ -253,68 +247,84 @@ impl Default for Nip46Config { #[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] #[serde(rename_all = "snake_case")] -pub enum PublishProxyRelayUrlPolicy { +pub enum NostrRelayUrlPolicy { Public, Localhost, } #[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] -pub struct PublishProxyConfig { - #[serde(default = "default_publish_proxy_enabled")] +#[serde(deny_unknown_fields)] +pub struct TransportPublishNostrConfig { + #[serde(default = "default_nostr_relay_url_policy")] + pub relay_url_policy: NostrRelayUrlPolicy, + #[serde(default)] + pub author_relay_discovery_relays: Vec<String>, + #[serde(default)] + pub daemon_default_relays: Vec<String>, +} + +impl Default for TransportPublishNostrConfig { + fn default() -> Self { + Self { + relay_url_policy: default_nostr_relay_url_policy(), + author_relay_discovery_relays: Vec::new(), + daemon_default_relays: Vec::new(), + } + } +} + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct TransportPublishConfig { + #[serde(default = "default_transport_publish_enabled")] pub enabled: bool, - #[serde(default = "default_publish_proxy_connect_timeout_secs")] + #[serde(default = "default_transport_publish_connect_timeout_secs")] pub connect_timeout_secs: u64, - #[serde(default = "default_publish_proxy_max_event_bytes")] + #[serde(default = "default_transport_publish_max_event_bytes")] pub max_event_bytes: usize, - #[serde(default = "default_publish_proxy_max_relays_per_request")] - pub max_relays_per_request: usize, - #[serde(default = "default_publish_proxy_job_list_limit")] + #[serde(default = "default_transport_publish_max_targets_per_request")] + pub max_targets_per_request: usize, + #[serde(default = "default_transport_publish_job_list_limit")] pub job_list_limit: usize, - #[serde(default = "default_publish_proxy_max_concurrent_publish_jobs")] + #[serde(default = "default_transport_publish_max_concurrent_publish_jobs")] pub max_concurrent_publish_jobs: usize, - #[serde(default = "default_publish_proxy_database_path")] + #[serde(default = "default_transport_publish_database_path")] pub database_path: PathBuf, - #[serde(default = "default_publish_proxy_relay_url_policy")] - pub relay_url_policy: PublishProxyRelayUrlPolicy, #[serde(default)] - pub author_relay_discovery_relays: Vec<String>, - #[serde(default)] - pub daemon_default_publish_relays: Vec<String>, + pub nostr: TransportPublishNostrConfig, } -impl Default for PublishProxyConfig { +impl Default for TransportPublishConfig { fn default() -> Self { Self { - enabled: default_publish_proxy_enabled(), - connect_timeout_secs: default_publish_proxy_connect_timeout_secs(), - max_event_bytes: default_publish_proxy_max_event_bytes(), - max_relays_per_request: default_publish_proxy_max_relays_per_request(), - job_list_limit: default_publish_proxy_job_list_limit(), - max_concurrent_publish_jobs: default_publish_proxy_max_concurrent_publish_jobs(), - database_path: default_publish_proxy_database_path(), - relay_url_policy: default_publish_proxy_relay_url_policy(), - author_relay_discovery_relays: Vec::new(), - daemon_default_publish_relays: Vec::new(), + enabled: default_transport_publish_enabled(), + connect_timeout_secs: default_transport_publish_connect_timeout_secs(), + max_event_bytes: default_transport_publish_max_event_bytes(), + max_targets_per_request: default_transport_publish_max_targets_per_request(), + job_list_limit: default_transport_publish_job_list_limit(), + max_concurrent_publish_jobs: default_transport_publish_max_concurrent_publish_jobs(), + database_path: default_transport_publish_database_path(), + nostr: TransportPublishNostrConfig::default(), } } } -impl PublishProxyConfig { +impl TransportPublishConfig { pub fn validate(&self) -> Result<()> { if self.max_event_bytes == 0 { - bail!("publish_proxy max_event_bytes must be greater than zero"); + bail!("transport_publish max_event_bytes must be greater than zero"); } - if self.max_relays_per_request == 0 { - bail!("publish_proxy max_relays_per_request must be greater than zero"); + if self.max_targets_per_request == 0 { + bail!("transport_publish max_targets_per_request must be greater than zero"); } if self.job_list_limit == 0 { - bail!("publish_proxy job_list_limit must be greater than zero"); + bail!("transport_publish job_list_limit must be greater than zero"); } if self.max_concurrent_publish_jobs == 0 { - bail!("publish_proxy max_concurrent_publish_jobs must be greater than zero"); + bail!("transport_publish max_concurrent_publish_jobs must be greater than zero"); } if self.connect_timeout_secs == 0 { - bail!("publish_proxy connect_timeout_secs must be greater than zero"); + bail!("transport_publish connect_timeout_secs must be greater than zero"); } Ok(()) } @@ -363,9 +373,7 @@ pub struct Configuration { #[serde(default)] pub nip46: Nip46Config, #[serde(default)] - pub publish_proxy: PublishProxyConfig, - #[serde(default, skip_serializing)] - pub(crate) obsolete_bridge_config_present: bool, + pub transport_publish: TransportPublishConfig, } impl Configuration { @@ -374,10 +382,7 @@ impl Configuration { } pub fn validate(&self) -> Result<()> { - if self.obsolete_bridge_config_present { - bail!("config.bridge is obsolete; use config.publish_proxy"); - } - self.publish_proxy.validate()?; + self.transport_publish.validate()?; Ok(()) } } @@ -399,7 +404,7 @@ mod tests { use std::path::PathBuf; use super::{ - Configuration, Nip46Config, PublishProxyConfig, PublishProxyRelayUrlPolicy, RpcConfig, + Configuration, Nip46Config, NostrRelayUrlPolicy, RpcConfig, TransportPublishConfig, load_settings_from_path_with_resolver, }; use crate::app::paths::{ @@ -465,19 +470,19 @@ mod tests { } #[test] - fn publish_proxy_defaults_are_expected() { + fn transport_publish_defaults_are_expected() { let paths = default_runtime_paths_for_process().expect("resolve process runtime paths"); - let cfg = PublishProxyConfig::default(); + let cfg = TransportPublishConfig::default(); assert!(cfg.enabled); assert_eq!(cfg.connect_timeout_secs, 10); assert_eq!(cfg.max_event_bytes, 128 * 1024); - assert_eq!(cfg.max_relays_per_request, 20); + assert_eq!(cfg.max_targets_per_request, 20); assert_eq!(cfg.job_list_limit, 100); assert_eq!(cfg.max_concurrent_publish_jobs, 8); - assert_eq!(cfg.database_path, paths.publish_proxy_database_path); - assert_eq!(cfg.relay_url_policy, PublishProxyRelayUrlPolicy::Public); - assert!(cfg.author_relay_discovery_relays.is_empty()); - assert!(cfg.daemon_default_publish_relays.is_empty()); + assert_eq!(cfg.database_path, paths.transport_publish_database_path); + assert_eq!(cfg.nostr.relay_url_policy, NostrRelayUrlPolicy::Public); + assert!(cfg.nostr.author_relay_discovery_relays.is_empty()); + assert!(cfg.nostr.daemon_default_relays.is_empty()); } #[test] @@ -490,8 +495,7 @@ mod tests { }, rpc_addr: None, nip46: Nip46Config::default(), - publish_proxy: PublishProxyConfig::default(), - obsolete_bridge_config_present: false, + transport_publish: TransportPublishConfig::default(), }; assert_eq!(cfg.rpc_addr(), "127.0.0.1:1111"); cfg.rpc_addr = Some("127.0.0.1:2222".to_string()); @@ -499,20 +503,20 @@ mod tests { } #[test] - fn publish_proxy_validation_rejects_zero_limits() { - let mut cfg = PublishProxyConfig::default(); + fn transport_publish_validation_rejects_zero_limits() { + let mut cfg = TransportPublishConfig::default(); cfg.max_event_bytes = 0; assert!(cfg.validate().is_err()); - let mut cfg = PublishProxyConfig::default(); - cfg.max_relays_per_request = 0; + let mut cfg = TransportPublishConfig::default(); + cfg.max_targets_per_request = 0; assert!(cfg.validate().is_err()); - let mut cfg = PublishProxyConfig::default(); + let mut cfg = TransportPublishConfig::default(); cfg.job_list_limit = 0; assert!(cfg.validate().is_err()); - let mut cfg = PublishProxyConfig::default(); + let mut cfg = TransportPublishConfig::default(); cfg.max_concurrent_publish_jobs = 0; assert!(cfg.validate().is_err()); - let mut cfg = PublishProxyConfig::default(); + let mut cfg = TransportPublishConfig::default(); cfg.connect_timeout_secs = 0; assert!(cfg.validate().is_err()); } @@ -541,8 +545,10 @@ mod tests { ) ); assert_eq!( - paths.publish_proxy_database_path, - PathBuf::from("/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite") + paths.transport_publish_database_path, + PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite" + ) ); } @@ -568,8 +574,8 @@ mod tests { PathBuf::from("/etc/radroots/secrets/services/radrootsd/identity.secret.json") ); assert_eq!( - paths.publish_proxy_database_path, - PathBuf::from("/var/lib/radroots/services/radrootsd/publish_proxy.sqlite") + paths.transport_publish_database_path, + PathBuf::from("/var/lib/radroots/services/radrootsd/transport_publish.sqlite") ); } @@ -596,8 +602,8 @@ mod tests { repo_local_root.join("secrets/services/radrootsd/identity.secret.json") ); assert_eq!( - paths.publish_proxy_database_path, - repo_local_root.join("data/services/radrootsd/publish_proxy.sqlite") + paths.transport_publish_database_path, + repo_local_root.join("data/services/radrootsd/transport_publish.sqlite") ); } @@ -633,13 +639,15 @@ addr = "127.0.0.1:7070" "/home/treesap/.radroots/logs/services/radrootsd" ); assert_eq!( - settings.config.publish_proxy.database_path, - PathBuf::from("/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite") + settings.config.transport_publish.database_path, + PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite" + ) ); } #[test] - fn obsolete_config_is_rejected() { + fn obsolete_transport_publish_config_is_rejected() { let temp = tempfile::tempdir().expect("tempdir"); let config_path = temp.path().join("radrootsd.toml"); std::fs::write( @@ -651,8 +659,8 @@ name = "radrootsd-test" [config] relays = [] -[config.bridge] -enabled = true +[config.transport_publish] +relay_url_policy = "localhost" "#, ) .expect("write config"); @@ -663,8 +671,10 @@ enabled = true RadrootsPathProfile::InteractiveUser, None, ) - .expect_err("obsolete config should fail"); - assert!(err.to_string().contains("config.bridge")); + .expect_err("obsolete transport_publish config should fail"); + let err_chain = format!("{err:?}"); + assert!(err_chain.contains("unknown field")); + assert!(err_chain.contains("relay_url_policy")); } #[test] @@ -681,12 +691,14 @@ enabled = true contract.path_overrides.subordinate_path_override_keys, vec![ "config.service.logs_dir".to_owned(), - "config.publish_proxy.database_path".to_owned(), + "config.transport_publish.database_path".to_owned(), ] ); assert_eq!( - contract.canonical_publish_proxy_database_path, - PathBuf::from("/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite") + contract.canonical_transport_publish_database_path, + PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite" + ) ); } } diff --git a/src/app/paths.rs b/src/app/paths.rs @@ -9,7 +9,7 @@ use radroots_runtime_paths::{ use serde::Serialize; const RADROOTSD_RUNTIME_ID: &str = "radrootsd"; -const PUBLISH_PROXY_DATABASE_FILE_NAME: &str = "publish_proxy.sqlite"; +const TRANSPORT_PUBLISH_DATABASE_FILE_NAME: &str = "transport_publish.sqlite"; const RADROOTSD_PATHS_PROFILE_ENV: &str = "RADROOTSD_PATHS_PROFILE"; const RADROOTSD_PATHS_REPO_LOCAL_ROOT_ENV: &str = "RADROOTSD_PATHS_REPO_LOCAL_ROOT"; const RADROOTSD_DEFAULT_SHARED_SECRET_BACKEND: &str = "encrypted_file"; @@ -18,7 +18,7 @@ const RADROOTSD_ALLOWED_SHARED_SECRET_BACKENDS: [&str; 1] = ["encrypted_file"]; const SUBORDINATE_PATH_OVERRIDE_SOURCE: &str = "config_artifact"; const SUBORDINATE_PATH_OVERRIDE_KEYS: [&str; 2] = [ "config.service.logs_dir", - "config.publish_proxy.database_path", + "config.transport_publish.database_path", ]; #[derive(Debug, Clone, PartialEq, Eq)] @@ -26,7 +26,7 @@ pub(crate) struct RadrootsdRuntimePaths { pub(crate) config_path: PathBuf, pub(crate) logs_dir: PathBuf, pub(crate) identity_path: PathBuf, - pub(crate) publish_proxy_database_path: PathBuf, + pub(crate) transport_publish_database_path: PathBuf, } #[derive(Debug, Clone, PartialEq, Eq, Serialize)] @@ -39,7 +39,7 @@ pub struct RadrootsdRuntimeContractOutput { pub canonical_config_path: PathBuf, pub canonical_logs_dir: PathBuf, pub canonical_identity_path: PathBuf, - pub canonical_publish_proxy_database_path: PathBuf, + pub canonical_transport_publish_database_path: PathBuf, } pub type RadrootsdRuntimePathOverrideContractOutput = RadrootsRuntimeSelectionOverrideContract; @@ -77,7 +77,7 @@ pub(crate) fn resolve_runtime_paths_with_resolver( config_path: namespaced.config.join(DEFAULT_CONFIG_FILE_NAME), logs_dir: namespaced.logs, identity_path: namespaced.secrets.join(DEFAULT_SERVICE_IDENTITY_FILE_NAME), - publish_proxy_database_path: namespaced.data.join(PUBLISH_PROXY_DATABASE_FILE_NAME), + transport_publish_database_path: namespaced.data.join(TRANSPORT_PUBLISH_DATABASE_FILE_NAME), }) } @@ -90,10 +90,10 @@ pub(crate) fn default_runtime_paths_for_process() -> Result<RadrootsdRuntimePath ) } -pub(crate) fn default_publish_proxy_database_path() -> PathBuf { +pub(crate) fn default_transport_publish_database_path() -> PathBuf { default_runtime_paths_for_process() .expect("resolve canonical radrootsd runtime paths") - .publish_proxy_database_path + .transport_publish_database_path } pub fn default_config_path_for_process() -> Result<PathBuf> { @@ -133,7 +133,7 @@ pub(crate) fn runtime_contract_with_selection( canonical_config_path: paths.config_path, canonical_logs_dir: paths.logs_dir, canonical_identity_path: paths.identity_path, - canonical_publish_proxy_database_path: paths.publish_proxy_database_path, + canonical_transport_publish_database_path: paths.transport_publish_database_path, }) } @@ -209,8 +209,10 @@ mod tests { ) ); assert_eq!( - contract.canonical_publish_proxy_database_path, - PathBuf::from("/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite") + contract.canonical_transport_publish_database_path, + PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite" + ) ); } } diff --git a/src/app/runtime.rs b/src/app/runtime.rs @@ -57,9 +57,9 @@ struct RadrootsdRuntimeStartupReport { identity_path: PathBuf, identity_path_source: String, canonical_identity_path: PathBuf, - publish_proxy_database_path: PathBuf, - publish_proxy_database_path_source: String, - canonical_publish_proxy_database_path: PathBuf, + transport_publish_database_path: PathBuf, + transport_publish_database_path_source: String, + canonical_transport_publish_database_path: PathBuf, path_overrides: paths::RadrootsdRuntimePathOverrideContractOutput, default_shared_secret_backend: String, allowed_shared_secret_backends: Vec<String>, @@ -194,13 +194,13 @@ fn runtime_startup_report( &contract.canonical_identity_path, ), canonical_identity_path: contract.canonical_identity_path.clone(), - publish_proxy_database_path: settings.config.publish_proxy.database_path.clone(), - publish_proxy_database_path_source: config_or_profile_path_source( - &settings.config.publish_proxy.database_path, - &contract.canonical_publish_proxy_database_path, + transport_publish_database_path: settings.config.transport_publish.database_path.clone(), + transport_publish_database_path_source: config_or_profile_path_source( + &settings.config.transport_publish.database_path, + &contract.canonical_transport_publish_database_path, ), - canonical_publish_proxy_database_path: contract - .canonical_publish_proxy_database_path + canonical_transport_publish_database_path: contract + .canonical_transport_publish_database_path .clone(), path_overrides: contract.path_overrides.clone(), default_shared_secret_backend: contract.default_shared_secret_backend.clone(), @@ -246,9 +246,9 @@ fn log_runtime_startup_report(report: &RadrootsdRuntimeStartupReport) { identity_path = %report.identity_path.display(), identity_path_source = report.identity_path_source.as_str(), canonical_identity_path = %report.canonical_identity_path.display(), - publish_proxy_database_path = %report.publish_proxy_database_path.display(), - publish_proxy_database_path_source = report.publish_proxy_database_path_source.as_str(), - canonical_publish_proxy_database_path = %report.canonical_publish_proxy_database_path.display(), + transport_publish_database_path = %report.transport_publish_database_path.display(), + transport_publish_database_path_source = report.transport_publish_database_path_source.as_str(), + canonical_transport_publish_database_path = %report.canonical_transport_publish_database_path.display(), default_shared_secret_backend = report.default_shared_secret_backend.as_str(), allowed_shared_secret_backends = ?report.allowed_shared_secret_backends, "radrootsd runtime contract" @@ -383,40 +383,55 @@ async fn wait_for_shutdown_or_stopped(handle: ServerHandle) -> RunWaitOutcome { async fn handle_command(command: cli::Command, settings: &config::Settings) -> Result<()> { match command { - cli::Command::PublishProxy(command) => match command.command { - cli::PublishProxySubcommand::Principal(command) => match command.command { + cli::Command::TransportPublish(command) => match command.command { + cli::TransportPublishSubcommand::Principal(command) => match command.command { cli::PrincipalSubcommand::Init(args) => { - let token = crate::core::publish_proxy::generate_bearer_token(); - let token_hash = crate::core::publish_proxy::hash_bearer_token(token.as_str()); - let store = crate::core::publish_proxy::PublishProxyStore::open( - settings.config.publish_proxy.database_path.clone(), + let token = crate::core::transport_publish::generate_bearer_token(); + let token_hash = + crate::core::transport_publish::hash_bearer_token(token.as_str()); + let store = crate::core::transport_publish::TransportPublishStore::open( + settings.config.transport_publish.database_path.clone(), )?; let principal = store.create_principal( - crate::core::publish_proxy::PublishPrincipalInit { + crate::core::transport_publish::PublishPrincipalInit { label: args.label, token_hash, allowed_pubkeys: args.allowed_pubkey, allowed_kinds: args.allowed_kind, - allowed_relay_policies: args - .allowed_relay_policy + allowed_target_policies: args + .allowed_target_policy .iter() .map(|policy| { - crate::core::publish_proxy::parse_relay_policy(policy.as_str()) + crate::core::transport_publish::parse_target_policy( + policy.as_str(), + ) }) .collect::<Result<Vec<_>, _>>()?, - allow_request_relays: args.allow_request_relays, + allowed_nostr_source_policies: args + .allowed_nostr_source_policy + .iter() + .map(|policy| { + crate::core::transport_publish::parse_nostr_source_policy( + policy.as_str(), + ) + }) + .collect::<Result<Vec<_>, _>>()?, + allow_request_targets: args.allow_request_targets, job_visibility: args.job_visibility.parse()?, expires_at_unix: None, }, )?; - crate::core::publish_proxy::write_token_file(&args.token_file, token.as_str())?; + crate::core::transport_publish::write_token_file( + &args.token_file, + token.as_str(), + )?; println!( "{}", serde_json::json!({ "principal_id": principal.principal_id, "label": principal.label, "token_file": args.token_file, - "database_path": settings.config.publish_proxy.database_path, + "database_path": settings.config.transport_publish.database_path, }) ); Ok(()) @@ -450,7 +465,7 @@ pub async fn run() -> Result<()> { let radrootsd = Radrootsd::new( identity.clone(), settings.metadata.clone(), - settings.config.publish_proxy.clone(), + settings.config.transport_publish.clone(), settings.config.nip46.clone(), ); let radrootsd = radrootsd?; @@ -577,8 +592,7 @@ mod tests { }, rpc_addr: Some("127.0.0.1:0".to_string()), nip46: config::Nip46Config::default(), - publish_proxy: config::PublishProxyConfig::default(), - obsolete_bridge_config_present: false, + transport_publish: config::TransportPublishConfig::default(), }, } } @@ -599,7 +613,7 @@ mod tests { subordinate_path_override_source: "config_artifact".to_string(), subordinate_path_override_keys: vec![ "config.service.logs_dir".to_string(), - "config.publish_proxy.database_path".to_string(), + "config.transport_publish.database_path".to_string(), ], }, default_shared_secret_backend: "encrypted_file".to_string(), @@ -611,8 +625,8 @@ mod tests { canonical_identity_path: PathBuf::from( "/home/treesap/.radroots/secrets/services/radrootsd/identity.secret.json", ), - canonical_publish_proxy_database_path: PathBuf::from( - "/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite", + canonical_transport_publish_database_path: PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite", ), } } @@ -622,7 +636,7 @@ mod tests { let state = Radrootsd::new( identity, settings.metadata.clone(), - settings.config.publish_proxy.clone(), + settings.config.transport_publish.clone(), settings.config.nip46.clone(), ) .expect("state"); @@ -834,8 +848,8 @@ mod tests { }; let mut settings = settings_with_relays(Vec::new()); settings.config.service.logs_dir = "/tmp/radrootsd/logs".to_string(); - settings.config.publish_proxy.database_path = - PathBuf::from("/tmp/radrootsd/publish_proxy.sqlite"); + settings.config.transport_publish.database_path = + PathBuf::from("/tmp/radrootsd/transport_publish.sqlite"); let contract = sample_runtime_contract(); let report = runtime_startup_report(&args, &settings, &contract); @@ -859,10 +873,12 @@ mod tests { canonical_identity_path: PathBuf::from( "/home/treesap/.radroots/secrets/services/radrootsd/identity.secret.json" ), - publish_proxy_database_path: PathBuf::from("/tmp/radrootsd/publish_proxy.sqlite"), - publish_proxy_database_path_source: "config_artifact".to_string(), - canonical_publish_proxy_database_path: PathBuf::from( - "/home/treesap/.radroots/data/services/radrootsd/publish_proxy.sqlite" + transport_publish_database_path: PathBuf::from( + "/tmp/radrootsd/transport_publish.sqlite" + ), + transport_publish_database_path_source: "config_artifact".to_string(), + canonical_transport_publish_database_path: PathBuf::from( + "/home/treesap/.radroots/data/services/radrootsd/transport_publish.sqlite" ), path_overrides: sample_runtime_contract().path_overrides, default_shared_secret_backend: "encrypted_file".to_string(), @@ -884,8 +900,8 @@ mod tests { let contract = sample_runtime_contract(); let mut settings = settings_with_relays(Vec::new()); settings.config.service.logs_dir = contract.canonical_logs_dir.display().to_string(); - settings.config.publish_proxy.database_path = - contract.canonical_publish_proxy_database_path.clone(); + settings.config.transport_publish.database_path = + contract.canonical_transport_publish_database_path.clone(); let report = runtime_startup_report(&args, &settings, &contract); @@ -896,10 +912,13 @@ mod tests { assert_eq!(report.identity_path, contract.canonical_identity_path); assert_eq!(report.identity_path_source, "profile_default"); assert_eq!( - report.publish_proxy_database_path, - contract.canonical_publish_proxy_database_path + report.transport_publish_database_path, + contract.canonical_transport_publish_database_path + ); + assert_eq!( + report.transport_publish_database_path_source, + "profile_default" ); - assert_eq!(report.publish_proxy_database_path_source, "profile_default"); assert_eq!(report.path_overrides, contract.path_overrides); assert_eq!(report.default_shared_secret_backend, "encrypted_file"); assert_eq!( diff --git a/src/core/mod.rs b/src/core/mod.rs @@ -1,5 +1,5 @@ pub mod nip46; -pub mod publish_proxy; pub mod state; +pub mod transport_publish; pub use state::Radrootsd; diff --git a/src/core/publish_proxy/mod.rs b/src/core/publish_proxy/mod.rs @@ -1,3018 +0,0 @@ -use std::collections::BTreeMap; -use std::fmt; -use std::future::Future; -use std::net::IpAddr; -use std::path::{Path, PathBuf}; -use std::pin::Pin; -use std::str::FromStr; -use std::sync::{Arc, Mutex}; -use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; - -use radroots_events::RadrootsNostrEvent; -use radroots_events::draft::{ - RadrootsDraftError, RadrootsSignedNostrEvent, RadrootsSignedNostrEventParts, -}; -use radroots_nostr::prelude::{ - RadrootsNostrClient, RadrootsNostrEventVerification, RadrootsNostrFilter, RadrootsNostrKind, - RadrootsNostrPublicKey, radroots_nostr_verify_event, -}; -use radroots_publish_proxy_protocol::{ - PublishDeliveryPolicy, PublishEventRequest, PublishEventResponse, PublishJobStatus, - PublishJobView, PublishRelayOutcome, PublishRelayOutcomeKind, PublishRelayPolicy, - PublishRelaySource, SignedNostrEventWire, -}; -use radroots_relay_transport::{ - RadrootsNostrClientPublishAdapter, RadrootsRelayOutcome, RadrootsRelayOutcomeKind, - RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, - RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, -}; -use rusqlite::types::Type; -use rusqlite::{Connection, OptionalExtension, Row, params}; -use serde::{Deserialize, Serialize}; -use sha2::{Digest, Sha256}; -use thiserror::Error; -use tokio::sync::{OwnedSemaphorePermit, Semaphore}; -use uuid::Uuid; - -use crate::app::config::PublishProxyConfig; - -const TOKEN_PREFIX: &str = "rrd_pp_"; -const TOKEN_HASH_PREFIX: &str = "sha256:"; -const SCHEMA_VERSION: i64 = 2; - -#[derive(Debug, Error)] -pub enum PublishProxyError { - #[error("publish proxy storage error: {0}")] - Sqlite(#[from] rusqlite::Error), - #[error("publish proxy json error: {0}")] - Json(#[from] serde_json::Error), - #[error("publish proxy io error: {0}")] - Io(#[from] std::io::Error), - #[error("invalid publish proxy scope: {0}")] - InvalidScope(String), - #[error("invalid signed Nostr event: {0}")] - InvalidSignedEvent(String), - #[error("signed Nostr event verification failed: {0:?}")] - SignedEventVerification(RadrootsNostrEventVerification), - #[error("signed Nostr event conversion error: {0}")] - Draft(#[from] RadrootsDraftError), - #[error("publish proxy relay error: {0}")] - Relay(#[from] RadrootsRelayTransportError), - #[error("publish proxy transport error: {0}")] - Transport(String), - #[error("publish proxy concurrency limit reached")] - ConcurrencyLimit, - #[error("publish proxy idempotency conflict for key `{0}`")] - IdempotencyConflict(String), -} - -#[derive(Clone)] -pub struct PublishProxy { - pub config: PublishProxyConfig, - pub store: PublishProxyStore, - publisher: Option<Arc<dyn RadrootsRelayPublishAdapter>>, - resolver: Arc<dyn PublishRelayResolver>, - author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, - publish_jobs: Arc<Semaphore>, -} - -impl PublishProxy { - pub fn open(config: PublishProxyConfig) -> Result<Self, PublishProxyError> { - let store = PublishProxyStore::open(config.database_path.clone())?; - let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); - Ok(Self { - config, - store, - publisher: None, - resolver: Arc::new(SystemPublishRelayResolver), - author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), - publish_jobs, - }) - } - - pub fn memory(config: PublishProxyConfig) -> Result<Self, PublishProxyError> { - let store = PublishProxyStore::memory()?; - let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); - Ok(Self { - config, - store, - publisher: None, - resolver: Arc::new(SystemPublishRelayResolver), - author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), - publish_jobs, - }) - } - - pub fn with_publisher(mut self, publisher: Arc<dyn RadrootsRelayPublishAdapter>) -> Self { - self.publisher = Some(publisher); - self - } - - #[cfg(test)] - pub(crate) fn with_relay_resolver(mut self, resolver: Arc<dyn PublishRelayResolver>) -> Self { - self.resolver = resolver; - self - } - - #[cfg(test)] - fn with_author_relay_discovery( - mut self, - author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, - ) -> Self { - self.author_relay_discovery = author_relay_discovery; - self - } - - fn acquire_publish_permit(&self) -> Result<OwnedSemaphorePermit, PublishProxyError> { - self.publish_jobs - .clone() - .try_acquire_owned() - .map_err(|_| PublishProxyError::ConcurrencyLimit) - } - - pub async fn publish_event( - &self, - principal: &PublishPrincipal, - request: PublishEventRequest, - ) -> Result<PublishEventResponse, PublishProxyError> { - request - .validate(self.config.max_relays_per_request) - .map_err(|error| { - PublishProxyError::InvalidSignedEvent(format!( - "publish request validation failed: {error}" - )) - })?; - principal.allows_event(&request)?; - let signed_event = signed_event_from_wire(&request.event)?; - if signed_event.raw_json.len() > self.config.max_event_bytes { - return Err(PublishProxyError::InvalidSignedEvent( - "signed event exceeds publish_proxy max_event_bytes".to_owned(), - )); - } - let effective_timeout_ms = effective_publish_timeout_ms(&self.config, request.timeout_ms)?; - let _permit = self.acquire_publish_permit()?; - let request_fingerprint = request_intent_fingerprint( - principal.principal_id.as_str(), - signed_event.raw_json.as_str(), - &request, - effective_timeout_ms, - )?; - let resolution = self - .resolve_relays_for_request(signed_event.pubkey.as_str(), &request) - .await?; - let response = self.store.record_publish_job(PublishJobInsert { - principal_id: principal.principal_id.clone(), - idempotency_key: request.idempotency_key.clone(), - request: request.clone(), - request_fingerprint, - effective_relay_count: resolution.targets.len(), - })?; - if response.deduplicated { - return Ok(response); - } - let completed = self - .complete_job_execution( - response.job.job_id.as_str(), - signed_event, - request.delivery_policy.clone(), - effective_timeout_ms, - resolution, - ) - .await?; - Ok(PublishEventResponse { - deduplicated: false, - job: completed, - }) - } - - pub async fn resolve_relays_for_request( - &self, - pubkey: &str, - request: &PublishEventRequest, - ) -> Result<PublishRelayResolution, PublishProxyError> { - match request.relay_policy { - PublishRelayPolicy::ExplicitOnly => self.resolve_request_relays(&request.relays).await, - PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault => { - if !request.relays.is_empty() { - self.resolve_request_relays(&request.relays).await - } else { - self.resolve_author_or_default_relays(pubkey).await - } - } - PublishRelayPolicy::AuthorWriteThenDaemonDefault => { - self.resolve_author_or_default_relays(pubkey).await - } - PublishRelayPolicy::DaemonDefaultOnly => self.resolve_daemon_default_relays().await, - } - } - - async fn resolve_author_or_default_relays( - &self, - pubkey: &str, - ) -> Result<PublishRelayResolution, PublishProxyError> { - let mut author_relays = self.resolve_author_write_relays(pubkey).await?; - if author_relays.targets.is_empty() { - let mut daemon_defaults = self.resolve_daemon_default_relays().await?; - daemon_defaults.outcomes.append(&mut author_relays.outcomes); - Ok(daemon_defaults) - } else { - Ok(author_relays) - } - } - - async fn resolve_request_relays( - &self, - relays: &[String], - ) -> Result<PublishRelayResolution, PublishProxyError> { - let mut targets = Vec::new(); - let mut outcomes = Vec::new(); - for relay in relays { - match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { - Ok(url) => { - self.push_checked_relay_target( - &mut targets, - &mut outcomes, - url, - PublishRelaySource::Request, - ) - .await; - } - Err(error) => outcomes.push(PublishRelayOutcome { - relay_url: relay.trim().to_owned(), - source: PublishRelaySource::Request, - attempted: false, - outcome_kind: PublishRelayOutcomeKind::RelayUrlRejected, - message: Some(error.to_string()), - latency_ms: None, - }), - } - } - Ok(PublishRelayResolution { targets, outcomes }) - } - - async fn resolve_author_write_relays( - &self, - pubkey: &str, - ) -> Result<PublishRelayResolution, PublishProxyError> { - let cached = self.store.cached_author_write_relays(pubkey)?; - let mut cached_resolution = self.resolve_author_relay_inputs(&cached).await?; - if !cached_resolution.targets.is_empty() { - return Ok(cached_resolution); - } - if self.config.author_relay_discovery_relays.is_empty() { - return Ok(cached_resolution); - } - let mut discovery_targets = self - .resolve_config_relays( - &self.config.author_relay_discovery_relays, - PublishRelaySource::DaemonDefault, - ) - .await?; - if discovery_targets.targets.is_empty() { - discovery_targets - .outcomes - .append(&mut cached_resolution.outcomes); - return Ok(discovery_targets); - } - let discovered = self - .author_relay_discovery - .fetch_author_write_relays( - pubkey, - std::mem::take(&mut discovery_targets.targets), - self.config.connect_timeout_secs, - ) - .await?; - self.store.cache_author_write_relays(pubkey, &discovered)?; - let mut discovered_resolution = self.resolve_author_relay_inputs(&discovered).await?; - discovered_resolution - .outcomes - .append(&mut cached_resolution.outcomes); - discovered_resolution - .outcomes - .append(&mut discovery_targets.outcomes); - Ok(discovered_resolution) - } - - async fn resolve_author_relay_inputs( - &self, - relays: &[String], - ) -> Result<PublishRelayResolution, PublishProxyError> { - let mut targets = Vec::new(); - let mut outcomes = Vec::new(); - for relay in relays { - match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { - Ok(url) => { - self.push_checked_relay_target( - &mut targets, - &mut outcomes, - url, - PublishRelaySource::AuthorWrite, - ) - .await; - } - Err(error) => outcomes.push(PublishRelayOutcome { - relay_url: relay.trim().to_owned(), - source: PublishRelaySource::AuthorWrite, - attempted: false, - outcome_kind: PublishRelayOutcomeKind::RelayUrlRejected, - message: Some(error.to_string()), - latency_ms: None, - }), - } - } - Ok(PublishRelayResolution { targets, outcomes }) - } - - async fn resolve_daemon_default_relays( - &self, - ) -> Result<PublishRelayResolution, PublishProxyError> { - self.resolve_config_relays( - &self.config.daemon_default_publish_relays, - PublishRelaySource::DaemonDefault, - ) - .await - } - - async fn resolve_config_relays( - &self, - relays: &[String], - source: PublishRelaySource, - ) -> Result<PublishRelayResolution, PublishProxyError> { - let mut targets = Vec::new(); - let mut outcomes = Vec::new(); - for relay in relays { - match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { - Ok(url) => { - self.push_checked_relay_target(&mut targets, &mut outcomes, url, source) - .await; - } - Err(error) => outcomes.push(PublishRelayOutcome { - relay_url: relay.trim().to_owned(), - source, - attempted: false, - outcome_kind: PublishRelayOutcomeKind::RelayUrlRejected, - message: Some(error.to_string()), - latency_ms: None, - }), - } - } - Ok(PublishRelayResolution { targets, outcomes }) - } - - async fn push_checked_relay_target( - &self, - targets: &mut Vec<ResolvedPublishRelay>, - outcomes: &mut Vec<PublishRelayOutcome>, - url: RadrootsRelayUrl, - source: PublishRelaySource, - ) { - if relay_url_policy(&self.config) == RadrootsRelayUrlPolicy::Localhost { - push_resolved_relay(targets, url, source); - return; - } - match self.resolver.resolve(&url).await { - Ok(addresses) if addresses.is_empty() => { - outcomes.push(relay_resolution_connection_failure( - url.as_str(), - source, - "dns lookup returned no addresses", - )); - } - Ok(addresses) => match url.validate_public_resolved_ip_addrs(addresses) { - Ok(()) => push_resolved_relay(targets, url, source), - Err(error) => outcomes.push(PublishRelayOutcome { - relay_url: url.as_str().to_owned(), - source, - attempted: false, - outcome_kind: PublishRelayOutcomeKind::RelayUrlRejected, - message: Some(error.to_string()), - latency_ms: None, - }), - }, - Err(error) => outcomes.push(relay_resolution_connection_failure( - url.as_str(), - source, - format!("dns lookup failed: {error}"), - )), - } - } - - async fn complete_job_execution( - &self, - job_id: &str, - signed_event: RadrootsSignedNostrEvent, - delivery_policy: PublishDeliveryPolicy, - timeout_ms: u64, - resolution: PublishRelayResolution, - ) -> Result<PublishJobView, PublishProxyError> { - if resolution.targets.is_empty() { - let status = if resolution - .outcomes - .iter() - .any(|outcome| outcome.outcome_kind.is_retryable()) - { - PublishJobStatus::DeliveryUnsatisfiedRetryable - } else { - PublishJobStatus::Rejected - }; - let last_error = if status == PublishJobStatus::DeliveryUnsatisfiedRetryable { - "delivery_unsatisfied" - } else { - "no_publish_relays" - }; - self.store.complete_publish_job( - job_id, - status, - resolution.outcomes, - Some(last_error.to_owned()), - )?; - return self.store.job_by_id(job_id); - } - let required_ack_count = delivery_policy.required_ack_count(resolution.targets.len()); - if required_ack_count > resolution.targets.len() { - self.store.complete_publish_job( - job_id, - PublishJobStatus::Rejected, - resolution.outcomes, - Some("delivery_quorum_exceeds_relay_count".to_owned()), - )?; - return self.store.job_by_id(job_id); - } - let source_by_relay = resolution.source_by_relay(); - let target_set = RadrootsRelayTargetSet::from_urls( - resolution - .targets - .iter() - .map(|target| target.url.clone()) - .collect(), - )?; - let publish_request = - RadrootsRelayPublishRequest::new(signed_event, target_set, current_unix_millis()) - .with_accepted_quorum(required_ack_count); - let started = Instant::now(); - let publish_timeout = Duration::from_millis(timeout_ms); - let receipts = - match tokio::time::timeout(publish_timeout, self.publish_with_adapter(publish_request)) - .await - { - Ok(Ok(receipts)) => receipts, - Ok(Err(error)) => transport_error_receipts(&resolution.targets, error), - Err(_) => timeout_receipts(&resolution.targets), - }; - let latency_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX); - let mut outcomes = resolution.outcomes; - outcomes.extend(receipts.into_iter().map(|receipt| { - publish_outcome_from_receipt(receipt, &source_by_relay, Some(latency_ms)) - })); - let status = delivery_status(&delivery_policy, resolution.targets.len(), &outcomes); - let last_error = if status == PublishJobStatus::DeliverySatisfied { - None - } else { - Some("delivery_unsatisfied".to_owned()) - }; - self.store - .complete_publish_job(job_id, status, outcomes, last_error)?; - self.store.job_by_id(job_id) - } - - async fn publish_with_adapter( - &self, - request: RadrootsRelayPublishRequest, - ) -> Result<Vec<RadrootsRelayPublishRelayReceipt>, PublishProxyError> { - if let Some(publisher) = &self.publisher { - return publisher - .publish(request) - .await - .map_err(PublishProxyError::Relay); - } - let adapter = RadrootsNostrClientPublishAdapter::new(RadrootsNostrClient::new_signerless()); - adapter - .publish(request) - .await - .map_err(PublishProxyError::Relay) - } -} - -#[derive(Clone)] -pub struct PublishProxyStore { - inner: Arc<Mutex<Connection>>, -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum PublishJobVisibility { - Own, - Admin, -} - -impl FromStr for PublishJobVisibility { - type Err = PublishProxyError; - - fn from_str(value: &str) -> Result<Self, Self::Err> { - match value { - "own" => Ok(Self::Own), - "admin" => Ok(Self::Admin), - other => Err(PublishProxyError::InvalidScope(format!( - "unknown job visibility `{other}`" - ))), - } - } -} - -impl fmt::Display for PublishJobVisibility { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - Self::Own => f.write_str("own"), - Self::Admin => f.write_str("admin"), - } - } -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PublishPrincipalInit { - pub label: String, - pub token_hash: String, - pub allowed_pubkeys: Vec<String>, - pub allowed_kinds: Vec<u32>, - pub allowed_relay_policies: Vec<PublishRelayPolicy>, - pub allow_request_relays: bool, - pub job_visibility: PublishJobVisibility, - pub expires_at_unix: Option<i64>, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PublishPrincipal { - pub principal_id: String, - pub label: String, - pub allowed_pubkeys: Vec<String>, - pub allowed_kinds: Vec<u32>, - pub allowed_relay_policies: Vec<PublishRelayPolicy>, - pub allow_request_relays: bool, - pub job_visibility: PublishJobVisibility, - pub expires_at_unix: Option<i64>, -} - -impl PublishPrincipal { - pub fn allows_event(&self, request: &PublishEventRequest) -> Result<(), PublishProxyError> { - ensure_lower_hex("pubkey", request.event.pubkey.as_str(), 64)?; - if !self - .allowed_pubkeys - .iter() - .any(|pubkey| pubkey == &request.event.pubkey) - { - return Err(PublishProxyError::InvalidScope( - "principal is not allowed to publish for event pubkey".to_owned(), - )); - } - if !self.allowed_kinds.contains(&request.event.kind) { - return Err(PublishProxyError::InvalidScope( - "principal is not allowed to publish event kind".to_owned(), - )); - } - if !self.allowed_relay_policies.contains(&request.relay_policy) { - return Err(PublishProxyError::InvalidScope( - "principal is not allowed to use requested relay policy".to_owned(), - )); - } - if !self.allow_request_relays && !request.relays.is_empty() { - return Err(PublishProxyError::InvalidScope( - "principal is not allowed to provide request relays".to_owned(), - )); - } - Ok(()) - } - - fn can_read_job(&self, principal_id: &str) -> bool { - self.job_visibility == PublishJobVisibility::Admin || self.principal_id == principal_id - } -} - -#[derive(Debug, Clone)] -pub struct PublishJobInsert { - pub principal_id: String, - pub idempotency_key: Option<String>, - pub request: PublishEventRequest, - pub request_fingerprint: String, - pub effective_relay_count: usize, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct ResolvedPublishRelay { - pub url: RadrootsRelayUrl, - pub source: PublishRelaySource, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PublishRelayResolution { - pub targets: Vec<ResolvedPublishRelay>, - pub outcomes: Vec<PublishRelayOutcome>, -} - -impl PublishRelayResolution { - fn source_by_relay(&self) -> BTreeMap<String, PublishRelaySource> { - self.targets - .iter() - .map(|target| (target.url.as_str().to_owned(), target.source)) - .collect() - } -} - -pub(crate) type PublishRelayResolveFuture<'a> = - Pin<Box<dyn Future<Output = Result<Vec<IpAddr>, std::io::Error>> + Send + 'a>>; - -pub(crate) trait PublishRelayResolver: Send + Sync { - fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a>; -} - -type PublishAuthorRelayDiscoveryFuture<'a> = - Pin<Box<dyn Future<Output = Result<Vec<String>, PublishProxyError>> + Send + 'a>>; - -trait PublishAuthorRelayDiscovery: Send + Sync { - fn fetch_author_write_relays<'a>( - &'a self, - pubkey: &'a str, - discovery_targets: Vec<ResolvedPublishRelay>, - connect_timeout_secs: u64, - ) -> PublishAuthorRelayDiscoveryFuture<'a>; -} - -#[derive(Debug)] -struct SystemPublishRelayResolver; - -impl PublishRelayResolver for SystemPublishRelayResolver { - fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a> { - Box::pin(async move { - let (host, port) = relay_socket_target(url)?; - let addrs = tokio::net::lookup_host((host.as_str(), port)).await?; - Ok(addrs.map(|addr| addr.ip()).collect()) - }) - } -} - -#[derive(Debug)] -struct NostrPublishAuthorRelayDiscovery; - -impl PublishAuthorRelayDiscovery for NostrPublishAuthorRelayDiscovery { - fn fetch_author_write_relays<'a>( - &'a self, - pubkey: &'a str, - discovery_targets: Vec<ResolvedPublishRelay>, - connect_timeout_secs: u64, - ) -> PublishAuthorRelayDiscoveryFuture<'a> { - Box::pin(async move { - let Ok(public_key) = RadrootsNostrPublicKey::from_hex(pubkey) else { - return Ok(Vec::new()); - }; - let client = RadrootsNostrClient::new_signerless(); - for target in discovery_targets { - if client.add_read_relay(target.url.as_str()).await.is_err() { - return Ok(Vec::new()); - } - } - let filter = RadrootsNostrFilter::new() - .author(public_key) - .kind(RadrootsNostrKind::Custom(10_002)) - .limit(10); - let timeout = Duration::from_secs(connect_timeout_secs); - let Ok(events) = client.fetch_events(filter, timeout).await else { - return Ok(Vec::new()); - }; - let Some(event) = events.into_iter().max_by(|left, right| { - left.created_at - .as_secs() - .cmp(&right.created_at.as_secs()) - .then_with(|| left.id.to_hex().cmp(&right.id.to_hex())) - }) else { - return Ok(Vec::new()); - }; - Ok(author_write_relays_from_nip65_event(&event)) - }) - } -} - -impl PublishProxyStore { - pub fn open(path: PathBuf) -> Result<Self, PublishProxyError> { - if let Some(parent) = path - .parent() - .filter(|parent| !parent.as_os_str().is_empty()) - { - std::fs::create_dir_all(parent)?; - } - let connection = Connection::open(path)?; - Self::from_connection(connection) - } - - pub fn memory() -> Result<Self, PublishProxyError> { - Self::from_connection(Connection::open_in_memory()?) - } - - fn from_connection(connection: Connection) -> Result<Self, PublishProxyError> { - connection.execute_batch( - r#" - PRAGMA foreign_keys = ON; - CREATE TABLE IF NOT EXISTS publish_proxy_principals ( - principal_id TEXT PRIMARY KEY NOT NULL, - label TEXT NOT NULL, - token_hash TEXT NOT NULL UNIQUE, - allowed_pubkeys_json TEXT NOT NULL, - allowed_kinds_json TEXT NOT NULL, - allowed_relay_policies_json TEXT NOT NULL, - allow_request_relays INTEGER NOT NULL, - job_visibility TEXT NOT NULL, - expires_at_unix INTEGER, - revoked_at_unix INTEGER, - created_at_unix INTEGER NOT NULL - ); - CREATE TABLE IF NOT EXISTS publish_proxy_jobs ( - job_id TEXT PRIMARY KEY NOT NULL, - principal_id TEXT NOT NULL, - idempotency_key TEXT, - request_fingerprint TEXT NOT NULL, - status TEXT NOT NULL, - event_id TEXT NOT NULL, - event_pubkey TEXT NOT NULL, - event_kind INTEGER NOT NULL, - relay_policy_json TEXT NOT NULL, - delivery_policy_json TEXT NOT NULL, - requested_relay_count INTEGER NOT NULL, - effective_relay_count INTEGER NOT NULL, - request_json TEXT NOT NULL, - requested_at_ms INTEGER NOT NULL, - updated_at_ms INTEGER NOT NULL, - completed_at_ms INTEGER, - last_error TEXT, - FOREIGN KEY(principal_id) REFERENCES publish_proxy_principals(principal_id) - ); - CREATE UNIQUE INDEX IF NOT EXISTS publish_proxy_jobs_principal_idempotency_idx - ON publish_proxy_jobs(principal_id, idempotency_key) - WHERE idempotency_key IS NOT NULL; - CREATE TABLE IF NOT EXISTS publish_proxy_relay_results ( - job_id TEXT NOT NULL, - relay_url TEXT NOT NULL, - source TEXT NOT NULL, - attempted INTEGER NOT NULL, - outcome_kind TEXT NOT NULL, - message TEXT, - latency_ms INTEGER, - updated_at_ms INTEGER NOT NULL, - PRIMARY KEY(job_id, relay_url), - FOREIGN KEY(job_id) REFERENCES publish_proxy_jobs(job_id) - ); - CREATE TABLE IF NOT EXISTS publish_proxy_relay_list_cache ( - pubkey TEXT PRIMARY KEY NOT NULL, - relays_json TEXT NOT NULL, - updated_at_ms INTEGER NOT NULL - ); - "#, - )?; - migrate_schema(&connection)?; - recover_interrupted_publish_jobs(&connection)?; - connection.pragma_update(None, "user_version", SCHEMA_VERSION)?; - Ok(Self { - inner: Arc::new(Mutex::new(connection)), - }) - } - - pub fn create_principal( - &self, - input: PublishPrincipalInit, - ) -> Result<PublishPrincipal, PublishProxyError> { - validate_principal_init(&input)?; - let principal_id = Uuid::new_v4().to_string(); - let now = current_unix_secs(); - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - connection.execute( - r#" - INSERT INTO publish_proxy_principals ( - principal_id, - label, - token_hash, - allowed_pubkeys_json, - allowed_kinds_json, - allowed_relay_policies_json, - allow_request_relays, - job_visibility, - expires_at_unix, - revoked_at_unix, - created_at_unix - ) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10) - "#, - params![ - principal_id, - input.label.trim(), - input.token_hash, - serde_json::to_string(&input.allowed_pubkeys)?, - serde_json::to_string(&input.allowed_kinds)?, - serde_json::to_string(&input.allowed_relay_policies)?, - input.allow_request_relays, - input.job_visibility.to_string(), - input.expires_at_unix, - now, - ], - )?; - drop(connection); - self.principal_by_id(principal_id.as_str())? - .ok_or_else(|| PublishProxyError::InvalidScope("created principal missing".to_owned())) - } - - pub fn principal_for_token_hash( - &self, - token_hash: &str, - ) -> Result<Option<PublishPrincipal>, PublishProxyError> { - let now = current_unix_secs(); - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let principal = connection - .query_row( - r#" - SELECT - principal_id, - label, - allowed_pubkeys_json, - allowed_kinds_json, - allowed_relay_policies_json, - allow_request_relays, - job_visibility, - expires_at_unix - FROM publish_proxy_principals - WHERE token_hash = ?1 - AND revoked_at_unix IS NULL - AND (expires_at_unix IS NULL OR expires_at_unix > ?2) - "#, - params![token_hash, now], - principal_from_row, - ) - .optional()?; - Ok(principal) - } - - pub fn principal_by_id( - &self, - principal_id: &str, - ) -> Result<Option<PublishPrincipal>, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let principal = connection - .query_row( - r#" - SELECT - principal_id, - label, - allowed_pubkeys_json, - allowed_kinds_json, - allowed_relay_policies_json, - allow_request_relays, - job_visibility, - expires_at_unix - FROM publish_proxy_principals - WHERE principal_id = ?1 - "#, - params![principal_id], - principal_from_row, - ) - .optional()?; - Ok(principal) - } - - pub fn record_publish_job( - &self, - insert: PublishJobInsert, - ) -> Result<PublishEventResponse, PublishProxyError> { - if let Some(idempotency_key) = insert.idempotency_key.as_deref() { - if let Some(existing) = - self.job_for_principal_id_and_key(insert.principal_id.as_str(), idempotency_key)? - { - if existing.request_fingerprint != insert.request_fingerprint { - return Err(PublishProxyError::IdempotencyConflict( - idempotency_key.to_owned(), - )); - } - return Ok(PublishEventResponse { - deduplicated: true, - job: existing.view, - }); - } - } - - let job_id = Uuid::new_v4().to_string(); - let now = current_unix_millis(); - let request_json = serde_json::to_string(&insert.request)?; - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let insert_result = connection.execute( - r#" - INSERT INTO publish_proxy_jobs ( - job_id, - principal_id, - idempotency_key, - request_fingerprint, - status, - event_id, - event_pubkey, - event_kind, - relay_policy_json, - delivery_policy_json, - requested_relay_count, - effective_relay_count, - request_json, - requested_at_ms, - updated_at_ms - ) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15) - "#, - params![ - job_id, - insert.principal_id, - insert.idempotency_key, - insert.request_fingerprint, - serde_json::to_string(&PublishJobStatus::Publishing)?, - insert.request.event.id, - insert.request.event.pubkey, - insert.request.event.kind, - serde_json::to_string(&insert.request.relay_policy)?, - serde_json::to_string(&insert.request.delivery_policy)?, - insert.request.relays.len(), - insert.effective_relay_count, - request_json, - now, - now, - ], - ); - match insert_result { - Ok(_) => {} - Err(rusqlite::Error::SqliteFailure(error, _)) - if error.code == rusqlite::ErrorCode::ConstraintViolation => - { - return Err(PublishProxyError::IdempotencyConflict( - "idempotency key conflicts with an existing publish job".to_owned(), - )); - } - Err(error) => return Err(error.into()), - } - drop(connection); - let job = self.job_by_id(job_id.as_str())?; - Ok(PublishEventResponse { - deduplicated: false, - job, - }) - } - - pub fn job_by_id_for_principal( - &self, - job_id: &str, - principal: &PublishPrincipal, - ) -> Result<Option<PublishJobView>, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let sql = job_select_sql("WHERE job_id = ?1"); - let row = connection - .query_row(sql.as_str(), params![job_id], job_from_row) - .optional()?; - drop(connection); - let Some(mut job) = row else { - return Ok(None); - }; - if !principal.can_read_job(job.principal_id.as_str()) { - return Ok(None); - } - job.view.relays = self.relay_outcomes(job.view.job_id.as_str())?; - finalize_job_view(&mut job.view); - Ok(Some(job.view)) - } - - pub fn list_jobs_for_principal( - &self, - principal: &PublishPrincipal, - limit: usize, - ) -> Result<Vec<PublishJobView>, PublishProxyError> { - let limit = i64::try_from(limit.clamp(1, 200)).unwrap_or(200); - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let sql = if principal.job_visibility == PublishJobVisibility::Admin { - job_select_sql("ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?1") - } else { - job_select_sql( - "WHERE principal_id = ?1 ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?2", - ) - }; - let mut stmt = connection.prepare(sql.as_str())?; - let rows = if principal.job_visibility == PublishJobVisibility::Admin { - stmt.query_map(params![limit], job_from_row)? - .collect::<Result<Vec<_>, _>>()? - } else { - stmt.query_map(params![principal.principal_id, limit], job_from_row)? - .collect::<Result<Vec<_>, _>>()? - }; - drop(stmt); - drop(connection); - - rows.into_iter() - .map(|mut row| { - row.view.relays = self.relay_outcomes(row.view.job_id.as_str())?; - finalize_job_view(&mut row.view); - Ok(row.view) - }) - .collect() - } - - fn job_for_principal_id_and_key( - &self, - principal_id: &str, - idempotency_key: &str, - ) -> Result<Option<PublishJobRow>, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let sql = job_select_sql("WHERE principal_id = ?1 AND idempotency_key = ?2"); - let row = connection - .query_row( - sql.as_str(), - params![principal_id, idempotency_key], - job_from_row, - ) - .optional()?; - drop(connection); - let Some(mut job) = row else { - return Ok(None); - }; - job.view.relays = self.relay_outcomes(job.view.job_id.as_str())?; - finalize_job_view(&mut job.view); - Ok(Some(job)) - } - - pub fn job_by_id(&self, job_id: &str) -> Result<PublishJobView, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let sql = job_select_sql("WHERE job_id = ?1"); - let row = connection - .query_row(sql.as_str(), params![job_id], job_from_row) - .optional()?; - drop(connection); - let Some(mut job) = row else { - return Err(PublishProxyError::InvalidScope( - "unknown publish job".to_owned(), - )); - }; - job.view.relays = self.relay_outcomes(job.view.job_id.as_str())?; - finalize_job_view(&mut job.view); - Ok(job.view) - } - - pub fn complete_publish_job( - &self, - job_id: &str, - status: PublishJobStatus, - outcomes: Vec<PublishRelayOutcome>, - last_error: Option<String>, - ) -> Result<(), PublishProxyError> { - let now = current_unix_millis(); - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - connection.execute( - r#" - UPDATE publish_proxy_jobs - SET status = ?2, - updated_at_ms = ?3, - completed_at_ms = ?4, - last_error = ?5 - WHERE job_id = ?1 - "#, - params![ - job_id, - serde_json::to_string(&status)?, - now, - now, - last_error, - ], - )?; - connection.execute( - "DELETE FROM publish_proxy_relay_results WHERE job_id = ?1", - params![job_id], - )?; - for outcome in outcomes { - connection.execute( - r#" - INSERT OR REPLACE INTO publish_proxy_relay_results ( - job_id, - relay_url, - source, - attempted, - outcome_kind, - message, - latency_ms, - updated_at_ms - ) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) - "#, - params![ - job_id, - outcome.relay_url, - serde_json::to_string(&outcome.source)?, - outcome.attempted, - serde_json::to_string(&outcome.outcome_kind)?, - outcome.message, - outcome - .latency_ms - .and_then(|value| i64::try_from(value).ok()), - now, - ], - )?; - } - Ok(()) - } - - pub fn cached_author_write_relays( - &self, - pubkey: &str, - ) -> Result<Vec<String>, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let relays_json = connection - .query_row( - "SELECT relays_json FROM publish_proxy_relay_list_cache WHERE pubkey = ?1", - params![pubkey], - |row| row.get::<_, String>(0), - ) - .optional()?; - relays_json - .map(|value| serde_json::from_str(value.as_str()).map_err(PublishProxyError::from)) - .unwrap_or_else(|| Ok(Vec::new())) - } - - pub fn cache_author_write_relays( - &self, - pubkey: &str, - relays: &[String], - ) -> Result<(), PublishProxyError> { - let now = current_unix_millis(); - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - connection.execute( - r#" - INSERT INTO publish_proxy_relay_list_cache (pubkey, relays_json, updated_at_ms) - VALUES (?1, ?2, ?3) - ON CONFLICT(pubkey) DO UPDATE SET - relays_json = excluded.relays_json, - updated_at_ms = excluded.updated_at_ms - "#, - params![pubkey, serde_json::to_string(relays)?, now], - )?; - Ok(()) - } - - fn relay_outcomes(&self, job_id: &str) -> Result<Vec<PublishRelayOutcome>, PublishProxyError> { - let connection = self - .inner - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let mut stmt = connection.prepare( - r#" - SELECT relay_url, source, attempted, outcome_kind, message, latency_ms - FROM publish_proxy_relay_results - WHERE job_id = ?1 - ORDER BY relay_url - "#, - )?; - let outcomes = stmt - .query_map(params![job_id], relay_outcome_from_row)? - .collect::<Result<Vec<_>, _>>()?; - Ok(outcomes) - } -} - -struct PublishJobRow { - principal_id: String, - request_fingerprint: String, - view: PublishJobView, -} - -fn migrate_schema(connection: &Connection) -> Result<(), PublishProxyError> { - let version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; - if version < 2 { - if !table_has_column(connection, "publish_proxy_jobs", "request_fingerprint")? { - connection.execute( - "ALTER TABLE publish_proxy_jobs ADD COLUMN request_fingerprint TEXT NOT NULL DEFAULT ''", - [], - )?; - } - if !table_has_column(connection, "publish_proxy_jobs", "effective_relay_count")? { - connection.execute( - "ALTER TABLE publish_proxy_jobs ADD COLUMN effective_relay_count INTEGER NOT NULL DEFAULT 0", - [], - )?; - connection.execute( - "UPDATE publish_proxy_jobs SET effective_relay_count = requested_relay_count WHERE effective_relay_count = 0", - [], - )?; - } - } - Ok(()) -} - -fn recover_interrupted_publish_jobs(connection: &Connection) -> Result<(), PublishProxyError> { - let now = current_unix_millis(); - connection.execute( - r#" - UPDATE publish_proxy_jobs - SET status = ?1, - updated_at_ms = ?2, - completed_at_ms = ?3, - last_error = ?4 - WHERE status = ?5 - "#, - params![ - serde_json::to_string(&PublishJobStatus::DeliveryUnsatisfiedRetryable)?, - now, - now, - "publish_attempt_interrupted", - serde_json::to_string(&PublishJobStatus::Publishing)?, - ], - )?; - Ok(()) -} - -fn table_has_column( - connection: &Connection, - table: &str, - column: &str, -) -> Result<bool, PublishProxyError> { - let mut stmt = connection.prepare(format!("PRAGMA table_info({table})").as_str())?; - let columns = stmt - .query_map([], |row| row.get::<_, String>(1))? - .collect::<Result<Vec<_>, _>>()?; - Ok(columns.iter().any(|existing| existing == column)) -} - -fn job_select_sql(tail: &str) -> String { - format!( - r#" - SELECT - job_id, - principal_id, - request_fingerprint, - status, - event_id, - event_pubkey, - event_kind, - relay_policy_json, - delivery_policy_json, - effective_relay_count, - requested_at_ms, - completed_at_ms, - last_error - FROM publish_proxy_jobs - {tail} - "# - ) -} - -fn principal_from_row(row: &Row<'_>) -> Result<PublishPrincipal, rusqlite::Error> { - let visibility: String = row.get(6)?; - Ok(PublishPrincipal { - principal_id: row.get(0)?, - label: row.get(1)?, - allowed_pubkeys: json_column(row, 2)?, - allowed_kinds: json_column(row, 3)?, - allowed_relay_policies: json_column(row, 4)?, - allow_request_relays: row.get(5)?, - job_visibility: PublishJobVisibility::from_str(visibility.as_str()) - .map_err(|error| conversion_error(6, error))?, - expires_at_unix: row.get(7)?, - }) -} - -fn job_from_row(row: &Row<'_>) -> Result<PublishJobRow, rusqlite::Error> { - let status: PublishJobStatus = json_text(row, 3)?; - let relay_policy: PublishRelayPolicy = json_text(row, 7)?; - let delivery_policy: PublishDeliveryPolicy = json_text(row, 8)?; - let relay_count: i64 = row.get(9)?; - Ok(PublishJobRow { - principal_id: row.get(1)?, - request_fingerprint: row.get(2)?, - view: PublishJobView { - job_id: row.get(0)?, - status, - terminal: false, - delivery_satisfied: false, - event_id: row.get(4)?, - pubkey: row.get(5)?, - event_kind: row.get::<_, i64>(6)? as u32, - relay_policy, - delivery_policy, - relay_count: usize::try_from(relay_count).unwrap_or(0), - acknowledged_count: 0, - retryable_count: 0, - terminal_count: 0, - requested_at_ms: row.get(10)?, - completed_at_ms: row.get(11)?, - last_error: row.get(12)?, - relays: Vec::new(), - }, - }) -} - -fn relay_outcome_from_row(row: &Row<'_>) -> Result<PublishRelayOutcome, rusqlite::Error> { - let source: PublishRelaySource = json_text(row, 1)?; - let outcome_kind: PublishRelayOutcomeKind = json_text(row, 3)?; - Ok(PublishRelayOutcome { - relay_url: row.get(0)?, - source, - attempted: row.get(2)?, - outcome_kind, - message: row.get(4)?, - latency_ms: row - .get::<_, Option<i64>>(5)? - .map(|latency| u64::try_from(latency).unwrap_or(0)), - }) -} - -fn finalize_job_view(view: &mut PublishJobView) { - view.acknowledged_count = view - .relays - .iter() - .filter(|relay| relay.outcome_kind.counts_toward_quorum()) - .count(); - view.retryable_count = view - .relays - .iter() - .filter(|relay| relay.outcome_kind.is_retryable()) - .count(); - view.terminal_count = view - .relays - .iter() - .filter(|relay| relay.outcome_kind.is_terminal_failure()) - .count(); - view.terminal = matches!( - view.status, - PublishJobStatus::DeliverySatisfied - | PublishJobStatus::DeliveryUnsatisfiedTerminal - | PublishJobStatus::Rejected - ); - view.delivery_satisfied = view.status == PublishJobStatus::DeliverySatisfied; -} - -fn validate_principal_init(input: &PublishPrincipalInit) -> Result<(), PublishProxyError> { - if input.label.trim().is_empty() { - return Err(PublishProxyError::InvalidScope( - "principal label must not be empty".to_owned(), - )); - } - if !input.token_hash.starts_with(TOKEN_HASH_PREFIX) { - return Err(PublishProxyError::InvalidScope( - "principal token hash must use sha256 prefix".to_owned(), - )); - } - if input.allowed_pubkeys.is_empty() { - return Err(PublishProxyError::InvalidScope( - "principal must include at least one allowed pubkey".to_owned(), - )); - } - for pubkey in &input.allowed_pubkeys { - ensure_lower_hex("allowed_pubkey", pubkey, 64)?; - } - if input.allowed_kinds.is_empty() { - return Err(PublishProxyError::InvalidScope( - "principal must include at least one allowed kind".to_owned(), - )); - } - if input - .allowed_kinds - .iter() - .any(|kind| *kind > u16::MAX as u32) - { - return Err(PublishProxyError::InvalidScope( - "allowed kind exceeds publish proxy range".to_owned(), - )); - } - if input.allowed_relay_policies.is_empty() { - return Err(PublishProxyError::InvalidScope( - "principal must include at least one allowed relay policy".to_owned(), - )); - } - Ok(()) -} - -pub fn generate_bearer_token() -> String { - let bytes: [u8; 32] = rand::random(); - format!("{TOKEN_PREFIX}{}", hex_lower(&bytes)) -} - -pub fn hash_bearer_token(token: &str) -> String { - let mut hasher = Sha256::new(); - hasher.update(token.as_bytes()); - format!("{TOKEN_HASH_PREFIX}{}", hex_lower(&hasher.finalize())) -} - -fn hex_lower(bytes: &[u8]) -> String { - let mut output = String::with_capacity(bytes.len() * 2); - for byte in bytes { - use std::fmt::Write; - let _ = write!(&mut output, "{byte:02x}"); - } - output -} - -pub fn parse_relay_policy(value: &str) -> Result<PublishRelayPolicy, PublishProxyError> { - match value { - "explicit_only" => Ok(PublishRelayPolicy::ExplicitOnly), - "request_then_author_write_then_daemon_default" => { - Ok(PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault) - } - "author_write_then_daemon_default" => Ok(PublishRelayPolicy::AuthorWriteThenDaemonDefault), - "daemon_default_only" => Ok(PublishRelayPolicy::DaemonDefaultOnly), - other => Err(PublishProxyError::InvalidScope(format!( - "unknown relay policy `{other}`" - ))), - } -} - -fn signed_event_from_wire( - event: &SignedNostrEventWire, -) -> Result<RadrootsSignedNostrEvent, PublishProxyError> { - event - .validate() - .map_err(|error| PublishProxyError::InvalidSignedEvent(error.to_string()))?; - let created_at = u32::try_from(event.created_at).map_err(|_| { - PublishProxyError::InvalidSignedEvent( - "signed event created_at exceeds daemon-supported range".to_owned(), - ) - })?; - let raw_json = serde_json::to_string(event)?; - let radroots_event = RadrootsNostrEvent { - id: event.id.clone(), - author: event.pubkey.clone(), - created_at, - kind: event.kind, - tags: event.tags.clone(), - content: event.content.clone(), - sig: event.sig.clone(), - }; - match radroots_nostr_verify_event(&radroots_event) { - RadrootsNostrEventVerification::Verified => {} - verification => return Err(PublishProxyError::SignedEventVerification(verification)), - } - RadrootsSignedNostrEvent::new(RadrootsSignedNostrEventParts { - id: event.id.clone(), - pubkey: event.pubkey.clone(), - created_at, - kind: event.kind, - tags: event.tags.clone(), - content: event.content.clone(), - sig: event.sig.clone(), - raw_json, - }) - .map_err(PublishProxyError::from) -} - -fn request_intent_fingerprint( - principal_id: &str, - canonical_event_json: &str, - request: &PublishEventRequest, - effective_timeout_ms: u64, -) -> Result<String, PublishProxyError> { - #[derive(Serialize)] - struct FingerprintInput<'a> { - principal_id: &'a str, - canonical_event_json: &'a str, - relays: Vec<String>, - relay_policy: &'a PublishRelayPolicy, - delivery_policy: &'a PublishDeliveryPolicy, - effective_timeout_ms: u64, - } - - let input = FingerprintInput { - principal_id, - canonical_event_json, - relays: request - .relays - .iter() - .map(|relay| relay.trim().to_owned()) - .collect(), - relay_policy: &request.relay_policy, - delivery_policy: &request.delivery_policy, - effective_timeout_ms, - }; - let bytes = serde_json::to_vec(&input)?; - let mut hasher = Sha256::new(); - hasher.update(bytes); - Ok(hex_lower(&hasher.finalize())) -} - -fn effective_publish_timeout_ms( - config: &PublishProxyConfig, - timeout_ms: Option<u64>, -) -> Result<u64, PublishProxyError> { - let max_timeout_ms = config.connect_timeout_secs.saturating_mul(1_000); - match timeout_ms { - Some(0) => Err(PublishProxyError::InvalidSignedEvent( - "timeout_ms must be greater than zero".to_owned(), - )), - Some(timeout_ms) if timeout_ms > max_timeout_ms => { - Err(PublishProxyError::InvalidSignedEvent(format!( - "timeout_ms must be at most {max_timeout_ms}" - ))) - } - Some(timeout_ms) => Ok(timeout_ms), - None => Ok(max_timeout_ms), - } -} - -fn push_resolved_relay( - targets: &mut Vec<ResolvedPublishRelay>, - url: RadrootsRelayUrl, - source: PublishRelaySource, -) { - if !targets.iter().any(|target| target.url == url) { - targets.push(ResolvedPublishRelay { url, source }); - } -} - -fn relay_resolution_connection_failure( - relay_url: impl Into<String>, - source: PublishRelaySource, - message: impl Into<String>, -) -> PublishRelayOutcome { - PublishRelayOutcome { - relay_url: relay_url.into(), - source, - attempted: false, - outcome_kind: PublishRelayOutcomeKind::ConnectionFailed, - message: Some(message.into()), - latency_ms: None, - } -} - -fn relay_socket_target(url: &RadrootsRelayUrl) -> Result<(String, u16), std::io::Error> { - let parsed = url::Url::parse(url.as_str()) - .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?; - let host = parsed - .host_str() - .filter(|host| !host.is_empty()) - .ok_or_else(|| { - std::io::Error::new( - std::io::ErrorKind::InvalidInput, - "relay URL must include a DNS host", - ) - })? - .to_owned(); - let port = parsed.port_or_known_default().ok_or_else(|| { - std::io::Error::new( - std::io::ErrorKind::InvalidInput, - "relay URL scheme must have a default port", - ) - })?; - Ok((host, port)) -} - -fn relay_url_policy(config: &PublishProxyConfig) -> RadrootsRelayUrlPolicy { - match config.relay_url_policy { - crate::app::config::PublishProxyRelayUrlPolicy::Public => RadrootsRelayUrlPolicy::Public, - crate::app::config::PublishProxyRelayUrlPolicy::Localhost => { - RadrootsRelayUrlPolicy::Localhost - } - } -} - -fn author_write_relays_from_nip65_event( - event: &radroots_nostr::prelude::RadrootsNostrEvent, -) -> Vec<String> { - event - .tags - .iter() - .filter_map(|tag| { - let values = tag.as_slice(); - if values.first().map(String::as_str) != Some("r") { - return None; - } - let relay = values.get(1)?.trim(); - if relay.is_empty() { - return None; - } - if values.get(2).map(String::as_str) == Some("read") { - return None; - } - Some(relay.to_owned()) - }) - .collect() -} - -fn publish_outcome_from_receipt( - receipt: RadrootsRelayPublishRelayReceipt, - source_by_relay: &BTreeMap<String, PublishRelaySource>, - latency_ms: Option<u64>, -) -> PublishRelayOutcome { - let source = source_by_relay - .get(receipt.relay_url.as_str()) - .copied() - .unwrap_or(PublishRelaySource::DaemonDefault); - PublishRelayOutcome { - relay_url: receipt.relay_url, - source, - attempted: receipt.attempted, - outcome_kind: publish_outcome_kind(receipt.outcome.kind), - message: receipt.outcome.message, - latency_ms, - } -} - -fn publish_outcome_kind(kind: RadrootsRelayOutcomeKind) -> PublishRelayOutcomeKind { - match kind { - RadrootsRelayOutcomeKind::Accepted => PublishRelayOutcomeKind::Accepted, - RadrootsRelayOutcomeKind::DuplicateAccepted => PublishRelayOutcomeKind::DuplicateAccepted, - RadrootsRelayOutcomeKind::Blocked => PublishRelayOutcomeKind::Blocked, - RadrootsRelayOutcomeKind::RateLimited => PublishRelayOutcomeKind::RateLimited, - RadrootsRelayOutcomeKind::Invalid => PublishRelayOutcomeKind::Invalid, - RadrootsRelayOutcomeKind::PowRequired => PublishRelayOutcomeKind::PowRequired, - RadrootsRelayOutcomeKind::Restricted => PublishRelayOutcomeKind::Restricted, - RadrootsRelayOutcomeKind::AuthRequired => PublishRelayOutcomeKind::AuthRequired, - RadrootsRelayOutcomeKind::Muted => PublishRelayOutcomeKind::Muted, - RadrootsRelayOutcomeKind::Unsupported => PublishRelayOutcomeKind::Unsupported, - RadrootsRelayOutcomeKind::PaymentRequired => PublishRelayOutcomeKind::PaymentRequired, - RadrootsRelayOutcomeKind::Error => PublishRelayOutcomeKind::Error, - RadrootsRelayOutcomeKind::Timeout => PublishRelayOutcomeKind::Timeout, - RadrootsRelayOutcomeKind::ConnectionFailed => PublishRelayOutcomeKind::ConnectionFailed, - RadrootsRelayOutcomeKind::RelayUrlRejected => PublishRelayOutcomeKind::RelayUrlRejected, - RadrootsRelayOutcomeKind::SkippedAlreadyAccepted => { - PublishRelayOutcomeKind::SkippedAlreadyAccepted - } - RadrootsRelayOutcomeKind::Unknown => PublishRelayOutcomeKind::Unknown, - } -} - -fn delivery_status( - delivery_policy: &PublishDeliveryPolicy, - relay_count: usize, - outcomes: &[PublishRelayOutcome], -) -> PublishJobStatus { - let required = delivery_policy.required_ack_count(relay_count); - let acknowledged = outcomes - .iter() - .filter(|outcome| outcome.outcome_kind.counts_toward_quorum()) - .count(); - if acknowledged >= required { - return PublishJobStatus::DeliverySatisfied; - } - if outcomes - .iter() - .any(|outcome| outcome.outcome_kind.is_retryable()) - { - PublishJobStatus::DeliveryUnsatisfiedRetryable - } else { - PublishJobStatus::DeliveryUnsatisfiedTerminal - } -} - -fn timeout_receipts(targets: &[ResolvedPublishRelay]) -> Vec<RadrootsRelayPublishRelayReceipt> { - targets - .iter() - .map(|target| { - RadrootsRelayPublishRelayReceipt::attempted( - target.url.as_str(), - RadrootsRelayOutcome::timeout("timeout: publish attempt exceeded daemon bound"), - ) - }) - .collect() -} - -fn transport_error_receipts( - targets: &[ResolvedPublishRelay], - error: PublishProxyError, -) -> Vec<RadrootsRelayPublishRelayReceipt> { - let message = format!("error: {error}"); - targets - .iter() - .map(|target| { - RadrootsRelayPublishRelayReceipt::attempted( - target.url.as_str(), - RadrootsRelayOutcome::connection_failed(message.clone()), - ) - }) - .collect() -} - -pub fn write_token_file(path: &Path, token: &str) -> Result<(), PublishProxyError> { - if let Some(parent) = path - .parent() - .filter(|parent| !parent.as_os_str().is_empty()) - { - std::fs::create_dir_all(parent)?; - } - let mut options = std::fs::OpenOptions::new(); - options.write(true).create_new(true); - #[cfg(unix)] - { - use std::os::unix::fs::OpenOptionsExt; - options.mode(0o600); - } - use std::io::Write; - let mut file = options.open(path)?; - file.write_all(token.as_bytes())?; - file.write_all(b"\n")?; - Ok(()) -} - -fn ensure_lower_hex( - field: &str, - value: &str, - expected_len: usize, -) -> Result<(), PublishProxyError> { - if value.len() == expected_len - && value - .bytes() - .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f')) - { - Ok(()) - } else { - Err(PublishProxyError::InvalidScope(format!( - "{field} must be {expected_len} lowercase hex characters" - ))) - } -} - -fn json_column<T: for<'de> Deserialize<'de>>( - row: &Row<'_>, - index: usize, -) -> Result<T, rusqlite::Error> { - let value: String = row.get(index)?; - serde_json::from_str(value.as_str()).map_err(|error| conversion_error(index, error)) -} - -fn json_text<T: for<'de> Deserialize<'de>>( - row: &Row<'_>, - index: usize, -) -> Result<T, rusqlite::Error> { - let value: String = row.get(index)?; - serde_json::from_str(value.as_str()).map_err(|error| conversion_error(index, error)) -} - -fn conversion_error<E>(index: usize, error: E) -> rusqlite::Error -where - E: std::error::Error + Send + Sync + 'static, -{ - rusqlite::Error::FromSqlConversionFailure(index, Type::Text, Box::new(error)) -} - -fn current_unix_secs() -> i64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|duration| duration.as_secs() as i64) - .unwrap_or_default() -} - -fn current_unix_millis() -> i64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|duration| duration.as_millis() as i64) - .unwrap_or_default() -} - -#[cfg(test)] -mod tests { - use super::{ - PublishJobInsert, PublishJobVisibility, PublishPrincipal, PublishPrincipalInit, - PublishProxy, PublishProxyError, PublishProxyStore, generate_bearer_token, - hash_bearer_token, parse_relay_policy, - }; - use crate::app::config::{PublishProxyConfig, PublishProxyRelayUrlPolicy}; - use nostr::JsonUtil; - use radroots_identity::RadrootsIdentity; - use radroots_nostr::prelude::{ - RadrootsNostrEventVerification, RadrootsNostrTimestamp, radroots_nostr_build_event, - }; - use radroots_publish_proxy_protocol::{ - PublishDeliveryPolicy, PublishEventRequest, PublishJobStatus, PublishRelayOutcomeKind, - PublishRelayPolicy, PublishRelaySource, SignedNostrEventWire, - }; - use radroots_relay_transport::{RadrootsMockRelayPublishAdapter, RadrootsRelayOutcome}; - use std::collections::BTreeMap; - use std::net::{IpAddr, Ipv4Addr}; - use std::sync::Arc; - - const RELAY_PRIMARY: &str = "wss://relay.example.com"; - const RELAY_SECONDARY: &str = "wss://relay-2.example.com"; - const RELAY_FORBIDDEN: &str = "wss://forbidden-relay.example.com"; - - fn event(pubkey: &str, kind: u32) -> SignedNostrEventWire { - SignedNostrEventWire { - id: "0".repeat(64), - pubkey: pubkey.to_owned(), - created_at: 1_700_000_000, - kind, - tags: vec![vec!["d".to_owned(), "listing-1".to_owned()]], - content: "{}".to_owned(), - sig: "1".repeat(128), - } - } - - fn request(pubkey: &str, kind: u32) -> PublishEventRequest { - PublishEventRequest { - event: event(pubkey, kind), - relays: Vec::new(), - relay_policy: PublishRelayPolicy::DaemonDefaultOnly, - delivery_policy: PublishDeliveryPolicy::Any, - idempotency_key: Some("idem-1".to_owned()), - timeout_ms: None, - } - } - - fn signed_event(identity: &RadrootsIdentity, content: &str) -> SignedNostrEventWire { - let event = radroots_nostr_build_event( - 30_402, - content, - vec![vec!["d".to_owned(), "listing-1".to_owned()]], - ) - .expect("event builder") - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(identity.keys()) - .expect("signed event"); - serde_json::from_str(event.as_json().as_str()).expect("event wire") - } - - fn publish_request( - event: SignedNostrEventWire, - relays: Vec<String>, - relay_policy: PublishRelayPolicy, - delivery_policy: PublishDeliveryPolicy, - idempotency_key: Option<&str>, - ) -> PublishEventRequest { - PublishEventRequest { - event, - relays, - relay_policy, - delivery_policy, - idempotency_key: idempotency_key.map(str::to_owned), - timeout_ms: Some(5_000), - } - } - - fn publish_proxy( - config: PublishProxyConfig, - ) -> (PublishProxy, RadrootsMockRelayPublishAdapter) { - publish_proxy_with_resolver(config, Arc::new(StaticPublishRelayResolver::new())) - } - - fn publish_proxy_with_resolver( - config: PublishProxyConfig, - resolver: Arc<dyn super::PublishRelayResolver>, - ) -> (PublishProxy, RadrootsMockRelayPublishAdapter) { - let adapter = RadrootsMockRelayPublishAdapter::new(); - let proxy = PublishProxy::memory(config) - .expect("proxy") - .with_relay_resolver(resolver) - .with_publisher(Arc::new(adapter.clone())); - (proxy, adapter) - } - - fn principal( - proxy: &PublishProxy, - pubkey: String, - policies: Vec<PublishRelayPolicy>, - allow_request_relays: bool, - visibility: PublishJobVisibility, - ) -> PublishPrincipal { - proxy - .store - .create_principal(PublishPrincipalInit { - label: "tester".to_owned(), - token_hash: hash_bearer_token(generate_bearer_token().as_str()), - allowed_pubkeys: vec![pubkey], - allowed_kinds: vec![30_402], - allowed_relay_policies: policies, - allow_request_relays, - job_visibility: visibility, - expires_at_unix: None, - }) - .expect("principal") - } - - fn config_with_defaults(relays: Vec<&str>) -> PublishProxyConfig { - PublishProxyConfig { - daemon_default_publish_relays: relays.into_iter().map(str::to_owned).collect(), - ..PublishProxyConfig::default() - } - } - - #[derive(Default)] - struct StaticPublishRelayResolver { - results: BTreeMap<String, Result<Vec<IpAddr>, String>>, - } - - impl StaticPublishRelayResolver { - fn new() -> Self { - Self::default() - } - - fn with_addresses(mut self, url: &str, addresses: Vec<IpAddr>) -> Self { - self.results.insert(url.to_owned(), Ok(addresses)); - self - } - - fn with_failure(mut self, url: &str, error: &str) -> Self { - self.results.insert(url.to_owned(), Err(error.to_owned())); - self - } - } - - impl super::PublishRelayResolver for StaticPublishRelayResolver { - fn resolve<'a>( - &'a self, - url: &'a radroots_relay_transport::RadrootsRelayUrl, - ) -> super::PublishRelayResolveFuture<'a> { - Box::pin(async move { - match self.results.get(url.as_str()) { - Some(Ok(addresses)) => Ok(addresses.clone()), - Some(Err(error)) => Err(std::io::Error::other(error.clone())), - None => Ok(vec![IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))]), - } - }) - } - } - - struct StaticPublishAuthorRelayDiscovery { - relays: Vec<String>, - } - - impl StaticPublishAuthorRelayDiscovery { - fn new(relays: Vec<&str>) -> Self { - Self { - relays: relays.into_iter().map(str::to_owned).collect(), - } - } - } - - impl super::PublishAuthorRelayDiscovery for StaticPublishAuthorRelayDiscovery { - fn fetch_author_write_relays<'a>( - &'a self, - _pubkey: &'a str, - _discovery_targets: Vec<super::ResolvedPublishRelay>, - _connect_timeout_secs: u64, - ) -> super::PublishAuthorRelayDiscoveryFuture<'a> { - let relays = self.relays.clone(); - Box::pin(async move { Ok(relays) }) - } - } - - #[test] - fn token_generation_and_hashing_do_not_store_plaintext() { - let token = generate_bearer_token(); - assert!(token.starts_with("rrd_pp_")); - let hash = hash_bearer_token(token.as_str()); - assert!(hash.starts_with("sha256:")); - assert!(!hash.contains(token.as_str())); - } - - #[test] - fn relay_policy_parser_accepts_contract_values() { - assert_eq!( - parse_relay_policy("explicit_only").expect("policy"), - PublishRelayPolicy::ExplicitOnly - ); - assert!(parse_relay_policy("unknown").is_err()); - } - - #[test] - fn storage_authenticates_hashed_tokens_and_scopes_jobs() { - let store = PublishProxyStore::memory().expect("store"); - let token = generate_bearer_token(); - let token_hash = hash_bearer_token(token.as_str()); - let principal = store - .create_principal(PublishPrincipalInit { - label: "tester".to_owned(), - token_hash: token_hash.clone(), - allowed_pubkeys: vec!["a".repeat(64)], - allowed_kinds: vec![30_402], - allowed_relay_policies: vec![PublishRelayPolicy::DaemonDefaultOnly], - allow_request_relays: false, - job_visibility: PublishJobVisibility::Own, - expires_at_unix: None, - }) - .expect("principal"); - assert_eq!( - store - .principal_for_token_hash(token_hash.as_str()) - .expect("lookup") - .expect("principal") - .principal_id, - principal.principal_id - ); - let denied = request("b".repeat(64).as_str(), 30_402); - assert!(principal.allows_event(&denied).is_err()); - - let accepted = request("a".repeat(64).as_str(), 30_402); - principal.allows_event(&accepted).expect("scope"); - let response = store - .record_publish_job(PublishJobInsert { - principal_id: principal.principal_id.clone(), - idempotency_key: Some("idem-1".to_owned()), - request: accepted.clone(), - request_fingerprint: "fingerprint-1".to_owned(), - effective_relay_count: 1, - }) - .expect("record job"); - assert!(!response.deduplicated); - let duplicate = store - .record_publish_job(PublishJobInsert { - principal_id: principal.principal_id.clone(), - idempotency_key: Some("idem-1".to_owned()), - request: accepted, - request_fingerprint: "fingerprint-1".to_owned(), - effective_relay_count: 1, - }) - .expect("dedupe"); - assert!(duplicate.deduplicated); - assert_eq!(duplicate.job.job_id, response.job.job_id); - assert_eq!( - store - .list_jobs_for_principal(&principal, 50) - .expect("jobs") - .len(), - 1 - ); - } - - #[test] - fn store_open_recovers_interrupted_publishing_jobs() { - let directory = tempfile::tempdir().expect("tempdir"); - let database_path = directory.path().join("publish-proxy.sqlite"); - let token_hash = hash_bearer_token(generate_bearer_token().as_str()); - let pubkey = "a".repeat(64); - let request = request(pubkey.as_str(), 30_402); - let job_id = { - let store = PublishProxyStore::open(database_path.clone()).expect("store"); - let principal = store - .create_principal(PublishPrincipalInit { - label: "tester".to_owned(), - token_hash, - allowed_pubkeys: vec![pubkey], - allowed_kinds: vec![30_402], - allowed_relay_policies: vec![PublishRelayPolicy::DaemonDefaultOnly], - allow_request_relays: false, - job_visibility: PublishJobVisibility::Own, - expires_at_unix: None, - }) - .expect("principal"); - let response = store - .record_publish_job(PublishJobInsert { - principal_id: principal.principal_id, - idempotency_key: Some("idem-interrupted".to_owned()), - request, - request_fingerprint: "fingerprint-interrupted".to_owned(), - effective_relay_count: 1, - }) - .expect("record job"); - assert_eq!(response.job.status, PublishJobStatus::Publishing); - response.job.job_id - }; - - let reopened = PublishProxyStore::open(database_path).expect("reopen store"); - let recovered = reopened.job_by_id(job_id.as_str()).expect("recovered job"); - assert_eq!( - recovered.status, - PublishJobStatus::DeliveryUnsatisfiedRetryable - ); - assert_eq!( - recovered.last_error.as_deref(), - Some("publish_attempt_interrupted") - ); - assert!(recovered.completed_at_ms.is_some()); - assert!(recovered.relays.is_empty()); - } - - #[tokio::test] - async fn publish_event_verifies_and_records_daemon_default_outcome() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let event = signed_event(&identity, "{}"); - let raw_event = serde_json::to_string(&event).expect("raw event"); - let response = proxy - .publish_event( - &principal, - publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-valid"), - ), - ) - .await - .expect("publish"); - - assert!(!response.deduplicated); - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - assert_eq!(response.job.relay_count, 1); - assert_eq!(response.job.acknowledged_count, 1); - assert_eq!(response.job.relays[0].relay_url, RELAY_PRIMARY); - assert_eq!( - response.job.relays[0].source, - PublishRelaySource::DaemonDefault - ); - assert_eq!(adapter.captured_raw_events(), vec![raw_event]); - } - - #[tokio::test] - async fn publish_event_rejects_tampered_content_before_publish() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let mut event = signed_event(&identity, "trusted"); - event.content = "tampered".to_owned(); - let error = proxy - .publish_event( - &principal, - publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect_err("tampered event should fail"); - - assert!(matches!( - error, - PublishProxyError::SignedEventVerification(RadrootsNostrEventVerification::IdMismatch) - )); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_rejects_wrong_signature_before_publish() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let mut event = signed_event(&identity, "{}"); - let replacement = if event.sig.starts_with('0') { "1" } else { "0" }; - event.sig.replace_range(0..1, replacement); - let error = proxy - .publish_event( - &principal, - publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect_err("wrong signature should fail"); - - assert!(matches!( - error, - PublishProxyError::SignedEventVerification( - RadrootsNostrEventVerification::SignatureInvalid - ) - )); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_rejects_malformed_wire_fields() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let mut event = signed_event(&identity, "{}"); - event.id = event.id.to_uppercase(); - let error = proxy - .publish_event( - &principal, - publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect_err("malformed field should fail"); - - assert!(matches!(error, PublishProxyError::InvalidSignedEvent(_))); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_uses_explicit_request_relays_when_allowed() { - let identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_SECONDARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault], - true, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - vec![RELAY_PRIMARY.to_owned()], - PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - assert_eq!(response.job.relays[0].relay_url, RELAY_PRIMARY); - assert_eq!(response.job.relays[0].source, PublishRelaySource::Request); - } - - #[tokio::test] - async fn publish_event_uses_cached_nip65_author_write_before_defaults() { - let identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_SECONDARY])); - proxy - .store - .cache_author_write_relays( - identity.public_key_hex().as_str(), - &[RELAY_PRIMARY.to_owned()], - ) - .expect("cache author relays"); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::AuthorWriteThenDaemonDefault], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::AuthorWriteThenDaemonDefault, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.relays[0].relay_url, RELAY_PRIMARY); - assert_eq!( - response.job.relays[0].source, - PublishRelaySource::AuthorWrite - ); - } - - #[tokio::test] - async fn publish_event_records_invalid_cached_author_write_relay() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_SECONDARY])); - proxy - .store - .cache_author_write_relays( - identity.public_key_hex().as_str(), - &[RELAY_PRIMARY.to_owned(), "not a cached relay".to_owned()], - ) - .expect("cache author relays"); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::AuthorWriteThenDaemonDefault], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::AuthorWriteThenDaemonDefault, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - let accepted = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == RELAY_PRIMARY) - .expect("accepted author relay"); - assert_eq!(accepted.source, PublishRelaySource::AuthorWrite); - assert!(accepted.attempted); - let rejected = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == "not a cached relay") - .expect("rejected cached author relay"); - assert_eq!(rejected.source, PublishRelaySource::AuthorWrite); - assert_eq!( - rejected.outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!rejected.attempted); - assert_eq!(adapter.captured_raw_events().len(), 1); - } - - #[tokio::test] - async fn publish_event_preserves_author_and_discovery_rejections_through_relay_selection() { - let identity = RadrootsIdentity::generate(); - let mut config = config_with_defaults(vec![RELAY_SECONDARY]); - config.author_relay_discovery_relays = vec!["not a discovery relay".to_owned()]; - let (proxy, adapter) = publish_proxy(config); - proxy - .store - .cache_author_write_relays( - identity.public_key_hex().as_str(), - &["not a cached relay".to_owned()], - ) - .expect("cache author relays"); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::AuthorWriteThenDaemonDefault], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::AuthorWriteThenDaemonDefault, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - let daemon_default = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == RELAY_SECONDARY) - .expect("daemon default relay"); - assert_eq!(daemon_default.source, PublishRelaySource::DaemonDefault); - assert!(daemon_default.attempted); - let cached = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == "not a cached relay") - .expect("cached author rejection"); - assert_eq!(cached.source, PublishRelaySource::AuthorWrite); - assert_eq!( - cached.outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!cached.attempted); - let discovery = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == "not a discovery relay") - .expect("discovery relay rejection"); - assert_eq!(discovery.source, PublishRelaySource::DaemonDefault); - assert_eq!( - discovery.outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!discovery.attempted); - assert_eq!(adapter.captured_raw_events().len(), 1); - } - - #[tokio::test] - async fn publish_event_preserves_discovery_and_discovered_author_rejections() { - let identity = RadrootsIdentity::generate(); - let mut config = config_with_defaults(vec![RELAY_PRIMARY]); - config.author_relay_discovery_relays = - vec![RELAY_PRIMARY.to_owned(), RELAY_FORBIDDEN.to_owned()]; - let resolver = StaticPublishRelayResolver::new().with_addresses( - RELAY_FORBIDDEN, - vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))], - ); - let adapter = RadrootsMockRelayPublishAdapter::new(); - let proxy = PublishProxy::memory(config) - .expect("proxy") - .with_relay_resolver(Arc::new(resolver)) - .with_author_relay_discovery(Arc::new(StaticPublishAuthorRelayDiscovery::new(vec![ - "not a discovered author relay", - RELAY_SECONDARY, - ]))) - .with_publisher(Arc::new(adapter.clone())); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::AuthorWriteThenDaemonDefault], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::AuthorWriteThenDaemonDefault, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - let accepted = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == RELAY_SECONDARY) - .expect("discovered author relay"); - assert_eq!(accepted.source, PublishRelaySource::AuthorWrite); - assert!(accepted.attempted); - let discovered = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == "not a discovered author relay") - .expect("discovered author rejection"); - assert_eq!(discovered.source, PublishRelaySource::AuthorWrite); - assert_eq!( - discovered.outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!discovered.attempted); - let discovery = response - .job - .relays - .iter() - .find(|relay| relay.relay_url == RELAY_FORBIDDEN) - .expect("discovery relay rejection"); - assert_eq!(discovery.source, PublishRelaySource::DaemonDefault); - assert_eq!( - discovery.outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!discovery.attempted); - assert_eq!(adapter.captured_raw_events().len(), 1); - } - - #[tokio::test] - async fn publish_event_records_no_publish_relays_failure() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(PublishProxyConfig::default()); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::Rejected); - assert_eq!( - response.job.last_error.as_deref(), - Some("no_publish_relays") - ); - assert!(response.job.relays.is_empty()); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_records_unsafe_request_relay_rejection() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(PublishProxyConfig::default()); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::ExplicitOnly], - true, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - vec!["wss://127.0.0.1:7777".to_owned()], - PublishRelayPolicy::ExplicitOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::Rejected); - assert_eq!(response.job.relays.len(), 1); - assert_eq!( - response.job.relays[0].outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!response.job.relays[0].attempted); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_rejects_forbidden_public_dns_destination_before_publish() { - let identity = RadrootsIdentity::generate(); - let resolver = StaticPublishRelayResolver::new() - .with_addresses(RELAY_PRIMARY, vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))]); - let (proxy, adapter) = publish_proxy_with_resolver( - config_with_defaults(vec![RELAY_PRIMARY]), - Arc::new(resolver), - ); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::Rejected); - assert_eq!(response.job.relays.len(), 1); - assert_eq!( - response.job.relays[0].outcome_kind, - PublishRelayOutcomeKind::RelayUrlRejected - ); - assert!(!response.job.relays[0].attempted); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_records_dns_failure_as_unattempted_retryable_outcome() { - let identity = RadrootsIdentity::generate(); - let resolver = StaticPublishRelayResolver::new().with_failure(RELAY_PRIMARY, "no records"); - let (proxy, adapter) = publish_proxy_with_resolver( - config_with_defaults(vec![RELAY_PRIMARY]), - Arc::new(resolver), - ); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!( - response.job.status, - PublishJobStatus::DeliveryUnsatisfiedRetryable - ); - assert_eq!( - response.job.last_error.as_deref(), - Some("delivery_unsatisfied") - ); - assert_eq!(response.job.relays.len(), 1); - assert_eq!( - response.job.relays[0].outcome_kind, - PublishRelayOutcomeKind::ConnectionFailed - ); - assert!(!response.job.relays[0].attempted); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_localhost_policy_skips_public_dns_guard() { - let identity = RadrootsIdentity::generate(); - let mut config = config_with_defaults(vec!["ws://localhost:7777"]); - config.relay_url_policy = PublishProxyRelayUrlPolicy::Localhost; - let resolver = StaticPublishRelayResolver::new() - .with_failure("ws://localhost:7777", "localhost resolution should not run"); - let (proxy, adapter) = publish_proxy_with_resolver(config, Arc::new(resolver)); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!(response.job.status, PublishJobStatus::DeliverySatisfied); - assert_eq!(response.job.relays[0].relay_url, "ws://localhost:7777"); - assert!(!adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_deduplicates_same_intent_and_conflicts_different_intent() { - let identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let request = publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-conflict"), - ); - let first = proxy - .publish_event(&principal, request.clone()) - .await - .expect("first"); - let duplicate = proxy - .publish_event(&principal, request) - .await - .expect("duplicate"); - - assert!(!first.deduplicated); - assert!(duplicate.deduplicated); - assert_eq!(duplicate.job.job_id, first.job.job_id); - - let conflict = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "changed"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-conflict"), - ), - ) - .await - .expect_err("conflict"); - assert!(matches!( - conflict, - PublishProxyError::IdempotencyConflict(_) - )); - } - - #[tokio::test] - async fn publish_event_rejects_zero_and_excessive_timeout_before_job_creation() { - let identity = RadrootsIdentity::generate(); - let (proxy, adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let mut zero = publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-zero-timeout"), - ); - zero.timeout_ms = Some(0); - let zero_error = proxy - .publish_event(&principal, zero) - .await - .expect_err("zero timeout should fail"); - assert!(matches!( - zero_error, - PublishProxyError::InvalidSignedEvent(_) - )); - - let mut excessive = publish_request( - signed_event(&identity, "changed"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-excessive-timeout"), - ); - excessive.timeout_ms = Some(10_001); - let excessive_error = proxy - .publish_event(&principal, excessive) - .await - .expect_err("excessive timeout should fail"); - assert!(matches!( - excessive_error, - PublishProxyError::InvalidSignedEvent(_) - )); - assert!( - proxy - .store - .list_jobs_for_principal(&principal, 50) - .expect("jobs") - .is_empty() - ); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_event_default_timeout_fingerprints_as_effective_timeout() { - let identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let event = signed_event(&identity, "{}"); - let mut default_timeout = publish_request( - event.clone(), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-default-timeout"), - ); - default_timeout.timeout_ms = None; - let mut explicit_default = publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-default-timeout"), - ); - explicit_default.timeout_ms = Some(10_000); - - let first = proxy - .publish_event(&principal, default_timeout) - .await - .expect("first"); - let duplicate = proxy - .publish_event(&principal, explicit_default) - .await - .expect("duplicate"); - assert!(!first.deduplicated); - assert!(duplicate.deduplicated); - assert_eq!(duplicate.job.job_id, first.job.job_id); - } - - #[tokio::test] - async fn publish_event_fingerprint_conflicts_on_different_effective_timeout() { - let identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let event = signed_event(&identity, "{}"); - let first = publish_request( - event.clone(), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-timeout-conflict"), - ); - let mut conflict = publish_request( - event, - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-timeout-conflict"), - ); - conflict.timeout_ms = Some(6_000); - - proxy.publish_event(&principal, first).await.expect("first"); - let error = proxy - .publish_event(&principal, conflict) - .await - .expect_err("timeout conflict"); - assert!(matches!(error, PublishProxyError::IdempotencyConflict(_))); - } - - #[tokio::test] - async fn publish_event_concurrency_limit_rejects_without_job_creation() { - let identity = RadrootsIdentity::generate(); - let mut config = config_with_defaults(vec![RELAY_PRIMARY]); - config.max_concurrent_publish_jobs = 1; - let (proxy, adapter) = publish_proxy(config); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let _permit = proxy.acquire_publish_permit().expect("permit"); - let error = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - Some("idem-concurrency"), - ), - ) - .await - .expect_err("concurrency limit"); - assert!(matches!(error, PublishProxyError::ConcurrencyLimit)); - assert!( - proxy - .store - .list_jobs_for_principal(&principal, 50) - .expect("jobs") - .is_empty() - ); - assert!(adapter.captured_raw_events().is_empty()); - } - - #[tokio::test] - async fn publish_jobs_respect_own_and_admin_visibility() { - let identity = RadrootsIdentity::generate(); - let other_identity = RadrootsIdentity::generate(); - let (proxy, _adapter) = publish_proxy(config_with_defaults(vec![RELAY_PRIMARY])); - let owner = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let other = principal( - &proxy, - other_identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let admin = principal( - &proxy, - other_identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Admin, - ); - let response = proxy - .publish_event( - &owner, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert!( - proxy - .store - .job_by_id_for_principal(response.job.job_id.as_str(), &other) - .expect("other read") - .is_none() - ); - assert!( - proxy - .store - .job_by_id_for_principal(response.job.job_id.as_str(), &admin) - .expect("admin read") - .is_some() - ); - } - - #[tokio::test] - async fn publish_event_records_retryable_relay_failures() { - let identity = RadrootsIdentity::generate(); - let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( - RELAY_PRIMARY, - RadrootsRelayOutcome::connection_failed("error: unavailable"), - ); - let proxy = PublishProxy::memory(config_with_defaults(vec![RELAY_PRIMARY])) - .expect("proxy") - .with_publisher(Arc::new(adapter)); - let principal = principal( - &proxy, - identity.public_key_hex(), - vec![PublishRelayPolicy::DaemonDefaultOnly], - false, - PublishJobVisibility::Own, - ); - let response = proxy - .publish_event( - &principal, - publish_request( - signed_event(&identity, "{}"), - Vec::new(), - PublishRelayPolicy::DaemonDefaultOnly, - PublishDeliveryPolicy::Any, - None, - ), - ) - .await - .expect("publish"); - - assert_eq!( - response.job.status, - PublishJobStatus::DeliveryUnsatisfiedRetryable - ); - assert_eq!(response.job.retryable_count, 1); - } -} diff --git a/src/core/state.rs b/src/core/state.rs @@ -4,8 +4,8 @@ use radroots_nostr::prelude::{ RadrootsNostrClient, RadrootsNostrKeys, RadrootsNostrMetadata, RadrootsNostrPublicKey, }; -use crate::app::config::{Nip46Config, PublishProxyConfig}; -use crate::core::publish_proxy::PublishProxy; +use crate::app::config::{Nip46Config, TransportPublishConfig}; +use crate::core::transport_publish::TransportPublish; #[derive(Clone)] pub struct Radrootsd { @@ -14,7 +14,7 @@ pub struct Radrootsd { pub pubkey: RadrootsNostrPublicKey, pub metadata: RadrootsNostrMetadata, pub info: serde_json::Value, - pub publish_proxy: PublishProxy, + pub transport_publish: TransportPublish, pub(crate) nip46_sessions: crate::core::nip46::session::Nip46SessionStore, pub nip46_config: Nip46Config, } @@ -23,7 +23,7 @@ impl Radrootsd { pub fn new( identity: RadrootsIdentity, metadata: RadrootsNostrMetadata, - publish_proxy_config: PublishProxyConfig, + transport_publish_config: TransportPublishConfig, nip46_config: Nip46Config, ) -> Result<Self> { let keys: RadrootsNostrKeys = identity.keys().clone(); @@ -34,9 +34,9 @@ impl Radrootsd { "build": option_env!("GIT_HASH").unwrap_or("unknown"), }); #[cfg(test)] - let publish_proxy = PublishProxy::memory(publish_proxy_config)?; + let transport_publish = TransportPublish::memory(transport_publish_config)?; #[cfg(not(test))] - let publish_proxy = PublishProxy::open(publish_proxy_config)?; + let transport_publish = TransportPublish::open(transport_publish_config)?; let nip46_sessions = crate::core::nip46::session::Nip46SessionStore::new(); Ok(Self { @@ -45,7 +45,7 @@ impl Radrootsd { pubkey, metadata, info, - publish_proxy, + transport_publish, nip46_sessions, nip46_config, }) @@ -55,7 +55,7 @@ impl Radrootsd { #[cfg(test)] mod tests { use super::Radrootsd; - use crate::app::config::{Nip46Config, PublishProxyConfig}; + use crate::app::config::{Nip46Config, TransportPublishConfig}; use radroots_identity::RadrootsIdentity; use radroots_nostr::prelude::RadrootsNostrMetadata; @@ -64,12 +64,12 @@ mod tests { let identity = RadrootsIdentity::generate(); let metadata: RadrootsNostrMetadata = serde_json::from_str(r#"{"name":"radrootsd-test"}"#).expect("metadata"); - let publish_proxy_cfg = PublishProxyConfig::default(); + let transport_publish_cfg = TransportPublishConfig::default(); let cfg = Nip46Config::default(); let state = Radrootsd::new( identity.clone(), metadata.clone(), - publish_proxy_cfg.clone(), + transport_publish_cfg.clone(), cfg.clone(), ) .expect("state"); @@ -77,8 +77,8 @@ mod tests { assert_eq!(state.pubkey, identity.public_key()); assert_eq!(state.metadata, metadata); assert_eq!( - state.publish_proxy.config.enabled, - publish_proxy_cfg.enabled + state.transport_publish.config.enabled, + transport_publish_cfg.enabled ); assert_eq!(state.nip46_config.session_ttl_secs, cfg.session_ttl_secs); assert_eq!(state.nip46_config.perms, cfg.perms); diff --git a/src/core/transport_publish.rs b/src/core/transport_publish.rs @@ -0,0 +1,3270 @@ +use std::collections::BTreeMap; +use std::fmt; +use std::future::Future; +use std::net::IpAddr; +use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::str::FromStr; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +use radroots_events::RadrootsNostrEvent; +use radroots_events::draft::{ + RadrootsDraftError, RadrootsSignedNostrEvent, RadrootsSignedNostrEventParts, +}; +use radroots_nostr::prelude::{ + RadrootsNostrClient, RadrootsNostrEventVerification, RadrootsNostrFilter, RadrootsNostrKind, + RadrootsNostrPublicKey, radroots_nostr_verify_event, +}; +use radroots_relay_transport::{ + RadrootsNostrClientPublishAdapter, RadrootsRelayOutcome, RadrootsRelayOutcomeKind, + RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, + RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, +}; +use radroots_transport::RadrootsTransportSatisfactionPolicy; +use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, SignedNostrEventWire, TransportPublishDeliveryPolicy, + TransportPublishEventRequest, TransportPublishEventResponse, TransportPublishJobStatus, + TransportPublishJobView, TransportPublishOutcomeKind, TransportPublishPreviewBehavior, + TransportPublishTarget, TransportPublishTargetOutcome, TransportPublishTargetPolicy, + TransportPublishTargetPolicyName, TransportPublishTargetSource, +}; +use rusqlite::types::Type; +use rusqlite::{Connection, OptionalExtension, Row, params}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use thiserror::Error; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; +use uuid::Uuid; + +use crate::app::config::TransportPublishConfig; + +const TOKEN_PREFIX: &str = "rrd_tp_"; +const TOKEN_HASH_PREFIX: &str = "sha256:"; +const SCHEMA_VERSION: i64 = 1; +const TRANSPORT_KIND_NOSTR: &str = "nostr"; +const TRANSPORT_KIND_RETICULUM: &str = "reticulum"; + +#[derive(Debug, Error)] +pub enum TransportPublishError { + #[error("transport publish storage error: {0}")] + Sqlite(#[from] rusqlite::Error), + #[error("transport publish json error: {0}")] + Json(#[from] serde_json::Error), + #[error("transport publish io error: {0}")] + Io(#[from] std::io::Error), + #[error("invalid transport publish scope: {0}")] + InvalidScope(String), + #[error("invalid signed Nostr event: {0}")] + InvalidSignedEvent(String), + #[error("signed Nostr event verification failed: {0:?}")] + SignedEventVerification(RadrootsNostrEventVerification), + #[error("signed Nostr event conversion error: {0}")] + Draft(#[from] RadrootsDraftError), + #[error("transport publish relay error: {0}")] + Relay(#[from] RadrootsRelayTransportError), + #[error("transport publish transport error: {0}")] + Transport(String), + #[error("transport publish concurrency limit reached")] + ConcurrencyLimit, + #[error("transport publish idempotency conflict for key `{0}`")] + IdempotencyConflict(String), +} + +#[derive(Clone)] +pub struct TransportPublish { + pub config: TransportPublishConfig, + pub store: TransportPublishStore, + publisher: Option<Arc<dyn RadrootsRelayPublishAdapter>>, + resolver: Arc<dyn PublishRelayResolver>, + author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, + publish_jobs: Arc<Semaphore>, +} + +impl TransportPublish { + pub fn open(config: TransportPublishConfig) -> Result<Self, TransportPublishError> { + let store = TransportPublishStore::open(config.database_path.clone())?; + let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); + Ok(Self { + config, + store, + publisher: None, + resolver: Arc::new(SystemPublishRelayResolver), + author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), + publish_jobs, + }) + } + + pub fn memory(config: TransportPublishConfig) -> Result<Self, TransportPublishError> { + let store = TransportPublishStore::memory()?; + let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); + Ok(Self { + config, + store, + publisher: None, + resolver: Arc::new(SystemPublishRelayResolver), + author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), + publish_jobs, + }) + } + + pub fn with_publisher(mut self, publisher: Arc<dyn RadrootsRelayPublishAdapter>) -> Self { + self.publisher = Some(publisher); + self + } + + #[cfg(test)] + pub(crate) fn with_relay_resolver(mut self, resolver: Arc<dyn PublishRelayResolver>) -> Self { + self.resolver = resolver; + self + } + + #[cfg(test)] + fn with_author_relay_discovery( + mut self, + author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, + ) -> Self { + self.author_relay_discovery = author_relay_discovery; + self + } + + fn acquire_publish_permit(&self) -> Result<OwnedSemaphorePermit, TransportPublishError> { + self.publish_jobs + .clone() + .try_acquire_owned() + .map_err(|_| TransportPublishError::ConcurrencyLimit) + } + + pub async fn publish_event( + &self, + principal: &PublishPrincipal, + request: TransportPublishEventRequest, + ) -> Result<TransportPublishEventResponse, TransportPublishError> { + request + .validate(self.config.max_targets_per_request) + .map_err(|error| { + TransportPublishError::InvalidSignedEvent(format!( + "publish request validation failed: {error}" + )) + })?; + principal.allows_event(&request)?; + let signed_event = signed_event_from_wire(&request.event)?; + if signed_event.raw_json.len() > self.config.max_event_bytes { + return Err(TransportPublishError::InvalidSignedEvent( + "signed event exceeds transport_publish max_event_bytes".to_owned(), + )); + } + let effective_timeout_ms = effective_publish_timeout_ms(&self.config, request.timeout_ms)?; + let _permit = self.acquire_publish_permit()?; + let request_fingerprint = request_intent_fingerprint( + principal.principal_id.as_str(), + signed_event.raw_json.as_str(), + &request, + effective_timeout_ms, + )?; + let resolution = self + .resolve_targets_for_request(signed_event.pubkey.as_str(), &request) + .await?; + let response = self.store.record_publish_job(PublishJobInsert { + principal_id: principal.principal_id.clone(), + idempotency_key: request.idempotency_key.clone(), + request: request.clone(), + request_fingerprint, + effective_target_count: resolution.target_count(), + })?; + if response.deduplicated { + return Ok(response); + } + let completed = self + .complete_job_execution( + response.job.job_id.as_str(), + signed_event, + request.delivery_policy.clone(), + effective_timeout_ms, + resolution, + ) + .await?; + Ok(TransportPublishEventResponse { + deduplicated: false, + job: completed, + }) + } + + pub async fn resolve_targets_for_request( + &self, + pubkey: &str, + request: &TransportPublishEventRequest, + ) -> Result<PublishRelayResolution, TransportPublishError> { + match &request.target_policy { + TransportPublishTargetPolicy::ExplicitTargets { targets } => { + self.resolve_explicit_targets(targets).await + } + TransportPublishTargetPolicy::Nostr { + source_policy, + relay_urls, + } => match source_policy { + NostrPublishTargetSourcePolicy::ExplicitOnly => { + self.resolve_request_relays(relay_urls).await + } + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault => { + if !relay_urls.is_empty() { + self.resolve_request_relays(relay_urls).await + } else { + self.resolve_author_or_default_relays(pubkey).await + } + } + NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault => { + self.resolve_author_or_default_relays(pubkey).await + } + NostrPublishTargetSourcePolicy::DaemonDefaultOnly => { + self.resolve_daemon_default_relays().await + } + }, + } + } + + async fn resolve_explicit_targets( + &self, + targets: &[TransportPublishTarget], + ) -> Result<PublishRelayResolution, TransportPublishError> { + let mut resolved = Vec::new(); + let mut outcomes = Vec::new(); + for target in targets { + match target.transport_kind.as_str() { + TRANSPORT_KIND_NOSTR => { + self.resolve_request_target(&mut resolved, &mut outcomes, target) + .await; + } + TRANSPORT_KIND_RETICULUM => { + outcomes.push(reticulum_preview_outcome(target)); + } + _ => outcomes.push(unsupported_transport_outcome(target)), + } + } + Ok(PublishRelayResolution { + targets: resolved, + outcomes, + }) + } + + async fn resolve_request_target( + &self, + targets: &mut Vec<ResolvedPublishRelay>, + outcomes: &mut Vec<TransportPublishTargetOutcome>, + target: &TransportPublishTarget, + ) { + match RadrootsRelayUrl::parse(target.endpoint_uri.as_str(), relay_url_policy(&self.config)) + { + Ok(url) => { + self.push_checked_relay_target( + targets, + outcomes, + url, + TransportPublishTargetSource::Request, + ) + .await; + } + Err(error) => outcomes.push(TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: target.endpoint_uri.trim().to_owned(), + source: TransportPublishTargetSource::Request, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::TargetRejected, + message: Some(error.to_string()), + latency_ms: None, + }), + } + } + async fn resolve_author_or_default_relays( + &self, + pubkey: &str, + ) -> Result<PublishRelayResolution, TransportPublishError> { + let mut author_relays = self.resolve_author_write_relays(pubkey).await?; + if author_relays.targets.is_empty() { + let mut daemon_defaults = self.resolve_daemon_default_relays().await?; + daemon_defaults.outcomes.append(&mut author_relays.outcomes); + Ok(daemon_defaults) + } else { + Ok(author_relays) + } + } + + async fn resolve_request_relays( + &self, + relays: &[String], + ) -> Result<PublishRelayResolution, TransportPublishError> { + let mut targets = Vec::new(); + let mut outcomes = Vec::new(); + for relay in relays { + match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { + Ok(url) => { + self.push_checked_relay_target( + &mut targets, + &mut outcomes, + url, + TransportPublishTargetSource::Request, + ) + .await; + } + Err(error) => outcomes.push(TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: relay.trim().to_owned(), + source: TransportPublishTargetSource::Request, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::TargetRejected, + message: Some(error.to_string()), + latency_ms: None, + }), + } + } + Ok(PublishRelayResolution { targets, outcomes }) + } + + async fn resolve_author_write_relays( + &self, + pubkey: &str, + ) -> Result<PublishRelayResolution, TransportPublishError> { + let cached = self.store.cached_author_write_relays(pubkey)?; + let mut cached_resolution = self.resolve_author_relay_inputs(&cached).await?; + if !cached_resolution.targets.is_empty() { + return Ok(cached_resolution); + } + if self.config.nostr.author_relay_discovery_relays.is_empty() { + return Ok(cached_resolution); + } + let mut discovery_targets = self + .resolve_config_relays( + &self.config.nostr.author_relay_discovery_relays, + TransportPublishTargetSource::DaemonDefault, + ) + .await?; + if discovery_targets.targets.is_empty() { + discovery_targets + .outcomes + .append(&mut cached_resolution.outcomes); + return Ok(discovery_targets); + } + let discovered = self + .author_relay_discovery + .fetch_author_write_relays( + pubkey, + std::mem::take(&mut discovery_targets.targets), + self.config.connect_timeout_secs, + ) + .await?; + self.store.cache_author_write_relays(pubkey, &discovered)?; + let mut discovered_resolution = self.resolve_author_relay_inputs(&discovered).await?; + discovered_resolution + .outcomes + .append(&mut cached_resolution.outcomes); + discovered_resolution + .outcomes + .append(&mut discovery_targets.outcomes); + Ok(discovered_resolution) + } + + async fn resolve_author_relay_inputs( + &self, + relays: &[String], + ) -> Result<PublishRelayResolution, TransportPublishError> { + let mut targets = Vec::new(); + let mut outcomes = Vec::new(); + for relay in relays { + match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { + Ok(url) => { + self.push_checked_relay_target( + &mut targets, + &mut outcomes, + url, + TransportPublishTargetSource::NostrAuthorWrite, + ) + .await; + } + Err(error) => outcomes.push(TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: relay.trim().to_owned(), + source: TransportPublishTargetSource::NostrAuthorWrite, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::TargetRejected, + message: Some(error.to_string()), + latency_ms: None, + }), + } + } + Ok(PublishRelayResolution { targets, outcomes }) + } + + async fn resolve_daemon_default_relays( + &self, + ) -> Result<PublishRelayResolution, TransportPublishError> { + self.resolve_config_relays( + &self.config.nostr.daemon_default_relays, + TransportPublishTargetSource::DaemonDefault, + ) + .await + } + + async fn resolve_config_relays( + &self, + relays: &[String], + source: TransportPublishTargetSource, + ) -> Result<PublishRelayResolution, TransportPublishError> { + let mut targets = Vec::new(); + let mut outcomes = Vec::new(); + for relay in relays { + match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { + Ok(url) => { + self.push_checked_relay_target(&mut targets, &mut outcomes, url, source) + .await; + } + Err(error) => outcomes.push(TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: relay.trim().to_owned(), + source, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::TargetRejected, + message: Some(error.to_string()), + latency_ms: None, + }), + } + } + Ok(PublishRelayResolution { targets, outcomes }) + } + + async fn push_checked_relay_target( + &self, + targets: &mut Vec<ResolvedPublishRelay>, + outcomes: &mut Vec<TransportPublishTargetOutcome>, + url: RadrootsRelayUrl, + source: TransportPublishTargetSource, + ) { + if relay_url_policy(&self.config) == RadrootsRelayUrlPolicy::Localhost { + push_resolved_relay(targets, url, source); + return; + } + match self.resolver.resolve(&url).await { + Ok(addresses) if addresses.is_empty() => { + outcomes.push(relay_resolution_connection_failure( + url.as_str(), + source, + "dns lookup returned no addresses", + )); + } + Ok(addresses) => match url.validate_public_resolved_ip_addrs(addresses) { + Ok(()) => push_resolved_relay(targets, url, source), + Err(error) => outcomes.push(TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: url.as_str().to_owned(), + source, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::TargetRejected, + message: Some(error.to_string()), + latency_ms: None, + }), + }, + Err(error) => outcomes.push(relay_resolution_connection_failure( + url.as_str(), + source, + format!("dns lookup failed: {error}"), + )), + } + } + + async fn complete_job_execution( + &self, + job_id: &str, + signed_event: RadrootsSignedNostrEvent, + delivery_policy: TransportPublishDeliveryPolicy, + timeout_ms: u64, + resolution: PublishRelayResolution, + ) -> Result<TransportPublishJobView, TransportPublishError> { + let target_count = resolution.target_count(); + if resolution.targets.is_empty() { + let status = if resolution + .outcomes + .iter() + .any(|outcome| outcome.outcome_kind.is_retryable()) + { + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + } else if resolution.outcomes.is_empty() { + TransportPublishJobStatus::Rejected + } else { + TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + }; + let last_error = if status == TransportPublishJobStatus::DeliveryUnsatisfiedRetryable { + "delivery_unsatisfied" + } else if resolution.outcomes.is_empty() { + "no_transport_publish_targets" + } else { + "delivery_unsatisfied" + }; + self.store.complete_publish_job( + job_id, + status, + resolution.outcomes, + Some(last_error.to_owned()), + )?; + return self.store.job_by_id(job_id); + } + let required_target_count = delivery_policy.required_target_count(target_count); + if required_target_count > target_count { + self.store.complete_publish_job( + job_id, + TransportPublishJobStatus::Rejected, + resolution.outcomes, + Some("delivery_quorum_exceeds_target_count".to_owned()), + )?; + return self.store.job_by_id(job_id); + } + let source_by_relay = resolution.source_by_relay(); + let target_set = RadrootsRelayTargetSet::from_urls( + resolution + .targets + .iter() + .map(|target| target.url.clone()) + .collect(), + )?; + let satisfaction_policy = satisfaction_policy_from_delivery_policy( + &delivery_policy, + target_count, + resolution.targets.len(), + ); + let publish_request = + RadrootsRelayPublishRequest::new(signed_event, target_set, current_unix_millis()) + .with_satisfaction_policy(satisfaction_policy); + let started = Instant::now(); + let publish_timeout = Duration::from_millis(timeout_ms); + let receipts = + match tokio::time::timeout(publish_timeout, self.publish_with_adapter(publish_request)) + .await + { + Ok(Ok(receipts)) => receipts, + Ok(Err(error)) => transport_error_receipts(&resolution.targets, error), + Err(_) => timeout_receipts(&resolution.targets), + }; + let latency_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX); + let mut outcomes = resolution.outcomes; + outcomes.extend(receipts.into_iter().map(|receipt| { + publish_outcome_from_receipt(receipt, &source_by_relay, Some(latency_ms)) + })); + let status = delivery_status(&delivery_policy, target_count, &outcomes); + let last_error = if status == TransportPublishJobStatus::DeliverySatisfied { + None + } else { + Some("delivery_unsatisfied".to_owned()) + }; + self.store + .complete_publish_job(job_id, status, outcomes, last_error)?; + self.store.job_by_id(job_id) + } + + async fn publish_with_adapter( + &self, + request: RadrootsRelayPublishRequest, + ) -> Result<Vec<RadrootsRelayPublishRelayReceipt>, TransportPublishError> { + if let Some(publisher) = &self.publisher { + return publisher + .publish(request) + .await + .map_err(TransportPublishError::Relay); + } + let adapter = RadrootsNostrClientPublishAdapter::new(RadrootsNostrClient::new_signerless()); + adapter + .publish(request) + .await + .map_err(TransportPublishError::Relay) + } +} + +#[derive(Clone)] +pub struct TransportPublishStore { + inner: Arc<Mutex<Connection>>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum PublishJobVisibility { + Own, + Admin, +} + +impl FromStr for PublishJobVisibility { + type Err = TransportPublishError; + + fn from_str(value: &str) -> Result<Self, Self::Err> { + match value { + "own" => Ok(Self::Own), + "admin" => Ok(Self::Admin), + other => Err(TransportPublishError::InvalidScope(format!( + "unknown job visibility `{other}`" + ))), + } + } +} + +impl fmt::Display for PublishJobVisibility { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Own => f.write_str("own"), + Self::Admin => f.write_str("admin"), + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PublishPrincipalInit { + pub label: String, + pub token_hash: String, + pub allowed_pubkeys: Vec<String>, + pub allowed_kinds: Vec<u32>, + pub allowed_target_policies: Vec<TransportPublishTargetPolicyName>, + pub allowed_nostr_source_policies: Vec<NostrPublishTargetSourcePolicy>, + pub allow_request_targets: bool, + pub job_visibility: PublishJobVisibility, + pub expires_at_unix: Option<i64>, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PublishPrincipal { + pub principal_id: String, + pub label: String, + pub allowed_pubkeys: Vec<String>, + pub allowed_kinds: Vec<u32>, + pub allowed_target_policies: Vec<TransportPublishTargetPolicyName>, + pub allowed_nostr_source_policies: Vec<NostrPublishTargetSourcePolicy>, + pub allow_request_targets: bool, + pub job_visibility: PublishJobVisibility, + pub expires_at_unix: Option<i64>, +} + +impl PublishPrincipal { + pub fn allows_event( + &self, + request: &TransportPublishEventRequest, + ) -> Result<(), TransportPublishError> { + ensure_lower_hex("pubkey", request.event.pubkey.as_str(), 64)?; + if !self + .allowed_pubkeys + .iter() + .any(|pubkey| pubkey == &request.event.pubkey) + { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to publish for event pubkey".to_owned(), + )); + } + if !self.allowed_kinds.contains(&request.event.kind) { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to publish event kind".to_owned(), + )); + } + match &request.target_policy { + TransportPublishTargetPolicy::ExplicitTargets { targets } => { + if !self + .allowed_target_policies + .contains(&TransportPublishTargetPolicyName::ExplicitTargets) + { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to use explicit transport targets".to_owned(), + )); + } + if !self.allow_request_targets && !targets.is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to provide request targets".to_owned(), + )); + } + } + TransportPublishTargetPolicy::Nostr { + source_policy, + relay_urls, + } => { + if !self + .allowed_target_policies + .contains(&TransportPublishTargetPolicyName::Nostr) + { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to use Nostr target policy".to_owned(), + )); + } + if !self.allowed_nostr_source_policies.contains(source_policy) { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to use requested Nostr source policy".to_owned(), + )); + } + if !self.allow_request_targets && !relay_urls.is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal is not allowed to provide request targets".to_owned(), + )); + } + } + } + Ok(()) + } + + fn can_read_job(&self, principal_id: &str) -> bool { + self.job_visibility == PublishJobVisibility::Admin || self.principal_id == principal_id + } +} + +#[derive(Debug, Clone)] +pub struct PublishJobInsert { + pub principal_id: String, + pub idempotency_key: Option<String>, + pub request: TransportPublishEventRequest, + pub request_fingerprint: String, + pub effective_target_count: usize, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ResolvedPublishRelay { + pub url: RadrootsRelayUrl, + pub source: TransportPublishTargetSource, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PublishRelayResolution { + pub targets: Vec<ResolvedPublishRelay>, + pub outcomes: Vec<TransportPublishTargetOutcome>, +} + +impl PublishRelayResolution { + fn target_count(&self) -> usize { + self.targets.len() + self.outcomes.len() + } + + fn source_by_relay(&self) -> BTreeMap<String, TransportPublishTargetSource> { + self.targets + .iter() + .map(|target| (target.url.as_str().to_owned(), target.source)) + .collect() + } +} + +pub(crate) type PublishRelayResolveFuture<'a> = + Pin<Box<dyn Future<Output = Result<Vec<IpAddr>, std::io::Error>> + Send + 'a>>; + +pub(crate) trait PublishRelayResolver: Send + Sync { + fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a>; +} + +type PublishAuthorRelayDiscoveryFuture<'a> = + Pin<Box<dyn Future<Output = Result<Vec<String>, TransportPublishError>> + Send + 'a>>; + +trait PublishAuthorRelayDiscovery: Send + Sync { + fn fetch_author_write_relays<'a>( + &'a self, + pubkey: &'a str, + discovery_targets: Vec<ResolvedPublishRelay>, + connect_timeout_secs: u64, + ) -> PublishAuthorRelayDiscoveryFuture<'a>; +} + +#[derive(Debug)] +struct SystemPublishRelayResolver; + +impl PublishRelayResolver for SystemPublishRelayResolver { + fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a> { + Box::pin(async move { + let (host, port) = relay_socket_target(url)?; + let addrs = tokio::net::lookup_host((host.as_str(), port)).await?; + Ok(addrs.map(|addr| addr.ip()).collect()) + }) + } +} + +#[derive(Debug)] +struct NostrPublishAuthorRelayDiscovery; + +impl PublishAuthorRelayDiscovery for NostrPublishAuthorRelayDiscovery { + fn fetch_author_write_relays<'a>( + &'a self, + pubkey: &'a str, + discovery_targets: Vec<ResolvedPublishRelay>, + connect_timeout_secs: u64, + ) -> PublishAuthorRelayDiscoveryFuture<'a> { + Box::pin(async move { + let Ok(public_key) = RadrootsNostrPublicKey::from_hex(pubkey) else { + return Ok(Vec::new()); + }; + let client = RadrootsNostrClient::new_signerless(); + for target in discovery_targets { + if client.add_read_relay(target.url.as_str()).await.is_err() { + return Ok(Vec::new()); + } + } + let filter = RadrootsNostrFilter::new() + .author(public_key) + .kind(RadrootsNostrKind::Custom(10_002)) + .limit(10); + let timeout = Duration::from_secs(connect_timeout_secs); + let Ok(events) = client.fetch_events(filter, timeout).await else { + return Ok(Vec::new()); + }; + let Some(event) = events.into_iter().max_by(|left, right| { + left.created_at + .as_secs() + .cmp(&right.created_at.as_secs()) + .then_with(|| left.id.to_hex().cmp(&right.id.to_hex())) + }) else { + return Ok(Vec::new()); + }; + Ok(author_write_relays_from_nip65_event(&event)) + }) + } +} + +impl TransportPublishStore { + pub fn open(path: PathBuf) -> Result<Self, TransportPublishError> { + if let Some(parent) = path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + { + std::fs::create_dir_all(parent)?; + } + let connection = Connection::open(path)?; + Self::from_connection(connection) + } + + pub fn memory() -> Result<Self, TransportPublishError> { + Self::from_connection(Connection::open_in_memory()?) + } + + fn from_connection(connection: Connection) -> Result<Self, TransportPublishError> { + connection.execute_batch( + r#" + PRAGMA foreign_keys = ON; + CREATE TABLE IF NOT EXISTS transport_publish_principals ( + principal_id TEXT PRIMARY KEY NOT NULL, + label TEXT NOT NULL, + token_hash TEXT NOT NULL UNIQUE, + allowed_pubkeys_json TEXT NOT NULL, + allowed_kinds_json TEXT NOT NULL, + allowed_target_policies_json TEXT NOT NULL, + allowed_nostr_source_policies_json TEXT NOT NULL, + allow_request_targets INTEGER NOT NULL, + job_visibility TEXT NOT NULL, + expires_at_unix INTEGER, + revoked_at_unix INTEGER, + created_at_unix INTEGER NOT NULL + ); + CREATE TABLE IF NOT EXISTS transport_publish_jobs ( + job_id TEXT PRIMARY KEY NOT NULL, + principal_id TEXT NOT NULL, + idempotency_key TEXT, + request_fingerprint TEXT NOT NULL, + status TEXT NOT NULL, + event_id TEXT NOT NULL, + event_pubkey TEXT NOT NULL, + event_kind INTEGER NOT NULL, + target_policy_json TEXT NOT NULL, + delivery_policy_json TEXT NOT NULL, + requested_target_count INTEGER NOT NULL, + effective_target_count INTEGER NOT NULL, + request_json TEXT NOT NULL, + requested_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL, + completed_at_ms INTEGER, + last_error TEXT, + FOREIGN KEY(principal_id) REFERENCES transport_publish_principals(principal_id) + ); + CREATE UNIQUE INDEX IF NOT EXISTS transport_publish_jobs_principal_idempotency_idx + ON transport_publish_jobs(principal_id, idempotency_key) + WHERE idempotency_key IS NOT NULL; + CREATE TABLE IF NOT EXISTS transport_publish_target_results ( + job_id TEXT NOT NULL, + transport_kind TEXT NOT NULL, + endpoint_uri TEXT NOT NULL, + source TEXT NOT NULL, + attempted INTEGER NOT NULL, + outcome_kind TEXT NOT NULL, + message TEXT, + latency_ms INTEGER, + updated_at_ms INTEGER NOT NULL, + PRIMARY KEY(job_id, transport_kind, endpoint_uri), + FOREIGN KEY(job_id) REFERENCES transport_publish_jobs(job_id) + ); + CREATE TABLE IF NOT EXISTS transport_publish_nostr_author_cache ( + pubkey TEXT PRIMARY KEY NOT NULL, + relays_json TEXT NOT NULL, + updated_at_ms INTEGER NOT NULL + ); + "#, + )?; + recover_interrupted_publish_jobs(&connection)?; + connection.pragma_update(None, "user_version", SCHEMA_VERSION)?; + Ok(Self { + inner: Arc::new(Mutex::new(connection)), + }) + } + + pub fn create_principal( + &self, + input: PublishPrincipalInit, + ) -> Result<PublishPrincipal, TransportPublishError> { + validate_principal_init(&input)?; + let principal_id = Uuid::new_v4().to_string(); + let now = current_unix_secs(); + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + connection.execute( + r#" + INSERT INTO transport_publish_principals ( + principal_id, + label, + token_hash, + allowed_pubkeys_json, + allowed_kinds_json, + allowed_target_policies_json, + allowed_nostr_source_policies_json, + allow_request_targets, + job_visibility, + expires_at_unix, + revoked_at_unix, + created_at_unix + ) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL, ?11) + "#, + params![ + principal_id, + input.label.trim(), + input.token_hash, + serde_json::to_string(&input.allowed_pubkeys)?, + serde_json::to_string(&input.allowed_kinds)?, + serde_json::to_string(&input.allowed_target_policies)?, + serde_json::to_string(&input.allowed_nostr_source_policies)?, + input.allow_request_targets, + input.job_visibility.to_string(), + input.expires_at_unix, + now, + ], + )?; + drop(connection); + self.principal_by_id(principal_id.as_str())?.ok_or_else(|| { + TransportPublishError::InvalidScope("created principal missing".to_owned()) + }) + } + + pub fn principal_for_token_hash( + &self, + token_hash: &str, + ) -> Result<Option<PublishPrincipal>, TransportPublishError> { + let now = current_unix_secs(); + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let principal = connection + .query_row( + r#" + SELECT + principal_id, + label, + allowed_pubkeys_json, + allowed_kinds_json, + allowed_target_policies_json, + allowed_nostr_source_policies_json, + allow_request_targets, + job_visibility, + expires_at_unix + FROM transport_publish_principals + WHERE token_hash = ?1 + AND revoked_at_unix IS NULL + AND (expires_at_unix IS NULL OR expires_at_unix > ?2) + "#, + params![token_hash, now], + principal_from_row, + ) + .optional()?; + Ok(principal) + } + + pub fn principal_by_id( + &self, + principal_id: &str, + ) -> Result<Option<PublishPrincipal>, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let principal = connection + .query_row( + r#" + SELECT + principal_id, + label, + allowed_pubkeys_json, + allowed_kinds_json, + allowed_target_policies_json, + allowed_nostr_source_policies_json, + allow_request_targets, + job_visibility, + expires_at_unix + FROM transport_publish_principals + WHERE principal_id = ?1 + "#, + params![principal_id], + principal_from_row, + ) + .optional()?; + Ok(principal) + } + + pub fn record_publish_job( + &self, + insert: PublishJobInsert, + ) -> Result<TransportPublishEventResponse, TransportPublishError> { + if let Some(idempotency_key) = insert.idempotency_key.as_deref() { + if let Some(existing) = + self.job_for_principal_id_and_key(insert.principal_id.as_str(), idempotency_key)? + { + if existing.request_fingerprint != insert.request_fingerprint { + return Err(TransportPublishError::IdempotencyConflict( + idempotency_key.to_owned(), + )); + } + return Ok(TransportPublishEventResponse { + deduplicated: true, + job: existing.view, + }); + } + } + + let job_id = Uuid::new_v4().to_string(); + let now = current_unix_millis(); + let request_json = serde_json::to_string(&insert.request)?; + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let insert_result = connection.execute( + r#" + INSERT INTO transport_publish_jobs ( + job_id, + principal_id, + idempotency_key, + request_fingerprint, + status, + event_id, + event_pubkey, + event_kind, + target_policy_json, + delivery_policy_json, + requested_target_count, + effective_target_count, + request_json, + requested_at_ms, + updated_at_ms + ) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15) + "#, + params![ + job_id, + insert.principal_id, + insert.idempotency_key, + insert.request_fingerprint, + serde_json::to_string(&TransportPublishJobStatus::Publishing)?, + insert.request.event.id, + insert.request.event.pubkey, + insert.request.event.kind, + serde_json::to_string(&insert.request.target_policy)?, + serde_json::to_string(&insert.request.delivery_policy)?, + insert.request.target_policy.request_target_count(), + insert.effective_target_count, + request_json, + now, + now, + ], + ); + match insert_result { + Ok(_) => {} + Err(rusqlite::Error::SqliteFailure(error, _)) + if error.code == rusqlite::ErrorCode::ConstraintViolation => + { + return Err(TransportPublishError::IdempotencyConflict( + "idempotency key conflicts with an existing publish job".to_owned(), + )); + } + Err(error) => return Err(error.into()), + } + drop(connection); + let job = self.job_by_id(job_id.as_str())?; + Ok(TransportPublishEventResponse { + deduplicated: false, + job, + }) + } + + pub fn job_by_id_for_principal( + &self, + job_id: &str, + principal: &PublishPrincipal, + ) -> Result<Option<TransportPublishJobView>, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let sql = job_select_sql("WHERE job_id = ?1"); + let row = connection + .query_row(sql.as_str(), params![job_id], job_from_row) + .optional()?; + drop(connection); + let Some(mut job) = row else { + return Ok(None); + }; + if !principal.can_read_job(job.principal_id.as_str()) { + return Ok(None); + } + job.view.targets = self.target_outcomes(job.view.job_id.as_str())?; + finalize_job_view(&mut job.view); + Ok(Some(job.view)) + } + + pub fn list_jobs_for_principal( + &self, + principal: &PublishPrincipal, + limit: usize, + ) -> Result<Vec<TransportPublishJobView>, TransportPublishError> { + let limit = i64::try_from(limit.clamp(1, 200)).unwrap_or(200); + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let sql = if principal.job_visibility == PublishJobVisibility::Admin { + job_select_sql("ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?1") + } else { + job_select_sql( + "WHERE principal_id = ?1 ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?2", + ) + }; + let mut stmt = connection.prepare(sql.as_str())?; + let rows = if principal.job_visibility == PublishJobVisibility::Admin { + stmt.query_map(params![limit], job_from_row)? + .collect::<Result<Vec<_>, _>>()? + } else { + stmt.query_map(params![principal.principal_id, limit], job_from_row)? + .collect::<Result<Vec<_>, _>>()? + }; + drop(stmt); + drop(connection); + + rows.into_iter() + .map(|mut row| { + row.view.targets = self.target_outcomes(row.view.job_id.as_str())?; + finalize_job_view(&mut row.view); + Ok(row.view) + }) + .collect() + } + + fn job_for_principal_id_and_key( + &self, + principal_id: &str, + idempotency_key: &str, + ) -> Result<Option<PublishJobRow>, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let sql = job_select_sql("WHERE principal_id = ?1 AND idempotency_key = ?2"); + let row = connection + .query_row( + sql.as_str(), + params![principal_id, idempotency_key], + job_from_row, + ) + .optional()?; + drop(connection); + let Some(mut job) = row else { + return Ok(None); + }; + job.view.targets = self.target_outcomes(job.view.job_id.as_str())?; + finalize_job_view(&mut job.view); + Ok(Some(job)) + } + + pub fn job_by_id( + &self, + job_id: &str, + ) -> Result<TransportPublishJobView, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let sql = job_select_sql("WHERE job_id = ?1"); + let row = connection + .query_row(sql.as_str(), params![job_id], job_from_row) + .optional()?; + drop(connection); + let Some(mut job) = row else { + return Err(TransportPublishError::InvalidScope( + "unknown publish job".to_owned(), + )); + }; + job.view.targets = self.target_outcomes(job.view.job_id.as_str())?; + finalize_job_view(&mut job.view); + Ok(job.view) + } + + pub fn complete_publish_job( + &self, + job_id: &str, + status: TransportPublishJobStatus, + outcomes: Vec<TransportPublishTargetOutcome>, + last_error: Option<String>, + ) -> Result<(), TransportPublishError> { + let now = current_unix_millis(); + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + connection.execute( + r#" + UPDATE transport_publish_jobs + SET status = ?2, + updated_at_ms = ?3, + completed_at_ms = ?4, + last_error = ?5 + WHERE job_id = ?1 + "#, + params![ + job_id, + serde_json::to_string(&status)?, + now, + now, + last_error, + ], + )?; + connection.execute( + "DELETE FROM transport_publish_target_results WHERE job_id = ?1", + params![job_id], + )?; + for outcome in outcomes { + connection.execute( + r#" + INSERT OR REPLACE INTO transport_publish_target_results ( + job_id, + transport_kind, + endpoint_uri, + source, + attempted, + outcome_kind, + message, + latency_ms, + updated_at_ms + ) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) + "#, + params![ + job_id, + outcome.transport_kind, + outcome.endpoint_uri, + serde_json::to_string(&outcome.source)?, + outcome.attempted, + serde_json::to_string(&outcome.outcome_kind)?, + outcome.message, + outcome + .latency_ms + .and_then(|value| i64::try_from(value).ok()), + now, + ], + )?; + } + Ok(()) + } + + pub fn cached_author_write_relays( + &self, + pubkey: &str, + ) -> Result<Vec<String>, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let relays_json = connection + .query_row( + "SELECT relays_json FROM transport_publish_nostr_author_cache WHERE pubkey = ?1", + params![pubkey], + |row| row.get::<_, String>(0), + ) + .optional()?; + relays_json + .map(|value| serde_json::from_str(value.as_str()).map_err(TransportPublishError::from)) + .unwrap_or_else(|| Ok(Vec::new())) + } + + pub fn cache_author_write_relays( + &self, + pubkey: &str, + relays: &[String], + ) -> Result<(), TransportPublishError> { + let now = current_unix_millis(); + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + connection.execute( + r#" + INSERT INTO transport_publish_nostr_author_cache (pubkey, relays_json, updated_at_ms) + VALUES (?1, ?2, ?3) + ON CONFLICT(pubkey) DO UPDATE SET + relays_json = excluded.relays_json, + updated_at_ms = excluded.updated_at_ms + "#, + params![pubkey, serde_json::to_string(relays)?, now], + )?; + Ok(()) + } + + fn target_outcomes( + &self, + job_id: &str, + ) -> Result<Vec<TransportPublishTargetOutcome>, TransportPublishError> { + let connection = self + .inner + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let mut stmt = connection.prepare( + r#" + SELECT transport_kind, endpoint_uri, source, attempted, outcome_kind, message, latency_ms + FROM transport_publish_target_results + WHERE job_id = ?1 + ORDER BY transport_kind, endpoint_uri + "#, + )?; + let outcomes = stmt + .query_map(params![job_id], target_outcome_from_row)? + .collect::<Result<Vec<_>, _>>()?; + Ok(outcomes) + } +} + +struct PublishJobRow { + principal_id: String, + request_fingerprint: String, + view: TransportPublishJobView, +} + +fn recover_interrupted_publish_jobs(connection: &Connection) -> Result<(), TransportPublishError> { + let now = current_unix_millis(); + connection.execute( + r#" + UPDATE transport_publish_jobs + SET status = ?1, + updated_at_ms = ?2, + completed_at_ms = ?3, + last_error = ?4 + WHERE status = ?5 + "#, + params![ + serde_json::to_string(&TransportPublishJobStatus::DeliveryUnsatisfiedRetryable)?, + now, + now, + "publish_attempt_interrupted", + serde_json::to_string(&TransportPublishJobStatus::Publishing)?, + ], + )?; + Ok(()) +} + +fn job_select_sql(tail: &str) -> String { + format!( + r#" + SELECT + job_id, + principal_id, + request_fingerprint, + status, + event_id, + event_pubkey, + event_kind, + target_policy_json, + delivery_policy_json, + effective_target_count, + requested_at_ms, + completed_at_ms, + last_error + FROM transport_publish_jobs + {tail} + "# + ) +} + +fn principal_from_row(row: &Row<'_>) -> Result<PublishPrincipal, rusqlite::Error> { + let visibility: String = row.get(7)?; + Ok(PublishPrincipal { + principal_id: row.get(0)?, + label: row.get(1)?, + allowed_pubkeys: json_column(row, 2)?, + allowed_kinds: json_column(row, 3)?, + allowed_target_policies: json_column(row, 4)?, + allowed_nostr_source_policies: json_column(row, 5)?, + allow_request_targets: row.get(6)?, + job_visibility: PublishJobVisibility::from_str(visibility.as_str()) + .map_err(|error| conversion_error(7, error))?, + expires_at_unix: row.get(8)?, + }) +} + +fn job_from_row(row: &Row<'_>) -> Result<PublishJobRow, rusqlite::Error> { + let status: TransportPublishJobStatus = json_text(row, 3)?; + let target_policy: TransportPublishTargetPolicy = json_text(row, 7)?; + let delivery_policy: TransportPublishDeliveryPolicy = json_text(row, 8)?; + let target_count: i64 = row.get(9)?; + Ok(PublishJobRow { + principal_id: row.get(1)?, + request_fingerprint: row.get(2)?, + view: TransportPublishJobView { + job_id: row.get(0)?, + status, + terminal: false, + delivery_satisfied: false, + event_id: row.get(4)?, + pubkey: row.get(5)?, + event_kind: row.get::<_, i64>(6)? as u32, + target_policy, + delivery_policy, + target_count: usize::try_from(target_count).unwrap_or(0), + acknowledged_count: 0, + retryable_count: 0, + terminal_count: 0, + requested_at_ms: row.get(10)?, + completed_at_ms: row.get(11)?, + last_error: row.get(12)?, + targets: Vec::new(), + }, + }) +} + +fn target_outcome_from_row( + row: &Row<'_>, +) -> Result<TransportPublishTargetOutcome, rusqlite::Error> { + let source: TransportPublishTargetSource = json_text(row, 2)?; + let outcome_kind: TransportPublishOutcomeKind = json_text(row, 4)?; + Ok(TransportPublishTargetOutcome { + transport_kind: row.get(0)?, + endpoint_uri: row.get(1)?, + source, + attempted: row.get(3)?, + outcome_kind, + message: row.get(5)?, + latency_ms: row + .get::<_, Option<i64>>(6)? + .map(|latency| u64::try_from(latency).unwrap_or(0)), + }) +} + +fn finalize_job_view(view: &mut TransportPublishJobView) { + view.acknowledged_count = view + .targets + .iter() + .filter(|relay| relay.outcome_kind.counts_toward_satisfaction()) + .count(); + view.retryable_count = view + .targets + .iter() + .filter(|relay| relay.outcome_kind.is_retryable()) + .count(); + view.terminal_count = view + .targets + .iter() + .filter(|relay| relay.outcome_kind.is_terminal_failure()) + .count(); + view.terminal = matches!( + view.status, + TransportPublishJobStatus::DeliverySatisfied + | TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + | TransportPublishJobStatus::Rejected + ); + view.delivery_satisfied = view.status == TransportPublishJobStatus::DeliverySatisfied; +} + +fn validate_principal_init(input: &PublishPrincipalInit) -> Result<(), TransportPublishError> { + if input.label.trim().is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal label must not be empty".to_owned(), + )); + } + if !input.token_hash.starts_with(TOKEN_HASH_PREFIX) { + return Err(TransportPublishError::InvalidScope( + "principal token hash must use sha256 prefix".to_owned(), + )); + } + if input.allowed_pubkeys.is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal must include at least one allowed pubkey".to_owned(), + )); + } + for pubkey in &input.allowed_pubkeys { + ensure_lower_hex("allowed_pubkey", pubkey, 64)?; + } + if input.allowed_kinds.is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal must include at least one allowed kind".to_owned(), + )); + } + if input + .allowed_kinds + .iter() + .any(|kind| *kind > u16::MAX as u32) + { + return Err(TransportPublishError::InvalidScope( + "allowed kind exceeds transport publish range".to_owned(), + )); + } + if input.allowed_target_policies.is_empty() { + return Err(TransportPublishError::InvalidScope( + "principal must include at least one allowed target policy".to_owned(), + )); + } + if input + .allowed_target_policies + .contains(&TransportPublishTargetPolicyName::Nostr) + && input.allowed_nostr_source_policies.is_empty() + { + return Err(TransportPublishError::InvalidScope( + "principal must include at least one allowed Nostr source policy".to_owned(), + )); + } + Ok(()) +} + +pub fn generate_bearer_token() -> String { + let bytes: [u8; 32] = rand::random(); + format!("{TOKEN_PREFIX}{}", hex_lower(&bytes)) +} + +pub fn hash_bearer_token(token: &str) -> String { + let mut hasher = Sha256::new(); + hasher.update(token.as_bytes()); + format!("{TOKEN_HASH_PREFIX}{}", hex_lower(&hasher.finalize())) +} + +fn hex_lower(bytes: &[u8]) -> String { + let mut output = String::with_capacity(bytes.len() * 2); + for byte in bytes { + use std::fmt::Write; + let _ = write!(&mut output, "{byte:02x}"); + } + output +} + +pub fn parse_nostr_source_policy( + value: &str, +) -> Result<NostrPublishTargetSourcePolicy, TransportPublishError> { + match value { + "explicit_only" => Ok(NostrPublishTargetSourcePolicy::ExplicitOnly), + "request_then_author_write_then_daemon_default" => { + Ok(NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault) + } + "author_write_then_daemon_default" => { + Ok(NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault) + } + "daemon_default_only" => Ok(NostrPublishTargetSourcePolicy::DaemonDefaultOnly), + other => Err(TransportPublishError::InvalidScope(format!( + "unknown Nostr source policy `{other}`" + ))), + } +} + +pub fn parse_target_policy( + value: &str, +) -> Result<TransportPublishTargetPolicyName, TransportPublishError> { + match value { + "explicit_targets" => Ok(TransportPublishTargetPolicyName::ExplicitTargets), + "nostr" => Ok(TransportPublishTargetPolicyName::Nostr), + other => Err(TransportPublishError::InvalidScope(format!( + "unknown target policy `{other}`" + ))), + } +} + +fn signed_event_from_wire( + event: &SignedNostrEventWire, +) -> Result<RadrootsSignedNostrEvent, TransportPublishError> { + event + .validate() + .map_err(|error| TransportPublishError::InvalidSignedEvent(error.to_string()))?; + let created_at = u32::try_from(event.created_at).map_err(|_| { + TransportPublishError::InvalidSignedEvent( + "signed event created_at exceeds daemon-supported range".to_owned(), + ) + })?; + let raw_json = serde_json::to_string(event)?; + let radroots_event = RadrootsNostrEvent { + id: event.id.clone(), + author: event.pubkey.clone(), + created_at, + kind: event.kind, + tags: event.tags.clone(), + content: event.content.clone(), + sig: event.sig.clone(), + }; + match radroots_nostr_verify_event(&radroots_event) { + RadrootsNostrEventVerification::Verified => {} + verification => return Err(TransportPublishError::SignedEventVerification(verification)), + } + RadrootsSignedNostrEvent::new(RadrootsSignedNostrEventParts { + id: event.id.clone(), + pubkey: event.pubkey.clone(), + created_at, + kind: event.kind, + tags: event.tags.clone(), + content: event.content.clone(), + sig: event.sig.clone(), + raw_json, + }) + .map_err(TransportPublishError::from) +} + +fn request_intent_fingerprint( + principal_id: &str, + canonical_event_json: &str, + request: &TransportPublishEventRequest, + effective_timeout_ms: u64, +) -> Result<String, TransportPublishError> { + #[derive(Serialize)] + struct FingerprintInput<'a> { + principal_id: &'a str, + canonical_event_json: &'a str, + target_policy: &'a TransportPublishTargetPolicy, + delivery_policy: &'a TransportPublishDeliveryPolicy, + effective_timeout_ms: u64, + } + + let input = FingerprintInput { + principal_id, + canonical_event_json, + target_policy: &request.target_policy, + delivery_policy: &request.delivery_policy, + effective_timeout_ms, + }; + let bytes = serde_json::to_vec(&input)?; + let mut hasher = Sha256::new(); + hasher.update(bytes); + Ok(hex_lower(&hasher.finalize())) +} + +fn effective_publish_timeout_ms( + config: &TransportPublishConfig, + timeout_ms: Option<u64>, +) -> Result<u64, TransportPublishError> { + let max_timeout_ms = config.connect_timeout_secs.saturating_mul(1_000); + match timeout_ms { + Some(0) => Err(TransportPublishError::InvalidSignedEvent( + "timeout_ms must be greater than zero".to_owned(), + )), + Some(timeout_ms) if timeout_ms > max_timeout_ms => { + Err(TransportPublishError::InvalidSignedEvent(format!( + "timeout_ms must be at most {max_timeout_ms}" + ))) + } + Some(timeout_ms) => Ok(timeout_ms), + None => Ok(max_timeout_ms), + } +} + +fn push_resolved_relay( + targets: &mut Vec<ResolvedPublishRelay>, + url: RadrootsRelayUrl, + source: TransportPublishTargetSource, +) { + if !targets.iter().any(|target| target.url == url) { + targets.push(ResolvedPublishRelay { url, source }); + } +} + +fn reticulum_preview_outcome(target: &TransportPublishTarget) -> TransportPublishTargetOutcome { + let outcome_kind = match target.preview_behavior.unwrap_or_default() { + TransportPublishPreviewBehavior::RejectDeliveryAttempts => { + TransportPublishOutcomeKind::Unavailable + } + TransportPublishPreviewBehavior::DeferDeliveryPlans => { + TransportPublishOutcomeKind::Deferred + } + }; + TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_RETICULUM.to_owned(), + endpoint_uri: target.endpoint_uri.trim().to_owned(), + source: TransportPublishTargetSource::ReticulumPreview, + attempted: false, + outcome_kind, + message: Some( + "reticulum transport is registered for preview but not routable by radrootsd" + .to_owned(), + ), + latency_ms: None, + } +} + +fn unsupported_transport_outcome(target: &TransportPublishTarget) -> TransportPublishTargetOutcome { + TransportPublishTargetOutcome { + transport_kind: target.transport_kind.trim().to_owned(), + endpoint_uri: target.endpoint_uri.trim().to_owned(), + source: TransportPublishTargetSource::Request, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::Unsupported, + message: Some("transport kind is not supported by radrootsd transport publish".to_owned()), + latency_ms: None, + } +} + +fn relay_resolution_connection_failure( + relay_url: impl Into<String>, + source: TransportPublishTargetSource, + message: impl Into<String>, +) -> TransportPublishTargetOutcome { + TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: relay_url.into(), + source, + attempted: false, + outcome_kind: TransportPublishOutcomeKind::ConnectionFailed, + message: Some(message.into()), + latency_ms: None, + } +} + +fn relay_socket_target(url: &RadrootsRelayUrl) -> Result<(String, u16), std::io::Error> { + let parsed = url::Url::parse(url.as_str()) + .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?; + let host = parsed + .host_str() + .filter(|host| !host.is_empty()) + .ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "relay URL must include a DNS host", + ) + })? + .to_owned(); + let port = parsed.port_or_known_default().ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "relay URL scheme must have a default port", + ) + })?; + Ok((host, port)) +} + +fn relay_url_policy(config: &TransportPublishConfig) -> RadrootsRelayUrlPolicy { + match config.nostr.relay_url_policy { + crate::app::config::NostrRelayUrlPolicy::Public => RadrootsRelayUrlPolicy::Public, + crate::app::config::NostrRelayUrlPolicy::Localhost => RadrootsRelayUrlPolicy::Localhost, + } +} + +fn author_write_relays_from_nip65_event( + event: &radroots_nostr::prelude::RadrootsNostrEvent, +) -> Vec<String> { + event + .tags + .iter() + .filter_map(|tag| { + let values = tag.as_slice(); + if values.first().map(String::as_str) != Some("r") { + return None; + } + let relay = values.get(1)?.trim(); + if relay.is_empty() { + return None; + } + if values.get(2).map(String::as_str) == Some("read") { + return None; + } + Some(relay.to_owned()) + }) + .collect() +} + +fn publish_outcome_from_receipt( + receipt: RadrootsRelayPublishRelayReceipt, + source_by_relay: &BTreeMap<String, TransportPublishTargetSource>, + latency_ms: Option<u64>, +) -> TransportPublishTargetOutcome { + let source = source_by_relay + .get(receipt.relay_url.as_str()) + .copied() + .unwrap_or(TransportPublishTargetSource::DaemonDefault); + TransportPublishTargetOutcome { + transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), + endpoint_uri: receipt.relay_url, + source, + attempted: receipt.attempted, + outcome_kind: publish_outcome_kind(receipt.outcome.kind), + message: receipt.outcome.message, + latency_ms, + } +} + +fn publish_outcome_kind(kind: RadrootsRelayOutcomeKind) -> TransportPublishOutcomeKind { + match kind { + RadrootsRelayOutcomeKind::Accepted => TransportPublishOutcomeKind::Accepted, + RadrootsRelayOutcomeKind::DuplicateAccepted => { + TransportPublishOutcomeKind::DuplicateAccepted + } + RadrootsRelayOutcomeKind::Blocked => TransportPublishOutcomeKind::Blocked, + RadrootsRelayOutcomeKind::RateLimited => TransportPublishOutcomeKind::RateLimited, + RadrootsRelayOutcomeKind::Invalid => TransportPublishOutcomeKind::Invalid, + RadrootsRelayOutcomeKind::PowRequired => TransportPublishOutcomeKind::PowRequired, + RadrootsRelayOutcomeKind::Restricted => TransportPublishOutcomeKind::Restricted, + RadrootsRelayOutcomeKind::AuthRequired => TransportPublishOutcomeKind::AuthRequired, + RadrootsRelayOutcomeKind::Muted => TransportPublishOutcomeKind::Muted, + RadrootsRelayOutcomeKind::Unsupported => TransportPublishOutcomeKind::Unsupported, + RadrootsRelayOutcomeKind::PaymentRequired => TransportPublishOutcomeKind::PaymentRequired, + RadrootsRelayOutcomeKind::Error => TransportPublishOutcomeKind::Error, + RadrootsRelayOutcomeKind::Timeout => TransportPublishOutcomeKind::Timeout, + RadrootsRelayOutcomeKind::ConnectionFailed => TransportPublishOutcomeKind::ConnectionFailed, + RadrootsRelayOutcomeKind::RelayUrlRejected => TransportPublishOutcomeKind::TargetRejected, + RadrootsRelayOutcomeKind::SkippedAlreadyAccepted => { + TransportPublishOutcomeKind::SkippedAlreadyAccepted + } + RadrootsRelayOutcomeKind::Unknown => TransportPublishOutcomeKind::Unknown, + } +} + +fn satisfaction_policy_from_delivery_policy( + delivery_policy: &TransportPublishDeliveryPolicy, + target_count: usize, + nostr_target_count: usize, +) -> RadrootsTransportSatisfactionPolicy { + match delivery_policy { + TransportPublishDeliveryPolicy::Any => RadrootsTransportSatisfactionPolicy::AnyTarget, + TransportPublishDeliveryPolicy::All => RadrootsTransportSatisfactionPolicy::AllTargets, + TransportPublishDeliveryPolicy::Quorum { quorum } => { + let required = (*quorum).min(target_count).min(nostr_target_count).max(1); + RadrootsTransportSatisfactionPolicy::AtLeast( + u16::try_from(required).unwrap_or(u16::MAX), + ) + } + } +} + +fn delivery_status( + delivery_policy: &TransportPublishDeliveryPolicy, + target_count: usize, + outcomes: &[TransportPublishTargetOutcome], +) -> TransportPublishJobStatus { + let required = delivery_policy.required_target_count(target_count); + let acknowledged = outcomes + .iter() + .filter(|outcome| outcome.outcome_kind.counts_toward_satisfaction()) + .count(); + if acknowledged >= required { + return TransportPublishJobStatus::DeliverySatisfied; + } + if outcomes + .iter() + .any(|outcome| outcome.outcome_kind.is_retryable()) + { + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + } else { + TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + } +} + +fn timeout_receipts(targets: &[ResolvedPublishRelay]) -> Vec<RadrootsRelayPublishRelayReceipt> { + targets + .iter() + .map(|target| { + RadrootsRelayPublishRelayReceipt::attempted( + target.url.as_str(), + RadrootsRelayOutcome::timeout("timeout: publish attempt exceeded daemon bound"), + ) + }) + .collect() +} + +fn transport_error_receipts( + targets: &[ResolvedPublishRelay], + error: TransportPublishError, +) -> Vec<RadrootsRelayPublishRelayReceipt> { + let message = format!("error: {error}"); + targets + .iter() + .map(|target| { + RadrootsRelayPublishRelayReceipt::attempted( + target.url.as_str(), + RadrootsRelayOutcome::connection_failed(message.clone()), + ) + }) + .collect() +} + +pub fn write_token_file(path: &Path, token: &str) -> Result<(), TransportPublishError> { + if let Some(parent) = path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + { + std::fs::create_dir_all(parent)?; + } + let mut options = std::fs::OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + use std::io::Write; + let mut file = options.open(path)?; + file.write_all(token.as_bytes())?; + file.write_all(b"\n")?; + Ok(()) +} + +fn ensure_lower_hex( + field: &str, + value: &str, + expected_len: usize, +) -> Result<(), TransportPublishError> { + if value.len() == expected_len + && value + .bytes() + .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f')) + { + Ok(()) + } else { + Err(TransportPublishError::InvalidScope(format!( + "{field} must be {expected_len} lowercase hex characters" + ))) + } +} + +fn json_column<T: for<'de> Deserialize<'de>>( + row: &Row<'_>, + index: usize, +) -> Result<T, rusqlite::Error> { + let value: String = row.get(index)?; + serde_json::from_str(value.as_str()).map_err(|error| conversion_error(index, error)) +} + +fn json_text<T: for<'de> Deserialize<'de>>( + row: &Row<'_>, + index: usize, +) -> Result<T, rusqlite::Error> { + let value: String = row.get(index)?; + serde_json::from_str(value.as_str()).map_err(|error| conversion_error(index, error)) +} + +fn conversion_error<E>(index: usize, error: E) -> rusqlite::Error +where + E: std::error::Error + Send + Sync + 'static, +{ + rusqlite::Error::FromSqlConversionFailure(index, Type::Text, Box::new(error)) +} + +fn current_unix_secs() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_secs() as i64) + .unwrap_or_default() +} + +fn current_unix_millis() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_millis() as i64) + .unwrap_or_default() +} + +#[cfg(test)] +mod tests { + use super::{ + PublishJobInsert, PublishJobVisibility, PublishPrincipal, PublishPrincipalInit, + TransportPublish, TransportPublishError, TransportPublishStore, generate_bearer_token, + hash_bearer_token, parse_nostr_source_policy, + }; + use crate::app::config::{ + NostrRelayUrlPolicy, TransportPublishConfig, TransportPublishNostrConfig, + }; + use nostr::JsonUtil; + use radroots_identity::RadrootsIdentity; + use radroots_nostr::prelude::{ + RadrootsNostrEventVerification, RadrootsNostrTimestamp, radroots_nostr_build_event, + }; + use radroots_relay_transport::{RadrootsMockRelayPublishAdapter, RadrootsRelayOutcome}; + use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, SignedNostrEventWire, TransportPublishDeliveryPolicy, + TransportPublishEventRequest, TransportPublishJobStatus, TransportPublishOutcomeKind, + TransportPublishTargetPolicy, TransportPublishTargetPolicyName, + TransportPublishTargetSource, + }; + use std::collections::BTreeMap; + use std::net::{IpAddr, Ipv4Addr}; + use std::sync::Arc; + + const RELAY_PRIMARY: &str = "wss://relay.example.com"; + const RELAY_SECONDARY: &str = "wss://relay-2.example.com"; + const RELAY_FORBIDDEN: &str = "wss://forbidden-relay.example.com"; + + fn event(pubkey: &str, kind: u32) -> SignedNostrEventWire { + SignedNostrEventWire { + id: "0".repeat(64), + pubkey: pubkey.to_owned(), + created_at: 1_700_000_000, + kind, + tags: vec![vec!["d".to_owned(), "listing-1".to_owned()]], + content: "{}".to_owned(), + sig: "1".repeat(128), + } + } + + fn request(pubkey: &str, kind: u32) -> TransportPublishEventRequest { + TransportPublishEventRequest { + event: event(pubkey, kind), + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + Vec::new(), + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, + idempotency_key: Some("idem-1".to_owned()), + timeout_ms: None, + } + } + + fn signed_event(identity: &RadrootsIdentity, content: &str) -> SignedNostrEventWire { + let event = radroots_nostr_build_event( + 30_402, + content, + vec![vec!["d".to_owned(), "listing-1".to_owned()]], + ) + .expect("event builder") + .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) + .sign_with_keys(identity.keys()) + .expect("signed event"); + serde_json::from_str(event.as_json().as_str()).expect("event wire") + } + + fn publish_request( + event: SignedNostrEventWire, + relays: Vec<String>, + source_policy: NostrPublishTargetSourcePolicy, + delivery_policy: TransportPublishDeliveryPolicy, + idempotency_key: Option<&str>, + ) -> TransportPublishEventRequest { + TransportPublishEventRequest { + event, + target_policy: TransportPublishTargetPolicy::nostr(source_policy, relays), + delivery_policy, + idempotency_key: idempotency_key.map(str::to_owned), + timeout_ms: Some(5_000), + } + } + + fn transport_publish( + config: TransportPublishConfig, + ) -> (TransportPublish, RadrootsMockRelayPublishAdapter) { + transport_publish_with_resolver(config, Arc::new(StaticPublishRelayResolver::new())) + } + + fn transport_publish_with_resolver( + config: TransportPublishConfig, + resolver: Arc<dyn super::PublishRelayResolver>, + ) -> (TransportPublish, RadrootsMockRelayPublishAdapter) { + let adapter = RadrootsMockRelayPublishAdapter::new(); + let proxy = TransportPublish::memory(config) + .expect("proxy") + .with_relay_resolver(resolver) + .with_publisher(Arc::new(adapter.clone())); + (proxy, adapter) + } + + fn principal( + proxy: &TransportPublish, + pubkey: String, + nostr_source_policies: Vec<NostrPublishTargetSourcePolicy>, + allow_request_targets: bool, + visibility: PublishJobVisibility, + ) -> PublishPrincipal { + proxy + .store + .create_principal(PublishPrincipalInit { + label: "tester".to_owned(), + token_hash: hash_bearer_token(generate_bearer_token().as_str()), + allowed_pubkeys: vec![pubkey], + allowed_kinds: vec![30_402], + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: nostr_source_policies, + allow_request_targets, + job_visibility: visibility, + expires_at_unix: None, + }) + .expect("principal") + } + + fn config_with_defaults(relays: Vec<&str>) -> TransportPublishConfig { + TransportPublishConfig { + nostr: TransportPublishNostrConfig { + daemon_default_relays: relays.into_iter().map(str::to_owned).collect(), + ..TransportPublishNostrConfig::default() + }, + ..TransportPublishConfig::default() + } + } + + #[derive(Default)] + struct StaticPublishRelayResolver { + results: BTreeMap<String, Result<Vec<IpAddr>, String>>, + } + + impl StaticPublishRelayResolver { + fn new() -> Self { + Self::default() + } + + fn with_addresses(mut self, url: &str, addresses: Vec<IpAddr>) -> Self { + self.results.insert(url.to_owned(), Ok(addresses)); + self + } + + fn with_failure(mut self, url: &str, error: &str) -> Self { + self.results.insert(url.to_owned(), Err(error.to_owned())); + self + } + } + + impl super::PublishRelayResolver for StaticPublishRelayResolver { + fn resolve<'a>( + &'a self, + url: &'a radroots_relay_transport::RadrootsRelayUrl, + ) -> super::PublishRelayResolveFuture<'a> { + Box::pin(async move { + match self.results.get(url.as_str()) { + Some(Ok(addresses)) => Ok(addresses.clone()), + Some(Err(error)) => Err(std::io::Error::other(error.clone())), + None => Ok(vec![IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))]), + } + }) + } + } + + struct StaticPublishAuthorRelayDiscovery { + relays: Vec<String>, + } + + impl StaticPublishAuthorRelayDiscovery { + fn new(relays: Vec<&str>) -> Self { + Self { + relays: relays.into_iter().map(str::to_owned).collect(), + } + } + } + + impl super::PublishAuthorRelayDiscovery for StaticPublishAuthorRelayDiscovery { + fn fetch_author_write_relays<'a>( + &'a self, + _pubkey: &'a str, + _discovery_targets: Vec<super::ResolvedPublishRelay>, + _connect_timeout_secs: u64, + ) -> super::PublishAuthorRelayDiscoveryFuture<'a> { + let relays = self.relays.clone(); + Box::pin(async move { Ok(relays) }) + } + } + + #[test] + fn token_generation_and_hashing_do_not_store_plaintext() { + let token = generate_bearer_token(); + assert!(token.starts_with("rrd_tp_")); + let hash = hash_bearer_token(token.as_str()); + assert!(hash.starts_with("sha256:")); + assert!(!hash.contains(token.as_str())); + } + + #[test] + fn nostr_source_policy_parser_accepts_contract_values() { + assert_eq!( + parse_nostr_source_policy("explicit_only").expect("policy"), + NostrPublishTargetSourcePolicy::ExplicitOnly + ); + assert!(parse_nostr_source_policy("unknown").is_err()); + } + + #[test] + fn storage_authenticates_hashed_tokens_and_scopes_jobs() { + let store = TransportPublishStore::memory().expect("store"); + let token = generate_bearer_token(); + let token_hash = hash_bearer_token(token.as_str()); + let principal = store + .create_principal(PublishPrincipalInit { + label: "tester".to_owned(), + token_hash: token_hash.clone(), + allowed_pubkeys: vec!["a".repeat(64)], + allowed_kinds: vec![30_402], + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: vec![ + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, + job_visibility: PublishJobVisibility::Own, + expires_at_unix: None, + }) + .expect("principal"); + assert_eq!( + store + .principal_for_token_hash(token_hash.as_str()) + .expect("lookup") + .expect("principal") + .principal_id, + principal.principal_id + ); + let denied = request("b".repeat(64).as_str(), 30_402); + assert!(principal.allows_event(&denied).is_err()); + + let accepted = request("a".repeat(64).as_str(), 30_402); + principal.allows_event(&accepted).expect("scope"); + let response = store + .record_publish_job(PublishJobInsert { + principal_id: principal.principal_id.clone(), + idempotency_key: Some("idem-1".to_owned()), + request: accepted.clone(), + request_fingerprint: "fingerprint-1".to_owned(), + effective_target_count: 1, + }) + .expect("record job"); + assert!(!response.deduplicated); + let duplicate = store + .record_publish_job(PublishJobInsert { + principal_id: principal.principal_id.clone(), + idempotency_key: Some("idem-1".to_owned()), + request: accepted, + request_fingerprint: "fingerprint-1".to_owned(), + effective_target_count: 1, + }) + .expect("dedupe"); + assert!(duplicate.deduplicated); + assert_eq!(duplicate.job.job_id, response.job.job_id); + assert_eq!( + store + .list_jobs_for_principal(&principal, 50) + .expect("jobs") + .len(), + 1 + ); + } + + #[test] + fn store_open_recovers_interrupted_publishing_jobs() { + let directory = tempfile::tempdir().expect("tempdir"); + let database_path = directory.path().join("publish-proxy.sqlite"); + let token_hash = hash_bearer_token(generate_bearer_token().as_str()); + let pubkey = "a".repeat(64); + let request = request(pubkey.as_str(), 30_402); + let job_id = { + let store = TransportPublishStore::open(database_path.clone()).expect("store"); + let principal = store + .create_principal(PublishPrincipalInit { + label: "tester".to_owned(), + token_hash, + allowed_pubkeys: vec![pubkey], + allowed_kinds: vec![30_402], + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: vec![ + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, + job_visibility: PublishJobVisibility::Own, + expires_at_unix: None, + }) + .expect("principal"); + let response = store + .record_publish_job(PublishJobInsert { + principal_id: principal.principal_id, + idempotency_key: Some("idem-interrupted".to_owned()), + request, + request_fingerprint: "fingerprint-interrupted".to_owned(), + effective_target_count: 1, + }) + .expect("record job"); + assert_eq!(response.job.status, TransportPublishJobStatus::Publishing); + response.job.job_id + }; + + let reopened = TransportPublishStore::open(database_path).expect("reopen store"); + let recovered = reopened.job_by_id(job_id.as_str()).expect("recovered job"); + assert_eq!( + recovered.status, + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + ); + assert_eq!( + recovered.last_error.as_deref(), + Some("publish_attempt_interrupted") + ); + assert!(recovered.completed_at_ms.is_some()); + assert!(recovered.targets.is_empty()); + } + + #[tokio::test] + async fn publish_event_verifies_and_records_daemon_default_outcome() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let event = signed_event(&identity, "{}"); + let raw_event = serde_json::to_string(&event).expect("raw event"); + let response = proxy + .publish_event( + &principal, + publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-valid"), + ), + ) + .await + .expect("publish"); + + assert!(!response.deduplicated); + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + assert_eq!(response.job.target_count, 1); + assert_eq!(response.job.acknowledged_count, 1); + assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); + assert_eq!( + response.job.targets[0].source, + TransportPublishTargetSource::DaemonDefault + ); + assert_eq!(adapter.captured_raw_events(), vec![raw_event]); + } + + #[tokio::test] + async fn publish_event_rejects_tampered_content_before_publish() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let mut event = signed_event(&identity, "trusted"); + event.content = "tampered".to_owned(); + let error = proxy + .publish_event( + &principal, + publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect_err("tampered event should fail"); + + assert!(matches!( + error, + TransportPublishError::SignedEventVerification( + RadrootsNostrEventVerification::IdMismatch + ) + )); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_rejects_wrong_signature_before_publish() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let mut event = signed_event(&identity, "{}"); + let replacement = if event.sig.starts_with('0') { "1" } else { "0" }; + event.sig.replace_range(0..1, replacement); + let error = proxy + .publish_event( + &principal, + publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect_err("wrong signature should fail"); + + assert!(matches!( + error, + TransportPublishError::SignedEventVerification( + RadrootsNostrEventVerification::SignatureInvalid + ) + )); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_rejects_malformed_wire_fields() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let mut event = signed_event(&identity, "{}"); + event.id = event.id.to_uppercase(); + let error = proxy + .publish_event( + &principal, + publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect_err("malformed field should fail"); + + assert!(matches!( + error, + TransportPublishError::InvalidSignedEvent(_) + )); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_uses_explicit_request_relays_when_allowed() { + let identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault], + true, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + vec![RELAY_PRIMARY.to_owned()], + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); + assert_eq!( + response.job.targets[0].source, + TransportPublishTargetSource::Request + ); + } + + #[tokio::test] + async fn publish_event_uses_cached_nip65_author_write_before_defaults() { + let identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); + proxy + .store + .cache_author_write_relays( + identity.public_key_hex().as_str(), + &[RELAY_PRIMARY.to_owned()], + ) + .expect("cache author relays"); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); + assert_eq!( + response.job.targets[0].source, + TransportPublishTargetSource::NostrAuthorWrite + ); + } + + #[tokio::test] + async fn publish_event_records_invalid_cached_author_write_relay() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); + proxy + .store + .cache_author_write_relays( + identity.public_key_hex().as_str(), + &[RELAY_PRIMARY.to_owned(), "not a cached relay".to_owned()], + ) + .expect("cache author relays"); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + let accepted = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == RELAY_PRIMARY) + .expect("accepted author relay"); + assert_eq!( + accepted.source, + TransportPublishTargetSource::NostrAuthorWrite + ); + assert!(accepted.attempted); + let rejected = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == "not a cached relay") + .expect("rejected cached author relay"); + assert_eq!( + rejected.source, + TransportPublishTargetSource::NostrAuthorWrite + ); + assert_eq!( + rejected.outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!rejected.attempted); + assert_eq!(adapter.captured_raw_events().len(), 1); + } + + #[tokio::test] + async fn publish_event_preserves_author_and_discovery_rejections_through_relay_selection() { + let identity = RadrootsIdentity::generate(); + let mut config = config_with_defaults(vec![RELAY_SECONDARY]); + config.nostr.author_relay_discovery_relays = vec!["not a discovery relay".to_owned()]; + let (proxy, adapter) = transport_publish(config); + proxy + .store + .cache_author_write_relays( + identity.public_key_hex().as_str(), + &["not a cached relay".to_owned()], + ) + .expect("cache author relays"); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + let daemon_default = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == RELAY_SECONDARY) + .expect("daemon default relay"); + assert_eq!( + daemon_default.source, + TransportPublishTargetSource::DaemonDefault + ); + assert!(daemon_default.attempted); + let cached = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == "not a cached relay") + .expect("cached author rejection"); + assert_eq!( + cached.source, + TransportPublishTargetSource::NostrAuthorWrite + ); + assert_eq!( + cached.outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!cached.attempted); + let discovery = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == "not a discovery relay") + .expect("discovery relay rejection"); + assert_eq!( + discovery.source, + TransportPublishTargetSource::DaemonDefault + ); + assert_eq!( + discovery.outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!discovery.attempted); + assert_eq!(adapter.captured_raw_events().len(), 1); + } + + #[tokio::test] + async fn publish_event_preserves_discovery_and_discovered_author_rejections() { + let identity = RadrootsIdentity::generate(); + let mut config = config_with_defaults(vec![RELAY_PRIMARY]); + config.nostr.author_relay_discovery_relays = + vec![RELAY_PRIMARY.to_owned(), RELAY_FORBIDDEN.to_owned()]; + let resolver = StaticPublishRelayResolver::new().with_addresses( + RELAY_FORBIDDEN, + vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))], + ); + let adapter = RadrootsMockRelayPublishAdapter::new(); + let proxy = TransportPublish::memory(config) + .expect("proxy") + .with_relay_resolver(Arc::new(resolver)) + .with_author_relay_discovery(Arc::new(StaticPublishAuthorRelayDiscovery::new(vec![ + "not a discovered author relay", + RELAY_SECONDARY, + ]))) + .with_publisher(Arc::new(adapter.clone())); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::AuthorWriteThenDaemonDefault, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + let accepted = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == RELAY_SECONDARY) + .expect("discovered author relay"); + assert_eq!( + accepted.source, + TransportPublishTargetSource::NostrAuthorWrite + ); + assert!(accepted.attempted); + let discovered = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == "not a discovered author relay") + .expect("discovered author rejection"); + assert_eq!( + discovered.source, + TransportPublishTargetSource::NostrAuthorWrite + ); + assert_eq!( + discovered.outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!discovered.attempted); + let discovery = response + .job + .targets + .iter() + .find(|relay| relay.endpoint_uri == RELAY_FORBIDDEN) + .expect("discovery relay rejection"); + assert_eq!( + discovery.source, + TransportPublishTargetSource::DaemonDefault + ); + assert_eq!( + discovery.outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!discovery.attempted); + assert_eq!(adapter.captured_raw_events().len(), 1); + } + + #[tokio::test] + async fn publish_event_records_no_transport_publish_targets_failure() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!(response.job.status, TransportPublishJobStatus::Rejected); + assert_eq!( + response.job.last_error.as_deref(), + Some("no_transport_publish_targets") + ); + assert!(response.job.targets.is_empty()); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_records_unsafe_request_relay_rejection() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::ExplicitOnly], + true, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + vec!["wss://127.0.0.1:7777".to_owned()], + NostrPublishTargetSourcePolicy::ExplicitOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + ); + assert_eq!(response.job.targets.len(), 1); + assert_eq!( + response.job.targets[0].outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!response.job.targets[0].attempted); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_rejects_forbidden_public_dns_destination_before_publish() { + let identity = RadrootsIdentity::generate(); + let resolver = StaticPublishRelayResolver::new() + .with_addresses(RELAY_PRIMARY, vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))]); + let (proxy, adapter) = transport_publish_with_resolver( + config_with_defaults(vec![RELAY_PRIMARY]), + Arc::new(resolver), + ); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + ); + assert_eq!(response.job.targets.len(), 1); + assert_eq!( + response.job.targets[0].outcome_kind, + TransportPublishOutcomeKind::TargetRejected + ); + assert!(!response.job.targets[0].attempted); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_records_dns_failure_as_unattempted_retryable_outcome() { + let identity = RadrootsIdentity::generate(); + let resolver = StaticPublishRelayResolver::new().with_failure(RELAY_PRIMARY, "no records"); + let (proxy, adapter) = transport_publish_with_resolver( + config_with_defaults(vec![RELAY_PRIMARY]), + Arc::new(resolver), + ); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + ); + assert_eq!( + response.job.last_error.as_deref(), + Some("delivery_unsatisfied") + ); + assert_eq!(response.job.targets.len(), 1); + assert_eq!( + response.job.targets[0].outcome_kind, + TransportPublishOutcomeKind::ConnectionFailed + ); + assert!(!response.job.targets[0].attempted); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_localhost_policy_skips_public_dns_guard() { + let identity = RadrootsIdentity::generate(); + let mut config = config_with_defaults(vec!["ws://localhost:7777"]); + config.nostr.relay_url_policy = NostrRelayUrlPolicy::Localhost; + let resolver = StaticPublishRelayResolver::new() + .with_failure("ws://localhost:7777", "localhost resolution should not run"); + let (proxy, adapter) = transport_publish_with_resolver(config, Arc::new(resolver)); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliverySatisfied + ); + assert_eq!(response.job.targets[0].endpoint_uri, "ws://localhost:7777"); + assert!(!adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_deduplicates_same_intent_and_conflicts_different_intent() { + let identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let request = publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-conflict"), + ); + let first = proxy + .publish_event(&principal, request.clone()) + .await + .expect("first"); + let duplicate = proxy + .publish_event(&principal, request) + .await + .expect("duplicate"); + + assert!(!first.deduplicated); + assert!(duplicate.deduplicated); + assert_eq!(duplicate.job.job_id, first.job.job_id); + + let conflict = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "changed"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-conflict"), + ), + ) + .await + .expect_err("conflict"); + assert!(matches!( + conflict, + TransportPublishError::IdempotencyConflict(_) + )); + } + + #[tokio::test] + async fn publish_event_rejects_zero_and_excessive_timeout_before_job_creation() { + let identity = RadrootsIdentity::generate(); + let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let mut zero = publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-zero-timeout"), + ); + zero.timeout_ms = Some(0); + let zero_error = proxy + .publish_event(&principal, zero) + .await + .expect_err("zero timeout should fail"); + assert!(matches!( + zero_error, + TransportPublishError::InvalidSignedEvent(_) + )); + + let mut excessive = publish_request( + signed_event(&identity, "changed"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-excessive-timeout"), + ); + excessive.timeout_ms = Some(10_001); + let excessive_error = proxy + .publish_event(&principal, excessive) + .await + .expect_err("excessive timeout should fail"); + assert!(matches!( + excessive_error, + TransportPublishError::InvalidSignedEvent(_) + )); + assert!( + proxy + .store + .list_jobs_for_principal(&principal, 50) + .expect("jobs") + .is_empty() + ); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_event_default_timeout_fingerprints_as_effective_timeout() { + let identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let event = signed_event(&identity, "{}"); + let mut default_timeout = publish_request( + event.clone(), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-default-timeout"), + ); + default_timeout.timeout_ms = None; + let mut explicit_default = publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-default-timeout"), + ); + explicit_default.timeout_ms = Some(10_000); + + let first = proxy + .publish_event(&principal, default_timeout) + .await + .expect("first"); + let duplicate = proxy + .publish_event(&principal, explicit_default) + .await + .expect("duplicate"); + assert!(!first.deduplicated); + assert!(duplicate.deduplicated); + assert_eq!(duplicate.job.job_id, first.job.job_id); + } + + #[tokio::test] + async fn publish_event_fingerprint_conflicts_on_different_effective_timeout() { + let identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let event = signed_event(&identity, "{}"); + let first = publish_request( + event.clone(), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-timeout-conflict"), + ); + let mut conflict = publish_request( + event, + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-timeout-conflict"), + ); + conflict.timeout_ms = Some(6_000); + + proxy.publish_event(&principal, first).await.expect("first"); + let error = proxy + .publish_event(&principal, conflict) + .await + .expect_err("timeout conflict"); + assert!(matches!( + error, + TransportPublishError::IdempotencyConflict(_) + )); + } + + #[tokio::test] + async fn publish_event_concurrency_limit_rejects_without_job_creation() { + let identity = RadrootsIdentity::generate(); + let mut config = config_with_defaults(vec![RELAY_PRIMARY]); + config.max_concurrent_publish_jobs = 1; + let (proxy, adapter) = transport_publish(config); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let _permit = proxy.acquire_publish_permit().expect("permit"); + let error = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + Some("idem-concurrency"), + ), + ) + .await + .expect_err("concurrency limit"); + assert!(matches!(error, TransportPublishError::ConcurrencyLimit)); + assert!( + proxy + .store + .list_jobs_for_principal(&principal, 50) + .expect("jobs") + .is_empty() + ); + assert!(adapter.captured_raw_events().is_empty()); + } + + #[tokio::test] + async fn publish_jobs_respect_own_and_admin_visibility() { + let identity = RadrootsIdentity::generate(); + let other_identity = RadrootsIdentity::generate(); + let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); + let owner = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let other = principal( + &proxy, + other_identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let admin = principal( + &proxy, + other_identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Admin, + ); + let response = proxy + .publish_event( + &owner, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert!( + proxy + .store + .job_by_id_for_principal(response.job.job_id.as_str(), &other) + .expect("other read") + .is_none() + ); + assert!( + proxy + .store + .job_by_id_for_principal(response.job.job_id.as_str(), &admin) + .expect("admin read") + .is_some() + ); + } + + #[tokio::test] + async fn publish_event_records_retryable_relay_failures() { + let identity = RadrootsIdentity::generate(); + let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( + RELAY_PRIMARY, + RadrootsRelayOutcome::connection_failed("error: unavailable"), + ); + let proxy = TransportPublish::memory(config_with_defaults(vec![RELAY_PRIMARY])) + .expect("proxy") + .with_publisher(Arc::new(adapter)); + let principal = principal( + &proxy, + identity.public_key_hex(), + vec![NostrPublishTargetSourcePolicy::DaemonDefaultOnly], + false, + PublishJobVisibility::Own, + ); + let response = proxy + .publish_event( + &principal, + publish_request( + signed_event(&identity, "{}"), + Vec::new(), + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + TransportPublishDeliveryPolicy::Any, + None, + ), + ) + .await + .expect("publish"); + + assert_eq!( + response.job.status, + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + ); + assert_eq!(response.job.retryable_count, 1); + } +} diff --git a/src/transport/jsonrpc/auth.rs b/src/transport/jsonrpc/auth.rs @@ -2,26 +2,26 @@ use jsonrpsee::core::server::Extensions; -use crate::core::publish_proxy::{PublishPrincipal, PublishProxyStore, hash_bearer_token}; +use crate::core::transport_publish::{PublishPrincipal, TransportPublishStore, hash_bearer_token}; use super::RpcError; #[cfg(test)] -pub(crate) const PUBLISH_PROXY_AUTH_MODE: &str = "scoped_bearer_token"; +pub(crate) const TRANSPORT_PUBLISH_AUTH_MODE: &str = "scoped_bearer_token"; #[derive(Clone, Debug, PartialEq, Eq)] -pub(crate) enum PublishProxyAuthorization { +pub(crate) enum TransportPublishAuthorization { Authorized(PublishPrincipal), Missing, Invalid, } -pub(crate) fn authorize_publish_proxy_request( +pub(crate) fn authorize_transport_publish_request( authorization_header: Option<&str>, - store: &PublishProxyStore, -) -> PublishProxyAuthorization { + store: &TransportPublishStore, +) -> TransportPublishAuthorization { let Some(authorization_header) = authorization_header else { - return PublishProxyAuthorization::Missing; + return TransportPublishAuthorization::Missing; }; let mut parts = authorization_header.split_whitespace(); @@ -29,12 +29,12 @@ pub(crate) fn authorize_publish_proxy_request( let token = parts.next().unwrap_or_default(); if !scheme.eq_ignore_ascii_case("bearer") || token.is_empty() || parts.next().is_some() { - return PublishProxyAuthorization::Invalid; + return TransportPublishAuthorization::Invalid; } match store.principal_for_token_hash(hash_bearer_token(token).as_str()) { - Ok(Some(principal)) => PublishProxyAuthorization::Authorized(principal), - Ok(None) | Err(_) => PublishProxyAuthorization::Invalid, + Ok(Some(principal)) => TransportPublishAuthorization::Authorized(principal), + Ok(None) | Err(_) => TransportPublishAuthorization::Invalid, } } @@ -42,16 +42,16 @@ pub(crate) fn require_publish_principal( extensions: &Extensions, ) -> Result<PublishPrincipal, RpcError> { match extensions - .get::<PublishProxyAuthorization>() + .get::<TransportPublishAuthorization>() .cloned() - .unwrap_or(PublishProxyAuthorization::Missing) + .unwrap_or(TransportPublishAuthorization::Missing) { - PublishProxyAuthorization::Authorized(principal) => Ok(principal), - PublishProxyAuthorization::Missing => Err(RpcError::Unauthorized( - "publish proxy bearer token required".to_string(), + TransportPublishAuthorization::Authorized(principal) => Ok(principal), + TransportPublishAuthorization::Missing => Err(RpcError::Unauthorized( + "transport publish bearer token required".to_string(), )), - PublishProxyAuthorization::Invalid => Err(RpcError::Unauthorized( - "invalid publish proxy bearer token".to_string(), + TransportPublishAuthorization::Invalid => Err(RpcError::Unauthorized( + "invalid transport publish bearer token".to_string(), )), } } @@ -59,19 +59,21 @@ pub(crate) fn require_publish_principal( #[cfg(test)] mod tests { use jsonrpsee::core::server::Extensions; - use radroots_publish_proxy_protocol::PublishRelayPolicy; + use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, TransportPublishTargetPolicyName, + }; use super::{ - PUBLISH_PROXY_AUTH_MODE, PublishProxyAuthorization, authorize_publish_proxy_request, - require_publish_principal, + TRANSPORT_PUBLISH_AUTH_MODE, TransportPublishAuthorization, + authorize_transport_publish_request, require_publish_principal, }; - use crate::core::publish_proxy::{ - PublishJobVisibility, PublishPrincipalInit, PublishProxyStore, generate_bearer_token, + use crate::core::transport_publish::{ + PublishJobVisibility, PublishPrincipalInit, TransportPublishStore, generate_bearer_token, hash_bearer_token, }; - fn store_with_token() -> (PublishProxyStore, String) { - let store = PublishProxyStore::memory().expect("store"); + fn store_with_token() -> (TransportPublishStore, String) { + let store = TransportPublishStore::memory().expect("store"); let token = generate_bearer_token(); store .create_principal(PublishPrincipalInit { @@ -79,8 +81,11 @@ mod tests { token_hash: hash_bearer_token(token.as_str()), allowed_pubkeys: vec!["a".repeat(64)], allowed_kinds: vec![30_402], - allowed_relay_policies: vec![PublishRelayPolicy::DaemonDefaultOnly], - allow_request_relays: false, + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: vec![ + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, job_visibility: PublishJobVisibility::Own, expires_at_unix: None, }) @@ -89,28 +94,28 @@ mod tests { } #[test] - fn publish_proxy_auth_accepts_matching_bearer_token() { + fn transport_publish_auth_accepts_matching_bearer_token() { let (store, token) = store_with_token(); let header = format!("Bearer {token}"); - let auth = authorize_publish_proxy_request(Some(header.as_str()), &store); - assert!(matches!(auth, PublishProxyAuthorization::Authorized(_))); - assert_eq!(PUBLISH_PROXY_AUTH_MODE, "scoped_bearer_token"); + let auth = authorize_transport_publish_request(Some(header.as_str()), &store); + assert!(matches!(auth, TransportPublishAuthorization::Authorized(_))); + assert_eq!(TRANSPORT_PUBLISH_AUTH_MODE, "scoped_bearer_token"); } #[test] - fn publish_proxy_auth_rejects_missing_and_invalid_headers() { + fn transport_publish_auth_rejects_missing_and_invalid_headers() { let (store, _token) = store_with_token(); assert_eq!( - authorize_publish_proxy_request(None, &store), - PublishProxyAuthorization::Missing + authorize_transport_publish_request(None, &store), + TransportPublishAuthorization::Missing ); assert_eq!( - authorize_publish_proxy_request(Some("Basic secret"), &store), - PublishProxyAuthorization::Invalid + authorize_transport_publish_request(Some("Basic secret"), &store), + TransportPublishAuthorization::Invalid ); assert_eq!( - authorize_publish_proxy_request(Some("Bearer wrong"), &store), - PublishProxyAuthorization::Invalid + authorize_transport_publish_request(Some("Bearer wrong"), &store), + TransportPublishAuthorization::Invalid ); } @@ -118,7 +123,7 @@ mod tests { fn require_publish_principal_reads_authorized_extensions() { let (store, token) = store_with_token(); let header = format!("Bearer {token}"); - let auth = authorize_publish_proxy_request(Some(header.as_str()), &store); + let auth = authorize_transport_publish_request(Some(header.as_str()), &store); let mut extensions = Extensions::new(); extensions.insert(auth); require_publish_principal(&extensions).expect("authorized"); diff --git a/src/transport/jsonrpc/methods/mod.rs b/src/transport/jsonrpc/methods/mod.rs @@ -6,15 +6,15 @@ use jsonrpsee::server::RpcModule; use crate::transport::jsonrpc::{MethodRegistry, RpcContext}; pub mod nip46; -pub mod publish_proxy; +pub mod transport_publish; pub fn register_all( root: &mut RpcModule<RpcContext>, ctx: RpcContext, registry: MethodRegistry, ) -> Result<()> { - if ctx.state.publish_proxy.config.enabled { - root.merge(publish_proxy::module(ctx.clone(), registry.clone())?)?; + if ctx.state.transport_publish.config.enabled { + root.merge(transport_publish::module(ctx.clone(), registry.clone())?)?; } if ctx.state.nip46_config.public_jsonrpc_enabled { root.merge(nip46::module(ctx, registry)?)?; @@ -29,42 +29,41 @@ mod tests { use radroots_nostr::prelude::RadrootsNostrMetadata; use super::register_all; - use crate::app::config::{Nip46Config, PublishProxyConfig}; + use crate::app::config::{Nip46Config, TransportPublishConfig}; use crate::core::Radrootsd; - use crate::transport::jsonrpc::auth::PublishProxyAuthorization; + use crate::transport::jsonrpc::auth::TransportPublishAuthorization; use crate::transport::jsonrpc::{MethodRegistry, RpcContext}; mod removed_surface_fixtures { pub const BRIDGE_STATUS_METHOD: &str = "bridge.status"; } - fn state(publish_proxy_enabled: bool, nip46_public_jsonrpc_enabled: bool) -> Radrootsd { + fn state(transport_publish_enabled: bool, nip46_public_jsonrpc_enabled: bool) -> Radrootsd { let identity = RadrootsIdentity::generate(); let metadata: RadrootsNostrMetadata = serde_json::from_str(r#"{"name":"radrootsd-test"}"#).expect("metadata"); - let publish_proxy = PublishProxyConfig { - enabled: publish_proxy_enabled, - ..PublishProxyConfig::default() + let transport_publish = TransportPublishConfig { + enabled: transport_publish_enabled, + ..TransportPublishConfig::default() }; let nip46 = Nip46Config { public_jsonrpc_enabled: nip46_public_jsonrpc_enabled, ..Nip46Config::default() }; - Radrootsd::new(identity, metadata, publish_proxy, nip46).expect("state") + Radrootsd::new(identity, metadata, transport_publish, nip46).expect("state") } #[test] - fn register_all_exposes_publish_proxy_methods_by_default() { + fn register_all_exposes_transport_publish_methods_by_default() { let registry = MethodRegistry::default(); let ctx = RpcContext::new(state(true, false), registry.clone()); let mut root = RpcModule::new(ctx.clone()); register_all(&mut root, ctx, registry).expect("register"); - assert!(root.method("publish.capabilities").is_some()); - assert!(root.method("publish.event").is_some()); - assert!(root.method("publish.job.get").is_some()); - assert!(root.method("publish.job.list").is_some()); - assert!(root.method("publish.relays.resolve").is_some()); + assert!(root.method("transport.publish.capabilities").is_some()); + assert!(root.method("transport.publish.event").is_some()); + assert!(root.method("transport.publish.job.get").is_some()); + assert!(root.method("transport.publish.job.list").is_some()); assert!( root.method(removed_surface_fixtures::BRIDGE_STATUS_METHOD) .is_none() @@ -79,7 +78,7 @@ mod tests { let mut root = RpcModule::new(ctx.clone()); register_all(&mut root, ctx, registry).expect("register"); - assert!(root.method("publish.capabilities").is_some()); + assert!(root.method("transport.publish.capabilities").is_some()); assert!(root.method("nip46.connect").is_some()); } @@ -92,7 +91,7 @@ mod tests { let (response, _stream) = root .raw_json_request( - r#"{"jsonrpc":"2.0","method":"publish.capabilities","id":1}"#, + r#"{"jsonrpc":"2.0","method":"transport.publish.capabilities","id":1}"#, 1, ) .await @@ -106,29 +105,32 @@ mod tests { let ctx = RpcContext::new(state(true, false), registry.clone()); let principal = ctx .state - .publish_proxy + .transport_publish .store - .create_principal(crate::core::publish_proxy::PublishPrincipalInit { + .create_principal(crate::core::transport_publish::PublishPrincipalInit { label: "tester".to_owned(), - token_hash: crate::core::publish_proxy::hash_bearer_token("secret"), + token_hash: crate::core::transport_publish::hash_bearer_token("secret"), allowed_pubkeys: vec!["a".repeat(64)], allowed_kinds: vec![30_402], - allowed_relay_policies: vec![ - radroots_publish_proxy_protocol::PublishRelayPolicy::DaemonDefaultOnly, + allowed_target_policies: vec![ + radroots_transport_publish_protocol::TransportPublishTargetPolicyName::Nostr, ], - allow_request_relays: false, - job_visibility: crate::core::publish_proxy::PublishJobVisibility::Own, + allowed_nostr_source_policies: vec![ + radroots_transport_publish_protocol::NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, + job_visibility: crate::core::transport_publish::PublishJobVisibility::Own, expires_at_unix: None, }) .expect("principal"); let mut root = RpcModule::new(ctx.clone()); root.extensions_mut() - .insert(PublishProxyAuthorization::Authorized(principal)); + .insert(TransportPublishAuthorization::Authorized(principal)); register_all(&mut root, ctx, registry).expect("register"); let (response, _stream) = root .raw_json_request( - r#"{"jsonrpc":"2.0","method":"publish.capabilities","id":1}"#, + r#"{"jsonrpc":"2.0","method":"transport.publish.capabilities","id":1}"#, 1, ) .await diff --git a/src/transport/jsonrpc/methods/publish_proxy.rs b/src/transport/jsonrpc/methods/publish_proxy.rs @@ -1,526 +0,0 @@ -use anyhow::Result; -use jsonrpsee::server::RpcModule; -use radroots_publish_proxy_protocol::{ - METHOD_CAPABILITIES, METHOD_EVENT, METHOD_JOB_GET, METHOD_JOB_LIST, METHOD_RELAYS_RESOLVE, - PublishCapabilities, PublishDeliveryPolicy, PublishEventRequest, PublishRelayOutcome, - PublishRelaySource, -}; -use serde::{Deserialize, Serialize}; - -use crate::core::publish_proxy::PublishProxyError; -use crate::transport::jsonrpc::auth::require_publish_principal; -use crate::transport::jsonrpc::{MethodRegistry, RpcContext, RpcError}; - -#[derive(Debug, Deserialize)] -struct JobGetParams { - job_id: String, -} - -#[derive(Debug, Deserialize)] -#[serde(deny_unknown_fields)] -struct JobListParams { - limit: Option<usize>, -} - -#[derive(Debug, Deserialize)] -struct RelaysResolveParams { - event: radroots_publish_proxy_protocol::SignedNostrEventWire, - relay_policy: radroots_publish_proxy_protocol::PublishRelayPolicy, - #[serde(default)] - relays: Vec<String>, -} - -#[derive(Clone, Debug, Serialize)] -struct RelaysResolveResponse { - relays: Vec<ResolvedRelayResponseItem>, - rejected_relays: Vec<PublishRelayOutcome>, -} - -#[derive(Clone, Debug, Serialize)] -struct ResolvedRelayResponseItem { - relay_url: String, - source: PublishRelaySource, -} - -pub fn module(ctx: RpcContext, registry: MethodRegistry) -> Result<RpcModule<RpcContext>> { - let mut module = RpcModule::new(ctx); - register_capabilities(&mut module, &registry)?; - register_event(&mut module, &registry)?; - register_job_get(&mut module, &registry)?; - register_job_list(&mut module, &registry)?; - register_relays_resolve(&mut module, &registry)?; - Ok(module) -} - -fn register_capabilities( - module: &mut RpcModule<RpcContext>, - registry: &MethodRegistry, -) -> Result<()> { - registry.track(METHOD_CAPABILITIES); - module.register_async_method(METHOD_CAPABILITIES, |_params, ctx, extensions| async move { - require_publish_principal(&extensions)?; - Ok::<PublishCapabilities, RpcError>(PublishCapabilities::v1( - ctx.state.publish_proxy.config.max_event_bytes, - ctx.state.publish_proxy.config.max_relays_per_request, - )) - })?; - Ok(()) -} - -fn register_event(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { - registry.track(METHOD_EVENT); - module.register_async_method(METHOD_EVENT, |params, ctx, extensions| async move { - let principal = require_publish_principal(&extensions)?; - let request: PublishEventRequest = params - .parse() - .map_err(|error| RpcError::InvalidParams(error.to_string()))?; - ctx.state - .publish_proxy - .publish_event(&principal, request) - .await - .map_err(rpc_error_from_publish_proxy) - })?; - Ok(()) -} - -fn register_job_get(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { - registry.track(METHOD_JOB_GET); - module.register_async_method(METHOD_JOB_GET, |params, ctx, extensions| async move { - let principal = require_publish_principal(&extensions)?; - let params: JobGetParams = params - .parse() - .map_err(|error| RpcError::InvalidParams(error.to_string()))?; - let job_id = params.job_id.trim(); - if job_id.is_empty() { - return Err(RpcError::InvalidParams("missing job_id".to_owned())); - } - ctx.state - .publish_proxy - .store - .job_by_id_for_principal(job_id, &principal) - .map_err(|error| RpcError::Other(error.to_string()))? - .ok_or_else(|| RpcError::Other(format!("unknown publish job: {job_id}"))) - })?; - Ok(()) -} - -fn register_job_list(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { - registry.track(METHOD_JOB_LIST); - module.register_async_method(METHOD_JOB_LIST, |params, ctx, extensions| async move { - let principal = require_publish_principal(&extensions)?; - let params = if params.len_bytes() == 0 || params.as_str() == Some("[]") { - JobListParams { limit: None } - } else { - params - .parse::<JobListParams>() - .map_err(|error| RpcError::InvalidParams(error.to_string()))? - }; - if params.limit == Some(0) { - return Err(RpcError::InvalidParams( - "limit must be greater than zero".to_owned(), - )); - } - let configured_limit = ctx.state.publish_proxy.config.job_list_limit; - let limit = params - .limit - .unwrap_or(configured_limit) - .min(configured_limit); - ctx.state - .publish_proxy - .store - .list_jobs_for_principal(&principal, limit) - .map_err(|error| RpcError::Other(error.to_string())) - })?; - Ok(()) -} - -fn register_relays_resolve( - module: &mut RpcModule<RpcContext>, - registry: &MethodRegistry, -) -> Result<()> { - registry.track(METHOD_RELAYS_RESOLVE); - module.register_async_method( - METHOD_RELAYS_RESOLVE, - |params, ctx, extensions| async move { - let principal = require_publish_principal(&extensions)?; - let params: RelaysResolveParams = params - .parse() - .map_err(|error| RpcError::InvalidParams(error.to_string()))?; - params - .event - .validate() - .map_err(|error| RpcError::InvalidParams(error.to_string()))?; - let request = PublishEventRequest { - event: params.event, - relays: params.relays, - relay_policy: params.relay_policy, - delivery_policy: PublishDeliveryPolicy::Any, - idempotency_key: None, - timeout_ms: None, - }; - principal - .allows_event(&request) - .map_err(|error| RpcError::Unauthorized(error.to_string()))?; - let resolution = ctx - .state - .publish_proxy - .resolve_relays_for_request(request.event.pubkey.as_str(), &request) - .await - .map_err(rpc_error_from_publish_proxy)?; - Ok::<RelaysResolveResponse, RpcError>(RelaysResolveResponse { - relays: resolution - .targets - .into_iter() - .map(|target| ResolvedRelayResponseItem { - relay_url: target.url.into_string(), - source: target.source, - }) - .collect(), - rejected_relays: resolution.outcomes, - }) - }, - )?; - Ok(()) -} - -fn rpc_error_from_publish_proxy(error: PublishProxyError) -> RpcError { - match error { - PublishProxyError::InvalidScope(message) => RpcError::Unauthorized(message), - PublishProxyError::InvalidSignedEvent(message) => RpcError::InvalidParams(message), - PublishProxyError::SignedEventVerification(_) - | PublishProxyError::Draft(_) - | PublishProxyError::Relay(_) => RpcError::InvalidParams(error.to_string()), - PublishProxyError::IdempotencyConflict(_) => RpcError::Other(error.to_string()), - other => RpcError::Other(other.to_string()), - } -} - -#[cfg(test)] -mod tests { - use super::module; - use std::sync::Arc; - - use crate::app::config::{Nip46Config, PublishProxyConfig}; - use crate::core::Radrootsd; - use crate::core::publish_proxy::{ - PublishJobVisibility, PublishPrincipalInit, generate_bearer_token, hash_bearer_token, - }; - use crate::transport::jsonrpc::auth::{ - PublishProxyAuthorization, authorize_publish_proxy_request, - }; - use crate::transport::jsonrpc::{MethodRegistry, RpcContext}; - use jsonrpsee::server::RpcModule; - use nostr::JsonUtil; - use radroots_identity::RadrootsIdentity; - use radroots_nostr::prelude::{ - RadrootsNostrMetadata, RadrootsNostrTimestamp, radroots_nostr_build_event, - }; - use radroots_publish_proxy_protocol::{PublishRelayPolicy, SignedNostrEventWire}; - use radroots_relay_transport::RadrootsMockRelayPublishAdapter; - - fn signed_event(identity: &RadrootsIdentity) -> SignedNostrEventWire { - let event = radroots_nostr_build_event( - 30_402, - "{}", - vec![vec!["d".to_owned(), "listing-1".to_owned()]], - ) - .expect("event builder") - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(identity.keys()) - .expect("signed event"); - serde_json::from_str(event.as_json().as_str()).expect("event wire") - } - - fn module_with_principal_and_config( - admin: bool, - publish_proxy_config: PublishProxyConfig, - ) -> ( - RpcModule<RpcContext>, - RpcContext, - String, - SignedNostrEventWire, - ) { - let identity = RadrootsIdentity::generate(); - let signed_event = signed_event(&identity); - let metadata: RadrootsNostrMetadata = - serde_json::from_str(r#"{"name":"radrootsd-test"}"#).expect("metadata"); - let state = Radrootsd::new( - identity.clone(), - metadata, - publish_proxy_config, - Nip46Config::default(), - ) - .expect("state"); - let mut state = state; - state.publish_proxy = state - .publish_proxy - .clone() - .with_publisher(Arc::new(RadrootsMockRelayPublishAdapter::new())); - let token = generate_bearer_token(); - let principal = state - .publish_proxy - .store - .create_principal(PublishPrincipalInit { - label: "tester".to_owned(), - token_hash: hash_bearer_token(token.as_str()), - allowed_pubkeys: vec![identity.public_key_hex()], - allowed_kinds: vec![30_402], - allowed_relay_policies: vec![PublishRelayPolicy::DaemonDefaultOnly], - allow_request_relays: false, - job_visibility: if admin { - PublishJobVisibility::Admin - } else { - PublishJobVisibility::Own - }, - expires_at_unix: None, - }) - .expect("principal"); - let registry = MethodRegistry::default(); - let ctx = RpcContext::new(state, registry.clone()); - let mut module = module(ctx.clone(), registry).expect("module"); - module - .extensions_mut() - .insert(PublishProxyAuthorization::Authorized(principal)); - (module, ctx, token, signed_event) - } - - fn module_with_principal( - admin: bool, - ) -> ( - RpcModule<RpcContext>, - RpcContext, - String, - SignedNostrEventWire, - ) { - module_with_principal_and_config( - admin, - PublishProxyConfig { - daemon_default_publish_relays: vec!["wss://relay.example.com".to_owned()], - ..PublishProxyConfig::default() - }, - ) - } - - #[tokio::test] - async fn publish_event_records_job_and_deduplicates_idempotency() { - let (module, _ctx, _token, event) = module_with_principal(false); - let request = format!( - r#"{{ - "jsonrpc":"2.0", - "method":"publish.event", - "params":{{ - "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", - "delivery_policy":{{"mode":"any"}}, - "idempotency_key":"idem-1" - }}, - "id":1 - }}"#, - serde_json::to_string(&event).expect("event json") - ); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("request"); - assert!(response.get().contains("\"deduplicated\":false")); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("request"); - assert!(response.get().contains("\"deduplicated\":true")); - } - - #[tokio::test] - async fn publish_event_rejects_principal_scope_gap() { - let (module, _ctx, _token, _pubkey) = module_with_principal(false); - let other_identity = RadrootsIdentity::generate(); - let event = signed_event(&other_identity); - let request = format!( - r#"{{ - "jsonrpc":"2.0", - "method":"publish.event", - "params":{{ - "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", - "delivery_policy":{{"mode":"any"}} - }}, - "id":1 - }}"#, - serde_json::to_string(&event).expect("event json") - ); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("request"); - assert!(response.get().contains("unauthorized")); - } - - #[tokio::test] - async fn publish_job_list_rejects_malformed_and_zero_limits() { - let (module, _ctx, _token, _event) = module_with_principal(false); - let malformed = r#"{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "params":"bad", - "id":1 - }"#; - let (response, _stream) = module - .raw_json_request(malformed, 1) - .await - .expect("malformed request"); - assert!(response.get().contains("\"code\":-32602")); - - let zero = r#"{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "params":{"limit":0}, - "id":1 - }"#; - let (response, _stream) = module - .raw_json_request(zero, 1) - .await - .expect("zero request"); - assert!(response.get().contains("\"code\":-32602")); - assert!(response.get().contains("limit must be greater than zero")); - } - - #[tokio::test] - async fn publish_job_list_rejects_unknown_fields() { - let (module, _ctx, _token, _event) = module_with_principal(false); - for params in [ - r#"{"cursor":"next"}"#, - r#"{"status":"publishing"}"#, - r#"{"limit":1,"extra":true}"#, - ] { - let request = format!( - r#"{{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "params":{params}, - "id":1 - }}"# - ); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("unknown field request"); - assert!( - response.get().contains("\"code\":-32602"), - "{}", - response.get() - ); - } - } - - #[tokio::test] - async fn publish_job_list_uses_configured_limit_when_omitted_and_caps_positive_limits() { - let mut config = PublishProxyConfig { - daemon_default_publish_relays: vec!["wss://relay.example.com".to_owned()], - ..PublishProxyConfig::default() - }; - config.job_list_limit = 1; - let (module, _ctx, _token, event) = module_with_principal_and_config(false, config); - for idempotency_key in ["idem-list-1", "idem-list-2"] { - let request = format!( - r#"{{ - "jsonrpc":"2.0", - "method":"publish.event", - "params":{{ - "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", - "delivery_policy":{{"mode":"any"}}, - "idempotency_key":"{idempotency_key}" - }}, - "id":1 - }}"#, - serde_json::to_string(&event).expect("event json") - ); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("publish request"); - assert!(response.get().contains("\"deduplicated\":false")); - } - - let omitted = r#"{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "id":1 - }"#; - let (response, _stream) = module - .raw_json_request(omitted, 1) - .await - .expect("omitted request"); - let value: serde_json::Value = - serde_json::from_str(response.get()).expect("omitted response json"); - assert_eq!(value["result"].as_array().expect("jobs").len(), 1); - - let empty_array = r#"{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "params":[], - "id":1 - }"#; - let (response, _stream) = module - .raw_json_request(empty_array, 1) - .await - .expect("empty array request"); - let value: serde_json::Value = - serde_json::from_str(response.get()).expect("empty array response json"); - assert_eq!(value["result"].as_array().expect("jobs").len(), 1); - - let over_limit = r#"{ - "jsonrpc":"2.0", - "method":"publish.job.list", - "params":{"limit":50}, - "id":1 - }"#; - let (response, _stream) = module - .raw_json_request(over_limit, 1) - .await - .expect("over limit request"); - let value: serde_json::Value = - serde_json::from_str(response.get()).expect("over limit response json"); - assert_eq!(value["result"].as_array().expect("jobs").len(), 1); - } - - #[tokio::test] - async fn publish_relays_resolve_returns_daemon_default_targets() { - let (module, _ctx, _token, event) = module_with_principal(false); - let request = format!( - r#"{{ - "jsonrpc":"2.0", - "method":"publish.relays.resolve", - "params":{{ - "event":{}, - "relay_policy":"daemon_default_only", - "relays":[] - }}, - "id":1 - }}"#, - serde_json::to_string(&event).expect("event json") - ); - let (response, _stream) = module - .raw_json_request(request.as_str(), 1) - .await - .expect("request"); - assert!( - response - .get() - .contains("\"relay_url\":\"wss://relay.example.com\"") - ); - assert!(response.get().contains("\"source\":\"daemon_default\"")); - } - - #[test] - fn http_auth_finds_principal_from_hashed_token() { - let (_module, ctx, token, _pubkey) = module_with_principal(false); - let header = format!("Bearer {token}"); - let auth = - authorize_publish_proxy_request(Some(header.as_str()), &ctx.state.publish_proxy.store); - assert!(matches!(auth, PublishProxyAuthorization::Authorized(_))); - } -} diff --git a/src/transport/jsonrpc/methods/transport_publish.rs b/src/transport/jsonrpc/methods/transport_publish.rs @@ -0,0 +1,437 @@ +use anyhow::Result; +use jsonrpsee::server::RpcModule; +use radroots_transport_publish_protocol::{ + METHOD_CAPABILITIES, METHOD_EVENT, METHOD_JOB_GET, METHOD_JOB_LIST, + TransportPublishCapabilities, TransportPublishEventRequest, +}; +use serde::Deserialize; + +use crate::core::transport_publish::TransportPublishError; +use crate::transport::jsonrpc::auth::require_publish_principal; +use crate::transport::jsonrpc::{MethodRegistry, RpcContext, RpcError}; + +#[derive(Debug, Deserialize)] +struct JobGetParams { + job_id: String, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct JobListParams { + limit: Option<usize>, +} + +pub fn module(ctx: RpcContext, registry: MethodRegistry) -> Result<RpcModule<RpcContext>> { + let mut module = RpcModule::new(ctx); + register_capabilities(&mut module, &registry)?; + register_event(&mut module, &registry)?; + register_job_get(&mut module, &registry)?; + register_job_list(&mut module, &registry)?; + Ok(module) +} + +fn register_capabilities( + module: &mut RpcModule<RpcContext>, + registry: &MethodRegistry, +) -> Result<()> { + registry.track(METHOD_CAPABILITIES); + module.register_async_method(METHOD_CAPABILITIES, |_params, ctx, extensions| async move { + require_publish_principal(&extensions)?; + Ok::<TransportPublishCapabilities, RpcError>(TransportPublishCapabilities::v2( + ctx.state.transport_publish.config.max_event_bytes, + ctx.state.transport_publish.config.max_targets_per_request, + )) + })?; + Ok(()) +} + +fn register_event(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { + registry.track(METHOD_EVENT); + module.register_async_method(METHOD_EVENT, |params, ctx, extensions| async move { + let principal = require_publish_principal(&extensions)?; + let request: TransportPublishEventRequest = params + .parse() + .map_err(|error| RpcError::InvalidParams(error.to_string()))?; + ctx.state + .transport_publish + .publish_event(&principal, request) + .await + .map_err(rpc_error_from_transport_publish) + })?; + Ok(()) +} + +fn register_job_get(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { + registry.track(METHOD_JOB_GET); + module.register_async_method(METHOD_JOB_GET, |params, ctx, extensions| async move { + let principal = require_publish_principal(&extensions)?; + let params: JobGetParams = params + .parse() + .map_err(|error| RpcError::InvalidParams(error.to_string()))?; + let job_id = params.job_id.trim(); + if job_id.is_empty() { + return Err(RpcError::InvalidParams("missing job_id".to_owned())); + } + ctx.state + .transport_publish + .store + .job_by_id_for_principal(job_id, &principal) + .map_err(|error| RpcError::Other(error.to_string()))? + .ok_or_else(|| RpcError::Other(format!("unknown publish job: {job_id}"))) + })?; + Ok(()) +} + +fn register_job_list(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { + registry.track(METHOD_JOB_LIST); + module.register_async_method(METHOD_JOB_LIST, |params, ctx, extensions| async move { + let principal = require_publish_principal(&extensions)?; + let params = if params.len_bytes() == 0 || params.as_str() == Some("[]") { + JobListParams { limit: None } + } else { + params + .parse::<JobListParams>() + .map_err(|error| RpcError::InvalidParams(error.to_string()))? + }; + if params.limit == Some(0) { + return Err(RpcError::InvalidParams( + "limit must be greater than zero".to_owned(), + )); + } + let configured_limit = ctx.state.transport_publish.config.job_list_limit; + let limit = params + .limit + .unwrap_or(configured_limit) + .min(configured_limit); + ctx.state + .transport_publish + .store + .list_jobs_for_principal(&principal, limit) + .map_err(|error| RpcError::Other(error.to_string())) + })?; + Ok(()) +} + +fn rpc_error_from_transport_publish(error: TransportPublishError) -> RpcError { + match error { + TransportPublishError::InvalidScope(message) => RpcError::Unauthorized(message), + TransportPublishError::InvalidSignedEvent(message) => RpcError::InvalidParams(message), + TransportPublishError::SignedEventVerification(_) + | TransportPublishError::Draft(_) + | TransportPublishError::Relay(_) => RpcError::InvalidParams(error.to_string()), + TransportPublishError::IdempotencyConflict(_) => RpcError::Other(error.to_string()), + other => RpcError::Other(other.to_string()), + } +} + +#[cfg(test)] +mod tests { + use super::module; + use std::sync::Arc; + + use crate::app::config::{Nip46Config, TransportPublishConfig, TransportPublishNostrConfig}; + use crate::core::Radrootsd; + use crate::core::transport_publish::{ + PublishJobVisibility, PublishPrincipalInit, generate_bearer_token, hash_bearer_token, + }; + use crate::transport::jsonrpc::auth::{ + TransportPublishAuthorization, authorize_transport_publish_request, + }; + use crate::transport::jsonrpc::{MethodRegistry, RpcContext}; + use jsonrpsee::server::RpcModule; + use nostr::JsonUtil; + use radroots_identity::RadrootsIdentity; + use radroots_nostr::prelude::{ + RadrootsNostrMetadata, RadrootsNostrTimestamp, radroots_nostr_build_event, + }; + use radroots_relay_transport::RadrootsMockRelayPublishAdapter; + use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, SignedNostrEventWire, TransportPublishTargetPolicyName, + }; + + fn signed_event(identity: &RadrootsIdentity) -> SignedNostrEventWire { + let event = radroots_nostr_build_event( + 30_402, + "{}", + vec![vec!["d".to_owned(), "listing-1".to_owned()]], + ) + .expect("event builder") + .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) + .sign_with_keys(identity.keys()) + .expect("signed event"); + serde_json::from_str(event.as_json().as_str()).expect("event wire") + } + + fn module_with_principal_and_config( + admin: bool, + transport_publish_config: TransportPublishConfig, + ) -> ( + RpcModule<RpcContext>, + RpcContext, + String, + SignedNostrEventWire, + ) { + let identity = RadrootsIdentity::generate(); + let signed_event = signed_event(&identity); + let metadata: RadrootsNostrMetadata = + serde_json::from_str(r#"{"name":"radrootsd-test"}"#).expect("metadata"); + let state = Radrootsd::new( + identity.clone(), + metadata, + transport_publish_config, + Nip46Config::default(), + ) + .expect("state"); + let mut state = state; + state.transport_publish = state + .transport_publish + .clone() + .with_publisher(Arc::new(RadrootsMockRelayPublishAdapter::new())); + let token = generate_bearer_token(); + let principal = state + .transport_publish + .store + .create_principal(PublishPrincipalInit { + label: "tester".to_owned(), + token_hash: hash_bearer_token(token.as_str()), + allowed_pubkeys: vec![identity.public_key_hex()], + allowed_kinds: vec![30_402], + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: vec![ + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, + job_visibility: if admin { + PublishJobVisibility::Admin + } else { + PublishJobVisibility::Own + }, + expires_at_unix: None, + }) + .expect("principal"); + let registry = MethodRegistry::default(); + let ctx = RpcContext::new(state, registry.clone()); + let mut module = module(ctx.clone(), registry).expect("module"); + module + .extensions_mut() + .insert(TransportPublishAuthorization::Authorized(principal)); + (module, ctx, token, signed_event) + } + + fn module_with_principal( + admin: bool, + ) -> ( + RpcModule<RpcContext>, + RpcContext, + String, + SignedNostrEventWire, + ) { + module_with_principal_and_config( + admin, + TransportPublishConfig { + nostr: TransportPublishNostrConfig { + daemon_default_relays: vec!["wss://relay.example.com".to_owned()], + ..TransportPublishNostrConfig::default() + }, + ..TransportPublishConfig::default() + }, + ) + } + + #[tokio::test] + async fn publish_event_records_job_and_deduplicates_idempotency() { + let (module, _ctx, _token, event) = module_with_principal(false); + let request = format!( + r#"{{ + "jsonrpc":"2.0", + "method":"transport.publish.event", + "params":{{ + "event":{}, + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, + "delivery_policy":{{"mode":"any"}}, + "idempotency_key":"idem-1" + }}, + "id":1 + }}"#, + serde_json::to_string(&event).expect("event json") + ); + let (response, _stream) = module + .raw_json_request(request.as_str(), 1) + .await + .expect("request"); + assert!(response.get().contains("\"deduplicated\":false")); + let (response, _stream) = module + .raw_json_request(request.as_str(), 1) + .await + .expect("request"); + assert!(response.get().contains("\"deduplicated\":true")); + } + + #[tokio::test] + async fn publish_event_rejects_principal_scope_gap() { + let (module, _ctx, _token, _pubkey) = module_with_principal(false); + let other_identity = RadrootsIdentity::generate(); + let event = signed_event(&other_identity); + let request = format!( + r#"{{ + "jsonrpc":"2.0", + "method":"transport.publish.event", + "params":{{ + "event":{}, + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, + "delivery_policy":{{"mode":"any"}} + }}, + "id":1 + }}"#, + serde_json::to_string(&event).expect("event json") + ); + let (response, _stream) = module + .raw_json_request(request.as_str(), 1) + .await + .expect("request"); + assert!(response.get().contains("unauthorized")); + } + + #[tokio::test] + async fn publish_job_list_rejects_malformed_and_zero_limits() { + let (module, _ctx, _token, _event) = module_with_principal(false); + let malformed = r#"{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "params":"bad", + "id":1 + }"#; + let (response, _stream) = module + .raw_json_request(malformed, 1) + .await + .expect("malformed request"); + assert!(response.get().contains("\"code\":-32602")); + + let zero = r#"{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "params":{"limit":0}, + "id":1 + }"#; + let (response, _stream) = module + .raw_json_request(zero, 1) + .await + .expect("zero request"); + assert!(response.get().contains("\"code\":-32602")); + assert!(response.get().contains("limit must be greater than zero")); + } + + #[tokio::test] + async fn publish_job_list_rejects_unknown_fields() { + let (module, _ctx, _token, _event) = module_with_principal(false); + for params in [ + r#"{"cursor":"next"}"#, + r#"{"status":"publishing"}"#, + r#"{"limit":1,"extra":true}"#, + ] { + let request = format!( + r#"{{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "params":{params}, + "id":1 + }}"# + ); + let (response, _stream) = module + .raw_json_request(request.as_str(), 1) + .await + .expect("unknown field request"); + assert!( + response.get().contains("\"code\":-32602"), + "{}", + response.get() + ); + } + } + + #[tokio::test] + async fn publish_job_list_uses_configured_limit_when_omitted_and_caps_positive_limits() { + let mut config = TransportPublishConfig { + nostr: TransportPublishNostrConfig { + daemon_default_relays: vec!["wss://relay.example.com".to_owned()], + ..TransportPublishNostrConfig::default() + }, + ..TransportPublishConfig::default() + }; + config.job_list_limit = 1; + let (module, _ctx, _token, event) = module_with_principal_and_config(false, config); + for idempotency_key in ["idem-list-1", "idem-list-2"] { + let request = format!( + r#"{{ + "jsonrpc":"2.0", + "method":"transport.publish.event", + "params":{{ + "event":{}, + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, + "delivery_policy":{{"mode":"any"}}, + "idempotency_key":"{idempotency_key}" + }}, + "id":1 + }}"#, + serde_json::to_string(&event).expect("event json") + ); + let (response, _stream) = module + .raw_json_request(request.as_str(), 1) + .await + .expect("publish request"); + assert!(response.get().contains("\"deduplicated\":false")); + } + + let omitted = r#"{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "id":1 + }"#; + let (response, _stream) = module + .raw_json_request(omitted, 1) + .await + .expect("omitted request"); + let value: serde_json::Value = + serde_json::from_str(response.get()).expect("omitted response json"); + assert_eq!(value["result"].as_array().expect("jobs").len(), 1); + + let empty_array = r#"{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "params":[], + "id":1 + }"#; + let (response, _stream) = module + .raw_json_request(empty_array, 1) + .await + .expect("empty array request"); + let value: serde_json::Value = + serde_json::from_str(response.get()).expect("empty array response json"); + assert_eq!(value["result"].as_array().expect("jobs").len(), 1); + + let over_limit = r#"{ + "jsonrpc":"2.0", + "method":"transport.publish.job.list", + "params":{"limit":50}, + "id":1 + }"#; + let (response, _stream) = module + .raw_json_request(over_limit, 1) + .await + .expect("over limit request"); + let value: serde_json::Value = + serde_json::from_str(response.get()).expect("over limit response json"); + assert_eq!(value["result"].as_array().expect("jobs").len(), 1); + } + + #[test] + fn http_auth_finds_principal_from_hashed_token() { + let (_module, ctx, token, _pubkey) = module_with_principal(false); + let header = format!("Bearer {token}"); + let auth = authorize_transport_publish_request( + Some(header.as_str()), + &ctx.state.transport_publish.store, + ); + assert!(matches!(auth, TransportPublishAuthorization::Authorized(_))); + } +} diff --git a/src/transport/jsonrpc/mod.rs b/src/transport/jsonrpc/mod.rs @@ -27,14 +27,14 @@ pub async fn start_rpc( addr: SocketAddr, rpc_cfg: &RpcConfig, ) -> Result<ServerHandle> { - state.publish_proxy.config.validate()?; + state.transport_publish.config.validate()?; let registry = MethodRegistry::default(); let ctx = RpcContext::new(state, registry.clone()); - let publish_proxy_store = ctx.state.publish_proxy.store.clone(); + let transport_publish_store = ctx.state.transport_publish.store.clone(); let mut root = RpcModule::new(ctx.clone()); methods::register_all(&mut root, ctx, registry)?; - let handle = server::start_server(addr, rpc_cfg, publish_proxy_store, root).await?; + let handle = server::start_server(addr, rpc_cfg, transport_publish_store, root).await?; Ok(handle) } diff --git a/src/transport/jsonrpc/server.rs b/src/transport/jsonrpc/server.rs @@ -13,7 +13,7 @@ use jsonrpsee::server::{ use jsonrpsee::types::{ErrorObject, Id}; use crate::app::config::RpcConfig; -use crate::core::publish_proxy::PublishProxyStore; +use crate::core::transport_publish::TransportPublishStore; use crate::transport::jsonrpc::RpcContext; use crate::transport::jsonrpc::auth; @@ -57,12 +57,12 @@ where ) -> impl Future<Output = Self::NotificationResponse> + Send + 'a { let service = self.service.clone(); async move { - if notification.method_name().starts_with("publish.") { + if notification.method_name().starts_with("transport.publish.") { MethodResponse::error( Id::Null, ErrorObject::owned( -32600, - "publish notifications are not accepted", + "transport publish notifications are not accepted", None::<()>, ), ) @@ -77,10 +77,11 @@ where mod tests { use super::start_server; use crate::app::config::{ - Nip46Config, PublishProxyConfig, PublishProxyRelayUrlPolicy, RpcConfig, + Nip46Config, NostrRelayUrlPolicy, RpcConfig, TransportPublishConfig, + TransportPublishNostrConfig, }; use crate::core::Radrootsd; - use crate::core::publish_proxy::{ + use crate::core::transport_publish::{ PublishJobVisibility, PublishPrincipalInit, PublishRelayResolveFuture, PublishRelayResolver, generate_bearer_token, hash_bearer_token, }; @@ -92,8 +93,10 @@ mod tests { use radroots_nostr::prelude::{ RadrootsNostrMetadata, RadrootsNostrTimestamp, radroots_nostr_build_event, }; - use radroots_publish_proxy_protocol::PublishRelayPolicy; use radroots_relay_transport::RadrootsMockRelayPublishAdapter; + use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, TransportPublishTargetPolicyName, + }; use serde_json::Value; use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpListener}; use std::sync::Arc; @@ -139,7 +142,7 @@ mod tests { } fn publish_server_state_with_config( - publish_proxy_config: PublishProxyConfig, + transport_publish_config: TransportPublishConfig, resolver: Option<Arc<dyn PublishRelayResolver>>, ) -> ( Radrootsd, @@ -153,27 +156,30 @@ mod tests { let mut state = Radrootsd::new( identity.clone(), metadata, - publish_proxy_config, + transport_publish_config, Nip46Config::default(), ) .expect("state"); let adapter = RadrootsMockRelayPublishAdapter::new(); - let mut publish_proxy = state.publish_proxy.clone(); + let mut transport_publish = state.transport_publish.clone(); if let Some(resolver) = resolver { - publish_proxy = publish_proxy.with_relay_resolver(resolver); + transport_publish = transport_publish.with_relay_resolver(resolver); } - state.publish_proxy = publish_proxy.with_publisher(Arc::new(adapter.clone())); + state.transport_publish = transport_publish.with_publisher(Arc::new(adapter.clone())); let token = generate_bearer_token(); state - .publish_proxy + .transport_publish .store .create_principal(PublishPrincipalInit { label: "tester".to_owned(), token_hash: hash_bearer_token(token.as_str()), allowed_pubkeys: vec![identity.public_key_hex()], allowed_kinds: vec![30_402], - allowed_relay_policies: vec![PublishRelayPolicy::DaemonDefaultOnly], - allow_request_relays: false, + allowed_target_policies: vec![TransportPublishTargetPolicyName::Nostr], + allowed_nostr_source_policies: vec![ + NostrPublishTargetSourcePolicy::DaemonDefaultOnly, + ], + allow_request_targets: false, job_visibility: PublishJobVisibility::Own, expires_at_unix: None, }) @@ -188,10 +194,13 @@ mod tests { RadrootsMockRelayPublishAdapter, ) { publish_server_state_with_config( - PublishProxyConfig { - daemon_default_publish_relays: vec![RELAY_PRIMARY.to_owned()], - relay_url_policy: PublishProxyRelayUrlPolicy::Localhost, - ..PublishProxyConfig::default() + TransportPublishConfig { + nostr: TransportPublishNostrConfig { + daemon_default_relays: vec![RELAY_PRIMARY.to_owned()], + relay_url_policy: NostrRelayUrlPolicy::Localhost, + ..TransportPublishNostrConfig::default() + }, + ..TransportPublishConfig::default() }, None, ) @@ -223,7 +232,7 @@ mod tests { rpc_cfg: RpcConfig, ) -> (SocketAddr, jsonrpsee::server::ServerHandle) { let addr = unused_addr(); - let store = state.publish_proxy.store.clone(); + let store = state.transport_publish.store.clone(); let registry = MethodRegistry::default(); let ctx = RpcContext::new(state, registry.clone()); let mut root = RpcModule::new(ctx.clone()); @@ -247,11 +256,10 @@ mod tests { let publish = format!( r#"{{ "jsonrpc":"2.0", - "method":"publish.event", + "method":"transport.publish.event", "params":{{ "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, "delivery_policy":{{"mode":"any"}}, "idempotency_key":"raw-http-idem" }}, @@ -274,7 +282,7 @@ mod tests { let get = format!( r#"{{ "jsonrpc":"2.0", - "method":"publish.job.get", + "method":"transport.publish.job.get", "params":{{"job_id":"{job_id}"}}, "id":2 }}"# @@ -286,7 +294,7 @@ mod tests { let list = r#"{ "jsonrpc":"2.0", - "method":"publish.job.list", + "method":"transport.publish.job.list", "params":{"limit":10}, "id":3 }"#; @@ -303,10 +311,13 @@ mod tests { #[tokio::test] async fn raw_http_publish_event_rejects_public_relay_forbidden_dns_destination() { let (state, token, identity, adapter) = publish_server_state_with_config( - PublishProxyConfig { - daemon_default_publish_relays: vec![RELAY_PUBLIC.to_owned()], - relay_url_policy: PublishProxyRelayUrlPolicy::Public, - ..PublishProxyConfig::default() + TransportPublishConfig { + nostr: TransportPublishNostrConfig { + daemon_default_relays: vec![RELAY_PUBLIC.to_owned()], + relay_url_policy: NostrRelayUrlPolicy::Public, + ..TransportPublishNostrConfig::default() + }, + ..TransportPublishConfig::default() }, Some(Arc::new(StaticPublishRelayResolver::forbidden_localhost())), ); @@ -315,11 +326,10 @@ mod tests { let publish = format!( r#"{{ "jsonrpc":"2.0", - "method":"publish.event", + "method":"transport.publish.event", "params":{{ "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, "delivery_policy":{{"mode":"any"}}, "idempotency_key":"raw-http-public-dns-reject" }}, @@ -333,30 +343,29 @@ mod tests { let publish_value = json_response_body(publish_response.as_str()); let job = &publish_value["result"]["job"]; assert_eq!(publish_value["result"]["deduplicated"], false); - assert_eq!(job["status"], "rejected"); - assert_eq!(job["last_error"], "no_publish_relays"); - let relays = job["relays"].as_array().expect("relay outcomes"); - assert_eq!(relays.len(), 1); - assert_eq!(relays[0]["relay_url"], RELAY_PUBLIC); - assert_eq!(relays[0]["source"], "daemon_default"); - assert_eq!(relays[0]["outcome_kind"], "relay_url_rejected"); - assert_eq!(relays[0]["attempted"], false); + assert_eq!(job["status"], "delivery_unsatisfied_terminal"); + assert_eq!(job["last_error"], "delivery_unsatisfied"); + let targets = job["targets"].as_array().expect("target outcomes"); + assert_eq!(targets.len(), 1); + assert_eq!(targets[0]["endpoint_uri"], RELAY_PUBLIC); + assert_eq!(targets[0]["source"], "daemon_default"); + assert_eq!(targets[0]["outcome_kind"], "target_rejected"); + assert_eq!(targets[0]["attempted"], false); assert!(adapter.captured_raw_events().is_empty()); } #[tokio::test] async fn publish_notifications_do_not_create_jobs() { let (state, token, identity, _adapter) = publish_server_state(); - let store = state.publish_proxy.store.clone(); + let store = state.transport_publish.store.clone(); let (addr, handle) = start_publish_server(state, RpcConfig::default()).await; let notification = format!( r#"{{ "jsonrpc":"2.0", - "method":"publish.event", + "method":"transport.publish.event", "params":{{ "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, "delivery_policy":{{"mode":"any"}} }} }}"#, @@ -366,7 +375,7 @@ mod tests { handle.stop().expect("stop server"); assert!( - response.contains("publish notifications are not accepted") + response.contains("transport publish notifications are not accepted") || response.ends_with("\r\n\r\n") ); let principal = store @@ -384,16 +393,15 @@ mod tests { #[tokio::test] async fn batch_requests_are_disabled_by_default() { let (state, token, identity, _adapter) = publish_server_state(); - let store = state.publish_proxy.store.clone(); + let store = state.transport_publish.store.clone(); let (addr, handle) = start_publish_server(state, RpcConfig::default()).await; let batch = format!( r#"[{{ "jsonrpc":"2.0", - "method":"publish.event", + "method":"transport.publish.event", "params":{{ "event":{}, - "relays":[], - "relay_policy":"daemon_default_only", + "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, "delivery_policy":{{"mode":"any"}} }}, "id":1 @@ -423,7 +431,7 @@ mod tests { pub async fn start_server( addr: SocketAddr, rpc_cfg: &RpcConfig, - publish_proxy_store: PublishProxyStore, + transport_publish_store: TransportPublishStore, root: RpcModule<RpcContext>, ) -> Result<ServerHandle> { let mut builder = ServerConfigBuilder::new() @@ -449,14 +457,14 @@ pub async fn start_server( .set_rpc_middleware(rpc_middleware) .set_http_middleware(tower::ServiceBuilder::new().map_request( move |mut request: HttpRequest<HttpBody>| { - let publish_proxy_auth = auth::authorize_publish_proxy_request( + let transport_publish_auth = auth::authorize_transport_publish_request( request .headers() .get("authorization") .and_then(|value| value.to_str().ok()), - &publish_proxy_store, + &transport_publish_store, ); - request.extensions_mut().insert(publish_proxy_auth); + request.extensions_mut().insert(transport_publish_auth); request }, ))