Merge max/relay-mesh-membership: registry, gossip, status

Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>

* origin/max/relay-mesh-membership:
  feat(mesh): add membership registry and gossip
This commit is contained in:
npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d
2026-07-08 11:09:14 -04:00
co-authored by Tyler Longwell
7 changed files with 1132 additions and 0 deletions
Generated
+1
View File
@@ -1104,6 +1104,7 @@ dependencies = [
"hex",
"hmac 0.13.0",
"iroh",
"nostr",
"postcard",
"proptest",
"redis",
+1
View File
@@ -21,6 +21,7 @@ uuid = { workspace = true }
hmac = { workspace = true }
sha2 = { workspace = true }
hex = { workspace = true }
nostr = { workspace = true }
bytes = "1"
futures-util = { workspace = true }
+303
View File
@@ -0,0 +1,303 @@
//! Scuttlebutt-style membership gossip over the mesh control stream.
//!
//! Gossip answers liveness/dialability questions only. It never elects owners,
//! never transfers sessions, and never carries tunnel data bytes.
use std::collections::HashMap;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use serde::{Deserialize, Serialize};
use crate::{MeshError, RuntimeId};
pub const GOSSIP_PAYLOAD_VERSION: u8 = 1;
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct GossipRecord {
pub runtime_id: RuntimeId,
pub endpoint_addrs: Vec<String>,
pub proto_version: u16,
pub load: f32,
pub draining: bool,
pub capabilities: Vec<String>,
/// Per-runtime monotonic version. Only the owning runtime may increment its
/// own record; receivers apply last-version-wins.
pub version: u64,
pub heartbeat_millis: u64,
}
impl GossipRecord {
pub fn new(runtime_id: RuntimeId, endpoint_addrs: Vec<String>, proto_version: u16) -> Self {
Self {
runtime_id,
endpoint_addrs,
proto_version,
load: 0.0,
draining: false,
capabilities: Vec::new(),
version: 1,
heartbeat_millis: now_millis(),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct GossipDigestEntry {
pub runtime_id: RuntimeId,
pub version: u64,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub enum GossipMessage {
Digest {
version: u8,
entries: Vec<GossipDigestEntry>,
},
Delta {
version: u8,
records: Vec<GossipRecord>,
},
}
pub fn encode_message(message: &GossipMessage) -> Result<Vec<u8>, MeshError> {
postcard::to_extend(message, Vec::new()).map_err(MeshError::Encode)
}
pub fn decode_message(bytes: &[u8]) -> Result<GossipMessage, MeshError> {
let message: GossipMessage = postcard::from_bytes(bytes).map_err(MeshError::Decode)?;
let version = match &message {
GossipMessage::Digest { version, .. } | GossipMessage::Delta { version, .. } => *version,
};
if version != GOSSIP_PAYLOAD_VERSION {
return Err(MeshError::Transport(format!(
"unknown gossip payload version {version}"
)));
}
Ok(message)
}
/// Pure scuttlebutt state: digest exchange + delta application.
#[derive(Clone, Debug)]
pub struct GossipState {
records: HashMap<RuntimeId, GossipRecord>,
}
impl GossipState {
pub fn new(local: GossipRecord) -> Self {
let mut records = HashMap::new();
records.insert(local.runtime_id, local);
Self { records }
}
pub fn records(&self) -> impl Iterator<Item = &GossipRecord> {
self.records.values()
}
pub fn get(&self, runtime_id: RuntimeId) -> Option<&GossipRecord> {
self.records.get(&runtime_id)
}
pub fn update_local<F>(&mut self, runtime_id: RuntimeId, update: F) -> Option<GossipRecord>
where
F: FnOnce(&mut GossipRecord),
{
let record = self.records.get_mut(&runtime_id)?;
update(record);
record.version = record.version.saturating_add(1);
record.heartbeat_millis = now_millis();
Some(record.clone())
}
pub fn digest(&self) -> GossipMessage {
let mut entries: Vec<_> = self
.records
.values()
.map(|record| GossipDigestEntry {
runtime_id: record.runtime_id,
version: record.version,
})
.collect();
entries.sort_by_key(|entry| entry.runtime_id.to_hex());
GossipMessage::Digest {
version: GOSSIP_PAYLOAD_VERSION,
entries,
}
}
pub fn delta_for(&self, digest: &[GossipDigestEntry]) -> GossipMessage {
let remote_versions: HashMap<_, _> = digest
.iter()
.map(|entry| (entry.runtime_id, entry.version))
.collect();
let mut records: Vec<_> = self
.records
.values()
.filter(|record| {
remote_versions
.get(&record.runtime_id)
.is_none_or(|remote| *remote < record.version)
})
.cloned()
.collect();
records.sort_by_key(|record| record.runtime_id.to_hex());
GossipMessage::Delta {
version: GOSSIP_PAYLOAD_VERSION,
records,
}
}
/// Applies records whose version is newer than the local copy. Returns the
/// runtime ids that changed.
pub fn apply_delta(&mut self, records: Vec<GossipRecord>) -> Vec<RuntimeId> {
let mut changed = Vec::new();
for record in records {
let should_apply = self
.records
.get(&record.runtime_id)
.is_none_or(|existing| record.version > existing.version);
if should_apply {
changed.push(record.runtime_id);
self.records.insert(record.runtime_id, record);
}
}
changed
}
}
#[derive(Clone, Debug)]
pub struct PhiAccrual {
samples: Vec<Duration>,
last_heartbeat: Option<SystemTime>,
max_samples: usize,
}
impl Default for PhiAccrual {
fn default() -> Self {
Self::new(100)
}
}
impl PhiAccrual {
pub fn new(max_samples: usize) -> Self {
Self {
samples: Vec::new(),
last_heartbeat: None,
max_samples: max_samples.max(1),
}
}
pub fn observe(&mut self, at: SystemTime) {
if let Some(prev) = self.last_heartbeat {
if let Ok(interval) = at.duration_since(prev) {
if !interval.is_zero() {
self.samples.push(interval);
if self.samples.len() > self.max_samples {
self.samples.remove(0);
}
}
}
}
self.last_heartbeat = Some(at);
}
pub fn phi_at(&self, now: SystemTime) -> Option<f64> {
let last = self.last_heartbeat?;
if self.samples.is_empty() {
return None;
}
let elapsed = now.duration_since(last).ok()?.as_secs_f64();
let mean = self.mean_secs();
if mean <= f64::EPSILON {
return None;
}
// Exponential approximation: phi = -log10(e^(-elapsed/mean)).
Some((elapsed / mean) / std::f64::consts::LN_10)
}
pub fn mean_secs(&self) -> f64 {
let total: f64 = self.samples.iter().map(Duration::as_secs_f64).sum();
total / self.samples.len() as f64
}
}
pub fn now_millis() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.min(u128::from(u64::MAX)) as u64
}
pub fn system_time_from_millis(millis: u64) -> SystemTime {
UNIX_EPOCH + Duration::from_millis(millis)
}
#[cfg(test)]
mod tests {
use super::*;
fn rid(byte: u8) -> RuntimeId {
RuntimeId([byte; 32])
}
#[test]
fn digest_delta_only_sends_newer_records() {
let mut a = GossipState::new(GossipRecord::new(rid(1), vec!["a".into()], 1));
a.apply_delta(vec![GossipRecord::new(rid(2), vec!["b".into()], 1)]);
let b_digest = [GossipDigestEntry {
runtime_id: rid(1),
version: 1,
}];
let GossipMessage::Delta { records, .. } = a.delta_for(&b_digest) else {
panic!("expected delta")
};
assert_eq!(records.len(), 1);
assert_eq!(records[0].runtime_id, rid(2));
}
#[test]
fn apply_delta_ignores_stale_versions() {
let mut state = GossipState::new(GossipRecord::new(rid(1), vec![], 1));
let newer = GossipRecord {
version: 10,
..GossipRecord::new(rid(2), vec!["new".into()], 1)
};
assert_eq!(state.apply_delta(vec![newer.clone()]), vec![rid(2)]);
let stale = GossipRecord {
version: 9,
endpoint_addrs: vec!["stale".into()],
..newer
};
assert!(state.apply_delta(vec![stale]).is_empty());
assert_eq!(state.get(rid(2)).unwrap().endpoint_addrs, vec!["new"]);
}
#[test]
fn gossip_payload_roundtrips() {
let message = GossipMessage::Digest {
version: GOSSIP_PAYLOAD_VERSION,
entries: vec![GossipDigestEntry {
runtime_id: rid(9),
version: 3,
}],
};
assert_eq!(
decode_message(&encode_message(&message).unwrap()).unwrap(),
message
);
}
#[test]
fn phi_rises_as_heartbeats_age() {
let start = UNIX_EPOCH + Duration::from_secs(1_000);
let mut phi = PhiAccrual::default();
phi.observe(start);
phi.observe(start + Duration::from_secs(1));
phi.observe(start + Duration::from_secs(2));
let early = phi.phi_at(start + Duration::from_secs(3)).unwrap();
let late = phi.phi_at(start + Duration::from_secs(12)).unwrap();
assert!(late > early);
}
}
+8
View File
@@ -18,6 +18,10 @@
//! **The law:** mesh membership is a hint; the Redis fenced generation is the
//! arbiter. Nothing in this crate grants ownership — see [`wire::FencedHeader`].
pub mod gossip;
pub mod membership;
pub mod registry;
pub mod status;
pub mod wire;
// Lane modules — one owner per file (see the mesh thread for lane map):
@@ -32,6 +36,10 @@ use std::pin::Pin;
use bytes::Bytes;
pub use gossip::{GossipDigestEntry, GossipMessage, GossipRecord, GossipState, PhiAccrual};
pub use membership::MeshMembership;
pub use registry::{ReadyHeartbeat, ReadyRecord, ReadyRegistry, RuntimeAttestation};
pub use status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus};
pub use wire::{
FencedHeader, GoodbyeReason, MeshDatagram, MeshStreamFrame, Profile, RuntimeId, StreamHello,
StreamRole, ALPN, WIRE_VERSION,
+377
View File
@@ -0,0 +1,377 @@
//! In-memory mesh membership table fed by Redis seed records and gossip.
//!
//! This module implements the relay-facing [`RelayMeshMembership`] seam. It is
//! deliberately incapable of electing session owners: peers here are dial/routing
//! hints only, and liveness disagreement never performs takeover.
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use std::time::SystemTime;
use crate::gossip::{system_time_from_millis, GossipRecord, PhiAccrual};
use crate::registry::ReadyRecord;
use crate::status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus};
use crate::{PeerInfo, RelayMeshMembership, RuntimeId};
pub const DEFAULT_PHI_SUSPECT_THRESHOLD: f64 = 8.0;
#[derive(Clone, Debug)]
struct PeerState {
record: GossipRecord,
phi: PhiAccrual,
connection_state: ConnectionState,
counters: MeshPeerCounters,
}
/// Thread-safe membership view consumed by the relay.
#[derive(Clone, Debug)]
pub struct MeshMembership {
local_runtime_id: RuntimeId,
local_record: Arc<RwLock<GossipRecord>>,
peers: Arc<RwLock<HashMap<RuntimeId, PeerState>>>,
draining: Arc<AtomicBool>,
stale_generation_rejections: Arc<AtomicU64>,
phi_suspect_threshold: f64,
}
impl MeshMembership {
pub fn new(local_record: GossipRecord) -> Self {
Self {
local_runtime_id: local_record.runtime_id,
local_record: Arc::new(RwLock::new(local_record)),
peers: Arc::new(RwLock::new(HashMap::new())),
draining: Arc::new(AtomicBool::new(false)),
stale_generation_rejections: Arc::new(AtomicU64::new(0)),
phi_suspect_threshold: DEFAULT_PHI_SUSPECT_THRESHOLD,
}
}
pub fn with_phi_suspect_threshold(mut self, threshold: f64) -> Self {
self.phi_suspect_threshold = threshold;
self
}
pub fn local_record(&self) -> GossipRecord {
self.local_record
.read()
.expect("local record lock poisoned")
.clone()
}
/// Apply Redis bootstrap records. Existing gossip records win when they are
/// newer; ready-registry records enter as version 1 hints.
pub fn apply_ready_records(&self, records: impl IntoIterator<Item = ReadyRecord>) {
for ready in records {
if ready.runtime_id == self.local_runtime_id {
continue;
}
if let Err(err) = ready.verify_attestation() {
tracing::warn!(
runtime_id = %ready.runtime_id,
%err,
"mesh membership rejected unauthenticated ready seed"
);
continue;
}
let mut record =
GossipRecord::new(ready.runtime_id, ready.endpoint_addrs, ready.proto_version);
record.capabilities = ready.capabilities;
self.apply_gossip_record(record);
}
}
/// Apply a gossiped record if it is newer than the local copy.
pub fn apply_gossip_record(&self, record: GossipRecord) -> bool {
if record.runtime_id == self.local_runtime_id {
return false;
}
let heartbeat = system_time_from_millis(record.heartbeat_millis);
let mut peers = self.peers.write().expect("membership lock poisoned");
match peers.get_mut(&record.runtime_id) {
Some(peer) if record.version <= peer.record.version => false,
Some(peer) => {
peer.record = record;
peer.connection_state = ConnectionState::Connected;
peer.phi.observe(heartbeat);
true
}
None => {
let mut phi = PhiAccrual::default();
phi.observe(heartbeat);
peers.insert(
record.runtime_id,
PeerState {
counters: MeshPeerCounters {
runtime_id: record.runtime_id.to_string(),
..MeshPeerCounters::default()
},
record,
phi,
connection_state: ConnectionState::Connected,
},
);
true
}
}
}
pub fn mark_connection_state(&self, runtime_id: RuntimeId, state: ConnectionState) {
if let Some(peer) = self
.peers
.write()
.expect("membership lock poisoned")
.get_mut(&runtime_id)
{
peer.connection_state = state;
}
}
pub fn update_local<F>(&self, update: F) -> GossipRecord
where
F: FnOnce(&mut GossipRecord),
{
let mut local = self
.local_record
.write()
.expect("local record lock poisoned");
update(&mut local);
local.version = local.version.saturating_add(1);
local.heartbeat_millis = crate::gossip::now_millis();
local.clone()
}
pub fn is_draining(&self) -> bool {
self.draining.load(Ordering::Relaxed)
}
pub fn record_stream_opened(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.streams_opened = c.streams_opened.saturating_add(1)
});
}
pub fn record_stream_received(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.streams_received = c.streams_received.saturating_add(1)
});
}
pub fn record_datagram_sent(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.datagrams_sent = c.datagrams_sent.saturating_add(1)
});
}
pub fn record_datagram_received(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.datagrams_received = c.datagrams_received.saturating_add(1)
});
}
pub fn record_gossip_frame_sent(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.gossip_frames_sent = c.gossip_frames_sent.saturating_add(1)
});
}
pub fn record_gossip_frame_received(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| {
c.gossip_frames_received = c.gossip_frames_received.saturating_add(1)
});
}
pub fn record_stale_generation_rejection(&self, runtime_id: Option<RuntimeId>) {
self.stale_generation_rejections
.fetch_add(1, Ordering::Relaxed);
if let Some(runtime_id) = runtime_id {
self.update_peer_counters(runtime_id, |c| {
c.stale_generation_rejections = c.stale_generation_rejections.saturating_add(1)
});
}
}
pub fn status(&self) -> MeshStatus {
let now = SystemTime::now();
let local = self.local_record();
let mut peers = self.peer_statuses(now);
peers.sort_by(|a, b| a.runtime_id.cmp(&b.runtime_id));
let counters = MeshCounters {
stale_generation_rejections: self.stale_generation_rejections.load(Ordering::Relaxed),
peers: peers.iter().map(|peer| peer.counters.clone()).collect(),
};
MeshStatus {
enabled: true,
local_runtime_id: local.runtime_id.to_string(),
draining: self.is_draining(),
peer_count: peers.len(),
peers,
counters,
}
}
fn update_peer_counters<F>(&self, runtime_id: RuntimeId, update: F)
where
F: FnOnce(&mut MeshPeerCounters),
{
if let Some(peer) = self
.peers
.write()
.expect("membership lock poisoned")
.get_mut(&runtime_id)
{
update(&mut peer.counters);
}
}
fn peer_statuses(&self, now: SystemTime) -> Vec<MeshPeerStatus> {
self.peers
.read()
.expect("membership lock poisoned")
.values()
.map(|peer| {
let phi = peer.phi.phi_at(now);
let connection_state = if phi.is_some_and(|p| p >= self.phi_suspect_threshold) {
ConnectionState::Suspect
} else {
peer.connection_state
};
MeshPeerStatus {
runtime_id: peer.record.runtime_id.to_string(),
endpoint_addrs: peer.record.endpoint_addrs.clone(),
proto_version: peer.record.proto_version,
draining: peer.record.draining,
connection_state,
phi,
load: peer.record.load,
record_version: peer.record.version,
last_heartbeat_millis: peer.record.heartbeat_millis,
counters: peer.counters.clone(),
}
})
.collect()
}
}
impl RelayMeshMembership for MeshMembership {
fn peers(&self) -> Vec<PeerInfo> {
let now = SystemTime::now();
self.peers
.read()
.expect("membership lock poisoned")
.values()
.filter_map(|peer| {
let phi = peer.phi.phi_at(now);
if phi.is_some_and(|p| p >= self.phi_suspect_threshold) {
return None;
}
Some(PeerInfo {
runtime_id: peer.record.runtime_id,
draining: peer.record.draining,
phi,
load: peer.record.load,
})
})
.collect()
}
fn local_runtime_id(&self) -> RuntimeId {
self.local_runtime_id
}
fn begin_drain(&self) {
self.draining.store(true, Ordering::Relaxed);
self.update_local(|record| record.draining = true);
}
}
#[cfg(test)]
mod tests {
use std::time::{Duration, UNIX_EPOCH};
use super::*;
fn rid(byte: u8) -> RuntimeId {
RuntimeId([byte; 32])
}
fn record(byte: u8, version: u64, heartbeat_secs: u64) -> GossipRecord {
GossipRecord {
runtime_id: rid(byte),
endpoint_addrs: vec![format!("127.0.0.{byte}:3478")],
proto_version: 1,
load: 0.25,
draining: false,
capabilities: vec!["reliable-stream".to_string()],
version,
heartbeat_millis: (UNIX_EPOCH + Duration::from_secs(heartbeat_secs))
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis() as u64,
}
}
fn relay_keys() -> nostr::Keys {
nostr::Keys::generate()
}
fn ready_record(byte: u8, endpoint_addr: &str) -> ReadyRecord {
ReadyRecord::new(
rid(byte),
&relay_keys(),
vec![endpoint_addr.into()],
1,
vec![],
)
}
#[test]
fn ready_records_seed_peers_but_skip_self() {
let membership = MeshMembership::new(record(1, 1, 1));
membership.apply_ready_records([ready_record(1, "self"), ready_record(2, "peer")]);
let peers = membership.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].runtime_id, rid(2));
}
#[test]
fn ready_records_must_have_valid_attestation() {
let membership = MeshMembership::new(record(1, 1, 1));
let mut tampered = ready_record(2, "peer");
tampered.runtime_id = rid(3);
tampered.runtime_pubkey = rid(3).to_hex();
membership.apply_ready_records([tampered]);
assert!(membership.peers().is_empty());
}
#[test]
fn stale_gossip_record_is_ignored() {
let membership = MeshMembership::new(record(1, 1, 1));
assert!(membership.apply_gossip_record(record(2, 5, 1)));
assert!(!membership.apply_gossip_record(record(2, 4, 2)));
assert_eq!(membership.status().peers[0].record_version, 5);
}
#[test]
fn counters_are_reflected_in_status() {
let membership = MeshMembership::new(record(1, 1, 1));
membership.apply_gossip_record(record(2, 1, 1));
membership.record_datagram_sent(rid(2));
membership.record_stale_generation_rejection(Some(rid(2)));
let status = membership.status();
assert_eq!(status.counters.stale_generation_rejections, 1);
assert_eq!(status.peers[0].counters.datagrams_sent, 1);
assert_eq!(status.peers[0].counters.stale_generation_rejections, 1);
}
#[test]
fn begin_drain_updates_local_record() {
let membership = MeshMembership::new(record(1, 1, 1));
membership.begin_drain();
assert!(membership.is_draining());
assert!(membership.local_record().draining);
assert_eq!(membership.local_record().version, 2);
}
}
+385
View File
@@ -0,0 +1,385 @@
//! Redis ready-registry bootstrap for the relay mesh.
//!
//! The registry is only the way into the mesh. Entries are membership hints:
//! they tell a fresh runtime which peer endpoints to dial, but never decide
//! session ownership or takeover. The fenced Redis session directory remains
//! the arbiter for session generations.
use std::str::FromStr;
use std::time::Duration;
use nostr::secp256k1::schnorr::Signature;
use nostr::secp256k1::{Message, XOnlyPublicKey};
use nostr::PublicKey;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::{MeshError, RuntimeId};
pub const READY_KEY_PREFIX: &str = "mesh:ready:";
pub const DEFAULT_REGISTRY_REFRESH: Duration = Duration::from_secs(15);
pub const REGISTRY_EXPIRY_MULTIPLIER: u64 = 3;
pub const ATTESTATION_CONTEXT: &str = "buzz-relay-mesh-ready-v1";
/// Relay-key-signed binding for a boot-unique runtime endpoint pubkey.
///
/// The relay public key is the deployment Nostr/secp256k1 identity. It never
/// becomes the mesh runtime id; it only signs this Redis-published binding so
/// peers can reject unauthenticated endpoint ids before dialing/accepting.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct RuntimeAttestation {
/// Nostr/secp256k1 relay public key, hex encoded.
pub relay_pubkey: String,
/// Schnorr signature by `relay_pubkey` over [`attestation_preimage`].
pub relay_sig: String,
}
impl RuntimeAttestation {
pub fn new(relay_keys: &nostr::Keys, runtime_id: RuntimeId) -> Self {
let relay_pubkey = relay_keys.public_key().to_hex();
let message = attestation_message(runtime_id, &relay_pubkey);
let relay_sig = relay_keys.sign_schnorr(&message).to_string();
Self {
relay_pubkey,
relay_sig,
}
}
pub fn verify(&self, runtime_id: RuntimeId) -> Result<(), MeshError> {
verify_attestation(runtime_id, &self.relay_pubkey, &self.relay_sig)
}
}
fn verify_attestation(
runtime_id: RuntimeId,
relay_pubkey: &str,
relay_sig: &str,
) -> Result<(), MeshError> {
let relay_pubkey = PublicKey::from_hex(relay_pubkey).map_err(|err| {
MeshError::Transport(format!(
"ready registry attestation invalid relay_pubkey: {err}"
))
})?;
let xonly: XOnlyPublicKey = relay_pubkey.xonly().map_err(|err| {
MeshError::Transport(format!(
"ready registry attestation relay_pubkey xonly conversion failed: {err}"
))
})?;
let sig = Signature::from_str(relay_sig).map_err(|err| {
MeshError::Transport(format!(
"ready registry attestation invalid relay_sig: {err}"
))
})?;
let message = attestation_message(runtime_id, &relay_pubkey.to_hex());
nostr::secp256k1::SECP256K1
.verify_schnorr(&sig, &message, &xonly)
.map_err(|err| {
MeshError::Transport(format!(
"ready registry attestation signature verification failed: {err}"
))
})
}
/// Stable signed payload. Keep this textual and versioned so transport/relay
/// integration can reproduce it exactly without depending on JSON key order.
pub fn attestation_preimage(runtime_id: RuntimeId, relay_pubkey: &str) -> String {
format!(
"{ATTESTATION_CONTEXT}\nruntime_pubkey={}\nrelay_pubkey={relay_pubkey}",
runtime_id.to_hex()
)
}
fn attestation_message(runtime_id: RuntimeId, relay_pubkey: &str) -> Message {
let digest = Sha256::digest(attestation_preimage(runtime_id, relay_pubkey).as_bytes());
Message::from_digest(digest.into())
}
/// Value stored at `mesh:ready:{runtime_id}`.
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct ReadyRecord {
pub runtime_id: RuntimeId,
/// Explicit duplicate of `runtime_id` for the contract record shape: this
/// is the boot-unique ed25519/iroh endpoint pubkey being attested.
pub runtime_pubkey: String,
/// Nostr/secp256k1 relay public key that signs `runtime_pubkey`.
pub relay_pubkey: String,
/// Schnorr signature by `relay_pubkey` over [`attestation_preimage`].
pub relay_sig: String,
/// Dialable iroh endpoint addresses, serialized as strings so this layer
/// does not depend on transport internals.
pub endpoint_addrs: Vec<String>,
pub proto_version: u16,
pub capabilities: Vec<String>,
}
impl ReadyRecord {
pub fn new(
runtime_id: RuntimeId,
relay_keys: &nostr::Keys,
endpoint_addrs: Vec<String>,
proto_version: u16,
capabilities: Vec<String>,
) -> Self {
let attestation = RuntimeAttestation::new(relay_keys, runtime_id);
Self {
runtime_id,
runtime_pubkey: runtime_id.to_hex(),
relay_pubkey: attestation.relay_pubkey,
relay_sig: attestation.relay_sig,
endpoint_addrs,
proto_version,
capabilities,
}
}
pub fn key(&self) -> String {
ready_key(self.runtime_id)
}
pub fn verify_attestation(&self) -> Result<(), MeshError> {
if self.runtime_pubkey != self.runtime_id.to_hex() {
return Err(MeshError::Transport(format!(
"ready registry runtime_id/runtime_pubkey mismatch: {} != {}",
self.runtime_id, self.runtime_pubkey
)));
}
verify_attestation(self.runtime_id, &self.relay_pubkey, &self.relay_sig)
}
}
pub fn ready_key(runtime_id: RuntimeId) -> String {
format!("{READY_KEY_PREFIX}{runtime_id}")
}
pub fn expiry_for(refresh: Duration) -> Duration {
refresh.saturating_mul(REGISTRY_EXPIRY_MULTIPLIER as u32)
}
/// Redis-backed mesh bootstrap registry.
#[derive(Clone)]
pub struct ReadyRegistry {
pool: deadpool_redis::Pool,
refresh: Duration,
}
impl ReadyRegistry {
pub fn new(pool: deadpool_redis::Pool, refresh: Duration) -> Self {
Self { pool, refresh }
}
pub fn refresh_interval(&self) -> Duration {
self.refresh
}
pub fn expiry(&self) -> Duration {
expiry_for(self.refresh)
}
/// Publish this runtime as ready. Callers MUST only invoke this after the
/// relay would pass readiness (shutdown=false, Postgres reachable, Redis
/// reachable). This method deliberately has no hidden readiness probe so the
/// rule stays explicit at the relay boundary.
pub async fn publish_ready(&self, record: &ReadyRecord) -> Result<(), MeshError> {
record.verify_attestation()?;
let mut conn = self.conn().await?;
let payload = serde_json::to_string(record)
.map_err(|e| MeshError::Transport(format!("ready registry encode: {e}")))?;
let ttl_secs = self.expiry().as_secs().max(1);
redis::cmd("SET")
.arg(record.key())
.arg(payload)
.arg("EX")
.arg(ttl_secs)
.query_async::<()>(&mut conn)
.await?;
Ok(())
}
/// Remove this runtime on clean shutdown. A crash is handled by TTL expiry.
pub async fn clear_ready(&self, runtime_id: RuntimeId) -> Result<(), MeshError> {
let mut conn = self.conn().await?;
redis::cmd("DEL")
.arg(ready_key(runtime_id))
.query_async::<()>(&mut conn)
.await?;
Ok(())
}
/// Scan all ready records. Malformed/stale/unauthenticated values are
/// skipped with a warn: a bad registry entry must not prevent bootstrap
/// from healthy peers.
pub async fn scan_ready(&self) -> Result<Vec<ReadyRecord>, MeshError> {
let mut conn = self.conn().await?;
let mut cursor = 0u64;
let mut out = Vec::new();
loop {
let (next, keys): (u64, Vec<String>) = redis::cmd("SCAN")
.arg(cursor)
.arg("MATCH")
.arg(format!("{READY_KEY_PREFIX}*"))
.arg("COUNT")
.arg(100u32)
.query_async(&mut conn)
.await?;
for key in keys {
let raw: Option<String> =
redis::cmd("GET").arg(&key).query_async(&mut conn).await?;
let Some(raw) = raw else { continue };
match serde_json::from_str::<ReadyRecord>(&raw) {
Ok(record) if record.key() == key => match record.verify_attestation() {
Ok(()) => out.push(record),
Err(err) => tracing::warn!(
key,
runtime_id = %record.runtime_id,
%err,
"mesh ready registry attestation failed — skipping"
),
},
Ok(record) => tracing::warn!(
key,
runtime_id = %record.runtime_id,
"mesh ready registry key/runtime mismatch — skipping"
),
Err(err) => {
tracing::warn!(key, %err, "mesh ready registry decode failed — skipping")
}
}
}
if next == 0 {
break;
}
cursor = next;
}
Ok(out)
}
pub fn heartbeat(&self, record: ReadyRecord) -> ReadyHeartbeat {
ReadyHeartbeat {
registry: self.clone(),
record,
published: false,
}
}
async fn conn(&self) -> Result<deadpool_redis::Connection, MeshError> {
self.pool
.get()
.await
.map_err(|e| MeshError::Transport(format!("redis pool: {e}")))
}
}
/// Readiness-gated registry heartbeat.
///
/// The relay owns the readiness predicate; this helper owns the edge behavior:
/// publish only while ready, clear on ready→not-ready, and clear on shutdown.
pub struct ReadyHeartbeat {
registry: ReadyRegistry,
record: ReadyRecord,
published: bool,
}
impl ReadyHeartbeat {
pub fn record(&self) -> &ReadyRecord {
&self.record
}
pub fn published(&self) -> bool {
self.published
}
pub async fn tick(&mut self, ready: bool) -> Result<(), MeshError> {
if ready {
self.registry.publish_ready(&self.record).await?;
self.published = true;
} else if self.published {
self.registry.clear_ready(self.record.runtime_id).await?;
self.published = false;
}
Ok(())
}
pub async fn shutdown(&mut self) -> Result<(), MeshError> {
if self.published {
self.registry.clear_ready(self.record.runtime_id).await?;
self.published = false;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn rid(byte: u8) -> RuntimeId {
RuntimeId([byte; 32])
}
fn relay_keys() -> nostr::Keys {
nostr::Keys::generate()
}
fn ready_record(byte: u8) -> ReadyRecord {
ReadyRecord::new(rid(byte), &relay_keys(), vec![], 1, vec![])
}
#[test]
fn ready_key_is_stable_and_namespaced() {
assert_eq!(
ready_key(rid(0xAB)),
format!("mesh:ready:{}", "ab".repeat(32))
);
}
#[test]
fn expiry_is_three_refreshes() {
assert_eq!(expiry_for(Duration::from_secs(15)), Duration::from_secs(45));
}
#[test]
fn heartbeat_starts_unpublished() {
let pool = deadpool_redis::Config::from_url("redis://127.0.0.1:6379")
.create_pool(Some(deadpool_redis::Runtime::Tokio1))
.unwrap();
let registry = ReadyRegistry::new(pool, Duration::from_secs(15));
let heartbeat = registry.heartbeat(ready_record(1));
assert!(!heartbeat.published());
assert_eq!(heartbeat.record().runtime_id, rid(1));
}
#[test]
fn ready_record_roundtrips_json() {
let record = ReadyRecord::new(
rid(7),
&relay_keys(),
vec!["127.0.0.1:3478".to_string()],
1,
vec!["realtime-media".to_string()],
);
let raw = serde_json::to_string(&record).unwrap();
assert_eq!(serde_json::from_str::<ReadyRecord>(&raw).unwrap(), record);
}
#[test]
fn ready_record_attestation_verifies_and_binds_runtime_pubkey() {
let record = ready_record(9);
record.verify_attestation().unwrap();
let mut tampered = record.clone();
tampered.runtime_pubkey = rid(10).to_hex();
assert!(tampered.verify_attestation().is_err());
}
#[test]
fn attestation_rejects_signature_for_other_runtime() {
let mut record = ready_record(11);
record.runtime_id = rid(12);
record.runtime_pubkey = rid(12).to_hex();
assert!(record.verify_attestation().is_err());
}
}
+57
View File
@@ -0,0 +1,57 @@
//! `/_mesh` status data model.
//!
//! The relay's axum handler can serialize [`MeshStatus`] directly as JSON.
use serde::Serialize;
#[derive(Clone, Debug, Default, Serialize)]
pub struct MeshStatus {
pub enabled: bool,
pub local_runtime_id: String,
pub draining: bool,
pub peer_count: usize,
pub peers: Vec<MeshPeerStatus>,
pub counters: MeshCounters,
}
#[derive(Clone, Debug, Serialize)]
pub struct MeshPeerStatus {
pub runtime_id: String,
pub endpoint_addrs: Vec<String>,
pub proto_version: u16,
pub draining: bool,
pub connection_state: ConnectionState,
pub phi: Option<f64>,
pub load: f32,
pub record_version: u64,
pub last_heartbeat_millis: u64,
pub counters: MeshPeerCounters,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ConnectionState {
#[default]
Disconnected,
Connecting,
Connected,
Suspect,
}
#[derive(Clone, Debug, Default, Serialize)]
pub struct MeshCounters {
pub stale_generation_rejections: u64,
pub peers: Vec<MeshPeerCounters>,
}
#[derive(Clone, Debug, Default, Serialize)]
pub struct MeshPeerCounters {
pub runtime_id: String,
pub streams_opened: u64,
pub streams_received: u64,
pub datagrams_sent: u64,
pub datagrams_received: u64,
pub gossip_frames_sent: u64,
pub gossip_frames_received: u64,
pub stale_generation_rejections: u64,
}