use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use mcpmesh_local_api::{
BlobFetchResult, BlobPublishResult, BlobScopeList, InviteResult, PairResult, PeerAddParams,
PeerRemoveParams, PeerRenameParams, RegisterServiceParams, ScopeInfo,
};
use mcpmesh_net::errors::{ERR_UNREACHABLE, synthesized};
use mcpmesh_net::framing::{FrameReader, write_frame};
use serde_json::Value;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::allowlist::{PeerEntry, PeerStore};
use crate::audit::{AuditRecord, now_ts};
use crate::config::Config;
use crate::control::DaemonState;
use crate::pairing::Invite;
use crate::util::{blocking, epoch_now_u64};
use super::accept::reload_accept_loop;
use super::config_write::{
append_allow_to_config, remove_allow_from_config, rename_allow_in_config,
write_service_to_config,
};
use super::status::service_infos;
use super::{MeshState, dial_service, pipe_session};
const INVITE_TTL: Duration = Duration::from_secs(24 * 60 * 60);
const RELAY_READY_TIMEOUT: Duration = Duration::from_secs(3);
pub(crate) async fn blob_publish(
state: &DaemonState,
scope: String,
path: String,
) -> Result<BlobPublishResult> {
let mesh = state.mesh_required()?;
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_required()?;
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_required()?;
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_required()?;
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,
})
}
async fn reload_services_from_disk(mesh: &Arc<MeshState>, why: &str) -> Result<()> {
let cfg = Config::load(&mesh.config_path)
.map_err(|e| anyhow::anyhow!("reload config after {why}: {e}"))?;
let ephemeral = mesh
.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned")
.clone();
reload_accept_loop(
mesh,
crate::daemon::build_services_with_ephemeral(
&cfg,
&mesh.audit(),
&mesh.limits(),
&ephemeral,
),
)
.await;
Ok(())
}
pub(crate) async fn register_service(
state: &DaemonState,
params: RegisterServiceParams,
) -> Result<()> {
let mesh = state.mesh_required()?;
let _reload = mesh.reload_lock.lock().await;
let RegisterServiceParams {
name,
backend,
allow,
ephemeral,
} = params;
if ephemeral {
let cfg = Config::load(&mesh.config_path)
.map_err(|e| anyhow::anyhow!("config error in {}: {e}", mesh.config_path.display()))?;
if cfg.services.contains_key(&name) {
anyhow::bail!(
"service '{name}' is already registered persistently in config; \
use a different name for an ephemeral registration"
);
}
mesh.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned")
.insert(
name.clone(),
crate::daemon::EphemeralService {
backend,
allow: allow.clone(),
},
);
reload_services_from_disk(mesh, "register-ephemeral").await?;
tracing::info!(service = %name, "registered ephemeral service");
return Ok(());
}
let config_path = mesh.config_path.clone();
let (name_w, backend_w, allow_w) = (name.clone(), backend.clone(), allow.clone());
blocking("join config write", move || {
write_service_to_config(&config_path, &name_w, &backend_w, &allow_w)
})
.await??;
reload_services_from_disk(mesh, "register").await?;
tracing::info!(service = %name, "registered/updated service");
Ok(())
}
pub(crate) async fn unregister_ephemeral(mesh: &Arc<MeshState>, names: &[String]) {
if names.is_empty() {
return;
}
let _reload = mesh.reload_lock.lock().await;
{
let mut map = mesh
.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned");
for name in names {
map.remove(name);
}
}
if let Err(e) = reload_services_from_disk(mesh, "unregister-ephemeral").await {
tracing::warn!(%e, "reload after ephemeral unregister failed");
}
}
pub(crate) async fn add_peer(state: &DaemonState, params: PeerAddParams) -> Result<()> {
let mesh = state.mesh_required()?;
let PeerAddParams {
nickname,
endpoint_id,
allow,
} = params;
let endpoint_id = endpoint_id
.parse::<iroh::EndpointId>()
.map_err(|e| anyhow::anyhow!("peer_add: endpoint_id is not a valid EndpointId: {e}"))?;
let entry = PeerEntry {
endpoint_id: *endpoint_id.as_bytes(),
nickname: nickname.clone(),
services: allow,
paired_at: None,
user_id: None,
last_addr: None,
};
let store = mesh.store.clone();
blocking("join peer add", move || store.add(entry)).await??;
tracing::info!(peer = %nickname, "added peer to allowlist");
Ok(())
}
pub async fn remove_peer(state: &DaemonState, params: PeerRemoveParams) -> Result<()> {
let mesh = state.mesh_required()?;
let nickname = params.nickname;
let revoked = revoke_service_access(mesh, &nickname).await?;
let store = mesh.store.clone();
let nickname_w = nickname.clone();
let removed = blocking("join peer remove", move || store.remove(&nickname_w)).await??;
if !revoked && !removed {
anyhow::bail!("no paired peer named '{nickname}' — 'mcpmesh status' lists your peers");
}
tracing::info!(peer = %nickname, "unpaired peer");
mesh.audit().record(AuditRecord::trust(
now_ts(),
"unpair".into(),
Some(nickname.clone()),
));
Ok(())
}
struct RenamePlan {
targets: Vec<PeerEntry>,
old_nicknames: std::collections::BTreeSet<String>,
}
fn rename_plan(
store: &PeerStore,
config_path: &Path,
user_id: Option<&str>,
nickname: 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.nickname.as_str()) == nickname,
})
.cloned()
.collect();
if targets.is_empty() {
anyhow::bail!("peer_rename: no matching contact");
}
if targets.iter().all(|e| e.nickname == 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.nickname == 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.nickname == to);
if !backed_by_peer
&& crate::pairing::rendezvous::nickname_in_any_service_allow(config_path, to)?
{
anyhow::bail!("the nickname \"{to}\" is already granted access — pick another");
}
let old_nicknames = targets.iter().map(|e| e.nickname.clone()).collect();
Ok(Some(RenamePlan {
targets,
old_nicknames,
}))
}
pub async fn rename_peer(state: &DaemonState, params: PeerRenameParams) -> Result<()> {
let mesh = state.mesh_required()?;
let to = params.to.trim().to_string();
if to.is_empty() {
anyhow::bail!("peer_rename: the new nickname is empty");
}
let PeerRenameParams {
user_id, nickname, ..
} = params;
if user_id.is_none() && nickname.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(), nickname.clone(), to.clone());
let plan = blocking("join rename plan", move || {
rename_plan(
&store,
&config_path,
uid_c.as_deref(),
pn_c.as_deref(),
&to_c,
)
})
.await??;
let RenamePlan {
targets,
old_nicknames,
} = 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();
blocking("join rename mutate", move || {
for old in &old_nicknames {
if old != &to_c {
rename_allow_in_config(&config_path, old, &to_c)?;
}
}
for mut e in targets {
e.nickname = to_c.clone();
store.add(e)?;
}
anyhow::Ok(())
})
.await??;
reload_services_from_disk(mesh, "rename").await?;
tracing::info!(to = %to, "renamed contact");
Ok(())
}
fn unregistered_service_error(requested: &[String], served: &[String]) -> Option<String> {
let unknown: Vec<&String> = requested.iter().filter(|r| !served.contains(r)).collect();
let quoted: Vec<String> = unknown.iter().map(|n| format!("'{n}'")).collect();
let named = match quoted.as_slice() {
[] => return None,
[one] => format!("no service named {one}"),
many => format!("no services named {}", many.join(", ")),
};
Some(if served.is_empty() {
format!(
"{named} — nothing is served yet; register one with \
'mcpmesh serve <name> -- <command>'"
)
} else {
format!(
"{named} — you serve: {} (see 'mcpmesh status')",
served.join(", ")
)
})
}
pub(crate) async fn mint_invite(
services: Vec<String>,
app_label: Option<String>,
mesh: &MeshState,
) -> Result<InviteResult> {
use rand::RngCore;
if let Some(label) = &app_label
&& label.len() > crate::pairing::MAX_APP_LABEL_LEN
{
anyhow::bail!(
"app_label is {} bytes; the maximum is {}",
label.len(),
crate::pairing::MAX_APP_LABEL_LEN
);
}
if services.is_empty() {
anyhow::bail!(
"invite must name at least one registered service (an invite granting nothing is useless)"
);
}
let cfg = Config::load(&mesh.config_path)
.map_err(|e| anyhow::anyhow!("config error in {}: {e}", mesh.config_path.display()))?;
let ephemeral = mesh
.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned")
.clone();
let served: Vec<String> = service_infos(&cfg, &ephemeral)
.into_iter()
.map(|s| s.name)
.collect();
if let Some(msg) = unregistered_service_error(&services, &served) {
anyhow::bail!(msg);
}
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,
nickname: mesh.self_nickname(),
services: services.clone(),
expires_at_epoch,
app_label,
};
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_required()?;
crate::pairing::rendezvous::redeem_invite(
mesh.endpoint.clone(),
mesh.self_nickname(),
invite_line,
mesh.store.clone(),
mesh.self_binding(),
&mesh.config_path,
)
.await
}
pub async fn grant_service_access(
mesh: &Arc<MeshState>,
redeemer_nickname: &str,
services: &[String],
) -> Result<()> {
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let nickname = redeemer_nickname.to_string();
let services_w = services.to_vec();
let changed = blocking("join grant config write", move || {
append_allow_to_config(&config_path, &nickname, &services_w)
})
.await??;
if changed {
reload_services_from_disk(mesh, "grant").await?;
}
tracing::info!(peer = %redeemer_nickname, ?services, changed, "granted service access");
mesh.audit().record(AuditRecord::trust(
now_ts(),
"pair".into(),
Some(redeemer_nickname.to_string()),
));
Ok(())
}
pub(crate) async fn revoke_service_access(mesh: &Arc<MeshState>, nickname: &str) -> Result<bool> {
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let nickname_w = nickname.to_string();
let changed = blocking("join revoke config write", move || {
remove_allow_from_config(&config_path, &nickname_w)
})
.await??;
if changed {
reload_services_from_disk(mesh, "revoke").await?;
}
tracing::info!(peer = %nickname, 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) => {
mesh.audit().record(
AuditRecord::session_open(now_ts(), Some(peer.to_string()), service.to_string())
.with_status("error"),
);
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
}
#[cfg(test)]
mod tests {
use super::*;
use crate::daemon::testutil::hermetic_mesh;
#[test]
fn unregistered_service_error_message_shapes() {
let s = |names: &[&str]| -> Vec<String> { names.iter().map(|n| n.to_string()).collect() };
assert_eq!(
unregistered_service_error(&s(&["notes"]), &s(&["notes", "kb"])),
None
);
assert_eq!(
unregistered_service_error(&s(&["nosuchsvc"]), &s(&["notes", "code"])).unwrap(),
"no service named 'nosuchsvc' — you serve: notes, code (see 'mcpmesh status')"
);
assert_eq!(
unregistered_service_error(&s(&["a", "notes", "b"]), &s(&["notes"])).unwrap(),
"no services named 'a', 'b' — you serve: notes (see 'mcpmesh status')"
);
assert_eq!(
unregistered_service_error(&s(&["nosuchsvc"]), &[]).unwrap(),
"no service named 'nosuchsvc' — nothing is served yet; register one with \
'mcpmesh serve <name> -- <command>'"
);
}
#[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()
);
}
fn rename_params(user_id: Option<&str>, to: &str) -> PeerRenameParams {
PeerRenameParams {
user_id: user_id.map(str::to_string),
nickname: None,
to: to.into(),
}
}
#[tokio::test]
async fn rename_peer_renames_all_devices_and_carries_grants() {
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());
rename_peer(&state, rename_params(Some("b64u:ALICE"), "Alice"))
.await
.unwrap();
let names: Vec<String> = mesh
.store
.list()
.unwrap()
.into_iter()
.map(|e| e.nickname)
.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() {
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());
assert!(
rename_peer(&state, rename_params(Some("b64u:ALICE"), " "))
.await
.is_err()
);
assert!(rename_peer(&state, rename_params(None, "X")).await.is_err());
assert!(
rename_peer(&state, rename_params(Some("b64u:NOBODY"), "X"))
.await
.is_err()
);
assert!(
rename_peer(&state, rename_params(Some("b64u:ALICE"), "bob"))
.await
.is_err()
);
let names: std::collections::BTreeSet<String> = mesh
.store
.list()
.unwrap()
.into_iter()
.map(|e| e.nickname)
.collect();
assert!(
names.contains("alice") && names.contains("bob"),
"no rename should have occurred: {names:?}"
);
}
fn rename_entry(id: u8, nickname: &str, user_id: Option<&str>) -> PeerEntry {
PeerEntry {
endpoint_id: [id; 32],
nickname: nickname.into(),
services: Vec::new(),
paired_at: None,
user_id: user_id.map(str::to_string),
last_addr: None,
}
}
#[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_nicknames,
["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());
}
}