lib

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

exact_delivery.rs (12770B)


      1 use futures::{FutureExt, SinkExt, StreamExt};
      2 use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Kind, Tag, Timestamp};
      3 use radroots_transport::{
      4     DeliveryRequest, EventSink, EventSource, FetchRequest, TargetSet,
      5     policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
      6     sink::DeliveryPayload,
      7     source::FetchBounds,
      8 };
      9 use radroots_transport_nostr::{
     10     Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind,
     11     RelayUrlPolicy,
     12 };
     13 use serde_json::{Value, json};
     14 use std::time::{Duration, SystemTime, UNIX_EPOCH};
     15 use tokio::net::{TcpListener, TcpStream};
     16 use tokio_tungstenite::{WebSocketStream, accept_async, tungstenite::Message};
     17 
     18 const KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001";
     19 
     20 fn now() -> u64 {
     21     SystemTime::now()
     22         .duration_since(UNIX_EPOCH)
     23         .unwrap()
     24         .as_millis()
     25         .try_into()
     26         .unwrap()
     27 }
     28 
     29 fn raw_event() -> String {
     30     let event = EventBuilder::text_note("harvest café\n🌱")
     31         .custom_created_at(Timestamp::from_secs(1_800_000_000))
     32         .sign_with_keys(&Keys::parse(KEY).unwrap())
     33         .unwrap();
     34     let mut value: Value = serde_json::from_str(&event.as_json()).unwrap();
     35     value["extension"] = json!({"retained": true});
     36     format!(" \n{} \t", serde_json::to_string_pretty(&value).unwrap())
     37 }
     38 
     39 fn config(url: &str, timeout: u64) -> Config {
     40     Config::from_profile(
     41         RelayProfile::explicit(
     42             RelayProfileKind::Simulator,
     43             [RelayEndpoint::new(url, RelayUrlPolicy::Local, RelayAccess::ReadWrite).unwrap()],
     44         )
     45         .unwrap(),
     46     )
     47     .with_timeouts(1_000, timeout, 500)
     48     .unwrap()
     49 }
     50 
     51 fn targets(config: &Config) -> TargetSet {
     52     TargetSet::new(
     53         config
     54             .relays()
     55             .iter()
     56             .map(|relay| relay.to_target().unwrap())
     57             .collect(),
     58     )
     59     .unwrap()
     60 }
     61 
     62 fn request(config: &Config, raw: &str) -> DeliveryRequest {
     63     DeliveryRequest::new(
     64         "exact-wire",
     65         DeliveryPayload::new(radroots_event_codec::decode::signed_event(raw).unwrap()),
     66         targets(config),
     67         SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
     68         now() + 5_000,
     69     )
     70     .unwrap()
     71 }
     72 
     73 async fn text(socket: &mut WebSocketStream<TcpStream>) -> String {
     74     tokio::time::timeout(Duration::from_secs(5), async {
     75         loop {
     76             if let Message::Text(text) = socket.next().await.unwrap().unwrap() {
     77                 return text.to_string();
     78             }
     79         }
     80     })
     81     .await
     82     .unwrap()
     83 }
     84 
     85 async fn reply(socket: &mut WebSocketStream<TcpStream>, value: Value) {
     86     socket
     87         .send(Message::Text(value.to_string().into()))
     88         .await
     89         .unwrap();
     90 }
     91 
     92 #[tokio::test]
     93 async fn selected_delivery_never_connects_to_held_target_and_preserves_full_request() {
     94     let a = TcpListener::bind("127.0.0.1:0").await.unwrap();
     95     let b = TcpListener::bind("127.0.0.1:0").await.unwrap();
     96     let config = Config::from_profile(
     97         RelayProfile::explicit(
     98             RelayProfileKind::Simulator,
     99             [a.local_addr().unwrap(), b.local_addr().unwrap()].map(|address| {
    100                 RelayEndpoint::new(
    101                     format!("ws://{address}"),
    102                     RelayUrlPolicy::Local,
    103                     RelayAccess::ReadWrite,
    104                 )
    105                 .unwrap()
    106             }),
    107         )
    108         .unwrap(),
    109     )
    110     .with_timeouts(1_000, 2_000, 500)
    111     .unwrap();
    112     let raw = raw_event();
    113     let expected = format!("[\"EVENT\",{raw}]");
    114     let server = tokio::spawn(async move {
    115         let (tcp, _) = a.accept().await.unwrap();
    116         let mut socket = accept_async(tcp).await.unwrap();
    117         let wire = text(&mut socket).await;
    118         assert_eq!(wire, expected);
    119         let event: Value = serde_json::from_str(&wire).unwrap();
    120         reply(&mut socket, json!(["OK", event[1]["id"], true, ""])).await;
    121     });
    122     let transport = NostrTransport::new(config.clone());
    123     let request = request(&config, &raw);
    124     let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap();
    125     let result = transport
    126         .deliver_selected(request.clone(), selected)
    127         .await
    128         .unwrap();
    129     result.validate_for_request(&request).unwrap();
    130     assert!(result.target_receipts()[0].was_attempted());
    131     assert!(
    132         result.target_receipts()[0]
    133             .outcome()
    134             .satisfies(SatisfactionClass::Accepted)
    135     );
    136     assert!(!result.target_receipts()[1].was_attempted());
    137     assert_eq!(
    138         result.target_receipts()[1].outcome().code(),
    139         Some("target_not_selected")
    140     );
    141     assert!(!result.is_satisfied(&request).unwrap());
    142     server.await.unwrap();
    143     assert!(
    144         tokio::time::timeout(Duration::from_millis(150), b.accept())
    145             .await
    146             .is_err()
    147     );
    148 }
    149 
    150 #[tokio::test]
    151 async fn exact_signed_bytes_read_requests_and_auth_use_the_same_connection() {
    152     let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
    153     let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 2_000);
    154     let raw = raw_event();
    155     let expected = format!("[\"EVENT\",{raw}]");
    156     let (auth_sent, auth_seen) = tokio::sync::oneshot::channel();
    157     let server = tokio::spawn(async move {
    158         let (tcp, _) = listener.accept().await.unwrap();
    159         let mut socket = accept_async(tcp).await.unwrap();
    160         let req: Value = serde_json::from_str(&text(&mut socket).await).unwrap();
    161         assert_eq!(req[0], "REQ");
    162         reply(&mut socket, json!(["EOSE", req[1]])).await;
    163         loop {
    164             let message: Value = serde_json::from_str(&text(&mut socket).await).unwrap();
    165             if message[0] == "AUTH" {
    166                 assert_eq!(message[1]["kind"], 22242);
    167                 auth_sent.send(()).unwrap();
    168                 break;
    169             }
    170             assert_eq!(message[0], "CLOSE");
    171         }
    172         let wire = text(&mut socket).await;
    173         assert_eq!(wire, expected);
    174         let event: Value = serde_json::from_str(&wire).unwrap();
    175         reply(
    176             &mut socket,
    177             json!(["OK", "11".repeat(32), true, "unrelated event"]),
    178         )
    179         .await;
    180         reply(&mut socket, json!(["OK", event[1]["id"], true, ""])).await;
    181     });
    182     let transport = NostrTransport::new(config.clone());
    183     let page = transport
    184         .fetch(
    185             FetchRequest::new(
    186                 "same-socket-read",
    187                 targets(&config),
    188                 FetchBounds::new(1, now() + 5_000).unwrap(),
    189             )
    190             .unwrap(),
    191         )
    192         .await
    193         .unwrap();
    194     assert!(page.events().is_empty());
    195     let relay = &config.relays()[0];
    196     let at = now() / 1_000 * 1_000;
    197     transport
    198         .begin_authentication(relay, "challenge", at, at + 5_000)
    199         .unwrap();
    200     let auth = EventBuilder::new(Kind::Authentication, "")
    201         .tags([
    202             Tag::parse(["relay", relay.as_str()]).unwrap(),
    203             Tag::parse(["challenge", "challenge"]).unwrap(),
    204         ])
    205         .custom_created_at(Timestamp::from_secs(at / 1_000))
    206         .sign_with_keys(&Keys::parse(KEY).unwrap())
    207         .unwrap();
    208     transport
    209         .complete_authentication(relay, "challenge", Some(&auth.as_json()), at)
    210         .await
    211         .unwrap();
    212     tokio::time::timeout(Duration::from_secs(5), auth_seen)
    213         .await
    214         .unwrap()
    215         .unwrap();
    216     let receipt = transport.deliver(request(&config, &raw)).await.unwrap();
    217     assert!(
    218         receipt.target_receipts()[0]
    219             .outcome()
    220             .satisfies(SatisfactionClass::Accepted)
    221     );
    222     server.await.unwrap();
    223 }
    224 
    225 #[tokio::test]
    226 async fn rejection_lost_ack_and_disconnect_never_invent_acceptance() {
    227     for mode in ["reject", "lost", "disconnect"] {
    228         let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
    229         let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 250);
    230         let raw = raw_event();
    231         let expected = format!("[\"EVENT\",{raw}]");
    232         let server = tokio::spawn(async move {
    233             let (tcp, _) = listener.accept().await.unwrap();
    234             let mut socket = accept_async(tcp).await.unwrap();
    235             let wire = text(&mut socket).await;
    236             assert_eq!(wire, expected);
    237             let event: Value = serde_json::from_str(&wire).unwrap();
    238             match mode {
    239                 "reject" => {
    240                     reply(
    241                         &mut socket,
    242                         json!(["OK", event[1]["id"], false, "blocked: fixture"]),
    243                     )
    244                     .await
    245                 }
    246                 "lost" => {
    247                     reply(
    248                         &mut socket,
    249                         json!(["OK", "22".repeat(32), true, "wrong ID"]),
    250                     )
    251                     .await;
    252                     tokio::time::sleep(Duration::from_millis(500)).await;
    253                 }
    254                 "disconnect" => socket.close(None).await.unwrap(),
    255                 _ => unreachable!(),
    256             }
    257         });
    258         let receipt = NostrTransport::new(config.clone())
    259             .deliver(request(&config, &raw))
    260             .await
    261             .unwrap();
    262         assert!(
    263             !receipt.target_receipts()[0]
    264                 .outcome()
    265                 .satisfies(SatisfactionClass::Accepted),
    266             "{mode}"
    267         );
    268         assert!(receipt.target_receipts()[0].was_attempted());
    269         server.await.unwrap();
    270     }
    271 }
    272 
    273 #[tokio::test]
    274 async fn quota_refusal_is_terminal_and_redacted_at_the_receipt_boundary() {
    275     let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
    276     let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 2_000);
    277     let server = tokio::spawn(async move {
    278         let (tcp, _) = listener.accept().await.unwrap();
    279         let mut socket = accept_async(tcp).await.unwrap();
    280         let event: Value = serde_json::from_str(&text(&mut socket).await).unwrap();
    281         reply(
    282             &mut socket,
    283             json!([
    284                 "OK",
    285                 event[1]["id"],
    286                 false,
    287                 "quota exceeded: private-account-detail"
    288             ]),
    289         )
    290         .await;
    291     });
    292     let transport = NostrTransport::new(config.clone());
    293     let receipt = transport
    294         .deliver(request(&config, &raw_event()))
    295         .await
    296         .unwrap();
    297     let target = &receipt.target_receipts()[0];
    298     assert!(target.was_attempted());
    299     assert_eq!(target.outcome().code(), Some("quota_exceeded"));
    300     assert_eq!(
    301         target.outcome().retryability(),
    302         radroots_transport::outcome::Retryability::Terminal
    303     );
    304     assert_eq!(
    305         target.outcome().kind(),
    306         radroots_transport::outcome::DeliveryOutcomeKind::Rejected
    307     );
    308     assert_eq!(
    309         target.outcome().message(),
    310         Some("relay quota was exhausted")
    311     );
    312     assert_eq!(
    313         transport.relay_status().relays()[0]
    314             .write()
    315             .last_failure_retryable(),
    316         Some(false)
    317     );
    318     server.await.unwrap();
    319 }
    320 
    321 #[tokio::test]
    322 async fn queued_delivery_targets_cannot_start_after_the_shared_deadline() {
    323     let first = TcpListener::bind("127.0.0.1:0").await.unwrap();
    324     let second = TcpListener::bind("127.0.0.1:0").await.unwrap();
    325     let endpoints = [&first, &second].map(|listener| {
    326         RelayEndpoint::new(
    327             format!("ws://{}", listener.local_addr().unwrap()),
    328             RelayUrlPolicy::Local,
    329             RelayAccess::ReadWrite,
    330         )
    331         .unwrap()
    332     });
    333     let config = Config::from_profile(
    334         RelayProfile::explicit(RelayProfileKind::Simulator, endpoints).unwrap(),
    335     )
    336     .with_timeouts(1_000, 250, 500)
    337     .unwrap()
    338     .with_max_connections(1)
    339     .unwrap();
    340     let server = tokio::spawn(async move {
    341         let (tcp, _) = first.accept().await.unwrap();
    342         let mut socket = accept_async(tcp).await.unwrap();
    343         let _published = text(&mut socket).await;
    344         futures::future::pending::<()>().await;
    345     });
    346     let transport = NostrTransport::new(config.clone());
    347     let receipt = transport
    348         .deliver(request(&config, &raw_event()))
    349         .await
    350         .unwrap();
    351     assert_eq!(receipt.target_receipts().len(), 2);
    352     assert!(
    353         receipt
    354             .target_receipts()
    355             .iter()
    356             .all(|target| !target.outcome().satisfies(SatisfactionClass::Accepted))
    357     );
    358     assert!(second.accept().now_or_never().is_none());
    359     assert!(receipt.target_receipts()[0].was_attempted());
    360     assert!(
    361         !receipt.target_receipts()[1].was_attempted(),
    362         "a target expired while queued never entered remote publication"
    363     );
    364     server.abort();
    365     assert!(server.await.unwrap_err().is_cancelled());
    366 }