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 }