mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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:
co-authored by
Tyler Longwell
parent
59a7ce3617
commit
8b077fdb44
@@ -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);
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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.
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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());
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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() {
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user