use std::future::Future;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use mcpmesh_local_api::transport::{LocalListener, LocalStream};
use mcpmesh_local_api::{
API_NAME, API_VERSION, AuditListParams, AuditPruneParams, BlobFetchCancelParams,
BlobFetchParams, BlobGrantParams, BlobPublishParams, BlobRepublishParams, BlobRevokeParams,
BlobUnpublishParams, Hello, InviteParams, OpenSessionParams, OrgJoinParams, PairParams,
PeerServicesParams, RosterInstallParams, ServiceAllowParams, SetAppMetadataParams,
SetNicknameParams, SetRelaysParams, SetRosterUrlParams, StatusResult, UnregisterServiceParams,
method_of,
};
use mcpmesh_net::framing::{FrameReader, Inbound, write_frame};
use serde_json::{Value, json};
use tokio::sync::Notify;
use tokio::task::JoinHandle;
use crate::daemon::MeshState;
use crate::ipc::{self, MAX_FRAME_BYTES};
pub struct DaemonState {
pub stack_version: String,
pub(crate) mesh: Option<Arc<MeshState>>,
shutdown: Notify,
control_tasks: std::sync::Mutex<Vec<JoinHandle<()>>>,
}
impl DaemonState {
pub fn new(stack_version: impl Into<String>) -> Self {
Self {
stack_version: stack_version.into(),
mesh: None,
shutdown: Notify::new(),
control_tasks: std::sync::Mutex::new(Vec::new()),
}
}
pub fn with_mesh(stack_version: impl Into<String>, mesh: Arc<MeshState>) -> Self {
Self {
stack_version: stack_version.into(),
mesh: Some(mesh),
shutdown: Notify::new(),
control_tasks: std::sync::Mutex::new(Vec::new()),
}
}
pub(crate) async fn shutdown_requested(&self) {
self.shutdown.notified().await;
}
pub(crate) fn request_shutdown(&self) {
self.shutdown.notify_one();
}
pub(crate) fn mesh(&self) -> Option<&Arc<MeshState>> {
self.mesh.as_ref()
}
pub(crate) fn mesh_required(&self) -> Result<&Arc<MeshState>> {
self.mesh()
.context("daemon has no mesh (control-only mode)")
}
pub(crate) fn track_control_task(&self, handle: JoinHandle<()>) {
let mut tasks = self
.control_tasks
.lock()
.expect("control_tasks lock not poisoned");
tasks.retain(|h| !h.is_finished());
tasks.push(handle);
}
pub(crate) fn abort_control_tasks(&self) {
let tasks = std::mem::take(
&mut *self
.control_tasks
.lock()
.expect("control_tasks lock not poisoned"),
);
for task in tasks {
task.abort();
}
}
}
pub async fn serve_control(mut listener: LocalListener, state: Arc<DaemonState>) -> Result<()> {
loop {
tokio::select! {
() = state.shutdown.notified() => {
tracing::info!("shutdown requested; control server stopping");
return Ok(());
}
accepted = listener.accept() => {
let stream = match accepted {
Ok(s) => s,
Err(e) => {
tracing::warn!(%e, "control accept failed; backing off");
tokio::time::sleep(Duration::from_millis(50)).await;
continue;
}
};
let conn_state = state.clone();
let handle = tokio::spawn(async move {
if let Err(e) = handle_conn(stream, conn_state).await {
tracing::debug!(%e, "control connection ended");
}
});
state.track_control_task(handle);
}
}
}
}
async fn handle_conn(stream: LocalStream, state: Arc<DaemonState>) -> Result<()> {
if let Err(e) = ipc::check_peer(&stream) {
tracing::warn!(%e, "refused unauthorized control connection");
return Ok(());
}
let (read_half, write_half) = mcpmesh_local_api::transport::split_local(stream);
serve_control_io(read_half, write_half, state).await
}
pub async fn serve_control_io<W>(
read_half: impl tokio::io::AsyncRead + Unpin + Send + 'static,
write_half: W,
state: Arc<DaemonState>,
) -> Result<()>
where
W: tokio::io::AsyncWrite + Unpin + Send + 'static,
{
let hello = Hello {
api: API_NAME.into(),
api_version: API_VERSION.into(),
api_minor: mcpmesh_local_api::API_MINOR,
stack_version: state.stack_version.clone(),
};
let writer = Arc::new(tokio::sync::Mutex::new(write_half));
write_frame(&mut *writer.lock().await, &serde_json::to_value(&hello)?).await?;
let reader = FrameReader::new(tokio::io::BufReader::new(read_half), MAX_FRAME_BYTES);
let ephemeral_registered = Arc::new(std::sync::Mutex::new(Vec::<String>::new()));
let loop_state = state.clone();
let eph = ephemeral_registered.clone();
let outcome: Result<()> = async move {
let mut reader = reader;
let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
let inflight = Arc::new(tokio::sync::Semaphore::new(
mcpmesh_local_api::MAX_INFLIGHT,
));
loop {
while tasks.try_join_next().is_some() {}
match reader.next().await? {
None => return Ok(()), Some(Inbound::Violation(v)) => {
let resp = error(Value::Null, -32700, format!("invalid request frame: {v:?}"));
write_frame(&mut *writer.lock().await, &resp).await?;
}
Some(Inbound::Frame(req)) => {
if method_of(&req) == Some("shutdown") {
loop_state.shutdown.notify_one();
let resp = dispatch(&req, &loop_state);
let _ = write_frame(&mut *writer.lock().await, &resp).await;
loop_state.abort_control_tasks();
return Ok(());
}
if method_of(&req) == Some("open_session") {
let params = req.get("params").cloned().unwrap_or(Value::Null);
let p: OpenSessionParams = match params_of(¶ms) {
Ok(p) => p,
Err(e) => {
let id = req.get("id").cloned().unwrap_or(Value::Null);
let resp = error(id, -32602, format!("open_session failed: {e}"));
write_frame(&mut *writer.lock().await, &resp).await?;
continue;
}
};
let write_half = reclaim_writer(&mut tasks, writer).await?;
return crate::daemon::open_session(
&loop_state,
&p.peer,
&p.service,
reader,
write_half,
)
.await;
}
if method_of(&req) == Some("subscribe") {
let write_half = reclaim_writer(&mut tasks, writer).await?;
return run_subscription(&loop_state, write_half).await;
}
let permit = match inflight.clone().try_acquire_owned() {
Ok(p) => p,
Err(_) => {
let id = req.get("id").cloned().unwrap_or(Value::Null);
let resp = error(
id,
mcpmesh_local_api::ERR_TOO_MANY_INFLIGHT,
format!(
"connection already has {} requests in flight; retry after one completes, or use a second control connection",
mcpmesh_local_api::MAX_INFLIGHT
),
);
write_frame(&mut *writer.lock().await, &resp).await?;
continue;
}
};
let pending_ephemeral = ephemeral_name(&req).map(|name| {
eph.lock()
.expect("ephemeral_registered lock not poisoned")
.push(name.to_string());
name.to_string()
});
let task_state = loop_state.clone();
let task_writer = writer.clone();
let task_eph = eph.clone();
tasks.spawn(async move {
let _permit = permit;
let resp = match n0_future::FutureExt::catch_unwind(
std::panic::AssertUnwindSafe(handle_request(&req, &task_state)),
)
.await
{
Ok(resp) => resp,
Err(_) => {
let id = req.get("id").cloned().unwrap_or(Value::Null);
error(
id,
-32603,
"internal error: the request handler panicked (see the daemon log)",
)
}
};
if let Some(name) = pending_ephemeral
&& resp.get("result").is_none()
{
let mut held = task_eph
.lock()
.expect("ephemeral_registered lock not poisoned");
if let Some(i) = held.iter().rposition(|n| *n == name) {
held.remove(i);
}
}
let _ = write_frame(&mut *task_writer.lock().await, &resp).await;
});
}
}
}
}
.await;
if let Some(mesh) = state.mesh() {
let names = ephemeral_registered
.lock()
.expect("ephemeral_registered lock not poisoned")
.clone();
crate::daemon::unregister_ephemeral(mesh, &names).await;
}
outcome
}
fn ephemeral_name(req: &Value) -> Option<&str> {
if method_of(req) != Some("register_service") {
return None;
}
let params = req.get("params")?;
params
.get("ephemeral")
.and_then(|v| v.as_bool())
.unwrap_or(false)
.then(|| params.get("name").and_then(|v| v.as_str()))
.flatten()
}
async fn reclaim_writer<W>(
tasks: &mut tokio::task::JoinSet<()>,
writer: Arc<tokio::sync::Mutex<W>>,
) -> Result<W> {
while tasks.join_next().await.is_some() {}
Arc::try_unwrap(writer)
.map(tokio::sync::Mutex::into_inner)
.map_err(|_| {
anyhow::anyhow!("control writer still shared after draining in-flight requests")
})
}
async fn run_subscription(
state: &Arc<DaemonState>,
mut w: impl tokio::io::AsyncWrite + Unpin,
) -> Result<()> {
use crate::stream::StreamFrame;
let (audit, mesh) = match state.mesh() {
Some(mesh) => (mesh.audit(), Some(mesh)),
None => (crate::audit::AuditSink::disabled(), None),
};
let rx = audit.subscribe();
let mut reach_rx = mesh.map(|m| m.reach_bcast.subscribe());
let mut self_rx = mesh.map(|m| m.self_net_bcast.subscribe());
let mut blob_rx = mesh.map(|m| m.blob_bcast.subscribe());
let snapshot = StreamFrame::Snapshot {
active_sessions: audit.active_sessions(),
reachability: mesh.map(crate::daemon::reachability_of).unwrap_or_default(),
self_network: mesh.map(|m| {
let stamp = *m
.self_net_change
.lock()
.expect("self_net_change lock not poisoned");
crate::daemon::self_network_now(m, stamp)
}),
};
write_frame(&mut w, &serde_json::to_value(&snapshot)?).await?;
let mut rx = rx;
use tokio::sync::broadcast::error::RecvError;
fn audit_frame(
r: Result<crate::audit::AuditRecord, RecvError>,
closed: &mut bool,
) -> Option<StreamFrame> {
match r {
Ok(record) => Some(StreamFrame::Event {
record: Box::new(record),
}),
Err(RecvError::Lagged(n)) => Some(StreamFrame::Lagged { dropped: n }),
Err(RecvError::Closed) => {
*closed = true;
None
}
}
}
fn blob_frame(
r: Result<crate::daemon::BlobTransfer, RecvError>,
closed: &mut bool,
) -> Option<StreamFrame> {
match r {
Ok(t) => Some(StreamFrame::BlobTransfer {
direction: t.direction,
hash: t.hash,
bytes_done: t.bytes_done,
bytes_total: t.bytes_total,
state: t.state,
peer: t.peer,
}),
Err(RecvError::Lagged(n)) => Some(StreamFrame::Lagged { dropped: n }),
Err(RecvError::Closed) => {
*closed = true;
None
}
}
}
fn reach_frame(
r: Result<crate::daemon::ReachTransition, RecvError>,
closed: &mut bool,
) -> Option<StreamFrame> {
match r {
Ok(t) => Some(StreamFrame::Reachability {
peer: t.peer,
source: t.source,
}),
Err(RecvError::Lagged(n)) => Some(StreamFrame::Lagged { dropped: n }),
Err(RecvError::Closed) => {
*closed = true;
None
}
}
}
fn self_net_frame(
r: Result<mcpmesh_local_api::SelfNetwork, RecvError>,
closed: &mut bool,
) -> Option<StreamFrame> {
match r {
Ok(self_network) => Some(StreamFrame::SelfNetwork { self_network }),
Err(RecvError::Lagged(n)) => Some(StreamFrame::Lagged { dropped: n }),
Err(RecvError::Closed) => {
*closed = true;
None
}
}
}
async fn tap<T: Clone>(
rx: &mut Option<tokio::sync::broadcast::Receiver<T>>,
) -> Result<T, RecvError> {
match rx {
Some(rx) => rx.recv().await,
None => std::future::pending().await,
}
}
let (mut closed_audit, mut closed_reach, mut closed_self) = (false, false, false);
let mut closed_blob = false;
loop {
if rx.is_none() && reach_rx.is_none() && self_rx.is_none() && blob_rx.is_none() {
return Ok(());
}
let frame = tokio::select! {
r = tap(&mut rx) => audit_frame(r, &mut closed_audit),
r = tap(&mut reach_rx) => reach_frame(r, &mut closed_reach),
r = tap(&mut self_rx) => self_net_frame(r, &mut closed_self),
r = tap(&mut blob_rx) => blob_frame(r, &mut closed_blob),
};
if closed_audit {
rx = None;
closed_audit = false;
}
if closed_reach {
reach_rx = None;
closed_reach = false;
}
if closed_blob {
blob_rx = None;
closed_blob = false;
}
if closed_self {
self_rx = None;
closed_self = false;
}
let Some(frame) = frame else { continue };
if write_frame(&mut w, &serde_json::to_value(&frame)?)
.await
.is_err()
{
return Ok(()); }
}
}
pub(crate) async fn handle_request(req: &Value, state: &DaemonState) -> Value {
let id = req.get("id").cloned().unwrap_or(Value::Null);
let params = req.get("params").cloned().unwrap_or(Value::Null);
match method_of(req) {
Some("register_service") => respond(
id,
"register_service",
with_params(¶ms, |p| crate::daemon::register_service(state, p))
.await
.map(unit),
),
Some("peer_add") => respond(
id,
"peer_add",
with_params(¶ms, |p| crate::daemon::add_peer(state, p))
.await
.map(unit),
),
Some("peer_introduce") => respond(
id,
"peer_introduce",
with_params(¶ms, |p| crate::daemon::introduce_peer(state, p))
.await
.map(unit),
),
Some("peer_endorse") => respond(
id,
"peer_endorse",
with_params(¶ms, |p| crate::daemon::endorse_peer(state, p)).await,
),
Some("peer_remove") => respond(
id,
"peer_remove",
with_params(¶ms, |p| crate::daemon::remove_peer(state, p))
.await
.map(unit),
),
Some("peer_rename") => respond(
id,
"peer_rename",
with_params(¶ms, |p| crate::daemon::rename_peer(state, p))
.await
.map(unit),
),
Some("invite") => {
let mesh = match state.mesh_required() {
Ok(mesh) => mesh,
Err(e) => return error(id, -32000, e.to_string()),
};
respond(
id,
"invite",
with_params(¶ms, |p: InviteParams| {
crate::daemon::mint_invite(
p.services,
p.app_label,
p.max_uses,
p.peer_nickname,
p.as_self,
mesh,
)
})
.await,
)
}
Some("pair") => respond(
id,
"pair",
with_params(¶ms, |p: PairParams| {
crate::daemon::redeem(state, p.invite_line, p.as_nickname)
})
.await,
),
Some("roster_install") => respond(
id,
"roster_install",
with_params(¶ms, |p: RosterInstallParams| {
crate::daemon::install_roster(state, p.path, p.org_root_pk)
})
.await,
),
Some("org_join") => respond(
id,
"org_join",
with_params(¶ms, |p: OrgJoinParams| {
crate::daemon::org_join(state, p.org_id, p.org_root_pk, p.user_id, p.user_key)
})
.await,
),
Some("set_roster_url") => respond(
id,
"set_roster_url",
with_params(¶ms, |p: SetRosterUrlParams| {
crate::daemon::set_roster_url(state, p.url)
})
.await
.map(unit),
),
Some("set_nickname") => respond(
id,
"set_nickname",
with_params(¶ms, |p: SetNicknameParams| {
crate::daemon::set_nickname(state, p.nickname)
})
.await
.map(unit),
),
Some("set_relays") => respond(
id,
"set_relays",
with_params(¶ms, |p: SetRelaysParams| {
crate::daemon::set_relays(state, p.relay_urls)
})
.await,
),
Some("peer_services") => {
let p: PeerServicesParams = match params_of(¶ms) {
Ok(p) => p,
Err(e) => return error(id, -32602, format!("peer_services: {e}")),
};
respond(
id,
"peer_services",
crate::daemon::peer_services(state, p.peer).await,
)
}
Some("peer_diagnostics") => {
let p: mcpmesh_local_api::PeerDiagnosticsParams = match params_of(¶ms) {
Ok(p) => p,
Err(e) => return error(id, -32602, format!("peer_diagnostics: {e}")),
};
respond(
id,
"peer_diagnostics",
crate::daemon::peer_diagnostics(state, &p.peer).await,
)
}
Some("unregister_service") => respond(
id,
"unregister_service",
with_params(¶ms, |p: UnregisterServiceParams| {
crate::daemon::unregister_service(state, p.name)
})
.await
.map(unit),
),
Some("service_allow_grant") => respond(
id,
"service_allow_grant",
with_params(¶ms, |p: ServiceAllowParams| {
crate::daemon::service_allow_grant(state, p.service, p.principal)
})
.await
.map(unit),
),
Some("service_allow_revoke") => respond(
id,
"service_allow_revoke",
with_params(¶ms, |p: ServiceAllowParams| {
crate::daemon::service_allow_revoke(state, p.service, p.principal)
})
.await
.map(unit),
),
Some("set_app_metadata") => respond(
id,
"set_app_metadata",
with_params(¶ms, |p: SetAppMetadataParams| {
crate::daemon::set_app_metadata(state, p.metadata)
})
.await
.map(unit),
),
Some("blob_publish") => respond(
id,
"blob_publish",
with_params(¶ms, |p: BlobPublishParams| {
crate::daemon::blob_publish(state, p.scope, p.path)
})
.await,
),
Some("blob_grant") => respond(
id,
"blob_grant",
with_params(¶ms, |p: BlobGrantParams| {
crate::daemon::blob_grant(state, p.scope, p.principal)
})
.await
.map(unit),
),
Some("blob_revoke") => respond(
id,
"blob_revoke",
with_params(¶ms, |p: BlobRevokeParams| {
crate::daemon::blob_revoke(state, p.scope, p.principals)
})
.await
.map(unit),
),
Some("blob_unpublish") => respond(
id,
"blob_unpublish",
with_params(¶ms, |p: BlobUnpublishParams| {
crate::daemon::blob_unpublish(state, p.scope, p.hash)
})
.await
.map(unit),
),
Some("blob_republish") => respond(
id,
"blob_republish",
with_params(¶ms, |p: BlobRepublishParams| {
crate::daemon::blob_republish(state, p.scope, p.hash)
})
.await,
),
Some("blob_list") => respond(
id,
"blob_list",
with_params(¶ms, |p: mcpmesh_local_api::BlobListParams| {
crate::daemon::blob_list(state, p)
})
.await,
),
Some("blob_fetch") => respond(
id,
"blob_fetch",
with_params(¶ms, |p: BlobFetchParams| {
crate::daemon::blob_fetch(state, p.ticket, p.dest_path)
})
.await,
),
Some("blob_fetch_cancel") => respond(
id,
"blob_fetch_cancel",
with_params(¶ms, |p: BlobFetchCancelParams| async move {
crate::daemon::blob_fetch_cancel(state, &p.hash)
})
.await,
),
Some("audit_summary") => {
let sink_dir = state
.mesh
.as_ref()
.and_then(|m| m.audit().dir().map(std::path::Path::to_path_buf));
match tokio::task::spawn_blocking(move || {
let dir = match sink_dir {
Some(d) => d,
None => mcpmesh_trust::paths::default_audit_dir()?,
};
crate::audit::read_all_records(&dir)
.map(|recs| crate::audit::summarize_sessions(&recs))
})
.await
{
Ok(r) => respond(id, "audit_summary", r.map_err(anyhow::Error::from)),
Err(e) => error(id, -32000, format!("audit_summary task failed: {e}")),
}
}
Some("audit_prune") => {
let sink_dir = state
.mesh
.as_ref()
.and_then(|m| m.audit().dir().map(std::path::Path::to_path_buf));
let r = with_params(¶ms, move |p: AuditPruneParams| async move {
anyhow::ensure!(
crate::audit::valid_month_key(&p.before),
"before must be a zero-padded YYYY-MM month key, got '{}'",
p.before
);
let dir = sink_dir.context(
"audit_prune requires a daemon with a live audit writer — refusing to \
guess a directory for a destructive operation",
)?;
tokio::task::spawn_blocking(move || crate::audit::prune_before(&dir, &p.before))
.await
.map_err(|e| anyhow::anyhow!("audit_prune task failed: {e}"))?
.map(|deleted_months| mcpmesh_local_api::AuditPruneResult { deleted_months })
.map_err(anyhow::Error::from)
})
.await;
respond(id, "audit_prune", r)
}
Some("audit_list") => {
let sink_dir = state
.mesh
.as_ref()
.and_then(|m| m.audit().dir().map(std::path::Path::to_path_buf));
let r = with_params(¶ms, move |p: AuditListParams| async move {
for (name, bound) in [("since", &p.since), ("until", &p.until)] {
if let Some(m) = bound {
anyhow::ensure!(
crate::audit::valid_month_key(m),
"{name} must be a zero-padded YYYY-MM month key, got '{m}'"
);
}
}
let kind = match &p.kind {
Some(s) => Some(crate::audit::parse_kind(s).ok_or_else(|| {
anyhow::anyhow!(
"unknown kind '{s}' — one of session_open, session_close, request, \
blob_fetch, trust"
)
})?),
None => None,
};
let limit = p.limit.unwrap_or(500).min(1000) as usize;
let offset = p.offset.unwrap_or(0) as usize;
tokio::task::spawn_blocking(move || {
let dir = match sink_dir {
Some(d) => d,
None => mcpmesh_trust::paths::default_audit_dir()?,
};
crate::audit::list_page(
&dir,
p.since.as_deref(),
p.until.as_deref(),
kind,
p.peer.as_deref(),
limit,
offset,
)
})
.await
.map_err(|e| anyhow::anyhow!("audit_list task failed: {e}"))?
.map_err(anyhow::Error::from)
})
.await;
respond(id, "audit_list", r)
}
#[cfg(test)]
Some("__test_panic") => panic!("deliberate test panic"),
#[cfg(test)]
Some("__test_block") => {
let gate = params.get("gate").and_then(|v| v.as_u64()).unwrap_or(0);
tests::test_block(gate).await;
ok(id, json!({}))
}
_ => dispatch(req, state),
}
}
#[derive(Debug)]
pub(crate) struct InvalidParams(pub(crate) String);
impl std::fmt::Display for InvalidParams {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
impl std::error::Error for InvalidParams {}
fn respond<T: serde::Serialize>(id: Value, method: &str, r: anyhow::Result<T>) -> Value {
match r {
Ok(v) => ok(
id,
serde_json::to_value(v).expect("control result serializes"),
),
Err(e) if e.downcast_ref::<InvalidParams>().is_some() => {
error(id, -32602, format!("{method} failed: {e}"))
}
Err(e) if e.downcast_ref::<crate::daemon::BlobWithdrawn>().is_some() => error(
id,
mcpmesh_local_api::ERR_BLOB_WITHDRAWN,
format!("{method} failed: {e}"),
),
Err(e) if e.downcast_ref::<crate::daemon::NoSuchBlob>().is_some() => error(
id,
mcpmesh_local_api::ERR_NO_SUCH_BLOB,
format!("{method} failed: {e}"),
),
Err(e) if e.downcast_ref::<crate::daemon::Cancelled>().is_some() => error(
id,
mcpmesh_local_api::ERR_CANCELLED,
format!("{method} failed: {e}"),
),
Err(e)
if e.downcast_ref::<crate::pairing::rendezvous::PairRefusal>()
.is_some() =>
{
let refusal = e
.downcast_ref::<crate::pairing::rendezvous::PairRefusal>()
.expect("checked by the guard");
error(id, refusal.code(), format!("{method} failed: {e}"))
}
Err(e)
if e.downcast_ref::<crate::pairing::rendezvous::NicknameTaken>()
.is_some() =>
{
error(
id,
mcpmesh_local_api::ERR_NICKNAME_TAKEN,
format!("{method} failed: {e}"),
)
}
Err(e)
if e.downcast_ref::<crate::daemon::NoSuchService>().is_some()
|| e.downcast_ref::<crate::daemon::NoSuchBlobScope>().is_some() =>
{
error(
id,
mcpmesh_local_api::ERR_NO_SUCH_SERVICE,
format!("{method} failed: {e}"),
)
}
Err(e) => error(id, -32000, format!("{method} failed: {e}")),
}
}
fn unit((): ()) -> Value {
json!({})
}
fn dispatch(req: &Value, state: &DaemonState) -> Value {
let id = req.get("id").cloned().unwrap_or(Value::Null);
match method_of(req) {
Some("status") => respond(id, "status", status_result(state)),
Some("shutdown") => ok(id, json!({})),
Some(other) => error(id, -32601, format!("unknown method: {other}")),
None => error(id, -32600, "request is missing a `method`"),
}
}
pub(crate) fn status_result(state: &DaemonState) -> Result<StatusResult> {
let (services, peers, roster) = match state.mesh() {
Some(mesh) => {
let cfg = crate::config::Config::load(&mesh.config_path).map_err(|e| {
anyhow::anyhow!("config unreadable at {}: {e}", mesh.config_path.display())
})?;
let roster = crate::daemon::roster_status(mesh, Some(&cfg));
{
let entries = mesh.store.list().unwrap_or_default();
(
crate::daemon::service_infos(&mesh.live_services(), &entries),
crate::daemon::peer_infos(&mesh.store),
roster,
)
}
}
None => (Vec::new(), Vec::new(), None),
};
let presence = state
.mesh()
.map(crate::daemon::presence_peers)
.unwrap_or_default();
let self_user_id = state
.mesh()
.and_then(|mesh| mesh.self_binding())
.map(|binding| binding.user_pk);
let recent_pairings = state
.mesh()
.map(|mesh| mesh.recent_pairings())
.unwrap_or_default();
let reachability = state
.mesh()
.map(crate::daemon::reachability_of)
.unwrap_or_default();
let storage = state.mesh().map(|mesh| {
let audit_bytes = mesh
.audit()
.dir()
.and_then(|d| crate::audit::list_month_files(d).ok())
.map(|files| files.iter().map(|(_, _, size)| size).sum())
.unwrap_or(0);
let redb_bytes = std::fs::metadata(mesh.store.path())
.map(|m| m.len())
.unwrap_or(0);
let blobs_bytes = mesh
.blobs_dir()
.map(crate::util::dir_size_bytes)
.unwrap_or(0);
mcpmesh_local_api::StorageInfo {
audit_bytes,
redb_bytes,
blobs_bytes,
}
});
let self_network = state.mesh().map(|mesh| {
let stamp = *mesh
.self_net_change
.lock()
.expect("self_net_change lock not poisoned");
crate::daemon::self_network_now(mesh, stamp)
});
Ok(StatusResult {
stack_version: state.stack_version.clone(),
services,
peers,
roster,
presence,
self_user_id,
recent_pairings,
reachability,
self_nickname: state
.mesh()
.map(|mesh| mesh.self_nickname())
.unwrap_or_default(),
storage,
self_network,
})
}
fn params_of<T: serde::de::DeserializeOwned>(params: &Value) -> anyhow::Result<T> {
let v = match params {
Value::Null => json!({}),
p => p.clone(),
};
serde_json::from_value(v)
.map_err(|e| anyhow::Error::new(InvalidParams(format!("invalid params: {e}"))))
}
async fn with_params<P, R, F>(params: &Value, f: impl FnOnce(P) -> F) -> anyhow::Result<R>
where
P: serde::de::DeserializeOwned,
F: Future<Output = anyhow::Result<R>>,
{
f(params_of(params)?).await
}
fn ok(id: Value, result: Value) -> Value {
json!({ "jsonrpc": "2.0", "id": id, "result": result })
}
fn error(id: Value, code: i64, message: impl Into<String>) -> Value {
json!({ "jsonrpc": "2.0", "id": id, "error": { "code": code, "message": message.into() } })
}
#[cfg(test)]
mod tests {
use super::*;
fn control_only() -> Arc<DaemonState> {
Arc::new(DaemonState::new("0.1.0-test"))
}
fn req(method: &str, params: Value) -> Value {
json!({ "jsonrpc": "2.0", "id": 1, "method": method, "params": params })
}
#[test]
fn each_onboarding_refusal_carries_its_own_code() {
use crate::pairing::rendezvous::PairRefusal;
let coded = |code: i64| -> Value {
let e: anyhow::Result<()> = Err(anyhow::Error::new(PairRefusal::new(code, "why")));
respond(json!(1), "pair", e)
};
for code in [
mcpmesh_local_api::ERR_INVITE_EXPIRED,
mcpmesh_local_api::ERR_INVITE_NOT_LIVE,
mcpmesh_local_api::ERR_INVITER_UNREACHABLE,
mcpmesh_local_api::ERR_INVITER_MISMATCH,
mcpmesh_local_api::ERR_INVITE_NAME_CONFLICT,
mcpmesh_local_api::ERR_INVITE_REFUSED,
] {
let v = coded(code);
assert_eq!(
v["error"]["code"], code,
"each condition keeps its own code: {v}"
);
}
let all = [
mcpmesh_local_api::ERR_INVITE_EXPIRED,
mcpmesh_local_api::ERR_INVITE_NOT_LIVE,
mcpmesh_local_api::ERR_INVITER_UNREACHABLE,
mcpmesh_local_api::ERR_INVITER_MISMATCH,
mcpmesh_local_api::ERR_INVITE_NAME_CONFLICT,
mcpmesh_local_api::ERR_INVITE_REFUSED,
mcpmesh_local_api::ERR_NICKNAME_TAKEN,
];
let unique: std::collections::BTreeSet<i64> = all.iter().copied().collect();
assert_eq!(unique.len(), all.len(), "codes must not collide: {all:?}");
let recoverable = [
mcpmesh_local_api::ERR_INVITE_EXPIRED,
mcpmesh_local_api::ERR_INVITE_NOT_LIVE,
mcpmesh_local_api::ERR_INVITER_UNREACHABLE,
mcpmesh_local_api::ERR_INVITE_NAME_CONFLICT,
mcpmesh_local_api::ERR_INVITE_REFUSED,
mcpmesh_local_api::ERR_NICKNAME_TAKEN,
];
assert!(
!recoverable.contains(&mcpmesh_local_api::ERR_INVITER_MISMATCH),
"the address-swap refusal must not share a code with any recoverable one — an app \
that renders every pairing failure as a friendly retry would paper over exactly the \
attack that check exists to catch"
);
let plain: anyhow::Result<()> = Err(anyhow::anyhow!("something else"));
assert_eq!(respond(json!(1), "pair", plain)["error"]["code"], -32000);
}
#[test]
fn a_nickname_collision_answers_err_nickname_taken() {
let reason = "pairing refused: nickname 'studio-mac' is already taken by another paired \
peer; the invite was NOT consumed — rename this node and redeem the same \
invite again";
let typed: anyhow::Result<()> = Err(anyhow::Error::new(
crate::pairing::rendezvous::NicknameTaken(reason.into()),
));
let v = respond(json!(1), "pair", typed);
assert_eq!(
v["error"]["code"],
mcpmesh_local_api::ERR_NICKNAME_TAKEN,
"the collision must be branchable, not -32000: {v}"
);
assert_eq!(
v["error"]["code"], -32043,
"the value is the wire contract: {v}"
);
let msg = v["error"]["message"].as_str().unwrap_or_default();
assert!(
msg.contains("rename this node") && !msg.contains("set_nickname"),
"the message carries the inviter's reworded prose verbatim: {v}"
);
let generic: anyhow::Result<()> = Err(anyhow::anyhow!("pairing refused: pairing refused"));
let v = respond(json!(1), "pair", generic);
assert_eq!(v["error"]["code"], -32000, "got {v}");
}
#[tokio::test]
async fn serve_control_io_speaks_the_protocol_over_a_duplex() {
let state = control_only();
let (client_io, server_io) = tokio::io::duplex(64 * 1024);
let (sr, sw) = tokio::io::split(server_io);
tokio::spawn(serve_control_io(sr, sw, state));
let (cr, cw) = tokio::io::split(client_io);
let mut client = mcpmesh_local_api::connect_control_io(cr, cw)
.await
.expect("hello handshake");
assert_eq!(client.hello().stack_version, "0.1.0-test");
let status = client.status().await.expect("status");
assert_eq!(status.stack_version, "0.1.0-test");
assert!(status.services.is_empty());
}
#[test]
fn dispatch_status_answers_empty_lists_without_a_mesh() {
let st = control_only();
let r = dispatch(&req("status", json!({})), &st);
assert_eq!(r["result"]["stack_version"], "0.1.0-test");
assert!(r["result"]["peers"].as_array().unwrap().is_empty());
assert!(r["result"]["services"].as_array().unwrap().is_empty());
assert!(r["result"]["roster"].is_null());
}
#[test]
fn dispatch_status_tolerates_any_params_shape() {
let st = control_only();
for p in [json!({}), Value::Null, json!({ "junk": true })] {
assert!(dispatch(&req("status", p), &st).get("result").is_some());
}
let omitted = json!({ "jsonrpc": "2.0", "id": 1, "method": "status" });
assert!(dispatch(&omitted, &st).get("result").is_some());
}
#[test]
fn dispatch_shutdown_acks_and_unknown_methods_error() {
let st = control_only();
assert_eq!(
dispatch(&req("shutdown", json!({})), &st)["result"],
json!({})
);
assert_eq!(
dispatch(&req("frobnicate", json!({})), &st)["error"]["code"],
-32601
);
let no_method = json!({ "jsonrpc": "2.0", "id": 1 });
assert_eq!(dispatch(&no_method, &st)["error"]["code"], -32600);
}
#[tokio::test]
async fn mesh_methods_error_gracefully_without_a_mesh() {
let st = control_only();
for method in [
"register_service",
"peer_add",
"peer_introduce",
"peer_endorse",
"peer_remove",
"peer_rename",
"invite",
"pair",
"roster_install",
"org_join",
"set_roster_url",
"blob_publish",
"blob_grant",
"blob_list",
"blob_fetch",
] {
let r = handle_request(&req(method, json!({})), &st).await;
let code = r["error"]["code"].as_i64();
assert!(
matches!(code, Some(-32000) | Some(-32602)),
"method {method} should error gracefully in control-only mode, got {r}"
);
assert!(
r.get("result").is_none(),
"method {method} must not succeed: {r}"
);
}
}
#[tokio::test]
async fn malformed_params_answer_an_invalid_params_error() {
let st = control_only();
let r = handle_request(&req("peer_remove", json!({ "nickname": 42 })), &st).await;
assert_eq!(r["error"]["code"], -32602);
assert!(
r["error"]["message"]
.as_str()
.unwrap()
.contains("invalid params"),
"message names the params problem: {r}"
);
let r = handle_request(&req("peer_rename", json!({ "user_id": "u" })), &st).await;
assert_eq!(r["error"]["code"], -32602);
let r = handle_request(
&req("peer_remove", json!({ "nickname": "a", "extra": true })),
&st,
)
.await;
assert_eq!(
r["error"]["code"], -32602,
"unknown params field is rejected: {r}"
);
}
#[tokio::test]
async fn audit_summary_works_in_control_only_mode() {
let st = control_only();
let r = handle_request(&req("audit_summary", json!({})), &st).await;
assert!(
r.get("result").is_some(),
"audit_summary should succeed: {r}"
);
}
#[tokio::test]
async fn handle_request_delegates_status_to_dispatch() {
let st = control_only();
let r = handle_request(&req("status", json!({})), &st).await;
assert_eq!(r["result"]["stack_version"], "0.1.0-test");
}
pub(super) struct Gate {
released: crate::cancel::CancelToken,
entered: std::sync::atomic::AtomicUsize,
completed: std::sync::atomic::AtomicUsize,
}
static GATES: std::sync::LazyLock<
std::sync::Mutex<std::collections::HashMap<u64, std::sync::Arc<Gate>>>,
> = std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
static NEXT_GATE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
fn new_gate() -> (u64, std::sync::Arc<Gate>) {
let id = NEXT_GATE.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let gate = std::sync::Arc::new(Gate {
released: crate::cancel::CancelToken::new(),
entered: std::sync::atomic::AtomicUsize::new(0),
completed: std::sync::atomic::AtomicUsize::new(0),
});
GATES.lock().unwrap().insert(id, gate.clone());
(id, gate)
}
pub(super) async fn test_block(id: u64) {
let gate = GATES
.lock()
.unwrap()
.get(&id)
.cloned()
.expect("__test_block names a registered gate");
gate.entered
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
gate.released.cancelled().await;
gate.completed
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
fn load(c: &std::sync::atomic::AtomicUsize) -> usize {
c.load(std::sync::atomic::Ordering::SeqCst)
}
async fn until(what: &str, mut cond: impl FnMut() -> bool) {
let deadline = tokio::time::Instant::now() + Duration::from_secs(20);
while !cond() {
assert!(
tokio::time::Instant::now() < deadline,
"timed out waiting for: {what}"
);
tokio::time::sleep(Duration::from_millis(2)).await;
}
}
type ClientReader =
FrameReader<tokio::io::BufReader<tokio::io::ReadHalf<tokio::io::DuplexStream>>>;
type ClientWriter = tokio::io::WriteHalf<tokio::io::DuplexStream>;
async fn raw_conn(
state: Arc<DaemonState>,
) -> (
ClientReader,
ClientWriter,
tokio::task::JoinHandle<Result<()>>,
) {
let (client_io, server_io) = tokio::io::duplex(256 * 1024);
let (sr, sw) = tokio::io::split(server_io);
let server = tokio::spawn(serve_control_io(sr, sw, state));
let (cr, cw) = tokio::io::split(client_io);
let mut reader = FrameReader::new(tokio::io::BufReader::new(cr), MAX_FRAME_BYTES);
let hello = next_frame(&mut reader).await;
assert_eq!(hello["api"], API_NAME, "server speaks first with a Hello");
(reader, cw, server)
}
async fn next_frame(r: &mut ClientReader) -> Value {
let read = tokio::time::timeout(Duration::from_secs(10), r.next())
.await
.expect("a frame should arrive within 10s");
match read.expect("read a frame") {
Some(Inbound::Frame(v)) => v,
other => panic!("expected a frame, got {other:?}"),
}
}
async fn send(w: &mut ClientWriter, v: &Value) {
write_frame(w, v).await.expect("write a request frame");
}
fn blocking_req(id: u64, gate: u64) -> Value {
json!({ "id": id, "method": "__test_block", "params": { "gate": gate } })
}
#[tokio::test(flavor = "multi_thread")]
async fn a_slow_request_does_not_stall_the_ones_behind_it() {
let (gate_id, gate) = new_gate();
let (mut r, mut w, _server) = raw_conn(control_only()).await;
send(&mut w, &blocking_req(1, gate_id)).await;
send(&mut w, &json!({ "id": 2, "method": "status" })).await;
let first = next_frame(&mut r).await;
assert_eq!(
first["id"], 2,
"the fast request must answer first: {first}"
);
assert!(
first.get("result").is_some(),
"status should succeed: {first}"
);
assert_eq!(
load(&gate.completed),
0,
"the slow request is still running"
);
gate.released.cancel();
let second = next_frame(&mut r).await;
assert_eq!(
second["id"], 1,
"the released request answers second: {second}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn over_the_inflight_cap_a_request_is_refused_not_queued() {
let (gate_id, gate) = new_gate();
let (mut r, mut w, _server) = raw_conn(control_only()).await;
let cap = mcpmesh_local_api::MAX_INFLIGHT as u64;
for id in 1..=cap {
send(&mut w, &blocking_req(id, gate_id)).await;
}
send(&mut w, &blocking_req(cap + 1, gate_id)).await;
let refusal = next_frame(&mut r).await;
assert_eq!(
refusal["id"],
cap + 1,
"the overflow request is the one refused"
);
assert_eq!(
refusal["error"]["code"],
mcpmesh_local_api::ERR_TOO_MANY_INFLIGHT,
"over the cap must be branchable, not -32000: {refusal}"
);
gate.released.cancel();
let mut answered = std::collections::HashSet::new();
for _ in 0..cap {
let f = next_frame(&mut r).await;
answered.insert(f["id"].as_u64().expect("id is a number"));
}
assert_eq!(
answered,
(1..=cap).collect::<std::collections::HashSet<_>>(),
"every accepted request answers"
);
send(&mut w, &json!({ "id": 999, "method": "status" })).await;
let after = next_frame(&mut r).await;
assert_eq!(after["id"], 999);
assert!(
after.get("result").is_some(),
"usable after the cap clears: {after}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn an_upgrade_waits_for_in_flight_responses_before_it_takes_the_writer() {
let (gate_id, gate) = new_gate();
let (mut r, mut w, _server) = raw_conn(control_only()).await;
send(&mut w, &blocking_req(1, gate_id)).await;
send(&mut w, &json!({ "method": "subscribe" })).await;
tokio::time::sleep(Duration::from_millis(50)).await;
gate.released.cancel();
let first = next_frame(&mut r).await;
assert_eq!(
first["id"], 1,
"the in-flight response lands first: {first}"
);
let snapshot = next_frame(&mut r).await;
assert!(
snapshot.get("id").is_none() && snapshot["type"] == "snapshot",
"the subscription snapshot follows it: {snapshot}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn open_session_drains_in_flight_requests_before_it_takes_the_writer() {
let (gate_id, gate) = new_gate();
let (mut r, mut w, _server) = raw_conn(control_only()).await;
send(&mut w, &blocking_req(7, gate_id)).await;
send(
&mut w,
&json!({ "id": 8, "method": "open_session", "params": { "peer": "bob", "service": "kb" } }),
)
.await;
tokio::time::sleep(Duration::from_millis(50)).await;
gate.released.cancel();
let first = next_frame(&mut r).await;
assert_eq!(
first["id"], 7,
"the in-flight response lands first: {first}"
);
let second = next_frame(&mut r).await;
assert!(
second.get("error").is_some(),
"a mesh-less daemon answers open_session with an error frame: {second}"
);
}
fn ephemeral_register(id: u64, name: &str) -> Value {
json!({
"id": id,
"method": "register_service",
"params": {
"name": name,
"backend": { "socket": { "path": "/run/nowhere.sock" } },
"allow": [],
"ephemeral": true,
}
})
}
async fn mesh_state(
bulk: usize,
) -> (
tempfile::TempDir,
Arc<crate::daemon::MeshState>,
Arc<DaemonState>,
) {
let dir = tempfile::tempdir().unwrap();
let config_path = dir.path().join("config.toml");
let mut cfg = String::new();
for i in 0..bulk {
cfg.push_str(&format!(
"[services.pad{i}]\nsocket = \"/run/pad{i}.sock\"\nallow = []\n"
));
}
std::fs::write(&config_path, cfg).unwrap();
let mesh = crate::daemon::testutil::hermetic_mesh(config_path).await;
let state = Arc::new(DaemonState::with_mesh("test", mesh.clone()));
(dir, mesh, state)
}
fn ephemeral_names(mesh: &Arc<crate::daemon::MeshState>) -> Vec<String> {
let mut v: Vec<String> = mesh
.ephemeral_services
.lock()
.expect("ephemeral_services lock not poisoned")
.keys()
.cloned()
.collect();
v.sort();
v
}
#[tokio::test(flavor = "multi_thread")]
async fn an_ephemeral_registration_is_torn_down_even_if_the_connection_closes_mid_register() {
let (_dir, mesh, state) = mesh_state(300).await;
let (r, mut w, server) = raw_conn(state).await;
send(&mut w, &ephemeral_register(1, "leaky")).await;
until("the ephemeral registration to appear", || {
!ephemeral_names(&mesh).is_empty()
})
.await;
drop(w);
drop(r);
server
.await
.expect("connection task not panicked")
.expect("connection ends cleanly");
assert!(
ephemeral_names(&mesh).is_empty(),
"a registration must not outlive the connection that made it, however it ended: {:?}",
ephemeral_names(&mesh)
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_refused_register_does_not_tear_down_a_name_another_connection_holds() {
let (_dir, mesh, state) = mesh_state(0).await;
let (mut r1, mut w1, _holder) = raw_conn(state.clone()).await;
send(&mut w1, &ephemeral_register(1, "contested")).await;
let ack = next_frame(&mut r1).await;
assert!(
ack.get("result").is_some(),
"the first register wins: {ack}"
);
assert_eq!(ephemeral_names(&mesh), vec!["contested".to_string()]);
let (mut r2, mut w2, loser) = raw_conn(state).await;
let mut bad = ephemeral_register(2, "contested");
bad["params"]["rate_limit_per_min"] = json!(0);
send(&mut w2, &bad).await;
let refused = next_frame(&mut r2).await;
assert!(
refused.get("error").is_some(),
"the second register loses: {refused}"
);
drop(w2);
drop(r2);
loser
.await
.expect("connection task not panicked")
.expect("connection ends cleanly");
assert_eq!(
ephemeral_names(&mesh),
vec!["contested".to_string()],
"the holder's service must survive the loser closing"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_panicking_handler_answers_instead_of_hanging_the_caller() {
let (mut r, mut w, _server) = raw_conn(control_only()).await;
send(&mut w, &json!({ "id": 5, "method": "__test_panic" })).await;
let f = next_frame(&mut r).await;
assert_eq!(f["id"], 5, "the panicking request still answers: {f}");
assert_eq!(f["error"]["code"], -32603, "as an internal error: {f}");
send(&mut w, &json!({ "id": 6, "method": "status" })).await;
let after = next_frame(&mut r).await;
assert_eq!(after["id"], 6);
assert!(after.get("result").is_some(), "still usable: {after}");
}
#[test]
fn a_cancelled_request_answers_err_cancelled() {
let r = respond::<()>(
json!(1),
"blob_fetch",
Err(crate::daemon::Cancelled("blake3:beef".into()).into()),
);
assert_eq!(
r["error"]["code"],
mcpmesh_local_api::ERR_CANCELLED,
"a cancel is branchable, not a generic failure: {r}"
);
assert!(
r["error"]["message"]
.as_str()
.expect("a message")
.contains("blake3:beef"),
"and it names what was cancelled: {r}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn closing_the_connection_aborts_its_in_flight_requests() {
let (gate_id, gate) = new_gate();
let (r, mut w, server) = raw_conn(control_only()).await;
send(&mut w, &blocking_req(1, gate_id)).await;
until("the handler to start", || load(&gate.entered) > 0).await;
drop(w);
drop(r);
server
.await
.expect("connection task not panicked")
.expect("connection ends cleanly on client close");
gate.released.cancel();
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(load(&gate.entered), 1, "the handler did start");
assert_eq!(
load(&gate.completed),
0,
"an in-flight request must be aborted when its connection closes"
);
}
}