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 }