use std::sync::Arc;
use std::time::Duration;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::config::Config;
use mcpmesh::daemon::{
MeshState, build_services, revoke_service_access, revoke_service_allow, spawn_accept_loop,
};
use mcpmesh::pairing::LiveInvites;
use mcpmesh::roster::gate::RosterGate;
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, MAX_FRAME_BYTES, SessionTransport, TrustGate, framing::write_frame};
use serde_json::json;
use tokio::time::timeout;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
async fn server_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 server endpoint")
}
async fn client_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 client endpoint")
}
fn initialize_frame(service: &str) -> serde_json::Value {
json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {
"protocolVersion": "2025-11-25",
"_meta": {"mcpmesh/service": service},
"capabilities": {}, "clientInfo": {"name": "tester", "version": "0"}
}
})
}
fn tools_call_frame(text: &str) -> serde_json::Value {
json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"name": "echo", "arguments": {"text": text}}
})
}
struct Peer {
endpoint: iroh::Endpoint,
principal: String,
}
async fn paired_mesh(
nicknames: &[&str],
) -> (
Arc<MeshState>,
iroh::EndpointAddr,
Arc<ConnRegistry>,
Vec<Peer>,
tempfile::TempDir,
) {
let dir = tempfile::tempdir().unwrap();
let store = Arc::new(PeerStore::open(&dir.path().join("state.redb")).unwrap());
let mut peers = Vec::new();
for nickname in nicknames {
let endpoint = client_endpoint().await;
let id = *endpoint.id().as_bytes();
store
.add(PeerEntry {
endpoint_id: id,
nickname: (*nickname).to_string(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
peers.push(Peer {
principal: format!("eid:{}", endpoint.id()),
endpoint,
});
}
let allow = peers
.iter()
.map(|p| format!("\"{}\"", p.principal))
.collect::<Vec<_>>()
.join(", ");
let config_path = dir.path().join("config.toml");
let toml = format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [{allow}]\n\
\n[services.later]\nrun = ['{STUB}']\nallow = []\n"
);
std::fs::write(&config_path, &toml).unwrap();
let cfg = Config::from_toml_str(&toml).expect("parse config");
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(store.clone()));
let conn_registry = Arc::new(ConnRegistry::new());
let server = server_endpoint().await;
let addr = server.addr();
let mesh = MeshState::new(
server,
gate,
store,
Arc::new(LiveInvites::new()),
"server".into(),
config_path,
Arc::new(RosterGate::empty()),
conn_registry.clone(),
None,
None,
None,
None,
);
mesh.set_accept_task(spawn_accept_loop(
mesh.clone(),
Arc::new(build_services(&cfg)),
))
.await;
(mesh, addr, conn_registry, peers, dir)
}
async fn dial(client: &iroh::Endpoint, addr: iroh::EndpointAddr) -> iroh::endpoint::Connection {
client
.connect(addr, ALPN_MCP)
.await
.expect("dial the mesh ALPN")
}
async fn open_session(
conn: &iroh::endpoint::Connection,
service: &str,
) -> Option<SessionTransport> {
let (mut send, recv) = conn.open_bi().await.ok()?;
write_frame(&mut send, &initialize_frame(service))
.await
.ok()?;
Some(SessionTransport::new(recv, send, MAX_FRAME_BYTES))
}
async fn session_served(transport: Option<&mut SessionTransport>) -> bool {
let Some(transport) = transport else {
return false;
};
match timeout(Duration::from_secs(5), transport.recv_value()).await {
Ok(Ok(Some(v))) => v["result"]["serverInfo"]["name"] == "echo-stub",
_ => false,
}
}
#[tokio::test]
async fn a_grant_is_live_on_an_already_open_connection() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice"]).await;
let alice = &peers[0];
let conn = dial(&alice.endpoint, addr).await;
let mut before = open_session(&conn, "later").await;
assert!(
!session_served(before.as_mut()).await,
"`later` admits nobody yet, so this session must be refused (setup)"
);
mcpmesh::daemon::grant_service_access(
&mesh,
&alice.principal,
&alice.principal,
&["later".to_string()],
)
.await
.expect("grant succeeds");
assert!(
conn.close_reason().is_none(),
"a grant must not disturb the connection — if it closed, this test is not isolating \
the live registry"
);
let mut after = open_session(&conn, "later").await;
assert!(
session_served(after.as_mut()).await,
"a grant must be visible to the NEXT session on an ALREADY-OPEN connection"
);
})
.await
.expect("grant-goes-live test timed out");
}
#[tokio::test]
async fn a_new_session_on_an_open_connection_is_refused_after_revoke() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice"]).await;
let alice = &peers[0];
let conn = dial(&alice.endpoint, addr).await;
let mut first = open_session(&conn, "echo").await;
assert!(
session_served(first.as_mut()).await,
"the granted peer must be served BEFORE the revoke (setup)"
);
revoke_service_allow(&mesh, "echo".into(), alice.principal.clone())
.await
.expect("revoke succeeds");
let mut second = open_session(&conn, "echo").await;
assert!(
!session_served(second.as_mut()).await,
"a revoked peer must not open a NEW session on its already-open connection"
);
})
.await
.expect("new-session-after-revoke test timed out");
}
#[tokio::test]
async fn revoke_severs_the_live_connection() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice"]).await;
let alice = &peers[0];
let conn = dial(&alice.endpoint, addr).await;
let mut session = open_session(&conn, "echo").await;
assert!(
session_served(session.as_mut()).await,
"served before revoke"
);
revoke_service_allow(&mesh, "echo".into(), alice.principal.clone())
.await
.expect("revoke succeeds");
timeout(Duration::from_secs(5), conn.closed())
.await
.expect("a revoke must SEVER the live connection, not leave it running");
})
.await
.expect("sever test timed out");
}
#[tokio::test]
async fn revoke_does_not_sever_an_unrelated_peer() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice", "bob"]).await;
let (alice, bob) = (&peers[0], &peers[1]);
let alice_conn = dial(&alice.endpoint, addr.clone()).await;
let mut alice_session = open_session(&alice_conn, "echo").await;
assert!(session_served(alice_session.as_mut()).await, "alice served");
let bob_conn = dial(&bob.endpoint, addr).await;
let mut bob_session = open_session(&bob_conn, "echo").await;
assert!(session_served(bob_session.as_mut()).await, "bob served");
revoke_service_allow(&mesh, "echo".into(), alice.principal.clone())
.await
.expect("revoke alice");
let bob_session = bob_session.as_mut().expect("bob's session opened");
bob_session
.send_value(tools_call_frame("still-alive"))
.await
.expect("bob's session is still live");
let reply = timeout(Duration::from_secs(5), bob_session.recv_value())
.await
.expect("bob's kept session must answer promptly")
.expect("bob transport ok")
.expect("bob reply frame");
assert_eq!(
reply["result"]["content"][0]["text"], "still-alive",
"revoking alice must NOT sever bob's live session: {reply}"
);
let mut bob_second = open_session(&bob_conn, "echo").await;
assert!(
session_served(bob_second.as_mut()).await,
"revoking alice must not affect bob's new sessions"
);
})
.await
.expect("unrelated-peer regression timed out");
}
#[tokio::test]
async fn removing_a_peer_severs_its_live_connection() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice", "bob"]).await;
let (alice, bob) = (&peers[0], &peers[1]);
let alice_conn = dial(&alice.endpoint, addr.clone()).await;
let mut alice_session = open_session(&alice_conn, "echo").await;
assert!(session_served(alice_session.as_mut()).await, "alice served");
let bob_conn = dial(&bob.endpoint, addr).await;
let mut bob_session = open_session(&bob_conn, "echo").await;
assert!(session_served(bob_session.as_mut()).await, "bob served");
mcpmesh::daemon::revoke_service_access(&mesh, "alice")
.await
.expect("revoke alice's authorization");
timeout(Duration::from_secs(5), alice_conn.closed())
.await
.expect("removing a peer must sever its live connection");
let bob_session = bob_session.as_mut().expect("bob's session opened");
bob_session
.send_value(tools_call_frame("bob-ok"))
.await
.expect("bob's session is still live");
let reply = timeout(Duration::from_secs(5), bob_session.recv_value())
.await
.expect("bob answers")
.expect("bob transport ok")
.expect("bob reply frame");
assert_eq!(reply["result"]["content"][0]["text"], "bob-ok");
})
.await
.expect("peer-remove sever test timed out");
}
#[tokio::test]
async fn the_unpair_path_swaps_the_registry_before_it_severs() {
timeout(Duration::from_secs(60), async {
let (mesh, addr, _registry, peers, _dir) = paired_mesh(&["alice"]).await;
let alice = &peers[0];
let conn = dial(&alice.endpoint, addr).await;
let mut session = open_session(&conn, "echo").await;
assert!(
session_served(session.as_mut()).await,
"served before revoke"
);
let seen: Arc<std::sync::Mutex<Vec<Vec<String>>>> = Arc::new(std::sync::Mutex::new(vec![]));
let sink = seen.clone();
mesh.set_sever_observer(move |live| {
sink.lock().expect("observer sink not poisoned").push(
live.get("echo")
.map(|e| e.allow.clone())
.unwrap_or_default(),
);
});
revoke_service_access(&mesh, "alice")
.await
.expect("unpair revoke succeeds");
let observed = seen.lock().expect("observer sink not poisoned").clone();
assert!(
!observed.is_empty(),
"the unpair must have severed, firing the observer"
);
for at_sever in &observed {
assert!(
!at_sever.contains(&alice.principal),
"the rebuilt registry must be installed BEFORE every sever — at one sever `echo` \
still admitted {at_sever:?}, so the peer that sever cut could have redialled \
straight back in (all observations: {observed:?})"
);
}
})
.await
.expect("unpair swap-before-sever test timed out");
}