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 }