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.
This commit is contained in:
Johnathan Corgan
2026-08-23 11:44:45 +01:00
parent f48113a808
commit 759fee50f5
5 changed files with 210 additions and 13 deletions
+18
View File
@@ -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
+179 -12
View File
@@ -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> {
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<Vec<DiscoveredPeer>>,
peers: Mutex<Buffered>,
}
/// 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<DiscoveredPeer> {
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);
}
}
+3 -1
View File
@@ -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),
+9
View File
@@ -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,
}
+1
View File
@@ -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,
})