lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

exact_delivery.rs (4981B)


      1 //! Exact retained EVENT framing and connection-local acknowledgement handling.
      2 
      3 use crate::{relay::MAX_WIRE_MESSAGE_BYTES, socket_write::WriterRegistry, status};
      4 use nostr_relay_pool::RelayNotification;
      5 use nostr_sdk::{EventId, RelayMessage};
      6 use radroots_transport::{DeliveryRequest, outcome::DeliveryOutcome};
      7 use std::sync::Arc;
      8 
      9 #[derive(Clone)]
     10 pub(crate) struct ExactEvent {
     11     id: EventId,
     12     frame: Arc<str>,
     13 }
     14 
     15 impl ExactEvent {
     16     pub(crate) fn from_request(request: &DeliveryRequest) -> Option<Self> {
     17         let signed = request.payload().event();
     18         let raw = signed.raw_json();
     19         let event = radroots_nostr::event::to_nostr(signed.envelope()).ok()?;
     20         Some(Self {
     21             id: event.id,
     22             frame: frame(raw)?,
     23         })
     24     }
     25 
     26     pub(crate) async fn publish(
     27         &self,
     28         client: &nostr_sdk::Client,
     29         writers: &WriterRegistry,
     30         url: &str,
     31     ) -> Result<DeliveryOutcome, String> {
     32         let relay = client.relay(url).await.map_err(|error| error.to_string())?;
     33         // Subscribe before sending on the same connection; no success from enqueue alone.
     34         let mut notifications = relay.notifications();
     35         let writer = writers.get(url).map_err(|error| error.to_string())?;
     36         writer
     37             .send(async_wsocket::Message::Text(self.frame.to_string()))
     38             .await
     39             .map_err(|error| error.to_string())?;
     40         loop {
     41             let notification = notifications
     42                 .recv()
     43                 .await
     44                 .map_err(|_| "relay acknowledgement unavailable".to_owned())?;
     45             if let Some(outcome) = acknowledgement(self.id, notification) {
     46                 return Ok(outcome);
     47             }
     48         }
     49     }
     50 }
     51 
     52 fn frame(raw: &str) -> Option<Arc<str>> {
     53     // Bound before allocation; preserve whitespace, field order and extensions.
     54     (raw.len() <= MAX_WIRE_MESSAGE_BYTES - "[\"EVENT\",]".len())
     55         .then(|| format!("[\"EVENT\",{raw}]").into())
     56 }
     57 
     58 fn acknowledgement(id: EventId, notification: RelayNotification) -> Option<DeliveryOutcome> {
     59     match notification {
     60         RelayNotification::Message {
     61             message:
     62                 RelayMessage::Ok {
     63                     event_id,
     64                     status: accepted,
     65                     message,
     66                 },
     67         } if event_id == id => Some(if accepted {
     68             DeliveryOutcome::accepted()
     69         } else {
     70             status::delivery_failure(&message)
     71         }),
     72         RelayNotification::RelayStatus {
     73             status:
     74                 nostr_sdk::RelayStatus::Disconnected
     75                 | nostr_sdk::RelayStatus::Terminated
     76                 | nostr_sdk::RelayStatus::Banned,
     77         } => Some(status::delivery_failure("relay disconnected")),
     78         RelayNotification::Shutdown => Some(status::delivery_failure("relay shutdown")),
     79         _ => None,
     80     }
     81 }
     82 
     83 #[cfg(test)]
     84 mod tests {
     85     use super::*;
     86 
     87     #[test]
     88     fn framing_bounds_complete_message_bytes_before_allocation() {
     89         let limit = MAX_WIRE_MESSAGE_BYTES - "[\"EVENT\",]".len();
     90         let exact = " ".repeat(limit);
     91         assert_eq!(frame(&exact).unwrap().len(), MAX_WIRE_MESSAGE_BYTES);
     92         assert!(frame(&(exact + " ")).is_none());
     93         let unicode = "🌱".repeat(limit / 4 + 1);
     94         assert!(frame(&unicode).is_none());
     95         assert_eq!(
     96             &*frame(" \n{\"extension\":true} \t").unwrap(),
     97             "[\"EVENT\", \n{\"extension\":true} \t]"
     98         );
     99     }
    100 
    101     #[test]
    102     fn only_matching_ok_can_accept_and_disconnects_remain_failures() {
    103         let id = EventId::from_byte_array([1; 32]);
    104         let other = EventId::from_byte_array([2; 32]);
    105         let ok = |event_id, status, message: &'static str| RelayNotification::Message {
    106             message: RelayMessage::Ok {
    107                 event_id,
    108                 status,
    109                 message: message.into(),
    110             },
    111         };
    112         assert!(acknowledgement(id, ok(other, true, "")).is_none());
    113         assert_eq!(
    114             acknowledgement(id, ok(id, true, "")).unwrap(),
    115             DeliveryOutcome::accepted()
    116         );
    117         assert!(!status::delivery_succeeded(
    118             &acknowledgement(id, ok(id, false, "blocked: denied")).unwrap()
    119         ));
    120         for status in [
    121             nostr_sdk::RelayStatus::Disconnected,
    122             nostr_sdk::RelayStatus::Terminated,
    123             nostr_sdk::RelayStatus::Banned,
    124         ] {
    125             assert!(!super::status::delivery_succeeded(
    126                 &acknowledgement(id, RelayNotification::RelayStatus { status }).unwrap()
    127             ));
    128         }
    129         assert!(
    130             acknowledgement(
    131                 id,
    132                 RelayNotification::RelayStatus {
    133                     status: nostr_sdk::RelayStatus::Connected
    134                 }
    135             )
    136             .is_none()
    137         );
    138         assert!(!status::delivery_succeeded(
    139             &acknowledgement(id, RelayNotification::Shutdown).unwrap()
    140         ));
    141     }
    142 }