use std::sync::Arc;
use std::time::Duration;
use anyhow::Result;
use mcpmesh_net::ALPN_PING;
use mcpmesh_net::framing::{FrameReader, Inbound, write_frame};
use crate::util::epoch_now_i64;
use super::MeshState;
#[derive(Clone)]
pub struct ReachEntry {
pub reachable: bool,
pub rtt_ms: Option<u64>,
pub probed_at: i64,
pub meta: String,
pub services: Vec<String>,
pub seq: u64,
pub path: mcpmesh_local_api::PeerPath,
}
pub const REACH_TTL_SECS: i64 = 20;
const PROBE_TIMEOUT: Duration = Duration::from_secs(3);
pub async fn probe_peer(mesh: &Arc<MeshState>, endpoint_id: [u8; 32]) -> ReachEntry {
let seq = mesh
.probe_seq
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let started = std::time::Instant::now();
let outcome = tokio::time::timeout(PROBE_TIMEOUT, probe_once(mesh, endpoint_id)).await;
let (reachable, meta, services, path) = match outcome {
Ok(Ok((meta, services, path))) => (true, meta, services, path),
_ => (
false,
String::new(),
Vec::new(),
mcpmesh_local_api::PeerPath::Unknown,
),
};
let entry = ReachEntry {
reachable,
rtt_ms: reachable.then(|| started.elapsed().as_millis() as u64),
probed_at: epoch_now_i64(),
meta,
services,
seq,
path,
};
let outcome = {
let mut cache = mesh
.reachability
.lock()
.expect("reachability lock not poisoned");
match cache.get(&endpoint_id) {
Some(newer) if !supersedes(seq, newer) => Outcome::Superseded(newer.clone()),
other => {
let previous = other.cloned();
cache.insert(endpoint_id, entry.clone());
Outcome::Committed(previous)
}
}
};
let previous = match outcome {
Outcome::Superseded(newer) => return newer,
Outcome::Committed(previous) => previous,
};
if is_transition(previous.as_ref(), &entry) {
if let Some(row) = stored_row(mesh, endpoint_id, &entry) {
let _ = mesh.reach_bcast.send(row);
}
}
entry
}
fn supersedes(seq: u64, existing: &ReachEntry) -> bool {
seq >= existing.seq
}
enum Outcome {
Committed(Option<ReachEntry>),
Superseded(ReachEntry),
}
const PATH_SETTLE: Duration = Duration::from_millis(600);
async fn settled_path(conn: &iroh::endpoint::Connection) -> mcpmesh_local_api::PeerPath {
let deadline = tokio::time::Instant::now() + PATH_SETTLE;
loop {
let path = selected_path(conn);
if path == mcpmesh_local_api::PeerPath::Direct || tokio::time::Instant::now() >= deadline {
return path;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
fn selected_path(conn: &iroh::endpoint::Connection) -> mcpmesh_local_api::PeerPath {
let paths = conn.paths();
for path in &paths {
if !path.is_selected() {
continue;
}
return match path.remote_addr() {
iroh::TransportAddr::Relay(url) => mcpmesh_local_api::PeerPath::Relay {
url: Some(sanitize_relay_url(url)),
},
iroh::TransportAddr::Ip(_) => mcpmesh_local_api::PeerPath::Direct,
_ => mcpmesh_local_api::PeerPath::Unknown,
};
}
mcpmesh_local_api::PeerPath::Unknown
}
fn sanitize_relay_url(url: &iroh::RelayUrl) -> String {
let u: &url::Url = url; match (u.host_str(), u.port()) {
(Some(host), Some(port)) => format!("{}://{host}:{port}", u.scheme()),
(Some(host), None) => format!("{}://{host}", u.scheme()),
(None, _) => u.scheme().to_string(),
}
}
fn is_transition(previous: Option<&ReachEntry>, current: &ReachEntry) -> bool {
match previous {
None => current.reachable,
Some(prev) => prev.reachable != current.reachable,
}
}
fn reachability_row(
nickname: String,
endpoint_id: [u8; 32],
entry: Option<&ReachEntry>,
age_secs: Option<u64>,
) -> mcpmesh_local_api::PeerReachability {
mcpmesh_local_api::PeerReachability {
path: entry.map(|e| e.path.clone()).unwrap_or_default(),
name: nickname,
reachable: entry.is_some_and(|e| e.reachable),
rtt_ms: entry.and_then(|e| e.rtt_ms),
age_secs,
meta: entry.map(|e| e.meta.clone()).unwrap_or_default(),
principal: Some(mcpmesh_net::EndpointId::from_bytes(endpoint_id).principal()),
}
}
fn stored_row(
mesh: &Arc<MeshState>,
endpoint_id: [u8; 32],
entry: &ReachEntry,
) -> Option<mcpmesh_local_api::PeerReachability> {
let nickname = mesh.store.resolve(&endpoint_id).ok().flatten()?.nickname;
Some(reachability_row(
nickname,
endpoint_id,
Some(entry),
Some(0),
))
}
async fn probe_once(
mesh: &Arc<MeshState>,
endpoint_id: [u8; 32],
) -> Result<(String, Vec<String>, mcpmesh_local_api::PeerPath)> {
let id = iroh::EndpointId::from_bytes(&endpoint_id)
.map_err(|e| anyhow::anyhow!("invalid endpoint id: {e}"))?;
let store = mesh.store.clone();
let last_addr = tokio::task::spawn_blocking(move || store.resolve(&endpoint_id))
.await
.map_err(|e| anyhow::anyhow!("join peer resolve for probe: {e}"))?
.ok()
.flatten()
.and_then(|e| e.last_addr);
let addr = super::dial::stored_dial_addr(last_addr.as_deref(), id);
let conn = mesh.endpoint.connect(addr, ALPN_PING).await?;
let (mut send, recv) = conn.open_bi().await?;
write_frame(&mut send, &serde_json::json!({ "ping": true })).await?;
let _ = send.finish();
let mut reader = FrameReader::new(
tokio::io::BufReader::new(recv),
mcpmesh_net::framing::MAX_FRAME_BYTES,
);
match reader.next().await? {
Some(Inbound::Frame(v)) => {
Ok((pong_meta(&v), pong_services(&v), settled_path(&conn).await))
}
_ => anyhow::bail!("no pong from peer"),
}
}
fn pong_services(pong: &serde_json::Value) -> Vec<String> {
pong.get("services")
.and_then(|s| s.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default()
}
fn pong_meta(pong: &serde_json::Value) -> String {
pong.get("meta")
.and_then(|m| m.as_str())
.filter(|s| s.len() <= crate::roster::presence::APP_METADATA_MAX_BYTES)
.unwrap_or_default()
.to_string()
}
pub(crate) fn caller_admitted_services(
mesh: &Arc<MeshState>,
identity: &mcpmesh_net::PeerIdentity,
) -> Vec<String> {
use std::collections::HashSet;
let eid = identity.endpoint.principal();
let principals: HashSet<&str> =
mcpmesh_local_api::principal_set(Some(&eid), identity.user_id.as_deref(), &identity.groups)
.into_iter()
.collect();
let admits = |allow: &[String]| allow.iter().any(|a| principals.contains(a.as_str()));
let mut out: Vec<String> = Vec::new();
if let Ok(cfg) = crate::config::Config::load(&mesh.config_path) {
for (name, svc) in &cfg.services {
if admits(&svc.allow) {
out.push(name.clone());
}
}
}
for (name, eph) in mesh
.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned")
.iter()
{
if admits(&eph.allow) && !out.contains(name) {
out.push(name.clone());
}
}
out.sort();
out
}
pub fn reachability_of(mesh: &Arc<MeshState>) -> Vec<mcpmesh_local_api::PeerReachability> {
let now = epoch_now_i64();
let peers: Vec<(String, [u8; 32])> = mesh
.store
.list()
.unwrap_or_default()
.into_iter()
.map(|e| (e.nickname, e.endpoint_id))
.collect();
let cache = mesh
.reachability
.lock()
.expect("reachability lock not poisoned")
.clone();
let mut stale: Vec<[u8; 32]> = Vec::new();
let mut out = Vec::with_capacity(peers.len());
for (nickname, eid) in peers {
match cache.get(&eid) {
Some(e) => {
let age = (now - e.probed_at).max(0);
if age > REACH_TTL_SECS {
stale.push(eid);
}
out.push(reachability_row(nickname, eid, Some(e), Some(age as u64)));
}
None => {
stale.push(eid);
out.push(reachability_row(nickname, eid, None, None));
}
}
}
for eid in stale {
let mesh = mesh.clone();
tokio::spawn(async move {
probe_peer(&mesh, eid).await;
});
}
out
}
#[cfg(test)]
mod tests {
use super::pong_meta;
use super::{ReachEntry, is_transition, sanitize_relay_url, supersedes};
use crate::roster::presence::APP_METADATA_MAX_BYTES;
fn entry(reachable: bool, rtt_ms: Option<u64>) -> ReachEntry {
ReachEntry {
reachable,
rtt_ms,
probed_at: 1_700_000_000,
meta: String::new(),
services: Vec::new(),
seq: 0,
path: mcpmesh_local_api::PeerPath::Unknown,
}
}
#[test]
fn sanitize_relay_url_drops_credentials_and_path() {
let u = |s: &str| -> iroh::RelayUrl { s.parse().expect("relay url") };
assert_eq!(
sanitize_relay_url(&u("https://user:token@relay.internal/")),
"https://relay.internal",
"userinfo must never reach the wire"
);
assert_eq!(
sanitize_relay_url(&u("https://relay.example:4433/some/path?x=1")),
"https://relay.example:4433",
"port is kept (it names the relay); path/query are not"
);
assert_eq!(
sanitize_relay_url(&u("https://relay.example/")),
"https://relay.example"
);
assert_eq!(
sanitize_relay_url(&u("http://192.168.1.5:4433/")),
"http://192.168.1.5:4433"
);
}
#[test]
fn a_probe_never_overwrites_a_newer_one() {
let mut newer = entry(true, Some(50));
newer.seq = 7;
assert!(
!supersedes(3, &newer),
"an older probe (ticket 3) must not overwrite ticket 7's result"
);
assert!(
supersedes(9, &newer),
"a newer probe (ticket 9) must overwrite"
);
assert!(
supersedes(7, &newer),
"equal tickets cannot happen, but must not wedge the guard"
);
}
#[test]
fn only_a_change_in_the_reachable_verdict_is_a_transition() {
assert!(
is_transition(None, &entry(true, Some(9))),
"first probe finds the peer UP — news"
);
assert!(
!is_transition(None, &entry(false, None)),
"first probe confirms DOWN — the snapshot already said so"
);
assert!(
is_transition(Some(&entry(false, None)), &entry(true, Some(9))),
"came back online"
);
assert!(
is_transition(Some(&entry(true, Some(9))), &entry(false, None)),
"went offline"
);
assert!(
!is_transition(Some(&entry(true, Some(9))), &entry(true, Some(9))),
"unchanged refresh"
);
assert!(
!is_transition(Some(&entry(true, Some(9))), &entry(true, Some(120))),
"rtt drift alone is not a transition"
);
assert!(
!is_transition(Some(&entry(false, None)), &entry(false, None)),
"still offline"
);
let mut with_meta = entry(true, Some(9));
with_meta.meta = "status: away".into();
assert!(
!is_transition(Some(&entry(true, Some(9))), &with_meta),
"meta drift alone is not a transition"
);
}
#[test]
fn pong_services_parses_the_array_and_tolerates_hostile_shapes() {
use super::pong_services;
assert_eq!(
pong_services(&serde_json::json!({"services": ["notes", "kb"]})),
vec!["notes".to_string(), "kb".to_string()]
);
assert!(pong_services(&serde_json::json!({"stack_version": "1"})).is_empty());
assert!(pong_services(&serde_json::json!({"services": 42})).is_empty());
assert!(pong_services(&serde_json::json!({"services": [1, {"x": 2}]})).is_empty());
}
#[tokio::test(flavor = "multi_thread")]
async fn caller_admitted_services_returns_only_admitted() {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
let caller_eid = mcpmesh_net::EndpointId::from_bytes([7u8; 32]).principal();
std::fs::write(
&config_path,
format!(
"[services.shared]\nsocket = \"/run/a.sock\"\nallow = [\"{caller_eid}\"]\n [services.grouped]\nsocket = \"/run/b.sock\"\nallow = [\"team-eng\"]\n [services.private]\nsocket = \"/run/c.sock\"\nallow = [\"eid:other\"]\n"
),
)
.unwrap();
let mesh = crate::daemon::testutil::hermetic_mesh(config_path).await;
let identity = mcpmesh_net::PeerIdentity {
endpoint: mcpmesh_net::EndpointId::from_bytes([7u8; 32]),
name: "bob".into(),
user_id: None,
groups: vec!["team-eng".into()],
};
let admitted = super::caller_admitted_services(&mesh, &identity);
assert_eq!(admitted, vec!["grouped".to_string(), "shared".to_string()]);
assert!(
!admitted.contains(&"private".to_string()),
"never a non-admitted service"
);
}
#[test]
fn pong_meta_extracts_within_cap_and_drops_the_rest() {
assert_eq!(
pong_meta(&serde_json::json!({"stack_version": "1", "meta": "v=1.2.3"})),
"v=1.2.3"
);
assert_eq!(pong_meta(&serde_json::json!({"stack_version": "1"})), "");
assert_eq!(pong_meta(&serde_json::json!({"meta": 42})), "");
assert_eq!(pong_meta(&serde_json::json!({"meta": {"x": 1}})), "");
let at = "x".repeat(APP_METADATA_MAX_BYTES);
assert_eq!(pong_meta(&serde_json::json!({"meta": at.clone()})), at);
let over = "x".repeat(APP_METADATA_MAX_BYTES + 1);
assert_eq!(
pong_meta(&serde_json::json!({"meta": over})),
"",
"oversized meta dropped"
);
}
}