mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
Add iroh transport core for relay mesh
Co-authored-by: Tyler Longwell <tlongwell@block.xyz> Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
9a24af8aa9
commit
b9feb446c0
+1
-1
@@ -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"
|
||||
|
||||
@@ -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, MeshError> {
|
||||
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<Self, MeshError> {
|
||||
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<Option<crate::peer::MeshPeer>, 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<crate::peer::MeshPeer, MeshError> {
|
||||
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, MeshError> {
|
||||
PublicKey::from_bytes(&runtime_id.0).map_err(|err| MeshError::Transport(err.to_string()))
|
||||
}
|
||||
|
||||
pub fn direct_addr(runtime_id: RuntimeId, addr: SocketAddr) -> Result<EndpointAddr, MeshError> {
|
||||
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::<Vec<_>>());
|
||||
}
|
||||
}
|
||||
@@ -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):
|
||||
|
||||
@@ -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<PeerCountersInner>,
|
||||
}
|
||||
|
||||
impl MeshPeer {
|
||||
pub(crate) fn from_connection(
|
||||
endpoint: iroh::Endpoint,
|
||||
conn: iroh::endpoint::Connection,
|
||||
) -> Result<Self, MeshError> {
|
||||
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<usize> {
|
||||
self.conn.max_datagram_size()
|
||||
}
|
||||
|
||||
pub fn counters(&self) -> PeerCounters {
|
||||
self.counters.snapshot()
|
||||
}
|
||||
|
||||
pub async fn open_bi(&self) -> Result<MeshStream, MeshError> {
|
||||
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<MeshStream, MeshError> {
|
||||
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<MeshDatagram, MeshError> {
|
||||
let bytes = self
|
||||
.conn
|
||||
.read_datagram()
|
||||
.await
|
||||
.map_err(|err| MeshError::Transport(err.to_string()))?;
|
||||
let dgram = wire::decode::<MeshDatagram>(&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<Option<MeshStreamFrame>, 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::<MeshStreamFrame>(&bytes).map(Some)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl MeshStream {
|
||||
pub(crate) fn new(send: Box<dyn StreamSendHalf>, recv: Box<dyn StreamRecvHalf>) -> Self {
|
||||
Self { send, recv }
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user