From 5819fde83942b572ebceac976f3c492ecb581cd9 Mon Sep 17 00:00:00 2001 From: npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d <011987e296fd5006292d2f930b574be47c7801048d1983c46c425d3c95f0cffd@sprout-oss.stage.blox.sqprod.co> Date: Wed, 8 Jul 2026 10:04:10 -0400 Subject: [PATCH] =?UTF-8?q?feat(mesh):=20buzz-relay-mesh=20crate=20skeleto?= =?UTF-8?q?n=20=E2=80=94=20frozen=20wire=20contract=20+=20relay=20seams?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Signed-off-by: Tyler Longwell --- Cargo.lock | 55 ++++++ Cargo.toml | 5 + crates/buzz-relay-mesh/Cargo.toml | 29 ++++ crates/buzz-relay-mesh/src/lib.rs | 179 ++++++++++++++++++++ crates/buzz-relay-mesh/src/wire.rs | 261 +++++++++++++++++++++++++++++ 5 files changed, 529 insertions(+) create mode 100644 crates/buzz-relay-mesh/Cargo.toml create mode 100644 crates/buzz-relay-mesh/src/lib.rs create mode 100644 crates/buzz-relay-mesh/src/wire.rs diff --git a/Cargo.lock b/Cargo.lock index afb474a2d..ebe0a69e1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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", ] diff --git a/Cargo.toml b/Cargo.toml index 17267207a..b752a466e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" diff --git a/crates/buzz-relay-mesh/Cargo.toml b/crates/buzz-relay-mesh/Cargo.toml new file mode 100644 index 000000000..33f901994 --- /dev/null +++ b/crates/buzz-relay-mesh/Cargo.toml @@ -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 } diff --git a/crates/buzz-relay-mesh/src/lib.rs b/crates/buzz-relay-mesh/src/lib.rs new file mode 100644 index 000000000..cbe922c32 --- /dev/null +++ b/crates/buzz-relay-mesh/src/lib.rs @@ -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, + /// Advisory load factor gossiped by the peer (0.0..). + pub load: f32, +} + +type BoxFuture<'a, T> = Pin + 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; + /// 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>; + + /// Register the handler invoked for inbound datagrams / session streams. + /// Called once at relay startup. + fn set_inbound(&self, handler: Box); +} + +/// 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, + pub(crate) recv: Box, +} + +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, 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, 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 { + 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)) +} diff --git a/crates/buzz-relay-mesh/src/wire.rs b/crates/buzz-relay-mesh/src/wire.rs new file mode 100644 index 000000000..8aeb540fb --- /dev/null +++ b/crates/buzz-relay-mesh/src/wire.rs @@ -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, +} + +/// 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, + }, + /// 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 }, +} + +/// 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(frame: &T) -> Result, 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 { + 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::(&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" + ); + } +}