use anyhow::{Context, Result};
use mcpmesh_net::{EndpointId, PeerIdentity, TrustGate};
use redb::{Database, ReadableTable, TableDefinition};
use std::path::Path;
use std::sync::Arc;
const PEERS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("peers");
const REVOKED: TableDefinition<&[u8], &[u8]> = TableDefinition::new("revoked");
const REVOKED_USERS: TableDefinition<&str, &[u8]> = TableDefinition::new("revoked_users");
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct RevokedEntry {
pub endpoint_id: [u8; 32],
pub revoked_at: u64,
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub source: String,
#[serde(default)]
pub signer_user_id: Option<String>,
#[serde(default)]
pub issued_at: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct PeerEntry {
pub endpoint_id: [u8; 32],
pub nickname: String,
pub services: Vec<String>,
#[serde(default)]
pub paired_at: Option<String>,
#[serde(default)]
pub user_id: Option<String>,
#[serde(default)]
pub last_addr: Option<String>,
}
pub struct PeerStore {
db: Database,
path: std::path::PathBuf,
}
impl PeerStore {
pub fn open(path: &Path) -> Result<Self> {
let db = Database::create(path)
.with_context(|| format!("open peer store {}", path.display()))?;
let txn = db.begin_write()?;
txn.open_table(PEERS)?;
txn.open_table(REVOKED)?;
txn.open_table(REVOKED_USERS)?;
txn.commit()?;
Ok(Self {
db,
path: path.to_path_buf(),
})
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn is_revoked(&self, endpoint_id: &[u8; 32]) -> bool {
match self.revoked_entry(endpoint_id) {
Ok(v) => v.is_some(),
Err(e) => {
tracing::warn!(
%e,
"revocation lookup failed; treating this endpoint as REVOKED (fail-closed)"
);
true
}
}
}
pub fn revoked_entry(&self, endpoint_id: &[u8; 32]) -> Result<Option<RevokedEntry>> {
let txn = self.db.begin_read()?;
let table = txn.open_table(REVOKED)?;
let Some(v) = table.get(endpoint_id.as_slice())? else {
return Ok(None);
};
Ok(Some(serde_json::from_slice(v.value()).unwrap_or_else(|e| {
tracing::warn!(%e, "unreadable revocation row; still treating the endpoint as revoked");
RevokedEntry {
endpoint_id: *endpoint_id,
revoked_at: 0,
reason: Some("unreadable revocation row".into()),
source: "unknown".into(),
signer_user_id: None,
issued_at: None,
}
})))
}
pub fn is_user_revoked(&self, user_id: &str) -> bool {
let read = || -> Result<bool> {
let txn = self.db.begin_read()?;
let table = txn.open_table(REVOKED_USERS)?;
Ok(table.get(user_id)?.is_some())
};
match read() {
Ok(v) => v,
Err(e) => {
tracing::warn!(
%e,
"identity-revocation lookup failed; treating as REVOKED (fail-closed)"
);
true
}
}
}
pub fn revoke_user(&self, user_id: &str, e: &RevokedEntry) -> Result<()> {
let bytes = serde_json::to_vec(e)?;
let txn = self.db.begin_write()?;
{
let mut table = txn.open_table(REVOKED_USERS)?;
table.insert(user_id, bytes.as_slice())?;
}
txn.commit()?;
Ok(())
}
pub fn unrevoke_user(&self, user_id: &str) -> Result<bool> {
let txn = self.db.begin_write()?;
let removed = {
let mut table = txn.open_table(REVOKED_USERS)?;
table.remove(user_id)?.is_some()
};
txn.commit()?;
Ok(removed)
}
pub fn list_revoked_users(&self) -> Result<Vec<(String, RevokedEntry)>> {
let txn = self.db.begin_read()?;
let table = txn.open_table(REVOKED_USERS)?;
let mut out = Vec::new();
for row in table.iter()? {
let (k, v) = row?;
let uid = k.value().to_string();
out.push((
uid.clone(),
serde_json::from_slice(v.value()).unwrap_or(RevokedEntry {
endpoint_id: [0u8; 32],
revoked_at: 0,
reason: Some("unreadable revocation row".into()),
source: "unknown".into(),
signer_user_id: None,
issued_at: None,
}),
));
}
Ok(out)
}
pub fn revoke(&self, e: RevokedEntry) -> Result<()> {
let bytes = serde_json::to_vec(&e)?;
let txn = self.db.begin_write()?;
{
let mut table = txn.open_table(REVOKED)?;
table.insert(e.endpoint_id.as_slice(), bytes.as_slice())?;
}
txn.commit()?;
Ok(())
}
pub fn unrevoke(&self, endpoint_id: &[u8; 32]) -> Result<bool> {
let txn = self.db.begin_write()?;
let removed = {
let mut table = txn.open_table(REVOKED)?;
table.remove(endpoint_id.as_slice())?.is_some()
};
txn.commit()?;
Ok(removed)
}
pub fn list_revoked(&self) -> Result<Vec<RevokedEntry>> {
let txn = self.db.begin_read()?;
let table = txn.open_table(REVOKED)?;
let mut out = Vec::new();
for row in table.iter()? {
let (k, v) = row?;
let mut eid = [0u8; 32];
if k.value().len() == 32 {
eid.copy_from_slice(k.value());
}
out.push(serde_json::from_slice(v.value()).unwrap_or(RevokedEntry {
endpoint_id: eid,
revoked_at: 0,
reason: Some("unreadable revocation row".into()),
source: "unknown".into(),
signer_user_id: None,
issued_at: None,
}));
}
Ok(out)
}
pub fn add(&self, e: PeerEntry) -> Result<()> {
let bytes = serde_json::to_vec(&e)?;
let txn = self.db.begin_write()?;
{
let mut table = txn.open_table(PEERS)?;
table.insert(e.endpoint_id.as_slice(), bytes.as_slice())?;
}
txn.commit()?;
Ok(())
}
pub fn set_last_addr(&self, endpoint_id: &[u8; 32], last_addr: &str) -> Result<bool> {
let txn = self.db.begin_write()?;
let changed = {
let mut table = txn.open_table(PEERS)?;
let Some(existing) = table.get(endpoint_id.as_slice())? else {
return Ok(false); };
let mut entry: PeerEntry = serde_json::from_slice(existing.value())?;
drop(existing);
if entry.last_addr.as_deref() == Some(last_addr) {
return Ok(false);
} else {
entry.last_addr = Some(last_addr.to_string());
let bytes = serde_json::to_vec(&entry)?;
table.insert(endpoint_id.as_slice(), bytes.as_slice())?;
true
}
};
txn.commit()?;
Ok(changed)
}
pub fn resolve(&self, endpoint_id: &[u8; 32]) -> Result<Option<PeerEntry>> {
let txn = self.db.begin_read()?;
let table = txn.open_table(PEERS)?;
match table.get(endpoint_id.as_slice())? {
Some(v) => match serde_json::from_slice::<PeerEntry>(v.value()) {
Ok(entry) => Ok(Some(entry)),
Err(e) => {
tracing::warn!(
key_prefix = ?&endpoint_id[..8],
error = %e,
"corrupt peer entry for queried key; treating as unresolved (deny)"
);
Ok(None)
}
},
None => Ok(None),
}
}
pub fn entry_for(&self, nickname: &str) -> Result<Option<PeerEntry>> {
Ok(self.list()?.into_iter().find(|e| e.nickname == nickname))
}
pub fn entries_for_user(&self, user_id: &str) -> Result<Vec<PeerEntry>> {
Ok(self
.list()?
.into_iter()
.filter(|e| e.user_id.as_deref() == Some(user_id))
.collect())
}
pub fn list(&self) -> Result<Vec<PeerEntry>> {
let txn = self.db.begin_read()?;
let table = txn.open_table(PEERS)?;
let mut out = Vec::new();
for row in table.iter()? {
let (k, v) = row?;
match serde_json::from_slice::<PeerEntry>(v.value()) {
Ok(entry) => out.push(entry),
Err(e) => {
let kb = k.value();
tracing::warn!(
key_prefix = ?&kb[..kb.len().min(8)],
error = %e,
"skipping corrupt peer entry during list"
);
}
}
}
Ok(out)
}
pub fn remove(&self, nickname: &str) -> Result<bool> {
let txn = self.db.begin_write()?;
let removed = {
let mut table = txn.open_table(PEERS)?;
let victims: Vec<Vec<u8>> = {
let mut v = Vec::new();
for row in table.iter()? {
let (k, val) = row?;
match serde_json::from_slice::<PeerEntry>(val.value()) {
Ok(entry) if entry.nickname == nickname => v.push(k.value().to_vec()),
Ok(_) => {}
Err(e) => {
let kb = k.value();
tracing::warn!(
key_prefix = ?&kb[..kb.len().min(8)],
error = %e,
"skipping corrupt peer entry during remove"
);
}
}
}
v
};
for k in &victims {
table.remove(k.as_slice())?;
}
!victims.is_empty()
};
txn.commit()?;
Ok(removed)
}
}
pub struct AllowlistGate {
store: Arc<PeerStore>,
}
impl AllowlistGate {
pub fn new(store: Arc<PeerStore>) -> Self {
Self { store }
}
}
impl TrustGate for AllowlistGate {
fn resolve(&self, endpoint: &EndpointId) -> Option<PeerIdentity> {
if self.store.is_revoked(endpoint.as_bytes()) {
return None;
}
match self.store.resolve(endpoint.as_bytes()) {
Ok(Some(e)) => Some(PeerIdentity {
endpoint: *endpoint,
user_id: e.user_id, name: e.nickname,
groups: vec![],
}),
Ok(None) => None,
Err(e) => {
tracing::warn!(%e, "peer store read failed; refusing (default-deny)");
None
}
}
}
fn is_revoked(&self, endpoint: &EndpointId) -> bool {
self.store.is_revoked(endpoint.as_bytes())
}
fn should_sever_now(&self, endpoint: &EndpointId, _roster_user: Option<&str>) -> bool {
self.store.is_revoked(endpoint.as_bytes())
}
}
#[cfg(test)]
mod tests {
#[test]
fn an_unchanged_set_last_addr_aborts_instead_of_committing() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("p.redb")).unwrap();
let eid = [5u8; 32];
let addr = r#"{"id":"x","addrs":[]}"#;
store
.add(PeerEntry {
endpoint_id: eid,
nickname: "bob".into(),
services: vec![],
paired_at: None,
user_id: None,
last_addr: Some(addr.to_string()),
})
.unwrap();
let started = std::time::Instant::now();
for _ in 0..1000 {
assert!(
!store.set_last_addr(&eid, addr).unwrap(),
"an unchanged hint must report no write"
);
}
let elapsed = started.elapsed();
assert!(
elapsed < std::time::Duration::from_secs(2),
"1000 unchanged refreshes took {elapsed:?} — committing an empty txn per call is a \
~6ms fsync holding redb's GLOBAL writer lock, which blocks pairing and peer add on \
every path event (#124)"
);
assert_eq!(
store.resolve(&eid).unwrap().unwrap().last_addr.as_deref(),
Some(addr)
);
assert!(
store
.set_last_addr(&eid, r#"{"id":"y","addrs":[]}"#)
.unwrap()
);
assert!(!store.set_last_addr(&[7u8; 32], addr).unwrap());
assert!(store.resolve(&[7u8; 32]).unwrap().is_none());
}
use super::*;
fn entry(eid: [u8; 32], nickname: &str, services: &[&str]) -> PeerEntry {
PeerEntry {
endpoint_id: eid,
nickname: nickname.into(),
services: services.iter().map(|s| s.to_string()).collect(),
paired_at: None,
user_id: None,
last_addr: None,
}
}
fn inject_raw(store: &PeerStore, eid: &[u8; 32], bytes: &[u8]) {
let txn = store.db.begin_write().unwrap();
{
let mut table = txn.open_table(PEERS).unwrap();
table.insert(eid.as_slice(), bytes).unwrap();
}
txn.commit().unwrap();
}
#[test]
fn gate_resolves_known_nickname_refuses_unknown() {
use mcpmesh_net::TrustGate;
use std::sync::Arc;
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let known_eid = [7u8; 32];
store.add(entry(known_eid, "bob", &["notes"])).unwrap();
let gate = AllowlistGate::new(Arc::new(store));
let id = gate.resolve(&known_eid.into()).unwrap();
assert_eq!(id.name, "bob");
assert_eq!(id.user_id, None);
assert!(id.groups.is_empty());
assert!(gate.resolve(&[9u8; 32].into()).is_none());
}
#[test]
fn add_then_resolve_and_list() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [7u8; 32];
store.add(entry(eid, "bob", &["notes"])).unwrap();
assert_eq!(store.resolve(&eid).unwrap().unwrap().nickname, "bob");
assert!(store.resolve(&[9u8; 32]).unwrap().is_none());
assert_eq!(store.list().unwrap().len(), 1);
}
#[test]
fn entry_persists_across_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("state.redb");
let eid = [42u8; 32];
{
let store = PeerStore::open(&path).unwrap();
store.add(entry(eid, "alice", &["notes", "kb"])).unwrap();
} let store = PeerStore::open(&path).unwrap();
let got = store.resolve(&eid).unwrap().unwrap();
assert_eq!(got.nickname, "alice");
assert_eq!(got.services, vec!["notes".to_string(), "kb".to_string()]);
}
#[test]
fn add_upserts_same_endpoint_id() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [1u8; 32];
store.add(entry(eid, "bob", &["notes"])).unwrap();
store.add(entry(eid, "bob-renamed", &["kb"])).unwrap();
let all = store.list().unwrap();
assert_eq!(all.len(), 1);
assert_eq!(all[0].nickname, "bob-renamed");
assert_eq!(all[0].services, vec!["kb".to_string()]);
}
#[test]
fn remove_deletes_match_and_is_a_noop_for_absent() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [3u8; 32];
store.add(entry(eid, "carol", &[])).unwrap();
assert!(
!store.remove("nobody").unwrap(),
"removing an absent nickname removes nothing"
);
assert!(store.resolve(&eid).unwrap().is_some());
assert!(
store.remove("carol").unwrap(),
"removing a present nickname reports the deletion"
);
assert!(store.resolve(&eid).unwrap().is_none());
}
#[test]
fn remove_persists_across_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("state.redb");
let eid = [5u8; 32];
{
let store = PeerStore::open(&path).unwrap();
store.add(entry(eid, "dave", &[])).unwrap();
store.remove("dave").unwrap();
}
let store = PeerStore::open(&path).unwrap();
assert!(store.resolve(&eid).unwrap().is_none());
}
#[test]
fn remove_deletes_all_entries_sharing_a_nickname() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
store.add(entry([10u8; 32], "dup", &[])).unwrap();
store.add(entry([11u8; 32], "dup", &[])).unwrap();
assert_eq!(store.list().unwrap().len(), 2);
store.remove("dup").unwrap();
assert_eq!(store.list().unwrap().len(), 0);
}
#[test]
fn old_row_without_paired_at_still_resolves_defaulting_to_none() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [7u8; 32];
let old_shape = serde_json::json!({
"endpoint_id": eid.to_vec(),
"nickname": "old",
"services": ["notes"],
});
inject_raw(&store, &eid, &serde_json::to_vec(&old_shape).unwrap());
let got = store.resolve(&eid).unwrap().unwrap();
assert_eq!(got.nickname, "old");
assert_eq!(got.services, vec!["notes".to_string()]);
assert_eq!(got.paired_at, None); }
#[test]
fn paired_at_round_trips_when_set() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [8u8; 32];
let mut e = entry(eid, "bob", &["notes"]);
e.paired_at = Some("1751760000".into());
store.add(e).unwrap();
let got = store.resolve(&eid).unwrap().unwrap();
assert_eq!(got.paired_at.as_deref(), Some("1751760000"));
}
#[test]
fn old_row_without_last_addr_still_resolves_defaulting_to_none() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [9u8; 32];
let old_shape = serde_json::json!({
"endpoint_id": eid.to_vec(),
"nickname": "old",
"services": ["notes"],
"paired_at": "1751760000",
"user_id": null,
});
inject_raw(&store, &eid, &serde_json::to_vec(&old_shape).unwrap());
let got = store.resolve(&eid).unwrap().unwrap();
assert_eq!(got.nickname, "old");
assert_eq!(got.last_addr, None); }
#[test]
fn last_addr_round_trips_when_set() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [10u8; 32];
let mut e = entry(eid, "bob", &["notes"]);
e.last_addr = Some(r#"{"id":"whatever","addrs":[]}"#.into());
store.add(e).unwrap();
let got = store.resolve(&eid).unwrap().unwrap();
assert_eq!(
got.last_addr.as_deref(),
Some(r#"{"id":"whatever","addrs":[]}"#)
);
}
#[test]
fn entry_for_returns_the_full_entry() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let eid = [11u8; 32];
let mut e = entry(eid, "alice", &["echo"]);
e.last_addr = Some("{}".into());
store.add(e).unwrap();
let got = store.entry_for("alice").unwrap().unwrap();
assert_eq!(got.endpoint_id, eid);
assert_eq!(got.last_addr.as_deref(), Some("{}"));
assert!(store.entry_for("nobody").unwrap().is_none());
}
#[test]
fn entries_for_user_groups_a_persons_devices() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let mut laptop = entry([1u8; 32], "alice", &["notes"]);
laptop.user_id = Some("b64u:ALICE".into());
let mut phone = entry([2u8; 32], "alice-phone", &["notes"]);
phone.user_id = Some("b64u:ALICE".into());
let mut bob = entry([3u8; 32], "bob", &["kb"]);
bob.user_id = Some("b64u:BOB".into());
let legacy = entry([4u8; 32], "carol", &["x"]); for e in [laptop, phone, bob, legacy] {
store.add(e).unwrap();
}
let alice = store.entries_for_user("b64u:ALICE").unwrap();
assert_eq!(alice.len(), 2, "both of alice's devices match her user_id");
let mut eids: Vec<_> = alice.iter().map(|e| e.endpoint_id).collect();
eids.sort();
assert_eq!(eids, vec![[1u8; 32], [2u8; 32]]);
assert_eq!(store.entries_for_user("b64u:BOB").unwrap().len(), 1);
assert!(store.entries_for_user("b64u:NOBODY").unwrap().is_empty());
}
#[test]
fn corrupt_row_is_skipped_on_list_and_denied_on_resolve() {
let dir = tempfile::tempdir().unwrap();
let store = PeerStore::open(&dir.path().join("state.redb")).unwrap();
let good = [1u8; 32];
let bad = [2u8; 32];
store.add(entry(good, "good", &["notes"])).unwrap();
inject_raw(&store, &bad, b"not json at all");
let all = store.list().unwrap();
assert_eq!(all.len(), 1);
assert_eq!(all[0].nickname, "good");
assert!(store.resolve(&bad).unwrap().is_none());
assert_eq!(store.resolve(&good).unwrap().unwrap().nickname, "good");
store.remove("good").unwrap();
assert!(store.resolve(&good).unwrap().is_none());
}
#[test]
fn a_revoked_endpoint_is_refused_even_with_a_live_pair_row() {
use mcpmesh_net::TrustGate;
let dir = tempfile::tempdir().unwrap();
let store = Arc::new(PeerStore::open(&dir.path().join("p.redb")).unwrap());
let eid = [5u8; 32];
store
.add(PeerEntry {
endpoint_id: eid,
nickname: "bob".into(),
services: vec!["notes".into()],
paired_at: None,
user_id: Some("b64u:BOB".into()),
last_addr: None,
})
.unwrap();
let gate = AllowlistGate::new(store.clone());
let id: EndpointId = eid.into();
assert!(
gate.resolve(&id).is_some(),
"precondition: the peer resolves before revocation"
);
assert!(!gate.is_revoked(&id));
assert!(!gate.should_sever_now(&id, None));
store
.revoke(RevokedEntry {
endpoint_id: eid,
revoked_at: 1_754_300_000,
reason: Some("laptop stolen".into()),
source: "local".into(),
signer_user_id: None,
issued_at: None,
})
.unwrap();
assert!(
gate.resolve(&id).is_none(),
"a revoked endpoint must not resolve, even though its pair row is untouched"
);
assert!(
gate.is_revoked(&id),
"the check-register recheck must see it — that is the TOCTOU close (#54)"
);
assert!(
gate.should_sever_now(&id, None),
"and a LIVE session must be severed, not left to end on its own"
);
assert!(
store.resolve(&eid).unwrap().is_some(),
"revocation must not delete the pair row — unrevoking has to restore the peer"
);
assert!(store.unrevoke(&eid).unwrap(), "the revocation was present");
assert!(
gate.resolve(&id).is_some(),
"unrevoking restores the peer, since the pair row survived"
);
assert!(!store.unrevoke(&eid).unwrap(), "idempotent");
}
#[test]
fn a_revocation_survives_a_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("p.redb");
let eid = [9u8; 32];
{
let store = PeerStore::open(&path).unwrap();
store
.revoke(RevokedEntry {
endpoint_id: eid,
revoked_at: 42,
reason: Some("stolen".into()),
source: "signed".into(),
signer_user_id: Some("b64u:BOB".into()),
issued_at: Some(1000),
})
.unwrap();
}
let store = PeerStore::open(&path).unwrap();
assert!(store.is_revoked(&eid));
let e = store.revoked_entry(&eid).unwrap().expect("row survives");
assert_eq!(
(e.source.as_str(), e.signer_user_id.as_deref(), e.issued_at),
("signed", Some("b64u:BOB"), Some(1000)),
"the PROVENANCE survives too — an operator has to be able to tell a signed revocation \
from their own local one after a restart"
);
assert_eq!(store.list_revoked().unwrap().len(), 1);
}
#[test]
fn an_unreadable_revocation_row_still_revokes() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("p.redb");
let eid = [4u8; 32];
{
let db = Database::create(&path).unwrap();
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(REVOKED).unwrap();
t.insert(eid.as_slice(), b"{not json".as_slice()).unwrap();
}
txn.commit().unwrap();
}
let store = PeerStore::open(&path).unwrap();
assert!(
store.is_revoked(&eid),
"an unreadable revocation row must still REVOKE — failing open here undoes the only \
remedy for a stolen device"
);
let listed = store.list_revoked().unwrap();
assert_eq!(listed.len(), 1, "and it must still be VISIBLE in status");
assert_eq!(
listed[0].source, "unknown",
"…rendered as unknown-provenance rather than dropped, so the list cannot disagree with \
the gate"
);
}
}