#![cfg(unix)]
use std::sync::Arc;
use std::time::Duration;
use iroh::address_lookup::MemoryLookup;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::client::connect_control;
use mcpmesh::config::Config;
use mcpmesh::control::{DaemonState, serve_control};
use mcpmesh::daemon::{
MeshState, STACK_VERSION, build_services, probe_peer, reachability_of, spawn_accept_loop,
};
use mcpmesh::pairing::LiveInvites;
use mcpmesh::roster::gate::RosterGate;
use mcpmesh::{Request, StatusResult};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, ALPN_PAIR, ALPN_PING, TrustGate};
use tokio::time::timeout;
async fn target_endpoint() -> iroh::Endpoint {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![
ALPN_MCP.to_vec(),
ALPN_PAIR.to_vec(),
ALPN_PING.to_vec(),
])
.bind()
.await
.expect("bind target endpoint")
}
async fn dialer_endpoint() -> iroh::Endpoint {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![ALPN_MCP.to_vec()])
.bind()
.await
.expect("bind dialer endpoint")
}
fn seed_lookup(dialer: &iroh::Endpoint, target_addr: iroh::EndpointAddr) {
let mem = MemoryLookup::new();
mem.add_endpoint_info(target_addr);
dialer
.address_lookup()
.expect("address lookup services")
.add(mem);
}
fn seed_peer(store: &PeerStore, endpoint_id: [u8; 32], nickname: &str) {
store
.add(PeerEntry {
endpoint_id,
nickname: nickname.into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
}
fn seed_peer_with_addr(
store: &PeerStore,
endpoint_id: [u8; 32],
nickname: &str,
addr: &iroh::EndpointAddr,
) {
store
.add(PeerEntry {
endpoint_id,
nickname: nickname.into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: Some(serde_json::to_string(addr).expect("addr serializes")),
})
.unwrap();
}
fn assemble_mesh(
endpoint: iroh::Endpoint,
store: Arc<PeerStore>,
config_path: std::path::PathBuf,
) -> Arc<MeshState> {
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(store.clone()));
MeshState::new(
endpoint,
gate,
store,
Arc::new(LiveInvites::new()),
"self".into(),
config_path,
Arc::new(RosterGate::empty()),
Arc::new(ConnRegistry::new()),
None,
None,
None,
None,
)
}
#[tokio::test(flavor = "multi_thread")]
async fn ping_probe_reports_paired_peer_reachable_stranger_and_down_peer_not() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_ep_handle = a_ep.clone(); let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice"); seed_peer(&a_store, b_id, "beacon-b");
let c_ep = dialer_endpoint().await;
seed_lookup(&c_ep, a_addr.clone());
let c_store = Arc::new(PeerStore::open(&dir.path().join("c.redb")).unwrap());
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let c_mesh = assemble_mesh(c_ep, c_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "a paired peer's probe must be reachable");
assert!(
entry.rtt_ms.is_some(),
"a reachable probe records a round-trip time"
);
let list = reachability_of(&b_mesh);
let alice = list
.iter()
.find(|r| r.name == "alice")
.expect("reachability_of surfaces the paired peer by nickname");
assert!(alice.reachable, "the cached probe result is surfaced");
assert!(alice.rtt_ms.is_some(), "the cached RTT is surfaced");
let stranger = probe_peer(&c_mesh, a_id).await;
assert!(
!stranger.reachable,
"an unpaired peer gets no pong (trust gate closed the connection)"
);
assert!(stranger.rtt_ms.is_none());
accept.abort();
a_ep_handle.close().await;
let down = probe_peer(&b_mesh, a_id).await;
assert!(
!down.reachable,
"a probe of a down peer must be unreachable"
);
std::mem::forget(dir);
})
.await
.expect("reachability test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn status_includes_reachability() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
seed_peer(&a_store, b_id, "beacon-b");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let entry = probe_peer(&b_mesh, a_id).await;
assert!(
entry.reachable,
"precondition: the paired peer must probe reachable"
);
let socket = dir.path().join("control.sock");
let listener = mcpmesh::ipc::bind_control_socket(&socket).await.unwrap();
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, b_mesh.clone()));
let control = tokio::spawn(serve_control(listener, state));
let mut client = connect_control(&socket)
.await
.expect("raw connect_control to B");
let value = client
.request(Request::Status)
.await
.expect("status over mcpmesh-local/1");
let status: StatusResult =
serde_json::from_value(value).expect("StatusResult deserializes");
assert!(
status.reachability.iter().any(|r| r.name == "alice"),
"status.reachability must surface the paired peer by nickname: {:?}",
status.reachability
);
control.abort();
accept.abort();
std::mem::forget(dir);
})
.await
.expect("status reachability test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn cold_probe_uses_the_pairing_proven_address_hint_without_discovery() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer_with_addr(&b_store, a_id, "alice", &a_addr);
seed_peer(&a_store, b_id, "beacon-b");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
tokio::time::sleep(Duration::from_millis(200)).await;
let entry = probe_peer(&b_mesh, a_id).await;
assert!(
entry.reachable,
"a cold probe must reach the peer using the pairing-proven `last_addr` hint rather \
than depending on discovery to resolve the bare endpoint-id"
);
assert!(
entry.rtt_ms.is_some(),
"a reachable probe reports a measured RTT"
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("cold-probe address-hint test timed out");
}