#![cfg(unix)]
use std::sync::Arc;
use std::time::Duration;
use iroh::address_lookup::MemoryLookup;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::client::connect_control;
use mcpmesh::config::Config;
use mcpmesh::control::{DaemonState, serve_control};
use mcpmesh::daemon::{
MeshState, PresenceMode, STACK_VERSION, build_services, probe_peer, reachability_of,
spawn_accept_loop,
};
use mcpmesh::pairing::LiveInvites;
use mcpmesh::roster::gate::RosterGate;
use mcpmesh::{Request, StatusResult};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, ALPN_PAIR, ALPN_PING, TrustGate};
use tokio::time::timeout;
static SERIAL: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
async fn target_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(),
ALPN_PING.to_vec(),
])
.bind()
.await
.expect("bind target endpoint")
}
async fn dialer_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 dialer endpoint")
}
fn seed_lookup(dialer: &iroh::Endpoint, target_addr: iroh::EndpointAddr) {
let mem = MemoryLookup::new();
mem.add_endpoint_info(target_addr);
dialer
.address_lookup()
.expect("address lookup services")
.add(mem);
}
fn seed_peer(store: &PeerStore, endpoint_id: [u8; 32], nickname: &str) {
store
.add(PeerEntry {
endpoint_id,
nickname: nickname.into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
}
fn seed_peer_with_addr(
store: &PeerStore,
endpoint_id: [u8; 32],
nickname: &str,
addr: &iroh::EndpointAddr,
) {
store
.add(PeerEntry {
endpoint_id,
nickname: nickname.into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: Some(serde_json::to_string(addr).expect("addr serializes")),
})
.unwrap();
}
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,
)
}
async fn ping_refusal_reason(dialer: &iroh::Endpoint, target: [u8; 32]) -> Option<Vec<u8>> {
let conn = match dialer
.connect(iroh::EndpointId::from_bytes(&target).unwrap(), ALPN_PING)
.await
{
Ok(c) => c,
Err(e) => return Some(format!("{e}").into_bytes()),
};
let Ok((mut send, mut recv)) = conn.open_bi().await else {
return close_reason_bytes(&conn);
};
let _ = mcpmesh_net::framing::write_frame(&mut send, &serde_json::json!({"ping": 1})).await;
let _ = send.finish();
let mut buf = Vec::new();
match tokio::io::AsyncReadExt::read_to_end(&mut recv, &mut buf).await {
Ok(_) if !buf.is_empty() => None, _ => close_reason_bytes(&conn),
}
}
fn close_reason_bytes(conn: &iroh::endpoint::Connection) -> Option<Vec<u8>> {
match conn.close_reason() {
Some(iroh::endpoint::ConnectionError::ApplicationClosed(ac)) => Some(
format!(
"code={} reason={}",
ac.error_code,
String::from_utf8_lossy(&ac.reason)
)
.into_bytes(),
),
other => Some(format!("{other:?}").into_bytes()),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_hidden_node_never_leaks_presence_through_the_throttle_close() {
timeout(Duration::from_secs(120), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("ta.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
seed_peer(&a_store, b_id, "flooder");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let a_limits =
mcpmesh::limits::MeshLimiters::from_config(&mcpmesh::config::LimitsCfg::default());
a_mesh.set_limits(a_limits.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let mut throttled = None;
for _ in 0..120 {
if let Some(reason) = ping_refusal_reason(&b_ep, a_id).await
&& String::from_utf8_lossy(&reason).contains("ping rate limited")
{
throttled = Some(reason);
break;
}
}
assert!(
throttled.is_some(),
"the bucket must actually be exhausted, or the ordering below proves nothing"
);
a_mesh.set_presence_mode(PresenceMode::Off);
let hidden = ping_refusal_reason(&b_ep, a_id)
.await
.expect("an off node must refuse");
assert_eq!(
String::from_utf8_lossy(&hidden),
"code=401 reason=unauthorized",
"a hidden node with an EXHAUSTED bucket must still answer like the trust gate. Got \
{:?} — if this is the throttle close, the policy is being consulted after the limiter \
and 'appear offline' announces itself to anyone who floods first.",
String::from_utf8_lossy(&hidden)
);
let refused_before = a_limits.pings_refused();
for _ in 0..20 {
let _ = ping_refusal_reason(&b_ep, a_id).await;
}
assert!(
a_limits.pings_refused() > refused_before,
"a hidden node must still consult the limiter for refused probes — the counter did \
not move across 20 of them ({refused_before} → {}), so a revoked peer can flood a \
hidden node for free",
a_limits.pings_refused()
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("throttle-ordering test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn presence_mode_controls_who_gets_a_pong_and_never_says_why() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("pa.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("pb.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
seed_peer(&a_store, b_id, "granted-peer");
let d_ep = dialer_endpoint().await;
let d_id = *d_ep.id().as_bytes();
seed_lookup(&d_ep, a_addr.clone());
let d_store = Arc::new(PeerStore::open(&dir.path().join("pd.redb")).unwrap());
seed_peer(&d_store, a_id, "alice");
seed_peer(&a_store, d_id, "revoked-peer");
let c_ep = dialer_endpoint().await;
seed_lookup(&c_ep, a_addr.clone());
let b_principal = mcpmesh_net::EndpointId::from_bytes(b_id).principal();
let toml = format!("[services.notes]\nrun = [\"true\"]\nallow = [\"{b_principal}\"]\n");
let (b_dial, d_dial) = (b_ep.clone(), d_ep.clone());
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let d_mesh = assemble_mesh(d_ep, d_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str(&toml).unwrap())),
);
let stranger_refusal = ping_refusal_reason(&c_ep, a_id)
.await
.expect("an unpaired stranger must never get a pong");
assert_eq!(
String::from_utf8_lossy(&stranger_refusal),
"code=401 reason=unauthorized",
"the reference refusal must be the trust gate's application close, not a dial failure"
);
assert_eq!(a_mesh.presence_mode(), PresenceMode::Paired, "the default");
assert!(
probe_peer(&d_mesh, a_id).await.reachable,
"presence_mode = paired must pong a paired peer holding no grant (today's behaviour)"
);
a_mesh.set_presence_mode(PresenceMode::Granted);
assert!(
probe_peer(&b_mesh, a_id).await.reachable,
"presence_mode = granted must still pong a caller holding a service grant"
);
let revoked_refusal = ping_refusal_reason(&d_dial, a_id)
.await
.expect("granted must refuse a paired caller holding NO grant — this is the whole ask");
a_mesh.set_presence_mode(PresenceMode::Off);
let off_refusal = ping_refusal_reason(&b_dial, a_id)
.await
.expect("presence_mode = off must refuse even a granted caller");
assert_eq!(
revoked_refusal, stranger_refusal,
"a granted-mode refusal must be byte-identical to the trust gate's, or the prober \
learns the peer is online and merely ungranted"
);
assert_eq!(
off_refusal, stranger_refusal,
"an off-mode refusal must be byte-identical to the trust gate's, or 'appear offline' \
announces itself"
);
assert!(
!String::from_utf8_lossy(&off_refusal).contains("ping rate limited"),
"a hidden node must not leak presence through the throttle close"
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("presence_mode test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn ping_probe_reports_paired_peer_reachable_stranger_and_down_peer_not() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_ep_handle = a_ep.clone(); let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice"); seed_peer(&a_store, b_id, "beacon-b");
let c_ep = dialer_endpoint().await;
seed_lookup(&c_ep, a_addr.clone());
let c_store = Arc::new(PeerStore::open(&dir.path().join("c.redb")).unwrap());
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let c_mesh = assemble_mesh(c_ep, c_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "a paired peer's probe must be reachable");
assert!(
entry.rtt_ms.is_some(),
"a reachable probe records a round-trip time"
);
let list = reachability_of(&b_mesh);
let alice = list
.iter()
.find(|r| r.name == "alice")
.expect("reachability_of surfaces the paired peer by nickname");
assert!(alice.reachable, "the cached probe result is surfaced");
assert!(alice.rtt_ms.is_some(), "the cached RTT is surfaced");
let stranger = probe_peer(&c_mesh, a_id).await;
assert!(
!stranger.reachable,
"an unpaired peer gets no pong (trust gate closed the connection)"
);
assert!(stranger.rtt_ms.is_none());
accept.abort();
a_ep_handle.close().await;
let down = probe_peer(&b_mesh, a_id).await;
assert!(
!down.reachable,
"a probe of a down peer must be unreachable"
);
std::mem::forget(dir);
})
.await
.expect("reachability test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn status_includes_reachability() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
seed_peer(&a_store, b_id, "beacon-b");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let entry = probe_peer(&b_mesh, a_id).await;
assert!(
entry.reachable,
"precondition: the paired peer must probe reachable"
);
let socket = dir.path().join("control.sock");
let listener = mcpmesh::ipc::bind_control_socket(&socket).await.unwrap();
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, b_mesh.clone()));
let control = tokio::spawn(serve_control(listener, state));
let mut client = connect_control(&socket)
.await
.expect("raw connect_control to B");
let value = client
.request(Request::Status)
.await
.expect("status over mcpmesh-local/1");
let status: StatusResult =
serde_json::from_value(value).expect("StatusResult deserializes");
assert!(
status.reachability.iter().any(|r| r.name == "alice"),
"status.reachability must surface the paired peer by nickname: {:?}",
status.reachability
);
control.abort();
accept.abort();
std::mem::forget(dir);
})
.await
.expect("status reachability test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn cold_probe_uses_the_pairing_proven_address_hint_without_discovery() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer_with_addr(&b_store, a_id, "alice", &a_addr);
seed_peer(&a_store, b_id, "beacon-b");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
tokio::time::sleep(Duration::from_millis(200)).await;
let entry = probe_peer(&b_mesh, a_id).await;
assert!(
entry.reachable,
"a cold probe must reach the peer using the pairing-proven `last_addr` hint rather \
than depending on discovery to resolve the bare endpoint-id"
);
assert!(
entry.rtt_ms.is_some(),
"a reachable probe reports a measured RTT"
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("cold-probe address-hint test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn probe_carries_peer_app_metadata_into_reachability() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_eid = format!("eid:{}", a_ep.id());
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
seed_peer(&a_store, b_id, "beacon-b");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let a_socket = dir.path().join("a-control.sock");
let a_listener = mcpmesh::ipc::bind_control_socket(&a_socket).await.unwrap();
let a_state = Arc::new(DaemonState::with_mesh(STACK_VERSION, a_mesh.clone()));
let a_control = tokio::spawn(serve_control(a_listener, a_state));
connect_control(&a_socket)
.await
.expect("connect A control")
.set_app_metadata("v=4.2.0")
.await
.expect("A sets its app metadata");
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "the paired peer is reachable");
assert_eq!(
entry.meta, "v=4.2.0",
"the probe carried the peer's app metadata off the pong"
);
let list = reachability_of(&b_mesh);
let alice = list.iter().find(|r| r.name == "alice").expect("A surfaced");
assert_eq!(
alice.meta, "v=4.2.0",
"reachability surfaces the peer's app metadata"
);
assert_eq!(
alice.principal.as_deref(),
Some(a_eid.as_str()),
"reachability row carries the peer's eid principal"
);
a_control.abort();
accept.abort();
std::mem::forget(dir);
})
.await
.expect("probe-metadata test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn probe_surfaces_the_services_the_peer_grants_the_caller() {
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_eid = format!("eid:{}", b_ep.id());
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
let config = dir.path().join("a-config.toml");
std::fs::write(
&config,
format!(
"[services.shared]\nsocket = \"/run/s.sock\"\nallow = [\"{b_eid}\"]\n\
[services.private]\nsocket = \"/run/p.sock\"\nallow = [\"eid:someoneelse\"]\n"
),
)
.unwrap();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
seed_peer(&a_store, b_id, "beacon-b");
let a_cfg = Config::load(&config).expect("A's config parses");
let a_mesh = assemble_mesh(a_ep, a_store, config);
let b_mesh = assemble_mesh(b_ep, b_store, dir.path().join("b-config.toml"));
let accept = spawn_accept_loop(a_mesh.clone(), Arc::new(build_services(&a_cfg)));
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable);
assert_eq!(
entry.services,
vec!["shared".to_string()],
"probe surfaces exactly the caller-admitted services (#52)"
);
assert!(
!entry.services.contains(&"private".to_string()),
"never a service the peer does not grant the caller"
);
accept.abort();
std::mem::forget(dir);
})
.await
.expect("peer-services probe test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_throttled_probe_is_refused_but_never_reports_the_peer_offline() {
let _serial = SERIAL.lock().await;
timeout(Duration::from_secs(120), async {
let dir = tempfile::tempdir().unwrap();
let config = dir.path().join("config.toml");
std::fs::write(&config, "").unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_store = Arc::new(PeerStore::open(&dir.path().join("fa.redb")).unwrap());
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("fb.redb")).unwrap());
seed_peer(&a_store, b_id, "b");
seed_peer(&b_store, a_id, "a");
let a_mesh = assemble_mesh(a_ep, a_store, config.clone());
let a_limits =
mcpmesh::limits::MeshLimiters::from_config(&mcpmesh::config::LimitsCfg::default());
a_mesh.set_limits(a_limits.clone());
let b_mesh = assemble_mesh(b_ep, b_store, config.clone());
let _accept = spawn_accept_loop(
a_mesh.clone(),
Arc::new(build_services(&Config::from_toml_str("").unwrap())),
);
let mut reachable = 0usize;
let mut fresh = 0usize;
let mut stale = 0usize;
let mut last_seq: Option<u64> = None;
for _ in 0..90 {
let entry = probe_peer(&b_mesh, a_id).await;
if entry.reachable {
reachable += 1;
}
if last_seq == Some(entry.seq) {
stale += 1;
} else {
fresh += 1;
}
last_seq = Some(entry.seq);
}
assert_eq!(
reachable, 90,
"a rate-limit refusal is not evidence the peer is down: every probe of a live, paired \
peer must report reachable — the refused ones from the still-fresh cache. A false \
count here means a refusal wrote `reachable: false` (PR #142 gate, HIGH). \
reachable={reachable} fresh={fresh} stale={stale}"
);
assert!(
stale > 0,
"a paired peer flooding past the cap must eventually be REFUSED (#89), observable as \
the cached entry returned unchanged. Zero stale answers means the accept arm never \
consulted the limiter — the unmetered pong-flood this issue reports. \
fresh={fresh} stale={stale}"
);
assert!(
a_limits.pings_refused() > 0,
"the responder must COUNT its refusals (#89 defect 3: unmetered AND unrecorded) — \
the count is the refusal's only footprint besides the debug log"
);
})
.await
.expect("ping flood test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn peer_services_answers_from_the_fresh_cache_without_probing() {
let _serial = SERIAL.lock().await;
timeout(Duration::from_secs(60), async {
let dir = tempfile::tempdir().unwrap();
let a_ep = target_endpoint().await;
let a_id = *a_ep.id().as_bytes();
let a_addr = a_ep.addr();
let a_ep_handle = a_ep.clone();
let b_ep = dialer_endpoint().await;
let b_id = *b_ep.id().as_bytes();
let b_eid = format!("eid:{}", b_ep.id());
seed_lookup(&b_ep, a_addr.clone());
let b_store = Arc::new(PeerStore::open(&dir.path().join("b.redb")).unwrap());
seed_peer(&b_store, a_id, "alice");
let config = dir.path().join("a-config.toml");
std::fs::write(
&config,
format!("[services.shared]\nsocket = \"/run/s.sock\"\nallow = [\"{b_eid}\"]\n"),
)
.unwrap();
let a_store = Arc::new(PeerStore::open(&dir.path().join("a.redb")).unwrap());
seed_peer(&a_store, b_id, "beacon-b");
let a_cfg = Config::load(&config).expect("A's config parses");
let a_mesh = assemble_mesh(a_ep, a_store, config);
let b_mesh = assemble_mesh(b_ep, b_store, dir.path().join("b-config.toml"));
let accept = spawn_accept_loop(a_mesh.clone(), Arc::new(build_services(&a_cfg)));
let entry = probe_peer(&b_mesh, a_id).await;
assert!(entry.reachable, "precondition: A must probe reachable");
assert_eq!(
entry.services,
vec!["shared".to_string()],
"precondition: the cached entry carries the granted service"
);
let socket = dir.path().join("control.sock");
let listener = mcpmesh::ipc::bind_control_socket(&socket).await.unwrap();
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, b_mesh.clone()));
let control = tokio::spawn(serve_control(listener, state));
let mut client = connect_control(&socket).await.expect("connect B control");
accept.abort();
a_ep_handle.close().await;
let services = client.peer_services("alice").await.expect(
"peer_services must answer from the fresh cache rather than probing — an \
unconditional probe is the #142-gate shape that reported healthy peers offline",
);
assert_eq!(services, vec!["shared".to_string()]);
control.abort();
std::mem::forget(dir);
})
.await
.expect("peer_services cache-freshness test timed out");
}