#![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;
static SERIAL: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
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");
}
#[tokio::test(flavor = "multi_thread")]
async fn probe_carries_peer_app_metadata_into_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_eid = format!("eid:{}", a_ep.id());
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 a_socket = dir.path().join("a-control.sock");
let a_listener = mcpmesh::ipc::bind_control_socket(&a_socket).await.unwrap();
let a_state = Arc::new(DaemonState::with_mesh(STACK_VERSION, a_mesh.clone()));
let a_control = tokio::spawn(serve_control(a_listener, a_state));
connect_control(&a_socket)
.await
.expect("connect A control")
.set_app_metadata("v=4.2.0")
.await
.expect("A sets its app metadata");
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "the paired peer is reachable");
assert_eq!(
entry.meta, "v=4.2.0",
"the probe carried the peer's app metadata off the pong"
);
let list = reachability_of(&b_mesh);
let alice = list.iter().find(|r| r.name == "alice").expect("A surfaced");
assert_eq!(
alice.meta, "v=4.2.0",
"reachability surfaces the peer's app metadata"
);
assert_eq!(
alice.principal.as_deref(),
Some(a_eid.as_str()),
"reachability row carries the peer's eid principal"
);
a_control.abort();
accept.abort();
std::mem::forget(dir);
})
.await
.expect("probe-metadata test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn probe_surfaces_the_services_the_peer_grants_the_caller() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_eid = format!("eid:{}", b_ep.id());
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");
let config = dir.path().join("a-config.toml");
std::fs::write(
&config,
format!(
"[services.shared]\nsocket = \"/run/s.sock\"\nallow = [\"{b_eid}\"]\n\
[services.private]\nsocket = \"/run/p.sock\"\nallow = [\"eid:someoneelse\"]\n"
),
)
.unwrap();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
seed_peer(&a_store, b_id, "beacon-b");
let a_cfg = Config::load(&config).expect("A's config parses");
let a_mesh = assemble_mesh(a_ep, a_store, config);
let b_mesh = assemble_mesh(b_ep, b_store, dir.path().join("b-config.toml"));
let accept = spawn_accept_loop(a_mesh.clone(), Arc::new(build_services(&a_cfg)));
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable);
assert_eq!(
entry.services,
vec!["shared".to_string()],
"probe surfaces exactly the caller-admitted services (#52)"
);
assert!(
!entry.services.contains(&"private".to_string()),
"never a service the peer does not grant the caller"
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("peer-services probe test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_throttled_probe_is_refused_but_never_reports_the_peer_offline() {
let _serial = SERIAL.lock().await;
timeout(Duration::from_secs(120), 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("fa.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("fb.redb")).unwrap());
seed_peer(&a_store, b_id, "b");
seed_peer(&b_store, a_id, "a");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let a_limits =
mcpmesh::limits::MeshLimiters::from_config(&mcpmesh::config::LimitsCfg::default());
a_mesh.set_limits(a_limits.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 mut reachable = 0usize;
let mut fresh = 0usize;
let mut stale = 0usize;
let mut last_seq: Option<u64> = None;
for _ in 0..90 {
let entry = probe_peer(&b_mesh, a_id).await;
if entry.reachable {
reachable += 1;
}
if last_seq == Some(entry.seq) {
stale += 1;
} else {
fresh += 1;
}
last_seq = Some(entry.seq);
}
assert_eq!(
reachable, 90,
"a rate-limit refusal is not evidence the peer is down: every probe of a live, paired \
peer must report reachable — the refused ones from the still-fresh cache. A false \
count here means a refusal wrote `reachable: false` (PR #142 gate, HIGH). \
reachable={reachable} fresh={fresh} stale={stale}"
);
assert!(
stale > 0,
"a paired peer flooding past the cap must eventually be REFUSED (#89), observable as \
the cached entry returned unchanged. Zero stale answers means the accept arm never \
consulted the limiter — the unmetered pong-flood this issue reports. \
fresh={fresh} stale={stale}"
);
assert!(
a_limits.pings_refused() > 0,
"the responder must COUNT its refusals (#89 defect 3: unmetered AND unrecorded) — \
the count is the refusal's only footprint besides the debug log"
);
})
.await
.expect("ping flood test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn peer_services_answers_from_the_fresh_cache_without_probing() {
let _serial = SERIAL.lock().await;
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().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 b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_eid = format!("eid:{}", b_ep.id());
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");
let config = dir.path().join("a-config.toml");
std::fs::write(
&config,
format!("[services.shared]\nsocket = \"/run/s.sock\"\nallow = [\"{b_eid}\"]\n"),
)
.unwrap();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
seed_peer(&a_store, b_id, "beacon-b");
let a_cfg = Config::load(&config).expect("A's config parses");
let a_mesh = assemble_mesh(a_ep, a_store, config);
let b_mesh = assemble_mesh(b_ep, b_store, dir.path().join("b-config.toml"));
let accept = spawn_accept_loop(a_mesh.clone(), Arc::new(build_services(&a_cfg)));
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "precondition: A must probe reachable");
assert_eq!(
entry.services,
vec!["shared".to_string()],
"precondition: the cached entry carries the granted service"
);
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("connect B control");
accept.abort();
a_ep_handle.close().await;
let services = client.peer_services("alice").await.expect(
"peer_services must answer from the fresh cache rather than probing — an \
unconditional probe is the #142-gate shape that reported healthy peers offline",
);
assert_eq!(services, vec!["shared".to_string()]);
control.abort();
std::mem::forget(dir);
})
.await
.expect("peer_services cache-freshness test timed out");
}