feat(mesh): buzz-relay-mesh crate skeleton — frozen wire contract + relay seams

The contract every mesh lane builds against:

- wire.rs: ALPN buzz/mesh/1, WIRE_VERSION byte + postcard framing;
  FencedHeader {session_id, generation, owner_runtime_id} on every
  session-bearing frame (fencing law); MeshDatagram (realtime-media),
  MeshStreamFrame (Hello/Data/Goodbye/Gossip) with u32-LE length-delimited
  bi-stream framing; Profile {ReliableStream, RealtimeMedia, HuddleControl}.
- lib.rs: MeshConfig (BUZZ_MESH kill switch), MeshError taxonomy, and the
  two relay seams — RelayMeshMembership (hints only) and RelayPeerTransport
  (fenced bytes; size/version checks in transport, generation fencing in the
  session layer at every hop).

Unit tests pin roundtrips, unknown-version rejection, and a 64B datagram
header overhead budget (Opus @20ms + header must clear the ~1200B QUIC
datagram floor).

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 10:04:10 -04:00
co-authored by Tyler Longwell
parent 5dce80aff5
commit 5819fde839
5 changed files with 529 additions and 0 deletions
Generated
+55
View File
@@ -366,6 +366,15 @@ version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ef49f5882e4b6afaac09ad239a4f8c70a24b8f2b0897edb1f706008efd109cf4"
[[package]]
name = "atomic-polyfill"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8cf2bce30dfe09ef0bfaef228b9d414faaf7e563035494d7fe092dba54b300f4"
dependencies = [
"critical-section",
]
[[package]]
name = "atomic-waker"
version = "1.1.2"
@@ -1084,6 +1093,28 @@ dependencies = [
"uuid",
]
[[package]]
name = "buzz-relay-mesh"
version = "0.1.0"
dependencies = [
"bytes",
"deadpool-redis",
"futures-util",
"hex",
"hmac 0.13.0",
"iroh",
"postcard",
"proptest",
"redis",
"serde",
"serde_json",
"sha2 0.11.0",
"thiserror 2.0.18",
"tokio",
"tracing",
"uuid",
]
[[package]]
name = "buzz-sdk"
version = "0.1.0"
@@ -2851,6 +2882,15 @@ dependencies = [
"tracing",
]
[[package]]
name = "hash32"
version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b0c35f58762feb77d74ebe43bdbc3210f09be9fe6742234d573bacc26ed92b67"
dependencies = [
"byteorder",
]
[[package]]
name = "hashbag"
version = "0.1.13"
@@ -2909,6 +2949,20 @@ version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0049b265b7f201ca9ab25475b22b47fe444060126a51abe00f77d986fc5cc52e"
[[package]]
name = "heapless"
version = "0.7.17"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cdc6457c0eb62c71aac4bc17216026d8410337c4126773b9c5daba343f17964f"
dependencies = [
"atomic-polyfill",
"hash32",
"rustc_version",
"serde",
"spin 0.9.8",
"stable_deref_trait",
]
[[package]]
name = "heck"
version = "0.5.0"
@@ -5909,6 +5963,7 @@ dependencies = [
"cobs",
"embedded-io 0.4.0",
"embedded-io 0.6.1",
"heapless",
"postcard-derive",
"serde",
]
+5
View File
@@ -23,6 +23,7 @@ members = [
"crates/git-credential-nostr",
"crates/git-sign-nostr",
"crates/buzz-pair-relay",
"crates/buzz-relay-mesh",
"crates/buzz-dev-mcp",
"examples/countdown-bot",
]
@@ -60,6 +61,10 @@ nostr = { version = "0.44", features = ["nip44", "nip98"] }
# Serialization
serde = { version = "1", features = ["derive"] }
postcard = { version = "1", features = ["use-std"] }
# Inter-relay mesh transport (buzz-relay-mesh)
iroh = { version = "1.0.0-rc.0", default-features = false }
serde_json = "1"
serde_yaml = "0.9"
evalexpr = "11"
+29
View File
@@ -0,0 +1,29 @@
[package]
name = "buzz-relay-mesh"
description = "Inter-relay QUIC mesh: transport, membership, and the fenced wire contract"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
repository.workspace = true
[dependencies]
tokio = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
postcard = { workspace = true }
iroh = { workspace = true }
redis = { workspace = true }
deadpool-redis = { workspace = true }
thiserror = { workspace = true }
tracing = { workspace = true }
uuid = { workspace = true }
hmac = { workspace = true }
sha2 = { workspace = true }
hex = { workspace = true }
bytes = "1"
futures-util = { workspace = true }
[dev-dependencies]
tokio = { workspace = true, features = ["test-util"] }
proptest = { workspace = true }
+179
View File
@@ -0,0 +1,179 @@
//! buzz-relay-mesh — the inter-relay QUIC mesh.
//!
//! One iroh endpoint per relay runtime (identity = the relay's signing key),
//! a warm full mesh of authenticated connections, scuttlebutt membership
//! gossip on a control substream, and a fenced wire contract that carries
//! tunnel traffic (reliable streams + realtime datagrams) between pods.
//!
//! The relay consumes this crate exclusively through two seams:
//!
//! - [`RelayMeshMembership`] — "who is alive / draining / dialable?"
//! - [`RelayPeerTransport`] — "move these bytes to that runtime."
//!
//! The seams are what keep single-instance deployments and same-pod sessions
//! mesh-free: when `BUZZ_MESH=off` or no peers exist, the relay never
//! constructs a mesh and the in-process fast path is untouched.
//!
//! **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 wire;
// Lane modules — one owner per file (see the mesh thread for lane map):
// endpoint.rs, peer.rs — Mari (transport core)
// registry.rs, gossip.rs,
// membership.rs, status.rs — Max (membership + /_mesh)
// Session directory + tunnel routing live relay-side (Perci), consuming the
// seams below; huddle fan-out lives in buzz-relay's audio module (Dawn).
use std::future::Future;
use std::pin::Pin;
use bytes::Bytes;
pub use wire::{
FencedHeader, GoodbyeReason, MeshDatagram, MeshStreamFrame, Profile, RuntimeId, StreamHello,
StreamRole, ALPN, WIRE_VERSION,
};
/// Mesh configuration, resolved from env by the relay.
#[derive(Clone, Debug)]
pub struct MeshConfig {
/// `BUZZ_MESH` — `on` (default when replicas can exist) | `off` kill
/// switch. When off, the relay must behave exactly like single-instance.
pub enabled: bool,
/// UDP bind for the iroh endpoint (`BUZZ_MESH_BIND_ADDR`, default
/// `0.0.0.0:3478`). Excluded from istio sidecar capture in k8s.
pub bind_addr: std::net::SocketAddr,
/// Ready-registry heartbeat refresh (default 15s; expiry is 3x).
pub registry_refresh: std::time::Duration,
}
#[derive(Debug, thiserror::Error)]
pub enum MeshError {
#[error("frame encode: {0}")]
Encode(#[source] postcard::Error),
#[error("frame decode: {0}")]
Decode(#[source] postcard::Error),
#[error("unknown wire version {0}")]
UnknownWireVersion(u8),
#[error("empty frame")]
EmptyFrame,
#[error("frame exceeds max size ({size} > {max})")]
FrameTooLarge { size: usize, max: usize },
#[error("datagram exceeds connection max_datagram_size ({size} > {max})")]
DatagramTooLarge { size: usize, max: usize },
#[error("peer {0} not connected")]
PeerNotConnected(RuntimeId),
#[error("peer {0} is draining")]
PeerDraining(RuntimeId),
#[error("stale generation for session {session_id}: frame {frame_generation} < known {known_generation}")]
StaleGeneration {
session_id: uuid::Uuid,
frame_generation: u64,
known_generation: u64,
},
#[error("mesh is disabled (BUZZ_MESH=off)")]
Disabled,
#[error("transport: {0}")]
Transport(String),
#[error("redis: {0}")]
Redis(#[from] redis::RedisError),
}
/// A peer as membership sees it. Everything here is a routing HINT.
#[derive(Clone, Debug)]
pub struct PeerInfo {
pub runtime_id: RuntimeId,
pub draining: bool,
/// Phi-accrual suspicion; `None` until enough heartbeats observed.
pub phi: Option<f64>,
/// Advisory load factor gossiped by the peer (0.0..).
pub load: f32,
}
type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
/// Seam 1: membership. Answers "who can I route to?" — never "who owns what."
pub trait RelayMeshMembership: Send + Sync + 'static {
/// Live, non-suspect peers (self excluded).
fn peers(&self) -> Vec<PeerInfo>;
/// This runtime's mesh identity.
fn local_runtime_id(&self) -> RuntimeId;
/// Begin drain: gossip `draining=true`, stop accepting new sessions.
fn begin_drain(&self);
}
/// Seam 2: transport. Moves fenced bytes to a specific runtime.
///
/// Implementations perform the datagram-size and wire-version checks; they do
/// NOT perform generation fencing — that belongs to the session layer on both
/// ends (fencing at every hop means every consumer checks, not the pipe).
pub trait RelayPeerTransport: Send + Sync + 'static {
/// Fire-and-forget realtime datagram (drop-on-full, never blocks on old
/// audio). Errors only for disconnected peer / oversize frame.
fn send_datagram(&self, to: RuntimeId, dgram: MeshDatagram) -> Result<(), MeshError>;
/// Open a reliable bi-stream to a peer for a session (`ReliableStream`
/// or `HuddleControl` profile). Sends the `Hello` before returning.
fn open_session_stream(
&self,
to: RuntimeId,
hello: StreamHello,
) -> BoxFuture<'_, Result<MeshStream, MeshError>>;
/// Register the handler invoked for inbound datagrams / session streams.
/// Called once at relay startup.
fn set_inbound(&self, handler: Box<dyn InboundHandler>);
}
/// Inbound mesh traffic, delivered after wire decode + Hello validation.
pub trait InboundHandler: Send + Sync + 'static {
fn on_datagram(&self, from: RuntimeId, dgram: MeshDatagram);
fn on_session_stream(&self, from: RuntimeId, hello: StreamHello, stream: MeshStream);
}
/// A reliable mesh stream: length-delimited `MeshStreamFrame`s over QUIC.
/// Concrete type (not a trait) so lanes share one framing implementation.
pub struct MeshStream {
// Mari: wrap iroh SendStream/RecvStream with the u32-LE length framing
// from `wire`. Placeholder halves keep the seam compilable pre-transport.
pub(crate) send: Box<dyn StreamSendHalf>,
pub(crate) recv: Box<dyn StreamRecvHalf>,
}
pub trait StreamSendHalf: Send + 'static {
fn send_frame(&mut self, frame: MeshStreamFrame) -> BoxFuture<'_, Result<(), MeshError>>;
fn finish(&mut self) -> Result<(), MeshError>;
}
pub trait StreamRecvHalf: Send + 'static {
fn recv_frame(&mut self) -> BoxFuture<'_, Result<Option<MeshStreamFrame>, MeshError>>;
}
impl MeshStream {
pub fn send_frame(&mut self, frame: MeshStreamFrame) -> BoxFuture<'_, Result<(), MeshError>> {
self.send.send_frame(frame)
}
pub fn recv_frame(&mut self) -> BoxFuture<'_, Result<Option<MeshStreamFrame>, MeshError>> {
self.recv.recv_frame()
}
pub fn finish(&mut self) -> Result<(), MeshError> {
self.send.finish()
}
}
/// Raw bytes helper used by transport internals.
pub fn encode_datagram_checked(
dgram: &MeshDatagram,
max_datagram_size: usize,
) -> Result<Bytes, MeshError> {
let bytes = wire::encode(dgram)?;
if bytes.len() > max_datagram_size {
return Err(MeshError::DatagramTooLarge {
size: bytes.len(),
max: max_datagram_size,
});
}
Ok(Bytes::from(bytes))
}
+261
View File
@@ -0,0 +1,261 @@
//! The mesh wire contract — FROZEN surface.
//!
//! Every byte that crosses the mesh is one of the frames in this module,
//! postcard-encoded behind a one-byte protocol version. This file is the
//! contract between all mesh lanes: transport (endpoint/peer), membership
//! (gossip/registry), the session directory, and the media fan-out all build
//! against these types. **Changes here require a post in the mesh thread
//! before the edit** — two lanes compiling against different frame layouts is
//! the failure mode this file exists to prevent.
//!
//! ## The fencing law (non-negotiable)
//!
//! Every session-bearing frame carries the fenced tuple
//! [`FencedHeader`] `{session_id, generation, owner_runtime_id}`. Receivers
//! MUST reject frames whose generation is stale for that session, at every
//! hop. Mesh membership is a hint; the fenced generation (Redis CAS lease)
//! is the arbiter. The mesh may say "don't dial" — it may never say "take
//! over."
//!
//! ## Framing
//!
//! - **Datagrams** (realtime-media): one [`MeshDatagram`] per QUIC datagram,
//! postcard-encoded, no length prefix (the datagram boundary is the frame
//! boundary). Senders MUST check the encoded size against the connection's
//! `max_datagram_size()` and fail loud, never truncate.
//! - **Bi-streams** (reliable-stream + gossip control): length-delimited
//! postcard. Each frame is a u32-LE length followed by that many bytes of
//! postcard-encoded [`MeshStreamFrame`]. Max frame size: [`MAX_STREAM_FRAME`].
//! The first frame on any stream MUST be `Hello`; a non-`Hello` first frame
//! is a protocol error and the stream is reset.
use serde::{Deserialize, Serialize};
use uuid::Uuid;
/// ALPN for the mesh QUIC endpoint. Version bumps get a new ALPN so old and
/// new pods never half-speak to each other during a rolling deploy.
pub const ALPN: &[u8] = b"buzz/mesh/1";
/// Wire protocol version, first byte of every encoded frame (datagram or
/// stream frame). Receivers MUST reject unknown versions loudly (count it,
/// log it) rather than guessing.
pub const WIRE_VERSION: u8 = 1;
/// Hard cap on a single length-delimited stream frame (16 MiB). Anything
/// larger is a protocol error, not a bigger buffer.
pub const MAX_STREAM_FRAME: u32 = 16 * 1024 * 1024;
/// A relay runtime's mesh identity: the ed25519 public key of its signing
/// key, which is also its iroh endpoint id. Mesh authentication IS relay
/// identity — there is no separate credential.
#[derive(Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct RuntimeId(pub [u8; 32]);
impl RuntimeId {
pub fn to_hex(&self) -> String {
hex::encode(self.0)
}
}
impl std::fmt::Debug for RuntimeId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "RuntimeId({}…)", &self.to_hex()[..8])
}
}
impl std::fmt::Display for RuntimeId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.to_hex())
}
}
/// The fenced tuple. Present on every session-bearing frame; checked at
/// every hop against the Redis lease.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct FencedHeader {
pub session_id: Uuid,
/// Monotonic lease generation from the Redis CAS. A receiver that has
/// observed generation G for a session rejects any frame with < G.
pub generation: u64,
/// The runtime the sender believes owns the session. Advisory for
/// routing/diagnostics; the generation is what fences.
pub owner_runtime_id: RuntimeId,
}
/// Tunnel profile, fixed at session establishment.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum Profile {
/// Ordered, reliable, backpressured (goose/berd). Rides `open_bi()`.
ReliableStream,
/// Lossy-by-design realtime media (huddle Opus). Rides QUIC datagrams.
RealtimeMedia,
/// Huddle roster/join/leave control. State-bearing — a dropped roster
/// delta is an unrecoverable peer-index desync, so this rides a reliable
/// stream like `ReliableStream`, never datagrams. Separate variant so
/// routing intent and `/_mesh` counters stay legible.
HuddleControl,
}
/// One QUIC datagram: realtime media only.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct MeshDatagram {
pub fenced: FencedHeader,
/// Sender-scoped monotonic sequence for loss/reorder observability.
/// Receivers tolerate gaps and reordering; they never wait.
pub seq: u64,
/// Opaque, NIP-44-encrypted between client endpoints. The relay never
/// holds plaintext.
pub payload: Vec<u8>,
}
/// One length-delimited frame on a mesh bi-stream.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum MeshStreamFrame {
/// MUST be the first frame on every stream, in both directions.
Hello(StreamHello),
/// Opaque tunnel bytes for a reliable-stream session.
Data {
fenced: FencedHeader,
payload: Vec<u8>,
},
/// Clean close: the sender will send no more `Data` for this session.
/// Distinct from a QUIC reset — receivers treat reset as abnormal.
Goodbye {
fenced: FencedHeader,
reason: GoodbyeReason,
},
/// Membership gossip on the control stream (one per peer connection).
/// Payload is the gossip lane's postcard-encoded digest/delta exchange —
/// opaque at this layer so gossip can evolve without a wire bump here.
Gossip { payload: Vec<u8> },
}
/// Stream role, declared in the Hello.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum StreamRole {
/// The per-connection control stream (gossip + liveness). Exactly one
/// per peer connection, opened by the dialer immediately after connect.
Control,
/// A reliable-stream tunnel session.
Session {
fenced: FencedHeader,
profile: Profile,
},
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct StreamHello {
pub sender: RuntimeId,
pub role: StreamRole,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum GoodbyeReason {
/// Client closed / session ended normally.
SessionEnded,
/// This runtime is draining (SIGTERM) — re-establish elsewhere.
Draining,
/// The sender observed a newer generation and is fencing itself out.
StaleGeneration,
}
/// Encode a frame: version byte + postcard.
pub fn encode<T: Serialize>(frame: &T) -> Result<Vec<u8>, crate::MeshError> {
let buf = vec![WIRE_VERSION];
postcard::to_extend(frame, buf).map_err(crate::MeshError::Encode)
}
/// Decode a frame: check version byte, then postcard.
pub fn decode<'a, T: Deserialize<'a>>(bytes: &'a [u8]) -> Result<T, crate::MeshError> {
match bytes.split_first() {
Some((&WIRE_VERSION, rest)) => postcard::from_bytes(rest).map_err(crate::MeshError::Decode),
Some((&v, _)) => Err(crate::MeshError::UnknownWireVersion(v)),
None => Err(crate::MeshError::EmptyFrame),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn fenced() -> FencedHeader {
FencedHeader {
session_id: Uuid::from_u128(0xDEAD_BEEF),
generation: 42,
owner_runtime_id: RuntimeId([7u8; 32]),
}
}
#[test]
fn datagram_roundtrip() {
let d = MeshDatagram {
fenced: fenced(),
seq: 9001,
payload: vec![1, 2, 3],
};
let bytes = encode(&d).unwrap();
assert_eq!(bytes[0], WIRE_VERSION);
let back: MeshDatagram = decode(&bytes).unwrap();
assert_eq!(back, d);
}
#[test]
fn stream_frame_roundtrip() {
for f in [
MeshStreamFrame::Hello(StreamHello {
sender: RuntimeId([1u8; 32]),
role: StreamRole::Session {
fenced: fenced(),
profile: Profile::ReliableStream,
},
}),
MeshStreamFrame::Data {
fenced: fenced(),
payload: b"opaque".to_vec(),
},
MeshStreamFrame::Goodbye {
fenced: fenced(),
reason: GoodbyeReason::Draining,
},
MeshStreamFrame::Gossip {
payload: vec![0xAA; 16],
},
] {
let back: MeshStreamFrame = decode(&encode(&f).unwrap()).unwrap();
assert_eq!(back, f);
}
}
#[test]
fn unknown_version_rejected() {
let d = MeshDatagram {
fenced: fenced(),
seq: 1,
payload: vec![],
};
let mut bytes = encode(&d).unwrap();
bytes[0] = 99;
assert!(matches!(
decode::<MeshDatagram>(&bytes),
Err(crate::MeshError::UnknownWireVersion(99))
));
}
/// Opus @ 20ms worst case (~160B) + header must clear the conservative
/// QUIC datagram floor (~1200B path MTU minus QUIC overhead). This pins
/// the header overhead so it can't silently grow past the budget.
#[test]
fn datagram_header_overhead_within_budget() {
let payload = vec![0u8; 160];
let d = MeshDatagram {
fenced: fenced(),
seq: u64::MAX,
payload: payload.clone(),
};
let overhead = encode(&d).unwrap().len() - payload.len();
assert!(
overhead <= 64,
"datagram header overhead {overhead}B exceeds 64B budget"
);
}
}