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 <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d
2026-07-08 12:43:32 -04:00
co-authored by Tyler Longwell
parent 59a7ce3617
commit 8b077fdb44
9 changed files with 1240 additions and 0 deletions
+14
View File
@@ -56,6 +56,20 @@ impl MeshEndpoint {
self.endpoint.addr() 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<SocketAddr> {
self.endpoint
.addr()
.addrs
.iter()
.filter_map(|ta| match ta {
TransportAddr::Ip(sock) => Some(*sock),
_ => None,
})
.collect()
}
pub async fn accept(&self) -> Result<Option<crate::peer::MeshPeer>, MeshError> { pub async fn accept(&self) -> Result<Option<crate::peer::MeshPeer>, MeshError> {
let Some(incoming) = self.endpoint.accept().await else { let Some(incoming) = self.endpoint.accept().await else {
return Ok(None); return Ok(None);
+2
View File
@@ -23,6 +23,7 @@ pub mod gossip;
pub mod membership; pub mod membership;
pub mod peer; pub mod peer;
pub mod registry; pub mod registry;
pub mod runtime;
pub mod status; pub mod status;
pub mod wire; pub mod wire;
@@ -41,6 +42,7 @@ use bytes::Bytes;
pub use gossip::{GossipDigestEntry, GossipMessage, GossipRecord, GossipState, PhiAccrual}; pub use gossip::{GossipDigestEntry, GossipMessage, GossipRecord, GossipState, PhiAccrual};
pub use membership::MeshMembership; pub use membership::MeshMembership;
pub use registry::{ReadyHeartbeat, ReadyRecord, ReadyRegistry, RuntimeAttestation}; pub use registry::{ReadyHeartbeat, ReadyRecord, ReadyRegistry, RuntimeAttestation};
pub use runtime::MeshRuntime;
pub use status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus}; pub use status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus};
pub use wire::{ pub use wire::{
FencedHeader, GoodbyeReason, MeshDatagram, MeshStreamFrame, Profile, RuntimeId, StreamHello, FencedHeader, GoodbyeReason, MeshDatagram, MeshStreamFrame, Profile, RuntimeId, StreamHello,
+65
View File
@@ -146,6 +146,71 @@ impl MeshMembership {
self.draining.load(Ordering::Relaxed) 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<GossipRecord> {
let mut records: Vec<GossipRecord> = 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) { pub fn record_stream_opened(&self, runtime_id: RuntimeId) {
self.update_peer_counters(runtime_id, |c| { self.update_peer_counters(runtime_id, |c| {
c.streams_opened = c.streams_opened.saturating_add(1) c.streams_opened = c.streams_opened.saturating_add(1)
+896
View File
@@ -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<mpsc::Sender<MeshStreamFrame>>,
tasks: Vec<JoinHandle<()>>,
}
impl PeerEntry {
fn abort(&self) {
for task in &self.tasks {
task.abort();
}
}
}
struct Inner {
endpoint: MeshEndpoint,
membership: MeshMembership,
registry: Option<ReadyRegistry>,
peers: RwLock<HashMap<RuntimeId, PeerEntry>>,
handler: Mutex<Option<Arc<dyn InboundHandler>>>,
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<Inner>,
loops: Arc<Mutex<Vec<JoinHandle<()>>>>,
}
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<ReadyRegistry>,
) -> 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<ReadyRegistry>,
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<RuntimeId> {
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<MeshStream, MeshError>> {
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<dyn InboundHandler>) {
*self.inner.handler.lock().expect("handler lock poisoned") = Some(Arc::from(handler));
}
}
fn inbound_handler(inner: &Inner) -> Option<Arc<dyn InboundHandler>> {
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<Inner>, 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<Inner>) {
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<Inner>, 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<Inner>) {
loop {
reconcile_once(&inner).await;
tokio::time::sleep(inner.reconcile_interval).await;
}
}
async fn reconcile_once(inner: &Arc<Inner>) {
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<Inner>, 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<Inner>, 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<Inner>, 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<Inner>,
peer: MeshPeer,
rx: mpsc::Receiver<MeshStreamFrame>,
) {
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<Inner>,
remote: RuntimeId,
stream: MeshStream,
mut rx: mpsc::Receiver<MeshStreamFrame>,
) {
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<Inner>) {
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<RuntimeId> = {
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<dyn Fn() -> 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<String>) {
let endpoint = MeshEndpoint::bind_with_secret_key(
SecretKey::from_bytes(&[key_byte; 32]),
loopback_any(),
)
.await
.unwrap();
let addrs: Vec<String> = 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<Vec<(RuntimeId, MeshDatagram)>>,
streams: StdMutex<Vec<(RuntimeId, StreamHello)>>,
}
impl RecordingHandler {
fn new() -> Arc<Self> {
Arc::new(Self {
datagrams: StdMutex::new(Vec::new()),
streams: StdMutex::new(Vec::new()),
})
}
}
impl InboundHandler for Arc<RecordingHandler> {
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));
}
}
+2
View File
@@ -19,6 +19,8 @@ pub mod connection;
pub mod error; pub mod error;
/// WebSocket message handlers for NIP-01 client commands. /// WebSocket message handlers for NIP-01 client commands.
pub mod handlers; pub mod handlers;
/// Inter-relay mesh startup wiring (`BUZZ_MESH` seam).
pub mod mesh_boot;
/// Relay-signed mesh-LLM status publisher. /// Relay-signed mesh-LLM status publisher.
pub mod mesh_status_publisher; pub mod mesh_status_publisher;
/// Prometheus metrics: recorder, upkeep, HTTP middleware. /// Prometheus metrics: recorder, upkeep, HTTP middleware.
+20
View File
@@ -334,6 +334,26 @@ async fn main() -> anyhow::Result<()> {
); );
let state = Arc::new(app_state); 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 // Git-on-object-storage: admit the configured S3/MinIO backend against the
// linearizable conditional-write axiom (A3) before serving git traffic. // linearizable conditional-write axiom (A3) before serving git traffic.
// Failure is fatal: a backend that cannot satisfy pointer CAS invalidates // Failure is fatal: a backend that cannot satisfy pointer CAS invalidates
+214
View File
@@ -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<dyn RelayPeerTransport>,
/// Routing hints: who is alive / draining / dialable.
pub membership: Arc<dyn RelayMeshMembership>,
/// 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<String> {
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<String> {
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<AtomicBool>,
) -> anyhow::Result<Option<MeshHandle>> {
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<dyn RelayMeshMembership> = Arc::new(runtime.membership().clone());
let transport: Arc<dyn RelayPeerTransport> = 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());
}
}
+13
View File
@@ -129,6 +129,7 @@ pub fn build_health_router(state: Arc<AppState>) -> Router {
.route("/_liveness", get(liveness_handler)) .route("/_liveness", get(liveness_handler))
.route("/_readiness", get(readiness_handler)) .route("/_readiness", get(readiness_handler))
.route("/_status", get(status_handler)) .route("/_status", get(status_handler))
.route("/_mesh", get(mesh_status_handler))
.with_state(state) .with_state(state)
} }
@@ -251,6 +252,18 @@ async fn status_handler(State(state): State<Arc<AppState>>) -> 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<Arc<AppState>>) -> 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. /// Build a CORS layer from the configured origins list.
fn build_cors_layer(cors_origins: &[String]) -> CorsLayer { fn build_cors_layer(cors_origins: &[String]) -> CorsLayer {
if cors_origins.is_empty() { if cors_origins.is_empty() {
+14
View File
@@ -318,6 +318,13 @@ pub struct AppState {
/// See `crates/buzz-conformance/` and `crate::conformance` for the /// See `crates/buzz-conformance/` and `crate::conformance` for the
/// schema, emitter helpers, and the independent checker. /// schema, emitter helpers, and the independent checker.
pub tracer: Arc<dyn buzz_conformance::Tracer>, pub tracer: Arc<dyn buzz_conformance::Tracer>,
/// 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<std::sync::OnceLock<crate::mesh_boot::MeshHandle>>,
} }
impl AppState { impl AppState {
@@ -453,6 +460,7 @@ impl AppState {
// construction (see test helpers in // construction (see test helpers in
// `crates/buzz-test-client` once those land). // `crates/buzz-test-client` once those land).
tracer: Arc::new(crate::conformance::NoopTracer), tracer: Arc::new(crate::conformance::NoopTracer),
mesh: Arc::new(std::sync::OnceLock::new()),
}; };
( (
state, 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 /// Record an event ID as locally-published for dedup, scoped to the
/// community it was fanned out in. Called before Redis publish so 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 /// multi-node consumer can skip the echo for *this* community only — a