use std::sync::Arc;
use std::time::Duration;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::config::Config;
use mcpmesh::daemon::build_services;
use mcpmesh_net::{SessionTransport, TrustGate, connect, serve};
use serde_json::{Value, json};
use tokio::time::timeout;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
async fn local_endpoint() -> iroh::Endpoint {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![mcpmesh_net::ALPN_MCP.to_vec()])
.bind()
.await
.expect("bind localhost endpoint")
}
#[tokio::test]
async fn daemon_serves_run_service_and_injects_caller_identity_over_the_mesh() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"tester\"]\n"
))
.expect("parse config");
let client = local_endpoint().await;
let client_id = *client.id().as_bytes();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
store
.add(PeerEntry {
endpoint_id: client_id,
nickname: "tester".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(Arc::new(store)));
let server = local_endpoint().await;
let addr = server.addr();
let _handle = serve(
server,
gate,
build_services(&cfg),
Arc::new(mcpmesh_net::ConnRegistry::new()),
);
let mut transport = connect(&client, addr, "echo").await.unwrap();
transport
.send_value(json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {
"protocolVersion": "2025-06-18",
"_meta": {"mcpmesh/service": "echo"},
"capabilities": {}, "clientInfo": {"name": "tester", "version": "0"}
}
}))
.await
.unwrap();
let init_res = transport.recv_value().await.unwrap().unwrap();
assert_eq!(
init_res["result"]["serverInfo"]["name"], "echo-stub",
"the served child answered initialize over the mesh"
);
transport
.send_value(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"name": "echo", "arguments": {"text": "over the mesh"}}
}))
.await
.unwrap();
let call_res = transport.recv_value().await.unwrap().unwrap();
assert_eq!(
call_res["result"]["content"][0]["text"], "over the mesh",
"the served child echoed the tools/call payload"
);
assert_eq!(
call_res["result"]["peer_name"], "tester",
"the child's MCPMESH_PEER_NAME carried the gate-resolved nickname across the mesh"
);
})
.await
.expect("daemon serve integration test timed out");
}
async fn send_initialize(transport: &mut SessionTransport, service: &str) -> Value {
transport
.send_value(json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {
"protocolVersion": "2025-06-18",
"_meta": {"mcpmesh/service": service},
"capabilities": {}
}
}))
.await
.unwrap();
transport.recv_value().await.unwrap().unwrap()
}
#[tokio::test]
async fn hot_reload_serves_a_newly_registered_service_over_the_mesh() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let before = local_endpoint().await;
let after = local_endpoint().await;
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
for c in [&before, &after] {
store
.add(PeerEntry {
endpoint_id: *c.id().as_bytes(),
nickname: "tester".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
}
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(Arc::new(store)));
let server = local_endpoint().await;
let handle = serve(
server.clone(),
gate.clone(),
build_services(&Config::from_toml_str("").unwrap()),
Arc::new(mcpmesh_net::ConnRegistry::new()),
);
let mut t_before = connect(&before, server.addr(), "echo").await.unwrap();
let refused = send_initialize(&mut t_before, "echo").await;
assert_eq!(
refused["error"]["code"], -32054,
"echo must be UNSERVED before the reload: {refused}"
);
handle.shutdown();
let cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"tester\"]\n"
))
.unwrap();
let _new_handle = serve(
server.clone(),
gate.clone(),
build_services(&cfg),
Arc::new(mcpmesh_net::ConnRegistry::new()),
);
let mut t_after = connect(&after, server.addr(), "echo").await.unwrap();
let init = send_initialize(&mut t_after, "echo").await;
assert_eq!(
init["result"]["serverInfo"]["name"], "echo-stub",
"echo must be SERVED after the reload: {init}"
);
t_after
.send_value(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"arguments": {"text": "after reload"}}
}))
.await
.unwrap();
let call = t_after.recv_value().await.unwrap().unwrap();
assert_eq!(call["result"]["content"][0]["text"], "after reload");
assert_eq!(
call["result"]["peer_name"], "tester",
"identity injection still holds through the reloaded serve loop"
);
})
.await
.expect("hot-reload serve test timed out");
}