Fix startup-wiring review blockers; add inbound profile dispatcher

Wren's review of 8b077fdb raised two correctness blockers, both fixed:

1. BUZZ_MESH now defaults OFF. Mesh forms only on explicit
   BUZZ_MESH=on|true|1 — an image upgrade with untouched env binds no
   UDP port and writes no Redis key (strict rollout no-regression).
   Pinned by mesh_defaults_off_when_env_absent.

2. Ready-record acceptance is anchored to the deployment's relay
   identity: MeshMembership::with_expected_relay_pubkey (set from the
   relay signing key in boot_mesh) rejects seeds attested by any other
   key — possession of some relay key is not authorization. Unanchored
   membership is fail-closed (admits nothing). Rejections are counted
   as foreign_relay_rejections in /_mesh.

Plus the single-slot inbound dispatcher (thread-agreed contract):
MeshInboundDispatcher in mesh_boot.rs implements InboundHandler,
installed once by boot_mesh; consumers register per-profile
entrypoints via MeshHandle.dispatcher (register_huddle_control /
register_reliable_stream / register_datagrams, first-registration
wins). HuddleControl/ReliableStream streams fan out by hello profile;
RealtimeMedia as a stream is rejected (datagram-only); traffic before
registration is logged and dropped (bounded boot-window race, fencing
makes retry safe). MeshStream::new and BoxFuture are now public so
consumer crates can stub streams/transports in tests.

cargo test -p buzz-relay-mesh: 32 passed; -p buzz-relay: 494 passed,
2 ignored; clippy both packages clean; fmt clean.

Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d
2026-07-08 13:06:49 -04:00
co-authored by Tyler Longwell
parent 8b077fdb44
commit 785653dae1
6 changed files with 346 additions and 27 deletions
+4 -1
View File
@@ -135,7 +135,10 @@ pub struct PeerInfo {
pub load: f32,
}
type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
/// Boxed future used across the seam traits. Public because implementors of
/// [`StreamSendHalf`]/[`StreamRecvHalf`]/[`RelayPeerTransport`] outside this
/// crate must name it.
pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
/// Seam 1: membership. Answers "who can I route to?" — never "who owns what."
pub trait RelayMeshMembership: Send + Sync + 'static {
+71 -12
View File
@@ -32,6 +32,13 @@ pub struct MeshMembership {
peers: Arc<RwLock<HashMap<RuntimeId, PeerState>>>,
draining: Arc<AtomicBool>,
stale_generation_rejections: Arc<AtomicU64>,
foreign_relay_rejections: Arc<AtomicU64>,
/// The relay identity ready records must be attested by. All pods in one
/// deployment share the relay signing key, so a valid seed is one signed
/// by *our* key — "signed by some relay key" is possession, not
/// authorization. `None` (never set) rejects every ready record: the
/// unanchored state is fail-closed, not accept-any.
expected_relay_pubkey: Option<String>,
phi_suspect_threshold: f64,
}
@@ -43,10 +50,19 @@ impl MeshMembership {
peers: Arc::new(RwLock::new(HashMap::new())),
draining: Arc::new(AtomicBool::new(false)),
stale_generation_rejections: Arc::new(AtomicU64::new(0)),
foreign_relay_rejections: Arc::new(AtomicU64::new(0)),
expected_relay_pubkey: None,
phi_suspect_threshold: DEFAULT_PHI_SUSPECT_THRESHOLD,
}
}
/// Anchor ready-record acceptance to this relay identity (hex pubkey).
/// Without an anchor, [`Self::apply_ready_records`] admits nothing.
pub fn with_expected_relay_pubkey(mut self, pubkey_hex: String) -> Self {
self.expected_relay_pubkey = Some(pubkey_hex);
self
}
pub fn with_phi_suspect_threshold(mut self, threshold: f64) -> Self {
self.phi_suspect_threshold = threshold;
self
@@ -61,11 +77,30 @@ impl MeshMembership {
/// Apply Redis bootstrap records. Existing gossip records win when they are
/// newer; ready-registry records enter as version 1 hints.
///
/// A record is admitted only when its `relay_pubkey` matches the expected
/// relay identity AND its attestation signature verifies. Matching first
/// makes the authorization question explicit: a record signed by a key we
/// don't recognize is foreign no matter how valid its signature is.
pub fn apply_ready_records(&self, records: impl IntoIterator<Item = ReadyRecord>) {
for ready in records {
if ready.runtime_id == self.local_runtime_id {
continue;
}
match self.expected_relay_pubkey.as_deref() {
Some(expected) if ready.relay_pubkey == expected => {}
anchor => {
self.foreign_relay_rejections
.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
runtime_id = %ready.runtime_id,
record_relay_pubkey = %ready.relay_pubkey,
anchored = anchor.is_some(),
"mesh membership rejected ready seed not attested by expected relay identity"
);
continue;
}
}
if let Err(err) = ready.verify_attestation() {
tracing::warn!(
runtime_id = %ready.runtime_id,
@@ -264,6 +299,7 @@ impl MeshMembership {
peers.sort_by(|a, b| a.runtime_id.cmp(&b.runtime_id));
let counters = MeshCounters {
stale_generation_rejections: self.stale_generation_rejections.load(Ordering::Relaxed),
foreign_relay_rejections: self.foreign_relay_rejections.load(Ordering::Relaxed),
peers: peers.iter().map(|peer| peer.counters.clone()).collect(),
};
MeshStatus {
@@ -381,20 +417,19 @@ mod tests {
nostr::Keys::generate()
}
fn ready_record(byte: u8, endpoint_addr: &str) -> ReadyRecord {
ReadyRecord::new(
rid(byte),
&relay_keys(),
vec![endpoint_addr.into()],
1,
vec![],
)
fn ready_record_signed(byte: u8, endpoint_addr: &str, keys: &nostr::Keys) -> ReadyRecord {
ReadyRecord::new(rid(byte), keys, vec![endpoint_addr.into()], 1, vec![])
}
#[test]
fn ready_records_seed_peers_but_skip_self() {
let membership = MeshMembership::new(record(1, 1, 1));
membership.apply_ready_records([ready_record(1, "self"), ready_record(2, "peer")]);
let keys = relay_keys();
let membership = MeshMembership::new(record(1, 1, 1))
.with_expected_relay_pubkey(keys.public_key().to_hex());
membership.apply_ready_records([
ready_record_signed(1, "self", &keys),
ready_record_signed(2, "peer", &keys),
]);
let peers = membership.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].runtime_id, rid(2));
@@ -402,8 +437,10 @@ mod tests {
#[test]
fn ready_records_must_have_valid_attestation() {
let membership = MeshMembership::new(record(1, 1, 1));
let mut tampered = ready_record(2, "peer");
let keys = relay_keys();
let membership = MeshMembership::new(record(1, 1, 1))
.with_expected_relay_pubkey(keys.public_key().to_hex());
let mut tampered = ready_record_signed(2, "peer", &keys);
tampered.runtime_id = rid(3);
tampered.runtime_pubkey = rid(3).to_hex();
@@ -411,6 +448,28 @@ mod tests {
assert!(membership.peers().is_empty());
}
#[test]
fn ready_records_from_foreign_relay_identity_are_rejected() {
let ours = relay_keys();
let theirs = relay_keys();
let membership = MeshMembership::new(record(1, 1, 1))
.with_expected_relay_pubkey(ours.public_key().to_hex());
// Validly signed, but by a key that isn't our deployment's identity.
membership.apply_ready_records([ready_record_signed(2, "peer", &theirs)]);
assert!(membership.peers().is_empty());
assert_eq!(membership.status().counters.foreign_relay_rejections, 1);
}
#[test]
fn unanchored_membership_rejects_all_ready_records() {
let keys = relay_keys();
let membership = MeshMembership::new(record(1, 1, 1));
membership.apply_ready_records([ready_record_signed(2, "peer", &keys)]);
assert!(membership.peers().is_empty());
assert_eq!(membership.status().counters.foreign_relay_rejections, 1);
}
#[test]
fn stale_gossip_record_is_ignored() {
let membership = MeshMembership::new(record(1, 1, 1));
+4 -1
View File
@@ -193,7 +193,10 @@ impl StreamRecvHalf for IrohRecvHalf {
}
impl MeshStream {
pub(crate) fn new(send: Box<dyn StreamSendHalf>, recv: Box<dyn StreamRecvHalf>) -> Self {
/// Assemble a stream from framing halves. Public so consumer crates can
/// build in-memory streams over stub halves in tests; production streams
/// only come from the transport (`MeshPeer::open_bi` / accept loop).
pub fn new(send: Box<dyn StreamSendHalf>, recv: Box<dyn StreamRecvHalf>) -> Self {
Self { send, recv }
}
}
+3
View File
@@ -41,6 +41,9 @@ pub enum ConnectionState {
#[derive(Clone, Debug, Default, Serialize)]
pub struct MeshCounters {
pub stale_generation_rejections: u64,
/// Ready-registry seeds rejected because their `relay_pubkey` did not
/// match this deployment's relay identity (or no anchor was configured).
pub foreign_relay_rejections: u64,
pub peers: Vec<MeshPeerCounters>,
}
+11 -9
View File
@@ -94,10 +94,10 @@ pub struct Config {
pub huddle_audio_available: bool,
/// Inter-relay mesh configuration (`BUZZ_MESH`, `BUZZ_MESH_BIND_ADDR`).
/// `mesh.enabled=false` (`BUZZ_MESH=off`) is the incident kill switch: the
/// relay must behave exactly like a single-instance deployment. Even when
/// enabled, the mesh only forms once peer registry records exist, so
/// N=1 deployments carry no mesh at runtime.
/// Opt-in: mesh forms only when `BUZZ_MESH=on` is explicit. The default
/// (absent/off) is exact single-instance behavior — no bind, no Redis
/// registry write — so an image upgrade with untouched env is a strict
/// no-regression rollout.
pub mesh: buzz_relay_mesh::MeshConfig,
/// Optional hex-encoded pubkey of the relay owner.
@@ -243,12 +243,14 @@ impl Config {
.map(|v| !(v == "false" || v == "0"))
.unwrap_or(true);
// Mesh kill switch: default enabled — the mesh is inert until peer
// registry records exist, so N=1 deployments pay nothing. `off`
// hard-disables for incidents (single-instance behavior guaranteed).
// Mesh opt-in: default OFF. Strict rollout no-regression — an image
// upgrade with untouched env must not bind a new UDP port or write a
// new Redis key. Horizontally-scaled deployments explicitly set
// `BUZZ_MESH=on`; anything else (absent, `off`, other values) keeps
// exact single-instance behavior.
let mesh_enabled = std::env::var("BUZZ_MESH")
.map(|v| !(v.eq_ignore_ascii_case("off") || v == "false" || v == "0"))
.unwrap_or(true);
.map(|v| v.eq_ignore_ascii_case("on") || v == "true" || v == "1")
.unwrap_or(false);
let mesh_bind_addr = std::env::var("BUZZ_MESH_BIND_ADDR")
.map(|raw| {
raw.parse::<SocketAddr>().map_err(|e| {
+253 -4
View File
@@ -20,18 +20,114 @@
//! "behave exactly like a single-instance relay."
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::sync::{Arc, OnceLock};
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,
InboundHandler, MeshDatagram, MeshMembership, MeshRuntime, MeshStatus, MeshStream, Profile,
RelayMeshMembership, RelayPeerTransport, RuntimeId, StreamHello, StreamRole,
};
use crate::config::Config;
use crate::tunnel::directory::SessionDirectory;
/// Handler for one inbound session-stream profile. Called on the accept task;
/// implementations must hand off promptly (spawn) rather than block.
pub type SessionStreamHandler = Box<dyn Fn(RuntimeId, StreamHello, MeshStream) + Send + Sync>;
/// Handler for inbound realtime-media datagrams.
pub type DatagramHandler = Box<dyn Fn(RuntimeId, MeshDatagram) + Send + Sync>;
/// The single [`InboundHandler`] slot owner: fans inbound mesh traffic out to
/// per-profile consumers.
///
/// The transport has exactly one inbound slot (`set_inbound`), but two lanes
/// consume session streams (`HuddleControl`, `ReliableStream`) and one
/// consumes datagrams (`RealtimeMedia`). This dispatcher is installed once by
/// [`boot_mesh`]; consumers register their entrypoints afterwards via the
/// `register_*` methods on [`MeshHandle`]'s dispatcher. Traffic arriving
/// before a slot is registered is logged and dropped — a bounded boot-window
/// race; fencing makes the peer's retry safe.
#[derive(Clone, Default)]
pub struct MeshInboundDispatcher {
slots: Arc<DispatcherSlots>,
}
#[derive(Default)]
struct DispatcherSlots {
huddle_control: OnceLock<SessionStreamHandler>,
reliable_stream: OnceLock<SessionStreamHandler>,
datagrams: OnceLock<DatagramHandler>,
}
impl MeshInboundDispatcher {
/// Register the `HuddleControl` session-stream consumer (huddle join lane).
/// First registration wins; later calls are logged and ignored.
pub fn register_huddle_control(&self, handler: SessionStreamHandler) {
if self.slots.huddle_control.set(handler).is_err() {
tracing::warn!("mesh dispatcher: huddle_control handler already registered — ignored");
}
}
/// Register the `ReliableStream` session-stream consumer (goose/berd lane).
pub fn register_reliable_stream(&self, handler: SessionStreamHandler) {
if self.slots.reliable_stream.set(handler).is_err() {
tracing::warn!("mesh dispatcher: reliable_stream handler already registered — ignored");
}
}
/// Register the realtime-media datagram consumer (huddle audio fan-out).
pub fn register_datagrams(&self, handler: DatagramHandler) {
if self.slots.datagrams.set(handler).is_err() {
tracing::warn!("mesh dispatcher: datagram handler already registered — ignored");
}
}
}
impl InboundHandler for MeshInboundDispatcher {
fn on_datagram(&self, from: RuntimeId, dgram: MeshDatagram) {
match self.slots.datagrams.get() {
Some(handler) => handler(from, dgram),
None => tracing::warn!(
peer = %from,
"mesh dispatcher: datagram before handler registration — dropped"
),
}
}
fn on_session_stream(&self, from: RuntimeId, hello: StreamHello, stream: MeshStream) {
let StreamRole::Session { profile, .. } = &hello.role else {
// Control streams are consumed inside the runtime and never reach
// the inbound slot; anything else here is a peer bug.
tracing::warn!(peer = %from, "mesh dispatcher: non-session stream role — dropped");
return;
};
let slot = match profile {
Profile::HuddleControl => &self.slots.huddle_control,
Profile::ReliableStream => &self.slots.reliable_stream,
Profile::RealtimeMedia => {
// Datagram-only profile: a *stream* claiming it is a protocol
// violation, never a valid session.
tracing::warn!(
peer = %from,
"mesh dispatcher: RealtimeMedia arrived as a stream (datagram-only profile) — rejected"
);
return;
}
};
match slot.get() {
Some(handler) => handler(from, hello, stream),
None => tracing::warn!(
peer = %from,
?profile,
"mesh dispatcher: session stream before handler registration — dropped"
),
}
}
}
/// Everything a mesh consumer needs, as one bundle.
#[derive(Clone)]
pub struct MeshHandle {
@@ -43,6 +139,11 @@ pub struct MeshHandle {
pub membership: Arc<dyn RelayMeshMembership>,
/// This runtime's boot-unique mesh identity.
pub local_runtime_id: RuntimeId,
/// Per-profile inbound registration: consumers call
/// `dispatcher.register_*` to receive their profile's traffic. The
/// dispatcher itself is already installed as the transport's single
/// inbound slot by [`boot_mesh`].
pub dispatcher: MeshInboundDispatcher,
/// The running mesh (status snapshots, shutdown).
runtime: MeshRuntime,
}
@@ -106,7 +207,7 @@ pub async fn boot_mesh(
shutting_down: Arc<AtomicBool>,
) -> anyhow::Result<Option<MeshHandle>> {
if !config.mesh.enabled {
tracing::info!("mesh disabled (BUZZ_MESH=off) — single-instance behavior");
tracing::info!("mesh disabled (BUZZ_MESH is not 'on') — single-instance behavior");
return Ok(None);
}
@@ -129,7 +230,11 @@ pub async fn boot_mesh(
let mut local_record = GossipRecord::new(runtime_id, addrs.clone(), PROTO_VERSION);
local_record.capabilities = capabilities();
let membership = MeshMembership::new(local_record);
// Anchor ready-record acceptance to this deployment's relay identity: all
// pods share the relay signing key, so a seed attested by any other key is
// foreign and rejected (Wren's review — possession is not authorization).
let membership = MeshMembership::new(local_record)
.with_expected_relay_pubkey(relay_keypair.public_key().to_hex());
let registry = ReadyRegistry::new(redis_pool.clone(), config.mesh.registry_refresh);
let ready_record = ReadyRecord::new(
@@ -181,11 +286,18 @@ pub async fn boot_mesh(
let membership_arc: Arc<dyn RelayMeshMembership> = Arc::new(runtime.membership().clone());
let transport: Arc<dyn RelayPeerTransport> = Arc::new(runtime.clone());
// Install the profile dispatcher as the transport's single inbound slot.
// Consumers (huddle control, reliable-stream) register their entrypoints
// on the handle's dispatcher after AppState wiring.
let dispatcher = MeshInboundDispatcher::default();
transport.set_inbound(Box::new(dispatcher.clone()));
Ok(Some(MeshHandle {
directory: SessionDirectory::new(redis_pool),
transport,
membership: membership_arc,
local_runtime_id: runtime_id,
dispatcher,
runtime,
}))
}
@@ -211,4 +323,141 @@ mod tests {
.expect("off path is never an error");
assert!(handle.is_none());
}
/// Blocker fix (Wren review of 8b077fdb): absent `BUZZ_MESH`, the mesh is
/// OFF — an env-untouched image upgrade must not bind or write Redis.
#[test]
fn mesh_defaults_off_when_env_absent() {
// `Config::from_env` in the test env has no BUZZ_MESH set unless a
// caller exported it; assert the fail-safe reading.
if std::env::var("BUZZ_MESH").is_ok() {
return; // externally forced — skip rather than assert a lie
}
let config = crate::config::Config::from_env().expect("default config loads");
assert!(!config.mesh.enabled, "BUZZ_MESH absent must mean mesh off");
}
use std::sync::Mutex;
use buzz_relay_mesh::{
BoxFuture, FencedHeader, MeshError, MeshStreamFrame, StreamRecvHalf, StreamSendHalf,
};
struct StubSend;
impl StreamSendHalf for StubSend {
fn send_frame(&mut self, _frame: MeshStreamFrame) -> BoxFuture<'_, Result<(), MeshError>> {
Box::pin(async { Ok(()) })
}
fn finish(&mut self) -> Result<(), MeshError> {
Ok(())
}
}
struct StubRecv;
impl StreamRecvHalf for StubRecv {
fn recv_frame(&mut self) -> BoxFuture<'_, Result<Option<MeshStreamFrame>, MeshError>> {
Box::pin(async { Ok(None) })
}
}
fn stub_stream() -> MeshStream {
MeshStream::new(Box::new(StubSend), Box::new(StubRecv))
}
fn rid(byte: u8) -> RuntimeId {
RuntimeId([byte; 32])
}
fn session_hello(sender: RuntimeId, profile: Profile) -> StreamHello {
StreamHello {
sender,
role: buzz_relay_mesh::StreamRole::Session {
fenced: FencedHeader {
session_id: uuid::Uuid::nil(),
generation: 1,
owner_runtime_id: sender,
},
profile,
},
}
}
#[test]
fn dispatcher_routes_session_streams_by_profile() {
let dispatcher = MeshInboundDispatcher::default();
let huddle_hits: Arc<Mutex<Vec<RuntimeId>>> = Arc::new(Mutex::new(vec![]));
let reliable_hits: Arc<Mutex<Vec<RuntimeId>>> = Arc::new(Mutex::new(vec![]));
let h = Arc::clone(&huddle_hits);
dispatcher.register_huddle_control(Box::new(move |from, _hello, _stream| {
h.lock().unwrap().push(from);
}));
let r = Arc::clone(&reliable_hits);
dispatcher.register_reliable_stream(Box::new(move |from, _hello, _stream| {
r.lock().unwrap().push(from);
}));
dispatcher.on_session_stream(
rid(1),
session_hello(rid(1), Profile::HuddleControl),
stub_stream(),
);
dispatcher.on_session_stream(
rid(2),
session_hello(rid(2), Profile::ReliableStream),
stub_stream(),
);
// Datagram-only profile arriving as a stream: rejected, routed nowhere.
dispatcher.on_session_stream(
rid(3),
session_hello(rid(3), Profile::RealtimeMedia),
stub_stream(),
);
assert_eq!(*huddle_hits.lock().unwrap(), vec![rid(1)]);
assert_eq!(*reliable_hits.lock().unwrap(), vec![rid(2)]);
}
#[test]
fn dispatcher_drops_traffic_before_registration_and_keeps_first_handler() {
let dispatcher = MeshInboundDispatcher::default();
// Pre-registration traffic must not panic — logged and dropped.
dispatcher.on_session_stream(
rid(1),
session_hello(rid(1), Profile::HuddleControl),
stub_stream(),
);
dispatcher.on_datagram(
rid(1),
MeshDatagram {
fenced: FencedHeader {
session_id: uuid::Uuid::nil(),
generation: 1,
owner_runtime_id: rid(1),
},
seq: 0,
payload: vec![],
},
);
let first: Arc<Mutex<u32>> = Arc::new(Mutex::new(0));
let f = Arc::clone(&first);
dispatcher.register_datagrams(Box::new(move |_, _| *f.lock().unwrap() += 1));
// Second registration is ignored; the first handler keeps the slot.
dispatcher.register_datagrams(Box::new(|_, _| panic!("second handler must not win")));
dispatcher.on_datagram(
rid(2),
MeshDatagram {
fenced: FencedHeader {
session_id: uuid::Uuid::nil(),
generation: 1,
owner_runtime_id: rid(2),
},
seq: 1,
payload: vec![],
},
);
assert_eq!(*first.lock().unwrap(), 1);
}
}