relay_publish.rs (9798B)
1 //! Daemon-owned adaptation from the V5 publish RPC to bounded Nostr attempts. 2 3 use core::future::Future; 4 use core::pin::Pin; 5 use core::time::Duration; 6 #[cfg(test)] 7 use std::collections::BTreeMap; 8 use std::net::IpAddr; 9 #[cfg(test)] 10 use std::sync::{Arc, Mutex}; 11 12 use nostr::JsonUtil; 13 use radroots_event::SignedEvent; 14 15 use crate::host_nostr::DaemonNostrClient; 16 17 pub(crate) type PublishFuture<'a> = 18 Pin<Box<dyn Future<Output = Result<Vec<RelayPublishReceipt>, RelayPublishError>> + Send + 'a>>; 19 20 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] 21 pub(crate) enum RelayUrlPolicy { 22 Public, 23 Localhost, 24 } 25 26 impl RelayUrlPolicy { 27 const fn native(self) -> radroots_transport_nostr::RelayUrlPolicy { 28 match self { 29 Self::Public => radroots_transport_nostr::RelayUrlPolicy::Public, 30 Self::Localhost => radroots_transport_nostr::RelayUrlPolicy::Local, 31 } 32 } 33 } 34 35 #[derive(Clone, Debug, Eq, PartialEq, Ord, PartialOrd)] 36 pub(crate) struct RelayUrl { 37 inner: radroots_transport_nostr::RelayUrl, 38 policy: RelayUrlPolicy, 39 } 40 41 impl RelayUrl { 42 pub(crate) fn parse( 43 value: impl AsRef<str>, 44 policy: RelayUrlPolicy, 45 ) -> Result<Self, RelayPublishError> { 46 let inner = radroots_transport_nostr::RelayUrl::parse(value, policy.native()) 47 .map_err(|error| RelayPublishError(error.to_string()))?; 48 Ok(Self { inner, policy }) 49 } 50 51 pub(crate) fn as_str(&self) -> &str { 52 self.inner.as_str() 53 } 54 55 pub(crate) fn validate_public_resolved_ip_addrs( 56 &self, 57 addresses: impl IntoIterator<Item = IpAddr>, 58 ) -> Result<(), RelayPublishError> { 59 self.inner 60 .validate_resolved_addresses(self.policy.native(), addresses) 61 .map_err(|error| RelayPublishError(error.to_string())) 62 } 63 } 64 65 #[derive(Clone, Debug)] 66 pub(crate) struct RelayTargetSet { 67 relays: Vec<RelayUrl>, 68 } 69 70 impl RelayTargetSet { 71 pub(crate) fn from_urls(mut relays: Vec<RelayUrl>) -> Result<Self, RelayPublishError> { 72 if relays.is_empty() { 73 return Err(RelayPublishError( 74 "relay target set must not be empty".to_owned(), 75 )); 76 } 77 relays.sort(); 78 relays.dedup(); 79 Ok(Self { relays }) 80 } 81 82 pub(crate) fn relays(&self) -> &[RelayUrl] { 83 self.relays.as_slice() 84 } 85 } 86 87 #[derive(Clone, Debug)] 88 pub(crate) struct RelayPublishRequest { 89 signed_event: SignedEvent, 90 targets: RelayTargetSet, 91 now_ms: i64, 92 } 93 94 impl RelayPublishRequest { 95 pub(crate) fn new(signed_event: SignedEvent, targets: RelayTargetSet, now_ms: i64) -> Self { 96 Self { 97 signed_event, 98 targets, 99 now_ms, 100 } 101 } 102 } 103 104 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 105 pub(crate) enum RelayOutcomeKind { 106 Accepted, 107 DuplicateAccepted, 108 Blocked, 109 RateLimited, 110 Invalid, 111 PowRequired, 112 Restricted, 113 AuthRequired, 114 Muted, 115 Unsupported, 116 PaymentRequired, 117 Error, 118 Timeout, 119 ConnectionFailed, 120 Unknown, 121 } 122 123 #[derive(Clone, Debug, Eq, PartialEq)] 124 pub(crate) struct RelayOutcome { 125 pub(crate) kind: RelayOutcomeKind, 126 pub(crate) message: Option<String>, 127 } 128 129 impl RelayOutcome { 130 pub(crate) const fn accepted() -> Self { 131 Self { 132 kind: RelayOutcomeKind::Accepted, 133 message: None, 134 } 135 } 136 137 pub(crate) fn classify(message: impl AsRef<str>) -> Self { 138 let message = message.as_ref().trim().to_owned(); 139 let lower = message.to_ascii_lowercase(); 140 let kind = if lower.starts_with("duplicate:") { 141 RelayOutcomeKind::DuplicateAccepted 142 } else if lower.starts_with("blocked:") { 143 RelayOutcomeKind::Blocked 144 } else if lower.starts_with("rate-limited:") || lower.contains("rate limit") { 145 RelayOutcomeKind::RateLimited 146 } else if lower.starts_with("invalid:") { 147 RelayOutcomeKind::Invalid 148 } else if lower.starts_with("pow:") { 149 RelayOutcomeKind::PowRequired 150 } else if lower.starts_with("restricted:") { 151 RelayOutcomeKind::Restricted 152 } else if lower.starts_with("auth-required:") || lower.contains("auth required") { 153 RelayOutcomeKind::AuthRequired 154 } else if lower.starts_with("mute:") { 155 RelayOutcomeKind::Muted 156 } else if lower.starts_with("unsupported:") { 157 RelayOutcomeKind::Unsupported 158 } else if lower.starts_with("payment-required:") { 159 RelayOutcomeKind::PaymentRequired 160 } else if lower.starts_with("timeout:") || lower.contains("timeout") { 161 RelayOutcomeKind::Timeout 162 } else if lower.starts_with("error:") { 163 RelayOutcomeKind::Error 164 } else { 165 RelayOutcomeKind::Unknown 166 }; 167 Self { 168 kind, 169 message: Some(message), 170 } 171 } 172 173 pub(crate) fn timeout(message: impl Into<String>) -> Self { 174 Self { 175 kind: RelayOutcomeKind::Timeout, 176 message: Some(message.into()), 177 } 178 } 179 180 pub(crate) fn connection_failed(message: impl Into<String>) -> Self { 181 Self { 182 kind: RelayOutcomeKind::ConnectionFailed, 183 message: Some(message.into()), 184 } 185 } 186 } 187 188 #[derive(Clone, Debug)] 189 pub(crate) struct RelayPublishReceipt { 190 pub(crate) relay_url: String, 191 pub(crate) outcome: RelayOutcome, 192 pub(crate) attempted: bool, 193 } 194 195 impl RelayPublishReceipt { 196 pub(crate) fn attempted(relay_url: impl Into<String>, outcome: RelayOutcome) -> Self { 197 Self { 198 relay_url: relay_url.into(), 199 outcome, 200 attempted: true, 201 } 202 } 203 } 204 205 #[derive(Debug, thiserror::Error)] 206 #[error("{0}")] 207 pub(crate) struct RelayPublishError(pub(crate) String); 208 209 pub(crate) trait RelayPublishAdapter: Send + Sync { 210 fn publish(&self, request: RelayPublishRequest) -> PublishFuture<'_>; 211 } 212 213 #[derive(Clone)] 214 pub(crate) struct LiveRelayPublishAdapter { 215 client: DaemonNostrClient, 216 } 217 218 impl LiveRelayPublishAdapter { 219 pub(crate) const fn new(client: DaemonNostrClient) -> Self { 220 Self { client } 221 } 222 } 223 224 impl RelayPublishAdapter for LiveRelayPublishAdapter { 225 fn publish(&self, request: RelayPublishRequest) -> PublishFuture<'_> { 226 Box::pin(async move { 227 if request.now_ms < 0 { 228 return Err(RelayPublishError( 229 "publish timestamp must be nonnegative".to_owned(), 230 )); 231 } 232 let event = nostr::Event::from_json(request.signed_event.raw_json()) 233 .map_err(|error| RelayPublishError(error.to_string()))?; 234 let client = self.client.clone().into_inner(); 235 let mut receipts = Vec::with_capacity(request.targets.relays().len()); 236 for relay in request.targets.relays() { 237 let relay_url = relay.as_str().to_owned(); 238 let attempt = async { 239 client.add_relay(relay_url.as_str()).await?; 240 client 241 .try_connect_relay(relay_url.as_str(), Duration::from_secs(10)) 242 .await?; 243 client.send_event_to([relay_url.as_str()], &event).await 244 }; 245 let outcome = match attempt.await { 246 Ok(output) 247 if output.success.iter().any(|url| { 248 url.to_string().trim_end_matches('/') == relay_url.trim_end_matches('/') 249 }) => 250 { 251 RelayOutcome::accepted() 252 } 253 Ok(output) => output 254 .failed 255 .values() 256 .next() 257 .map(RelayOutcome::classify) 258 .unwrap_or_else(|| RelayOutcome::connection_failed("relay omitted result")), 259 Err(error) => RelayOutcome::connection_failed(error.to_string()), 260 }; 261 receipts.push(RelayPublishReceipt::attempted(relay_url, outcome)); 262 } 263 Ok(receipts) 264 }) 265 } 266 } 267 268 #[cfg(test)] 269 #[derive(Clone, Default)] 270 pub(crate) struct MockRelayPublishAdapter { 271 outcomes: BTreeMap<String, RelayOutcome>, 272 captured_raw_events: Arc<Mutex<Vec<String>>>, 273 } 274 275 #[cfg(test)] 276 impl MockRelayPublishAdapter { 277 pub(crate) fn new() -> Self { 278 Self::default() 279 } 280 281 pub(crate) fn with_outcome( 282 mut self, 283 relay_url: impl Into<String>, 284 outcome: RelayOutcome, 285 ) -> Self { 286 self.outcomes.insert(relay_url.into(), outcome); 287 self 288 } 289 290 pub(crate) fn captured_raw_events(&self) -> Vec<String> { 291 self.captured_raw_events 292 .lock() 293 .unwrap_or_else(std::sync::PoisonError::into_inner) 294 .clone() 295 } 296 } 297 298 #[cfg(test)] 299 impl RelayPublishAdapter for MockRelayPublishAdapter { 300 fn publish(&self, request: RelayPublishRequest) -> PublishFuture<'_> { 301 Box::pin(async move { 302 self.captured_raw_events 303 .lock() 304 .map_err(|_| RelayPublishError("captured event lock poisoned".to_owned()))? 305 .push(request.signed_event.raw_json().to_owned()); 306 Ok(request 307 .targets 308 .relays() 309 .iter() 310 .map(|relay| { 311 RelayPublishReceipt::attempted( 312 relay.as_str(), 313 self.outcomes 314 .get(relay.as_str()) 315 .cloned() 316 .unwrap_or_else(RelayOutcome::accepted), 317 ) 318 }) 319 .collect()) 320 }) 321 } 322 }