#![cfg(unix)]
use std::sync::Arc;
use std::time::Duration;
use iroh::address_lookup::MemoryLookup;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::audit::{AuditLog, AuditSink};
use mcpmesh::client::connect_control;
use mcpmesh::config::Config;
use mcpmesh::control::{DaemonState, serve_control};
use mcpmesh::daemon::{self, MeshState, STACK_VERSION, build_services_audited};
use mcpmesh::limits::MeshLimiters;
use mcpmesh::pairing::LiveInvites;
use mcpmesh::roster::gate::RosterGate;
use mcpmesh_net::framing::{FrameReader, Inbound, write_frame};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, TrustGate, serve};
use serde_json::{Value, json};
use tokio::io::BufReader;
use tokio::net::UnixStream;
use tokio::net::unix::{OwnedReadHalf, OwnedWriteHalf};
use tokio::time::timeout;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
const MAX_FRAME: usize = 16 * 1024 * 1024;
async fn local_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 localhost endpoint")
}
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,
)
}
struct SubClient {
reader: FrameReader<BufReader<OwnedReadHalf>>,
_write_half: OwnedWriteHalf,
}
impl SubClient {
async fn connect(socket: &std::path::Path) -> Self {
let stream = UnixStream::connect(socket).await.expect("connect control");
let (read_half, mut write_half) = stream.into_split();
let mut reader = FrameReader::new(BufReader::new(read_half), MAX_FRAME);
match reader.next().await.expect("hello read") {
Some(Inbound::Frame(_hello)) => {}
other => panic!("expected Hello, got {other:?}"),
}
write_frame(&mut write_half, &json!({ "method": "subscribe" }))
.await
.expect("send subscribe");
Self {
reader,
_write_half: write_half,
}
}
async fn next(&mut self) -> Value {
match timeout(Duration::from_secs(5), self.reader.next())
.await
.expect("stream frame within timeout")
.expect("stream read")
{
Some(Inbound::Frame(v)) => v,
other => panic!("expected a stream frame, got {other:?}"),
}
}
}
#[tokio::test(flavor = "multi_thread")]
async fn subscribe_pushes_snapshot_then_live_session_events() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let server_ep = local_endpoint().await;
let server_id = *server_ep.id().as_bytes();
let server_addr = server_ep.addr();
let daemon_ep = local_endpoint().await;
let daemon_id = *daemon_ep.id().as_bytes();
let daemon_eid = format!("eid:{}", daemon_ep.id());
let server_cfg = Config::from_toml_str(&format!(
"[services.echo]\nrun = ['{STUB}']\nallow = [\"{daemon_eid}\"]\n"
))
.expect("parse server config");
let server_store = Arc::new(PeerStore::open(&dir.path().join("server.redb")).unwrap());
server_store
.add(PeerEntry {
endpoint_id: daemon_id,
nickname: "daemon".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let audit = AuditSink::new(AuditLog::spawn(dir.path().join("audit")));
let limiters = MeshLimiters::unlimited();
let server_gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(server_store.clone()));
let _serve = serve(
server_ep.clone(),
server_gate,
build_services_audited(&server_cfg, &audit, &limiters),
Arc::new(ConnRegistry::new()),
);
let s_mesh = assemble_mesh(server_ep, server_store, config.clone());
s_mesh.set_audit(audit.clone());
let s_socket = dir.path().join("s.sock");
let s_listener = mcpmesh::ipc::bind_control_socket(&s_socket).await.unwrap();
let s_state = Arc::new(DaemonState::with_mesh(STACK_VERSION, s_mesh));
let s_control = tokio::spawn(serve_control(s_listener, s_state));
let mem = MemoryLookup::new();
mem.add_endpoint_info(server_addr);
daemon_ep
.address_lookup()
.expect("address lookup services")
.add(mem);
let daemon_store = Arc::new(PeerStore::open(&dir.path().join("daemon.redb")).unwrap());
daemon_store
.add(PeerEntry {
endpoint_id: server_id,
nickname: "tester".into(),
services: vec!["echo".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let d_socket = dir.path().join("d.sock");
let d_listener = mcpmesh::ipc::bind_control_socket(&d_socket).await.unwrap();
let d_state = daemon::serving_state(daemon_ep, daemon_store);
let d_control = tokio::spawn(serve_control(d_listener, d_state));
let mut sub = SubClient::connect(&s_socket).await;
let snapshot = sub.next().await;
assert_eq!(
snapshot["type"], "snapshot",
"the first pushed frame must be a snapshot: {snapshot}"
);
{
let client = connect_control(&d_socket)
.await
.expect("connect to D control");
let (mut reader, mut writer) = client
.open_session("tester".into(), "echo".into())
.await
.expect("open_session upgrade");
write_frame(
&mut writer,
&json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": "2025-11-25", "capabilities": {},
"clientInfo": {"name": "ai", "version": "0"}}
}),
)
.await
.expect("send initialize");
let init = match timeout(Duration::from_secs(10), reader.next())
.await
.expect("initialize response within timeout")
.expect("initialize read")
{
Some(Inbound::Frame(v)) => v,
other => panic!("expected initialize response, got {other:?}"),
};
assert_eq!(
init["result"]["serverInfo"]["name"], "echo-stub",
"the served child answered initialize over the mesh: {init}"
);
}
let mut saw_open = false;
let mut saw_close = false;
for _ in 0..50 {
let f = sub.next().await;
if f["type"] == "event" && f["record"]["kind"] == "session_open" {
saw_open = true;
}
if f["type"] == "event" && f["record"]["kind"] == "session_close" {
saw_close = true;
break;
}
}
assert!(
saw_open && saw_close,
"the live stream must carry session_open then session_close events (open={saw_open}, close={saw_close})"
);
s_control.abort();
d_control.abort();
std::mem::forget(dir);
})
.await
.expect("subscribe test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn dial_failure_emits_error_event() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let ep = local_endpoint().await;
let store = Arc::new(PeerStore::open(&dir.path().join("node.redb")).unwrap());
let mesh = assemble_mesh(ep, store, config.clone());
let audit = AuditSink::new(AuditLog::spawn(dir.path().join("audit")));
mesh.set_audit(audit.clone());
let socket = dir.path().join("node.sock");
let listener = mcpmesh::ipc::bind_control_socket(&socket).await.unwrap();
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, mesh));
let control = tokio::spawn(serve_control(listener, state));
let mut sub = SubClient::connect(&socket).await;
let snapshot = sub.next().await;
assert_eq!(
snapshot["type"], "snapshot",
"the first pushed frame must be a snapshot: {snapshot}"
);
{
let client = connect_control(&socket).await.expect("connect control");
let (mut reader, _writer) = client
.open_session("ghost".into(), "nope".into())
.await
.expect("open_session upgrade");
let _ = timeout(Duration::from_secs(5), reader.next()).await;
}
let mut saw_error = false;
for _ in 0..50 {
let f = sub.next().await;
if f["type"] == "event"
&& f["record"]["kind"] == "session_open"
&& f["record"]["status"] == "error"
{
assert_eq!(
f["record"]["peer"], "ghost",
"the error record must name the requested dial target: {f}"
);
saw_error = true;
break;
}
}
assert!(
saw_error,
"a failed dial must emit a session_open event with status=error on the live stream"
);
control.abort();
std::mem::forget(dir);
})
.await
.expect("dial-failure test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn reachability_flips_are_pushed_to_subscribers() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let peer_ep = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![ALPN_MCP.to_vec(), mcpmesh_net::ALPN_PING.to_vec()])
.bind()
.await
.expect("bind peer endpoint");
let peer_id = *peer_ep.id().as_bytes();
let peer_addr = peer_ep.addr();
let our_ep = local_endpoint().await;
let our_id = *our_ep.id().as_bytes();
let mem = MemoryLookup::new();
mem.add_endpoint_info(peer_addr);
let our_ep = our_ep;
let peer_store = Arc::new(PeerStore::open(&dir.path().join("peer.redb")).unwrap());
peer_store
.add(PeerEntry {
endpoint_id: our_id,
nickname: "us".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let our_store = Arc::new(PeerStore::open(&dir.path().join("our.redb")).unwrap());
our_store
.add(PeerEntry {
endpoint_id: peer_id,
nickname: "bob".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: Some(
serde_json::to_string(&peer_ep.addr()).expect("serialize peer addr"),
),
})
.unwrap();
let peer_mesh = assemble_mesh(peer_ep, peer_store, dir.path().join("peer.toml"));
std::fs::write(dir.path().join("peer.toml"), "").unwrap();
let peer_accept = daemon::spawn_accept_loop(
peer_mesh.clone(),
Arc::new(build_services_audited(
&Config::default(),
&AuditSink::disabled(),
&MeshLimiters::unlimited(),
)),
);
let mesh = assemble_mesh(our_ep, our_store, config);
mesh.set_audit(AuditSink::new(AuditLog::spawn(dir.path().join("audit"))));
let socket = dir.path().join("s.sock");
let listener = mcpmesh::ipc::bind_control_socket(&socket).await.unwrap();
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, mesh.clone()));
let _control = tokio::spawn(serve_control(listener, state));
let mut sub = SubClient::connect(&socket).await;
assert_eq!(sub.next().await["type"], "snapshot");
let up = daemon::probe_peer(&mesh, peer_id).await;
assert!(up.reachable, "the live peer must probe reachable");
let mut path = up.path.clone();
for _ in 0..10 {
if path == mcpmesh_local_api::PeerPath::Direct {
break;
}
assert_ne!(
path,
mcpmesh_local_api::PeerPath::Relay { url: None },
"no relay is configured, so a relay verdict would be wrong"
);
tokio::time::sleep(Duration::from_millis(300)).await;
path = daemon::probe_peer(&mesh, peer_id).await.path;
}
assert_eq!(
path,
mcpmesh_local_api::PeerPath::Direct,
"a loopback peer with relays disabled must settle on Direct"
);
let frame = sub.next().await;
assert_eq!(frame["type"], "reachability", "got {frame}");
assert_eq!(frame["peer"]["reachable"], true, "came online: {frame}");
assert_eq!(frame["peer"]["name"], "bob", "got {frame}");
assert_ne!(frame["peer"]["path"]["kind"], "relay", "got {frame}");
assert_eq!(
frame["peer"]["principal"],
format!(
"eid:{}",
peer_id
.iter()
.map(|b| format!("{b:02x}"))
.collect::<String>()
),
"got {frame}"
);
let _ = daemon::probe_peer(&mesh, peer_id).await;
assert!(
timeout(Duration::from_millis(600), sub.reader.next())
.await
.is_err(),
"an unchanged re-probe must push nothing"
);
peer_accept.abort();
drop(peer_mesh);
let down = daemon::probe_peer(&mesh, peer_id).await;
assert!(!down.reachable, "the dead peer must probe unreachable");
let frame = sub.next().await;
assert_eq!(frame["type"], "reachability", "got {frame}");
assert_eq!(frame["peer"]["reachable"], false, "went offline: {frame}");
assert_eq!(frame["peer"]["name"], "bob", "got {frame}");
let _ = mem;
})
.await
.expect("reachability flip test timed out");
}