use std::sync::Arc;
use ed25519_dalek::{Signature, SigningKey, VerifyingKey};
use mcpmesh_trust::roster::validate::RosterView;
use serde::{Deserialize, Serialize};
use crate::roster::gate::RosterGate;
use crate::roster::transport::{self, RosterGossip};
const PRESENCE_DOMAIN: &[u8] = b"mcpmesh/presence/1";
pub const APP_METADATA_MAX_BYTES: usize = 256;
pub const PRESENCE_SKEW_SECS: i64 = 120;
pub const PRESENCE_TTL_SECS: i64 = 180;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Presence {
pub t: String, pub endpoint_id: String, pub user_id: String,
pub ts: i64, pub roster_serial: u64, #[serde(default, skip_serializing_if = "String::is_empty")]
pub meta: String,
pub sig: String, }
fn preimage(endpoint_id: &[u8; 32], user_id: &str, ts: i64, serial: u64, meta: &str) -> Vec<u8> {
let mut m = Vec::with_capacity(PRESENCE_DOMAIN.len() + 32 + user_id.len() + 16 + meta.len());
m.extend_from_slice(PRESENCE_DOMAIN);
m.extend_from_slice(endpoint_id);
m.extend_from_slice(user_id.as_bytes());
m.extend_from_slice(&ts.to_le_bytes());
m.extend_from_slice(&serial.to_le_bytes());
if !meta.is_empty() {
m.extend_from_slice(meta.as_bytes());
}
m
}
impl Presence {
pub fn signed(
device_key: &SigningKey,
endpoint_id: &[u8; 32],
user_id: &str,
ts: i64,
serial: u64,
meta: &str,
) -> Self {
use ed25519_dalek::Signer;
let sig = device_key.sign(&preimage(endpoint_id, user_id, ts, serial, meta));
Self {
t: "presence".into(),
endpoint_id: mcpmesh_trust::roster::encode_b64u(endpoint_id),
user_id: user_id.to_string(),
ts,
roster_serial: serial,
meta: meta.to_string(),
sig: mcpmesh_trust::roster::encode_b64u(&sig.to_bytes()),
}
}
pub fn verify(&self, now: i64) -> Option<[u8; 32]> {
if self.t != "presence"
|| now.saturating_sub(self.ts).saturating_abs() >= PRESENCE_SKEW_SECS
|| self.meta.len() > APP_METADATA_MAX_BYTES
{
return None;
}
let eid = mcpmesh_trust::roster::decode_endpoint_id(&self.endpoint_id).ok()?;
let vk = VerifyingKey::from_bytes(&eid).ok()?;
let sig =
Signature::from_slice(&mcpmesh_trust::roster::decode_b64u(&self.sig).ok()?).ok()?;
vk.verify_strict(
&preimage(&eid, &self.user_id, self.ts, self.roster_serial, &self.meta),
&sig,
)
.ok()?;
Some(eid)
}
pub fn accept(&self, now: i64, view: &RosterView) -> Option<[u8; 32]> {
let eid = self.verify(now)?;
let resolved = view.resolve(&eid)?; if resolved.user_id == self.user_id {
Some(eid)
} else {
None
} }
pub fn to_bytes(&self) -> Vec<u8> {
serde_json::to_vec(self).expect("Presence serializes")
}
pub fn from_bytes(b: &[u8]) -> Option<Self> {
serde_json::from_slice(b).ok()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PresenceEntry {
pub user_id: String,
pub ts: i64,
pub meta: String,
}
#[derive(Debug, Default)]
pub struct PresenceTable {
inner: std::sync::Mutex<std::collections::HashMap<[u8; 32], PresenceEntry>>,
}
impl PresenceTable {
pub fn new() -> Self {
Self::default()
}
pub fn record(&self, eid: [u8; 32], user_id: String, ts: i64, meta: String) {
let mut g = self.inner.lock().expect("presence table mutex poisoned");
match g.get(&eid) {
Some(e) if e.ts >= ts => {} _ => {
g.insert(eid, PresenceEntry { user_id, ts, meta });
}
}
}
pub fn active(&self, now: i64) -> Vec<([u8; 32], PresenceEntry)> {
let g = self.inner.lock().expect("presence table mutex poisoned");
g.iter()
.filter(|(_, e)| now - e.ts < PRESENCE_TTL_SECS)
.map(|(eid, e)| (*eid, e.clone()))
.collect()
}
pub fn endpoints_for_user_by_recency(&self, user_id: &str) -> Vec<[u8; 32]> {
let g = self.inner.lock().expect("presence table mutex poisoned");
let mut hits: Vec<(&[u8; 32], &PresenceEntry)> =
g.iter().filter(|(_, e)| e.user_id == user_id).collect();
hits.sort_by_key(|(_, e)| std::cmp::Reverse(e.ts)); hits.into_iter().map(|(eid, _)| *eid).collect()
}
}
pub struct PresenceCtx {
pub roster: Arc<RosterGate>,
pub table: Arc<PresenceTable>,
pub topic: Arc<tokio::sync::Mutex<Option<RosterGossip>>>,
pub app_metadata: Arc<std::sync::RwLock<String>>,
}
impl PresenceCtx {
async fn sender(&self) -> Option<iroh_gossip::api::GossipSender> {
self.topic.lock().await.as_ref().map(|g| g.sender.clone())
}
async fn take_receiver(&self) -> Option<iroh_gossip::api::GossipReceiver> {
self.topic
.lock()
.await
.as_mut()
.and_then(|g| g.receiver.take())
}
}
pub fn publish_loop(
ctx: PresenceCtx,
device_key: SigningKey,
user_id: String,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let endpoint_id = device_key.verifying_key().to_bytes();
let Some(sender) = ctx.sender().await else {
return; };
loop {
let now = crate::util::epoch_now_i64();
let serial = ctx.roster.view().map(|v| v.serial()).unwrap_or(0);
let meta = ctx
.app_metadata
.read()
.expect("app_metadata lock not poisoned")
.clone();
let beat = Presence::signed(&device_key, &endpoint_id, &user_id, now, serial, &meta);
if let Err(e) = transport::broadcast(&sender, beat.to_bytes()).await {
tracing::debug!(%e, "presence heartbeat broadcast failed; will retry next beat");
}
tokio::time::sleep(std::time::Duration::from_secs(next_beat_secs())).await;
}
})
}
pub fn track_loop(ctx: PresenceCtx) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let Some(mut receiver) = ctx.take_receiver().await else {
return;
};
while let Some(content) = transport::next_message(&mut receiver).await {
let Some(p) = Presence::from_bytes(&content) else {
tracing::trace!("malformed presence payload dropped");
continue;
};
let now = crate::util::epoch_now_i64();
let Some(view) = ctx.roster.view() else {
continue; };
if let Some(eid) = p.accept(now, &view) {
ctx.table.record(eid, p.user_id, p.ts, p.meta);
}
}
})
}
fn next_beat_secs() -> u64 {
use rand::RngCore;
let mut buf = [0u8; 4];
rand::rngs::OsRng.fill_bytes(&mut buf);
let span = (2 * PRESENCE_JITTER_SECS + 1) as u32; let jitter = (u32::from_le_bytes(buf) % span) as i64 - PRESENCE_JITTER_SECS;
(PRESENCE_PERIOD_SECS + jitter).max(1) as u64
}
const PRESENCE_PERIOD_SECS: i64 = 60;
const PRESENCE_JITTER_SECS: i64 = 15;
#[cfg(test)]
mod tests {
use super::*;
use ed25519_dalek::SigningKey;
fn legacy_preimage(endpoint_id: &[u8; 32], user_id: &str, ts: i64, serial: u64) -> Vec<u8> {
let mut m = Vec::new();
m.extend_from_slice(PRESENCE_DOMAIN);
m.extend_from_slice(endpoint_id);
m.extend_from_slice(user_id.as_bytes());
m.extend_from_slice(&ts.to_le_bytes());
m.extend_from_slice(&serial.to_le_bytes());
m
}
fn decode_b64u_sig(s: &str) -> Signature {
Signature::from_slice(&mcpmesh_trust::roster::decode_b64u(s).unwrap()).unwrap()
}
use mcpmesh_trust::roster::sign::mint_signed;
use mcpmesh_trust::roster::validate::load_installed;
use mcpmesh_trust::roster::{Roster, RosterDevice, RosterUser, encode_b64u};
#[test]
fn presence_sign_verify_round_trips_and_forgeries_fail() {
let dk = SigningKey::from_bytes(&[4u8; 32]);
let eid = dk.verifying_key().to_bytes();
let now = 1_760_000_000;
let p = Presence::signed(&dk, &eid, "alice", now, 7, "");
assert_eq!(p.verify(now + 5), Some(eid));
let sig_no_meta = decode_b64u_sig(&p.sig);
let dk_v = dk.verifying_key();
use ed25519_dalek::Verifier;
assert!(
dk_v.verify(&legacy_preimage(&eid, "alice", now, 7), &sig_no_meta)
.is_ok(),
"empty-meta signed bytes must equal the pre-#39 preimage"
);
let pm = Presence::signed(&dk, &eid, "alice", now, 7, "v=1.2.3");
assert_eq!(pm.verify(now + 5), Some(eid));
assert_eq!(pm.meta, "v=1.2.3");
let mut tampered = pm.clone();
tampered.meta = "v=9.9.9".into();
assert!(
tampered.verify(now).is_none(),
"meta is signed — a swap fails"
);
assert!(
p.verify(now + PRESENCE_SKEW_SECS).is_none(),
"outside ±120s rejected"
);
let mut bad = p.clone();
bad.user_id = "bob".into();
assert!(bad.verify(now).is_none()); let mut swapped = p.clone();
swapped.endpoint_id = encode_b64u(&[9u8; 32]);
assert!(swapped.verify(now).is_none()); assert!(serde_json::to_vec(&p).unwrap().len() < 768);
let full = Presence::signed(&dk, &eid, "alice", now, 7, &"x".repeat(256));
assert!(serde_json::to_vec(&full).unwrap().len() < 768);
}
#[test]
fn verify_does_not_overflow_on_a_crafted_extreme_ts() {
let dk = SigningKey::from_bytes(&[4u8; 32]);
let eid = dk.verifying_key().to_bytes();
let mut p = Presence::signed(&dk, &eid, "alice", 1_760_000_000, 7, "");
p.ts = i64::MIN;
assert!(p.verify(1_760_000_000).is_none());
p.ts = i64::MAX;
assert!(p.verify(i64::MIN).is_none());
}
#[test]
fn table_surfaces_freshest_meta_and_does_not_regress() {
let table = PresenceTable::new();
let eid = [3u8; 32];
table.record(eid, "b64u:A".into(), 1000, "v=1.0.0".into());
let got = table
.active(1000)
.into_iter()
.find(|(e, _)| *e == eid)
.unwrap()
.1;
assert_eq!(got.meta, "v=1.0.0");
table.record(eid, "b64u:A".into(), 1050, "v=1.1.0".into());
let got = table
.active(1050)
.into_iter()
.find(|(e, _)| *e == eid)
.unwrap()
.1;
assert_eq!(got.meta, "v=1.1.0");
table.record(eid, "b64u:A".into(), 1010, "v=0.9.0".into());
let got = table
.active(1050)
.into_iter()
.find(|(e, _)| *e == eid)
.unwrap()
.1;
assert_eq!(got.meta, "v=1.1.0", "older beat must not regress meta");
}
#[test]
fn receive_path_records_capped_meta_and_drops_oversized() {
let dk = SigningKey::from_bytes(&[5u8; 32]);
let eid = dk.verifying_key().to_bytes();
let root = SigningKey::from_bytes(&[9u8; 32]);
let signed = mint_signed(
&root,
Roster {
format: "mcpmesh-roster/1".into(),
org_id: "acme".into(),
serial: 1,
issued_at: "2000-01-01T00:00:00Z".into(),
expires_at: "2999-01-01T00:00:00Z".into(),
groups: vec![],
users: vec![RosterUser {
user_id: "alice".into(),
display_name: "Alice".into(),
user_pk: encode_b64u(&[1u8; 32]),
groups: vec![],
devices: vec![RosterDevice {
endpoint_id: encode_b64u(&eid),
label: "laptop".into(),
role: "primary".into(),
}],
}],
revoked_endpoints: vec![],
sig: String::new(),
},
);
let view = load_installed(&signed, &root.verifying_key()).unwrap();
let now = 1_760_000_000;
let table = PresenceTable::new();
let good = Presence::signed(&dk, &eid, "alice", now, 1, "v=1.2.3");
let got = good
.accept(now, &view)
.expect("valid rostered beat accepts");
assert_eq!(got, eid);
table.record(got, "alice".into(), good.ts, good.meta.clone());
assert_eq!(
table
.active(now)
.into_iter()
.find(|(e, _)| *e == eid)
.unwrap()
.1
.meta,
"v=1.2.3"
);
let big = Presence::signed(&dk, &eid, "alice", now, 1, &"x".repeat(257));
assert!(
big.verify(now).is_none(),
"oversized meta must drop the beat"
);
assert!(big.accept(now, &view).is_none());
}
#[test]
fn accept_binds_user_id_to_the_roster_authoritative_user() {
let root = SigningKey::from_bytes(&[9u8; 32]);
let dk = SigningKey::from_bytes(&[4u8; 32]);
let eid = dk.verifying_key().to_bytes();
let roster = mint_signed(
&root,
Roster {
format: "mcpmesh-roster/1".into(),
org_id: "acme".into(),
serial: 5,
issued_at: "2000-01-01T00:00:00Z".into(),
expires_at: "2999-01-01T00:00:00Z".into(),
groups: vec!["all".into()],
users: vec![RosterUser {
user_id: "alice".into(),
display_name: "Alice".into(),
user_pk: encode_b64u(&[1u8; 32]),
groups: vec!["all".into()],
devices: vec![RosterDevice {
endpoint_id: encode_b64u(&eid),
label: "l".into(),
role: "primary".into(),
}],
}],
revoked_endpoints: vec![],
sig: String::new(),
},
);
let view = load_installed(&roster, &root.verifying_key()).unwrap();
let now = 1_760_000_000;
assert_eq!(
Presence::signed(&dk, &eid, "alice", now, 5, "").accept(now, &view),
Some(eid)
);
assert!(
Presence::signed(&dk, &eid, "bob", now, 5, "")
.accept(now, &view)
.is_none()
);
let stranger = SigningKey::from_bytes(&[7u8; 32]);
let seid = stranger.verifying_key().to_bytes();
assert!(
Presence::signed(&stranger, &seid, "alice", now, 5, "")
.accept(now, &view)
.is_none()
);
}
}