use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant};
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::roster::transport::{self, RosterBlobs};
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 roster_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(),
transport::GOSSIP_ALPN.to_vec(),
transport::BLOB_ALPN.to_vec(),
])
.bind()
.await
.expect("bind roster-mode endpoint")
}
fn mint_signed_roster(
root: &SigningKey,
serial: u64,
users: &[([u8; 32], &str, &str)],
revoked: &[[u8; 32]],
) -> (Vec<u8>, RosterView) {
let roster_users = users
.iter()
.map(|(eid, uid, label)| 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: (*label).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(),
},
);
let bytes = serde_json::to_vec(&r).expect("serialize signed roster");
let view = load_installed(&r, &root.verifying_key()).expect("mint a valid roster view");
(bytes, view)
}
fn write_config(path: &Path, org_root_pk: &str, services_toml: &str) {
std::fs::write(
path,
format!(
"[network]\nrelay_mode = \"disabled\"\n[identity]\norg_root_pk = \"{org_root_pk}\"\norg_id = \"acme\"\n{services_toml}"
),
)
.expect("write roster config");
}
fn seed_addrs(ep: &iroh::Endpoint, peers: &[&iroh::Endpoint]) {
let mem = iroh::address_lookup::MemoryLookup::new();
for p in peers {
mem.add_endpoint_info(p.addr());
}
ep.address_lookup()
.expect("address lookup services")
.add(mem);
}
async fn build_node(
endpoint: iroh::Endpoint,
install_view: RosterView,
bootstrap: Vec<iroh::EndpointId>,
config_path: PathBuf,
store: Arc<PeerStore>,
conn_registry: Arc<ConnRegistry>,
) -> (Arc<MeshState>, Arc<RosterGate>) {
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let roster = Arc::new(RosterGate::empty());
roster.install(install_view);
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let gossip = transport::spawn_gossip(&endpoint);
let blobs = RosterBlobs::new(&endpoint);
let roster_gossip =
transport::subscribe(&gossip, transport::roster_topic_bytes("acme"), bootstrap)
.await
.expect("subscribe roster topic");
let mesh = MeshState::new(
endpoint,
gate,
store,
Arc::new(LiveInvites::new()),
"node".into(),
config_path,
roster.clone(),
conn_registry,
Some(gossip),
Some(blobs),
Some(roster_gossip),
None,
);
(mesh, roster)
}
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}}
})
}
async fn wait_for_len(registry: &ConnRegistry, target: usize) {
for _ in 0..100 {
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(flavor = "multi_thread")]
async fn revoked_device_cut_from_live_session_within_60s_across_nodes() {
timeout(Duration::from_secs(240), async {
let dir_o = tempfile::tempdir().unwrap();
let dir_a = tempfile::tempdir().unwrap();
let dir_b = tempfile::tempdir().unwrap();
let root = SigningKey::from_bytes(&[9u8; 32]);
let org_root_pk = encode_b64u(&root.verifying_key().to_bytes());
let ep_o = roster_alpn_endpoint().await;
let ep_a = roster_alpn_endpoint().await;
let ep_b = roster_alpn_endpoint().await;
let eid_o = ep_o.id();
let eid_a = ep_a.id();
let eid_b = ep_b.id();
let bytes_o = *eid_o.as_bytes();
let bytes_a = *eid_a.as_bytes();
let bytes_b = *eid_b.as_bytes();
let alice_dialer = ep_a.clone(); let b_addr = ep_b.addr();
seed_addrs(&ep_o, &[&ep_a, &ep_b]);
seed_addrs(&ep_a, &[&ep_o, &ep_b]);
seed_addrs(&ep_b, &[&ep_o, &ep_a]);
let users1 = [
(bytes_o, "operator", "console"),
(bytes_a, "alice", "laptop"),
(bytes_b, "hostb", "server"),
];
let (v1_bytes, v1_o) = mint_signed_roster(&root, 1, &users1, &[]);
let v1_a = mint_signed_roster(&root, 1, &users1, &[]).1;
let v1_b = mint_signed_roster(&root, 1, &users1, &[]).1;
let (v2_bytes, v2_o) = mint_signed_roster(&root, 2, &users1, &[bytes_a]);
let cfg_o_path = dir_o.path().join("config.toml");
let cfg_a_path = dir_a.path().join("config.toml");
let cfg_b_path = dir_b.path().join("config.toml");
write_config(&cfg_o_path, &org_root_pk, "");
write_config(&cfg_a_path, &org_root_pk, "");
write_config(
&cfg_b_path,
&org_root_pk,
&format!("[services.echo]\nrun = ['{STUB}']\nallow = [\"alice\", \"alice-pair\"]\n"),
);
std::fs::write(dir_o.path().join("roster.json"), &v1_bytes).unwrap();
std::fs::write(dir_b.path().join("roster.json"), &v1_bytes).unwrap();
let store_o = Arc::new(PeerStore::open(&dir_o.path().join("state.redb")).unwrap());
let store_a = Arc::new(PeerStore::open(&dir_a.path().join("state.redb")).unwrap());
let store_b = Arc::new(PeerStore::open(&dir_b.path().join("state.redb")).unwrap());
store_b
.add(PeerEntry {
endpoint_id: bytes_a,
nickname: "alice-pair".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let reg_o = Arc::new(ConnRegistry::new());
let reg_a = Arc::new(ConnRegistry::new());
let reg_b = Arc::new(ConnRegistry::new());
let (mesh_o, roster_o) = build_node(
ep_o,
v1_o,
vec![eid_a, eid_b],
cfg_o_path.clone(),
store_o,
reg_o,
)
.await;
let (mesh_a, _roster_a) = build_node(
ep_a,
v1_a,
vec![eid_o, eid_b],
cfg_a_path.clone(),
store_a,
reg_a,
)
.await;
let (mesh_b, roster_b) = build_node(
ep_b,
v1_b,
vec![eid_o, eid_a],
cfg_b_path.clone(),
store_b,
reg_b.clone(),
)
.await;
let cfg_o = Config::load(&cfg_o_path).expect("load O config");
let cfg_a = Config::load(&cfg_a_path).expect("load A config");
let cfg_b = Config::load(&cfg_b_path).expect("load B config");
let _task_o = spawn_accept_loop(mesh_o.clone(), Arc::new(build_services(&cfg_o)));
let _task_a = spawn_accept_loop(mesh_a.clone(), Arc::new(build_services(&cfg_a)));
let _task_b = spawn_accept_loop(mesh_b.clone(), Arc::new(build_services(&cfg_b)));
let _recv_b = mcpmesh::roster::distribute::spawn_receive_loop(mesh_b.clone());
let mut alice_t = connect(&alice_dialer, b_addr.clone(), "echo")
.await
.expect("A dials B's echo service").0;
alice_t.send_value(initialize_frame("echo")).await.unwrap();
let init = alice_t.recv_value().await.unwrap().unwrap();
assert_eq!(
init["result"]["serverInfo"]["name"], "echo-stub",
"A must complete a LIVE session to B before the revoke: {init}"
);
alice_t
.send_value(tools_call_frame("live-before-revoke"))
.await
.unwrap();
let live = timeout(Duration::from_secs(10), alice_t.recv_value())
.await
.expect("A's pre-revoke frame must round-trip promptly")
.expect("A transport ok")
.expect("A reply frame");
assert_eq!(
live["result"]["content"][0]["text"], "live-before-revoke",
"the pre-revoke session must GENUINELY round-trip a data frame: {live}"
);
wait_for_len(®_b, 1).await;
std::fs::write(dir_o.path().join("roster.json"), &v2_bytes).unwrap();
let _severed_o = install_roster_view_and_sever(&mesh_o, v2_o);
assert_eq!(
roster_o.view().unwrap().serial(),
2,
"O bumped its own view to serial 2 before announcing"
);
let propagation_start = Instant::now();
let deadline = propagation_start + Duration::from_secs(150);
loop {
mcpmesh::roster::distribute::announce_roster(&mesh_o)
.await
.expect("O announces serial 2 on the roster topic");
if roster_b.view().map(|v| v.serial()).unwrap_or(0) >= 2 {
break;
}
assert!(
Instant::now() < deadline,
"B did NOT converge to serial 2 via gossip propagation within the deadline"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
let elapsed = propagation_start.elapsed();
eprintln!("[T13] B converged to serial 2 via gossip in {elapsed:?}");
assert_eq!(
roster_b.view().expect("B installed a roster").serial(),
2,
"B converged to serial 2 (received → fetched → validated → installed) via gossip"
);
let alice_after = timeout(Duration::from_secs(10), alice_t.recv_value())
.await
.expect("A's severed session must close promptly after B's gossip-received revocation, not hang");
assert!(
!matches!(alice_after, Ok(Some(_))),
"A's live session must be CUT (EOF/close) by B's gossip-received revocation, got a frame: {alice_after:?}"
);
match connect(&alice_dialer, b_addr, "echo").await.map(|(t, _)| t) {
Err(_) => {} Ok(mut t) => {
let _ = t.send_value(initialize_frame("echo")).await;
let res = timeout(Duration::from_secs(10), t.recv_value())
.await
.expect("a revoked re-dial must close promptly, not hang");
assert!(
!matches!(res, Ok(Some(_))),
"a revoked endpoint holding a stale pair entry must be refused pre-MCP, got: {res:?}"
);
}
}
wait_for_len(®_b, 0).await;
})
.await
.expect("3-node revocation propagation exceeded 90s");
}