lib

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

local_io.rs (8067B)


      1 use futures::{SinkExt, StreamExt};
      2 use radroots_transport::{
      3     EventSource, EventSubscriber, FetchRequest, SubscriptionNext, SubscriptionRequest, TargetSet,
      4     capability::Availability,
      5     outcome::FetchTargetState,
      6     source::{FetchBounds, SubscriptionBounds, SubscriptionEndReason},
      7 };
      8 use radroots_transport_nostr::{
      9     Config, NostrTransport, RelayAccess, RelayAggregateState, RelayEndpoint, RelayEvidenceState,
     10     RelayProfile, RelayProfileKind, RelayUrlPolicy,
     11 };
     12 use serde_json::Value;
     13 use std::time::{Duration, SystemTime, UNIX_EPOCH};
     14 use tokio::net::TcpListener;
     15 use tokio_tungstenite::{accept_async, tungstenite::Message};
     16 
     17 const FIXTURE_SECRET_KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001";
     18 
     19 #[tokio::test(flavor = "multi_thread")]
     20 async fn simulator_profile_proves_read_capability_against_a_real_loopback_socket() {
     21     let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind relay");
     22     let address = listener.local_addr().expect("relay address");
     23     let relay_url = format!("ws://{address}");
     24     let server = tokio::spawn(async move {
     25         let (stream, _) = listener.accept().await.expect("accept relay client");
     26         let mut websocket = accept_async(stream).await.expect("websocket handshake");
     27         while let Some(message) = websocket.next().await {
     28             let Message::Text(message) = message.expect("client message") else {
     29                 continue;
     30             };
     31             let parsed: Value = serde_json::from_str(message.as_str()).expect("Nostr message");
     32             let Some(values) = parsed.as_array() else {
     33                 continue;
     34             };
     35             let [Value::String(kind), Value::String(subscription), ..] = values.as_slice() else {
     36                 continue;
     37             };
     38             if kind == "REQ" {
     39                 websocket
     40                     .send(Message::Text(
     41                         serde_json::to_string(&("EOSE", subscription))
     42                             .expect("EOSE message")
     43                             .into(),
     44                     ))
     45                     .await
     46                     .expect("send EOSE");
     47                 break;
     48             }
     49         }
     50     });
     51 
     52     let endpoint = RelayEndpoint::new(
     53         relay_url.as_str(),
     54         RelayUrlPolicy::Local,
     55         RelayAccess::ReadWrite,
     56     )
     57     .expect("loopback endpoint");
     58     let profile =
     59         RelayProfile::explicit(RelayProfileKind::Simulator, [endpoint]).expect("simulator profile");
     60     let config = Config::from_profile(profile)
     61         .with_timeouts(1_000, 2_000, 500)
     62         .expect("timeouts");
     63     let targets = TargetSet::new(
     64         config
     65             .read_relays()
     66             .map(|relay| relay.to_target())
     67             .collect::<Result<Vec<_>, _>>()
     68             .expect("targets"),
     69     )
     70     .expect("target set");
     71     let request = FetchRequest::new(
     72         "local-real-io",
     73         targets,
     74         FetchBounds::new(1, unix_time_ms() + 5_000).expect("bounds"),
     75     )
     76     .expect("request");
     77     let transport = NostrTransport::new(config);
     78     let page = tokio::time::timeout(Duration::from_secs(5), transport.fetch(request))
     79         .await
     80         .expect("bounded fetch")
     81         .expect("fetch page");
     82 
     83     assert!(page.events().is_empty());
     84     assert_eq!(page.target_outcomes().len(), 1);
     85     assert_eq!(
     86         page.target_outcomes()[0].state(),
     87         FetchTargetState::Complete
     88     );
     89     let status = transport.relay_status();
     90     assert_eq!(status.state(), RelayAggregateState::ReadOnly);
     91     assert_eq!(status.read_availability(), Availability::Available);
     92     assert_eq!(status.write_availability(), Availability::Unavailable);
     93     assert_eq!(
     94         status.relays()[0].read().state(),
     95         RelayEvidenceState::Available
     96     );
     97     assert_eq!(
     98         status.relays()[0].write().state(),
     99         RelayEvidenceState::Unobserved
    100     );
    101     server.await.expect("relay task");
    102 }
    103 
    104 #[tokio::test(flavor = "multi_thread")]
    105 async fn simulator_profile_streams_and_cancels_a_real_live_subscription() {
    106     use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Timestamp};
    107 
    108     let event = EventBuilder::text_note("bounded live event")
    109         .custom_created_at(Timestamp::from_secs(1_800_000_100))
    110         .sign_with_keys(&Keys::parse(FIXTURE_SECRET_KEY).expect("fixture keys"))
    111         .expect("signed event");
    112     let expected_id = event.id.to_hex();
    113     let event: Value = serde_json::from_str(event.as_json().as_str()).expect("event JSON");
    114     let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind relay");
    115     let address = listener.local_addr().expect("relay address");
    116     let relay_url = format!("ws://{address}");
    117     let server = tokio::spawn(async move {
    118         let (stream, _) = listener.accept().await.expect("accept relay client");
    119         let mut websocket = accept_async(stream).await.expect("websocket handshake");
    120         let mut subscription_id = None;
    121         while let Some(message) = websocket.next().await {
    122             let Message::Text(message) = message.expect("client message") else {
    123                 continue;
    124             };
    125             let parsed: Value = serde_json::from_str(message.as_str()).expect("Nostr message");
    126             let Some(values) = parsed.as_array() else {
    127                 continue;
    128             };
    129             match values.as_slice() {
    130                 [Value::String(kind), Value::String(subscription), ..] if kind == "REQ" => {
    131                     subscription_id = Some(subscription.clone());
    132                     websocket
    133                         .send(Message::Text(
    134                             serde_json::to_string(&("EVENT", subscription, &event))
    135                                 .expect("EVENT message")
    136                                 .into(),
    137                         ))
    138                         .await
    139                         .expect("send EVENT");
    140                 }
    141                 [Value::String(kind), Value::String(subscription)]
    142                     if kind == "CLOSE" && subscription_id.as_ref() == Some(subscription) =>
    143                 {
    144                     return;
    145                 }
    146                 _ => {}
    147             }
    148         }
    149         panic!("client disconnected without CLOSE");
    150     });
    151 
    152     let endpoint = RelayEndpoint::new(
    153         relay_url.as_str(),
    154         RelayUrlPolicy::Local,
    155         RelayAccess::ReadWrite,
    156     )
    157     .expect("loopback endpoint");
    158     let profile =
    159         RelayProfile::explicit(RelayProfileKind::Simulator, [endpoint]).expect("simulator profile");
    160     let config = Config::from_profile(profile)
    161         .with_timeouts(1_000, 2_000, 500)
    162         .expect("timeouts");
    163     let targets = TargetSet::new(
    164         config
    165             .read_relays()
    166             .map(|relay| relay.to_target())
    167             .collect::<Result<Vec<_>, _>>()
    168             .expect("targets"),
    169     )
    170     .expect("target set");
    171     let request = SubscriptionRequest::new(
    172         "local-live-io",
    173         targets,
    174         SubscriptionBounds::new(2, unix_time_ms() + 5_000).expect("bounds"),
    175     )
    176     .expect("request");
    177     let transport = NostrTransport::new(config);
    178     let mut subscription =
    179         tokio::time::timeout(Duration::from_secs(5), transport.subscribe(request))
    180             .await
    181             .expect("bounded subscribe")
    182             .expect("subscription");
    183     let SubscriptionNext::Event(observed) =
    184         tokio::time::timeout(Duration::from_secs(5), subscription.next())
    185             .await
    186             .expect("bounded event")
    187             .expect("event")
    188     else {
    189         panic!("event expected");
    190     };
    191     assert_eq!(observed.observed().event().id_str(), expected_id);
    192     assert_eq!(
    193         subscription.cancel().await.expect("cancel").reason(),
    194         SubscriptionEndReason::Cancelled
    195     );
    196     tokio::time::timeout(Duration::from_secs(5), server)
    197         .await
    198         .expect("server deadline")
    199         .expect("relay task");
    200 }
    201 
    202 fn unix_time_ms() -> u64 {
    203     SystemTime::now()
    204         .duration_since(UNIX_EPOCH)
    205         .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
    206         .unwrap_or_default()
    207 }