From b9feb446c042c2c34d687b0f006ba9cb4a34bae7 Mon Sep 17 00:00:00 2001 From: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co> Date: Wed, 8 Jul 2026 11:03:07 -0400 Subject: [PATCH] Add iroh transport core for relay mesh Co-authored-by: Tyler Longwell Signed-off-by: Tyler Longwell --- Cargo.toml | 2 +- crates/buzz-relay-mesh/src/endpoint.rs | 279 +++++++++++++++++++++++++ crates/buzz-relay-mesh/src/lib.rs | 2 + crates/buzz-relay-mesh/src/peer.rs | 199 ++++++++++++++++++ 4 files changed, 481 insertions(+), 1 deletion(-) create mode 100644 crates/buzz-relay-mesh/src/endpoint.rs create mode 100644 crates/buzz-relay-mesh/src/peer.rs diff --git a/Cargo.toml b/Cargo.toml index 908521c17..be5b29c7c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -64,7 +64,7 @@ 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 } +iroh = { version = "1.0.0-rc.0", default-features = false, features = ["tls-ring"] } serde_json = "1" serde_yaml = "0.9" evalexpr = "11" diff --git a/crates/buzz-relay-mesh/src/endpoint.rs b/crates/buzz-relay-mesh/src/endpoint.rs new file mode 100644 index 000000000..97fd84644 --- /dev/null +++ b/crates/buzz-relay-mesh/src/endpoint.rs @@ -0,0 +1,279 @@ +use std::net::SocketAddr; + +use iroh::{Endpoint, EndpointAddr, PublicKey, RelayMode, SecretKey, TransportAddr}; + +use crate::{MeshError, RuntimeId, ALPN}; + +/// Local iroh endpoint for the relay mesh. +/// +/// Identity is the iroh/ed25519 public key of a boot-unique keypair generated +/// at process start. +#[derive(Debug, Clone)] +pub struct MeshEndpoint { + endpoint: Endpoint, + runtime_id: RuntimeId, +} + +impl MeshEndpoint { + /// Generate a boot-unique mesh keypair and bind a mesh endpoint on `bind_addr`. + pub async fn bind(bind_addr: SocketAddr) -> Result { + Self::bind_with_secret_key(SecretKey::generate(), bind_addr).await + } + + /// Bind with an explicit keypair. Production should use [`Self::bind`] so + /// every process boot gets a fresh RuntimeId; tests use this for stable + /// identities. + pub async fn bind_with_secret_key( + secret_key: SecretKey, + bind_addr: SocketAddr, + ) -> Result { + let runtime_id = runtime_id_from_public_key(secret_key.public()); + let endpoint = Endpoint::builder(iroh::endpoint::presets::Minimal) + .secret_key(secret_key) + .alpns(vec![ALPN.to_vec()]) + .relay_mode(RelayMode::Disabled) + .bind_addr(bind_addr) + .map_err(|err| MeshError::Transport(err.to_string()))? + .bind() + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + + Ok(Self { + endpoint, + runtime_id, + }) + } + + pub fn runtime_id(&self) -> RuntimeId { + self.runtime_id + } + + pub fn endpoint(&self) -> Endpoint { + self.endpoint.clone() + } + + pub fn addr(&self) -> EndpointAddr { + self.endpoint.addr() + } + + pub async fn accept(&self) -> Result, MeshError> { + let Some(incoming) = self.endpoint.accept().await else { + return Ok(None); + }; + let conn = incoming + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + crate::peer::MeshPeer::from_connection(self.endpoint.clone(), conn).map(Some) + } + + pub async fn connect( + &self, + peer_addr: EndpointAddr, + ) -> Result { + let conn = self + .endpoint + .connect(peer_addr, ALPN) + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + crate::peer::MeshPeer::from_connection(self.endpoint.clone(), conn) + } +} + +pub fn runtime_id_from_public_key(public_key: PublicKey) -> RuntimeId { + RuntimeId(*public_key.as_bytes()) +} + +pub fn public_key_from_runtime_id(runtime_id: RuntimeId) -> Result { + PublicKey::from_bytes(&runtime_id.0).map_err(|err| MeshError::Transport(err.to_string())) +} + +pub fn direct_addr(runtime_id: RuntimeId, addr: SocketAddr) -> Result { + Ok(EndpointAddr::from_parts( + public_key_from_runtime_id(runtime_id)?, + [TransportAddr::Ip(addr)], + )) +} + +#[cfg(test)] +mod tests { + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + use std::time::Duration; + + use iroh::SecretKey; + use tokio::time::timeout; + use uuid::Uuid; + + use super::MeshEndpoint; + + use crate::{ + wire, FencedHeader, GoodbyeReason, MeshDatagram, MeshError, MeshStreamFrame, Profile, + RuntimeId, StreamHello, StreamRole, + }; + + fn loopback_any() -> SocketAddr { + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) + } + + fn fenced(owner_runtime_id: RuntimeId) -> FencedHeader { + FencedHeader { + session_id: Uuid::from_u128(0xABCD), + generation: 7, + owner_runtime_id, + } + } + + async fn endpoint_pair() -> (MeshEndpoint, MeshEndpoint) { + let a = + MeshEndpoint::bind_with_secret_key(SecretKey::from_bytes(&[1u8; 32]), loopback_any()) + .await + .unwrap(); + let b = + MeshEndpoint::bind_with_secret_key(SecretKey::from_bytes(&[2u8; 32]), loopback_any()) + .await + .unwrap(); + (a, b) + } + + async fn connected_pair() -> ( + crate::peer::MeshPeer, + crate::peer::MeshPeer, + RuntimeId, + RuntimeId, + ) { + let (a, b) = endpoint_pair().await; + let a_runtime_id = a.runtime_id(); + let b_runtime_id = b.runtime_id(); + + let b_addr = b.addr(); + let accept = tokio::spawn(async move { b.accept().await.unwrap().unwrap() }); + let a_peer = a.connect(b_addr).await.unwrap(); + let b_peer = accept.await.unwrap(); + + (a_peer, b_peer, a_runtime_id, b_runtime_id) + } + + #[tokio::test] + async fn two_endpoints_connect_with_alpn_and_authenticated_identity() { + let (a_peer, b_peer, a_runtime_id, b_runtime_id) = connected_pair().await; + + assert_eq!(a_peer.runtime_id(), b_runtime_id); + assert_eq!(b_peer.runtime_id(), a_runtime_id); + assert!(a_peer.max_datagram_size().expect("datagrams enabled") > 0); + } + + #[tokio::test] + async fn reliable_stream_roundtrip_carries_mesh_stream_frame() { + let (a_peer, b_peer, _a_runtime_id, b_runtime_id) = connected_pair().await; + let fenced = fenced(b_runtime_id); + let hello = MeshStreamFrame::Hello(StreamHello { + sender: RuntimeId([9u8; 32]), + role: StreamRole::Session { + fenced, + profile: Profile::ReliableStream, + }, + }); + let data = MeshStreamFrame::Data { + fenced, + payload: b"goose bytes".to_vec(), + }; + let goodbye = MeshStreamFrame::Goodbye { + fenced, + reason: GoodbyeReason::SessionEnded, + }; + + let recv = tokio::spawn(async move { + let mut stream = b_peer.accept_bi().await.unwrap(); + let first = stream.recv_frame().await.unwrap().unwrap(); + let second = stream.recv_frame().await.unwrap().unwrap(); + let third = stream.recv_frame().await.unwrap().unwrap(); + (first, second, third) + }); + + let mut stream = a_peer.open_bi().await.unwrap(); + stream.send_frame(hello.clone()).await.unwrap(); + stream.send_frame(data.clone()).await.unwrap(); + stream.send_frame(goodbye.clone()).await.unwrap(); + stream.finish().unwrap(); + + let (got_hello, got_data, got_goodbye) = timeout(Duration::from_secs(5), recv) + .await + .unwrap() + .unwrap(); + assert_eq!(got_hello, hello); + assert_eq!(got_data, data); + assert_eq!(got_goodbye, goodbye); + } + + #[tokio::test] + async fn datagram_roundtrip_carries_mesh_datagram() { + let (a_peer, b_peer, _a_runtime_id, b_runtime_id) = connected_pair().await; + let dgram = MeshDatagram { + fenced: fenced(b_runtime_id), + seq: 1, + payload: vec![13, 37, 42], + }; + + a_peer.send_datagram(&dgram).unwrap(); + let got = timeout(Duration::from_secs(5), b_peer.recv_datagram()) + .await + .unwrap() + .unwrap(); + assert_eq!(got, dgram); + } + + #[tokio::test] + async fn oversized_datagram_is_rejected_before_send() { + let (a_peer, _b_peer, _a_runtime_id, b_runtime_id) = connected_pair().await; + let max = a_peer.max_datagram_size().expect("datagrams enabled"); + let dgram = MeshDatagram { + fenced: fenced(b_runtime_id), + seq: 1, + payload: vec![0u8; max + 1], + }; + + let err = a_peer.send_datagram(&dgram).unwrap_err(); + assert!(matches!( + err, + MeshError::DatagramTooLarge { size, max: limit } if size > limit + )); + } + + #[tokio::test] + async fn opus_sized_datagrams_clear_empirical_local_loss_gate() { + let (a_peer, b_peer, _a_runtime_id, b_runtime_id) = connected_pair().await; + let payload_len = 1 /* Dawn huddle peer_index */ + 8 /* v2 audio header */ + 160; + let encoded_len = wire::encode(&MeshDatagram { + fenced: fenced(b_runtime_id), + seq: 0, + payload: vec![0u8; payload_len], + }) + .unwrap() + .len(); + assert!(encoded_len <= a_peer.max_datagram_size().expect("datagrams enabled")); + + let count = 64u64; + for seq in 0..count { + a_peer + .send_datagram(&MeshDatagram { + fenced: fenced(b_runtime_id), + seq, + payload: vec![seq as u8; payload_len], + }) + .unwrap(); + tokio::task::yield_now().await; + } + + let mut got = Vec::new(); + for _ in 0..count { + got.push( + timeout(Duration::from_secs(5), b_peer.recv_datagram()) + .await + .unwrap() + .unwrap() + .seq, + ); + } + got.sort_unstable(); + assert_eq!(got, (0..count).collect::>()); + } +} diff --git a/crates/buzz-relay-mesh/src/lib.rs b/crates/buzz-relay-mesh/src/lib.rs index f4c5e10e1..6a940c556 100644 --- a/crates/buzz-relay-mesh/src/lib.rs +++ b/crates/buzz-relay-mesh/src/lib.rs @@ -18,6 +18,8 @@ //! **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 endpoint; +pub mod peer; pub mod wire; // Lane modules — one owner per file (see the mesh thread for lane map): diff --git a/crates/buzz-relay-mesh/src/peer.rs b/crates/buzz-relay-mesh/src/peer.rs new file mode 100644 index 000000000..2a7326b6e --- /dev/null +++ b/crates/buzz-relay-mesh/src/peer.rs @@ -0,0 +1,199 @@ +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; + +use crate::{ + encode_datagram_checked, wire, MeshDatagram, MeshError, MeshStream, MeshStreamFrame, RuntimeId, + StreamRecvHalf, StreamSendHalf, ALPN, +}; + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct PeerCounters { + pub streams_opened: u64, + pub streams_accepted: u64, + pub datagrams_sent: u64, + pub datagrams_received: u64, +} + +#[derive(Debug, Default)] +struct PeerCountersInner { + streams_opened: AtomicU64, + streams_accepted: AtomicU64, + datagrams_sent: AtomicU64, + datagrams_received: AtomicU64, +} + +impl PeerCountersInner { + fn snapshot(&self) -> PeerCounters { + PeerCounters { + streams_opened: self.streams_opened.load(Ordering::Relaxed), + streams_accepted: self.streams_accepted.load(Ordering::Relaxed), + datagrams_sent: self.datagrams_sent.load(Ordering::Relaxed), + datagrams_received: self.datagrams_received.load(Ordering::Relaxed), + } + } +} + +/// Authenticated iroh connection to one peer runtime. +#[derive(Debug, Clone)] +pub struct MeshPeer { + _endpoint: iroh::Endpoint, + conn: iroh::endpoint::Connection, + runtime_id: RuntimeId, + counters: Arc, +} + +impl MeshPeer { + pub(crate) fn from_connection( + endpoint: iroh::Endpoint, + conn: iroh::endpoint::Connection, + ) -> Result { + if conn.alpn() != ALPN { + return Err(MeshError::Transport(format!( + "unexpected mesh ALPN {}", + String::from_utf8_lossy(conn.alpn()) + ))); + } + + Ok(Self { + _endpoint: endpoint, + runtime_id: crate::endpoint::runtime_id_from_public_key(conn.remote_id()), + conn, + counters: Arc::default(), + }) + } + + pub fn runtime_id(&self) -> RuntimeId { + self.runtime_id + } + + pub fn max_datagram_size(&self) -> Option { + self.conn.max_datagram_size() + } + + pub fn counters(&self) -> PeerCounters { + self.counters.snapshot() + } + + pub async fn open_bi(&self) -> Result { + let (send, recv) = self + .conn + .open_bi() + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + self.counters.streams_opened.fetch_add(1, Ordering::Relaxed); + Ok(MeshStream::new( + Box::new(IrohSendHalf(send)), + Box::new(IrohRecvHalf(recv)), + )) + } + + pub async fn accept_bi(&self) -> Result { + let (send, recv) = self + .conn + .accept_bi() + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + self.counters + .streams_accepted + .fetch_add(1, Ordering::Relaxed); + Ok(MeshStream::new( + Box::new(IrohSendHalf(send)), + Box::new(IrohRecvHalf(recv)), + )) + } + + pub fn send_datagram(&self, dgram: &MeshDatagram) -> Result<(), MeshError> { + let max = self + .conn + .max_datagram_size() + .ok_or_else(|| MeshError::Transport("peer does not support QUIC datagrams".into()))?; + let bytes = encode_datagram_checked(dgram, max)?; + self.conn + .send_datagram(bytes) + .map_err(|err| MeshError::Transport(err.to_string()))?; + self.counters.datagrams_sent.fetch_add(1, Ordering::Relaxed); + Ok(()) + } + + pub async fn recv_datagram(&self) -> Result { + let bytes = self + .conn + .read_datagram() + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + let dgram = wire::decode::(&bytes)?; + self.counters + .datagrams_received + .fetch_add(1, Ordering::Relaxed); + Ok(dgram) + } +} + +struct IrohSendHalf(iroh::endpoint::SendStream); +struct IrohRecvHalf(iroh::endpoint::RecvStream); + +impl StreamSendHalf for IrohSendHalf { + fn send_frame( + &mut self, + frame: MeshStreamFrame, + ) -> crate::BoxFuture<'_, Result<(), MeshError>> { + Box::pin(async move { + let bytes = wire::encode(&frame)?; + if bytes.len() > wire::MAX_STREAM_FRAME as usize { + return Err(MeshError::FrameTooLarge { + size: bytes.len(), + max: wire::MAX_STREAM_FRAME as usize, + }); + } + self.0 + .write_all(&(bytes.len() as u32).to_le_bytes()) + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + self.0 + .write_all(&bytes) + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + Ok(()) + }) + } + + fn finish(&mut self) -> Result<(), MeshError> { + self.0 + .finish() + .map_err(|err| MeshError::Transport(err.to_string())) + } +} + +impl StreamRecvHalf for IrohRecvHalf { + fn recv_frame(&mut self) -> crate::BoxFuture<'_, Result, MeshError>> { + Box::pin(async move { + let mut len = [0u8; 4]; + match self.0.read_exact(&mut len).await { + Ok(_) => {} + Err(iroh::endpoint::ReadExactError::FinishedEarly(0)) => return Ok(None), + Err(err) => return Err(MeshError::Transport(err.to_string())), + } + + let len = u32::from_le_bytes(len); + if len > wire::MAX_STREAM_FRAME { + return Err(MeshError::FrameTooLarge { + size: len as usize, + max: wire::MAX_STREAM_FRAME as usize, + }); + } + + let mut bytes = vec![0u8; len as usize]; + self.0 + .read_exact(&mut bytes) + .await + .map_err(|err| MeshError::Transport(err.to_string()))?; + wire::decode::(&bytes).map(Some) + }) + } +} + +impl MeshStream { + pub(crate) fn new(send: Box, recv: Box) -> Self { + Self { send, recv } + } +}