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 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;
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 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()))?,
};
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 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::new());
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());
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);
if roster_mode {
let scopes_path = paths.blob_scopes_path.clone();
match blocking("join app-blob scopes load", move || {
crate::blobs::scope::ScopeStore::load(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(),
)
.await
{
Ok(provider) => 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")
}
}
}
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,
},
}
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 mut alpns = vec![ALPN_MCP.to_vec(), ALPN_PAIR.to_vec(), ALPN_PING.to_vec()];
if roster_mode {
alpns.push(crate::roster::transport::GOSSIP_ALPN.to_vec());
alpns.push(crate::roster::transport::BLOB_ALPN.to_vec());
alpns.push(crate::blobs::APP_BLOB_ALPN.to_vec());
}
builder
.secret_key(secret)
.alpns(alpns)
.bind()
.await
.context("bind iroh endpoint")
}
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 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(),
};
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"
);
}
#[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()],
};
let ep = build_endpoint(iroh::SecretKey::from_bytes(&[9u8; 32]), &net, false)
.await
.expect("custom relay+discovery endpoint binds offline");
ep.close().await;
}
}