#![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, RosterDevice, RosterUser, encode_b64u};
use serde_json::{Value, 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")
}
fn signed_roster(root: &SigningKey, serial: u64) -> Roster {
mint_signed(
root,
Roster {
format: "mcpmesh-roster/1".into(),
org_id: "acme".into(),
serial,
issued_at: "2020-01-01T00:00:00Z".into(),
expires_at: "2999-01-01T00:00:00Z".into(),
groups: vec!["team-eng".into()],
users: vec![RosterUser {
user_id: "alice".into(),
display_name: "Alice".into(),
user_pk: encode_b64u(&[1u8; 32]),
groups: vec!["team-eng".into()],
devices: vec![RosterDevice {
endpoint_id: encode_b64u(&[2u8; 32]),
label: "laptop".into(),
role: "primary".into(),
}],
}],
revoked_endpoints: vec![],
sig: String::new(),
},
)
}
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 persisted_serial(config_dir: &Path) -> Option<u64> {
let bytes = std::fs::read(config_dir.join("roster.json")).ok()?;
let v: Value = serde_json::from_slice(&bytes).ok()?;
v.get("serial").and_then(Value::as_u64)
}
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 pinned_org_root_pk(config_dir: &Path) -> Option<String> {
Config::load(&config_dir.join("config.toml"))
.expect("reload config")
.identity
.org_root_pk
}
#[tokio::test(flavor = "multi_thread")]
async fn roster_install_pins_omits_rejects_wrong_pk_and_rejects_rollback() {
tokio::time::timeout(Duration::from_secs(45), async {
let tmp = tempfile::tempdir().unwrap();
let (exe, socket, config_dir, env) = launch_in(tmp.path());
let root = SigningKey::from_bytes(&[9u8; 32]);
let org_root_pk = encode_b64u(&root.verifying_key().to_bytes());
let f42 = write_roster(tmp.path(), "roster-42.json", &signed_roster(&root, 42));
let out = run_cmd(
&exe,
&env,
&[
"internal",
"roster",
"install",
f42.to_str().unwrap(),
"--org-root-pk",
&org_root_pk,
],
);
let stdout = String::from_utf8_lossy(&out.stdout).into_owned();
assert!(
out.status.success(),
"first `roster install` must exit 0; stderr: {}",
String::from_utf8_lossy(&out.stderr)
);
assert_eq!(
stdout.trim_end(),
"Installed roster for org 'acme' (serial 42). Severed 0 live sessions.",
"unexpected confirmation:\n{stdout}"
);
assert_eq!(persisted_serial(&config_dir), Some(42));
assert_eq!(
pinned_org_root_pk(&config_dir).as_deref(),
Some(org_root_pk.as_str()),
"org_root_pk must be pinned in config"
);
assert_eq!(
Config::load(&config_dir.join("config.toml"))
.unwrap()
.identity
.org_id
.as_deref(),
Some("acme"),
"org_id must be pinned in config"
);
let f43 = write_roster(tmp.path(), "roster-43.json", &signed_roster(&root, 43));
let out = run_cmd(
&exe,
&env,
&["internal", "roster", "install", f43.to_str().unwrap()],
);
assert!(
out.status.success(),
"an omitted --org-root-pk must reuse the pinned pk and succeed; stderr: {}",
String::from_utf8_lossy(&out.stderr)
);
assert_eq!(
String::from_utf8_lossy(&out.stdout).trim_end(),
"Installed roster for org 'acme' (serial 43). Severed 0 live sessions."
);
assert_eq!(persisted_serial(&config_dir), Some(43));
let wrong_pk = encode_b64u(
&SigningKey::from_bytes(&[8u8; 32])
.verifying_key()
.to_bytes(),
);
let f44 = write_roster(tmp.path(), "roster-44.json", &signed_roster(&root, 44));
let out = run_cmd(
&exe,
&env,
&[
"internal",
"roster",
"install",
f44.to_str().unwrap(),
"--org-root-pk",
&wrong_pk,
],
);
assert!(
!out.status.success(),
"a wrong --org-root-pk must exit non-zero; stdout: {}",
String::from_utf8_lossy(&out.stdout)
);
assert_eq!(
persisted_serial(&config_dir),
Some(43),
"a failed (wrong-pk) install must not touch the on-disk roster"
);
assert_eq!(
pinned_org_root_pk(&config_dir).as_deref(),
Some(org_root_pk.as_str()),
"pin-after-validate: a failed install must NOT re-pin the wrong anchor"
);
let f41 = write_roster(tmp.path(), "roster-41.json", &signed_roster(&root, 41));
let out = run_cmd(
&exe,
&env,
&["internal", "roster", "install", f41.to_str().unwrap()],
);
assert!(
!out.status.success(),
"a rolled-back serial must exit non-zero; stdout: {}",
String::from_utf8_lossy(&out.stdout)
);
assert_eq!(
persisted_serial(&config_dir),
Some(43),
"the on-disk roster must be untouched by the rejected rollback"
);
shutdown_daemon(&socket).await;
})
.await
.expect("roster install test timed out");
}
fn roster_owning_our_device(
root: &SigningKey,
user_id: &str,
device_endpoint: &[u8; 32],
) -> Roster {
mint_signed(
root,
Roster {
format: "mcpmesh-roster/1".into(),
org_id: "acme".into(),
serial: 7,
issued_at: "2020-01-01T00:00:00Z".into(),
expires_at: "2999-01-01T00:00:00Z".into(),
groups: vec!["team-eng".into()],
users: vec![RosterUser {
user_id: user_id.into(),
display_name: "Alice".into(),
user_pk: encode_b64u(&[1u8; 32]),
groups: vec!["team-eng".into()],
devices: vec![RosterDevice {
endpoint_id: encode_b64u(device_endpoint),
label: "laptop".into(),
role: "primary".into(),
}],
}],
revoked_endpoints: vec![],
sig: String::new(),
},
)
}
#[tokio::test(flavor = "multi_thread")]
async fn roster_install_reconciles_proposed_user_id_to_the_authoritative_value() {
tokio::time::timeout(Duration::from_secs(45), async {
let tmp = tempfile::tempdir().unwrap();
let (exe, socket, config_dir, env) = launch_in(tmp.path());
let device_secret = [7u8; 32];
let device_key_path = config_dir.join("device.key");
std::fs::write(&device_key_path, device_secret).unwrap();
let device_endpoint = SigningKey::from_bytes(&device_secret)
.verifying_key()
.to_bytes();
let config_path = config_dir.join("config.toml");
std::fs::write(
&config_path,
format!(
"[network]\nrelay_mode = \"disabled\"\n[identity]\nuser_id = \"alice-proposed\"\ndevice_key = \"{}\"\n",
device_key_path.display()
),
)
.unwrap();
let root = SigningKey::from_bytes(&[9u8; 32]);
let org_root_pk = encode_b64u(&root.verifying_key().to_bytes());
let roster = roster_owning_our_device(&root, "alice-authoritative", &device_endpoint);
let file = write_roster(tmp.path(), "roster-7.json", &roster);
assert_eq!(
Config::load(&config_path).unwrap().identity.user_id.as_deref(),
Some("alice-proposed"),
"config must start with the proposed user_id"
);
let out = run_cmd(
&exe,
&env,
&[
"internal",
"roster",
"install",
file.to_str().unwrap(),
"--org-root-pk",
&org_root_pk,
],
);
assert!(
out.status.success(),
"roster install must exit 0; stderr: {}",
String::from_utf8_lossy(&out.stderr)
);
assert_eq!(
Config::load(&config_path).unwrap().identity.user_id.as_deref(),
Some("alice-authoritative"),
"config user_id must be reconciled to the roster's authoritative value"
);
let mut client = connect_control(&socket).await.expect("connect control");
let status = client
.request_value(&json!({ "method": "status", "params": {} }))
.await
.expect("status over the control API");
assert_eq!(
status["roster"]["state"], "approved",
"roster_status must flip pending→approved once the view holds this device: {status}"
);
shutdown_daemon(&socket).await;
})
.await
.expect("reconcile test timed out");
}