diff --git a/crates/buzz-relay-mesh/src/lib.rs b/crates/buzz-relay-mesh/src/lib.rs index 676d73247..be031cb37 100644 --- a/crates/buzz-relay-mesh/src/lib.rs +++ b/crates/buzz-relay-mesh/src/lib.rs @@ -135,7 +135,10 @@ pub struct PeerInfo { pub load: f32, } -type BoxFuture<'a, T> = Pin + 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 + Send + 'a>>; /// Seam 1: membership. Answers "who can I route to?" — never "who owns what." pub trait RelayMeshMembership: Send + Sync + 'static { diff --git a/crates/buzz-relay-mesh/src/membership.rs b/crates/buzz-relay-mesh/src/membership.rs index f1d330bcb..9efffb1d8 100644 --- a/crates/buzz-relay-mesh/src/membership.rs +++ b/crates/buzz-relay-mesh/src/membership.rs @@ -32,6 +32,13 @@ pub struct MeshMembership { peers: Arc>>, draining: Arc, stale_generation_rejections: Arc, + foreign_relay_rejections: Arc, + /// 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, 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) { 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)); diff --git a/crates/buzz-relay-mesh/src/peer.rs b/crates/buzz-relay-mesh/src/peer.rs index 2a7326b6e..59ee8308f 100644 --- a/crates/buzz-relay-mesh/src/peer.rs +++ b/crates/buzz-relay-mesh/src/peer.rs @@ -193,7 +193,10 @@ impl StreamRecvHalf for IrohRecvHalf { } impl MeshStream { - pub(crate) fn new(send: Box, recv: Box) -> 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, recv: Box) -> Self { Self { send, recv } } } diff --git a/crates/buzz-relay-mesh/src/status.rs b/crates/buzz-relay-mesh/src/status.rs index a05cb276f..895cb7bea 100644 --- a/crates/buzz-relay-mesh/src/status.rs +++ b/crates/buzz-relay-mesh/src/status.rs @@ -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, } diff --git a/crates/buzz-relay/src/config.rs b/crates/buzz-relay/src/config.rs index a18cf1ee8..05b59f9fa 100644 --- a/crates/buzz-relay/src/config.rs +++ b/crates/buzz-relay/src/config.rs @@ -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::().map_err(|e| { diff --git a/crates/buzz-relay/src/mesh_boot.rs b/crates/buzz-relay/src/mesh_boot.rs index bd7490f71..ead21392b 100644 --- a/crates/buzz-relay/src/mesh_boot.rs +++ b/crates/buzz-relay/src/mesh_boot.rs @@ -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; + +/// Handler for inbound realtime-media datagrams. +pub type DatagramHandler = Box; + +/// 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, +} + +#[derive(Default)] +struct DispatcherSlots { + huddle_control: OnceLock, + reliable_stream: OnceLock, + datagrams: OnceLock, +} + +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, /// 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, ) -> anyhow::Result> { 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 = Arc::new(runtime.membership().clone()); let transport: Arc = 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, 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>> = Arc::new(Mutex::new(vec![])); + let reliable_hits: Arc>> = 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> = 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); + } }