radrootsd

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

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 }