mod common; use std::time::Duration; use common::{ challenge, client, publish, query_envelope, register, spawn_relay, spawn_relay_with_capacity, }; use frxd::crypto::Keypair; use futures_util::StreamExt; async fn open_stream(http: &reqwest::Client, relay: &str, key: &Keypair) -> reqwest::Response { let member = key.public_hex(); let nonce = challenge(http, relay, &member).await.expect("challenge"); let sig = key.sign(&frxd::crypto::poll_signing_bytes(&member, &nonce)); http.get(format!( "{relay}/v1/stream?member={member}&nonce={nonce}&sig={sig}" )) .send() .await .unwrap() } async fn read_until(mut stream: S, needle: &str, timeout: Duration) -> (String, S) where S: futures_util::Stream> + Unpin, { let mut buffer = String::new(); let deadline = tokio::time::Instant::now() + timeout; loop { let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); let chunk = tokio::time::timeout(remaining, stream.next()) .await .expect("timed out waiting for stream data"); match chunk { Some(Ok(bytes)) => { buffer.push_str(&String::from_utf8_lossy(&bytes)); if buffer.contains(needle) { return (buffer, stream); } } Some(Err(error)) => panic!("stream error: {error}"), None => panic!("stream closed before '{needle}'"), } } } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn stream_delivers_envelopes_live() { let relay = spawn_relay().await; let http = client(); let member = Keypair::generate(); let publisher = Keypair::generate(); let response = open_stream(&http, &relay, &member).await; assert!(response.status().is_success()); assert!( response .headers() .get("content-type") .and_then(|value| value.to_str().ok()) .unwrap_or("") .contains("text/event-stream") ); let envelope = query_envelope(&publisher, &publisher.public_hex(), "live stream", 5); assert!( publish(&http, &relay, &envelope) .await .status() .is_success() ); let (buffer, _) = read_until( response.bytes_stream(), "live stream", Duration::from_secs(3), ) .await; assert!(buffer.contains("event: envelope"), "{buffer}"); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn stream_reports_lag_then_recovers() { let relay = spawn_relay_with_capacity(1).await; let http = client(); let member = Keypair::generate(); let publisher = Keypair::generate(); register(&http, &relay, &member).await; let first = query_envelope(&publisher, &publisher.public_hex(), "first", 5); assert!(publish(&http, &relay, &first).await.status().is_success()); let second = query_envelope(&publisher, &publisher.public_hex(), "second", 5); assert!(publish(&http, &relay, &second).await.status().is_success()); let response = open_stream(&http, &relay, &member).await; assert!(response.status().is_success()); let (buffer, stream) = read_until( response.bytes_stream(), "\"missed\"", Duration::from_secs(3), ) .await; assert!(buffer.contains("event: lag"), "{buffer}"); let third = query_envelope(&publisher, &publisher.public_hex(), "third", 5); assert!(publish(&http, &relay, &third).await.status().is_success()); let (buffer, _) = read_until(stream, "third", Duration::from_secs(3)).await; assert!(buffer.contains("event: envelope"), "{buffer}"); }