pub(crate) mod config_write;
mod dial;
use std::collections::HashMap;
use std::fs::File;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use mcpmesh_local_api::{
BackendKind, BackendSpec, BlobFetchResult, BlobPublishResult, BlobScopeList, InviteResult,
OrgJoinResult, PairResult, PeerInfo, PresencePeer, RosterInstallResult, RosterStatus,
ScopeInfo, ServiceInfo,
};
use mcpmesh_net::errors::{ERR_UNREACHABLE, synthesized};
use mcpmesh_net::framing::{FrameReader, write_frame};
use mcpmesh_net::registry::ConnRegistry;
use mcpmesh_net::{
ALPN_MCP, ALPN_PAIR, ServiceEntry, Services, SessionBackend, TrustGate, run_mesh_connection,
};
use mcpmesh_trust::roster::validate::RosterState;
use mcpmesh_trust::{DeviceKey, paths};
use serde_json::Value;
use tokio::io::{AsyncRead, AsyncWrite};
use tokio::sync::Semaphore;
use tokio::task::JoinHandle;
pub use dial::{dial_service, pipe_session, race_dial};
use self::config_write::{
append_allow_to_config, remove_allow_from_config, rename_allow_in_config, write_identity_pin,
write_identity_user_id, write_join_pin, write_roster_url, write_service_to_config,
};
use crate::util::{epoch_now_i64, epoch_now_u64};
use crate::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use crate::audit::{AuditLog, AuditRecord, AuditSink, now_ts};
use crate::backends::socket::SocketBackend;
use crate::backends::spawn::SpawnBackend;
use crate::config::{Backend, Config};
use crate::control::{DaemonState, serve_control};
use crate::ipc;
use crate::pairing::{self, Invite, LiveInvites};
use crate::roster::RosterStore;
use crate::roster::freshness::FreshnessStore;
use crate::roster::gate::{ComposedGate, RosterGate};
pub const STACK_VERSION: &str = env!("CARGO_PKG_VERSION");
#[allow(dead_code)] const SPAWN_CONCURRENCY: usize = 4;
pub(crate) fn spawn_concurrency(cfg: &Config) -> usize {
(cfg.limits.max_sessions.max(1)) as usize
}
const INVITE_TTL: Duration = Duration::from_secs(24 * 60 * 60);
const RELAY_READY_TIMEOUT: Duration = Duration::from_secs(3);
pub struct MeshState {
pub(crate) endpoint: iroh::Endpoint,
pub(crate) gate: Arc<dyn TrustGate>,
pub(crate) store: Arc<PeerStore>,
pub(crate) invites: Arc<LiveInvites>,
pub(crate) self_petname: String,
pub(crate) accept_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
pub(crate) poll_loop: tokio::sync::Mutex<Option<JoinHandle<()>>>,
pub(crate) reload_lock: tokio::sync::Mutex<()>,
pub(crate) config_path: PathBuf,
pub(crate) roster: Arc<RosterGate>,
pub(crate) conn_registry: Arc<ConnRegistry>,
pub(crate) gossip: Option<iroh_gossip::net::Gossip>,
pub(crate) blobs: Option<crate::roster::transport::RosterBlobs>,
pub(crate) roster_topic: tokio::sync::Mutex<Option<crate::roster::transport::RosterGossip>>,
pub(crate) presence_topic: tokio::sync::Mutex<Option<crate::roster::transport::RosterGossip>>,
pub(crate) presence_table: Arc<crate::roster::presence::PresenceTable>,
pub(crate) app_blobs: tokio::sync::Mutex<Option<Arc<crate::blobs::provider::AppBlobs>>>,
pub(crate) audit: std::sync::OnceLock<AuditSink>,
pub(crate) limits: std::sync::OnceLock<Arc<crate::limits::MeshLimiters>>,
pub(crate) roster_addr_book:
std::sync::OnceLock<std::sync::Arc<crate::roster::transport::RosterAddrBook>>,
pub(crate) self_binding: std::sync::OnceLock<Option<crate::pairing::rendezvous::SelfBinding>>,
pub(crate) recent_pairings:
std::sync::Mutex<std::collections::VecDeque<mcpmesh_local_api::RecentPairing>>,
}
const RECENT_PAIRINGS_CAP: usize = 8;
impl MeshState {
#[allow(clippy::too_many_arguments)]
pub fn new(
endpoint: iroh::Endpoint,
gate: Arc<dyn TrustGate>,
store: Arc<PeerStore>,
invites: Arc<LiveInvites>,
self_petname: String,
config_path: PathBuf,
roster: Arc<RosterGate>,
conn_registry: Arc<ConnRegistry>,
gossip: Option<iroh_gossip::net::Gossip>,
blobs: Option<crate::roster::transport::RosterBlobs>,
roster_topic: Option<crate::roster::transport::RosterGossip>,
presence_topic: Option<crate::roster::transport::RosterGossip>,
) -> Arc<Self> {
Arc::new(Self {
endpoint,
gate,
store,
invites,
self_petname,
accept_task: tokio::sync::Mutex::new(None),
poll_loop: tokio::sync::Mutex::new(None),
reload_lock: tokio::sync::Mutex::new(()),
config_path,
roster,
conn_registry,
gossip,
blobs,
roster_topic: tokio::sync::Mutex::new(roster_topic),
presence_topic: tokio::sync::Mutex::new(presence_topic),
presence_table: Arc::new(crate::roster::presence::PresenceTable::new()),
app_blobs: tokio::sync::Mutex::new(None),
audit: std::sync::OnceLock::new(),
limits: std::sync::OnceLock::new(),
roster_addr_book: std::sync::OnceLock::new(),
self_binding: std::sync::OnceLock::new(),
recent_pairings: std::sync::Mutex::new(std::collections::VecDeque::new()),
})
}
pub(crate) fn record_pairing(
&self,
peer_petname: String,
sas_code: String,
paired_at_epoch: u64,
) {
let mut ring = self
.recent_pairings
.lock()
.expect("recent_pairings lock not poisoned");
if ring.len() >= RECENT_PAIRINGS_CAP {
ring.pop_front();
}
ring.push_back(mcpmesh_local_api::RecentPairing {
peer_petname,
sas_code,
paired_at_epoch,
});
}
pub(crate) fn recent_pairings(&self) -> Vec<mcpmesh_local_api::RecentPairing> {
self.recent_pairings
.lock()
.expect("recent_pairings lock not poisoned")
.iter()
.rev()
.cloned()
.collect()
}
pub async fn set_app_blobs(&self, provider: Arc<crate::blobs::provider::AppBlobs>) {
*self.app_blobs.lock().await = Some(provider);
}
pub async fn app_blobs(&self) -> Option<Arc<crate::blobs::provider::AppBlobs>> {
self.app_blobs.lock().await.clone()
}
pub fn set_audit(&self, sink: AuditSink) {
let _ = self.audit.set(sink);
}
pub(crate) fn audit(&self) -> AuditSink {
self.audit.get().cloned().unwrap_or_default()
}
pub fn set_self_binding(&self, binding: Option<crate::pairing::rendezvous::SelfBinding>) {
let _ = self.self_binding.set(binding);
}
pub(crate) fn self_binding(&self) -> Option<crate::pairing::rendezvous::SelfBinding> {
self.self_binding.get().cloned().flatten()
}
pub fn set_limits(&self, limits: Arc<crate::limits::MeshLimiters>) {
let _ = self.limits.set(limits);
}
pub(crate) fn limits(&self) -> Arc<crate::limits::MeshLimiters> {
self.limits
.get()
.cloned()
.unwrap_or_else(crate::limits::MeshLimiters::unlimited)
}
pub(crate) fn roster_addr_book(
&self,
) -> Option<std::sync::Arc<crate::roster::transport::RosterAddrBook>> {
self.roster_addr_book.get().cloned()
}
pub async fn roster_topic_sender(&self) -> Option<iroh_gossip::api::GossipSender> {
self.roster_topic
.lock()
.await
.as_ref()
.map(|g| g.sender.clone())
}
pub async fn take_roster_topic_receiver(&self) -> Option<iroh_gossip::api::GossipReceiver> {
self.roster_topic
.lock()
.await
.as_mut()
.and_then(|g| g.receiver.take())
}
pub async fn presence_topic_sender(&self) -> Option<iroh_gossip::api::GossipSender> {
self.presence_topic
.lock()
.await
.as_ref()
.map(|g| g.sender.clone())
}
pub async fn take_presence_topic_receiver(&self) -> Option<iroh_gossip::api::GossipReceiver> {
self.presence_topic
.lock()
.await
.as_mut()
.and_then(|g| g.receiver.take())
}
pub(crate) async fn confirm_roster_current(&self, now: i64) {
self.roster.set_last_confirmed(now);
let store = FreshnessStore::new(roster_confirmed_path(&self.config_path));
match tokio::task::spawn_blocking(move || store.store(now)).await {
Ok(Ok(())) => {}
Ok(Err(e)) => {
tracing::warn!(%e, "persist roster freshness (in-memory freshness still applied)")
}
Err(e) => tracing::warn!(%e, "join roster freshness persist"),
}
}
pub async fn set_accept_task(&self, handle: JoinHandle<()>) {
let mut guard = self.accept_task.lock().await;
if let Some(old) = guard.take() {
old.abort();
}
*guard = Some(handle);
}
}
pub fn run() -> Result<()> {
let runtime = paths::runtime_dir()?;
ipc::ensure_runtime_dir(&runtime)?;
let lock_path = runtime.join("mcpmesh.lock");
let Some(_lock) = acquire_singleton_lock(&lock_path)? else {
tracing::info!("another mcpmesh daemon already holds the singleton lock; exiting");
return Ok(());
};
let socket = paths::default_socket_path()?;
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.context("build daemon tokio runtime")?;
rt.block_on(async move { serve_forever(&socket).await })
}
async fn serve_forever(socket: &Path) -> Result<()> {
let _ = rustls::crypto::ring::default_provider().install_default();
let audit = AuditSink::new(AuditLog::spawn(paths::default_audit_dir()?));
let config_path = paths::default_config_path()?;
let cfg = Config::load(&config_path)
.map_err(|e| anyhow::anyhow!("config error in {}: {e}", config_path.display()))?;
let key_path = match cfg.identity.device_key.clone() {
Some(p) => p,
None => paths::default_device_key_path()?,
};
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::default_state_db_path()?;
if let Some(parent) = db_path.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("create data dir {}", parent.display()))?;
}
let store = tokio::task::spawn_blocking(move || PeerStore::open(&db_path))
.await
.context("join peer-store open")??;
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::default_roster_path()?);
match tokio::task::spawn_blocking(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)) => {
tracing::warn!(%e, "installed roster failed to load; refusing roster peers")
}
Err(e) => tracing::warn!(%e, "join roster load"),
}
}
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 tokio::task::spawn_blocking(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 tokio::task::spawn_blocking(move || fstore.store(now)).await {
Ok(Ok(())) => tracing::info!("roster freshness upgrade grace applied"),
Ok(Err(e)) => tracing::warn!(%e, "persist roster freshness upgrade grace"),
Err(e) => tracing::warn!(%e, "join roster freshness upgrade-grace persist"),
}
}
Ok(Ok(None)) => {}
Ok(Err(e)) => {
tracing::warn!(%e, "read roster freshness sidecar; treating as unconfirmed")
}
Err(e) => tracing::warn!(%e, "join roster freshness load"),
}
}
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 service_list = service_infos(&cfg);
let store_for_list = store.clone();
let peer_list = tokio::task::spawn_blocking(move || peer_infos(&store_for_list))
.await
.context("join peer list")?;
let invites = Arc::new(LiveInvites::new());
let self_petname = cfg
.identity
.petname
.clone()
.unwrap_or_else(|| default_self_petname(&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_petname,
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::default_user_key_path()?,
};
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::default_blob_scopes_path()?;
match tokio::task::spawn_blocking(move || {
crate::blobs::scope::ScopeStore::load(scopes_path)
})
.await
{
Ok(Ok(scopes)) => {
match crate::blobs::provider::AppBlobs::load(
paths::default_blobs_dir()?,
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)) => tracing::warn!(%e, "app-blob scopes failed to load; provider disabled"),
Err(e) => tracing::warn!(%e, "join app-blob scopes load"),
}
}
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);
let _converge = crate::roster::distribute::spawn_receive_loop(mesh.clone());
}
if roster_mode {
let _track = crate::roster::presence::track_loop(mesh.clone());
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());
let _publish =
crate::roster::presence::publish_loop(mesh.clone(), device_key, user_id);
}
None => tracing::debug!(
"presence publish skipped: no user_id for this node (track loop still runs)"
),
}
}
if roster_mode {
let _sweep = spawn_staleness_sweep(mesh.clone());
}
if roster_mode && let Some(url) = cfg.roster.url.clone() {
respawn_poll_loop(&mesh, url).await;
}
let listener = ipc::bind_control_socket(socket).await?;
let state = Arc::new(DaemonState::with_mesh(
STACK_VERSION,
mesh,
service_list,
peer_list,
));
tracing::info!(
endpoint_id = %our_id,
socket = %socket.display(),
"mcpmesh daemon serving mesh + control"
);
serve_control(listener, state).await
}
fn warn_if_degraded_grace(roster: &RosterGate) {
if roster.effective_state(epoch_now_i64()) == Some(RosterState::DegradedGrace) {
tracing::warn!(
"roster degraded (expired or stale); serving continues for the grace period — \
re-confirm currency or install a fresh roster"
);
}
}
pub(crate) fn roster_status(mesh: &Arc<MeshState>, cfg: Option<&Config>) -> Option<RosterStatus> {
let org_root_fingerprint = cfg
.and_then(|c| c.identity.org_root_pk.as_deref())
.and_then(|s| crate::roster::parse_org_root_pk(s).ok())
.map(|vk| pairing::sas::fingerprint_words(&vk.to_bytes()))
.unwrap_or_default();
match mesh.roster.view() {
Some(view) => {
let state = match mesh
.roster
.effective_state(epoch_now_i64())
.unwrap_or(RosterState::Approved)
{
RosterState::Approved => "approved",
RosterState::DegradedGrace => "degraded",
RosterState::DegradedStopped => "stopped",
};
Some(RosterStatus {
org_id: view.org_id().to_string(),
serial: view.serial(),
state: state.to_string(),
org_root_fingerprint,
})
}
None => {
let cfg = cfg?;
cfg.identity.org_root_pk.as_deref()?;
Some(RosterStatus {
org_id: cfg.identity.org_id.clone().unwrap_or_default(),
serial: 0,
state: "pending".to_string(),
org_root_fingerprint,
})
}
}
}
pub(crate) fn presence_peers(mesh: &Arc<MeshState>) -> Vec<PresencePeer> {
let Some(view) = mesh.roster.view() else {
return Vec::new();
};
let now = epoch_now_i64();
let online: std::collections::HashSet<[u8; 32]> = mesh
.presence_table
.active(now)
.into_iter()
.map(|(eid, _)| eid)
.collect();
let mut peers: Vec<PresencePeer> = view
.devices()
.map(|(eid, d)| PresencePeer {
user_id: d.user_id.clone(),
device_label: d.label.clone(),
role: d.role.clone(),
online: online.contains(eid),
})
.collect();
peers.sort_by(|a, b| {
a.user_id
.cmp(&b.user_id)
.then_with(|| dial::dial_role_rank(&a.role).cmp(&dial::dial_role_rank(&b.role)))
.then_with(|| a.device_label.cmp(&b.device_label))
});
peers
}
pub fn install_roster_view_and_sever(
mesh: &Arc<MeshState>,
view: mcpmesh_trust::roster::validate::RosterView,
) -> usize {
use std::collections::HashSet;
let (org_id, serial) = (view.org_id().to_string(), view.serial());
let revoked: HashSet<[u8; 32]> = view.revoked_endpoints().copied().collect();
let active_devices: HashSet<[u8; 32]> = view.device_endpoints().copied().collect();
mesh.roster.install(view);
warn_if_degraded_grace(&mesh.roster);
let severed = mesh.conn_registry.sever_matching(
mcpmesh_net::CLOSE_UNAUTHORIZED, b"roster revoked",
|eid, roster_user| mcpmesh_net::should_sever(eid, roster_user, &revoked, &active_devices),
);
mesh.audit().record(AuditRecord::trust(
now_ts(),
"roster_install".into(),
Some(format!("{org_id}/{serial}")),
));
severed
}
const STALENESS_SWEEP_INTERVAL: Duration = Duration::from_secs(60);
pub fn should_staleness_sever(state: Option<mcpmesh_trust::roster::validate::RosterState>) -> bool {
state == Some(mcpmesh_trust::roster::validate::RosterState::DegradedStopped)
}
pub fn staleness_sweep_once(mesh: &Arc<MeshState>, now: i64) -> usize {
if !should_staleness_sever(mesh.roster.effective_state(now)) {
return 0;
}
mesh.conn_registry.sever_matching(
mcpmesh_net::CLOSE_UNAUTHORIZED, b"roster stale",
|_eid, roster_user| roster_user.is_some(),
)
}
fn spawn_staleness_sweep(mesh: Arc<MeshState>) -> JoinHandle<()> {
tokio::spawn(async move {
loop {
tokio::time::sleep(STALENESS_SWEEP_INTERVAL).await;
let severed = staleness_sweep_once(&mesh, epoch_now_i64());
if severed > 0 {
tracing::warn!(
severed,
"staleness sweep cut roster-authorized sessions (degraded-stopped)"
);
}
}
})
}
pub(crate) async fn install_roster(
state: &DaemonState,
path: String,
org_root_pk: Option<String>,
) -> Result<RosterInstallResult> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let _reload = mesh.reload_lock.lock().await;
let pk_str = match &org_root_pk {
Some(s) => s.clone(),
None => mesh_config_org_root_pk(mesh)?
.context("no org root pinned; pass --org-root-pk on first install")?,
};
let pk = crate::roster::parse_org_root_pk(&pk_str)?;
let rstore = RosterStore::new(paths::default_roster_path()?);
let now = epoch_now_i64();
let file = PathBuf::from(path);
let view = tokio::task::spawn_blocking(move || rstore.install_from_file(&file, &pk, now))
.await
.context("join roster install")??;
let (org_id, serial) = (view.org_id().to_string(), view.serial());
if org_root_pk.is_some() {
let config_path = mesh.config_path.clone();
let (pk_w, oid_w) = (pk_str.clone(), org_id.clone());
tokio::task::spawn_blocking(move || write_identity_pin(&config_path, &pk_w, &oid_w))
.await
.context("join org-root pin config write")??;
tracing::info!(org_id = %org_id, "pinned org root");
}
mesh.confirm_roster_current(now).await;
let severed = install_roster_view_and_sever(mesh, view) as u32;
reconcile_user_id_from_roster(mesh).await;
tracing::info!(org_id = %org_id, serial, severed, "roster installed");
if let Err(e) = crate::roster::distribute::announce_roster(mesh).await {
tracing::warn!(%e, "roster announce-on-publish failed (install still applied)");
}
Ok(RosterInstallResult {
org_id,
serial,
severed,
})
}
pub(crate) async fn blob_publish(
state: &DaemonState,
scope: String,
path: String,
) -> Result<BlobPublishResult> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let provider = mesh
.app_blobs()
.await
.context("app-blob provider not enabled (roster mode only)")?;
let (ticket, hash) = provider
.publish_scope(&scope, Path::new(&path))
.await
.context("publish blob into scope")?;
Ok(BlobPublishResult { ticket, hash })
}
pub(crate) async fn blob_grant(
state: &DaemonState,
scope: String,
principal: String,
) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let provider = mesh
.app_blobs()
.await
.context("app-blob provider not enabled (roster mode only)")?;
provider.grant(&scope, &principal)
}
pub(crate) async fn blob_list(state: &DaemonState) -> Result<BlobScopeList> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let scopes = match mesh.app_blobs().await {
Some(provider) => provider
.list()
.into_iter()
.map(|(name, hashes, grants)| ScopeInfo {
name,
hashes,
grants,
})
.collect(),
None => Vec::new(),
};
Ok(BlobScopeList { scopes })
}
pub(crate) async fn blob_fetch(
state: &DaemonState,
ticket: String,
dest_path: String,
) -> Result<BlobFetchResult> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let provider = mesh
.app_blobs()
.await
.context("app-blob provider not enabled (roster mode only)")?;
let hash = provider.fetch(&ticket).await.context("fetch blob")?;
let bytes = provider
.read_bytes(hash)
.await
.context("read fetched blob")?;
let bytes_len = bytes.len() as u64;
let dest = PathBuf::from(dest_path);
tokio::fs::write(&dest, &bytes)
.await
.with_context(|| format!("write fetched blob to {}", dest.display()))?;
Ok(BlobFetchResult {
hash: hash.to_hex().to_string(),
bytes_len,
})
}
pub(crate) fn mesh_config_org_root_pk(mesh: &Arc<MeshState>) -> Result<Option<String>> {
let cfg = Config::load(&mesh.config_path)
.map_err(|e| anyhow::anyhow!("read config for org_root_pk: {e}"))?;
Ok(cfg.identity.org_root_pk)
}
pub(crate) fn installed_roster_path(mesh: &Arc<MeshState>) -> PathBuf {
mesh.config_path
.parent()
.map(|dir| dir.join("roster.json"))
.unwrap_or_else(|| PathBuf::from("roster.json")) }
fn roster_confirmed_path(config_path: &Path) -> PathBuf {
config_path
.parent()
.map(|dir| dir.join("roster.confirmed"))
.unwrap_or_else(|| PathBuf::from("roster.confirmed")) }
pub(crate) async fn reconcile_user_id_from_roster(mesh: &Arc<MeshState>) {
let our_id = mesh.endpoint.id();
let Some(roster_user_id) = mesh
.roster
.view()
.and_then(|v| v.resolve(our_id.as_bytes()).map(|d| d.user_id.clone()))
else {
return;
};
let config_path = mesh.config_path.clone();
let proposed = Config::load(&config_path)
.ok()
.and_then(|c| c.identity.user_id);
if proposed.as_deref() == Some(roster_user_id.as_str()) {
return;
}
let uid = roster_user_id.clone();
match tokio::task::spawn_blocking(move || write_identity_user_id(&config_path, &uid)).await {
Ok(Ok(())) => {
tracing::info!(user_id = %roster_user_id, "reconciled config user_id from the authoritative roster")
}
Ok(Err(e)) => tracing::warn!(%e, "reconcile config user_id write failed"),
Err(e) => tracing::warn!(%e, "join reconcile config user_id write"),
}
}
pub(crate) fn write_temp_roster(bytes: &[u8]) -> Result<PathBuf> {
use std::sync::atomic::{AtomicU64, Ordering};
static SEQ: AtomicU64 = AtomicU64::new(0);
let seq = SEQ.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"mcpmesh-roster-in.{}.{}.json",
std::process::id(),
seq
));
let mut f =
File::create(&path).with_context(|| format!("create temp roster {}", path.display()))?;
f.write_all(bytes).context("write temp roster")?;
f.sync_all().context("sync temp roster")?;
Ok(path)
}
pub(crate) async fn org_join(
state: &DaemonState,
org_id: String,
org_root_pk: String,
user_id: String,
user_key: String,
) -> Result<OrgJoinResult> {
crate::roster::parse_org_root_pk(&org_root_pk)?;
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let (oid, pk, uid, uk) = (org_id.clone(), org_root_pk, user_id, user_key);
tokio::task::spawn_blocking(move || write_join_pin(&config_path, &oid, &pk, &uid, &uk))
.await
.context("join org-root pin config write")??;
tracing::info!(org_id = %org_id, "pinned org root (join)");
Ok(OrgJoinResult { org_id })
}
pub(crate) async fn set_roster_url(state: &DaemonState, url: String) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let url_w = url.clone();
tokio::task::spawn_blocking(move || write_roster_url(&config_path, &url_w))
.await
.context("roster-url config write")??;
tracing::info!("pinned roster url");
if Config::load(&mesh.config_path)
.map(|c| c.identity.org_root_pk.is_some())
.unwrap_or(false)
{
respawn_poll_loop(mesh, url).await;
tracing::info!("roster URL poll loop (re)started");
}
Ok(())
}
#[derive(Debug)]
pub(crate) enum DiscoveryPlan {
N0,
Custom(Vec<url::Url>),
}
#[derive(Debug)]
pub(crate) enum NetPlan {
Hermetic,
Mesh {
relay: iroh::RelayMode,
discovery: DiscoveryPlan,
},
}
pub(crate) 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 })
}
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()];
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)
}
fn gate_and_register(
mesh: &Arc<MeshState>,
conn: &iroh::endpoint::Connection,
blob_conn_limit: bool,
) -> Option<mcpmesh_net::registry::Registration> {
let remote = *conn.remote_id().as_bytes();
if mesh.gate.resolve(&remote).is_none() {
conn.close(mcpmesh_net::CLOSE_UNAUTHORIZED.into(), b"unauthorized");
return None;
}
if blob_conn_limit && !mesh.limits().admit_blob_conn(&remote) {
conn.close(0u32.into(), b"blob rate limited");
return None;
}
let roster_user = mesh.gate.roster_user(&remote);
let registration = mesh
.conn_registry
.register_checked(conn, roster_user.clone(), |eid| {
mesh.gate.should_sever_now(eid, roster_user.as_deref())
});
if registration.is_none() {
conn.close(mcpmesh_net::CLOSE_UNAUTHORIZED.into(), b"unauthorized");
}
registration
}
pub fn spawn_accept_loop(mesh: Arc<MeshState>, services: Arc<Services>) -> JoinHandle<()> {
tokio::spawn(async move {
while let Some(incoming) = mesh.endpoint.accept().await {
let (mesh, services) = (mesh.clone(), services.clone());
tokio::spawn(async move {
let conn = match incoming.await {
Ok(conn) => conn,
Err(e) => {
tracing::debug!(%e, "inbound handshake failed");
return;
}
};
let alpn = conn.alpn().to_vec();
match alpn.as_slice() {
a if a == ALPN_MCP => {
run_mesh_connection(
conn,
mesh.gate.clone(),
services,
mesh.conn_registry.clone(),
)
.await;
}
a if a == ALPN_PAIR => {
if mesh.invites.count() == 0 {
conn.close(0u32.into(), b"no pairing in progress");
return;
}
if !mesh.limits().admit_pair_accept() {
conn.close(0u32.into(), b"pair rate limited");
return;
}
if let Err(e) =
pairing::rendezvous::handle_inviter_side(conn, mesh.clone()).await
{
tracing::debug!(%e, "pair rendezvous error");
}
}
a if a == crate::roster::transport::GOSSIP_ALPN => {
let Some(gossip) = mesh.gossip.clone() else {
conn.close(0u32.into(), b"gossip not enabled");
return;
};
let Some(_registration) = gate_and_register(&mesh, &conn, false) else {
return;
};
if let Err(e) = iroh::protocol::ProtocolHandler::accept(&gossip, conn).await
{
tracing::debug!(%e, "gossip accept error");
}
}
a if a == crate::roster::transport::BLOB_ALPN => {
let Some(blobs) = mesh.blobs.clone() else {
conn.close(0u32.into(), b"blobs not enabled");
return;
};
let Some(_registration) = gate_and_register(&mesh, &conn, false) else {
return;
};
let blob_proto = blobs.protocol();
if let Err(e) =
iroh::protocol::ProtocolHandler::accept(&blob_proto, conn).await
{
tracing::debug!(%e, "blob accept error");
}
}
a if a == crate::blobs::APP_BLOB_ALPN => {
let Some(app_blobs) = mesh.app_blobs().await else {
conn.close(0u32.into(), b"app blobs not enabled");
return;
};
let Some(_registration) = gate_and_register(&mesh, &conn, true) else {
return;
};
let blob_proto = app_blobs.protocol();
if let Err(e) =
iroh::protocol::ProtocolHandler::accept(&blob_proto, conn).await
{
tracing::debug!(%e, "app-blob accept error");
}
}
_ => conn.close(0u32.into(), b"unknown alpn"),
}
});
}
})
}
async fn reload_accept_loop(mesh: &Arc<MeshState>, services: Services) {
let mut guard = mesh.accept_task.lock().await;
if let Some(old) = guard.take() {
old.abort();
}
*guard = Some(spawn_accept_loop(mesh.clone(), Arc::new(services)));
}
async fn reload_services_from_disk(mesh: &Arc<MeshState>, why: &str) -> Result<Config> {
let cfg = Config::load(&mesh.config_path)
.map_err(|e| anyhow::anyhow!("reload config after {why}: {e}"))?;
reload_accept_loop(
mesh,
build_services_audited(&cfg, &mesh.audit(), &mesh.limits()),
)
.await;
Ok(cfg)
}
pub(crate) async fn respawn_poll_loop(mesh: &Arc<MeshState>, url: String) {
let interval = Config::load(&mesh.config_path)
.map(|c| c.roster.poll_interval_seconds())
.unwrap_or(3600);
let mut guard = mesh.poll_loop.lock().await;
if let Some(old) = guard.take() {
old.abort();
}
*guard = Some(crate::roster::distribute::spawn_poll_loop(
mesh.clone(),
url,
interval,
));
}
pub fn build_services(cfg: &Config) -> Services {
build_services_audited(
cfg,
&AuditSink::disabled(),
&crate::limits::MeshLimiters::unlimited(),
)
}
pub fn build_services_audited(
cfg: &Config,
audit: &AuditSink,
limiters: &Arc<crate::limits::MeshLimiters>,
) -> Services {
let mut map: HashMap<String, ServiceEntry> = HashMap::new();
for (name, svc) in &cfg.services {
let backend: Arc<dyn SessionBackend> = match svc.backend_result() {
Ok(Backend::Run(cmd)) => Arc::new(SpawnBackend {
cmd: cmd.to_vec(),
concurrency: Arc::new(Semaphore::new(spawn_concurrency(cfg))),
service: name.clone(),
audit: audit.clone(),
limiter: limiters.requests.clone(),
}),
Ok(Backend::Socket(path)) => Arc::new(SocketBackend {
path: path.to_string(),
service: name.clone(),
audit: audit.clone(),
limiter: limiters.requests.clone(),
}),
Err(e) => {
tracing::warn!(service = %name, %e, "skipping malformed service");
continue;
}
};
map.insert(
name.clone(),
ServiceEntry {
backend,
allow: svc.allow.clone(),
},
);
}
Services::new(map)
}
pub(crate) fn service_infos(cfg: &Config) -> Vec<ServiceInfo> {
cfg.services
.iter()
.filter_map(|(name, svc)| {
let backend = match svc.backend_result() {
Ok(Backend::Run(_)) => BackendKind::Run,
Ok(Backend::Socket(_)) => BackendKind::Socket,
Err(_) => return None,
};
Some(ServiceInfo {
name: name.clone(),
allow: svc.allow.clone(),
backend,
})
})
.collect()
}
pub(crate) fn peer_infos(store: &PeerStore) -> Vec<PeerInfo> {
store
.list()
.unwrap_or_default()
.into_iter()
.map(|e| PeerInfo {
name: e.petname,
services: e.services,
user_id: e.user_id,
})
.collect()
}
pub(crate) async fn register_service(state: &DaemonState, params: &Value) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let parsed: RegisterParams =
serde_json::from_value(params.clone()).context("register_service params")?;
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let (name, backend, allow) = (parsed.name, parsed.backend, parsed.allow);
let (name_w, backend_w, allow_w) = (name.clone(), backend.clone(), allow.clone());
tokio::task::spawn_blocking(move || {
write_service_to_config(&config_path, &name_w, &backend_w, &allow_w)
})
.await
.context("join config write")??;
let cfg = reload_services_from_disk(mesh, "register").await?;
*state.services.write().expect("services lock not poisoned") = service_infos(&cfg);
tracing::info!(service = %name, "registered/updated service");
Ok(())
}
pub(crate) async fn add_peer(state: &DaemonState, params: &Value) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let petname = params
.get("petname")
.and_then(Value::as_str)
.context("peer_add: missing petname")?
.to_string();
let eid_str = params
.get("endpoint_id")
.and_then(Value::as_str)
.context("peer_add: missing endpoint_id")?;
let endpoint_id = eid_str
.parse::<iroh::EndpointId>()
.map_err(|e| anyhow::anyhow!("peer_add: endpoint_id is not a valid EndpointId: {e}"))?;
let allow = params
.get("allow")
.and_then(Value::as_array)
.map(|a| {
a.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect()
})
.unwrap_or_default();
let entry = PeerEntry {
endpoint_id: *endpoint_id.as_bytes(),
petname: petname.clone(),
services: allow,
paired_at: None,
user_id: None,
};
let store = mesh.store.clone();
tokio::task::spawn_blocking(move || store.add(entry))
.await
.context("join peer add")??;
let store = mesh.store.clone();
let peers = tokio::task::spawn_blocking(move || peer_infos(&store))
.await
.context("join peer list")?;
*state.peers.write().expect("peers lock not poisoned") = peers;
tracing::info!(peer = %petname, "added peer to allowlist");
Ok(())
}
pub async fn remove_peer(state: &DaemonState, params: &Value) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let petname = params
.get("petname")
.and_then(Value::as_str)
.context("peer_remove: missing petname")?
.to_string();
let revoked = revoke_service_access(mesh, &petname).await?;
let store = mesh.store.clone();
let petname_w = petname.clone();
let removed = tokio::task::spawn_blocking(move || store.remove(&petname_w))
.await
.context("join peer remove")??;
tracing::info!(peer = %petname, "unpaired peer");
if revoked || removed {
mesh.audit().record(AuditRecord::trust(
now_ts(),
"unpair".into(),
Some(petname.clone()),
));
}
Ok(())
}
struct RenamePlan {
targets: Vec<PeerEntry>,
old_petnames: std::collections::BTreeSet<String>,
}
fn rename_plan(
store: &PeerStore,
config_path: &Path,
user_id: Option<&str>,
petname: Option<&str>,
to: &str,
) -> Result<Option<RenamePlan>> {
let all = store.list()?;
let targets: Vec<PeerEntry> = all
.iter()
.filter(|e| match user_id {
Some(u) => e.user_id.as_deref() == Some(u),
None => Some(e.petname.as_str()) == petname,
})
.cloned()
.collect();
if targets.is_empty() {
anyhow::bail!("peer_rename: no matching contact");
}
if targets.iter().all(|e| e.petname == to) {
return Ok(None); }
let target_ids: std::collections::BTreeSet<[u8; 32]> =
targets.iter().map(|e| e.endpoint_id).collect();
if all
.iter()
.any(|e| e.petname == to && !target_ids.contains(&e.endpoint_id))
{
anyhow::bail!("the nickname \"{to}\" is already used by another contact");
}
let backed_by_peer = all.iter().any(|e| e.petname == to);
if !backed_by_peer && crate::pairing::rendezvous::petname_in_any_service_allow(config_path, to)?
{
anyhow::bail!("the nickname \"{to}\" is already granted access — pick another");
}
let old_petnames = targets.iter().map(|e| e.petname.clone()).collect();
Ok(Some(RenamePlan {
targets,
old_petnames,
}))
}
pub async fn rename_peer(state: &DaemonState, params: &Value) -> Result<()> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
let to = params
.get("to")
.and_then(Value::as_str)
.unwrap_or("")
.trim()
.to_string();
if to.is_empty() {
anyhow::bail!("peer_rename: the new nickname is empty");
}
let user_id = params
.get("user_id")
.and_then(Value::as_str)
.map(str::to_string);
let petname = params
.get("petname")
.and_then(Value::as_str)
.map(str::to_string);
if user_id.is_none() && petname.is_none() {
anyhow::bail!("peer_rename: no contact identified");
}
let _reload = mesh.reload_lock.lock().await;
let store = mesh.store.clone();
let config_path = mesh.config_path.clone();
let (uid_c, pn_c, to_c) = (user_id.clone(), petname.clone(), to.clone());
let plan = tokio::task::spawn_blocking(move || {
rename_plan(
&store,
&config_path,
uid_c.as_deref(),
pn_c.as_deref(),
&to_c,
)
})
.await
.context("join rename plan")??;
let RenamePlan {
targets,
old_petnames,
} = match plan {
Some(p) => p,
None => return Ok(()), };
let store = mesh.store.clone();
let config_path = mesh.config_path.clone();
let to_c = to.clone();
tokio::task::spawn_blocking(move || {
for old in &old_petnames {
if old != &to_c {
rename_allow_in_config(&config_path, old, &to_c)?;
}
}
for mut e in targets {
e.petname = to_c.clone();
store.add(e)?;
}
anyhow::Ok(())
})
.await
.context("join rename mutate")??;
reload_services_from_disk(mesh, "rename").await?;
tracing::info!(to = %to, "renamed contact");
Ok(())
}
fn short_fingerprint(id: &iroh::EndpointId) -> String {
id.to_string().chars().take(8).collect()
}
fn default_self_petname(id: &iroh::EndpointId) -> String {
hostname_petname().unwrap_or_else(|| short_fingerprint(id))
}
fn hostname_petname() -> Option<String> {
let out = std::process::Command::new("hostname").output().ok()?;
sanitize_hostname(&String::from_utf8_lossy(&out.stdout))
}
fn sanitize_hostname(raw: &str) -> Option<String> {
let short = raw
.trim()
.split('.')
.next()
.unwrap_or("")
.to_ascii_lowercase();
let cleaned: String = short
.chars()
.filter(|c| c.is_ascii_alphanumeric() || *c == '-')
.collect();
(!cleaned.is_empty()).then_some(cleaned)
}
pub(crate) async fn mint_invite(services: Vec<String>, mesh: &MeshState) -> Result<InviteResult> {
use rand::RngCore;
let mut secret = [0u8; 32];
rand::rngs::OsRng.fill_bytes(&mut secret);
let inviter_id = *mesh.endpoint.id().as_bytes();
let _ = tokio::time::timeout(RELAY_READY_TIMEOUT, mesh.endpoint.online()).await;
let inviter_addr_json = serde_json::to_string(&mesh.endpoint.addr())
.context("serialize our own endpoint address for the invite")?;
let now = epoch_now_u64();
let expires_at_epoch = now + INVITE_TTL.as_secs();
let invite = Invite {
secret,
inviter_id,
inviter_addr_json,
petname: mesh.self_petname.clone(),
services: services.clone(),
expires_at_epoch,
};
let invite_line = invite.encode();
mesh.invites.remove_expired(now);
mesh.invites.mint(invite);
tracing::info!(?services, "invite minted");
Ok(InviteResult {
invite_line,
expires_at_epoch,
})
}
pub(crate) async fn redeem(state: &DaemonState, invite_line: String) -> Result<PairResult> {
let mesh = state
.mesh()
.context("daemon has no mesh (control-only mode)")?;
crate::pairing::rendezvous::redeem_invite(
mesh.endpoint.clone(),
mesh.self_petname.clone(),
invite_line,
mesh.store.clone(),
mesh.self_binding(),
)
.await
}
pub async fn grant_service_access(
mesh: &Arc<MeshState>,
redeemer_petname: &str,
services: &[String],
) -> Result<()> {
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let petname = redeemer_petname.to_string();
let services_w = services.to_vec();
let changed = tokio::task::spawn_blocking(move || {
append_allow_to_config(&config_path, &petname, &services_w)
})
.await
.context("join grant config write")??;
if changed {
reload_services_from_disk(mesh, "grant").await?;
}
tracing::info!(peer = %redeemer_petname, ?services, changed, "granted service access");
mesh.audit().record(AuditRecord::trust(
now_ts(),
"pair".into(),
Some(redeemer_petname.to_string()),
));
Ok(())
}
pub(crate) async fn revoke_service_access(mesh: &Arc<MeshState>, petname: &str) -> Result<bool> {
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let petname_w = petname.to_string();
let changed =
tokio::task::spawn_blocking(move || remove_allow_from_config(&config_path, &petname_w))
.await
.context("join revoke config write")??;
if changed {
reload_services_from_disk(mesh, "revoke").await?;
}
tracing::info!(peer = %petname, changed, "revoked service access");
Ok(changed)
}
pub(crate) async fn open_session<CR, CW>(
state: &DaemonState,
peer: &str,
service: &str,
control_reader: FrameReader<CR>,
mut control_writer: CW,
) -> Result<()>
where
CR: AsyncRead + Unpin + Send,
CW: AsyncWrite + Unpin + Send,
{
let Some(mesh) = state.mesh() else {
let _ = write_frame(
&mut control_writer,
&synthesized(Value::Null, ERR_UNREACHABLE, "daemon has no mesh"),
)
.await;
return Ok(());
};
let transport = match dial_service(mesh, peer, service).await {
Ok(t) => t,
Err(e) => {
tracing::warn!(peer, service, %e, "open_session dial failed; answering -32055");
let _ = write_frame(
&mut control_writer,
&synthesized(Value::Null, ERR_UNREACHABLE, "peer unreachable"),
)
.await;
return Ok(());
}
};
pipe_session(transport, service, control_reader, control_writer).await
}
pub fn serving_state(endpoint: iroh::Endpoint, store: Arc<PeerStore>) -> Arc<DaemonState> {
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new(store.clone()));
let self_petname = short_fingerprint(&endpoint.id());
let mesh = MeshState::new(
endpoint,
gate,
store,
Arc::new(LiveInvites::new()),
self_petname,
paths::default_config_path().unwrap_or_default(),
Arc::new(RosterGate::empty()),
Arc::new(ConnRegistry::new()),
None,
None,
None,
None,
);
Arc::new(DaemonState::with_mesh(
STACK_VERSION,
mesh,
Vec::new(),
Vec::new(),
))
}
#[derive(serde::Deserialize)]
struct RegisterParams {
name: String,
backend: BackendSpec,
allow: Vec<String>,
}
fn acquire_singleton_lock(lock_path: &Path) -> Result<Option<File>> {
use rustix::fs::{FlockOperation, flock};
let file = std::fs::OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(lock_path)
.with_context(|| format!("open singleton lock {}", lock_path.display()))?;
match flock(&file, FlockOperation::NonBlockingLockExclusive) {
Ok(()) => Ok(Some(file)),
Err(rustix::io::Errno::WOULDBLOCK) => Ok(None),
Err(e) => Err(anyhow::Error::new(e).context("flock singleton lock")),
}
}
#[cfg(test)]
mod tests {
use super::*;
use ed25519_dalek::SigningKey;
use mcpmesh_trust::roster::encode_b64u;
#[test]
fn sanitize_hostname_makes_a_friendly_petname() {
assert_eq!(sanitize_hostname("jetson\n").as_deref(), Some("jetson"));
assert_eq!(
sanitize_hostname("Johns-MacBook-Pro.local").as_deref(),
Some("johns-macbook-pro"),
"strip the domain, lowercase, keep dashes"
);
assert_eq!(
sanitize_hostname("nvidia jetson!").as_deref(),
Some("nvidiajetson"),
"drop spaces + punctuation"
);
assert_eq!(sanitize_hostname(" ").as_deref(), None);
assert_eq!(sanitize_hostname("").as_deref(), None);
assert_eq!(sanitize_hostname(".local").as_deref(), None);
}
#[tokio::test]
async fn blob_ops_error_without_a_mesh() {
let st = DaemonState::new("test");
assert!(blob_list(&st).await.is_err());
assert!(
blob_publish(&st, "scope".into(), "/tmp/x".into())
.await
.is_err()
);
assert!(blob_grant(&st, "scope".into(), "bob".into()).await.is_err());
assert!(
blob_fetch(&st, "ticket".into(), "/tmp/dst".into())
.await
.is_err()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn status_reads_live_config_and_store_not_the_stale_snapshot() {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
"[services.kb]\nsocket = \"/run/kb.sock\"\nallow = []\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path.clone()).await;
let cfg0 = Config::load(&config_path).unwrap();
let state = crate::control::DaemonState::with_mesh(
"test",
mesh.clone(),
service_infos(&cfg0),
peer_infos(&mesh.store),
);
append_allow_to_config(&config_path, "alice", &["kb".to_string()]).unwrap();
mesh.store
.add(PeerEntry {
endpoint_id: [9u8; 32],
petname: "alice".into(),
services: Vec::new(),
paired_at: None,
user_id: None,
})
.unwrap();
let status = crate::control::status_result(&state);
let kb = status
.services
.iter()
.find(|s| s.name == "kb")
.expect("kb service in status");
assert!(
kb.allow.contains(&"alice".to_string()),
"status must show the live grant, got allow={:?}",
kb.allow
);
assert!(
status.peers.iter().any(|p| p.name == "alice"),
"status must show the live peer, got peers={:?}",
status.peers
);
}
#[tokio::test]
async fn status_surfaces_self_and_peer_user_ids() {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
"[services.kb]\nsocket = \"/run/kb.sock\"\nallow = []\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path.clone()).await;
mesh.set_self_binding(Some(crate::pairing::rendezvous::SelfBinding {
user_pk: "b64u:selfpk".into(),
sig: "b64u:selfsig".into(),
}));
mesh.store
.add(PeerEntry {
endpoint_id: [1u8; 32],
petname: "alice".into(),
services: Vec::new(),
paired_at: Some("1".into()),
user_id: Some("b64u:alicepk".into()),
})
.unwrap();
mesh.store
.add(PeerEntry {
endpoint_id: [2u8; 32],
petname: "legacy".into(),
services: Vec::new(),
paired_at: None,
user_id: None,
})
.unwrap();
let cfg0 = Config::load(&config_path).unwrap();
let state = crate::control::DaemonState::with_mesh(
"test",
mesh.clone(),
service_infos(&cfg0),
peer_infos(&mesh.store),
);
let status = crate::control::status_result(&state);
assert_eq!(
status.self_user_id.as_deref(),
Some("b64u:selfpk"),
"status must surface this daemon's own self-sovereign user_id"
);
let alice = status
.peers
.iter()
.find(|p| p.name == "alice")
.expect("alice in status");
assert_eq!(
alice.user_id.as_deref(),
Some("b64u:alicepk"),
"a paired peer's PROVEN user_id must be surfaced in status"
);
let legacy = status
.peers
.iter()
.find(|p| p.name == "legacy")
.expect("legacy in status");
assert!(
legacy.user_id.is_none(),
"a petname-only peer stays user_id: None"
);
}
#[tokio::test]
async fn recent_pairings_ring_is_bounded_newest_first_and_surfaced_by_status() {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(&config_path, "").unwrap();
let mesh = hermetic_mesh(config_path).await;
for i in 0..10u64 {
mesh.record_pairing(format!("peer{i}"), format!("code-{i}"), i);
}
let recent = mesh.recent_pairings();
assert_eq!(recent.len(), 8, "the ring is capped at 8");
assert_eq!(recent[0].peer_petname, "peer9", "newest first");
assert_eq!(
recent[7].peer_petname, "peer2",
"the two oldest were dropped"
);
let state = crate::control::DaemonState::with_mesh("test", mesh, Vec::new(), Vec::new());
let status = crate::control::status_result(&state);
assert_eq!(status.recent_pairings.len(), 8);
assert_eq!(status.recent_pairings[0].sas_code, "code-9");
assert_eq!(status.recent_pairings[0].paired_at_epoch, 9);
}
#[tokio::test]
async fn rename_peer_renames_all_devices_and_carries_grants() {
use serde_json::json;
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
"[services.kb]\nsocket = \"/run/kb.sock\"\nallow = [\"alice-old\"]\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path.clone()).await;
mesh.store
.add(rename_entry(1, "alice-old", Some("b64u:ALICE")))
.unwrap();
mesh.store
.add(rename_entry(2, "alice-old", Some("b64u:ALICE")))
.unwrap();
let state =
crate::control::DaemonState::with_mesh("test", mesh.clone(), Vec::new(), Vec::new());
rename_peer(&state, &json!({ "user_id": "b64u:ALICE", "to": "Alice" }))
.await
.unwrap();
let names: Vec<String> = mesh
.store
.list()
.unwrap()
.into_iter()
.map(|e| e.petname)
.collect();
assert!(
names.iter().all(|n| n == "Alice"),
"all devices renamed, got {names:?}"
);
let doc: toml::Table =
toml::from_str(&std::fs::read_to_string(&config_path).unwrap()).unwrap();
let allow = doc["services"]["kb"]["allow"].as_array().unwrap();
assert!(allow.iter().any(|v| v.as_str() == Some("Alice")));
assert!(!allow.iter().any(|v| v.as_str() == Some("alice-old")));
}
#[tokio::test]
async fn rename_peer_guards_bad_requests_and_collisions() {
use serde_json::json;
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
"[services.kb]\nsocket = \"/run/kb.sock\"\nallow = []\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path).await;
mesh.store
.add(rename_entry(1, "alice", Some("b64u:ALICE")))
.unwrap();
mesh.store
.add(rename_entry(2, "bob", Some("b64u:BOB")))
.unwrap();
let state =
crate::control::DaemonState::with_mesh("test", mesh.clone(), Vec::new(), Vec::new());
assert!(
rename_peer(&state, &json!({ "user_id": "b64u:ALICE", "to": " " }))
.await
.is_err()
);
assert!(rename_peer(&state, &json!({ "to": "X" })).await.is_err());
assert!(
rename_peer(&state, &json!({ "user_id": "b64u:NOBODY", "to": "X" }))
.await
.is_err()
);
assert!(
rename_peer(&state, &json!({ "user_id": "b64u:ALICE", "to": "bob" }))
.await
.is_err()
);
let names: std::collections::BTreeSet<String> = mesh
.store
.list()
.unwrap()
.into_iter()
.map(|e| e.petname)
.collect();
assert!(
names.contains("alice") && names.contains("bob"),
"no rename should have occurred: {names:?}"
);
}
fn rename_entry(id: u8, petname: &str, user_id: Option<&str>) -> PeerEntry {
PeerEntry {
endpoint_id: [id; 32],
petname: petname.into(),
services: Vec::new(),
paired_at: None,
user_id: user_id.map(str::to_string),
}
}
#[test]
fn rename_plan_groups_by_user_id_and_guards_collisions() {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("config.toml");
std::fs::write(&cfg, "[services.kb]\nallow = [\"orphan\"]\n").unwrap();
let store = PeerStore::open(&dir.path().join("s.redb")).unwrap();
store
.add(rename_entry(1, "bob-phone", Some("b64u:BOB")))
.unwrap();
store
.add(rename_entry(2, "bob-laptop", Some("b64u:BOB")))
.unwrap();
store
.add(rename_entry(3, "carol", Some("b64u:CAROL")))
.unwrap();
let plan = rename_plan(&store, &cfg, Some("b64u:BOB"), None, "Bobby")
.unwrap()
.unwrap();
assert_eq!(plan.targets.len(), 2);
assert_eq!(
plan.old_petnames,
["bob-laptop".to_string(), "bob-phone".to_string()]
.into_iter()
.collect()
);
assert!(rename_plan(&store, &cfg, Some("b64u:BOB"), None, "carol").is_err());
assert!(rename_plan(&store, &cfg, Some("b64u:BOB"), None, "orphan").is_err());
store.add(rename_entry(4, "dave", None)).unwrap();
assert_eq!(
rename_plan(&store, &cfg, None, Some("dave"), "Dave")
.unwrap()
.unwrap()
.targets
.len(),
1
);
assert!(
rename_plan(&store, &cfg, Some("b64u:CAROL"), None, "carol")
.unwrap()
.is_none()
);
assert!(rename_plan(&store, &cfg, Some("b64u:NOBODY"), None, "x").is_err());
}
#[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;
}
async fn hermetic_mesh(config_path: PathBuf) -> Arc<MeshState> {
let dir = config_path.parent().unwrap();
let store = Arc::new(PeerStore::open(&dir.join("state.redb")).unwrap());
let pairs = Arc::new(AllowlistGate::new(store.clone()));
let roster = Arc::new(RosterGate::empty());
let gate: Arc<dyn TrustGate> = Arc::new(ComposedGate::new(roster.clone(), pairs));
let hermetic = crate::config::NetworkCfg {
relay_mode: "disabled".into(),
..Default::default()
};
let endpoint = build_endpoint(iroh::SecretKey::from_bytes(&[7u8; 32]), &hermetic, false)
.await
.unwrap();
MeshState::new(
endpoint,
gate,
store,
Arc::new(LiveInvites::new()),
"test".into(),
config_path,
roster,
Arc::new(ConnRegistry::new()),
None,
None,
None,
None,
)
}
#[tokio::test(flavor = "multi_thread")]
async fn org_join_pins_and_status_reports_pending() {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
"[network]\nrelay_mode = \"disabled\"\n\n[identity]\npetname = \"mydev\"\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path.clone()).await;
let state = DaemonState::with_mesh("0.0.0", mesh.clone(), vec![], vec![]);
let live_status =
|mesh: &Arc<MeshState>| roster_status(mesh, Config::load(&config_path).ok().as_ref());
assert!(
live_status(&mesh).is_none(),
"an unpinned daemon must surface no roster status"
);
assert!(
org_join(
&state,
"acme".into(),
"not-a-real-key".into(),
"alice".into(),
"/home/alice/user.key".into(),
)
.await
.is_err(),
"a garbage org_root_pk must be rejected"
);
assert!(
Config::load(&config_path)
.unwrap()
.identity
.org_root_pk
.is_none(),
"a rejected join must NOT pin a garbage anchor"
);
assert!(
live_status(&mesh).is_none(),
"a rejected join leaves status unpinned"
);
let root = SigningKey::from_bytes(&[9u8; 32]);
let pk_str = encode_b64u(root.verifying_key().as_bytes());
let res = org_join(
&state,
"acme".into(),
pk_str.clone(),
"alice".into(),
"/home/alice/user.key".into(),
)
.await
.expect("valid join pins");
assert_eq!(
res,
OrgJoinResult {
org_id: "acme".into()
}
);
let cfg = Config::load(&config_path).unwrap();
assert_eq!(cfg.identity.org_id.as_deref(), Some("acme"));
assert_eq!(cfg.identity.org_root_pk.as_deref(), Some(pk_str.as_str()));
assert_eq!(cfg.identity.user_id.as_deref(), Some("alice"));
assert_eq!(
cfg.identity.user_key.as_deref(),
Some(Path::new("/home/alice/user.key"))
);
assert_eq!(cfg.identity.petname.as_deref(), Some("mydev"));
let status = live_status(&mesh).expect("pinned org surfaces a pending status");
assert_eq!(status.state, "pending");
assert_eq!(status.serial, 0);
assert_eq!(status.org_id, "acme");
let expected_fp = pairing::sas::fingerprint_words(&root.verifying_key().to_bytes());
assert!(!expected_fp.is_empty());
assert_eq!(status.org_root_fingerprint, expected_fp);
}
#[tokio::test(flavor = "multi_thread")]
async fn url_poll_does_not_panic_when_the_ring_provider_is_installed() {
let _ = rustls::crypto::ring::default_provider().install_default();
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(&config_path, "[network]\nrelay_mode = \"disabled\"\n").unwrap();
let mesh = hermetic_mesh(config_path).await;
let result = crate::roster::distribute::poll_roster_url_once(
&mesh,
"http://127.0.0.1:1/roster.json",
)
.await;
assert!(
result.is_err(),
"an unreachable poll must return Err (not panic, not a silent Ok): {result:?}"
);
}
fn serve_roster_http(
body: Vec<u8>,
) -> (
String,
Arc<std::sync::atomic::AtomicBool>,
std::thread::JoinHandle<()>,
) {
use std::io::{Read, Write};
use std::sync::atomic::{AtomicBool, Ordering};
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind http server");
let port = listener.local_addr().unwrap().port();
listener.set_nonblocking(true).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
let handle = std::thread::spawn(move || {
while !stop_thread.load(Ordering::Relaxed) {
match listener.accept() {
Ok((mut stream, _)) => {
let _ = stream.set_read_timeout(Some(Duration::from_secs(2)));
let mut buf = [0u8; 2048];
let _ = stream.read(&mut buf);
let header = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
);
let _ = stream.write_all(header.as_bytes());
let _ = stream.write_all(&body);
let _ = stream.flush();
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
std::thread::sleep(Duration::from_millis(5));
}
Err(_) => break,
}
}
});
(format!("http://127.0.0.1:{port}/roster.json"), stop, handle)
}
#[tokio::test(flavor = "multi_thread")]
async fn runtime_set_roster_url_bootstraps_the_first_roster_without_a_restart() {
use mcpmesh_trust::roster::mutate;
use mcpmesh_trust::roster::sign::mint_signed;
use std::sync::atomic::Ordering;
let _ = rustls::crypto::ring::default_provider().install_default();
let root = SigningKey::from_bytes(&[9u8; 32]);
let pk_str = encode_b64u(&root.verifying_key().to_bytes());
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
std::fs::write(
&config_path,
format!(
"[network]\nrelay_mode = \"disabled\"\n[identity]\norg_root_pk = \"{pk_str}\"\norg_id = \"acme\"\n"
),
)
.unwrap();
let mesh = hermetic_mesh(config_path).await;
assert!(
mesh.roster.view().is_none(),
"pending joiner: no roster installed yet (D5)"
);
let state = DaemonState::with_mesh("0.0.0", mesh.clone(), vec![], vec![]);
let now = epoch_now_i64();
let roster = mint_signed(
&root,
mutate::empty_roster("acme", 1, now - 3600, now + 86_400),
);
let (url, stop, handle) = serve_roster_http(serde_json::to_vec(&roster).unwrap());
set_roster_url(&state, url.clone())
.await
.expect("set_roster_url pins the URL and (re)starts polling");
let deadline = std::time::Instant::now() + Duration::from_secs(60);
loop {
if mesh.roster.view().map(|v| v.serial()).unwrap_or(0) >= 1 {
break;
}
assert!(
std::time::Instant::now() < deadline,
"the first roster must install at RUNTIME (no daemon restart) after set_roster_url"
);
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert_eq!(
mesh.roster.view().expect("roster installed").serial(),
1,
"the runtime poll bootstrapped the joiner's first roster (D5) — gate hot-swapped to serial 1"
);
set_roster_url(&state, url)
.await
.expect("a repeat set_roster_url is idempotent (abort+replace, no stacking)");
assert_eq!(
mesh.roster.view().unwrap().serial(),
1,
"a repeat set_roster_url does not regress or double-install the roster"
);
stop.store(true, Ordering::Relaxed);
let _ = handle.join();
}
#[test]
fn spawn_concurrency_reads_max_sessions_with_a_safe_floor() {
let c = Config::from_toml_str("[limits]\nmax_sessions = 2\n").unwrap();
assert_eq!(super::spawn_concurrency(&c), 2);
let dflt = Config::from_toml_str("").unwrap();
assert_eq!(super::spawn_concurrency(&dflt), 4, "default max_sessions");
let zero = Config::from_toml_str("[limits]\nmax_sessions = 0\n").unwrap();
assert_eq!(
super::spawn_concurrency(&zero),
1,
"a 0 misconfig floors to 1, never no-permits"
);
assert_eq!(
super::SPAWN_CONCURRENCY as u32,
crate::config::LimitsCfg::default().max_sessions,
"the documented default matches the config default"
);
}
#[test]
fn staleness_sweep_fires_only_when_degraded_stopped() {
use mcpmesh_trust::roster::validate::RosterState;
assert!(super::should_staleness_sever(Some(
RosterState::DegradedStopped
)));
assert!(!super::should_staleness_sever(Some(
RosterState::DegradedGrace
)));
assert!(!super::should_staleness_sever(Some(RosterState::Approved)));
assert!(
!super::should_staleness_sever(None),
"no roster → never sweep"
);
}
}