use std::path::PathBuf;
use anyhow::Context;
use clap::{Parser, Subcommand};
use mcpmesh::{client, config, daemon, doctor, pairing, proxy, roster, util};
use mcpmesh_local_api::{
BackendKind, BackendSpec, BlobFetchResult, BlobPublishResult, BlobScopeList, Hello,
InviteResult, PairResult, PresencePeer, RecentPairing, Request, RosterInstallResult,
RosterStatus, StatusResult,
};
use mcpmesh_trust::{DeviceKey, paths};
use serde_json::json;
macro_rules! serve_example {
() => {
"mcpmesh serve notes -- npx -y @modelcontextprotocol/server-filesystem ~/notes"
};
}
const SERVE_EXAMPLE: &str = serve_example!();
#[derive(Parser)]
#[command(name = "mcpmesh", version)]
struct Cli {
#[command(subcommand)]
cmd: Option<Cmd>,
}
#[derive(Subcommand)]
enum Cmd {
Status,
Doctor,
#[command(after_help = concat!(
"Example — share a folder of notes (needs npx; no MCP server of your own required):\n ",
serve_example!(),
"\n\nThen `mcpmesh invite notes` mints an invite to send whoever you're sharing with."
))]
Serve {
name: String,
#[arg(long)]
allow: Option<String>,
#[arg(last = true, required = true)]
cmd: Vec<String>,
},
Connect {
target: String,
},
Invite {
services: Vec<String>,
},
Pair {
invite: Option<String>,
#[arg(long, value_name = "petname")]
remove: Option<String>,
},
Use {
target: String,
},
Join {
org_invite: String,
#[arg(long)]
name: Option<String>,
#[arg(long)]
user_id: Option<String>,
#[arg(long, default_value = "laptop")]
label: String,
},
Org {
#[command(subcommand)]
command: OrgCmd,
},
Devices {
#[command(subcommand)]
command: DevicesCmd,
},
Internal {
#[command(subcommand)]
command: Internal,
},
}
#[derive(Subcommand)]
enum OrgCmd {
Create {
name: String,
#[arg(long)]
expires: Option<String>,
#[arg(long)]
roster_url: Option<String>,
},
Approve {
join_code: String,
#[arg(long)]
groups: String,
#[arg(long)]
user_id: Option<String>,
},
Revoke {
target: String,
#[arg(long)]
user_key: bool,
},
}
#[derive(Subcommand)]
enum DevicesCmd {
Code {
#[arg(long, default_value = "desktop")]
label: String,
},
Add {
device_code: String,
},
}
#[derive(Subcommand)]
enum Internal {
Daemon,
Id,
Peer {
#[command(subcommand)]
command: PeerCmd,
},
Roster {
#[command(subcommand)]
command: RosterCmd,
},
Blob {
#[command(subcommand)]
command: BlobCmd,
},
Audit {
#[command(subcommand)]
command: AuditCmd,
},
Watch,
}
#[derive(Subcommand)]
enum BlobCmd {
Publish {
scope: String,
file: PathBuf,
},
Grant { scope: String, principal: String },
List,
Fetch {
ticket: String,
dest: PathBuf,
},
}
#[derive(clap::Subcommand)]
enum AuditCmd {
Tail {
#[arg(long, default_value_t = 20)]
lines: usize,
#[arg(long)]
kind: Option<String>,
#[arg(long)]
peer: Option<String>,
},
List,
Prune {
#[arg(long, value_name = "YYYY-MM")]
before: String,
},
}
#[derive(Subcommand)]
enum RosterCmd {
Install {
file: PathBuf,
#[arg(long)]
org_root_pk: Option<String>,
},
}
#[derive(Subcommand)]
enum PeerCmd {
Add {
petname: String,
endpoint_id: String,
#[arg(long)]
allow: Option<String>,
},
}
fn main() -> anyhow::Result<()> {
let cli = Cli::parse();
match cli.cmd {
Some(Cmd::Internal {
command: Internal::Daemon,
}) => daemon::run(),
Some(Cmd::Internal {
command: Internal::Id,
}) => run_internal_id(),
Some(Cmd::Serve { name, allow, cmd }) => run_serve(name, allow, cmd),
Some(Cmd::Connect { target }) => run_connect(target),
Some(Cmd::Invite { services }) => run_invite(services),
Some(Cmd::Pair { invite, remove }) => run_pair(invite, remove),
Some(Cmd::Use { target }) => run_use(target),
Some(Cmd::Join {
org_invite,
name,
user_id,
label,
}) => run_join(org_invite, name, user_id, label),
Some(Cmd::Org {
command:
OrgCmd::Create {
name,
expires,
roster_url,
},
}) => run_org_create(name, expires, roster_url),
Some(Cmd::Org {
command:
OrgCmd::Approve {
join_code,
groups,
user_id,
},
}) => run_org_approve(join_code, groups, user_id),
Some(Cmd::Org {
command: OrgCmd::Revoke { target, user_key },
}) => run_org_revoke(target, user_key),
Some(Cmd::Devices {
command: DevicesCmd::Code { label },
}) => run_devices_code(label),
Some(Cmd::Devices {
command: DevicesCmd::Add { device_code },
}) => run_devices_add(device_code),
Some(Cmd::Internal {
command:
Internal::Peer {
command:
PeerCmd::Add {
petname,
endpoint_id,
allow,
},
},
}) => run_peer_add(petname, endpoint_id, allow),
Some(Cmd::Internal {
command:
Internal::Roster {
command: RosterCmd::Install { file, org_root_pk },
},
}) => run_roster_install(file, org_root_pk),
Some(Cmd::Internal {
command: Internal::Blob { command },
}) => run_internal_blob(command),
Some(Cmd::Internal {
command: Internal::Audit { command },
}) => run_internal_audit(command),
Some(Cmd::Internal {
command: Internal::Watch,
}) => run_watch(),
Some(Cmd::Doctor) => doctor::run_doctor(),
Some(Cmd::Status) | None => run_status(),
}
}
fn with_daemon<T>(
f: impl AsyncFnOnce(client::ControlClient) -> anyhow::Result<T>,
) -> anyhow::Result<T> {
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(async move {
let client = client::ensure_daemon().await?;
f(client).await
})
}
fn run_serve(name: String, allow: Option<String>, cmd: Vec<String>) -> anyhow::Result<()> {
let allow = split_csv(allow);
with_daemon(async move |mut client| {
client
.request(Request::RegisterService {
name: name.clone(),
backend: BackendSpec::Run { cmd },
allow,
})
.await?;
println!("serving '{name}'");
println!(
"Next: run `mcpmesh invite {name}` to mint a one-time invite, and send it to the \
person you want to share it with."
);
Ok(())
})
}
fn run_connect(target: String) -> anyhow::Result<()> {
let (peer, service) = proxy::split_target(&target)?;
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(proxy::run(peer, service))
}
fn run_invite(services: Vec<String>) -> anyhow::Result<()> {
if services.is_empty() {
anyhow::bail!("specify at least one service to grant (e.g. `mcpmesh invite notes`)");
}
with_daemon(async move |mut client| {
let result = client
.request(Request::Invite {
services: services.clone(),
})
.await?;
let invite: InviteResult = serde_json::from_value(result).context("parse invite result")?;
for line in invite_lines(&invite, &services, util::epoch_now_u64()) {
println!("{line}");
}
Ok(())
})
}
fn run_pair(invite: Option<String>, remove: Option<String>) -> anyhow::Result<()> {
match (invite, remove) {
(Some(_), Some(_)) => {
anyhow::bail!("provide an invite to redeem OR --remove <petname>, not both")
}
(None, None) => {
anyhow::bail!("provide an invite to redeem, or --remove <petname> to unpair")
}
(Some(invite_line), None) => with_daemon(async move |mut client| {
let result = client.request(Request::Pair { invite_line }).await?;
let paired: PairResult = serde_json::from_value(result).context("parse pair result")?;
for line in pair_lines(&paired) {
println!("{line}");
}
Ok(())
}),
(None, Some(petname)) => with_daemon(async move |mut client| {
client
.request(Request::PeerRemove {
petname: petname.clone(),
})
.await?;
println!("Unpaired {petname}.");
Ok(())
}),
}
}
fn invite_lines(invite: &InviteResult, services: &[String], now: u64) -> Vec<String> {
vec![
format!(
"One-time invite (expires {}). Share it out-of-band:",
friendly_expiry(invite.expires_at_epoch, now)
),
format!(" {}", invite.invite_line),
format!("Whoever redeems it can access: {}", services.join(", ")),
String::new(),
"Next: send them that line over any channel. They redeem it with `mcpmesh pair <line>`,"
.to_string(),
"which prints a short safety code — run `mcpmesh status` to see yours and confirm the two"
.to_string(),
"match, out loud. Same words = the pairing is authentic.".to_string(),
]
}
fn pair_lines(result: &PairResult) -> Vec<String> {
let peer = &result.peer_petname;
let mut lines = vec![
format!("Paired with {peer} — code: {}", result.sas_code),
format!(
"Next: confirm this code matches what {peer} sees, out loud (they see it under \
`mcpmesh status`). Same words = the pairing is authentic."
),
];
if result.services.is_empty() {
return lines;
}
let mounts = result
.services
.iter()
.map(|s| format!("{peer}/{s}"))
.collect::<Vec<_>>()
.join(", ");
lines.push(String::new());
lines.push(format!("You can now use: {mounts}"));
lines.push(String::new());
lines.extend(proxy::client_instruction_lines(peer, &result.services));
lines
}
fn friendly_expiry(expires_at_epoch: u64, now: u64) -> String {
let remaining = expires_at_epoch.saturating_sub(now);
if remaining < 60 {
return "soon".to_string();
}
if remaining < 3600 {
let mins = (remaining + 30) / 60; return format!("in {mins}m");
}
let hours = (remaining + 1800) / 3600; format!("in {hours}h")
}
fn run_use(target: String) -> anyhow::Result<()> {
let (peer, service) = proxy::split_target(&target)?;
for line in proxy::client_instruction_lines(&peer, &[service]) {
println!("{line}");
}
Ok(())
}
fn run_peer_add(petname: String, endpoint_id: String, allow: Option<String>) -> anyhow::Result<()> {
let allow = split_csv(allow);
with_daemon(async move |mut client| {
client
.request_value(&json!({
"method": "peer_add",
"params": { "petname": petname, "endpoint_id": endpoint_id, "allow": allow }
}))
.await?;
println!("added peer '{petname}'");
Ok(())
})
}
fn run_roster_install(file: PathBuf, org_root_pk: Option<String>) -> anyhow::Result<()> {
let path = file.to_string_lossy().into_owned();
with_daemon(async move |mut client| {
let result = client
.request(Request::RosterInstall { path, org_root_pk })
.await?;
let installed: RosterInstallResult =
serde_json::from_value(result).context("parse roster install result")?;
println!("{}", roster_install_line(&installed));
Ok(())
})
}
fn run_internal_blob(command: BlobCmd) -> anyhow::Result<()> {
with_daemon(async move |mut client| {
match command {
BlobCmd::Publish { scope, file } => {
let path = file.to_string_lossy().into_owned();
let result = client.request(Request::BlobPublish { scope, path }).await?;
let r: BlobPublishResult =
serde_json::from_value(result).context("parse blob publish result")?;
println!("Published (hash {}).", r.hash);
println!("{}", r.ticket);
}
BlobCmd::Grant { scope, principal } => {
client
.request(Request::BlobGrant {
scope: scope.clone(),
principal: principal.clone(),
})
.await?;
println!("Granted scope '{scope}' to '{principal}'.");
}
BlobCmd::List => {
let result = client.request(Request::BlobList).await?;
let r: BlobScopeList =
serde_json::from_value(result).context("parse blob list result")?;
for s in r.scopes {
println!(
"{}: {} blob(s), granted to [{}]",
s.name,
s.hashes.len(),
s.grants.join(", ")
);
}
}
BlobCmd::Fetch { ticket, dest } => {
let dest_path = dest.to_string_lossy().into_owned();
let result = client
.request(Request::BlobFetch { ticket, dest_path })
.await?;
let r: BlobFetchResult =
serde_json::from_value(result).context("parse blob fetch result")?;
println!(
"Fetched {} bytes (hash {}) → {}",
r.bytes_len,
r.hash,
dest.display()
);
}
}
Ok(())
})
}
fn run_internal_audit(command: AuditCmd) -> anyhow::Result<()> {
use mcpmesh::audit;
let dir = paths::default_audit_dir()?;
match command {
AuditCmd::Tail { lines, kind, peer } => {
let kind_filter = match kind.as_deref() {
Some(s) => {
Some(audit::parse_kind(s).with_context(|| format!("unknown --kind '{s}'"))?)
}
None => None,
};
let all = audit::read_all_records(&dir)?;
let filtered = audit::filter_records(&all, kind_filter, peer.as_deref());
let start = filtered.len().saturating_sub(lines);
for rec in &filtered[start..] {
println!("{}", serde_json::to_string(rec)?);
}
}
AuditCmd::List => {
for (month, _, size) in audit::list_month_files(&dir)? {
println!("{month} {size} bytes");
}
}
AuditCmd::Prune { before } => {
let deleted = audit::prune_before(&dir, &before)?;
if deleted.is_empty() {
println!("Nothing to prune before {before}.");
} else {
println!("Pruned {} month(s): {}.", deleted.len(), deleted.join(", "));
}
}
}
Ok(())
}
fn roster_install_line(result: &RosterInstallResult) -> String {
let sessions = if result.severed == 1 {
"session"
} else {
"sessions"
};
format!(
"Installed roster for org '{}' (serial {}). Severed {} live {sessions}.",
result.org_id, result.serial, result.severed
)
}
const DEFAULT_EXPIRES_SECS: i64 = 90 * 86_400;
fn slug(name: &str) -> String {
let mut s = String::new();
let mut last_dash = true; for c in name.chars() {
if c.is_ascii_alphanumeric() {
s.push(c.to_ascii_lowercase());
last_dash = false;
} else if !last_dash {
s.push('-');
last_dash = true;
}
}
while s.ends_with('-') {
s.pop();
}
if s.is_empty() { "user".to_string() } else { s }
}
fn run_join(
org_invite: String,
name: Option<String>,
user_id: Option<String>,
label: String,
) -> anyhow::Result<()> {
use mcpmesh_trust::keys::UserKey;
use mcpmesh_trust::roster::encode_b64u;
use mcpmesh_trust::roster::sign::sign_device_binding;
let invite = roster::enroll::OrgInviteCode::decode(&org_invite)
.context("the org invite is not a valid mcpmesh-org: code")?;
let root_pk = mcpmesh_trust::roster::decode_endpoint_id(&invite.org_root_pk)
.context("org invite carries an invalid org_root_pk")?;
let display_name = name.unwrap_or_else(|| "user".to_string());
let requested_user_id = user_id.unwrap_or_else(|| slug(&display_name));
let user_key_path = paths::default_user_key_path()?;
let (user_key, _created) = UserKey::load_or_generate(&user_key_path)
.map_err(|e| anyhow::anyhow!("user key error at {}: {e}", user_key_path.display()))?;
let device_key = load_device_key()?;
let device_id = device_key.public_bytes();
let binding = sign_device_binding(user_key.signing_key(), &device_id);
let code_fp = pairing::sas::join_code_fingerprint(&user_key.public_bytes(), &device_id);
let join = roster::enroll::JoinCode {
display_name: display_name.clone(),
requested_user_id: requested_user_id.clone(),
user_pk: encode_b64u(&user_key.public_bytes()),
device_endpoint_id: encode_b64u(&device_id),
device_label: label,
binding_sig: encode_b64u(&binding),
}
.encode();
with_daemon(async |mut client| {
client
.request(Request::OrgJoin {
org_id: invite.org_id.clone(),
org_root_pk: invite.org_root_pk.clone(),
user_id: requested_user_id.clone(),
user_key: user_key_path.to_string_lossy().into_owned(),
})
.await?;
if let Some(url) = &invite.roster_url {
client
.request(Request::SetRosterUrl { url: url.clone() })
.await?;
}
Ok(())
})?;
let fingerprint = pairing::sas::fingerprint_words(&root_pk);
println!("Joined org '{}' as '{requested_user_id}'.", invite.org_id);
println!("Org root fingerprint: {fingerprint}");
println!(
" → Confirm this matches what the operator reads back, out-of-band, before they approve you."
);
println!("Send the operator your join code: {join}");
println!("Join code fingerprint: {code_fp}");
println!(
" → Read this back to your operator out-of-band so they confirm they received YOUR join code (not a substituted one)."
);
Ok(())
}
fn run_org_create(
name: String,
expires: Option<String>,
roster_url: Option<String>,
) -> anyhow::Result<()> {
use mcpmesh_trust::keys::OrgRootKey;
use mcpmesh_trust::roster::sign::mint_signed;
use mcpmesh_trust::roster::{encode_b64u, mutate};
let key_path = paths::default_org_root_key_path()?;
let (root, created) = OrgRootKey::load_or_generate(&key_path)
.map_err(|e| anyhow::anyhow!("org root key error at {}: {e}", key_path.display()))?;
if !created {
anyhow::bail!(
"this node already holds an org root key ({}); `org create` is one-time per node",
key_path.display()
);
}
let expires_secs = match &expires {
Some(s) => config::parse_duration(s).map_err(|e| anyhow::anyhow!("bad --expires: {e}"))?,
None => DEFAULT_EXPIRES_SECS,
};
let now = util::epoch_now_i64();
let roster = mint_signed(
root.signing_key(),
mutate::empty_roster(&name, 1, now, now.saturating_add(expires_secs)),
);
let org_root_pk = encode_b64u(&root.public_bytes());
let result = install_signed_roster(&roster, Some(org_root_pk.clone()))?;
if let Some(url) = &roster_url {
with_daemon(async |mut client| {
client
.request(Request::SetRosterUrl { url: url.clone() })
.await?;
Ok(())
})?;
}
let invite = roster::enroll::OrgInviteCode {
org_id: name.clone(),
org_root_pk,
roster_url: roster_url.clone(),
}
.encode();
let fingerprint = pairing::sas::fingerprint_words(&root.public_bytes());
println!(
"Created org '{}' (roster serial {}).",
result.org_id, result.serial
);
println!("Invite someone: {invite}");
println!("Org root fingerprint: {fingerprint} (read this aloud when you approve joiners)");
Ok(())
}
fn load_operator_roster() -> anyhow::Result<(
mcpmesh_trust::keys::OrgRootKey,
mcpmesh_trust::roster::Roster,
)> {
let key_path = paths::default_org_root_key_path()?;
if !key_path.exists() {
anyhow::bail!(
"this node is not an org operator (no org root key); run `mcpmesh org create` first"
);
}
let (root, _) = mcpmesh_trust::keys::OrgRootKey::load_or_generate(&key_path)
.map_err(|e| anyhow::anyhow!("org root key error at {}: {e}", key_path.display()))?;
let roster_path = paths::default_roster_path()?;
let bytes = std::fs::read(&roster_path).with_context(|| {
format!(
"no installed roster at {} — run `org create`",
roster_path.display()
)
})?;
let roster: mcpmesh_trust::roster::Roster =
serde_json::from_slice(&bytes).context("parse installed roster")?;
Ok((root, roster))
}
fn run_org_approve(
join_code: String,
groups: String,
user_id: Option<String>,
) -> anyhow::Result<()> {
use mcpmesh_trust::roster::sign::{sign, verify_device_binding};
use mcpmesh_trust::roster::{decode_endpoint_id, mutate};
let jc = roster::enroll::JoinCode::decode(&join_code)
.context("the join code is not a valid mcpmesh-join: code")?;
let user_pk = decode_endpoint_id(&jc.user_pk).context("join code has an invalid user_pk")?;
let device_id = decode_endpoint_id(&jc.device_endpoint_id)
.context("join code has an invalid device endpoint")?;
let sig = mcpmesh_trust::roster::decode_b64u(&jc.binding_sig)
.context("join code has an invalid signature")?;
verify_device_binding(&user_pk, &device_id, &sig).map_err(|_| {
anyhow::anyhow!("join code device binding failed — the code is forged or corrupt")
})?;
let (root, mut roster) = load_operator_roster()?;
let uid = user_id.unwrap_or(jc.requested_user_id);
let groups = split_csv(Some(groups));
let code_fp = pairing::sas::join_code_fingerprint(&user_pk, &device_id);
println!(
"Approving join code {code_fp} for '{}' as user '{uid}', groups [{}].",
jc.display_name,
groups.join(", ")
);
println!(
" → Verify {code_fp} matches what the joiner read back to you out-of-band; if it doesn't, \
run `org revoke` on this device."
);
roster.serial += 1;
mutate::upsert_member(
&mut roster,
&uid,
&jc.display_name,
&jc.user_pk, &groups,
&jc.device_endpoint_id, &jc.device_label,
)
.map_err(|e| anyhow::anyhow!("roster mutation rejected: {e}"))?;
sign(root.signing_key(), &mut roster).map_err(|e| anyhow::anyhow!("sign roster: {e}"))?;
let result = install_signed_roster(&roster, None)?; println!(
"Approved '{}' into [{}] (org '{}', serial {}).",
uid,
groups.join(", "),
result.org_id,
result.serial
);
Ok(())
}
fn run_org_revoke(target: String, user_key: bool) -> anyhow::Result<()> {
use mcpmesh_trust::roster::mutate;
use mcpmesh_trust::roster::sign::sign;
let (root, mut roster) = load_operator_roster()?;
roster.serial += 1;
let action: String = if user_key {
mutate::remove_user(&mut roster, &target, false).map_err(|e| anyhow::anyhow!("{e}"))?;
format!(
"Rotated '{target}': removed from the roster. They re-enroll with a fresh user key \
(same device), then re-approve with the same user_id"
)
} else if let Some((person, device)) = target.split_once('/') {
mutate::revoke_device(&mut roster, person, device).map_err(|e| anyhow::anyhow!("{e}"))?;
format!("Revoked device '{person}/{device}'")
} else {
mutate::remove_user(&mut roster, &target, true).map_err(|e| anyhow::anyhow!("{e}"))?;
format!("Revoked person '{target}' (all devices)")
};
sign(root.signing_key(), &mut roster).map_err(|e| anyhow::anyhow!("sign roster: {e}"))?;
let result = install_signed_roster(&roster, None)?;
println!(
"{action} (org '{}', serial {}). Severed {} live session{}.",
result.org_id,
result.serial,
result.severed,
if result.severed == 1 { "" } else { "s" }
);
Ok(())
}
fn run_devices_code(label: String) -> anyhow::Result<()> {
use mcpmesh_trust::roster::encode_b64u;
let device_id = load_device_key()?.public_bytes();
let code = roster::enroll::DeviceCode {
device_endpoint_id: encode_b64u(&device_id),
device_label: label,
}
.encode();
println!("Give this to an already-enrolled device (`mcpmesh devices add`): {code}");
Ok(())
}
fn run_devices_add(device_code: String) -> anyhow::Result<()> {
use mcpmesh_trust::keys::UserKey;
use mcpmesh_trust::roster::encode_b64u;
use mcpmesh_trust::roster::sign::sign_device_binding;
let dc = roster::enroll::DeviceCode::decode(&device_code)
.context("not a valid mcpmesh-device: code")?;
let new_device_id = mcpmesh_trust::roster::decode_endpoint_id(&dc.device_endpoint_id)
.context("device code has an invalid endpoint id")?;
let cfg = config::Config::load(&paths::default_config_path()?)
.map_err(|e| anyhow::anyhow!("config: {e}"))?;
let user_id = cfg
.identity
.user_id
.clone()
.context("this device is not enrolled (no user_id); run `mcpmesh join` first")?;
let user_key_path = match cfg.identity.user_key.clone() {
Some(p) => p,
None => paths::default_user_key_path()?,
};
if !user_key_path.exists() {
anyhow::bail!(
"this device is not enrolled (no user key at {}); run `mcpmesh join` first",
user_key_path.display()
);
}
let (user_key, _) = UserKey::load_or_generate(&user_key_path)
.map_err(|e| anyhow::anyhow!("user key error at {}: {e}", user_key_path.display()))?;
let user_pk = user_key.public_bytes();
let binding = sign_device_binding(user_key.signing_key(), &new_device_id);
let join = roster::enroll::JoinCode {
display_name: user_id.clone(),
requested_user_id: user_id,
user_pk: encode_b64u(&user_pk),
device_endpoint_id: dc.device_endpoint_id,
device_label: dc.device_label,
binding_sig: encode_b64u(&binding),
}
.encode();
let code_fp = pairing::sas::join_code_fingerprint(&user_pk, &new_device_id);
println!("Send the operator this join code to add the device: {join}");
println!("Join code fingerprint: {code_fp}");
println!(
" → Read this back to your operator out-of-band so they confirm they received THIS device's \
join code (not a substituted one)."
);
Ok(())
}
struct TempFile(std::path::PathBuf);
impl Drop for TempFile {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.0);
}
}
fn install_signed_roster(
roster: &mcpmesh_trust::roster::Roster,
org_root_pk: Option<String>,
) -> anyhow::Result<RosterInstallResult> {
use std::sync::atomic::{AtomicU64, Ordering};
static SEQ: AtomicU64 = AtomicU64::new(0);
let seq = SEQ.fetch_add(1, Ordering::Relaxed);
let temp = paths::config_dir()?.join(format!(
"roster.staging.{}.{}.json",
std::process::id(),
seq
));
let _guard = TempFile(temp.clone());
if let Some(parent) = temp.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(&temp, serde_json::to_vec(roster)?)
.with_context(|| format!("write staged roster {}", temp.display()))?;
let path = temp.to_string_lossy().into_owned();
with_daemon(async move |mut client| {
let value = client
.request(Request::RosterInstall { path, org_root_pk })
.await?;
serde_json::from_value::<RosterInstallResult>(value).context("parse roster install result")
})
}
fn split_csv(value: Option<String>) -> Vec<String> {
value
.map(|s| {
s.split(',')
.map(str::trim)
.filter(|x| !x.is_empty())
.map(String::from)
.collect()
})
.unwrap_or_default()
}
fn run_status() -> anyhow::Result<()> {
let fingerprint = load_device_key()?.fingerprint();
let has_roster_url = paths::default_config_path()
.ok()
.and_then(|p| config::Config::load(&p).ok())
.map(|c| c.roster.url.is_some())
.unwrap_or(false);
with_daemon(async move |mut client| {
let hello = client.hello().clone();
let result = client.request(Request::Status).await?;
let status: StatusResult = serde_json::from_value(result).context("parse status result")?;
render_status(&fingerprint, &hello, &status, has_roster_url);
Ok(())
})
}
fn render_status(fingerprint: &str, hello: &Hello, status: &StatusResult, has_roster_url: bool) {
println!(
"{} v{} · stack {}",
hello.api, hello.api_version, hello.stack_version
);
println!("device {fingerprint}");
if let Some(user_id) = &status.self_user_id {
println!("identity {user_id}");
}
println!();
if status.services.is_empty() {
println!("no services configured");
} else {
println!("serving:");
for svc in &status.services {
let kind = backend_kind_label(svc.backend);
let allowed = if svc.allow.is_empty() {
"no one yet".to_owned()
} else {
svc.allow.join(", ")
};
println!(" {} · {kind} · allowed: {allowed}", svc.name);
}
}
println!();
if status.peers.is_empty() {
println!("no peers yet");
} else {
println!("peers:");
for peer in &status.peers {
let services = if peer.services.is_empty() {
"none".to_owned()
} else {
peer.services.join(", ")
};
match &peer.user_id {
Some(user_id) => {
println!(" {} · services: {services} · {user_id}", peer.name)
}
None => println!(" {} · services: {services}", peer.name),
}
}
}
if !status.reachability.is_empty() {
println!();
println!("reachability:");
for r in &status.reachability {
let label = match (r.reachable, r.age_secs) {
(_, None) => "…", (true, _) => "online",
(false, _) => "offline",
};
match r.rtt_ms {
Some(ms) if r.reachable => println!(" {} · {label} · {ms}ms", r.name),
_ => println!(" {} · {label}", r.name),
}
}
}
if !status.recent_pairings.is_empty() {
println!();
println!("recent pairings (confirm the code with the other side):");
for line in recent_pairing_lines(&status.recent_pairings, util::epoch_now_u64()) {
println!("{line}");
}
}
if let Some(roster) = &status.roster {
println!();
for line in roster_status_lines(roster, has_roster_url) {
println!("{line}");
}
}
if !status.presence.is_empty() {
println!();
println!("reachable:");
for line in presence_lines(&status.presence) {
println!("{line}");
}
}
let next = next_steps_lines(status);
if !next.is_empty() {
println!();
for line in next {
println!("{line}");
}
}
}
fn next_steps_lines(status: &StatusResult) -> Vec<String> {
let mut steps = Vec::new();
if let Some((peer, service)) = status
.peers
.iter()
.find_map(|p| p.services.first().map(|s| (&p.name, s)))
{
steps.push(format!(
" Use {peer}/{service} from your AI client: `mcpmesh use {peer}/{service}`"
));
}
if status.services.is_empty() {
steps.push(
" Share one of your MCP servers: `mcpmesh serve <name> -- <command that runs it>`"
.to_string(),
);
steps.push(format!(
" No MCP server yet? Share a folder: `{SERVE_EXAMPLE}`"
));
} else if let Some(svc) = status.services.iter().find(|s| s.allow.is_empty()) {
steps.push(format!(
" Nobody can reach '{}' yet: `mcpmesh invite {}`",
svc.name, svc.name
));
}
if status.peers.is_empty() {
steps.push(" Someone sent you an invite? `mcpmesh pair mcpmesh-invite:…`".to_string());
}
if steps.is_empty() {
return steps;
}
let mut lines = vec!["next steps:".to_string()];
lines.extend(steps);
lines
}
fn recent_pairing_lines(pairings: &[RecentPairing], now: u64) -> Vec<String> {
pairings
.iter()
.map(|p| {
format!(
" {} · code: {} · {}",
p.peer_petname,
p.sas_code,
friendly_age(p.paired_at_epoch, now)
)
})
.collect()
}
fn friendly_age(epoch: u64, now: u64) -> String {
let elapsed = now.saturating_sub(epoch);
if elapsed < 60 {
return "just now".to_string();
}
if elapsed < 3600 {
return format!("{}m ago", elapsed / 60);
}
if elapsed < 24 * 3600 {
return format!("{}h ago", elapsed / 3600);
}
format!("{}d ago", elapsed / (24 * 3600))
}
fn presence_lines(presence: &[PresencePeer]) -> Vec<String> {
presence
.iter()
.map(|p| {
format!(
" {} · {} · {} · {}",
p.user_id,
p.device_label,
p.role,
if p.online { "online" } else { "offline" }
)
})
.collect()
}
fn roster_status_lines(roster: &RosterStatus, has_roster_url: bool) -> Vec<String> {
let mut lines = vec![format!(
"roster: org {} · serial {} · {}",
roster.org_id, roster.serial, roster.state
)];
if !roster.org_root_fingerprint.is_empty() {
lines.push(format!(
" org root: {} (confirm out-of-band)",
roster.org_root_fingerprint
));
}
if !has_roster_url {
lines.push(
"hint: no roster URL configured — this node degrades after max_staleness with no way \
to re-confirm currency; set [roster].url"
.to_string(),
);
}
lines
}
fn backend_kind_label(kind: BackendKind) -> &'static str {
match kind {
BackendKind::Run => "run",
BackendKind::Socket => "socket",
}
}
fn run_watch() -> anyhow::Result<()> {
with_daemon(async move |client| {
let (mut reader, _writer) = client.open_stream("subscribe").await?;
println!("watching the mesh — Ctrl-C to stop");
while let Some(inbound) = reader.next().await? {
if let mcpmesh_net::framing::Inbound::Frame(v) = inbound
&& let Some(line) = render_frame(&v)
{
println!("{line}");
}
}
Ok(())
})
}
fn render_frame(v: &serde_json::Value) -> Option<String> {
match v["type"].as_str()? {
"snapshot" => Some(format!(
"snapshot: {} active session(s), {} peer(s) known",
v["active_sessions"].as_array().map_or(0, |a| a.len()),
v["reachability"].as_array().map_or(0, |a| a.len()),
)),
"event" => {
let r = &v["record"];
let peer = r["peer"]
.as_str()
.map(|p| format!("{p} "))
.unwrap_or_default();
let service = r["service"]
.as_str()
.map(|s| format!("→ {s}"))
.unwrap_or_default();
let status = r["status"]
.as_str()
.map(|s| format!(" ({s})"))
.unwrap_or_default();
let line = format!(
"[{}] {} {peer}{service}",
r["ts"].as_str().unwrap_or(""),
r["kind"].as_str().unwrap_or("?"),
);
Some(format!("{}{status}", line.trim_end()))
}
"lagged" => Some(format!(
"(lagged {} events — reconnect for a fresh snapshot)",
v["dropped"].as_u64().unwrap_or(0)
)),
_ => None,
}
}
fn run_internal_id() -> anyhow::Result<()> {
let key = load_device_key()?;
let endpoint_id = iroh::SecretKey::from_bytes(&key.secret_bytes()).public();
println!("{endpoint_id}");
Ok(())
}
fn load_device_key() -> anyhow::Result<DeviceKey> {
let cfg_path = paths::default_config_path()?;
let cfg = config::Config::load(&cfg_path)
.map_err(|e| anyhow::anyhow!("config error in {}: {e}", cfg_path.display()))?;
let key_path = match cfg.identity.device_key.clone() {
Some(p) => p,
None => paths::default_device_key_path()?,
};
let (key, _created) = DeviceKey::load_or_generate(&key_path)
.map_err(|e| anyhow::anyhow!("device key error at {}: {e}", key_path.display()))?;
Ok(key)
}
#[cfg(test)]
mod tests {
use mcpmesh_local_api::{PeerInfo, ServiceInfo};
use super::*;
const DAY: u64 = 24 * 60 * 60;
#[test]
fn invite_block_has_the_expected_shape() {
let invite = InviteResult {
invite_line: "mcpmesh-invite:MFRGGZDF".into(),
expires_at_epoch: 1_000_000 + DAY,
};
let lines = invite_lines(&invite, &["notes".to_string()], 1_000_000);
assert_eq!(
lines[..3],
[
"One-time invite (expires in 24h). Share it out-of-band:".to_string(),
" mcpmesh-invite:MFRGGZDF".to_string(),
"Whoever redeems it can access: notes".to_string(),
]
);
let rendered = lines.join("\n");
assert!(
rendered.contains("Next:") && rendered.contains("mcpmesh pair"),
"the invite must name the redeemer's exact next command:\n{rendered}"
);
assert!(
rendered.contains("mcpmesh status"),
"the invite must point at where the inviter confirms the code:\n{rendered}"
);
}
#[test]
fn invite_block_lists_multiple_services() {
let invite = InviteResult {
invite_line: "mcpmesh-invite:X".into(),
expires_at_epoch: 500 + DAY,
};
let lines = invite_lines(&invite, &["notes".to_string(), "kb".to_string()], 500);
assert_eq!(lines[2], "Whoever redeems it can access: notes, kb");
assert!(lines[1].contains("mcpmesh-invite:"));
}
#[test]
fn pair_lines_render_the_sas_and_mount_targets() {
let result = PairResult {
peer_petname: "alice".into(),
sas_code: "tango-fig-42".into(),
services: vec!["notes".into()],
};
let lines = pair_lines(&result);
assert_eq!(lines[0], "Paired with alice — code: tango-fig-42");
assert!(lines[0].contains("code: tango-fig-42"));
let rendered = lines.join("\n");
assert!(
rendered.contains("Next: confirm this code matches what alice sees"),
"pair must name the ceremony as the next step:\n{rendered}"
);
assert!(
rendered.contains("You can now use: alice/notes"),
"pair must name the mount target:\n{rendered}"
);
assert!(
rendered.contains("claude mcp add alice-notes -- mcpmesh connect alice/notes")
&& rendered.contains("claude_desktop_config.json"),
"pair must print the client instructions inline:\n{rendered}"
);
}
#[test]
fn pair_lines_join_multiple_mount_targets_as_peer_slash_service() {
let result = PairResult {
peer_petname: "alice".into(),
sas_code: "a-b-c".into(),
services: vec!["notes".into(), "kb".into()],
};
let rendered = pair_lines(&result).join("\n");
assert!(
rendered.contains("You can now use: alice/notes, alice/kb"),
"both grants are named as mount targets:\n{rendered}"
);
assert!(
rendered.contains("claude mcp add alice-notes -- mcpmesh connect alice/notes")
&& rendered.contains("claude mcp add alice-kb -- mcpmesh connect alice/kb"),
"every granted service gets its own instruction:\n{rendered}"
);
}
#[test]
fn pair_lines_leak_no_endpoint_id() {
let alice_id = iroh::SecretKey::from_bytes(&[7u8; 32]).public().to_string();
let result = PairResult {
peer_petname: "alice".into(),
sas_code: "tango-fig-42".into(),
services: vec!["notes".into()],
};
let rendered = pair_lines(&result).join("\n");
assert!(
!rendered.contains(&alice_id),
"pair output must not leak an EndpointId: {rendered}"
);
for term in ["ALPN", "ticket", "mcpmesh/pair/1", "mcpmesh/mcp/1"] {
assert!(!rendered.contains(term), "pair output leaked '{term}'");
}
}
#[test]
fn pair_lines_tolerate_an_empty_service_grant() {
let result = PairResult {
peer_petname: "alice".into(),
sas_code: "a-b-c".into(),
services: vec![],
};
let lines = pair_lines(&result);
assert_eq!(lines[0], "Paired with alice — code: a-b-c");
let rendered = lines.join("\n");
assert!(rendered.contains("Next: confirm this code matches what alice sees"));
assert!(
!rendered.contains("You can now use") && !rendered.contains("claude mcp add"),
"no dangling mount/instruction block with nothing granted:\n{rendered}"
);
}
fn status_with(services: Vec<ServiceInfo>, peers: Vec<PeerInfo>) -> StatusResult {
StatusResult {
stack_version: "0".into(),
services,
peers,
roster: None,
presence: Vec::new(),
self_user_id: None,
recent_pairings: Vec::new(),
reachability: Vec::new(),
}
}
fn service(name: &str, allow: &[&str]) -> ServiceInfo {
ServiceInfo {
name: name.into(),
allow: allow.iter().map(|s| s.to_string()).collect(),
backend: BackendKind::Run,
}
}
fn peer(name: &str, services: &[&str]) -> PeerInfo {
PeerInfo {
name: name.into(),
services: services.iter().map(|s| s.to_string()).collect(),
user_id: None,
}
}
#[test]
fn next_steps_on_a_fresh_node_name_both_directions() {
let rendered = next_steps_lines(&status_with(vec![], vec![])).join("\n");
assert!(
rendered.contains("mcpmesh serve <name> --"),
"a fresh node must be told how to share:\n{rendered}"
);
assert!(
rendered.contains("mcpmesh pair"),
"a fresh node must be told how to redeem an invite:\n{rendered}"
);
}
#[test]
fn next_steps_offer_a_runnable_serve_example_to_someone_with_no_mcp_server() {
let rendered = next_steps_lines(&status_with(vec![], vec![])).join("\n");
assert!(
rendered.contains(SERVE_EXAMPLE),
"a fresh node must offer a runnable serve example:\n{rendered}"
);
assert!(
SERVE_EXAMPLE.contains("mcpmesh serve notes --")
&& SERVE_EXAMPLE.contains("@modelcontextprotocol/server-filesystem"),
"the example must be complete and copy-pasteable: {SERVE_EXAMPLE}"
);
}
#[test]
fn next_steps_point_a_served_but_ungranted_service_at_invite() {
let rendered =
next_steps_lines(&status_with(vec![service("notes", &[])], vec![])).join("\n");
assert!(
rendered.contains("mcpmesh invite notes"),
"an ungranted service must be pointed at `invite <name>`:\n{rendered}"
);
let granted = next_steps_lines(&status_with(
vec![service("notes", &["bob"])],
vec![peer("bob", &[])],
))
.join("\n");
assert!(
!granted.contains("mcpmesh invite notes"),
"a granted service needs no invite nag:\n{granted}"
);
}
#[test]
fn next_steps_point_a_reachable_peer_service_at_use() {
let rendered =
next_steps_lines(&status_with(vec![], vec![peer("alice", &["notes"])])).join("\n");
assert!(
rendered.contains("mcpmesh use alice/notes"),
"a reachable peer service must be pointed at `use`:\n{rendered}"
);
let bare = next_steps_lines(&status_with(vec![], vec![peer("alice", &[])])).join("\n");
assert!(
!bare.contains("mcpmesh use"),
"a peer with no grants offers no use step:\n{bare}"
);
}
#[test]
fn next_steps_are_silent_on_a_fully_configured_node() {
let lines = next_steps_lines(&status_with(
vec![service("notes", &["bob"])],
vec![peer("bob", &["code"])],
));
let rendered = lines.join("\n");
assert!(rendered.contains("mcpmesh use bob/code"));
assert!(!rendered.contains("mcpmesh serve") && !rendered.contains("mcpmesh pair"));
}
#[test]
fn roster_install_line_renders_org_serial_and_pluralized_sever_count() {
let one = RosterInstallResult {
org_id: "acme".into(),
serial: 42,
severed: 1,
};
assert_eq!(
roster_install_line(&one),
"Installed roster for org 'acme' (serial 42). Severed 1 live session."
);
let none = RosterInstallResult {
org_id: "acme".into(),
serial: 7,
severed: 0,
};
assert_eq!(
roster_install_line(&none),
"Installed roster for org 'acme' (serial 7). Severed 0 live sessions."
);
let many = RosterInstallResult {
org_id: "acme".into(),
serial: 100,
severed: 3,
};
assert_eq!(
roster_install_line(&many),
"Installed roster for org 'acme' (serial 100). Severed 3 live sessions."
);
}
#[test]
fn roster_install_line_leaks_no_transport_vocabulary() {
let result = RosterInstallResult {
org_id: "acme".into(),
serial: 42,
severed: 1,
};
let line = roster_install_line(&result);
for term in [
"b64u:",
"endpoint",
"EndpointId",
"ALPN",
"roster.json",
"/",
"key",
] {
assert!(
!line.contains(term),
"roster install output leaked '{term}': {line}"
);
}
}
#[test]
fn roster_status_lines_render_org_serial_state_and_fingerprint() {
let roster = RosterStatus {
org_id: "acme".into(),
serial: 42,
state: "approved".into(),
org_root_fingerprint: "tango-fig-cabbage-anchor".into(),
};
let lines = roster_status_lines(&roster, true); assert_eq!(lines[0], "roster: org acme · serial 42 · approved");
assert_eq!(
lines[1],
" org root: tango-fig-cabbage-anchor (confirm out-of-band)"
);
}
#[test]
fn roster_status_lines_omit_the_org_root_line_when_the_fingerprint_is_absent() {
let roster = RosterStatus {
org_id: "acme".into(),
serial: 7,
state: "degraded".into(),
org_root_fingerprint: String::new(),
};
let lines = roster_status_lines(&roster, true); assert_eq!(lines, vec!["roster: org acme · serial 7 · degraded"]);
}
#[test]
fn roster_status_lines_append_url_less_hint_when_no_roster_url() {
let roster = RosterStatus {
org_id: "acme".into(),
serial: 7,
state: "approved".into(),
org_root_fingerprint: String::new(),
};
let lines = roster_status_lines(&roster, false); assert!(
lines
.iter()
.any(|l| l.contains("no roster URL configured") && l.contains("set [roster].url")),
"expected the URL-less degrade hint: {lines:?}"
);
let lines = roster_status_lines(&roster, true);
assert!(
!lines.iter().any(|l| l.contains("hint:")),
"no hint when a roster url is configured: {lines:?}"
);
}
#[test]
fn presence_lines_render_user_label_role_and_online_flag() {
let presence = vec![
PresencePeer {
user_id: "alice".into(),
device_label: "laptop".into(),
role: "primary".into(),
online: true,
},
PresencePeer {
user_id: "alice".into(),
device_label: "desktop".into(),
role: "mirror".into(),
online: false,
},
];
let lines = presence_lines(&presence);
assert_eq!(lines.len(), 2);
assert!(
lines[0].contains("alice")
&& lines[0].contains("laptop")
&& lines[0].contains("primary")
&& lines[0].contains("online"),
"the online primary renders user·label·role·online: {lines:?}"
);
assert!(
lines[1].contains("desktop")
&& lines[1].contains("mirror")
&& lines[1].contains("offline"),
"the dead mirror renders offline: {lines:?}"
);
}
#[test]
fn presence_lines_leak_no_transport_vocabulary() {
let presence = vec![PresencePeer {
user_id: "alice".into(),
device_label: "laptop".into(),
role: "primary".into(),
online: true,
}];
let rendered = presence_lines(&presence).join("\n");
for term in ["b64u:", "EndpointId", "endpoint", "ALPN", "pubkey", "hash"] {
assert!(
!rendered.contains(term),
"presence output leaked '{term}': {rendered}"
);
}
}
#[test]
fn roster_status_lines_leak_no_transport_vocabulary() {
let roster = RosterStatus {
org_id: "acme".into(),
serial: 42,
state: "approved".into(),
org_root_fingerprint: "tango-fig-cabbage-anchor".into(),
};
let rendered = roster_status_lines(&roster, true).join("\n");
for term in [
"b64u:",
"EndpointId",
"endpoint",
"ALPN",
"ticket",
"roster.json",
] {
assert!(
!rendered.contains(term),
"roster status output leaked '{term}': {rendered}"
);
}
}
#[test]
fn org_invite_carries_and_round_trips_the_roster_url() {
let url = "https://intranet.acme.com/roster.json";
let code = roster::enroll::OrgInviteCode {
org_id: "acme".into(),
org_root_pk: "b64u:AAAA".into(),
roster_url: Some(url.to_string()),
};
let decoded = roster::enroll::OrgInviteCode::decode(&code.encode()).unwrap();
assert_eq!(decoded.roster_url.as_deref(), Some(url));
assert_eq!(decoded.org_id, "acme");
let bare = roster::enroll::OrgInviteCode {
org_id: "acme".into(),
org_root_pk: "b64u:AAAA".into(),
roster_url: None,
};
assert!(
roster::enroll::OrgInviteCode::decode(&bare.encode())
.unwrap()
.roster_url
.is_none()
);
}
#[test]
fn recent_pairing_lines_render_petname_code_and_age() {
let pairings = vec![
RecentPairing {
peer_petname: "bob".into(),
sas_code: "tango-fig-cabbage".into(),
paired_at_epoch: 1_000_000,
},
RecentPairing {
peer_petname: "carol".into(),
sas_code: "anchor-bean-cable".into(),
paired_at_epoch: 1_000_000 - 5 * 60,
},
];
let lines = recent_pairing_lines(&pairings, 1_000_010);
assert_eq!(lines[0], " bob · code: tango-fig-cabbage · just now");
assert_eq!(lines[1], " carol · code: anchor-bean-cable · 5m ago");
}
#[test]
fn recent_pairing_lines_leak_no_endpoint_id_or_transport_vocabulary() {
let bob_id = iroh::SecretKey::from_bytes(&[7u8; 32]).public().to_string();
let pairings = vec![RecentPairing {
peer_petname: "bob".into(),
sas_code: "tango-fig-cabbage".into(),
paired_at_epoch: 100,
}];
let rendered = recent_pairing_lines(&pairings, 200).join("\n");
assert!(
!rendered.contains(&bob_id),
"recent pairings must not leak an EndpointId: {rendered}"
);
for term in [
"endpoint",
"ticket",
"ALPN",
"iroh",
"pubkey",
"mcpmesh/pair/1",
] {
assert!(
!rendered.to_lowercase().contains(&term.to_lowercase()),
"recent pairings leaked '{term}': {rendered}"
);
}
}
#[test]
fn friendly_age_buckets_and_saturates_sensibly() {
assert_eq!(friendly_age(1_000, 1_030), "just now"); assert_eq!(friendly_age(1_000, 1_000 + 5 * 60), "5m ago");
assert_eq!(friendly_age(1_000, 1_000 + 3 * 3600), "3h ago");
assert_eq!(friendly_age(1_000, 1_000 + 2 * 24 * 3600), "2d ago");
assert_eq!(friendly_age(2_000, 1_000), "just now");
}
#[test]
fn friendly_expiry_rounds_and_degrades_sensibly() {
assert_eq!(friendly_expiry(1_000 + DAY, 1_002), "in 24h");
assert_eq!(friendly_expiry(3 * 3600, 0), "in 3h");
assert_eq!(friendly_expiry(45 * 60, 0), "in 45m");
assert_eq!(friendly_expiry(30, 0), "soon");
assert_eq!(friendly_expiry(100, 1_000), "soon");
}
#[test]
fn render_frame_summarizes_a_snapshot() {
let v = json!({
"type": "snapshot",
"active_sessions": [
{"peer": "bob", "service": "notes", "opened_at": 1},
{"peer": "carol", "service": "kb", "opened_at": 2},
],
"reachability": [{"name": "bob", "reachable": true}],
});
assert_eq!(
render_frame(&v).as_deref(),
Some("snapshot: 2 active session(s), 1 peer(s) known")
);
let empty = json!({ "type": "snapshot", "active_sessions": [], "reachability": [] });
assert_eq!(
render_frame(&empty).as_deref(),
Some("snapshot: 0 active session(s), 0 peer(s) known")
);
}
#[test]
fn render_frame_renders_an_event_line() {
let v = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "session_open",
"peer": "bob", "service": "notes" },
});
assert_eq!(
render_frame(&v).as_deref(),
Some("[2026-07-17T14:02:11.480Z] session_open bob → notes")
);
}
#[test]
fn render_frame_marks_a_failed_dial_with_its_status() {
let failed = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "session_open",
"peer": "bob", "service": "notes", "status": "error" },
});
assert_eq!(
render_frame(&failed).as_deref(),
Some("[2026-07-17T14:02:11.480Z] session_open bob → notes (error)")
);
let normal = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "session_open",
"peer": "bob", "service": "notes" },
});
assert!(!render_frame(&normal).unwrap().contains('('));
}
#[test]
fn render_frame_tolerates_a_bare_event_record() {
let v = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "trust", "event": "unpair" },
});
assert_eq!(
render_frame(&v).as_deref(),
Some("[2026-07-17T14:02:11.480Z] trust")
);
}
#[test]
fn render_frame_tolerates_asymmetric_event_records() {
let peer_only = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "blob_fetch", "peer": "bob" },
});
assert_eq!(
render_frame(&peer_only).as_deref(),
Some("[2026-07-17T14:02:11.480Z] blob_fetch bob")
);
let service_only = json!({
"type": "event",
"record": { "ts": "2026-07-17T14:02:11.480Z", "kind": "session_open", "service": "notes" },
});
assert_eq!(
render_frame(&service_only).as_deref(),
Some("[2026-07-17T14:02:11.480Z] session_open → notes")
);
}
#[test]
fn render_frame_renders_a_lagged_notice() {
let v = json!({ "type": "lagged", "dropped": 7 });
assert_eq!(
render_frame(&v).as_deref(),
Some("(lagged 7 events — reconnect for a fresh snapshot)")
);
}
#[test]
fn render_frame_skips_an_unknown_frame() {
assert_eq!(render_frame(&json!({ "type": "future_kind" })), None);
assert_eq!(render_frame(&json!({ "not_a_frame": true })), None);
}
}