use std::sync::Arc;
use std::time::Duration;
use ed25519_dalek::SigningKey;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::config::Config;
use mcpmesh::daemon::{
MeshState, build_services, install_roster_view_and_sever, spawn_accept_loop,
};
use mcpmesh::pairing::LiveInvites;
use mcpmesh::roster::gate::{ComposedGate, RosterGate};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, ALPN_PAIR, TrustGate, connect};
use mcpmesh_trust::roster::sign::mint_signed;
use mcpmesh_trust::roster::validate::{RosterView, load_installed};
use mcpmesh_trust::roster::{Roster, RosterDevice, RosterUser, encode_b64u};
use serde_json::json;
use tokio::time::timeout;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
async fn dual_alpn_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()])
.bind()
.await
.expect("bind dual-ALPN 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 mint_view(
root: &SigningKey,
serial: u64,
users: &[([u8; 32], &str)],
revoked: &[[u8; 32]],
) -> RosterView {
let roster_users = users
.iter()
.map(|(eid, uid)| RosterUser {
user_id: (*uid).into(),
display_name: (*uid).into(),
user_pk: encode_b64u(&[1u8; 32]),
groups: vec!["team-eng".into()],
devices: vec![RosterDevice {
endpoint_id: encode_b64u(eid),
label: "device".into(),
role: "primary".into(),
}],
})
.collect();
let r = mint_signed(
root,
Roster {
format: "mcpmesh-roster/1".into(),
org_id: "acme".into(),
serial,
issued_at: "2000-01-01T00:00:00Z".into(),
expires_at: "2999-01-01T00:00:00Z".into(),
groups: vec!["team-eng".into()],
users: roster_users,
revoked_endpoints: revoked.iter().map(|e| encode_b64u(e)).collect(),
sig: String::new(),
},
);
load_installed(&r, &root.verifying_key()).expect("mint a valid roster view")
}
fn initialize_frame(service: &str) -> serde_json::Value {
json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {
"protocolVersion": "2025-06-18",
"_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}}
})
}
async fn wait_for_len(registry: &ConnRegistry, target: usize) {
for _ in 0..50 {
if registry.len() == target {
return;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
panic!(
"conn registry len did not settle to {target} (still {})",
registry.len()
);
}
#[tokio::test]
async fn install_severs_a_revoked_roster_session_but_not_a_pairing_session() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"alice\", \"bob\"]\n"
))
.expect("parse config");
let root = SigningKey::from_bytes(&[9u8; 32]);
let alice_client = client_endpoint().await;
let bob_client = client_endpoint().await;
let alice_id = *alice_client.id().as_bytes();
let bob_id = *bob_client.id().as_bytes();
let store = Arc::new(PeerStore::open(&dir.path().join("state.redb")).unwrap());
store
.add(PeerEntry {
endpoint_id: alice_id,
nickname: "alice-stale".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
store
.add(PeerEntry {
endpoint_id: bob_id,
nickname: "bob".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let roster = Arc::new(RosterGate::empty());
roster.install(mint_view(&root, 1, &[(alice_id, "alice")], &[]));
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let conn_registry = Arc::new(ConnRegistry::new());
let server = dual_alpn_endpoint().await;
let addr = server.addr();
let mesh = MeshState::new(
server,
gate,
store,
Arc::new(LiveInvites::new()),
"server".into(),
dir.path().join("config.toml"),
roster.clone(),
conn_registry.clone(),
None,
None,
None,
None,
);
let _task = spawn_accept_loop(mesh.clone(), Arc::new(build_services(&cfg)));
let mut alice_t = connect(&alice_client, addr.clone(), "echo").await.unwrap();
alice_t.send_value(initialize_frame("echo")).await.unwrap();
let alice_init = alice_t.recv_value().await.unwrap().unwrap();
assert_eq!(
alice_init["result"]["serverInfo"]["name"], "echo-stub",
"rostered peer must complete a session under the composed gate: {alice_init}"
);
let mut bob_t = connect(&bob_client, addr.clone(), "echo").await.unwrap();
bob_t.send_value(initialize_frame("echo")).await.unwrap();
let bob_init = bob_t.recv_value().await.unwrap().unwrap();
assert_eq!(
bob_init["result"]["serverInfo"]["name"], "echo-stub",
"pairing peer must complete a session: {bob_init}"
);
wait_for_len(&conn_registry, 2).await;
let revoking = mint_view(&root, 2, &[(alice_id, "alice")], &[alice_id]);
let severed = install_roster_view_and_sever(&mesh, revoking);
assert_eq!(
severed, 1,
"exactly the revoked rostered session is severed"
);
let alice_after = timeout(Duration::from_secs(5), alice_t.recv_value())
.await
.expect("alice's severed session must close promptly, not hang");
assert!(
!matches!(alice_after, Ok(Some(_))),
"the revoked rostered session must be severed (closed), got a frame: {alice_after:?}"
);
bob_t
.send_value(tools_call_frame("still-alive"))
.await
.expect("bob's session is still live");
let bob_reply = timeout(Duration::from_secs(5), bob_t.recv_value())
.await
.expect("bob's kept session must answer promptly")
.expect("bob transport ok")
.expect("bob reply frame");
assert_eq!(
bob_reply["result"]["content"][0]["text"], "still-alive",
"the pairing-only session must NOT be severed by the roster install: {bob_reply}"
);
wait_for_len(&conn_registry, 1).await;
})
.await
.expect("D8 sever test timed out");
}
#[tokio::test]
async fn a_peer_the_roster_revokes_is_refused_pre_mcp() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"alice\"]\n"
))
.expect("parse config");
let root = SigningKey::from_bytes(&[9u8; 32]);
let client = client_endpoint().await;
let client_id = *client.id().as_bytes();
let store = Arc::new(PeerStore::open(&dir.path().join("state.redb")).unwrap());
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let roster = Arc::new(RosterGate::empty());
roster.install(mint_view(&root, 1, &[(client_id, "alice")], &[client_id]));
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let conn_registry = Arc::new(ConnRegistry::new());
let server = dual_alpn_endpoint().await;
let addr = server.addr();
let mesh = MeshState::new(
server,
gate,
store,
Arc::new(LiveInvites::new()),
"server".into(),
dir.path().join("config.toml"),
roster.clone(),
conn_registry.clone(),
None,
None,
None,
None,
);
let _task = spawn_accept_loop(mesh.clone(), Arc::new(build_services(&cfg)));
match connect(&client, addr, "echo").await {
Err(_) => {} Ok(mut transport) => {
let _ = transport.send_value(initialize_frame("echo")).await;
let res = timeout(Duration::from_secs(5), transport.recv_value())
.await
.expect("a revoked dial must close promptly, not hang");
assert!(
!matches!(res, Ok(Some(_))),
"a peer the roster revokes must be refused pre-MCP (no session frame), got: {res:?}"
);
}
}
wait_for_len(&conn_registry, 0).await;
})
.await
.expect("TOCTOU pre-MCP refusal test timed out");
}
#[tokio::test]
async fn install_severs_a_dropped_roster_session_but_keeps_a_still_listed_one() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"carol\", \"dave\"]\n"
))
.expect("parse config");
let root = SigningKey::from_bytes(&[9u8; 32]);
let carol_client = client_endpoint().await;
let dave_client = client_endpoint().await;
let carol_id = *carol_client.id().as_bytes();
let dave_id = *dave_client.id().as_bytes();
let store = Arc::new(PeerStore::open(&dir.path().join("state.redb")).unwrap());
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let roster = Arc::new(RosterGate::empty());
roster.install(mint_view(
&root,
1,
&[(carol_id, "carol"), (dave_id, "dave")],
&[],
));
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let conn_registry = Arc::new(ConnRegistry::new());
let server = dual_alpn_endpoint().await;
let addr = server.addr();
let mesh = MeshState::new(
server,
gate,
store,
Arc::new(LiveInvites::new()),
"server".into(),
dir.path().join("config.toml"),
roster.clone(),
conn_registry.clone(),
None,
None,
None,
None,
);
let _task = spawn_accept_loop(mesh.clone(), Arc::new(build_services(&cfg)));
let mut carol_t = connect(&carol_client, addr.clone(), "echo").await.unwrap();
carol_t.send_value(initialize_frame("echo")).await.unwrap();
let carol_init = carol_t.recv_value().await.unwrap().unwrap();
assert_eq!(carol_init["result"]["serverInfo"]["name"], "echo-stub");
let mut dave_t = connect(&dave_client, addr.clone(), "echo").await.unwrap();
dave_t.send_value(initialize_frame("echo")).await.unwrap();
let dave_init = dave_t.recv_value().await.unwrap().unwrap();
assert_eq!(dave_init["result"]["serverInfo"]["name"], "echo-stub");
wait_for_len(&conn_registry, 2).await;
let dropped = mint_view(&root, 2, &[(dave_id, "dave")], &[]);
let severed = install_roster_view_and_sever(&mesh, dropped);
assert_eq!(
severed, 1,
"the dropped (roster-resolved, now-absent, NOT revoked) session is severed"
);
let carol_after = timeout(Duration::from_secs(5), carol_t.recv_value())
.await
.expect("carol's dropped session must close promptly, not hang");
assert!(
!matches!(carol_after, Ok(Some(_))),
"a roster-resolved endpoint dropped from the new roster must be severed, got: {carol_after:?}"
);
dave_t
.send_value(tools_call_frame("still-listed"))
.await
.expect("dave's session is still live");
let dave_reply = timeout(Duration::from_secs(5), dave_t.recv_value())
.await
.expect("dave's kept session must answer promptly")
.expect("dave transport ok")
.expect("dave reply frame");
assert_eq!(
dave_reply["result"]["content"][0]["text"], "still-listed",
"a still-listed roster device must NOT be severed by the same install: {dave_reply}"
);
wait_for_len(&conn_registry, 1).await;
})
.await
.expect("dropped-branch sever test timed out");
}