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, })