diff --git a/Cargo.lock b/Cargo.lock index 39af7d3cf..2fc4fe836 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1104,6 +1104,7 @@ dependencies = [ "hex", "hmac 0.13.0", "iroh", + "nostr", "postcard", "proptest", "redis", diff --git a/crates/buzz-relay-mesh/Cargo.toml b/crates/buzz-relay-mesh/Cargo.toml index 33f901994..15d0c5ba3 100644 --- a/crates/buzz-relay-mesh/Cargo.toml +++ b/crates/buzz-relay-mesh/Cargo.toml @@ -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 } diff --git a/crates/buzz-relay-mesh/src/gossip.rs b/crates/buzz-relay-mesh/src/gossip.rs new file mode 100644 index 000000000..5da5cbbe5 --- /dev/null +++ b/crates/buzz-relay-mesh/src/gossip.rs @@ -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, + pub proto_version: u16, + pub load: f32, + pub draining: bool, + pub capabilities: Vec, + /// 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, 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, + }, + Delta { + version: u8, + records: Vec, + }, +} + +pub fn encode_message(message: &GossipMessage) -> Result, MeshError> { + postcard::to_extend(message, Vec::new()).map_err(MeshError::Encode) +} + +pub fn decode_message(bytes: &[u8]) -> Result { + 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, +} + +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 { + self.records.values() + } + + pub fn get(&self, runtime_id: RuntimeId) -> Option<&GossipRecord> { + self.records.get(&runtime_id) + } + + pub fn update_local(&mut self, runtime_id: RuntimeId, update: F) -> Option + 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) -> Vec { + 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, + last_heartbeat: Option, + 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 { + 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); + } +} diff --git a/crates/buzz-relay-mesh/src/lib.rs b/crates/buzz-relay-mesh/src/lib.rs index f4c5e10e1..42c02b67a 100644 --- a/crates/buzz-relay-mesh/src/lib.rs +++ b/crates/buzz-relay-mesh/src/lib.rs @@ -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, diff --git a/crates/buzz-relay-mesh/src/membership.rs b/crates/buzz-relay-mesh/src/membership.rs new file mode 100644 index 000000000..bfe2c3172 --- /dev/null +++ b/crates/buzz-relay-mesh/src/membership.rs @@ -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>, + peers: Arc>>, + draining: Arc, + stale_generation_rejections: Arc, + 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) { + 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(&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) { + 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(&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 { + 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 { + 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); + } +} diff --git a/crates/buzz-relay-mesh/src/registry.rs b/crates/buzz-relay-mesh/src/registry.rs new file mode 100644 index 000000000..635dff632 --- /dev/null +++ b/crates/buzz-relay-mesh/src/registry.rs @@ -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, + pub proto_version: u16, + pub capabilities: Vec, +} + +impl ReadyRecord { + pub fn new( + runtime_id: RuntimeId, + relay_keys: &nostr::Keys, + endpoint_addrs: Vec, + proto_version: u16, + capabilities: Vec, + ) -> 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, MeshError> { + let mut conn = self.conn().await?; + let mut cursor = 0u64; + let mut out = Vec::new(); + + loop { + let (next, keys): (u64, Vec) = 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 = + redis::cmd("GET").arg(&key).query_async(&mut conn).await?; + let Some(raw) = raw else { continue }; + match serde_json::from_str::(&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 { + 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::(&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()); + } +} diff --git a/crates/buzz-relay-mesh/src/status.rs b/crates/buzz-relay-mesh/src/status.rs new file mode 100644 index 000000000..a05cb276f --- /dev/null +++ b/crates/buzz-relay-mesh/src/status.rs @@ -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, + pub counters: MeshCounters, +} + +#[derive(Clone, Debug, Serialize)] +pub struct MeshPeerStatus { + pub runtime_id: String, + pub endpoint_addrs: Vec, + pub proto_version: u16, + pub draining: bool, + pub connection_state: ConnectionState, + pub phi: Option, + 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, +} + +#[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, +}