From 8b077fdb44903de6cda44e6ea23ef5fe0e208d88 Mon Sep 17 00:00:00 2001 From: npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d <011987e296fd5006292d2f930b574be47c7801048d1983c46c425d3c95f0cffd@sprout-oss.stage.blox.sqprod.co> Date: Wed, 8 Jul 2026 12:43:32 -0400 Subject: [PATCH] Wire mesh crate into relay startup behind BUZZ_MESH seam - MeshRuntime (buzz-relay-mesh): iroh endpoint ownership, warm peer dial/accept loops, per-peer control streams with gossip digest/delta exchange, registry-driven reconcile, drain flag, status snapshot. - boot_mesh (buzz-relay/mesh_boot.rs): BUZZ_MESH=off returns None with zero binds/spawns/Redis I/O (mesh-off byte-identical); enabled boots endpoint on BUZZ_MESH_BIND_ADDR, signs relay-key attestation into the ready registry, spawns runtime. Misconfig while enabled is fatal. - AppState.mesh via OnceLock set post-construction in main.rs; accessor state.mesh() -> Option<&MeshHandle>. No AppState::new call sites touched. - /_mesh status route on the health router; reports enabled:false when off. - membership: has_peer/records/digest/delta_for; endpoint: ip_addrs(). cargo test -p buzz-relay-mesh: 30 passed; -p buzz-relay: 491 passed, 2 ignored; clippy both packages clean. Co-authored-by: Tyler Longwell Signed-off-by: Tyler Longwell --- crates/buzz-relay-mesh/src/endpoint.rs | 14 + crates/buzz-relay-mesh/src/lib.rs | 2 + crates/buzz-relay-mesh/src/membership.rs | 65 ++ crates/buzz-relay-mesh/src/runtime.rs | 896 +++++++++++++++++++++++ crates/buzz-relay/src/lib.rs | 2 + crates/buzz-relay/src/main.rs | 20 + crates/buzz-relay/src/mesh_boot.rs | 214 ++++++ crates/buzz-relay/src/router.rs | 13 + crates/buzz-relay/src/state.rs | 14 + 9 files changed, 1240 insertions(+) create mode 100644 crates/buzz-relay-mesh/src/runtime.rs create mode 100644 crates/buzz-relay/src/mesh_boot.rs diff --git a/crates/buzz-relay-mesh/src/endpoint.rs b/crates/buzz-relay-mesh/src/endpoint.rs index 97fd84644..f7bedd878 100644 --- a/crates/buzz-relay-mesh/src/endpoint.rs +++ b/crates/buzz-relay-mesh/src/endpoint.rs @@ -56,6 +56,20 @@ impl MeshEndpoint { self.endpoint.addr() } + /// The endpoint's directly-dialable IP socket addrs (no relay paths). + /// Lets consumers build advertise records without depending on iroh types. + pub fn ip_addrs(&self) -> Vec { + self.endpoint + .addr() + .addrs + .iter() + .filter_map(|ta| match ta { + TransportAddr::Ip(sock) => Some(*sock), + _ => None, + }) + .collect() + } + pub async fn accept(&self) -> Result, MeshError> { let Some(incoming) = self.endpoint.accept().await else { return Ok(None); diff --git a/crates/buzz-relay-mesh/src/lib.rs b/crates/buzz-relay-mesh/src/lib.rs index 56a4c2deb..676d73247 100644 --- a/crates/buzz-relay-mesh/src/lib.rs +++ b/crates/buzz-relay-mesh/src/lib.rs @@ -23,6 +23,7 @@ pub mod gossip; pub mod membership; pub mod peer; pub mod registry; +pub mod runtime; pub mod status; pub mod wire; @@ -41,6 +42,7 @@ use bytes::Bytes; pub use gossip::{GossipDigestEntry, GossipMessage, GossipRecord, GossipState, PhiAccrual}; pub use membership::MeshMembership; pub use registry::{ReadyHeartbeat, ReadyRecord, ReadyRegistry, RuntimeAttestation}; +pub use runtime::MeshRuntime; pub use status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus}; pub use wire::{ FencedHeader, GoodbyeReason, MeshDatagram, MeshStreamFrame, Profile, RuntimeId, StreamHello, diff --git a/crates/buzz-relay-mesh/src/membership.rs b/crates/buzz-relay-mesh/src/membership.rs index bfe2c3172..f1d330bcb 100644 --- a/crates/buzz-relay-mesh/src/membership.rs +++ b/crates/buzz-relay-mesh/src/membership.rs @@ -146,6 +146,71 @@ impl MeshMembership { self.draining.load(Ordering::Relaxed) } + /// Whether `runtime_id` is present in the (attested) peer table. Used by + /// the runtime's accept loop to gate inbound connections — a dialability + /// hint, never an ownership statement. + pub fn has_peer(&self, runtime_id: RuntimeId) -> bool { + self.peers + .read() + .expect("membership lock poisoned") + .contains_key(&runtime_id) + } + + /// All known gossip records (local + peers), for reconcile/dial decisions. + pub fn records(&self) -> Vec { + let mut records: Vec = self + .peers + .read() + .expect("membership lock poisoned") + .values() + .map(|peer| peer.record.clone()) + .collect(); + records.push(self.local_record()); + records + } + + /// Scuttlebutt digest over every record this runtime knows (local + peers). + pub fn digest(&self) -> crate::gossip::GossipMessage { + let mut entries: Vec<_> = self + .records() + .into_iter() + .map(|record| crate::gossip::GossipDigestEntry { + runtime_id: record.runtime_id, + version: record.version, + }) + .collect(); + entries.sort_by_key(|entry| entry.runtime_id.to_hex()); + crate::gossip::GossipMessage::Digest { + version: crate::gossip::GOSSIP_PAYLOAD_VERSION, + entries, + } + } + + /// Records the remote digest is missing or behind on. + pub fn delta_for( + &self, + digest: &[crate::gossip::GossipDigestEntry], + ) -> crate::gossip::GossipMessage { + let remote: std::collections::HashMap<_, _> = digest + .iter() + .map(|entry| (entry.runtime_id, entry.version)) + .collect(); + let mut records: Vec<_> = self + .records() + .into_iter() + .filter(|record| { + remote + .get(&record.runtime_id) + .is_none_or(|version| *version < record.version) + }) + .collect(); + records.sort_by_key(|record| record.runtime_id.to_hex()); + crate::gossip::GossipMessage::Delta { + version: crate::gossip::GOSSIP_PAYLOAD_VERSION, + records, + } + } + 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) diff --git a/crates/buzz-relay-mesh/src/runtime.rs b/crates/buzz-relay-mesh/src/runtime.rs new file mode 100644 index 000000000..897de6ca8 --- /dev/null +++ b/crates/buzz-relay-mesh/src/runtime.rs @@ -0,0 +1,896 @@ +//! The live mesh runtime: warm peer manager, accept/dial loops, gossip +//! exchange, and the concrete [`RelayPeerTransport`] implementation. +//! +//! This is the piece that turns the lane modules into a running mesh: +//! +//! - [`MeshRuntime::start`] binds nothing itself — it takes an already-bound +//! [`MeshEndpoint`] plus a [`MeshMembership`] table and spawns the loops. +//! - **Accept loop**: inbound connections are admitted only when the remote +//! runtime id is present in the (attested) membership table; unknown ids get +//! one registry rescan before rejection. Membership is a hint — admission +//! here gates *dialability*, never session ownership. +//! - **Reconcile loop**: periodically rescans the Redis ready registry and +//! dials every known, non-draining peer we are not yet connected to. This is +//! what makes the mesh *warm*: failover is "next frame goes elsewhere," not +//! "wait for a handshake." +//! - **Control stream**: exactly one per peer connection, opened by the +//! dialer. Carries scuttlebutt gossip (`Digest` → `Delta`) both ways. +//! - **Simultaneous dial tie-break**: the connection dialed by the smaller +//! runtime id wins; the loser is dropped. Deterministic on both ends. +//! +//! The fencing law holds here too: nothing in this file consults or mutates +//! session ownership. Transport moves fenced bytes; the session layer on both +//! ends validates the fence against Redis. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex, RwLock}; +use std::time::Duration; + +use tokio::sync::mpsc; +use tokio::task::JoinHandle; + +use crate::endpoint::{direct_addr, MeshEndpoint}; +use crate::gossip::{decode_message, encode_message, GossipMessage}; +use crate::membership::MeshMembership; +use crate::peer::MeshPeer; +use crate::registry::{ReadyRecord, ReadyRegistry}; +use crate::status::ConnectionState; +use crate::wire::{MeshStreamFrame, StreamHello, StreamRole}; +use crate::{InboundHandler, MeshDatagram, MeshError, MeshStream, RelayPeerTransport, RuntimeId}; + +/// How often the reconcile loop rescans the registry and dials missing peers. +pub const DEFAULT_RECONCILE_INTERVAL: Duration = Duration::from_secs(5); +/// How often each side sends a gossip digest on every control stream. +pub const DEFAULT_GOSSIP_INTERVAL: Duration = Duration::from_secs(2); +/// Bound on queued control-stream frames per peer before backpressure. +const CONTROL_QUEUE_DEPTH: usize = 64; + +struct PeerEntry { + peer: MeshPeer, + /// Writer queue for the peer's control stream. Present once the control + /// stream is up (dialer opens it; acceptor receives it). + control_tx: Option>, + tasks: Vec>, +} + +impl PeerEntry { + fn abort(&self) { + for task in &self.tasks { + task.abort(); + } + } +} + +struct Inner { + endpoint: MeshEndpoint, + membership: MeshMembership, + registry: Option, + peers: RwLock>, + handler: Mutex>>, + gossip_interval: Duration, + reconcile_interval: Duration, +} + +/// Handle to the running mesh. Cheap to clone; dropping all clones does NOT +/// stop the loops — call [`MeshRuntime::shutdown`] for that. +#[derive(Clone)] +pub struct MeshRuntime { + inner: Arc, + loops: Arc>>>, +} + +impl MeshRuntime { + /// Spawn the mesh loops over an already-bound endpoint. + /// + /// `registry` is `None` in tests / single-instance shapes: the reconcile + /// loop then dials from the membership table alone (seeded by gossip or + /// test setup) and skips registry rescans. + pub fn start( + endpoint: MeshEndpoint, + membership: MeshMembership, + registry: Option, + ) -> Self { + Self::start_with_intervals( + endpoint, + membership, + registry, + DEFAULT_GOSSIP_INTERVAL, + DEFAULT_RECONCILE_INTERVAL, + ) + } + + pub fn start_with_intervals( + endpoint: MeshEndpoint, + membership: MeshMembership, + registry: Option, + gossip_interval: Duration, + reconcile_interval: Duration, + ) -> Self { + let inner = Arc::new(Inner { + endpoint, + membership, + registry, + peers: RwLock::new(HashMap::new()), + handler: Mutex::new(None), + gossip_interval, + reconcile_interval, + }); + + let accept = tokio::spawn(accept_loop(Arc::clone(&inner))); + let reconcile = tokio::spawn(reconcile_loop(Arc::clone(&inner))); + let gossip = tokio::spawn(gossip_tick_loop(Arc::clone(&inner))); + + Self { + inner, + loops: Arc::new(Mutex::new(vec![accept, reconcile, gossip])), + } + } + + pub fn membership(&self) -> &MeshMembership { + &self.inner.membership + } + + pub fn local_runtime_id(&self) -> RuntimeId { + self.inner.endpoint.runtime_id() + } + + /// Currently connected peer ids (either direction). + pub fn connected_peers(&self) -> Vec { + self.inner + .peers + .read() + .expect("peer lock poisoned") + .keys() + .copied() + .collect() + } + + /// Force one reconcile pass right now (bootstrap fast-path: dial the seed + /// records without waiting for the first interval tick). + pub async fn reconcile_now(&self) { + reconcile_once(&self.inner).await; + } + + /// Stop all loops and drop all peer connections. + pub fn shutdown(&self) { + for task in self.loops.lock().expect("loop lock poisoned").drain(..) { + task.abort(); + } + let mut peers = self.inner.peers.write().expect("peer lock poisoned"); + for (_, entry) in peers.drain() { + entry.abort(); + } + } +} + +impl RelayPeerTransport for MeshRuntime { + fn send_datagram(&self, to: RuntimeId, dgram: MeshDatagram) -> Result<(), MeshError> { + let peers = self.inner.peers.read().expect("peer lock poisoned"); + let entry = peers.get(&to).ok_or(MeshError::PeerNotConnected(to))?; + entry.peer.send_datagram(&dgram)?; + drop(peers); + self.inner.membership.record_datagram_sent(to); + Ok(()) + } + + fn open_session_stream( + &self, + to: RuntimeId, + hello: StreamHello, + ) -> crate::BoxFuture<'_, Result> { + Box::pin(async move { + let peer = { + let peers = self.inner.peers.read().expect("peer lock poisoned"); + peers + .get(&to) + .map(|entry| entry.peer.clone()) + .ok_or(MeshError::PeerNotConnected(to))? + }; + let mut stream = peer.open_bi().await?; + stream.send_frame(MeshStreamFrame::Hello(hello)).await?; + self.inner.membership.record_stream_opened(to); + Ok(stream) + }) + } + + fn set_inbound(&self, handler: Box) { + *self.inner.handler.lock().expect("handler lock poisoned") = Some(Arc::from(handler)); + } +} + +fn inbound_handler(inner: &Inner) -> Option> { + inner.handler.lock().expect("handler lock poisoned").clone() +} + +/// Simultaneous-dial tie-break: the connection dialed by the smaller runtime +/// id wins. Returns true when the NEW connection should replace the existing. +fn new_connection_wins(local: RuntimeId, remote: RuntimeId, new_dialed_by_us: bool) -> bool { + if local.0 < remote.0 { + // We are the canonical dialer: our outbound connection wins. + new_dialed_by_us + } else { + // The peer is the canonical dialer: their inbound connection wins. + !new_dialed_by_us + } +} + +/// Install a connected peer, spawning its datagram + stream accept loops. +/// Returns false when an existing connection won the tie-break. +fn install_peer(inner: &Arc, peer: MeshPeer, dialed_by_us: bool) -> bool { + let remote = peer.runtime_id(); + let local = inner.endpoint.runtime_id(); + let mut peers = inner.peers.write().expect("peer lock poisoned"); + + if let Some(existing) = peers.get(&remote) { + if !new_connection_wins(local, remote, dialed_by_us) { + tracing::debug!(peer = %remote, "mesh: kept existing connection (tie-break)"); + return false; + } + existing.abort(); + } + + let mut tasks = vec![ + tokio::spawn(datagram_recv_loop(Arc::clone(inner), peer.clone())), + tokio::spawn(stream_accept_loop(Arc::clone(inner), peer.clone())), + ]; + + // Dialer opens the control stream for the connection. + let control_tx = if dialed_by_us { + let (tx, rx) = mpsc::channel(CONTROL_QUEUE_DEPTH); + tasks.push(tokio::spawn(open_control_stream( + Arc::clone(inner), + peer.clone(), + rx, + ))); + Some(tx) + } else { + None + }; + + peers.insert( + remote, + PeerEntry { + peer, + control_tx, + tasks, + }, + ); + drop(peers); + inner + .membership + .mark_connection_state(remote, ConnectionState::Connected); + tracing::info!(peer = %remote, dialed_by_us, "mesh: peer connected"); + true +} + +fn remove_peer(inner: &Inner, runtime_id: RuntimeId) { + if let Some(entry) = inner + .peers + .write() + .expect("peer lock poisoned") + .remove(&runtime_id) + { + entry.abort(); + inner + .membership + .mark_connection_state(runtime_id, ConnectionState::Disconnected); + tracing::info!(peer = %runtime_id, "mesh: peer disconnected"); + } +} + +async fn accept_loop(inner: Arc) { + loop { + match inner.endpoint.accept().await { + Ok(Some(peer)) => { + let remote = peer.runtime_id(); + if !is_known_peer(&inner, remote).await { + tracing::warn!( + peer = %remote, + "mesh: rejected inbound connection from unattested runtime id" + ); + continue; + } + install_peer(&inner, peer, false); + } + Ok(None) => { + tracing::info!("mesh: endpoint closed, accept loop exiting"); + return; + } + Err(err) => { + tracing::warn!(%err, "mesh: inbound connection failed"); + } + } + } +} + +/// Admission check for inbound connections: the runtime id must appear in the +/// attested membership table. Unknown ids get one registry rescan (covers the +/// bootstrap race where a fresh pod dials us before our next reconcile tick). +async fn is_known_peer(inner: &Arc, runtime_id: RuntimeId) -> bool { + if inner.membership.has_peer(runtime_id) { + return true; + } + if let Some(registry) = &inner.registry { + match registry.scan_ready().await { + Ok(records) => inner.membership.apply_ready_records(records), + Err(err) => tracing::warn!(%err, "mesh: registry rescan on inbound failed"), + } + } + inner.membership.has_peer(runtime_id) +} + +async fn reconcile_loop(inner: Arc) { + loop { + reconcile_once(&inner).await; + tokio::time::sleep(inner.reconcile_interval).await; + } +} + +async fn reconcile_once(inner: &Arc) { + if let Some(registry) = &inner.registry { + match registry.scan_ready().await { + Ok(records) => inner.membership.apply_ready_records(records), + Err(err) => tracing::warn!(%err, "mesh: registry scan failed"), + } + } + + let local = inner.endpoint.runtime_id(); + let candidates: Vec<_> = inner + .membership + .records() + .into_iter() + .filter(|record| record.runtime_id != local && !record.draining) + .collect(); + + for record in candidates { + let already_connected = inner + .peers + .read() + .expect("peer lock poisoned") + .contains_key(&record.runtime_id); + if already_connected { + continue; + } + dial_peer(inner, &record).await; + } +} + +async fn dial_peer(inner: &Arc, record: &crate::gossip::GossipRecord) { + for addr in &record.endpoint_addrs { + let sock = match addr.parse() { + Ok(sock) => sock, + Err(err) => { + tracing::warn!(peer = %record.runtime_id, addr, %err, "mesh: bad peer addr"); + continue; + } + }; + let endpoint_addr = match direct_addr(record.runtime_id, sock) { + Ok(ea) => ea, + Err(err) => { + tracing::warn!(peer = %record.runtime_id, %err, "mesh: bad peer id"); + return; + } + }; + inner + .membership + .mark_connection_state(record.runtime_id, ConnectionState::Connecting); + match inner.endpoint.connect(endpoint_addr).await { + Ok(peer) => { + install_peer(inner, peer, true); + return; + } + Err(err) => { + tracing::warn!(peer = %record.runtime_id, addr, %err, "mesh: dial failed"); + } + } + } + inner + .membership + .mark_connection_state(record.runtime_id, ConnectionState::Disconnected); +} + +async fn datagram_recv_loop(inner: Arc, peer: MeshPeer) { + let remote = peer.runtime_id(); + loop { + match peer.recv_datagram().await { + Ok(dgram) => { + inner.membership.record_datagram_received(remote); + if let Some(handler) = inbound_handler(&inner) { + handler.on_datagram(remote, dgram); + } + } + Err(err) => { + tracing::debug!(peer = %remote, %err, "mesh: datagram loop ended"); + remove_peer(&inner, remote); + return; + } + } + } +} + +async fn stream_accept_loop(inner: Arc, peer: MeshPeer) { + let remote = peer.runtime_id(); + loop { + let mut stream = match peer.accept_bi().await { + Ok(stream) => stream, + Err(err) => { + tracing::debug!(peer = %remote, %err, "mesh: stream accept loop ended"); + remove_peer(&inner, remote); + return; + } + }; + + // The first frame on any stream MUST be Hello (wire contract). + let hello = match stream.recv_frame().await { + Ok(Some(MeshStreamFrame::Hello(hello))) => hello, + Ok(other) => { + tracing::warn!(peer = %remote, ?other, "mesh: stream without Hello — dropped"); + continue; + } + Err(err) => { + tracing::warn!(peer = %remote, %err, "mesh: stream Hello read failed"); + continue; + } + }; + + match hello.role { + StreamRole::Control => { + // Acceptor side of the per-connection control stream: register + // a writer queue and start the gossip exchange. + let (tx, rx) = mpsc::channel(CONTROL_QUEUE_DEPTH); + if let Some(entry) = inner + .peers + .write() + .expect("peer lock poisoned") + .get_mut(&remote) + { + entry.control_tx = Some(tx); + } + tokio::spawn(control_stream_exchange( + Arc::clone(&inner), + remote, + stream, + rx, + )); + } + StreamRole::Session { .. } => { + inner.membership.record_stream_received(remote); + if let Some(handler) = inbound_handler(&inner) { + handler.on_session_stream(remote, hello, stream); + } else { + tracing::warn!( + peer = %remote, + "mesh: session stream arrived before inbound handler was set — dropped" + ); + } + } + } + } +} + +/// Dialer side: open the control stream, send Hello{Control}, then exchange. +async fn open_control_stream( + inner: Arc, + peer: MeshPeer, + rx: mpsc::Receiver, +) { + let remote = peer.runtime_id(); + let mut stream = match peer.open_bi().await { + Ok(stream) => stream, + Err(err) => { + tracing::warn!(peer = %remote, %err, "mesh: control stream open failed"); + return; + } + }; + let hello = MeshStreamFrame::Hello(StreamHello { + sender: inner.endpoint.runtime_id(), + role: StreamRole::Control, + }); + if let Err(err) = stream.send_frame(hello).await { + tracing::warn!(peer = %remote, %err, "mesh: control Hello send failed"); + return; + } + control_stream_exchange(inner, remote, stream, rx).await; +} + +/// Both sides: pump queued outbound frames and dispatch inbound gossip. +/// +/// Scuttlebutt: a received `Digest` is answered with a `Delta` of records the +/// digest is missing/behind on; a received `Delta` is applied to membership. +async fn control_stream_exchange( + inner: Arc, + remote: RuntimeId, + stream: MeshStream, + mut rx: mpsc::Receiver, +) { + let MeshStream { mut send, mut recv } = stream; + + let send_inner = Arc::clone(&inner); + let send_task = tokio::spawn(async move { + while let Some(frame) = rx.recv().await { + if let Err(err) = send.send_frame(frame).await { + tracing::debug!(peer = %remote, %err, "mesh: control send ended"); + return; + } + send_inner.membership.record_gossip_frame_sent(remote); + } + }); + + loop { + match recv.recv_frame().await { + Ok(Some(MeshStreamFrame::Gossip { payload })) => { + inner.membership.record_gossip_frame_received(remote); + match decode_message(&payload) { + Ok(GossipMessage::Digest { entries, .. }) => { + let delta = inner.membership.delta_for(&entries); + if let Ok(payload) = encode_message(&delta) { + send_control_frame(&inner, remote, MeshStreamFrame::Gossip { payload }); + } + } + Ok(GossipMessage::Delta { records, .. }) => { + for record in records { + inner.membership.apply_gossip_record(record); + } + } + Err(err) => { + tracing::warn!(peer = %remote, %err, "mesh: bad gossip payload"); + } + } + } + Ok(Some(other)) => { + tracing::warn!(peer = %remote, ?other, "mesh: non-gossip frame on control stream"); + } + Ok(None) | Err(_) => { + tracing::debug!(peer = %remote, "mesh: control stream closed"); + send_task.abort(); + return; + } + } + } +} + +fn send_control_frame(inner: &Inner, remote: RuntimeId, frame: MeshStreamFrame) { + let peers = inner.peers.read().expect("peer lock poisoned"); + if let Some(tx) = peers.get(&remote).and_then(|e| e.control_tx.as_ref()) { + // try_send: gossip is periodic and idempotent — dropping a frame under + // backpressure is strictly better than blocking a recv loop. + let _ = tx.try_send(frame); + } +} + +/// Periodic gossip: refresh the local heartbeat and send a digest on every +/// control stream. Deltas flow back per the exchange loop. +async fn gossip_tick_loop(inner: Arc) { + loop { + tokio::time::sleep(inner.gossip_interval).await; + // Heartbeat: bump the local record so peers' phi accrual sees life. + inner.membership.update_local(|_| {}); + let digest = inner.membership.digest(); + let Ok(payload) = encode_message(&digest) else { + continue; + }; + let targets: Vec = { + let peers = inner.peers.read().expect("peer lock poisoned"); + peers + .iter() + .filter(|(_, e)| e.control_tx.is_some()) + .map(|(id, _)| *id) + .collect() + }; + for remote in targets { + send_control_frame( + &inner, + remote, + MeshStreamFrame::Gossip { + payload: payload.clone(), + }, + ); + } + } +} + +/// Readiness-gated registry heartbeat loop, spawned by the relay after boot. +/// `ready` is the relay-owned readiness predicate (shutdown flag et al.). +pub fn spawn_registry_heartbeat( + registry: ReadyRegistry, + record: ReadyRecord, + ready: Arc bool + Send + Sync>, +) -> JoinHandle<()> { + tokio::spawn(async move { + let mut heartbeat = registry.heartbeat(record); + let interval = registry.refresh_interval(); + loop { + if let Err(err) = heartbeat.tick(ready()).await { + tracing::warn!(%err, "mesh: registry heartbeat tick failed"); + } + tokio::time::sleep(interval).await; + } + }) +} + +#[cfg(test)] +mod tests { + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + use std::sync::Mutex as StdMutex; + use std::time::Duration; + + use iroh::SecretKey; + use tokio::time::timeout; + use uuid::Uuid; + + use super::*; + use crate::gossip::GossipRecord; + use crate::wire::{FencedHeader, Profile}; + + fn loopback_any() -> SocketAddr { + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) + } + + async fn runtime(key_byte: u8) -> (MeshRuntime, Vec) { + let endpoint = MeshEndpoint::bind_with_secret_key( + SecretKey::from_bytes(&[key_byte; 32]), + loopback_any(), + ) + .await + .unwrap(); + let addrs: Vec = endpoint + .addr() + .addrs + .iter() + .filter_map(|ta| match ta { + iroh::TransportAddr::Ip(sock) if sock.ip().is_loopback() => Some(sock.to_string()), + _ => None, + }) + .collect(); + assert!(!addrs.is_empty(), "endpoint must expose a loopback addr"); + let record = GossipRecord::new(endpoint.runtime_id(), addrs.clone(), 1); + let membership = MeshMembership::new(record); + let rt = MeshRuntime::start_with_intervals( + endpoint, + membership, + None, + Duration::from_millis(100), + Duration::from_millis(200), + ); + (rt, addrs) + } + + /// Seed b's record into a's membership so a dials b. + fn seed(a: &MeshRuntime, b: &MeshRuntime, b_addrs: &[String]) { + a.membership().apply_gossip_record(GossipRecord::new( + b.local_runtime_id(), + b_addrs.to_vec(), + 1, + )); + } + + async fn connected_pair() -> (MeshRuntime, MeshRuntime) { + let (a, a_addrs) = runtime(1).await; + let (b, b_addrs) = runtime(2).await; + // Both directions: with no registry to rescan, the acceptor's + // admission gate requires the dialer to already be in its membership + // table (production gets this from the attested ready registry). + seed(&a, &b, &b_addrs); + seed(&b, &a, &a_addrs); + a.reconcile_now().await; + // Wait for both sides to see the connection. + timeout(Duration::from_secs(5), async { + loop { + if a.connected_peers().contains(&b.local_runtime_id()) + && b.connected_peers().contains(&a.local_runtime_id()) + { + return; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("mesh pair should connect"); + (a, b) + } + + struct RecordingHandler { + datagrams: StdMutex>, + streams: StdMutex>, + } + + impl RecordingHandler { + fn new() -> Arc { + Arc::new(Self { + datagrams: StdMutex::new(Vec::new()), + streams: StdMutex::new(Vec::new()), + }) + } + } + + impl InboundHandler for Arc { + fn on_datagram(&self, from: RuntimeId, dgram: MeshDatagram) { + self.datagrams.lock().unwrap().push((from, dgram)); + } + fn on_session_stream(&self, from: RuntimeId, hello: StreamHello, _stream: MeshStream) { + self.streams.lock().unwrap().push((from, hello)); + } + } + + fn fenced(owner: RuntimeId) -> FencedHeader { + FencedHeader { + session_id: Uuid::from_u128(0xFEED), + generation: 3, + owner_runtime_id: owner, + } + } + + #[tokio::test] + async fn warm_pair_connects_and_gossips_membership() { + let (a, b) = connected_pair().await; + // Gossip heartbeats should keep flowing; wait for a to see a gossiped + // (version > 1) record from b. + timeout(Duration::from_secs(5), async { + loop { + let seen = a + .membership() + .records() + .into_iter() + .find(|r| r.runtime_id == b.local_runtime_id()); + if seen.is_some_and(|r| r.version > 1) { + return; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("gossip should refresh b's record on a"); + a.shutdown(); + b.shutdown(); + } + + #[tokio::test] + async fn transport_datagram_reaches_inbound_handler() { + let (a, b) = connected_pair().await; + let handler = RecordingHandler::new(); + b.set_inbound(Box::new(Arc::clone(&handler))); + + let dgram = MeshDatagram { + fenced: fenced(b.local_runtime_id()), + seq: 7, + payload: vec![1, 2, 3], + }; + a.send_datagram(b.local_runtime_id(), dgram.clone()) + .unwrap(); + + timeout(Duration::from_secs(5), async { + loop { + if !handler.datagrams.lock().unwrap().is_empty() { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("datagram should arrive"); + let got = handler.datagrams.lock().unwrap(); + assert_eq!(got[0].0, a.local_runtime_id()); + assert_eq!(got[0].1, dgram); + drop(got); + a.shutdown(); + b.shutdown(); + } + + #[tokio::test] + async fn transport_session_stream_reaches_inbound_handler() { + let (a, b) = connected_pair().await; + let handler = RecordingHandler::new(); + b.set_inbound(Box::new(Arc::clone(&handler))); + + let hello = StreamHello { + sender: a.local_runtime_id(), + role: StreamRole::Session { + fenced: fenced(b.local_runtime_id()), + profile: Profile::ReliableStream, + }, + }; + let mut stream = a + .open_session_stream(b.local_runtime_id(), hello.clone()) + .await + .unwrap(); + stream + .send_frame(MeshStreamFrame::Data { + fenced: fenced(b.local_runtime_id()), + payload: b"tunnel".to_vec(), + }) + .await + .unwrap(); + + timeout(Duration::from_secs(5), async { + loop { + if !handler.streams.lock().unwrap().is_empty() { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("session stream should arrive"); + let got = handler.streams.lock().unwrap(); + assert_eq!(got[0].0, a.local_runtime_id()); + assert_eq!(got[0].1, hello); + drop(got); + a.shutdown(); + b.shutdown(); + } + + #[tokio::test] + async fn send_to_unconnected_peer_is_typed_error() { + let (a, _addrs) = runtime(9).await; + let ghost = RuntimeId([42u8; 32]); + let err = a + .send_datagram( + ghost, + MeshDatagram { + fenced: fenced(ghost), + seq: 0, + payload: vec![], + }, + ) + .unwrap_err(); + assert!(matches!(err, MeshError::PeerNotConnected(id) if id == ghost)); + a.shutdown(); + } + + #[tokio::test] + async fn simultaneous_dial_converges_to_one_connection() { + let (a, a_addrs) = runtime(3).await; + let (b, b_addrs) = runtime(4).await; + seed(&a, &b, &b_addrs); + seed(&b, &a, &a_addrs); + // Both dial at once. + tokio::join!(a.reconcile_now(), b.reconcile_now()); + + timeout(Duration::from_secs(5), async { + loop { + if a.connected_peers().contains(&b.local_runtime_id()) + && b.connected_peers().contains(&a.local_runtime_id()) + { + return; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("simultaneous dial should converge"); + // Datagrams still flow after the tie-break. + let handler = RecordingHandler::new(); + b.set_inbound(Box::new(Arc::clone(&handler))); + // The surviving connection may need a beat to settle. + timeout(Duration::from_secs(5), async { + loop { + let dgram = MeshDatagram { + fenced: fenced(b.local_runtime_id()), + seq: 1, + payload: vec![9], + }; + let _ = a.send_datagram(b.local_runtime_id(), dgram); + if !handler.datagrams.lock().unwrap().is_empty() { + return; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + }) + .await + .expect("datagram should flow after tie-break"); + a.shutdown(); + b.shutdown(); + } + + #[test] + fn tie_break_is_symmetric() { + let small = RuntimeId([1u8; 32]); + let large = RuntimeId([2u8; 32]); + // small dials large: small's outbound wins, large's inbound wins. + assert!(new_connection_wins(small, large, true)); + assert!(new_connection_wins(large, small, false)); + // large dials small: loses on both ends. + assert!(!new_connection_wins(large, small, true)); + assert!(!new_connection_wins(small, large, false)); + } +} diff --git a/crates/buzz-relay/src/lib.rs b/crates/buzz-relay/src/lib.rs index 3a1c68f3f..b3ea65db5 100644 --- a/crates/buzz-relay/src/lib.rs +++ b/crates/buzz-relay/src/lib.rs @@ -19,6 +19,8 @@ pub mod connection; pub mod error; /// WebSocket message handlers for NIP-01 client commands. pub mod handlers; +/// Inter-relay mesh startup wiring (`BUZZ_MESH` seam). +pub mod mesh_boot; /// Relay-signed mesh-LLM status publisher. pub mod mesh_status_publisher; /// Prometheus metrics: recorder, upkeep, HTTP middleware. diff --git a/crates/buzz-relay/src/main.rs b/crates/buzz-relay/src/main.rs index 973b7f68b..d8202571c 100644 --- a/crates/buzz-relay/src/main.rs +++ b/crates/buzz-relay/src/main.rs @@ -334,6 +334,26 @@ async fn main() -> anyhow::Result<()> { ); let state = Arc::new(app_state); + // Inter-relay mesh (BUZZ_MESH seam). `boot_mesh` returns None when the + // kill switch is off — nothing is bound, published, or spawned, so the + // relay behaves byte-identically to a build without the mesh. When + // enabled, a misconfigured mesh is fatal here (bind/Redis failure): an + // operator who asked for the mesh gets it or gets told why not. + if let Some(handle) = buzz_relay::mesh_boot::boot_mesh( + &state.config, + state.redis_pool.clone(), + &state.relay_keypair, + Arc::clone(&state.shutting_down), + ) + .await? + { + let runtime_id = handle.local_runtime_id; + if state.mesh.set(handle).is_err() { + unreachable!("mesh handle is set exactly once, right here"); + } + info!(runtime_id = %runtime_id, "Inter-relay mesh started"); + } + // Git-on-object-storage: admit the configured S3/MinIO backend against the // linearizable conditional-write axiom (A3) before serving git traffic. // Failure is fatal: a backend that cannot satisfy pointer CAS invalidates diff --git a/crates/buzz-relay/src/mesh_boot.rs b/crates/buzz-relay/src/mesh_boot.rs new file mode 100644 index 000000000..bd7490f71 --- /dev/null +++ b/crates/buzz-relay/src/mesh_boot.rs @@ -0,0 +1,214 @@ +//! Relay startup wiring for the inter-relay mesh (`BUZZ_MESH` seam). +//! +//! [`boot_mesh`] is the ONLY place the relay constructs mesh machinery. It +//! returns `None` — and touches nothing — when `BUZZ_MESH=off`, so mesh-off +//! deployments stay byte-identical to a relay built before this module +//! existed. When enabled, it: +//! +//! 1. binds the iroh endpoint on `BUZZ_MESH_BIND_ADDR` (boot-unique keypair = +//! boot-unique `RuntimeId`), +//! 2. publishes a relay-key-attested [`ReadyRecord`] to the Redis ready +//! registry and starts the readiness-gated heartbeat, +//! 3. starts the [`MeshRuntime`] loops (accept, reconcile/dial, gossip) and +//! runs one immediate reconcile pass so seed peers are dialed at boot, +//! 4. spawns a drain watcher: when the relay's `shutting_down` flag flips, +//! membership gossips `draining=true` and the heartbeat clears the +//! registry record. +//! +//! Consumers (huddle control plane, reliable-stream tunnels) reach the mesh +//! exclusively through [`MeshHandle`] via `AppState::mesh()` — `None` means +//! "behave exactly like a single-instance relay." + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; + +use buzz_relay_mesh::endpoint::MeshEndpoint; +use buzz_relay_mesh::gossip::GossipRecord; +use buzz_relay_mesh::registry::{ReadyRecord, ReadyRegistry}; +use buzz_relay_mesh::{ + MeshMembership, MeshRuntime, MeshStatus, RelayMeshMembership, RelayPeerTransport, RuntimeId, +}; + +use crate::config::Config; +use crate::tunnel::directory::SessionDirectory; + +/// Everything a mesh consumer needs, as one bundle. +#[derive(Clone)] +pub struct MeshHandle { + /// Redis fenced session directory — the ownership arbiter. + pub directory: SessionDirectory, + /// Fenced byte transport to peer runtimes. + pub transport: Arc, + /// Routing hints: who is alive / draining / dialable. + pub membership: Arc, + /// This runtime's boot-unique mesh identity. + pub local_runtime_id: RuntimeId, + /// The running mesh (status snapshots, shutdown). + runtime: MeshRuntime, +} + +impl MeshHandle { + /// Live `/_mesh` status snapshot. + pub fn status(&self) -> MeshStatus { + self.runtime.membership().status() + } +} + +/// Wire protocol version advertised in registry/gossip records. +const PROTO_VERSION: u16 = buzz_relay_mesh::WIRE_VERSION as u16; + +/// Capabilities advertised by this build. All three tunnel profiles ship in +/// the same binary, so the list is static. +fn capabilities() -> Vec { + vec![ + "reliable-stream".to_string(), + "realtime-media".to_string(), + "huddle-control".to_string(), + ] +} + +/// Addresses peers should dial, in preference order: +/// `BUZZ_MESH_ADVERTISE_ADDR` (explicit, classic-LB shapes) → +/// `POD_IP` + actual bound port (k8s Downward API, zero RBAC) → +/// every IP transport addr the endpoint reports (dev/local). +fn advertise_addrs(endpoint: &MeshEndpoint) -> Vec { + if let Ok(addr) = std::env::var("BUZZ_MESH_ADVERTISE_ADDR") { + let addr = addr.trim().to_string(); + if !addr.is_empty() { + return vec![addr]; + } + } + + let ip_addrs = endpoint.ip_addrs(); + let bound_port = ip_addrs.first().map(|sock| sock.port()).unwrap_or(0); + + if let Ok(pod_ip) = std::env::var("POD_IP") { + let pod_ip = pod_ip.trim(); + if !pod_ip.is_empty() && bound_port != 0 { + return vec![format!("{pod_ip}:{bound_port}")]; + } + } + + ip_addrs.iter().map(|sock| sock.to_string()).collect() +} + +/// Boot the mesh, or return `None` when `BUZZ_MESH=off`. +/// +/// Never fatal to relay startup by policy? No — a *misconfigured* enabled mesh +/// fails loudly (bind failure, Redis unreachable at publish). An operator who +/// sets `BUZZ_MESH=on` wants the mesh or wants to know why not; silently +/// booting meshless would be the same class of bug as silently dropping to a +/// default tenant. +pub async fn boot_mesh( + config: &Config, + redis_pool: deadpool_redis::Pool, + relay_keypair: &nostr::Keys, + shutting_down: Arc, +) -> anyhow::Result> { + if !config.mesh.enabled { + tracing::info!("mesh disabled (BUZZ_MESH=off) — single-instance behavior"); + return Ok(None); + } + + let endpoint = MeshEndpoint::bind(config.mesh.bind_addr) + .await + .map_err(|e| { + anyhow::anyhow!( + "mesh endpoint bind on {} failed: {e}", + config.mesh.bind_addr + ) + })?; + let runtime_id = endpoint.runtime_id(); + let addrs = advertise_addrs(&endpoint); + tracing::info!( + runtime_id = %runtime_id, + bind_addr = %config.mesh.bind_addr, + advertise_addrs = ?addrs, + "mesh endpoint bound" + ); + + let mut local_record = GossipRecord::new(runtime_id, addrs.clone(), PROTO_VERSION); + local_record.capabilities = capabilities(); + let membership = MeshMembership::new(local_record); + + let registry = ReadyRegistry::new(redis_pool.clone(), config.mesh.registry_refresh); + let ready_record = ReadyRecord::new( + runtime_id, + relay_keypair, + addrs, + PROTO_VERSION, + capabilities(), + ); + + // First publish is part of boot: if Redis can't take the attested record, + // peers can never find us — fail loudly now, not quietly forever. + registry + .publish_ready(&ready_record) + .await + .map_err(|e| anyhow::anyhow!("mesh ready-registry publish failed: {e}"))?; + tracing::info!(runtime_id = %runtime_id, "mesh ready record published"); + + // Readiness-gated heartbeat: publishes while the relay would pass + // readiness, clears the record on ready→not-ready and on shutdown. + let hb_flag = Arc::clone(&shutting_down); + buzz_relay_mesh::runtime::spawn_registry_heartbeat( + registry.clone(), + ready_record, + Arc::new(move || !hb_flag.load(Ordering::Relaxed)), + ); + + let runtime = MeshRuntime::start(endpoint, membership, Some(registry)); + // Dial seed peers now rather than waiting for the first reconcile tick. + runtime.reconcile_now().await; + + // Drain watcher: SIGTERM flips `shutting_down`; gossip `draining=true` so + // peers stop routing new sessions here while in-flight ones drain. + { + let runtime = runtime.clone(); + let flag = shutting_down; + tokio::spawn(async move { + loop { + if flag.load(Ordering::Relaxed) { + runtime.membership().begin_drain(); + tracing::info!("mesh drain started (draining=true gossiped)"); + return; + } + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + } + }); + } + + let membership_arc: Arc = Arc::new(runtime.membership().clone()); + let transport: Arc = Arc::new(runtime.clone()); + + Ok(Some(MeshHandle { + directory: SessionDirectory::new(redis_pool), + transport, + membership: membership_arc, + local_runtime_id: runtime_id, + runtime, + })) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// BUZZ_MESH=off must be a hard no-op: no endpoint bind, no Redis write, + /// no background task — `boot_mesh` returns `None` before touching + /// anything. The Redis pool here points nowhere routable; if the off path + /// ever reached Redis this test would hang/fail. + #[tokio::test] + async fn mesh_off_boots_nothing() { + let mut config = crate::config::Config::from_env().expect("default config loads"); + config.mesh.enabled = false; + let pool = deadpool_redis::Config::from_url("redis://127.0.0.1:1") // unroutable + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .unwrap(); + let keys = nostr::Keys::generate(); + let handle = boot_mesh(&config, pool, &keys, Arc::new(AtomicBool::new(false))) + .await + .expect("off path is never an error"); + assert!(handle.is_none()); + } +} diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index fc9e1ec38..dbf69f92e 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -129,6 +129,7 @@ pub fn build_health_router(state: Arc) -> Router { .route("/_liveness", get(liveness_handler)) .route("/_readiness", get(readiness_handler)) .route("/_status", get(status_handler)) + .route("/_mesh", get(mesh_status_handler)) .with_state(state) } @@ -251,6 +252,18 @@ async fn status_handler(State(state): State>) -> impl IntoResponse })) } +/// `/_mesh` — live mesh status: peer table, connection/phi state, per-peer +/// counters, fence-rejection totals. Mesh-off reports `{"enabled": false}` so +/// operators can distinguish "off" from "on with zero peers". +async fn mesh_status_handler(State(state): State>) -> impl IntoResponse { + match state.mesh() { + Some(handle) => Json(serde_json::to_value(handle.status()).unwrap_or_else( + |e| json!({"enabled": true, "error": format!("status serialize: {e}")}), + )), + None => Json(json!({"enabled": false})), + } +} + /// Build a CORS layer from the configured origins list. fn build_cors_layer(cors_origins: &[String]) -> CorsLayer { if cors_origins.is_empty() { diff --git a/crates/buzz-relay/src/state.rs b/crates/buzz-relay/src/state.rs index a674e6d7c..adf273e64 100644 --- a/crates/buzz-relay/src/state.rs +++ b/crates/buzz-relay/src/state.rs @@ -318,6 +318,13 @@ pub struct AppState { /// See `crates/buzz-conformance/` and `crate::conformance` for the /// schema, emitter helpers, and the independent checker. pub tracer: Arc, + + /// Inter-relay mesh handle, set once by `main.rs` after `mesh_boot` (never + /// a constructor parameter, so `AppState::new` call sites are untouched). + /// `None`/unset ⇒ mesh-off / single-instance: consumers must behave + /// byte-identically to a relay without the mesh. Access via + /// [`AppState::mesh`]. + pub mesh: Arc>, } impl AppState { @@ -453,6 +460,7 @@ impl AppState { // construction (see test helpers in // `crates/buzz-test-client` once those land). tracer: Arc::new(crate::conformance::NoopTracer), + mesh: Arc::new(std::sync::OnceLock::new()), }; ( state, @@ -463,6 +471,12 @@ impl AppState { ) } + /// Inter-relay mesh handle. `None` ⇒ mesh-off / single-instance: callers + /// must no-op to today's behavior. Set once by `main.rs` after boot. + pub fn mesh(&self) -> Option<&crate::mesh_boot::MeshHandle> { + self.mesh.get() + } + /// Record an event ID as locally-published for dedup, scoped to the /// community it was fanned out in. Called before Redis publish so the /// multi-node consumer can skip the echo for *this* community only — a