lib

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

sink.rs (30071B)


      1 //! Nostr implementation of the transport event sink.
      2 
      3 use crate::exact_delivery::ExactEvent;
      4 use crate::{NostrTransport, RelayUrl, status};
      5 use core::{fmt, time::Duration};
      6 use futures::{StreamExt, stream};
      7 use radroots_transport::{
      8     BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, Target, TargetSet,
      9     outcome::DeliveryOutcome,
     10     sink::{DeliveryTargetReceipt, SinkStatus},
     11 };
     12 use std::collections::{BTreeMap, BTreeSet};
     13 
     14 #[derive(Clone, Debug)]
     15 pub(crate) struct RelayPublishResult {
     16     relay: RelayUrl,
     17     attempted: bool,
     18     outcome: DeliveryOutcome,
     19 }
     20 
     21 /// Sealed, no-I/O result of validating one delivery against this adapter.
     22 ///
     23 /// The value retains the exact request and signed event bytes. It is
     24 /// constructed only by [`NostrTransport::prepare_delivery`] and is consumed by
     25 /// [`NostrTransport::execute_prepared_delivery`]. Ordinary `Debug` never
     26 /// exposes event bytes, request identities, or relay destinations.
     27 #[must_use = "prepared delivery must be durably bound before execution or deliberately discarded"]
     28 pub struct PreparedDelivery {
     29     request: DeliveryRequest,
     30     config: crate::Config,
     31     event: ExactEvent,
     32     authorized: Vec<(RelayUrl, Target)>,
     33     skipped: Vec<DeliveryTargetReceipt>,
     34 }
     35 
     36 impl PreparedDelivery {
     37     /// Returns the exact validated request retained for persistence binding.
     38     #[must_use]
     39     pub const fn request(&self) -> &DeliveryRequest {
     40         &self.request
     41     }
     42 }
     43 
     44 impl fmt::Debug for PreparedDelivery {
     45     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     46         formatter.write_str("PreparedDelivery([redacted])")
     47     }
     48 }
     49 
     50 pub(crate) trait RelayClient: Send + Sync {
     51     fn publish<'a>(
     52         &'a self,
     53         relays: Vec<RelayUrl>,
     54         event: ExactEvent,
     55         max_connections: usize,
     56         connect_timeout: Duration,
     57         operation_timeout: Duration,
     58     ) -> BoxFuture<'a, Vec<RelayPublishResult>>;
     59 }
     60 
     61 #[derive(Clone, Debug)]
     62 pub(crate) struct LiveRelayClient {
     63     client: nostr_sdk::Client,
     64     writers: crate::socket_write::WriterRegistry,
     65 }
     66 
     67 impl LiveRelayClient {
     68     pub(crate) const fn new(
     69         client: nostr_sdk::Client,
     70         writers: crate::socket_write::WriterRegistry,
     71     ) -> Self {
     72         Self { client, writers }
     73     }
     74 
     75     #[cfg(test)]
     76     pub(crate) fn isolated() -> Self {
     77         let client = nostr_sdk::Client::default();
     78         client.automatic_authentication(false);
     79         Self::new(
     80             client,
     81             crate::socket_write::WriterRegistry::new(std::iter::empty()),
     82         )
     83     }
     84 }
     85 
     86 impl RelayClient for LiveRelayClient {
     87     // Socket scheduling is exercised by the real loopback suite. Deterministic
     88     // coverage owns framing, writer lifecycle and acknowledgement normalization.
     89     #[cfg_attr(coverage_nightly, coverage(off))]
     90     fn publish<'a>(
     91         &'a self,
     92         relays: Vec<RelayUrl>,
     93         event: ExactEvent,
     94         max_connections: usize,
     95         connect_timeout: Duration,
     96         operation_timeout: Duration,
     97     ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
     98         Box::pin(async move {
     99             let deadline = tokio::time::Instant::now() + operation_timeout;
    100             stream::iter(relays.into_iter().map(|relay| {
    101                 let event = event.clone();
    102                 async move {
    103                     let url = relay.as_str().to_owned();
    104                     let mut attempted = false;
    105                     let attempt = async {
    106                         if tokio::time::Instant::now() >= deadline {
    107                             return Err("timeout".to_owned());
    108                         }
    109                         attempted = true;
    110                         self.client
    111                             .add_relay(url.as_str())
    112                             .await
    113                             .map_err(|error| error.to_string())?;
    114                         self.client
    115                             .try_connect_relay(url.as_str(), connect_timeout)
    116                             .await
    117                             .map_err(|error| error.to_string())?;
    118                         event.publish(&self.client, &self.writers, &url).await
    119                     };
    120                     let outcome = match tokio::time::timeout_at(deadline, attempt).await {
    121                         Err(_) => status::delivery_failure("timeout"),
    122                         Ok(Err(error)) => status::delivery_failure(&error),
    123                         Ok(Ok(outcome)) => outcome,
    124                     };
    125                     RelayPublishResult {
    126                         relay,
    127                         attempted,
    128                         outcome,
    129                     }
    130                 }
    131             }))
    132             .buffered(max_connections)
    133             .collect()
    134             .await
    135         })
    136     }
    137 }
    138 
    139 impl NostrTransport {
    140     /// Validates and converts one delivery without reading a clock or performing relay I/O.
    141     pub fn prepare_delivery(
    142         &self,
    143         request: DeliveryRequest,
    144     ) -> Result<PreparedDelivery, Box<SinkFailure>> {
    145         self.prepare_delivery_inner(request, None)
    146     }
    147 
    148     /// Prepares only an exact subset, retaining the complete original request.
    149     /// Unselected targets remain explicit unattempted, retryable receipt rows.
    150     pub fn prepare_delivery_selected(
    151         &self,
    152         request: DeliveryRequest,
    153         selected: &TargetSet,
    154     ) -> Result<PreparedDelivery, Box<SinkFailure>> {
    155         request
    156             .validate_target_selection(selected)
    157             .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?;
    158         self.prepare_delivery_inner(request, Some(selected))
    159     }
    160 
    161     fn prepare_delivery_inner(
    162         &self,
    163         request: DeliveryRequest,
    164         selected: Option<&TargetSet>,
    165     ) -> Result<PreparedDelivery, Box<SinkFailure>> {
    166         let mut authorized = Vec::new();
    167         let mut skipped = Vec::new();
    168         for target in request.target_set().targets() {
    169             if selected.is_some_and(|targets| !targets.targets().contains(target)) {
    170                 skipped.push(
    171                     DeliveryTargetReceipt::skipped(
    172                         target.clone(),
    173                         DeliveryOutcome::unavailable()
    174                             .with_detail(
    175                                 "target_not_selected",
    176                                 "target is held for a later attempt",
    177                             )
    178                             .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?,
    179                     )
    180                     .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?,
    181                 );
    182                 continue;
    183             }
    184             match self.config().endpoint_for_target(target) {
    185                 Some(endpoint) if endpoint.access().can_write() => {
    186                     authorized.push((endpoint.url().clone(), target.clone()));
    187                 }
    188                 None | Some(_) => skipped.push(
    189                     DeliveryTargetReceipt::skipped(
    190                         target.clone(),
    191                         DeliveryOutcome::rejected()
    192                             .with_detail("target_denied", "target is not configured for this sink")
    193                             .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?,
    194                     )
    195                     .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?,
    196                 ),
    197             }
    198         }
    199         let event = ExactEvent::from_request(&request)
    200             .ok_or_else(|| Box::new(SinkFailure::invalid_contract(&request)))?;
    201         Ok(PreparedDelivery {
    202             request,
    203             config: self.config().clone(),
    204             event,
    205             authorized,
    206             skipped,
    207         })
    208     }
    209 
    210     /// Performs relay I/O for one exact prepared delivery and consumes its authority.
    211     pub fn execute_prepared_delivery(
    212         &self,
    213         prepared: PreparedDelivery,
    214     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    215         Box::pin(async move {
    216             let PreparedDelivery {
    217                 request,
    218                 config,
    219                 event,
    220                 authorized,
    221                 mut skipped,
    222             } = prepared;
    223             if config != *self.config() {
    224                 return Err(SinkFailure::invalid_contract(&request));
    225             }
    226             let now_unix_ms = unix_time_ms();
    227             let mut requested = Vec::new();
    228             for (relay, target) in authorized {
    229                 if self.status.may_write(&relay, now_unix_ms) {
    230                     requested.push((relay, target));
    231                 } else {
    232                     skipped.push(
    233                         DeliveryTargetReceipt::skipped(
    234                             target,
    235                             DeliveryOutcome::unavailable()
    236                                 .with_detail(
    237                                     "reconnect_backoff",
    238                                     "relay reconnect backoff is active",
    239                                 )
    240                                 .map_err(|_| SinkFailure::invalid_contract(&request))?,
    241                         )
    242                         .map_err(|_| SinkFailure::invalid_contract(&request))?,
    243                     );
    244                 }
    245             }
    246             let remaining_ms = request.deadline_unix_ms().saturating_sub(now_unix_ms);
    247             let operation_timeout_ms = remaining_ms.min(self.config().request_timeout_ms());
    248             if operation_timeout_ms == 0 {
    249                 for (relay, _) in &requested {
    250                     self.status.record_write(relay, false, true, now_unix_ms);
    251                 }
    252                 let timeout = status::delivery_failure("timeout");
    253                 skipped.extend(requested.into_iter().map(|(_, target)| {
    254                     DeliveryTargetReceipt::skipped(target, timeout.clone())
    255                         .expect("normalized timeout cannot satisfy delivery")
    256                 }));
    257                 return DeliveryReceipt::for_request(&request, skipped)
    258                     .map_err(|_| SinkFailure::invalid_contract(&request));
    259             }
    260             let expected: BTreeSet<_> = requested.iter().map(|(relay, _)| relay.clone()).collect();
    261             for (relay, _) in &requested {
    262                 self.status.begin_write(relay, now_unix_ms);
    263             }
    264             let results = self
    265                 .client
    266                 .publish(
    267                     requested.iter().map(|(relay, _)| relay.clone()).collect(),
    268                     event,
    269                     self.config().max_connections(),
    270                     Duration::from_millis(self.config().connect_timeout_ms()),
    271                     Duration::from_millis(operation_timeout_ms),
    272                 )
    273                 .await;
    274             let mut by_relay = BTreeMap::new();
    275             let observed_at_unix_ms = unix_time_ms().max(now_unix_ms);
    276             for result in results {
    277                 if !expected.contains(&result.relay)
    278                     || by_relay.contains_key(&result.relay)
    279                     || (!result.attempted && status::delivery_succeeded(&result.outcome))
    280                 {
    281                     return Err(SinkFailure::invalid_contract(&request));
    282                 }
    283                 let succeeded = status::delivery_succeeded(&result.outcome);
    284                 self.status.record_write(
    285                     &result.relay,
    286                     succeeded,
    287                     result.outcome.is_retryable(),
    288                     observed_at_unix_ms,
    289                 );
    290                 by_relay.insert(result.relay, (result.outcome, result.attempted));
    291             }
    292 
    293             let mut receipts = skipped;
    294             for (relay, target) in requested {
    295                 let (outcome, attempted) = by_relay.remove(&relay).unwrap_or_else(|| {
    296                     self.status
    297                         .record_write(&relay, false, true, observed_at_unix_ms);
    298                     (
    299                         DeliveryOutcome::unavailable()
    300                             .with_detail("missing_result", "relay returned no result")
    301                             .expect("static normalized outcome"),
    302                         true,
    303                     )
    304                 });
    305                 receipts.push(if attempted {
    306                     DeliveryTargetReceipt::attempted(target, outcome)
    307                 } else {
    308                     DeliveryTargetReceipt::skipped(target, outcome)
    309                         .map_err(|_| SinkFailure::invalid_contract(&request))?
    310                 });
    311             }
    312             DeliveryReceipt::for_request(&request, receipts)
    313                 .map_err(|_| SinkFailure::invalid_contract(&request))
    314         })
    315     }
    316 }
    317 
    318 impl EventSink for NostrTransport {
    319     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
    320         Box::pin(async move { Ok(status::sink_status(&self.status)) })
    321     }
    322 
    323     fn deliver(
    324         &self,
    325         request: DeliveryRequest,
    326     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    327         Box::pin(async move {
    328             let prepared = self.prepare_delivery(request).map_err(|failure| *failure)?;
    329             self.execute_prepared_delivery(prepared).await
    330         })
    331     }
    332 
    333     fn deliver_selected(
    334         &self,
    335         request: DeliveryRequest,
    336         selected: TargetSet,
    337     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    338         Box::pin(async move {
    339             let prepared = self
    340                 .prepare_delivery_selected(request, &selected)
    341                 .map_err(|failure| *failure)?;
    342             self.execute_prepared_delivery(prepared).await
    343         })
    344     }
    345 }
    346 
    347 #[cfg_attr(coverage_nightly, coverage(off))]
    348 fn unix_time_ms() -> u64 {
    349     std::time::SystemTime::now()
    350         .duration_since(std::time::UNIX_EPOCH)
    351         .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
    352         .unwrap_or_default()
    353 }
    354 
    355 #[cfg(test)]
    356 mod tests {
    357     use super::*;
    358     use crate::{Config, RelayUrlPolicy};
    359     use radroots_transport::{
    360         Target, TargetSet,
    361         outcome::DeliveryOutcomeKind,
    362         policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
    363         sink::DeliveryPayload,
    364     };
    365     use std::sync::{
    366         Arc,
    367         atomic::{AtomicUsize, Ordering},
    368     };
    369 
    370     #[derive(Debug)]
    371     struct MockRelayClient {
    372         outcomes: BTreeMap<RelayUrl, DeliveryOutcome>,
    373     }
    374 
    375     impl RelayClient for MockRelayClient {
    376         fn publish<'a>(
    377             &'a self,
    378             relays: Vec<RelayUrl>,
    379             _event: ExactEvent,
    380             _max_connections: usize,
    381             _connect_timeout: Duration,
    382             _operation_timeout: Duration,
    383         ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
    384             Box::pin(async move {
    385                 relays
    386                     .into_iter()
    387                     .map(|relay| RelayPublishResult {
    388                         attempted: true,
    389                         outcome: self
    390                             .outcomes
    391                             .get(&relay)
    392                             .cloned()
    393                             .unwrap_or_else(DeliveryOutcome::accepted),
    394                         relay,
    395                     })
    396                     .collect()
    397             })
    398         }
    399     }
    400 
    401     #[derive(Debug)]
    402     struct CountingRelayClient(Arc<AtomicUsize>);
    403 
    404     impl RelayClient for CountingRelayClient {
    405         fn publish<'a>(
    406             &'a self,
    407             _relays: Vec<RelayUrl>,
    408             _event: ExactEvent,
    409             _max_connections: usize,
    410             _connect_timeout: Duration,
    411             _operation_timeout: Duration,
    412         ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
    413             self.0.fetch_add(1, Ordering::SeqCst);
    414             Box::pin(async { Vec::new() })
    415         }
    416     }
    417 
    418     fn payload() -> DeliveryPayload {
    419         let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
    420         DeliveryPayload::new(radroots_event_codec::decode::signed_event(raw).expect("signed event"))
    421     }
    422 
    423     fn request() -> DeliveryRequest {
    424         request_with_deadline(1_800_000_000_000)
    425     }
    426 
    427     fn request_with_deadline(deadline_unix_ms: u64) -> DeliveryRequest {
    428         DeliveryRequest::new(
    429             "nostr-delivery",
    430             payload(),
    431             TargetSet::new(vec![
    432                 Target::nostr_relay("wss://one.example").expect("one"),
    433                 Target::nostr_relay("wss://two.example").expect("two"),
    434             ])
    435             .expect("targets"),
    436             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    437             deadline_unix_ms,
    438         )
    439         .expect("request")
    440     }
    441 
    442     #[test]
    443     fn sink_returns_normalized_per_relay_partial_success() {
    444         let config = Config::from_profile(
    445             crate::profile::test_profile(
    446                 crate::RelayProfileKind::Public,
    447                 RelayUrlPolicy::Public,
    448                 ["wss://one.example", "wss://two.example"],
    449             )
    450             .expect("profile"),
    451         );
    452         let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).expect("two");
    453         let client = MockRelayClient {
    454             outcomes: BTreeMap::from([(two, status::delivery_failure("rate limited"))]),
    455         };
    456         let transport = NostrTransport::with_client(config, Arc::new(client));
    457         let request = request();
    458         let receipt = futures::executor::block_on(transport.deliver(request.clone()))
    459             .expect("delivery receipt");
    460 
    461         assert_eq!(receipt.target_receipts().len(), 2);
    462         assert_eq!(
    463             receipt.target_receipts()[0].outcome().kind(),
    464             DeliveryOutcomeKind::Accepted
    465         );
    466         assert_eq!(
    467             receipt.target_receipts()[1].outcome().kind(),
    468             DeliveryOutcomeKind::Unavailable
    469         );
    470         assert!(!receipt.is_satisfied(&request).expect("satisfaction"));
    471     }
    472 
    473     #[test]
    474     fn selected_preparation_retains_binding_and_only_authorizes_selected_relays() {
    475         let config = Config::from_profile(
    476             crate::profile::test_profile(
    477                 crate::RelayProfileKind::Public,
    478                 RelayUrlPolicy::Public,
    479                 ["wss://one.example", "wss://two.example"],
    480             )
    481             .unwrap(),
    482         );
    483         let transport = NostrTransport::with_client(
    484             config,
    485             Arc::new(MockRelayClient {
    486                 outcomes: BTreeMap::new(),
    487             }),
    488         );
    489         let request = request();
    490         let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap();
    491         let prepared = transport
    492             .prepare_delivery_selected(request.clone(), &selected)
    493             .unwrap();
    494         assert_eq!(prepared.request(), &request);
    495         assert_eq!(prepared.authorized.len(), 1);
    496         assert_eq!(&prepared.authorized[0].1, &selected.targets()[0]);
    497         let receipt =
    498             futures::executor::block_on(transport.execute_prepared_delivery(prepared)).unwrap();
    499         receipt.validate_for_request(&request).unwrap();
    500         assert!(receipt.target_receipts()[0].was_attempted());
    501         assert_eq!(
    502             receipt.target_receipts()[0].outcome().kind(),
    503             DeliveryOutcomeKind::Accepted
    504         );
    505         assert!(!receipt.target_receipts()[1].was_attempted());
    506         assert_eq!(
    507             receipt.target_receipts()[1].outcome().code(),
    508             Some("target_not_selected")
    509         );
    510         assert!(receipt.target_receipts()[1].outcome().is_retryable());
    511         assert!(!receipt.is_satisfied(&request).unwrap());
    512         let foreign =
    513             TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap();
    514         let failure = transport
    515             .prepare_delivery_selected(request.clone(), &foreign)
    516             .unwrap_err();
    517         failure.validate_for_request(&request).unwrap();
    518         assert_eq!(failure.code(), "invalid_transport_contract");
    519         let all = futures::executor::block_on(
    520             transport.deliver_selected(request.clone(), request.target_set().clone()),
    521         )
    522         .unwrap();
    523         assert!(all.is_satisfied(&request).unwrap());
    524     }
    525 
    526     #[test]
    527     fn upstream_messages_map_to_stable_outcomes() {
    528         let cases = [
    529             ("duplicate: already have", DeliveryOutcomeKind::Accepted),
    530             ("blocked by policy", DeliveryOutcomeKind::Rejected),
    531             ("rate limited", DeliveryOutcomeKind::Unavailable),
    532             ("auth required", DeliveryOutcomeKind::Failed),
    533             ("connection timeout", DeliveryOutcomeKind::Unavailable),
    534             ("connection failed", DeliveryOutcomeKind::Unavailable),
    535         ];
    536         for (message, expected) in cases {
    537             assert_eq!(status::delivery_failure(message).kind(), expected);
    538         }
    539     }
    540 
    541     #[test]
    542     fn dropping_an_unpolled_delivery_performs_no_relay_work() {
    543         let calls = Arc::new(AtomicUsize::new(0));
    544         let config = Config::from_profile(
    545             crate::profile::test_profile(
    546                 crate::RelayProfileKind::Public,
    547                 RelayUrlPolicy::Public,
    548                 ["wss://one.example", "wss://two.example"],
    549             )
    550             .expect("profile"),
    551         );
    552         let transport =
    553             NostrTransport::with_client(config, Arc::new(CountingRelayClient(Arc::clone(&calls))));
    554         let delivery = transport.deliver(request());
    555         drop(delivery);
    556         assert_eq!(calls.load(Ordering::SeqCst), 0);
    557     }
    558 
    559     #[test]
    560     fn preparation_is_no_io_redacted_consuming_and_bound_to_exact_config() {
    561         let calls = Arc::new(AtomicUsize::new(0));
    562         let profile = crate::profile::test_profile(
    563             crate::RelayProfileKind::Public,
    564             RelayUrlPolicy::Public,
    565             ["wss://one.example", "wss://two.example"],
    566         )
    567         .expect("profile");
    568         let config = Config::from_profile(profile.clone());
    569         let transport =
    570             NostrTransport::with_client(config, Arc::new(CountingRelayClient(Arc::clone(&calls))));
    571         let request = request();
    572 
    573         let prepared = transport
    574             .prepare_delivery(request.clone())
    575             .expect("prepared delivery");
    576         assert_eq!(prepared.request(), &request);
    577         assert_eq!(format!("{prepared:?}"), "PreparedDelivery([redacted])");
    578         assert_eq!(calls.load(Ordering::SeqCst), 0);
    579 
    580         let receipt = futures::executor::block_on(transport.execute_prepared_delivery(prepared))
    581             .expect("executed delivery");
    582         assert_eq!(receipt.request_id(), request.request_id());
    583         assert_eq!(calls.load(Ordering::SeqCst), 1);
    584 
    585         let mismatch_calls = Arc::new(AtomicUsize::new(0));
    586         let mismatched = NostrTransport::with_client(
    587             Config::from_profile(profile)
    588                 .with_timeouts(5_000, 20_000, 2_000)
    589                 .expect("different bounded config"),
    590             Arc::new(CountingRelayClient(Arc::clone(&mismatch_calls))),
    591         );
    592         let prepared = transport
    593             .prepare_delivery(request)
    594             .expect("second prepared delivery");
    595         assert_eq!(
    596             futures::executor::block_on(mismatched.execute_prepared_delivery(prepared))
    597                 .expect_err("prepared authority is config-bound")
    598                 .code(),
    599             "invalid_transport_contract"
    600         );
    601         assert_eq!(mismatch_calls.load(Ordering::SeqCst), 0);
    602     }
    603 
    604     #[test]
    605     fn expired_delivery_deadline_performs_no_relay_work() {
    606         let calls = Arc::new(AtomicUsize::new(0));
    607         let config = Config::from_profile(
    608             crate::profile::test_profile(
    609                 crate::RelayProfileKind::Public,
    610                 RelayUrlPolicy::Public,
    611                 ["wss://one.example", "wss://two.example"],
    612             )
    613             .expect("profile"),
    614         );
    615         let transport =
    616             NostrTransport::with_client(config, Arc::new(CountingRelayClient(Arc::clone(&calls))));
    617         let receipt = futures::executor::block_on(transport.deliver(request_with_deadline(1)))
    618             .expect("bounded timeout receipt");
    619         assert_eq!(calls.load(Ordering::SeqCst), 0);
    620         assert!(
    621             receipt
    622                 .target_receipts()
    623                 .iter()
    624                 .all(|target| !target.was_attempted())
    625         );
    626     }
    627 
    628     #[derive(Debug)]
    629     struct ScriptedRelayClient(Vec<RelayPublishResult>);
    630 
    631     impl RelayClient for ScriptedRelayClient {
    632         fn publish<'a>(
    633             &'a self,
    634             _relays: Vec<RelayUrl>,
    635             _event: ExactEvent,
    636             _max_connections: usize,
    637             _connect_timeout: Duration,
    638             _operation_timeout: Duration,
    639         ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
    640             Box::pin(async move { self.0.clone() })
    641         }
    642     }
    643 
    644     fn scripted(results: Vec<RelayPublishResult>) -> NostrTransport {
    645         let config = Config::from_profile(
    646             crate::profile::test_profile(
    647                 crate::RelayProfileKind::Public,
    648                 RelayUrlPolicy::Public,
    649                 ["wss://one.example", "wss://two.example"],
    650             )
    651             .expect("profile"),
    652         );
    653         NostrTransport::with_client(config, Arc::new(ScriptedRelayClient(results)))
    654     }
    655 
    656     #[test]
    657     fn sink_handles_missing_duplicate_unexpected_and_denied_targets() {
    658         let one = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("one");
    659         let missing = futures::executor::block_on(
    660             scripted(vec![RelayPublishResult {
    661                 relay: one.clone(),
    662                 attempted: true,
    663                 outcome: DeliveryOutcome::accepted(),
    664             }])
    665             .deliver(request()),
    666         )
    667         .expect("missing result receipt");
    668         assert_eq!(missing.target_receipts().len(), 2);
    669 
    670         let duplicate = scripted(vec![
    671             RelayPublishResult {
    672                 relay: one.clone(),
    673                 attempted: true,
    674                 outcome: DeliveryOutcome::accepted(),
    675             },
    676             RelayPublishResult {
    677                 relay: one,
    678                 attempted: true,
    679                 outcome: DeliveryOutcome::accepted(),
    680             },
    681         ]);
    682         assert_eq!(
    683             futures::executor::block_on(duplicate.deliver(request()))
    684                 .expect_err("duplicate relay evidence")
    685                 .code(),
    686             "invalid_transport_contract"
    687         );
    688 
    689         let other = RelayUrl::parse("wss://other.example", RelayUrlPolicy::Public).expect("other");
    690         let unexpected = scripted(vec![RelayPublishResult {
    691             relay: other,
    692             attempted: true,
    693             outcome: DeliveryOutcome::accepted(),
    694         }]);
    695         assert_eq!(
    696             futures::executor::block_on(unexpected.deliver(request()))
    697                 .expect_err("unexpected relay evidence")
    698                 .code(),
    699             "invalid_transport_contract"
    700         );
    701 
    702         let denied_request = DeliveryRequest::new(
    703             "denied",
    704             payload(),
    705             TargetSet::new(vec![
    706                 Target::nostr_relay("wss://other.example").expect("other"),
    707             ])
    708             .expect("targets"),
    709             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    710             1_800_000_000_000,
    711         )
    712         .expect("request");
    713         let denied = futures::executor::block_on(scripted(vec![]).deliver(denied_request))
    714             .expect("denied receipt");
    715         assert!(!denied.target_receipts()[0].was_attempted());
    716         assert!(futures::executor::block_on(scripted(vec![]).status()).is_ok());
    717     }
    718 
    719     #[test]
    720     fn skipped_results_remain_skipped_and_cannot_claim_acceptance() {
    721         let one = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).unwrap();
    722         let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).unwrap();
    723         let receipt = futures::executor::block_on(
    724             scripted(vec![
    725                 RelayPublishResult {
    726                     relay: one.clone(),
    727                     attempted: false,
    728                     outcome: status::delivery_failure("timeout"),
    729                 },
    730                 RelayPublishResult {
    731                     relay: two,
    732                     attempted: true,
    733                     outcome: DeliveryOutcome::accepted(),
    734                 },
    735             ])
    736             .deliver(request()),
    737         )
    738         .unwrap();
    739         assert!(!receipt.target_receipts()[0].was_attempted());
    740         assert_eq!(
    741             receipt.target_receipts()[0].outcome().code(),
    742             Some("timeout")
    743         );
    744         assert!(receipt.target_receipts()[1].was_attempted());
    745         assert_eq!(
    746             receipt.target_receipts()[1].outcome(),
    747             &DeliveryOutcome::accepted()
    748         );
    749         for outcome in [DeliveryOutcome::accepted(), DeliveryOutcome::delivered()] {
    750             let transport = scripted(vec![RelayPublishResult {
    751                 relay: one.clone(),
    752                 attempted: false,
    753                 outcome,
    754             }]);
    755             let failure = futures::executor::block_on(transport.deliver(request())).unwrap_err();
    756             assert_eq!(failure.code(), "invalid_transport_contract");
    757             assert!(
    758                 transport
    759                     .relay_status()
    760                     .relays()
    761                     .iter()
    762                     .all(|relay| relay.write().last_success_unix_ms().is_none())
    763             );
    764         }
    765     }
    766 
    767     #[test]
    768     fn live_relay_client_accepts_an_empty_batch_without_io() {
    769         let client = LiveRelayClient::isolated();
    770         let results = futures::executor::block_on(client.publish(
    771             vec![],
    772             ExactEvent::from_request(&request()).expect("exact event"),
    773             1,
    774             Duration::from_millis(1),
    775             Duration::from_millis(1),
    776         ));
    777         assert!(results.is_empty());
    778     }
    779 }