#![cfg(unix)]
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use assert_cmd::cargo::cargo_bin;
use ed25519_dalek::SigningKey;
use mcpmesh::client::connect_control;
use mcpmesh::config::Config;
use mcpmesh_trust::roster::sign::mint_signed;
use mcpmesh_trust::roster::{Roster, encode_b64u, mutate};
use serde_json::json;
fn launch_in(dir: &Path) -> (PathBuf, PathBuf, PathBuf, Vec<(OsString, OsString)>) {
let runtime = dir.join("runtime");
let config = dir.join("config");
let data = dir.join("data");
let config_mcpmesh = config.join("mcpmesh");
std::fs::create_dir_all(&config_mcpmesh).unwrap();
std::fs::write(
config_mcpmesh.join("config.toml"),
"[network]\nrelay_mode = \"disabled\"\n",
)
.unwrap();
let socket = runtime.join("mcpmesh").join("mcpmesh.sock");
let env = vec![
(OsString::from("XDG_RUNTIME_DIR"), runtime.into_os_string()),
(OsString::from("XDG_CONFIG_HOME"), config.into_os_string()),
(OsString::from("XDG_DATA_HOME"), data.into_os_string()),
];
(cargo_bin("mcpmesh"), socket, config_mcpmesh, env)
}
fn run_cmd(exe: &Path, env: &[(OsString, OsString)], args: &[&str]) -> std::process::Output {
let mut cmd = std::process::Command::new(exe);
cmd.args(args);
for (k, v) in env {
cmd.env(k, v);
}
cmd.output().expect("run mcpmesh subcommand")
}
async fn shutdown_daemon(socket: &Path) {
if let Ok(mut client) = connect_control(socket).await {
let _ = client.request_value(&json!({ "method": "shutdown" })).await;
}
let deadline = Instant::now() + Duration::from_secs(5);
while connect_control(socket).await.is_ok() {
assert!(
Instant::now() < deadline,
"daemon still accepting connections after shutdown"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
fn staged_temps(config: &Path) -> Vec<PathBuf> {
std::fs::read_dir(config)
.into_iter()
.flatten()
.filter_map(Result::ok)
.map(|e| e.path())
.filter(|p| {
p.file_name()
.map(|n| n.to_string_lossy().starts_with("roster.staging."))
.unwrap_or(false)
})
.collect()
}
fn now_epoch_i64() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
fn write_roster(dir: &Path, name: &str, roster: &Roster) -> PathBuf {
let path = dir.join(name);
std::fs::write(&path, serde_json::to_vec(roster).unwrap()).unwrap();
path
}
fn fingerprint_after(label: &str, out: &str) -> String {
out.lines()
.find(|l| l.contains(label))
.and_then(|l| l.split_whitespace().find(|w| w.contains('-')))
.unwrap_or("")
.to_string()
}
fn org_invite_from(out: &std::process::Output) -> String {
String::from_utf8_lossy(&out.stdout)
.lines()
.find_map(|l| l.split_whitespace().find(|w| w.starts_with("mcpmesh-org:")))
.expect("an org invite code")
.to_string()
}
fn join_code_from(out: &std::process::Output) -> String {
String::from_utf8_lossy(&out.stdout)
.lines()
.find_map(|l| {
l.split_whitespace()
.find(|w| w.starts_with("mcpmesh-join:"))
})
.expect("a join code")
.to_string()
}
fn device_code_from(out: &std::process::Output) -> String {
String::from_utf8_lossy(&out.stdout)
.lines()
.find_map(|l| {
l.split_whitespace()
.find(|w| w.starts_with("mcpmesh-device:"))
})
.expect("a device code")
.to_string()
}
fn tamper_join_binding(code: &str) -> String {
let mut jc = mcpmesh::roster::enroll::JoinCode::decode(code).expect("decode join code");
let mut sig = mcpmesh_trust::roster::decode_b64u(&jc.binding_sig).expect("decode binding sig");
sig[0] ^= 0x01; jc.binding_sig = encode_b64u(&sig);
jc.encode()
}
#[tokio::test(flavor = "multi_thread")]
async fn org_create_mints_root_signs_empty_roster_and_prints_the_invite() {
tokio::time::timeout(Duration::from_secs(45), async {
let dir = tempfile::tempdir().unwrap();
let (exe, socket, config, env) = launch_in(dir.path());
let out = run_cmd(&exe, &env, &["org", "create", "acme", "--expires", "30d"]);
let stdout = String::from_utf8_lossy(&out.stdout);
assert!(
out.status.success(),
"org create failed: {}",
String::from_utf8_lossy(&out.stderr)
);
assert!(
stdout.contains("mcpmesh-org:"),
"prints the org invite code: {stdout}"
);
assert!(
stdout.contains("fingerprint"),
"prints the root fingerprint: {stdout}"
);
assert!(
!stdout.contains("b64u:") && !stdout.contains("org-root.key"),
"surface leak: {stdout}"
);
assert!(
stdout.contains("Created org 'acme' (roster serial 1)."),
"prints the confirmation: {stdout}"
);
let key = config.join("org-root.key");
assert!(key.exists(), "org root key minted");
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
assert_eq!(
std::fs::metadata(&key).unwrap().permissions().mode() & 0o777,
0o600
);
}
let roster: Roster =
serde_json::from_slice(&std::fs::read(config.join("roster.json")).unwrap()).unwrap();
assert_eq!(roster.serial, 1);
assert_eq!(roster.org_id, "acme");
assert!(roster.users.is_empty());
let cfg = Config::load(&config.join("config.toml")).expect("reload config");
assert_eq!(cfg.identity.org_id.as_deref(), Some("acme"));
assert!(
cfg.identity
.org_root_pk
.as_deref()
.is_some_and(|pk| pk.starts_with("b64u:")),
"org_root_pk pinned in config"
);
assert!(
staged_temps(&config).is_empty(),
"staged roster temp must be cleaned up on success, found {:?}",
staged_temps(&config)
);
let again = run_cmd(&exe, &env, &["org", "create", "acme2"]);
assert!(
!again.status.success(),
"second org create must refuse: {}",
String::from_utf8_lossy(&again.stdout)
);
shutdown_daemon(&socket).await;
})
.await
.expect("org create test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn org_create_cleans_up_the_staged_temp_when_the_install_is_rejected() {
tokio::time::timeout(Duration::from_secs(45), async {
let dir = tempfile::tempdir().unwrap();
let (exe, socket, config, env) = launch_in(dir.path());
let root_a = SigningKey::from_bytes(&[9u8; 32]);
let pk_a = encode_b64u(&root_a.verifying_key().to_bytes());
let now = now_epoch_i64();
let roster5 = mint_signed(
&root_a,
mutate::empty_roster("acme", 5, now - 3600, now + 86_400),
);
let f5 = write_roster(dir.path(), "roster-5.json", &roster5);
let seed = run_cmd(
&exe,
&env,
&[
"internal",
"roster",
"install",
f5.to_str().unwrap(),
"--org-root-pk",
&pk_a,
],
);
assert!(
seed.status.success(),
"seed install (serial 5) must succeed: {}",
String::from_utf8_lossy(&seed.stderr)
);
let out = run_cmd(&exe, &env, &["org", "create", "acme"]);
assert!(
!out.status.success(),
"org create must fail when its serial-1 install is a rollback; stdout: {}",
String::from_utf8_lossy(&out.stdout)
);
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(
stderr.contains("roster failed validation"),
"expected a daemon install rejection, got: {stderr}"
);
assert!(
staged_temps(&config).is_empty(),
"staged roster temp must be cleaned up on the install-error path, found {:?}",
staged_temps(&config)
);
shutdown_daemon(&socket).await;
})
.await
.expect("org create rejection test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn join_mints_user_key_pins_org_root_and_emits_a_join_code() {
tokio::time::timeout(Duration::from_secs(45), async {
let opdir = tempfile::tempdir().unwrap();
let (opexe, opsock, _opcfg, openv) = launch_in(opdir.path());
let create = run_cmd(&opexe, &openv, &["org", "create", "acme"]);
assert!(
create.status.success(),
"org create failed: {}",
String::from_utf8_lossy(&create.stderr)
);
let create_out = String::from_utf8_lossy(&create.stdout);
let invite = create_out
.lines()
.find_map(|l| l.split_whitespace().find(|w| w.starts_with("mcpmesh-org:")))
.expect("an org invite code")
.to_string();
let op_fp = fingerprint_after("Org root fingerprint:", &create_out);
assert!(!op_fp.is_empty(), "operator prints an org-root fingerprint");
let jdir = tempfile::tempdir().unwrap();
let (jexe, jsock, jcfg, jenv) = launch_in(jdir.path());
let join = run_cmd(&jexe, &jenv, &["join", &invite, "--name", "Alice Nguyen"]);
let jout = String::from_utf8_lossy(&join.stdout);
assert!(
join.status.success(),
"join failed: {}",
String::from_utf8_lossy(&join.stderr)
);
assert!(jout.contains("mcpmesh-join:"), "prints a join code: {jout}");
assert_eq!(
fingerprint_after("Org root fingerprint:", &jout),
op_fp,
"org-root fingerprint must match the operator's"
);
assert!(
jout.contains("Join code fingerprint:"),
"prints the join-code fingerprint: {jout}"
);
assert!(!fingerprint_after("Join code fingerprint:", &jout).is_empty());
assert!(
jout.contains("out-of-band"),
"prints the ceremony instruction"
);
let ukey = jcfg.join("user.key");
assert!(ukey.exists(), "user key minted");
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
assert_eq!(
std::fs::metadata(&ukey).unwrap().permissions().mode() & 0o777,
0o600
);
}
let cfg_txt = std::fs::read_to_string(jcfg.join("config.toml")).unwrap();
assert!(cfg_txt.contains("org_id = \"acme\""));
assert!(cfg_txt.contains("org_root_pk"));
assert!(cfg_txt.contains("user_id = \"alice-nguyen\""));
assert!(cfg_txt.contains("user_key"));
assert!(
!jout.contains("b64u:") && !jout.contains("user.key"),
"surface leak: {jout}"
);
let status = run_cmd(&jexe, &jenv, &["status"]);
assert!(
String::from_utf8_lossy(&status.stdout).contains(&op_fp),
"status shows the pinned org-root fingerprint"
);
shutdown_daemon(&opsock).await;
shutdown_daemon(&jsock).await;
})
.await
.expect("join test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn org_create_roster_url_lands_in_the_invite_and_config_and_join_stores_it() {
tokio::time::timeout(Duration::from_secs(45), async {
const URL: &str = "https://intranet.acme.com/roster.json";
let opdir = tempfile::tempdir().unwrap();
let (opexe, opsock, opcfg, openv) = launch_in(opdir.path());
let create = run_cmd(
&opexe,
&openv,
&["org", "create", "acme", "--roster-url", URL],
);
assert!(
create.status.success(),
"org create --roster-url failed: {}",
String::from_utf8_lossy(&create.stderr)
);
let invite = org_invite_from(&create);
let decoded =
mcpmesh::roster::enroll::OrgInviteCode::decode(&invite).expect("decode org invite");
assert_eq!(
decoded.roster_url.as_deref(),
Some(URL),
"the invite must carry the roster URL"
);
let opcfg_url = Config::load(&opcfg.join("config.toml"))
.expect("reload operator config")
.roster
.url;
assert_eq!(
opcfg_url.as_deref(),
Some(URL),
"operator config [roster].url must be pinned"
);
let jdir = tempfile::tempdir().unwrap();
let (jexe, jsock, jcfg, jenv) = launch_in(jdir.path());
let join = run_cmd(&jexe, &jenv, &["join", &invite, "--name", "Alice"]);
assert!(
join.status.success(),
"join failed: {}",
String::from_utf8_lossy(&join.stderr)
);
let jcfg_url = Config::load(&jcfg.join("config.toml"))
.expect("reload joiner config")
.roster
.url;
assert_eq!(
jcfg_url.as_deref(),
Some(URL),
"joiner config [roster].url must be stored from the invite (D5)"
);
let jcfg_txt = std::fs::read_to_string(jcfg.join("config.toml")).unwrap();
assert!(
jcfg_txt.contains("org_id = \"acme\""),
"org root still pinned"
);
shutdown_daemon(&opsock).await;
shutdown_daemon(&jsock).await;
})
.await
.expect("org create roster-url test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn approve_verifies_the_binding_and_adds_the_member_to_the_roster() {
tokio::time::timeout(Duration::from_secs(45), async {
let opdir = tempfile::tempdir().unwrap();
let (opexe, opsock, opcfg, openv) = launch_in(opdir.path());
let create = run_cmd(&opexe, &openv, &["org", "create", "acme"]);
assert!(
create.status.success(),
"org create failed: {}",
String::from_utf8_lossy(&create.stderr)
);
let invite = org_invite_from(&create);
let org_root_pk = {
let ic =
mcpmesh::roster::enroll::OrgInviteCode::decode(&invite).expect("decode invite");
let bytes = mcpmesh_trust::roster::decode_endpoint_id(&ic.org_root_pk)
.expect("decode org_root_pk");
ed25519_dalek::VerifyingKey::from_bytes(&bytes).expect("org root pk is a valid key")
};
let jdir = tempfile::tempdir().unwrap();
let (jexe, jsock, _jcfg, jenv) = launch_in(jdir.path());
let join_out = run_cmd(
&jexe,
&jenv,
&["join", &invite, "--name", "Alice", "--user-id", "alice"],
);
assert!(
join_out.status.success(),
"join failed: {}",
String::from_utf8_lossy(&join_out.stderr)
);
let join_code = join_code_from(&join_out);
let joiner_code_fp = fingerprint_after(
"Join code fingerprint:",
&String::from_utf8_lossy(&join_out.stdout),
);
assert!(
!joiner_code_fp.is_empty(),
"joiner prints a join-code fingerprint"
);
let jc = mcpmesh::roster::enroll::JoinCode::decode(&join_code).expect("decode join code");
let approve = run_cmd(
&opexe,
&openv,
&["org", "approve", &join_code, "--groups", "team-eng,all"],
);
let aout = String::from_utf8_lossy(&approve.stdout);
assert!(
approve.status.success(),
"approve failed: {}",
String::from_utf8_lossy(&approve.stderr)
);
assert!(
aout.contains("alice") && aout.contains("team-eng") && aout.contains("serial 2"),
"approve confirmation: {aout}"
);
assert!(
aout.contains("Approving join code"),
"approve surfaces the join-code fingerprint: {aout}"
);
assert_eq!(
fingerprint_after("Approving join code", &aout),
joiner_code_fp,
"operator + joiner must see the SAME join-code fingerprint"
);
let roster: Roster =
serde_json::from_slice(&std::fs::read(opcfg.join("roster.json")).unwrap()).unwrap();
assert_eq!(roster.serial, 2);
let alice = roster
.users
.iter()
.find(|u| u.user_id == "alice")
.expect("alice enrolled");
assert_eq!(alice.display_name, "Alice");
assert!(alice.groups.contains(&"team-eng".to_string()));
assert_eq!(alice.devices.len(), 1);
assert_eq!(alice.devices[0].endpoint_id, jc.device_endpoint_id);
assert!(
roster.groups.contains(&"team-eng".to_string()),
"team-eng declared top-level: {:?}",
roster.groups
);
mcpmesh_trust::roster::sign::verify(&roster, &org_root_pk)
.expect("re-signed roster verifies against the org root");
let forged = tamper_join_binding(&join_code);
let bad = run_cmd(
&opexe,
&openv,
&["org", "approve", &forged, "--groups", "team-eng"],
);
assert!(!bad.status.success(), "a forged binding must be refused");
assert!(
String::from_utf8_lossy(&bad.stderr).contains("device binding"),
"the refusal names the device-binding check: {}",
String::from_utf8_lossy(&bad.stderr)
);
let roster_after: Roster =
serde_json::from_slice(&std::fs::read(opcfg.join("roster.json")).unwrap()).unwrap();
assert_eq!(
roster_after.serial, 2,
"a forged approve must not bump serial"
);
assert!(
staged_temps(&opcfg).is_empty(),
"staged roster temp leak: {:?}",
staged_temps(&opcfg)
);
shutdown_daemon(&opsock).await;
shutdown_daemon(&jsock).await;
})
.await
.expect("approve test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn revoke_person_device_and_user_key_grammar() {
tokio::time::timeout(Duration::from_secs(45), async {
let opdir = tempfile::tempdir().unwrap();
let (opexe, opsock, opcfg, openv) = launch_in(opdir.path());
let invite = org_invite_from(&run_cmd(&opexe, &openv, &["org", "create", "acme"]));
let jdir = tempfile::tempdir().unwrap();
let (jexe, jsock, _jcfg, jenv) = launch_in(jdir.path());
let jc = join_code_from(&run_cmd(
&jexe,
&jenv,
&["join", &invite, "--name", "Alice", "--user-id", "alice"],
));
run_cmd(
&opexe,
&openv,
&["org", "approve", &jc, "--groups", "team-eng"],
);
let dev = run_cmd(&opexe, &openv, &["org", "revoke", "alice/laptop"]);
assert!(
dev.status.success(),
"device revoke: {}",
String::from_utf8_lossy(&dev.stderr)
);
let r3: Roster =
serde_json::from_slice(&std::fs::read(opcfg.join("roster.json")).unwrap()).unwrap();
assert_eq!(r3.serial, 3);
assert_eq!(r3.revoked_endpoints.len(), 1);
assert!(
r3.users
.iter()
.find(|u| u.user_id == "alice")
.unwrap()
.devices
.is_empty()
);
assert!(
!run_cmd(&opexe, &openv, &["org", "revoke", "ghost"])
.status
.success()
);
let rot = run_cmd(&opexe, &openv, &["org", "revoke", "alice", "--user-key"]);
assert!(rot.status.success());
assert!(String::from_utf8_lossy(&rot.stdout).contains("re-enroll"));
let r4: Roster =
serde_json::from_slice(&std::fs::read(opcfg.join("roster.json")).unwrap()).unwrap();
assert_eq!(r4.serial, 4);
assert!(r4.users.iter().all(|u| u.user_id != "alice"));
assert_eq!(r4.revoked_endpoints.len(), 1);
shutdown_daemon(&opsock).await;
shutdown_daemon(&jsock).await;
})
.await
.expect("revoke test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn devices_add_signs_a_new_device_binding_and_the_operator_appends_it() {
tokio::time::timeout(Duration::from_secs(45), async {
let opdir = tempfile::tempdir().unwrap();
let (opexe, opsock, opcfg, openv) = launch_in(opdir.path());
let invite = org_invite_from(&run_cmd(&opexe, &openv, &["org", "create", "acme"]));
let d1 = tempfile::tempdir().unwrap();
let (d1exe, d1sock, _d1cfg, d1env) = launch_in(d1.path());
let jc = join_code_from(&run_cmd(
&d1exe,
&d1env,
&["join", &invite, "--name", "Alice", "--user-id", "alice"],
));
run_cmd(
&opexe,
&openv,
&["org", "approve", &jc, "--groups", "team-eng"],
);
let jc1 = mcpmesh::roster::enroll::JoinCode::decode(&jc).expect("decode join code 1");
let d1_endpoint = jc1.device_endpoint_id.clone();
let alice_user_pk = jc1.user_pk.clone();
let d2 = tempfile::tempdir().unwrap();
let (d2exe, _d2sock, _d2cfg, d2env) = launch_in(d2.path());
let code_out = run_cmd(&d2exe, &d2env, &["devices", "code", "--label", "desktop"]);
assert!(
code_out.status.success(),
"devices code: {}",
String::from_utf8_lossy(&code_out.stderr)
);
let dcode = device_code_from(&code_out);
let dc = mcpmesh::roster::enroll::DeviceCode::decode(&dcode).expect("decode device code");
let d2_endpoint = dc.device_endpoint_id.clone();
assert_ne!(d2_endpoint, d1_endpoint, "device 2 has its own endpoint");
let code_stdout = String::from_utf8_lossy(&code_out.stdout);
assert!(
!code_stdout.contains("mcpmesh-join:") && !code_stdout.contains("user.key"),
"devices code carries no key material / join code: {code_stdout}"
);
let add = run_cmd(&d1exe, &d1env, &["devices", "add", &dcode]);
assert!(
add.status.success(),
"devices add: {}",
String::from_utf8_lossy(&add.stderr)
);
let add_stdout = String::from_utf8_lossy(&add.stdout);
assert!(
add_stdout.contains("Join code fingerprint:"),
"devices add prints the join-code fingerprint (ceremony): {add_stdout}"
);
let jc2 = join_code_from(&add);
let jc2d = mcpmesh::roster::enroll::JoinCode::decode(&jc2).expect("decode join code 2");
assert_eq!(
jc2d.user_pk, alice_user_pk,
"same user_pk (append, not new user)"
);
assert_eq!(jc2d.requested_user_id, "alice");
assert_eq!(jc2d.device_endpoint_id, d2_endpoint);
let ap = run_cmd(
&opexe,
&openv,
&["org", "approve", &jc2, "--groups", "team-eng"],
);
assert!(
ap.status.success(),
"second approve: {}",
String::from_utf8_lossy(&ap.stderr)
);
let r: Roster =
serde_json::from_slice(&std::fs::read(opcfg.join("roster.json")).unwrap()).unwrap();
assert_eq!(r.serial, 3);
assert_eq!(
r.users.iter().filter(|u| u.user_id == "alice").count(),
1,
"alice is a single user entry (append, not a new user)"
);
let alice = r.users.iter().find(|u| u.user_id == "alice").unwrap();
assert_eq!(
alice.devices.len(),
2,
"the second device is appended to alice"
);
assert_eq!(
alice.user_pk, alice_user_pk,
"same user_pk after the append"
);
let endpoints: Vec<&str> = alice
.devices
.iter()
.map(|d| d.endpoint_id.as_str())
.collect();
assert!(
endpoints.contains(&d1_endpoint.as_str()),
"device 1 endpoint present"
);
assert!(
endpoints.contains(&d2_endpoint.as_str()),
"device 2 endpoint present"
);
let fresh = tempfile::tempdir().unwrap();
let (fexe, _fsock, _fcfg, fenv) = launch_in(fresh.path());
let unenrolled = run_cmd(&fexe, &fenv, &["devices", "add", &dcode]);
assert!(
!unenrolled.status.success(),
"devices add on an unenrolled machine must refuse"
);
assert!(
String::from_utf8_lossy(&unenrolled.stderr).contains("not enrolled"),
"the refusal names the enrollment requirement: {}",
String::from_utf8_lossy(&unenrolled.stderr)
);
shutdown_daemon(&opsock).await;
shutdown_daemon(&d1sock).await;
})
.await
.expect("devices add test timed out");
}