From 759fee50f569c3e6157aa10034cf4823b0a77c46 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sat, 22 Aug 2026 20:20:54 +0100 Subject: [PATCH] Bound the Ethernet discovery buffer and drop its quadratic dedup Beacons are unauthenticated broadcast frames, and the discovery buffer deduplicated them by scanning a Vec for the source MAC with no cap on how many it held. Anything on the segment could name a fresh MAC per frame and drive both quadratic CPU in the receive loop and unbounded memory. The once-per-tick drain does not bound it either: a transport that is still receiving but not operational is skipped before draining. Key the buffer on source MAC and cap it at 1024 distinct MACs between drains. Keep the drain order oldest sighting first, because the reconcile layer spends a finite connect budget in that order and it must not depend on hash iteration order. A MAC already buffered is refreshed whether or not the buffer is full, so a flood of new MACs cannot crowd out a neighbour already seen. Count refused beacons in the transport stats as beacons_dropped, and log on the first drop and then on each power of ten, so the flooder does not choose the log rate. --- CHANGELOG.md | 18 +++ src/transport/ethernet/discovery.rs | 191 ++++++++++++++++++++++++++-- src/transport/ethernet/mod.rs | 4 +- src/transport/ethernet/stats.rs | 9 ++ src/transport/mod.rs | 1 + 5 files changed, 210 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index fe06b2db..9ac36e73 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -795,6 +795,24 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 #### Admission / peer caps +- The Ethernet transport's discovery buffer is now bounded and no longer costs + a linear scan per beacon. Beacons are unauthenticated broadcast frames, and + the buffer deduplicated by scanning a `Vec` for the source MAC and had no + cap, so anything on the segment could name a fresh MAC per frame and drive + both quadratic CPU in the receive loop and unbounded memory. It is drained + once per tick only while the transport is operational, so a transport that + is receiving but not operational was never drained at all. The buffer is now + a map keyed on source MAC, capped at 1024 distinct MACs between drains, with + the drain order still oldest sighting first so which neighbour gets dialed + under a connect budget does not depend on hash iteration order. A MAC already + buffered is always refreshed, so a flood of new MACs cannot crowd out a + neighbour already seen. Refused beacons are counted in the transport's stats + as `beacons_dropped` and reported in the log on the first drop and then on + each power-of-ten thereafter, so the flooder does not set the log rate. + **What this does not close**: a flood can still crowd out a neighbour not yet + seen in that tick, and anything able to flood raw frames on the segment can + already jam the beacon at L2 more cheaply. + - An accepted inbound TCP connection no longer holds a slot indefinitely without sending anything. The cap was tested at accept and the pool insert and counter bump followed with no read in between, while the frame reader's diff --git a/src/transport/ethernet/discovery.rs b/src/transport/ethernet/discovery.rs index 26ff75e9..d5528da7 100644 --- a/src/transport/ethernet/discovery.rs +++ b/src/transport/ethernet/discovery.rs @@ -7,7 +7,10 @@ use crate::transport::{DiscoveredPeer, TransportAddr, TransportId}; use secp256k1::XOnlyPublicKey; +use std::collections::HashMap; +use std::collections::hash_map::Entry; use std::sync::Mutex; +use tracing::warn; /// Discovery protocol version. pub const DISCOVERY_VERSION: u8 = 0x01; @@ -46,10 +49,33 @@ pub fn parse_beacon(data: &[u8]) -> Option { XOnlyPublicKey::from_slice(&data[2..34]).ok() } +/// Maximum distinct source MACs held between drains. +/// +/// Beacons are unauthenticated broadcast frames, so anything on the segment +/// can name as many source MACs as it likes; without a bound the buffer grows +/// with the flood rate, and it is not drained at all while the transport is +/// not operational. This caps it at roughly a thousand small structs, tens of +/// kilobytes. Raising it costs that much more memory per transport; lowering +/// it risks truncating discovery on a very large segment. A thousand distinct +/// beaconing FIPS neighbors within one tick is far outside anything a real +/// deployment produces. +const MAX_BUFFERED_PEERS: usize = 1024; + /// Buffer for discovered peers, drained by `discover()`. pub struct DiscoveryBuffer { transport_id: TransportId, - peers: Mutex>, + peers: Mutex, +} + +/// Peers keyed by source MAC, plus the sighting order `take()` restores. +#[derive(Default)] +struct Buffered { + by_mac: HashMap<[u8; 6], (u64, DiscoveredPeer)>, + seq: u64, + /// Beacons refused for want of room, cumulative and never reset. + dropped: u64, + /// Cumulative drop count that earns the next log record. + warn_at: u64, } impl DiscoveryBuffer { @@ -57,25 +83,82 @@ impl DiscoveryBuffer { pub fn new(transport_id: TransportId) -> Self { Self { transport_id, - peers: Mutex::new(Vec::new()), + peers: Mutex::new(Buffered::default()), } } /// Add a discovered peer from a received beacon. - pub fn add_peer(&self, src_mac: [u8; 6], pubkey: XOnlyPublicKey) { - let addr = TransportAddr::from_bytes(&src_mac); - let peer = DiscoveredPeer::with_hint(self.transport_id, addr, pubkey); - let mut peers = self.peers.lock().unwrap_or_else(|e| e.into_inner()); - // Deduplicate by MAC address — keep the latest - peers.retain(|p| p.addr.as_bytes() != src_mac); - peers.push(peer); + /// + /// Returns false when the beacon was refused because the buffer is full. + /// A MAC already buffered is always refreshed, so a flood of new MACs + /// cannot stop a known neighbor from being seen again. + pub fn add_peer(&self, src_mac: [u8; 6], pubkey: XOnlyPublicKey) -> bool { + let mut buffered = self.peers.lock().unwrap_or_else(|e| e.into_inner()); + buffered.seq += 1; + let seq = buffered.seq; + let full = buffered.by_mac.len() >= MAX_BUFFERED_PEERS; + let stored = match buffered.by_mac.entry(src_mac) { + // Refreshing moves the MAC to the end, as retain-then-push did. + Entry::Occupied(mut slot) => { + slot.insert((seq, self.peer(src_mac, pubkey))); + true + } + Entry::Vacant(slot) if !full => { + slot.insert((seq, self.peer(src_mac, pubkey))); + true + } + Entry::Vacant(_) => false, + }; + if !stored { + buffered.dropped += 1; + } + stored } - /// Drain all discovered peers since the last call. + /// Drain all discovered peers since the last call, oldest sighting first. pub fn take(&self) -> Vec { - let mut peers = self.peers.lock().unwrap_or_else(|e| e.into_inner()); - std::mem::take(&mut *peers) + let mut buffered = self.peers.lock().unwrap_or_else(|e| e.into_inner()); + let mut ordered: Vec<(u64, DiscoveredPeer)> = + buffered.by_mac.drain().map(|(_, entry)| entry).collect(); + // The reconcile layer spends a finite connect budget in this order, so + // which neighbor gets dialed must not depend on hash iteration order. + ordered.sort_unstable_by_key(|(seq, _)| *seq); + // Rate-limited: the drop rate is whatever the flooder chooses, and one + // record per drain would hand it the log volume too. + if buffered.dropped >= buffered.warn_at.max(1) { + warn!( + transport_id = %self.transport_id, + dropped = buffered.dropped, + cap = MAX_BUFFERED_PEERS, + "discovery buffer full, beacons from unseen neighbors refused" + ); + buffered.warn_at = next_decade(buffered.dropped); + } + ordered.into_iter().map(|(_, peer)| peer).collect() } + + /// Beacons refused for want of room since this buffer was created. + pub fn dropped(&self) -> u64 { + self.peers.lock().unwrap_or_else(|e| e.into_inner()).dropped + } + + /// Build the buffered peer record for one beacon. + fn peer(&self, src_mac: [u8; 6], pubkey: XOnlyPublicKey) -> DiscoveredPeer { + let addr = TransportAddr::from_bytes(&src_mac); + DiscoveredPeer::with_hint(self.transport_id, addr, pubkey) + } +} + +/// Smallest power of ten strictly greater than `n`, saturating at `u64::MAX`. +fn next_decade(n: u64) -> u64 { + let mut threshold = 1u64; + while threshold <= n { + match threshold.checked_mul(10) { + Some(next) => threshold = next, + None => return u64::MAX, + } + } + threshold } // ============================================================================ @@ -163,4 +246,88 @@ mod tests { let peers = buffer.take(); assert_eq!(peers.len(), 1); } + + /// Distinct MAC number `n`, for filling the buffer. + fn nth_mac(n: usize) -> [u8; 6] { + let bytes = (n as u64).to_be_bytes(); + [0x02, bytes[3], bytes[4], bytes[5], bytes[6], bytes[7]] + } + + #[test] + fn discovery_buffer_stops_buffering_past_the_cap() { + // The defect: an unauthenticated flood of source MACs grew the buffer + // without bound. Fails against the uncapped Vec, which returns all of + // them. + let buffer = DiscoveryBuffer::new(TransportId::new(1)); + let pubkey = test_pubkey(); + for n in 0..MAX_BUFFERED_PEERS + 50 { + buffer.add_peer(nth_mac(n), pubkey); + } + + let peers = buffer.take(); + assert_eq!(peers.len(), MAX_BUFFERED_PEERS); + // Drop-new keeps the earliest sightings. + assert_eq!(peers[0].addr.as_bytes(), &nth_mac(0)); + } + + #[test] + fn discovery_buffer_counts_dropped_beacons() { + let buffer = DiscoveryBuffer::new(TransportId::new(1)); + let pubkey = test_pubkey(); + for n in 0..MAX_BUFFERED_PEERS { + assert!(buffer.add_peer(nth_mac(n), pubkey)); + } + for n in MAX_BUFFERED_PEERS..MAX_BUFFERED_PEERS + 7 { + assert!(!buffer.add_peer(nth_mac(n), pubkey)); + } + + assert_eq!(buffer.dropped(), 7); + } + + #[test] + fn discovery_buffer_repeat_beacon_from_a_full_buffer_still_refreshes() { + let buffer = DiscoveryBuffer::new(TransportId::new(1)); + let pubkey = test_pubkey(); + for n in 0..MAX_BUFFERED_PEERS + 50 { + buffer.add_peer(nth_mac(n), pubkey); + } + // A neighbor already buffered must not be refused by a full buffer. + assert!(buffer.add_peer(nth_mac(0), pubkey)); + + let peers = buffer.take(); + assert_eq!(peers.len(), MAX_BUFFERED_PEERS); + assert_eq!(peers[peers.len() - 1].addr.as_bytes(), &nth_mac(0)); + } + + #[test] + fn discovery_buffer_drain_preserves_last_seen_order() { + // A regression pin on the map rewrite rather than a test of the + // defect: retain-then-push already produced this order. + let buffer = DiscoveryBuffer::new(TransportId::new(1)); + let pubkey = test_pubkey(); + let a = [0xaa; 6]; + let b = [0xbb; 6]; + let c = [0xcc; 6]; + + buffer.add_peer(a, pubkey); + buffer.add_peer(b, pubkey); + buffer.add_peer(c, pubkey); + buffer.add_peer(a, pubkey); + + let macs: Vec<_> = buffer + .take() + .iter() + .map(|p| p.addr.as_bytes().to_vec()) + .collect(); + assert_eq!(macs, vec![b.to_vec(), c.to_vec(), a.to_vec()]); + } + + #[test] + fn next_decade_steps_by_powers_of_ten() { + assert_eq!(next_decade(0), 1); + assert_eq!(next_decade(1), 10); + assert_eq!(next_decade(9), 10); + assert_eq!(next_decade(10), 100); + assert_eq!(next_decade(u64::MAX), u64::MAX); + } } diff --git a/src/transport/ethernet/mod.rs b/src/transport/ethernet/mod.rs index 60b0dbf1..ece1a804 100644 --- a/src/transport/ethernet/mod.rs +++ b/src/transport/ethernet/mod.rs @@ -446,7 +446,9 @@ async fn ethernet_receive_loop( stats.record_beacon_recv(); if discovery_enabled && let Some(pubkey) = parse_beacon(&buf[..len]) { - discovery_buffer.add_peer(src_mac, pubkey); + if !discovery_buffer.add_peer(src_mac, pubkey) { + stats.record_beacon_dropped(); + } trace!( transport_id = %transport_id, remote_mac = %format_mac(&src_mac), diff --git a/src/transport/ethernet/stats.rs b/src/transport/ethernet/stats.rs index c20492fd..02cb2d14 100644 --- a/src/transport/ethernet/stats.rs +++ b/src/transport/ethernet/stats.rs @@ -15,6 +15,7 @@ pub struct EthernetStats { pub recv_errors: AtomicU64, pub beacons_sent: AtomicU64, pub beacons_recv: AtomicU64, + pub beacons_dropped: AtomicU64, pub frames_too_short: AtomicU64, pub frames_too_long: AtomicU64, } @@ -31,6 +32,7 @@ impl EthernetStats { recv_errors: AtomicU64::new(0), beacons_sent: AtomicU64::new(0), beacons_recv: AtomicU64::new(0), + beacons_dropped: AtomicU64::new(0), frames_too_short: AtomicU64::new(0), frames_too_long: AtomicU64::new(0), } @@ -68,6 +70,11 @@ impl EthernetStats { self.beacons_recv.fetch_add(1, Ordering::Relaxed); } + /// Record a received beacon the discovery buffer had no room for. + pub fn record_beacon_dropped(&self) { + self.beacons_dropped.fetch_add(1, Ordering::Relaxed); + } + /// Take a snapshot of all counters. pub fn snapshot(&self) -> EthernetStatsSnapshot { EthernetStatsSnapshot { @@ -79,6 +86,7 @@ impl EthernetStats { recv_errors: self.recv_errors.load(Ordering::Relaxed), beacons_sent: self.beacons_sent.load(Ordering::Relaxed), beacons_recv: self.beacons_recv.load(Ordering::Relaxed), + beacons_dropped: self.beacons_dropped.load(Ordering::Relaxed), frames_too_short: self.frames_too_short.load(Ordering::Relaxed), frames_too_long: self.frames_too_long.load(Ordering::Relaxed), } @@ -102,6 +110,7 @@ pub struct EthernetStatsSnapshot { pub recv_errors: u64, pub beacons_sent: u64, pub beacons_recv: u64, + pub beacons_dropped: u64, pub frames_too_short: u64, pub frames_too_long: u64, } diff --git a/src/transport/mod.rs b/src/transport/mod.rs index a29cd34e..5fc7345a 100644 --- a/src/transport/mod.rs +++ b/src/transport/mod.rs @@ -1268,6 +1268,7 @@ impl TransportHandle { "recv_errors": snap.recv_errors, "beacons_sent": snap.beacons_sent, "beacons_recv": snap.beacons_recv, + "beacons_dropped": snap.beacons_dropped, "frames_too_short": snap.frames_too_short, "frames_too_long": snap.frames_too_long, })