use std::time::Duration;
use mcpmesh_local_api::PeerPath;
pub(crate) static LIVE_WATCHERS: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub(crate) const PATH_CHANGE_SETTLE: Duration = Duration::from_millis(600);
pub(crate) fn decide(observed: &PeerPath, cached: Option<&PeerPath>) -> Option<PeerPath> {
if matches!(observed, PeerPath::Unknown) {
return None;
}
match cached {
Some(known) if known == observed => None,
_ => Some(observed.clone()),
}
}
pub(crate) fn spawn(
mesh: std::sync::Arc<super::MeshState>,
endpoint_id: [u8; 32],
conn: &iroh::endpoint::Connection,
) -> tokio::task::JoinHandle<()> {
let mut events = conn.path_events();
let weak = conn.weak_handle();
LIVE_WATCHERS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tokio::spawn(async move {
let _guard = WatcherGuard;
use n0_future::StreamExt as _;
while let Some(event) = events.next().await {
match event {
iroh::endpoint::PathEvent::Selected { .. } => {}
iroh::endpoint::PathEvent::Lagged { missed, .. } => {
tracing::debug!(missed, "path events lagged; re-reading current path");
}
_ => continue,
}
let seq = mesh
.probe_seq
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let Some(strong) = weak.upgrade() else { break };
let observed =
super::reach::settle(PATH_CHANGE_SETTLE, || super::reach::selected_path(&strong))
.await;
drop(strong);
commit_observation(&mesh, endpoint_id, seq, &observed);
}
})
}
pub(crate) fn commit_observation(
mesh: &std::sync::Arc<super::MeshState>,
endpoint_id: [u8; 32],
seq: u64,
observed: &PeerPath,
) -> Option<mcpmesh_local_api::PeerReachability> {
let committed = {
let mut cache = mesh
.reachability
.lock()
.expect("reachability lock not poisoned");
if let Some(existing) = cache.get(&endpoint_id)
&& !super::reach::supersedes(seq, existing)
{
return None;
}
let cached = cache.get(&endpoint_id).map(|e| e.path.clone());
let path = decide(observed, cached.as_ref())?;
match cache.get_mut(&endpoint_id) {
Some(entry) => {
entry.path = path.clone();
entry.seq = seq;
entry.clone()
}
None => {
let entry = super::ReachEntry {
reachable: true,
rtt_ms: None,
probed_at: crate::util::epoch_now_i64(),
meta: String::new(),
services: Vec::new(),
seq,
path,
};
cache.insert(endpoint_id, entry.clone());
entry
}
}
};
let nickname = mesh.store.resolve(&endpoint_id).ok().flatten()?.nickname;
let row = super::reach::reachability_row(nickname, endpoint_id, Some(&committed), Some(0));
let _ = mesh.reach_bcast.send(row.clone());
Some(row)
}
struct WatcherGuard;
impl Drop for WatcherGuard {
fn drop(&mut self) {
LIVE_WATCHERS.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn only_a_differing_path_is_worth_emitting() {
let relay = PeerPath::Relay { url: None };
assert_eq!(
decide(&PeerPath::Direct, Some(&relay)),
Some(PeerPath::Direct),
"relay -> direct is the recovery a consumer must hear about"
);
assert_eq!(
decide(&relay, Some(&PeerPath::Direct)),
Some(relay.clone()),
"direct -> relay is the DEGRADATION — the privacy indicator just became wrong"
);
assert_eq!(
decide(&PeerPath::Direct, Some(&PeerPath::Direct)),
None,
"an unchanged path must stay quiet, or a stable session emits forever"
);
assert_eq!(
decide(&PeerPath::Direct, None),
Some(PeerPath::Direct),
"first knowledge of a live session's path is news"
);
}
#[test]
fn unknown_is_never_emitted() {
assert_eq!(decide(&PeerPath::Unknown, None), None);
assert_eq!(decide(&PeerPath::Unknown, Some(&PeerPath::Direct)), None);
assert_eq!(
decide(&PeerPath::Unknown, Some(&PeerPath::Relay { url: None })),
None
);
}
#[tokio::test(start_paused = true)]
async fn the_settle_window_waits_for_a_degradation_to_hold() {
let relay = PeerPath::Relay { url: None };
let settled = crate::daemon::reach::settle(PATH_CHANGE_SETTLE, || relay.clone()).await;
assert_eq!(
settled, relay,
"a degradation that holds for the whole window must be reported"
);
let mut polls = 0;
let settled = crate::daemon::reach::settle(Duration::ZERO, || {
polls += 1;
if polls > 1 {
PeerPath::Direct
} else {
relay.clone()
}
})
.await;
assert_eq!(
settled, relay,
"with no window the first observation wins — this assertion fails if \
PATH_CHANGE_SETTLE is ever zeroed"
);
let mut polls = 0;
let settled = crate::daemon::reach::settle(PATH_CHANGE_SETTLE, || {
polls += 1;
if polls > 2 {
PeerPath::Direct
} else {
relay.clone()
}
})
.await;
assert_eq!(
settled,
PeerPath::Direct,
"a recovered blip settles on Direct"
);
assert_eq!(
decide(&settled, Some(&PeerPath::Direct)),
None,
"and a flap that returns to where it started must emit NOTHING"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn an_older_writer_never_overwrites_a_newer_observation() {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("config.toml");
std::fs::write(&cfg, "").unwrap();
let mesh = crate::daemon::testutil::hermetic_mesh(cfg).await;
let eid = [9u8; 32];
mesh.store
.add(crate::allowlist::PeerEntry {
endpoint_id: eid,
nickname: "bob".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let row = commit_observation(&mesh, eid, 7, &PeerPath::Direct);
assert!(row.is_some(), "first observation commits");
let row = commit_observation(&mesh, eid, 3, &PeerPath::Relay { url: None });
assert!(
row.is_none(),
"an older writer must not overwrite a newer observation — this is the #58 defect \
class, and here it would push a stale path a consumer renders as a privacy claim"
);
let cached = mesh
.reachability
.lock()
.unwrap()
.get(&eid)
.map(|e| e.path.clone());
assert_eq!(
cached,
Some(PeerPath::Direct),
"the cache must still hold the NEWER value"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_seeded_entry_is_reachable_with_no_fabricated_rtt() {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("config.toml");
std::fs::write(&cfg, "").unwrap();
let mesh = crate::daemon::testutil::hermetic_mesh(cfg).await;
let eid = [11u8; 32];
mesh.store
.add(crate::allowlist::PeerEntry {
endpoint_id: eid,
nickname: "carol".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
let row = commit_observation(&mesh, eid, 1, &PeerPath::Direct).expect("seeds an entry");
assert!(row.reachable, "a live session IS reachability evidence");
assert_eq!(
row.rtt_ms, None,
"no RTT was measured — never fabricate one"
);
assert_eq!(row.path, PeerPath::Direct);
}
#[tokio::test(flavor = "multi_thread")]
async fn an_unchanged_path_commits_nothing() {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("config.toml");
std::fs::write(&cfg, "").unwrap();
let mesh = crate::daemon::testutil::hermetic_mesh(cfg).await;
let eid = [13u8; 32];
mesh.store
.add(crate::allowlist::PeerEntry {
endpoint_id: eid,
nickname: "dave".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
assert!(commit_observation(&mesh, eid, 1, &PeerPath::Direct).is_some());
assert!(
commit_observation(&mesh, eid, 2, &PeerPath::Direct).is_none(),
"a repeat observation is not news, however fresh its ticket"
);
}
#[test]
fn a_different_relay_url_is_a_change() {
let a = PeerPath::Relay {
url: Some("https://a.example".into()),
};
let b = PeerPath::Relay {
url: Some("https://b.example".into()),
};
assert_eq!(decide(&b, Some(&a)), Some(b.clone()));
assert_eq!(decide(&a, Some(&a)), None);
}
}