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::{OrgJoinResult, RosterInstallResult};
use mcpmesh_trust::roster::validate::RosterState;
use tokio::task::JoinHandle;
use crate::audit::{AuditRecord, now_ts};
use crate::config::Config;
use crate::control::DaemonState;
use crate::roster::RosterStore;
use crate::roster::gate::RosterGate;
use crate::util::{TempPathGuard, blocking, epoch_now_i64, unique_temp_path};
use super::MeshState;
use super::config_write::{
write_identity_pin, write_identity_user_id, write_join_pin, write_roster_url,
};
pub(crate) 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 fn install_roster_view_and_sever(
mesh: &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<mcpmesh_net::EndpointId> =
view.revoked_endpoints().map(|b| (*b).into()).collect();
let active_devices: HashSet<mcpmesh_net::EndpointId> =
view.device_endpoints().map(|b| (*b).into()).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
}
pub(crate) async fn converge_roster_bytes(
mesh: &MeshState,
bytes: &[u8],
serial: u64,
channel: &'static str,
) -> Result<bool> {
let _reload = mesh.reload_lock.lock().await;
if serial <= mesh.roster.view().map(|v| v.serial()).unwrap_or(0) {
return Ok(false);
}
let pk = crate::roster::parse_org_root_pk(
&mesh_config_org_root_pk(mesh)?.context("no pinned org root; cannot accept roster")?,
)?;
let tmp = write_temp_roster(bytes)?;
let rstore = RosterStore::new(installed_roster_path(mesh));
let now = epoch_now_i64();
let view = blocking("join roster install", move || {
rstore.install_from_file(tmp.path(), &pk, now)
})
.await??;
mesh.confirm_roster_current(now).await;
let severed = install_roster_view_and_sever(mesh, view);
reconcile_user_id_from_roster(mesh).await;
drop(_reload);
tracing::info!(serial, severed, channel, "installed roster");
Ok(true)
}
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: &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(),
)
}
pub(crate) 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_required()?;
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(installed_roster_path(mesh));
let now = epoch_now_i64();
let file = PathBuf::from(path);
let view = blocking("join roster install", move || {
rstore.install_from_file(&file, &pk, now)
})
.await??;
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());
blocking("join org-root pin config write", move || {
write_identity_pin(&config_path, &pk_w, &oid_w)
})
.await??;
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) fn mesh_config_org_root_pk(mesh: &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: &MeshState) -> PathBuf {
mesh.config_path
.parent()
.map(|dir| dir.join("roster.json"))
.unwrap_or_else(|| PathBuf::from("roster.json")) }
pub(crate) 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: &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 blocking("join reconcile config user_id write", 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)) | Err(e) => tracing::warn!(%e, "reconcile config user_id write failed"),
}
}
pub(crate) fn write_temp_roster(bytes: &[u8]) -> Result<TempPathGuard> {
let tmp = TempPathGuard::new(unique_temp_path(&std::env::temp_dir(), "mcpmesh-roster-in"));
let mut f = File::create(tmp.path())
.with_context(|| format!("create temp roster {}", tmp.path().display()))?;
f.write_all(bytes).context("write temp roster")?;
f.sync_all().context("sync temp roster")?;
Ok(tmp)
}
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_required()?;
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);
blocking("join org-root pin config write", move || {
write_join_pin(&config_path, &oid, &pk, &uid, &uk)
})
.await??;
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_required()?;
let _reload = mesh.reload_lock.lock().await;
let config_path = mesh.config_path.clone();
let url_w = url.clone();
blocking("roster-url config write", move || {
write_roster_url(&config_path, &url_w)
})
.await??;
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(())
}
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,
));
}
#[cfg(test)]
mod tests {
use super::*;
use crate::daemon::roster_status;
use crate::daemon::testutil::hermetic_mesh;
use crate::pairing;
use ed25519_dalek::SigningKey;
use mcpmesh_trust::roster::encode_b64u;
#[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]\nnickname = \"mydev\"\n",
)
.unwrap();
let mesh = hermetic_mesh(config_path.clone()).await;
let state = DaemonState::with_mesh("0.0.0", mesh.clone());
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.nickname.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[roster]\npoll_interval = \"2s\"\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());
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(180);
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 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"
);
}
}