mod common; use std::fs; use std::time::Duration; use common::{ client, collection, config_for, messages_of_type, poll_messages, publish, query_envelope, register, test_envelope, }; use frxd::crypto::{Keypair, now_ts}; use frxd::index::LocalIndex; use frxd::message::{EXPOSURE_FULL, TYPE_QUERY, TYPE_RESPONSE}; use frxd::node::Node; use frxd::registry::{self, KeyEntry, RegistryDoc, RegistryMember, SignedRegistry}; use frxd::relay::{self, RelayOptions}; async fn spawn_relay_with(options: RelayOptions) -> String { let (listener, addr) = relay::bind("127.0.0.1:0").await.unwrap(); tokio::spawn(async move { let _ = relay::run(listener, options).await; }); format!("http://{addr}") } async fn federation_mesh() -> [String; 3] { let mut listeners = Vec::new(); let mut urls = Vec::new(); for _ in 0..3 { let (listener, addr) = relay::bind("127.0.0.1:0").await.unwrap(); listeners.push(listener); urls.push(format!("http://{addr}")); } let peers: Vec> = vec![ vec![urls[1].clone(), urls[2].clone()], vec![urls[0].clone(), urls[2].clone()], vec![urls[0].clone(), urls[1].clone()], ]; for (index, listener) in listeners.into_iter().enumerate() { let options = RelayOptions { url: Some(urls[index].clone()), peers: peers[index].clone(), ..Default::default() }; tokio::spawn(async move { let _ = relay::run(listener, options).await; }); } [urls[0].clone(), urls[1].clone(), urls[2].clone()] } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn mesh_floods_queries_exactly_once() { let [r1, _r2, r3] = federation_mesh().await; let http = client(); let member = Keypair::generate(); let publisher = Keypair::generate(); register(&http, &r3, &member).await; let envelope = query_envelope(&publisher, &publisher.public_hex(), "mesh", 5); assert!(publish(&http, &r1, &envelope).await.status().is_success()); let messages = poll_messages(&http, &r3, &member, 900).await; assert_eq!(messages.len(), 1, "expected one flooded copy: {messages:?}"); tokio::time::sleep(Duration::from_millis(300)).await; let duplicates = poll_messages(&http, &r3, &member, 200).await; assert!( duplicates.is_empty(), "mesh delivered duplicates: {duplicates:?}" ); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn federation_rejects_unknown_origin() { let [r1, _r2, _r3] = federation_mesh().await; let http = client(); let key = Keypair::generate(); let envelope = query_envelope(&key, &key.public_hex(), "spoof", 5); let response = http .post(format!( "{r1}/v1/federation?origin=http://127.0.0.1:1&hops=1" )) .json(&envelope) .send() .await .unwrap(); assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); } fn registry_doc( ma: &Keypair, members: Vec, relays: Vec, version: u64, ) -> SignedRegistry { registry::sign_registry( RegistryDoc { version, issued_at: now_ts(), ma_key: String::new(), zone: "frx.invalid".to_string(), members, relays, }, ma, ) } fn member(id: &str, key: &Keypair, class: &str) -> RegistryMember { RegistryMember { id: id.to_string(), class: class.to_string(), keys: vec![KeyEntry { key: key.public_hex(), not_before: 0, not_after: None, }], enc_key: None, } } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn relay_admission_gates_unlisted_keys() { let root = tempfile::tempdir().unwrap(); let ma = Keypair::generate(); let alice = Keypair::generate(); let registry_path = root.path().join("registry.json"); registry::save_registry( ®istry_path, ®istry_doc( &ma, vec![member("alice.frx.example", &alice, "source")], vec![], 1, ), ) .unwrap(); let relay_url = spawn_relay_with(RelayOptions { registry: Some(registry_path.display().to_string()), ma_key: Some(ma.public_hex()), ..Default::default() }) .await; let http = client(); let accepted = query_envelope(&alice, "alice.frx.example", "rust", 5); assert_eq!( publish(&http, &relay_url, &accepted).await.status(), reqwest::StatusCode::ACCEPTED ); let stranger = Keypair::generate(); let denied = query_envelope(&stranger, &stranger.public_hex(), "rust", 5); assert_eq!( publish(&http, &relay_url, &denied).await.status(), reqwest::StatusCode::BAD_REQUEST, "unlisted key must not publish" ); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn node_discovers_relays_from_registry() { let root = tempfile::tempdir().unwrap(); let ma = Keypair::generate(); let alice = Keypair::generate(); let (listener, relay_addr) = relay::bind("127.0.0.1:0").await.unwrap(); tokio::spawn(async move { let _ = relay::run(listener, RelayOptions::default()).await; }); let relay_url = format!("http://{relay_addr}"); let docs = root.path().join("bob-docs"); fs::create_dir_all(&docs).unwrap(); fs::write(docs.join("doc.txt"), "registry discovery rust document").unwrap(); let mut bob_config = config_for(&root.path().join("bob"), "bob", "http://127.0.0.1:1"); bob_config.node.relays = Vec::new(); bob_config.node.id = Some("bob.frx.example".to_string()); let registry_path = root.path().join("registry.json"); bob_config.node.registry = Some(registry_path.display().to_string()); bob_config.node.ma_key = Some(ma.public_hex()); bob_config.node.dev_bootstrap = false; registry::save_registry( ®istry_path, ®istry_doc( &ma, vec![ member("alice.frx.example", &alice, "source"), member("bob.frx.example", &bob_config.load_key().unwrap(), "source"), ], vec![relay_url.clone()], 1, ), ) .unwrap(); { let index = LocalIndex::open(&bob_config.index_dir()).unwrap(); index .add_collection(&collection("docs", &docs, true, EXPOSURE_FULL)) .unwrap(); } let _bob = Node::start(bob_config).await.unwrap(); let http = client(); register(&http, &relay_url, &alice).await; let envelope = query_envelope(&alice, "alice.frx.example", "rust", 5); assert!( publish(&http, &relay_url, &envelope) .await .status() .is_success() ); let mut responses = Vec::new(); for _ in 0..5 { let messages = poll_messages(&http, &relay_url, &alice, 200).await; responses.extend(messages_of_type(&messages, TYPE_RESPONSE)); if !responses.is_empty() { break; } } assert_eq!( responses.len(), 1, "node did not discover its relay from the registry" ); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn unknown_types_never_flood() { let [r1, _r2, _r3] = federation_mesh().await; let http = client(); let key = Keypair::generate(); let response = test_envelope( &key, TYPE_RESPONSE, serde_json::json!({"qid": "q", "results": [], "truncated": false, "more_available": 0, "cursor": null}), ); let result = http .post(format!( "{r1}/v1/federation?origin={}&hops=1", "http://unknown" )) .json(&response) .send() .await .unwrap(); assert_eq!(result.status(), reqwest::StatusCode::BAD_REQUEST); let _ = TYPE_QUERY; }