use std::path::Path;
use std::sync::Arc;
use anyhow::{Context, Result};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{ALPN_MCP, ALPN_PAIR, ALPN_PING, TrustGate};
use mcpmesh_trust::DeviceKey;
use crate::allowlist::{AllowlistGate, PeerStore};
use crate::audit::{AuditLog, AuditSink};
use crate::config::Config;
use crate::control::{DaemonState, serve_control};
use crate::ipc;
use crate::node::StartError;
use std::time::Duration;
use crate::pairing::LiveInvites;
use crate::paths::NodePaths;
use crate::roster::RosterStore;
use crate::roster::freshness::FreshnessStore;
use crate::roster::gate::{ComposedGate, RosterGate};
use crate::util::{blocking, epoch_now_i64};
use super::accept::spawn_accept_loop;
use super::roster_install::{
respawn_poll_loop, roster_confirmed_path, spawn_staleness_sweep, warn_if_degraded_grace,
};
use super::{MeshState, STACK_VERSION, build_services_audited, default_self_nickname};
pub async fn serve_forever(socket: &Path, paths: NodePaths) -> Result<()> {
let listener = match ipc::bind_control_socket(socket).await {
Ok(l) => l,
Err(e)
if e.downcast_ref::<std::io::Error>()
.is_some_and(|io| io.kind() == std::io::ErrorKind::AddrInUse) =>
{
tracing::info!("another mcpmesh daemon already owns the control endpoint; exiting");
return Ok(());
}
Err(e) => return Err(e),
};
let booted = start_node(paths, None).await?;
let state = booted.state;
if let Some(mesh) = state.mesh() {
let conflict = std::sync::Arc::new(crate::diag::IdentityConflict::default());
mesh.adopt_identity_conflict(conflict.clone());
crate::diag::install_for_daemon(conflict);
}
drop(booted.background);
tracing::info!(
endpoint_id = %state.mesh_required()?.endpoint.id(),
socket = %socket.display(),
"mcpmesh daemon serving mesh + control"
);
serve_control(listener, state).await
}
pub(crate) struct BootedNode {
pub(crate) state: Arc<DaemonState>,
pub(crate) background: Vec<tokio::task::JoinHandle<()>>,
}
pub(crate) async fn shutdown_booted(booted: BootedNode) {
let state = &booted.state;
state.request_shutdown();
state.abort_control_tasks();
let Some(mesh) = state.mesh().cloned() else {
return;
};
if let Some(task) = mesh.accept_task.lock().await.take() {
task.abort();
}
if let Some(task) = mesh.poll_loop.lock().await.take() {
task.abort();
}
for task in booted.background {
task.abort();
}
if let Some(blobs) = mesh.app_blobs.lock().await.take() {
blobs.shutdown().await;
}
mesh.endpoint.close().await;
}
pub(crate) async fn start_node(
paths: NodePaths,
config: Option<Config>,
) -> Result<BootedNode, StartError> {
let config_path = paths.config_path.clone();
let db_path = paths.state_db_path.clone();
boot_node(paths, config)
.await
.map_err(|e| StartError::classify(e, &config_path, &db_path))
}
async fn boot_node(paths: NodePaths, config: Option<Config>) -> Result<BootedNode> {
let _ = rustls::crypto::ring::default_provider().install_default();
let mut background: Vec<tokio::task::JoinHandle<()>> = Vec::new();
let audit = AuditSink::new(AuditLog::spawn(paths.audit_dir.clone()));
let config_path = paths.config_path.clone();
let cfg = match config {
Some(c) => c,
None => Config::load(&config_path)
.with_context(|| format!("config error in {}", config_path.display()))?,
};
if let Some(cutoff) =
crate::audit::retention_cutoff(&crate::audit::now_ts()[..7], cfg.limits.audit_retain_months)
{
let dir = paths.audit_dir.clone();
match blocking("join audit retention prune", move || {
crate::audit::prune_before(&dir, &cutoff)
})
.await
{
Ok(Ok(deleted)) if !deleted.is_empty() => {
tracing::info!(?deleted, "audit retention pruned old months");
}
Ok(Ok(_)) => {}
Ok(Err(e)) => tracing::warn!(%e, "audit retention prune failed"),
Err(e) => tracing::warn!(%e, "audit retention prune failed"),
}
};
let key_path = match cfg.identity.device_key.clone() {
Some(p) => p,
None => paths.device_key_path.clone(),
};
let (key, _created) = DeviceKey::load_or_generate(&key_path)
.map_err(|e| anyhow::anyhow!("device key error at {}: {e}", key_path.display()))?;
let secret = iroh::SecretKey::from_bytes(&key.secret_bytes());
let roster_mode = cfg.identity.org_root_pk.is_some();
let presence = presence_mode(&cfg.network)?;
validate_service_rates(&cfg)?;
let endpoint = build_endpoint(secret, &cfg.network, roster_mode).await?;
let our_id = endpoint.id();
let db_path = paths.state_db_path.clone();
if let Some(parent) = db_path.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("create data dir {}", parent.display()))?;
}
let store = blocking("join peer-store open", move || PeerStore::open(&db_path)).await??;
let store = Arc::new(store);
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let grace = cfg.roster.grace_seconds();
let roster = Arc::new(RosterGate::with_freshness(
grace,
cfg.roster.max_staleness_seconds(),
));
if let Some(pk_str) = cfg.identity.org_root_pk.clone() {
match crate::roster::parse_org_root_pk(&pk_str) {
Ok(pk) => {
let rstore = RosterStore::new(paths.roster_path.clone());
match blocking("join roster load", move || rstore.load(&pk)).await {
Ok(Ok(Some(view))) => {
roster.install(view);
tracing::info!("installed roster loaded");
warn_if_degraded_grace(&roster);
}
Ok(Ok(None)) => {}
Ok(Err(e)) | Err(e) => {
tracing::warn!(%e, "installed roster failed to load; refusing roster peers")
}
}
}
Err(e) => tracing::warn!(%e, "pinned org_root_pk is invalid; roster mode disabled"),
}
}
{
let fpath = roster_confirmed_path(&config_path);
let fstore = FreshnessStore::new(fpath.clone());
match blocking("join roster freshness load", move || fstore.load()).await {
Ok(Ok(Some(lc))) => roster.set_last_confirmed(lc),
Ok(Ok(None)) if roster.view().is_some() => {
let now = epoch_now_i64();
roster.set_last_confirmed(now);
let fstore = FreshnessStore::new(fpath);
match blocking("join roster freshness upgrade-grace persist", move || {
fstore.store(now)
})
.await
{
Ok(Ok(())) => tracing::info!("roster freshness upgrade grace applied"),
Ok(Err(e)) | Err(e) => {
tracing::warn!(%e, "persist roster freshness upgrade grace")
}
}
}
Ok(Ok(None)) => {}
Ok(Err(e)) | Err(e) => {
tracing::warn!(%e, "read roster freshness sidecar; treating as unconfirmed")
}
}
}
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let limiters = crate::limits::MeshLimiters::from_config(&cfg.limits);
let services = Arc::new(build_services_audited(&cfg, &audit, &limiters));
let invites = Arc::new(LiveInvites::load(
paths.invites_path.clone(),
crate::util::epoch_now_u64(),
));
let self_nickname = cfg
.identity
.nickname
.clone()
.unwrap_or_else(|| default_self_nickname(&our_id));
let conn_registry = Arc::new(ConnRegistry::new());
let (gossip, blobs, roster_topic, presence_topic) =
compose_roster_transport(&endpoint, &roster, &cfg, roster_mode, &our_id).await;
let mesh = MeshState::new(
endpoint,
gate,
store,
invites,
self_nickname,
config_path,
roster,
conn_registry,
gossip,
blobs,
roster_topic,
presence_topic,
);
mesh.set_audit(audit.clone());
mesh.set_limits(limiters.clone());
mesh.set_applied_relays(&cfg.network.relay_mode, &cfg.network.relay_urls);
mesh.set_presence_mode(presence);
let user_key_path = match cfg.identity.user_key.clone() {
Some(p) => p,
None => paths.user_key_path.clone(),
};
let self_binding = match mcpmesh_trust::UserKey::load_or_generate(&user_key_path) {
Ok((user_key, _created)) => {
let (user_pk, sig) = mcpmesh_trust::binding::present(&user_key, our_id.as_bytes());
Some(crate::pairing::rendezvous::SelfBinding { user_pk, sig })
}
Err(e) => {
tracing::warn!(
%e,
path = %user_key_path.display(),
"no user key for pairing identity; paired peers will store this daemon without a user_id"
);
None
}
};
mesh.set_self_binding(self_binding);
{
let scopes_path = paths.blob_scopes_path.clone();
match blocking("join app-blob scopes load", move || {
crate::blobs::scope::ScopeStore::open(scopes_path)
})
.await
{
Ok(Ok(scopes)) => {
match crate::blobs::provider::AppBlobs::load(
paths.blobs_dir.clone(),
Arc::new(scopes),
mesh.gate.clone(),
mesh.endpoint.clone(),
audit.clone(),
limiters.clone(),
Some(mesh.blob_bcast.clone()),
)
.await
{
Ok(provider) => {
provider.enable_relay_wait();
mesh.set_blobs_dir(paths.blobs_dir.clone());
mesh.set_app_blobs(provider).await
}
Err(e) => tracing::warn!(%e, "app-blob provider disabled (build failed)"),
}
}
Ok(Err(e)) | Err(e) => {
tracing::warn!(%e, "app-blob scopes failed to load; provider disabled")
}
}
}
background.push(crate::daemon::spawn_self_net_watch(mesh.clone()));
let accept_task = spawn_accept_loop(mesh.clone(), services);
mesh.set_accept_task(accept_task).await;
if roster_mode {
let book = std::sync::Arc::new(crate::roster::transport::RosterAddrBook::register(
&mesh.endpoint,
256,
));
let _ = mesh.roster_addr_book.set(book);
background.push(crate::roster::distribute::spawn_receive_loop(mesh.clone()));
}
if roster_mode {
background.push(crate::roster::presence::track_loop(mesh.presence_ctx()));
let self_user_id = cfg.identity.user_id.clone().or_else(|| {
mesh.roster
.view()
.and_then(|v| v.resolve(our_id.as_bytes()).map(|d| d.user_id.clone()))
});
match self_user_id {
Some(user_id) => {
let device_key = ed25519_dalek::SigningKey::from_bytes(&key.secret_bytes());
background.push(crate::roster::presence::publish_loop(
mesh.presence_ctx(),
device_key,
user_id,
));
}
None => tracing::debug!(
"presence publish skipped: no user_id for this node (track loop still runs)"
),
}
}
if roster_mode {
background.push(spawn_staleness_sweep(mesh.clone()));
}
if roster_mode && let Some(url) = cfg.roster.url.clone() {
respawn_poll_loop(&mesh, url).await;
}
let state = Arc::new(DaemonState::with_mesh(STACK_VERSION, mesh));
Ok(BootedNode { state, background })
}
#[derive(Debug)]
pub enum DiscoveryPlan {
N0,
Custom(Vec<url::Url>),
}
#[derive(Debug)]
pub enum NetPlan {
Hermetic,
Mesh {
relay: iroh::RelayMode,
discovery: DiscoveryPlan,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum PresenceMode {
#[default]
Paired,
Granted,
Off,
}
impl PresenceMode {
pub fn as_str(self) -> &'static str {
match self {
Self::Paired => "paired",
Self::Granted => "granted",
Self::Off => "off",
}
}
}
pub fn presence_mode(net: &crate::config::NetworkCfg) -> Result<PresenceMode> {
match net.presence_mode.as_str() {
"paired" => Ok(PresenceMode::Paired),
"granted" => Ok(PresenceMode::Granted),
"off" => Ok(PresenceMode::Off),
other => anyhow::bail!(
"[network] unknown presence_mode {other:?} (expected \"paired\" | \"granted\" | \"off\")"
),
}
}
pub fn net_plan(net: &crate::config::NetworkCfg) -> Result<NetPlan> {
let relay = match net.relay_mode.as_str() {
"disabled" => return Ok(NetPlan::Hermetic),
"default" => iroh::RelayMode::Default,
"custom" => {
anyhow::ensure!(
!net.relay_urls.is_empty(),
"[network] relay_mode = \"custom\" requires at least one relay_urls entry"
);
let urls = net
.relay_urls
.iter()
.map(|u| {
u.parse::<iroh::RelayUrl>()
.map_err(|e| anyhow::anyhow!("[network] relay_urls entry {u:?}: {e}"))
})
.collect::<Result<Vec<_>>>()?;
iroh::RelayMode::custom(urls)
}
other => anyhow::bail!(
"[network] unknown relay_mode {other:?} (expected \"default\" | \"custom\" | \"disabled\")"
),
};
let discovery = match net.discovery_mode.as_str() {
"default" => DiscoveryPlan::N0,
"custom" => {
anyhow::ensure!(
!net.discovery_urls.is_empty(),
"[network] discovery_mode = \"custom\" requires at least one discovery_urls entry \
(a self-hosted pkarr relay, e.g. an iroh-dns-server)"
);
let urls = net
.discovery_urls
.iter()
.map(|u| {
u.parse::<url::Url>()
.map_err(|e| anyhow::anyhow!("[network] discovery_urls entry {u:?}: {e}"))
})
.collect::<Result<Vec<_>>>()?;
DiscoveryPlan::Custom(urls)
}
other => anyhow::bail!(
"[network] unknown discovery_mode {other:?} (expected \"default\" | \"custom\")"
),
};
Ok(NetPlan::Mesh { relay, discovery })
}
pub(crate) async fn build_endpoint(
secret: iroh::SecretKey,
net: &crate::config::NetworkCfg,
roster_mode: bool,
) -> Result<iroh::Endpoint> {
let builder = match net_plan(net)? {
NetPlan::Hermetic => iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled),
NetPlan::Mesh {
relay,
discovery: DiscoveryPlan::N0,
} => iroh::Endpoint::builder(iroh::endpoint::presets::N0).relay_mode(relay),
NetPlan::Mesh {
relay,
discovery: DiscoveryPlan::Custom(urls),
} => {
let mut b = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal).relay_mode(relay);
for u in urls {
b = b
.address_lookup(iroh::address_lookup::PkarrPublisher::builder(u.clone()))
.address_lookup(iroh::address_lookup::PkarrResolver::builder(u));
}
b
}
};
let alpns = alpns_for(roster_mode);
let builder = apply_relay_only(builder, net);
let builder = apply_transport_config(builder, net)?;
builder
.secret_key(secret)
.alpns(alpns)
.bind()
.await
.context("bind iroh endpoint")
}
pub(crate) const IROH_DEFAULT_IDLE_SECS: u64 = 30;
pub(crate) const IROH_MAX_PATH_KEEP_ALIVE_SECS: u64 = 5;
fn build_transport_config(
net: &crate::config::NetworkCfg,
) -> Result<Option<iroh::endpoint::QuicTransportConfig>> {
let (idle, keep) = (net.idle_timeout_secs, net.keep_alive_secs);
if idle.is_none() && keep.is_none() {
return Ok(None);
}
if let Some(k) = keep {
anyhow::ensure!(
k > 0,
"[network] keep_alive_secs = 0 is not \"disable keepalives\" — it is a zero-length \
timer that makes every packet emit a PING, saturating the link. There is no way to \
disable the transport keepalive; omit the key to get iroh's default of \
{IROH_MAX_PATH_KEEP_ALIVE_SECS}s"
);
anyhow::ensure!(
k <= IROH_MAX_PATH_KEEP_ALIVE_SECS,
"[network] keep_alive_secs ({k}) is above iroh's per-path keepalive cap of \
{IROH_MAX_PATH_KEEP_ALIVE_SECS}s, which it silently ignores — every path would keep \
pinging every {IROH_MAX_PATH_KEEP_ALIVE_SECS}s regardless, so this cannot reduce \
keepalive traffic. Lower it to make pings MORE frequent on a lossy link; there is no \
supported way to make them less frequent"
);
let effective = idle.unwrap_or(IROH_DEFAULT_IDLE_SECS);
anyhow::ensure!(
effective == 0 || k < effective,
"[network] keep_alive_secs ({k}) must be less than the idle timeout ({effective}s{}): \
a keepalive arriving after the peer's idle timer has fired severs sessions on a clock \
rather than keeping them open",
if idle.is_some() {
""
} else {
", iroh's default — set idle_timeout_secs to raise it"
}
);
}
let mut cfg = iroh::endpoint::QuicTransportConfig::builder();
if let Some(i) = idle {
let timeout = if i == 0 {
None
} else {
Some(
iroh::endpoint::IdleTimeout::try_from(Duration::from_secs(i)).map_err(|e| {
anyhow::anyhow!(
"[network] idle_timeout_secs {i} is out of the range QUIC can encode: {e}"
)
})?,
)
};
cfg = cfg.max_idle_timeout(timeout);
}
if let Some(k) = keep {
let d = Duration::from_secs(k);
cfg = cfg
.keep_alive_interval(d)
.default_path_keep_alive_interval(d);
}
Ok(Some(cfg.build()))
}
pub fn validate_service_rates(cfg: &crate::config::Config) -> Result<()> {
for (name, svc) in &cfg.services {
anyhow::ensure!(
svc.rate_limit_per_min != Some(0),
"[services.{name}] rate_limit_per_min must be at least 1 \
(omit it to inherit [limits].rate_limit_per_min; 0 does NOT mean unlimited here)"
);
}
Ok(())
}
pub fn validate_transport_config(net: &crate::config::NetworkCfg) -> Result<()> {
presence_mode(net)?;
build_transport_config(net).map(|_| ())
}
fn apply_transport_config(
builder: iroh::endpoint::Builder,
net: &crate::config::NetworkCfg,
) -> Result<iroh::endpoint::Builder> {
Ok(match build_transport_config(net)? {
None => builder,
Some(cfg) => builder.transport_config(cfg),
})
}
#[cfg(feature = "unstable-relay-only")]
#[derive(Debug)]
pub(crate) struct RelayOnlySelector;
#[cfg(feature = "unstable-relay-only")]
impl iroh::endpoint::transports::PathSelector for RelayOnlySelector {
fn select(
&self,
ctx: &iroh::endpoint::transports::PathSelectionContext<'_>,
) -> iroh::endpoint::transports::PathSelection {
let mut selection = iroh::endpoint::transports::PathSelection::none();
if let Some(p) = ctx.paths().find(|p| {
matches!(
p.network_path().remote(),
iroh::endpoint::transports::Addr::Relay(..)
)
}) {
selection.set(&p);
}
selection
}
}
#[cfg(feature = "unstable-relay-only")]
fn apply_relay_only(
builder: iroh::endpoint::Builder,
net: &crate::config::NetworkCfg,
) -> iroh::endpoint::Builder {
if net.relay_only {
tracing::warn!(
"[network] relay_only = true — TESTING POSTURE: no direct addresses are published or \
resolved, so peers can only reach this node through the relay. Not for production."
);
return builder
.addr_filter(iroh::address_lookup::AddrFilter::relay_only())
.path_selector(std::sync::Arc::new(RelayOnlySelector));
}
builder
}
#[cfg(not(feature = "unstable-relay-only"))]
fn apply_relay_only(
builder: iroh::endpoint::Builder,
net: &crate::config::NetworkCfg,
) -> iroh::endpoint::Builder {
if net.relay_only {
tracing::warn!(
"[network] relay_only = true is IGNORED — this binary was built without the \
`unstable-relay-only` cargo feature. Traffic will take whatever path iroh selects, \
which on a LAN with IPv6 is usually DIRECT. Rebuild with the feature, or do not rely \
on this test having exercised the relay."
);
}
builder
}
pub(crate) fn alpns_for(roster_mode: bool) -> Vec<Vec<u8>> {
let mut alpns = vec![
ALPN_MCP.to_vec(),
ALPN_PAIR.to_vec(),
ALPN_PING.to_vec(),
crate::blobs::APP_BLOB_ALPN.to_vec(),
];
if roster_mode {
alpns.push(crate::roster::transport::GOSSIP_ALPN.to_vec());
alpns.push(crate::roster::transport::BLOB_ALPN.to_vec());
}
alpns
}
async fn compose_roster_transport(
endpoint: &iroh::Endpoint,
roster: &Arc<RosterGate>,
cfg: &Config,
roster_mode: bool,
our_id: &iroh::EndpointId,
) -> (
Option<iroh_gossip::net::Gossip>,
Option<crate::roster::transport::RosterBlobs>,
Option<crate::roster::transport::RosterGossip>,
Option<crate::roster::transport::RosterGossip>,
) {
use crate::roster::transport;
if !roster_mode {
return (None, None, None, None);
}
let Some(org_id) = cfg
.identity
.org_id
.clone()
.or_else(|| roster.view().map(|v| v.org_id().to_string()))
else {
tracing::warn!("roster mode but no org_id known; gossip distribution disabled");
return (None, None, None, None);
};
let gossip = transport::spawn_gossip(endpoint);
let blobs = transport::RosterBlobs::new(endpoint);
let bootstrap: Vec<iroh::EndpointId> = roster
.view()
.map(|v| {
v.device_endpoints()
.filter(|d| *d != our_id.as_bytes())
.filter_map(|d| iroh::EndpointId::from_bytes(d).ok())
.collect()
})
.unwrap_or_default();
let roster_topic = match transport::subscribe(
&gossip,
transport::roster_topic_bytes(&org_id),
bootstrap.clone(),
)
.await
{
Ok(rg) => Some(rg),
Err(e) => {
tracing::warn!(%e, "roster-topic subscribe failed; distribution disabled");
None
}
};
let presence_topic =
match transport::subscribe(&gossip, transport::presence_topic_bytes(&org_id), bootstrap)
.await
{
Ok(rg) => Some(rg),
Err(e) => {
tracing::warn!(%e, "presence-topic subscribe failed; presence disabled");
None
}
};
(Some(gossip), Some(blobs), roster_topic, presence_topic)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn relay_only_parses_regardless_of_the_feature() {
let cfg: crate::config::NetworkCfg =
toml::from_str("relay_only = true\n").expect("the field parses on ANY build");
assert!(cfg.relay_only);
let plain: crate::config::NetworkCfg = toml::from_str("").unwrap();
assert!(!plain.relay_only, "default must be off");
}
#[cfg(feature = "unstable-relay-only")]
#[tokio::test(flavor = "multi_thread")]
async fn relay_only_keeps_data_on_the_relay_while_a_direct_path_exists() {
use std::time::Duration;
tokio::time::timeout(Duration::from_secs(90), async {
let (relay_map, _relay_url, _guard) = iroh::test_utils::run_relay_server()
.await
.expect("in-process relay");
let server = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Custom(relay_map.clone()))
.ca_tls_config(iroh_relay::tls::CaTlsConfig::insecure_skip_verify())
.alpns(vec![b"mcpmesh/relayonly/test".to_vec()])
.bind()
.await
.expect("bind server");
let server_addr = server.addr();
tokio::spawn(async move {
while let Some(incoming) = server.accept().await {
if let Ok(conn) = incoming.await
&& let Ok((mut send, _recv)) = conn.accept_bi().await
{
let _ = send.write_all(b"ok").await;
let _ = send.finish();
}
}
});
let client = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Custom(relay_map))
.ca_tls_config(iroh_relay::tls::CaTlsConfig::insecure_skip_verify())
.addr_filter(iroh::address_lookup::AddrFilter::relay_only())
.path_selector(std::sync::Arc::new(super::RelayOnlySelector))
.bind()
.await
.expect("bind client");
let relay_only_addr = iroh::EndpointAddr::from_parts(
server_addr.id,
server_addr
.addrs
.iter()
.filter(|a| matches!(a, iroh::TransportAddr::Relay(_)))
.cloned(),
);
let conn = client
.connect(relay_only_addr, b"mcpmesh/relayonly/test")
.await
.expect("connect over the relay");
let (mut send, mut recv) = conn.open_bi().await.expect("open bi");
send.write_all(b"hi").await.unwrap();
send.finish().unwrap();
let _ = recv.read_to_end(64).await;
let mut direct_available = false;
let mut selected_relay = false;
for _ in 0..40 {
for p in &conn.paths() {
if p.is_ip() {
direct_available = true;
}
if p.is_relay() && p.is_selected() {
selected_relay = true;
}
}
if direct_available && selected_relay {
break;
}
}
assert!(
direct_available,
"SETUP: loopback must offer a direct path, or this proves nothing — the selector \
would trivially pick the only path there is"
);
assert!(
selected_relay,
"relay_only must keep DATA on the relay while a direct path is available; the \
default selector switches to direct here (peer_path.rs asserts that)"
);
})
.await
.expect("relay-only e2e timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_booted_daemon_waits_for_the_relay_before_minting_tickets() {
let dir = tempfile::tempdir().unwrap();
let paths = crate::paths::NodePaths::under_root(dir.path());
std::fs::create_dir_all(paths.config_path.parent().unwrap()).unwrap();
std::fs::write(&paths.config_path, "[network]\nrelay_mode = \"disabled\"\n").unwrap();
let booted = super::boot_node(paths, None)
.await
.expect("the node boots in pairing mode");
let provider = booted
.state
.mesh_required()
.expect("mesh is up")
.app_blobs()
.await
.expect("the app-blob provider must build in pairing mode (#61) — if THIS fails it is a provider-build regression, not a relay-wait one");
assert!(
provider.relay_wait_enabled(),
"boot must enable the relay-ready wait — without it a ticket minted before the relay \
handshake carries direct addresses only: LAN-dialable and NAT-dead, and the sender \
cannot tell (#83 ask 3)"
);
super::shutdown_booted(booted).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn a_booted_daemon_installs_the_configured_presence_mode() {
let boot_with = async |cfg: &str| {
let dir = tempfile::tempdir().unwrap();
let paths = crate::paths::NodePaths::under_root(dir.path());
std::fs::create_dir_all(paths.config_path.parent().unwrap()).unwrap();
std::fs::write(&paths.config_path, cfg).unwrap();
let out = super::boot_node(paths, None).await;
(dir, out)
};
let (_d, booted) =
boot_with("[network]\nrelay_mode = \"disabled\"\npresence_mode = \"off\"\n").await;
let booted = booted.expect("the node boots with presence_mode = off");
assert_eq!(
booted
.state
.mesh_required()
.expect("mesh is up")
.presence_mode(),
PresenceMode::Off,
"boot must INSTALL the configured presence_mode — parsing it and dropping it on the \
floor leaves an operator who asked to be hidden pongging everyone, with nothing to see"
);
let reported = crate::daemon::self_net::read_current(
booted.state.mesh_required().expect("mesh is up"),
None,
);
assert_eq!(
reported.presence_mode.as_deref(),
Some("off"),
"status must report the LIVE presence mode, not a constant and not the on-disk config"
);
super::shutdown_booted(booted).await;
let (_d2, booted2) = boot_with("[network]\nrelay_mode = \"disabled\"\n").await;
let booted2 = booted2.expect("the node boots with no presence_mode set");
assert_eq!(
booted2
.state
.mesh_required()
.expect("mesh is up")
.presence_mode(),
PresenceMode::Paired,
"an unset presence_mode must leave today's behaviour untouched"
);
super::shutdown_booted(booted2).await;
let (_d3, refused) =
boot_with("[network]\nrelay_mode = \"disabled\"\npresence_mode = \"of\"\n").await;
let e = format!(
"{:#}",
refused
.err()
.expect("an unknown presence_mode must refuse to boot, never fall open")
);
assert!(
e.contains("presence_mode") && e.contains("of"),
"and the startup error must name the key and the typo: {e}"
);
}
#[test]
fn pairing_mode_advertises_app_blobs_but_not_gossip_or_roster_blobs() {
let pairing = super::alpns_for(false);
let roster = super::alpns_for(true);
let has = |v: &[Vec<u8>], a: &[u8]| v.iter().any(|x| x.as_slice() == a);
for alpn in [ALPN_MCP, ALPN_PAIR, ALPN_PING] {
assert!(
has(&pairing, alpn),
"every daemon advertises the base ALPNs"
);
}
assert!(
has(&pairing, crate::blobs::APP_BLOB_ALPN),
"a pairing-mode daemon MUST advertise mcpmesh/blob/1 — its scope gate is \
identity-generic and an eid: grant authorizes it (#61)"
);
assert!(
!has(&pairing, crate::roster::transport::GOSSIP_ALPN),
"gossip keys on org_id — never advertised without a roster"
);
assert!(
!has(&pairing, crate::roster::transport::BLOB_ALPN),
"the ROSTER blob transport is distinct from app blobs and stays roster-only"
);
assert!(has(&roster, crate::roster::transport::GOSSIP_ALPN));
assert!(has(&roster, crate::roster::transport::BLOB_ALPN));
assert!(has(&roster, crate::blobs::APP_BLOB_ALPN));
}
#[test]
fn net_plan_validates_the_shipped_network_surface() {
use crate::config::NetworkCfg;
let cfg = |relay: &str, relay_urls: &[&str], disc: &str, disc_urls: &[&str]| NetworkCfg {
relay_mode: relay.into(),
relay_urls: relay_urls.iter().map(|s| s.to_string()).collect(),
discovery_mode: disc.into(),
discovery_urls: disc_urls.iter().map(|s| s.to_string()).collect(),
relay_only: false,
..Default::default()
};
assert!(matches!(
net_plan(&NetworkCfg::default()).unwrap(),
NetPlan::Mesh {
relay: iroh::RelayMode::Default,
discovery: DiscoveryPlan::N0
}
));
assert!(matches!(
net_plan(&cfg("disabled", &[], "default", &[])).unwrap(),
NetPlan::Hermetic
));
let plan = net_plan(&cfg(
"custom",
&["https://relay.acme.com", "https://relay2.acme.com"],
"default",
&[],
))
.unwrap();
match plan {
NetPlan::Mesh {
relay: iroh::RelayMode::Custom(map),
discovery: DiscoveryPlan::N0,
} => {
assert_eq!(map.len(), 2, "both relay_urls land in the RelayMap");
}
other => panic!("expected a custom relay plan, got {other:?}"),
}
let plan = net_plan(&cfg(
"default",
&[],
"custom",
&["https://dns.acme.com/pkarr"],
))
.unwrap();
match plan {
NetPlan::Mesh {
relay: iroh::RelayMode::Default,
discovery: DiscoveryPlan::Custom(urls),
} => {
assert_eq!(urls.len(), 1);
assert_eq!(urls[0].as_str(), "https://dns.acme.com/pkarr");
}
other => panic!("expected a custom discovery plan, got {other:?}"),
}
assert!(net_plan(&cfg("custom", &[], "default", &[])).is_err());
assert!(net_plan(&cfg("custom", &["not a url"], "default", &[])).is_err());
assert!(net_plan(&cfg("default", &[], "custom", &[])).is_err());
assert!(net_plan(&cfg("default", &[], "custom", &["not a url"])).is_err());
assert!(net_plan(&cfg("relayless", &[], "default", &[])).is_err());
assert!(
net_plan(&cfg("default", &[], "local", &[])).is_err(),
"the never-implemented \"local\" mode is refused honestly"
);
}
#[test]
fn the_configured_values_reach_the_transport_config() {
let net = |idle: Option<u64>, keep: Option<u64>| crate::config::NetworkCfg {
idle_timeout_secs: idle,
keep_alive_secs: keep,
..Default::default()
};
let dbg = |idle, keep| {
format!(
"{:?}",
build_transport_config(&net(idle, keep))
.expect("valid config")
.expect("configured")
)
};
assert!(
build_transport_config(&net(None, None)).unwrap().is_none(),
"an unconfigured node must touch NOTHING — iroh's defaults verbatim"
);
let d = dbg(Some(60), Some(3));
assert!(d.contains("60000"), "idle must reach max_idle_timeout: {d}");
assert!(
d.contains("keep_alive_interval: Some(3s)"),
"keepalive must reach keep_alive_interval: {d}"
);
assert!(
d.contains("default_path_keep_alive_interval: Some(3s)"),
"and the PER-PATH keepalive too — every path pings on that value, so setting only the \
connection one makes the knob unable to REDUCE ping frequency: {d}"
);
let z = dbg(Some(0), None);
assert!(
z.contains("max_idle_timeout: None"),
"idle_timeout_secs = 0 must set no-timeout, not leave iroh's 30s: {z}"
);
let k = dbg(None, Some(3));
assert!(
k.contains("max_idle_timeout: Some(30000)"),
"a keepalive alone leaves iroh's idle timeout in place: {k}"
);
}
#[test]
fn iroh_transport_defaults_are_what_the_docs_claim() {
let d = format!("{:?}", iroh::endpoint::QuicTransportConfig::default());
for (needle, doc) in [
("max_idle_timeout: Some(30000)", "30s idle timeout"),
("keep_alive_interval: Some(5s)", "5s keepalive"),
(
"default_path_keep_alive_interval: Some(5s)",
"5s path keepalive",
),
(
"default_path_max_idle_timeout: Some(15s)",
"15s path idle timeout",
),
] {
assert!(
d.contains(needle),
"iroh's default changed: docs/config.md and node/src/config.rs claim {doc} for \
iroh 1.0.3. Re-measure and update BOTH, and note it in the release. NOTE: the \
relay-path idle timeout in that table is NOT pinned here — it never reaches \
QuicTransportConfig — so check it by hand too. Got: {d}"
);
}
assert_eq!(
IROH_DEFAULT_IDLE_SECS, 30,
"the constant a bare keep_alive_secs is validated against must track the real default"
);
let over = format!(
"{:?}",
iroh::endpoint::QuicTransportConfig::builder()
.default_path_keep_alive_interval(Duration::from_secs(
IROH_MAX_PATH_KEEP_ALIVE_SECS + 1
))
.build()
);
assert!(
over.contains(&format!(
"default_path_keep_alive_interval: Some({IROH_MAX_PATH_KEEP_ALIVE_SECS}s)"
)),
"iroh no longer caps the per-path keepalive at {IROH_MAX_PATH_KEEP_ALIVE_SECS}s. \
Raising keep_alive_secs may now genuinely reduce ping traffic — revisit the refusal \
in build_transport_config and the metered-link note in docs/config.md. Got: {over}"
);
}
#[test]
fn a_zero_service_rate_is_a_startup_error() {
let cfg = |rate: &str| {
crate::config::Config::from_toml_str(&format!(
"[services.kb]\nsocket = \"/run/kb.sock\"\nallow = []\n{rate}"
))
.unwrap()
};
validate_service_rates(&cfg("")).expect("an absent rate inherits the global");
validate_service_rates(&cfg("rate_limit_per_min = 1\n")).expect("1 is the floor and legal");
validate_service_rates(&cfg("rate_limit_per_min = 500\n"))
.expect("above the ceiling is CLAMPED later, not refused here");
let e = validate_service_rates(&cfg("rate_limit_per_min = 0\n"))
.unwrap_err()
.to_string();
assert!(
e.contains("kb") && e.contains("at least 1"),
"the error must name the service and the floor: {e}"
);
assert!(
e.contains("does NOT mean unlimited"),
"and correct the reading, since `blob_bytes_per_min = 0` in the same file DOES mean \
unlimited: {e}"
);
}
#[test]
fn an_unknown_presence_mode_is_a_startup_error() {
let with = |m: &str| crate::config::NetworkCfg {
presence_mode: m.into(),
..Default::default()
};
for good in ["paired", "granted", "off"] {
presence_mode(&with(good)).unwrap_or_else(|e| panic!("{good:?} must parse: {e:#}"));
}
assert_eq!(
presence_mode(&crate::config::NetworkCfg::default()).unwrap(),
PresenceMode::Paired,
"the DEFAULT must stay today's behaviour — this knob must not change anyone silently"
);
for bad in ["of", "OFF", "none", "disabled", "true", ""] {
let e = presence_mode(&with(bad)).unwrap_err().to_string();
assert!(
e.contains(bad) && e.contains("presence_mode"),
"the error must name the key AND the bad value: {e}"
);
assert!(
e.contains("paired") && e.contains("granted") && e.contains("off"),
"and list every legal value, or the operator guesses again: {e}"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_keepalive_at_or_above_the_idle_timeout_is_refused() {
let net = |idle: u64, keep: u64| crate::config::NetworkCfg {
relay_mode: "disabled".into(),
idle_timeout_secs: Some(idle),
keep_alive_secs: Some(keep),
..Default::default()
};
let key = || iroh::SecretKey::from_bytes(&[19u8; 32]);
let cfg = |idle: Option<u64>, keep: Option<u64>| crate::config::NetworkCfg {
relay_mode: "disabled".into(),
idle_timeout_secs: idle,
keep_alive_secs: keep,
..Default::default()
};
let e = build_transport_config(&cfg(Some(1200), Some(64)))
.expect_err("a keepalive above iroh's per-path cap must be refused, not silently sunk");
let msg = format!("{e:#}");
assert!(
msg.contains("64"),
"the error must name THEIR value, not just the cap: {msg}"
);
assert!(
msg.contains("cap of 5s"),
"and the cap itself, as a number they can act on: {msg}"
);
assert!(
msg.contains("cannot reduce"),
"and say plainly that raising it does NOT reduce keepalive traffic — that is the whole \
reason someone sets it: {msg}"
);
build_transport_config(&cfg(Some(1200), Some(6)))
.expect_err("6s is above iroh's cap — iroh would drop it, so boot must refuse it");
build_transport_config(&cfg(Some(1200), Some(IROH_MAX_PATH_KEEP_ALIVE_SECS)))
.expect("5s is exactly the cap — iroh keeps it, so refusing it would be wrong");
let z = format!(
"{:#}",
build_transport_config(&cfg(Some(1200), Some(0)))
.expect_err("keep_alive_secs = 0 arms a zero-length timer and must be refused")
);
assert!(
z.contains("not \"disable keepalives\"") && z.contains("omit the key"),
"and must say what 0 really does AND how to actually get the default: {z}"
);
build_transport_config(&cfg(Some(0), Some(3)))
.expect("no idle timeout means no keepalive can outlive it — this must be allowed");
for (idle, keep) in [(5, 5), (4, 5)] {
let e = build_endpoint(key(), &net(idle, keep), false)
.await
.expect_err("a keepalive that cannot arrive in time must be refused at boot");
let msg = format!("{e:#}");
assert!(
msg.contains(&keep.to_string()) && msg.contains(&idle.to_string()),
"the error must name BOTH values — an operator cannot fix what it does not \
identify: {msg}"
);
}
let ep = build_endpoint(key(), &net(30, 5), false)
.await
.expect("keepalive below the timeout is the working configuration");
ep.close().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn absent_transport_config_leaves_iroh_defaults_alone() {
let net = crate::config::NetworkCfg {
relay_mode: "disabled".into(),
..Default::default()
};
assert_eq!(net.idle_timeout_secs, None);
assert_eq!(net.keep_alive_secs, None);
let ep = build_endpoint(iroh::SecretKey::from_bytes(&[20u8; 32]), &net, false)
.await
.expect("an unconfigured node still binds");
ep.close().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn a_zero_idle_timeout_means_no_timeout_and_is_allowed() {
let net = crate::config::NetworkCfg {
relay_mode: "disabled".into(),
idle_timeout_secs: Some(0),
..Default::default()
};
let ep = build_endpoint(iroh::SecretKey::from_bytes(&[21u8; 32]), &net, false)
.await
.expect("0 = no idle timeout is a valid QUIC configuration");
ep.close().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn build_endpoint_binds_with_a_custom_relay_map() {
let net = crate::config::NetworkCfg {
relay_mode: "custom".into(),
relay_urls: vec!["https://relay.acme.com".into()],
discovery_mode: "custom".into(),
discovery_urls: vec!["https://dns.acme.com/pkarr".into()],
relay_only: false,
..Default::default()
};
let ep = build_endpoint(iroh::SecretKey::from_bytes(&[9u8; 32]), &net, false)
.await
.expect("custom relay+discovery endpoint binds offline");
ep.close().await;
}
}