Bound induced routing errors by the peer that induced them

The 100 ms suppression gate on CoordsRequired, PathBroken and MtuExceeded
was keyed on the failed datagram's destination address. That field is
chosen by whoever sent the datagram, so a fresh random destination per
packet was always a first sighting and the gate admitted every one. Each
admission inserted a key and then walked the whole map, making per-packet
cost grow with the flood rate while the sender's cost stayed flat, and
the emitted error is addressed to the datagram's source address, which
nothing binds to the sender either.

Add a per-link-peer token bucket, keyed on the AEAD-authenticated peer
the frame arrived over, and consult it ahead of the destination gate.
That peer is the only value at the emission point a sender cannot mint,
so it is the only one that can bound the emission or the growth of the
address-keyed map behind it. Default 20 signals a second sustained with
a burst of 50, both named constants.

The token is peeked and spent only once the destination gate has also
admitted, so one unroutable destination behind a high-fanout peer cannot
burn that peer's whole budget on signals nothing sends. The destination
map gets a hard 4096-entry ceiling and its expiry sweep is amortized to
once per eviction interval instead of running on every admission; at
capacity it admits without recording rather than refusing, because
refusing would turn a full map into node-wide silence during partition
healing, which is exactly when the map is largest and the signals are
most needed.

Also stop attaching the reporter's cached coordinates to a transit-
emitted PathBroken. The error goes to an unverified source address, so
the field answered a coordinate-cache read to anyone naming any address.
It is optional on the wire and no receiver reads it, so an unmodified
peer parses the frame unchanged.

Three counters make each outcome visible on the fipstop Routing tab.
This commit is contained in:
Johnathan Corgan
2026-08-23 11:46:15 +01:00
parent d95fc708e0
commit 42622c8efe
11 changed files with 637 additions and 30 deletions
+39
View File
@@ -939,6 +939,45 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
#### Data-plane / routing signals
- A transit node's induced routing errors are now bounded by the authenticated
link peer that induced them. The 100 ms suppression gate on
`CoordsRequired`, `PathBroken` and `MtuExceeded` was keyed on the failed
datagram's destination address, which is an envelope field the sender picks,
so a fresh random destination on every packet was always a first sighting and
every packet was admitted. Each admission also inserted a key and then walked
the whole map, so per-packet cost grew with the flood rate while the sender's
cost stayed flat, and the error itself is addressed to the datagram's source
address, which nothing binds to the sender either. A new per-link-peer token
bucket, 20 signals a second sustained with a burst of 50, is now consulted
first, keyed on the AEAD-authenticated peer the frame arrived over: the one
value at that point a sender cannot mint. The per-destination interval is
kept unchanged behind it, because it still does the aggregate suppression a
genuine outage needs, and no gate was added on the address the error is
returned to, which would have handed a sender a way to silence honest errors
toward a victim it names by keeping that victim's key hot.
Two ordering choices in there rather than left to be discovered. The peer's
token is peeked and only spent once the destination gate has also admitted,
so a single unroutable destination behind a high-fanout peer cannot burn that
peer's whole budget on signals nothing sends and silence every other
destination behind it. And the destination map now carries a hard ceiling of
4096 entries with its expiry sweep amortized to once per eviction interval
rather than run on every admission; when it is full it admits without
recording rather than refusing, because refusing would turn a full map into
node-wide silence exactly during partition healing, when many destinations
are legitimately unroutable at once. Emission stays bounded by the peer
budget in that state. Three counters, rendered on the fipstop Routing tab,
make each of the three outcomes visible instead of silent.
- A transit-emitted `PathBroken` no longer carries the reporter's cached
coordinates for the unreachable destination. The signal is returned to the
datagram's source address, so anyone able to reach the node could name any
address and have the node's coordinate cache read back to them, one entry per
packet. The field is optional on the wire and no receiver reads it, so this
is an emission change only: an unmodified peer parses the frame exactly as
before. Which of the two signals is emitted still discloses whether the entry
exists.
- The influence a remote party has over path MTU is now bounded, and the
per-destination path MTU cache has a way back. The `path_mtu` field is an
unsigned per-hop transit annotation carried outside the signed proof, and the
+3
View File
@@ -298,6 +298,9 @@ fn draw_routing_stats(
("Path Broken Refused", err("unbound_broken")),
("MTU Exceeded Refused", err("unbound_mtu")),
("Forged Pairing", err("unbound_forged")),
("Emit Over Peer Budget", err("emit_over_peer_budget")),
("Emit Over Dest Interval", err("emit_over_dest_interval")),
("Emit Limiter At Capacity", err("emit_limiter_at_capacity")),
],
));
right.push(Line::from(""));
+1 -1
View File
@@ -1151,7 +1151,7 @@ fn routing_focused_pane_scrolls() {
// column is the taller of the two, so scrolling fully to the bottom would
// over-scroll the right column past Congestion; this offset lands the
// Congestion region inside the short window instead.
app1.scroll_offsets.insert((Tab::Routing, 2), 28);
app1.scroll_offsets.insert((Tab::Routing, 2), 31);
let buf1 = testkit::render(100, 20, |frame, area| {
super::routing::draw(frame, &app1, area);
});
+3
View File
@@ -33,6 +33,9 @@
},
"error_signals": {
"coords_required": 0,
"emit_limiter_at_capacity": 0,
"emit_over_dest_interval": 0,
"emit_over_peer_budget": 0,
"lookup_resp_mtu_below_floor": 0,
"mtu_exceeded": 0,
"mtu_exceeded_below_floor": 0,
+74 -23
View File
@@ -8,6 +8,7 @@
use crate::NodeAddr;
use crate::node::reject::ForwardingReject;
use crate::node::routing_error_rate_limit::LimitVerdict;
use crate::node::session_wire::{
FSP_COMMON_PREFIX_SIZE, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED, FSP_PHASE_MSG1, FSP_PHASE_MSG2,
FspCommonPrefix, FspEncryptedHeader, parse_encrypted_coords,
@@ -101,7 +102,7 @@ impl Node {
bytes = payload.len(),
"Dropping transit SessionDatagram: no route to destination"
);
self.send_routing_error(&datagram).await;
self.send_routing_error(from, &datagram).await;
return;
}
};
@@ -145,7 +146,7 @@ impl Node {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::MtuExceeded, payload.len());
self.send_mtu_exceeded_error(&datagram, mtu).await;
self.send_mtu_exceeded_error(from, &datagram, mtu).await;
}
_ => {
self.metrics()
@@ -283,6 +284,43 @@ impl Node {
}
}
/// Decide whether this node may emit a routing error induced by `from`
/// about `dest`.
///
/// Two gates, in this order. The per-link-peer budget bounds what one
/// admitted peer can induce, and is the only one keyed on something the
/// sender cannot mint; it is consulted first so the per-destination map
/// only grows at the budget rate. The per-destination interval is the
/// aggregate suppressor during a genuine outage.
///
/// The budget token is peeked and only committed once the destination gate
/// has also admitted. Charging it on a suppressed signal would let a single
/// unroutable destination behind a high-fanout peer spend that peer's whole
/// budget on emissions nothing sends, silencing every other destination
/// behind it.
fn admit_error_emission(&mut self, from: &NodeAddr, dest: &NodeAddr) -> bool {
let now = Instant::now();
if !self.peer_error_budget.has_token(from, now) {
self.metrics().errors.emit_over_peer_budget.inc();
return false;
}
match self.routing_error_rate_limiter.check(dest, now) {
LimitVerdict::Suppress => {
self.metrics().errors.emit_over_dest_interval.inc();
return false;
}
LimitVerdict::AdmitAtCapacity => {
self.metrics().errors.emit_limiter_at_capacity.inc();
}
LimitVerdict::Admit => {}
}
self.peer_error_budget.commit(from, now);
true
}
/// Generate and send a routing error signal back to the datagram's source.
///
/// If we have cached coords for the destination, send PathBroken (we know
@@ -291,12 +329,11 @@ impl Node {
///
/// If we can't route the error back to the source either, drop silently.
/// No cascading errors.
async fn send_routing_error(&mut self, original: &SessionDatagram) {
// Rate limit: one error signal per destination per 100ms
if !self
.routing_error_rate_limiter
.should_send(&original.dest_addr)
{
///
/// `from` is the authenticated link peer the original datagram arrived
/// from, and is what the emission is charged against.
async fn send_routing_error(&mut self, from: &NodeAddr, original: &SessionDatagram) {
if !self.admit_error_emission(from, &original.dest_addr) {
return;
}
@@ -307,15 +344,21 @@ impl Node {
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let error_payload =
if let Some(coords) = self.coord_cache().get(&original.dest_addr, now_ms) {
let coords = coords.clone();
PathBroken::new(original.dest_addr, my_addr)
.with_last_coords(coords)
.encode()
} else {
CoordsRequired::new(original.dest_addr, my_addr).encode()
};
// The choice between the two signals still leaks whether this node
// holds coords for the destination, but the coordinates themselves
// are not attached: the address the error is returned to is the
// datagram's own src_addr, which nothing binds to the peer that sent
// it, so attaching them would answer a cache read to whoever names an
// address. No receiver reads the field.
let error_payload = if self
.coord_cache()
.get(&original.dest_addr, now_ms)
.is_some()
{
PathBroken::new(original.dest_addr, my_addr).encode()
} else {
CoordsRequired::new(original.dest_addr, my_addr).encode()
};
let error_dg = SessionDatagram::new(my_addr, original.src_addr, error_payload)
.with_ttl(self.config().node.session.default_ttl);
@@ -356,12 +399,20 @@ impl Node {
/// Called when `send_encrypted_link_message()` fails with
/// `NodeError::MtuExceeded` during forwarding. The signal tells the
/// source the bottleneck MTU so it can immediately reduce its path MTU.
async fn send_mtu_exceeded_error(&mut self, original: &SessionDatagram, bottleneck_mtu: u16) {
// Rate limit: reuse routing_error_rate_limiter keyed on dest_addr
if !self
.routing_error_rate_limiter
.should_send(&original.dest_addr)
{
///
/// `from` is the authenticated link peer the original datagram arrived
/// from, and is what the emission is charged against. MtuExceeded shares
/// the link peer's budget with the routing errors rather than holding its
/// own: a separate bucket would insulate path-MTU discovery from
/// routing-error pressure, at the cost of a second knob and of letting one
/// peer induce twice the total emission.
async fn send_mtu_exceeded_error(
&mut self,
from: &NodeAddr,
original: &SessionDatagram,
bottleneck_mtu: u16,
) {
if !self.admit_error_emission(from, &original.dest_addr) {
return;
}
+16
View File
@@ -508,6 +508,19 @@ pub struct ErrorMetrics {
/// annotation.
pub lookup_resp_mtu_below_floor: Counter,
pub unbound: UnboundSignals,
/// Routing errors this node declined to emit because the authenticated
/// link peer that induced them had spent its budget. A rising count is
/// either a peer flooding unroutable traffic or a hub relaying more
/// simultaneously-broken destinations than the budget allows.
pub emit_over_peer_budget: Counter,
/// Routing errors this node declined to emit because one for the same
/// destination went out within the per-destination interval. This is the
/// aggregate suppression a real outage produces.
pub emit_over_dest_interval: Counter,
/// Routing errors emitted without recording their destination, because
/// the per-destination limiter's map was full. The signal was still sent;
/// what was lost is interval suppression for that destination.
pub emit_limiter_at_capacity: Counter,
}
impl ErrorMetrics {
@@ -524,6 +537,9 @@ impl ErrorMetrics {
unbound_broken: self.unbound.broken.get(),
unbound_mtu: self.unbound.mtu.get(),
unbound_forged: self.unbound.forged.get(),
emit_over_peer_budget: self.emit_over_peer_budget.get(),
emit_over_dest_interval: self.emit_over_dest_interval.get(),
emit_limiter_at_capacity: self.emit_limiter_at_capacity.get(),
}
}
}
+8
View File
@@ -17,6 +17,7 @@ pub(crate) mod encrypt_worker;
mod handlers;
mod lifecycle;
pub(crate) mod metrics;
mod peer_error_budget;
mod rate_limit;
pub(crate) mod reject;
mod reloadable;
@@ -32,6 +33,7 @@ mod tree;
pub(crate) mod wire;
use self::discovery_rate_limit::{DiscoveryBackoff, DiscoveryForwardRateLimiter};
use self::peer_error_budget::PeerErrorBudget;
use self::rate_limit::{HandshakeRateLimiter, SessionSetupRateLimiter};
use self::reloadable::Reloadable;
use self::routing_error_rate_limit::RoutingErrorRateLimiter;
@@ -504,6 +506,10 @@ pub struct Node {
icmp_rate_limiter: IcmpRateLimiter,
/// Rate limiter for routing error signals (CoordsRequired / PathBroken).
routing_error_rate_limiter: RoutingErrorRateLimiter,
/// Budget for routing errors this node may be induced to emit, charged to
/// the authenticated link peer whose datagram induced them. The only
/// bound on the emission that a sender cannot escape by varying a field.
peer_error_budget: PeerErrorBudget,
/// Rate limiter for source-side CoordsRequired/PathBroken responses.
coords_response_rate_limiter: RoutingErrorRateLimiter,
/// Backoff for failed discovery lookups (originator-side).
@@ -799,6 +805,7 @@ impl Node {
setup_rate_limiter,
icmp_rate_limiter: IcmpRateLimiter::new(),
routing_error_rate_limiter: RoutingErrorRateLimiter::new(),
peer_error_budget: PeerErrorBudget::new(),
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
std::time::Duration::from_millis(coords_response_interval_ms),
),
@@ -959,6 +966,7 @@ impl Node {
setup_rate_limiter,
icmp_rate_limiter: IcmpRateLimiter::new(),
routing_error_rate_limiter: RoutingErrorRateLimiter::new(),
peer_error_budget: PeerErrorBudget::new(),
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
std::time::Duration::from_millis(coords_response_interval_ms),
),
+239
View File
@@ -0,0 +1,239 @@
//! Per-link-peer budget for induced routing-error emissions.
//!
//! A transit node synthesizes a routing error (CoordsRequired, PathBroken or
//! MtuExceeded) in response to a datagram it could not forward. Every field of
//! that datagram is chosen by whoever sent it, so a per-destination or
//! per-source gate can be escaped by varying the field it is keyed on. The one
//! value at the emission point an attacker cannot mint is the authenticated
//! link peer the frame arrived from, whose cardinality is bounded by the peer
//! table and by admission. This budget is keyed on it, and is consulted ahead
//! of any address-keyed structure so that those structures only grow at the
//! budget rate.
use crate::NodeAddr;
use std::collections::HashMap;
use std::time::{Duration, Instant};
/// Sustained rate, in signals per second, at which one authenticated link peer
/// may induce this node to emit routing errors.
///
/// Bounds the reflection an admitted peer can aim at a victim it names, and
/// bounds how fast that peer can grow the per-destination limiter's map.
/// Raising it costs proportionally more reflected traffic per peer; lowering it
/// silences a hub peer that relays many sources through a genuine outage
/// sooner, which costs those sources their CoordsRequired and their
/// path-MTU feedback.
pub const PEER_ERROR_RATE_PER_SEC: u32 = 20;
/// Number of routing errors one link peer may induce back to back before the
/// sustained rate applies.
///
/// Sized so that an ordinary burst of unroutable traffic behind one peer still
/// signals promptly. Raising it lets a peer front-load a larger reflection;
/// lowering it makes a legitimate convergence burst arrive as a trickle.
pub const PEER_ERROR_BURST: u32 = 50;
/// Tokens are carried in thousandths so the refill of a sub-millisecond
/// interval is not rounded away.
const MILLI: u64 = 1000;
/// A peer whose bucket has refilled to full carries no state worth keeping, so
/// entries are dropped once per this interval to bound the map across peer
/// churn.
const SWEEP_INTERVAL: Duration = Duration::from_secs(30);
/// One peer's token bucket.
struct Bucket {
/// Remaining tokens, in thousandths of a signal.
milli_tokens: u64,
/// When `milli_tokens` was last brought up to date.
last_refill: Instant,
}
/// Token-bucket budget for routing errors, keyed on the authenticated link
/// peer that induced them.
pub struct PeerErrorBudget {
buckets: HashMap<NodeAddr, Bucket>,
/// Refill rate in thousandths of a token per millisecond.
milli_per_ms: u64,
/// Bucket ceiling, in thousandths of a token.
capacity: u64,
last_sweep: Instant,
}
impl PeerErrorBudget {
/// Create a budget at the shipped rate and burst.
pub fn new() -> Self {
Self::with_rate(PEER_ERROR_RATE_PER_SEC, PEER_ERROR_BURST)
}
/// Create a budget with an explicit sustained rate and burst.
pub fn with_rate(per_sec: u32, burst: u32) -> Self {
Self {
buckets: HashMap::new(),
milli_per_ms: u64::from(per_sec),
capacity: u64::from(burst) * MILLI,
last_sweep: Instant::now(),
}
}
/// Whether `peer` has a token to spend, without spending it.
///
/// Separate from [`Self::commit`] so a signal that a later gate suppresses
/// does not consume budget: an outage behind a high-fanout peer would
/// otherwise spend that peer's whole budget on emissions the
/// per-destination interval discards, silencing every other destination
/// behind it.
pub fn has_token(&mut self, peer: &NodeAddr, now: Instant) -> bool {
self.refill(peer, now);
self.buckets
.get(peer)
.is_some_and(|b| b.milli_tokens >= MILLI)
}
/// Spend one token for `peer`. Call only on the path that actually emits.
pub fn commit(&mut self, peer: &NodeAddr, now: Instant) {
self.refill(peer, now);
if let Some(bucket) = self.buckets.get_mut(peer) {
bucket.milli_tokens = bucket.milli_tokens.saturating_sub(MILLI);
}
self.sweep(now);
}
/// Bring `peer`'s bucket up to date, creating a full one on first sighting.
fn refill(&mut self, peer: &NodeAddr, now: Instant) {
let capacity = self.capacity;
let milli_per_ms = self.milli_per_ms;
let bucket = self.buckets.entry(*peer).or_insert(Bucket {
milli_tokens: capacity,
last_refill: now,
});
let elapsed_ms = now
.saturating_duration_since(bucket.last_refill)
.as_millis() as u64;
if elapsed_ms > 0 {
bucket.milli_tokens = (bucket.milli_tokens + elapsed_ms * milli_per_ms).min(capacity);
bucket.last_refill = now;
}
}
/// Drop full buckets, at most once per [`SWEEP_INTERVAL`].
fn sweep(&mut self, now: Instant) {
if now.saturating_duration_since(self.last_sweep) < SWEEP_INTERVAL {
return;
}
self.last_sweep = now;
let capacity = self.capacity;
let milli_per_ms = self.milli_per_ms;
self.buckets.retain(|_, b| {
let elapsed_ms = now.saturating_duration_since(b.last_refill).as_millis() as u64;
b.milli_tokens + elapsed_ms * milli_per_ms < capacity
});
}
#[cfg(test)]
pub fn len(&self) -> usize {
self.buckets.len()
}
}
impl Default for PeerErrorBudget {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn addr(val: u8) -> NodeAddr {
let mut bytes = [0u8; 16];
bytes[0] = val;
NodeAddr::from_bytes(bytes)
}
/// Spend one token per admitted emission.
fn spend(budget: &mut PeerErrorBudget, peer: &NodeAddr, now: Instant) -> bool {
if !budget.has_token(peer, now) {
return false;
}
budget.commit(peer, now);
true
}
#[test]
fn a_peer_may_emit_its_full_burst_then_is_suppressed() {
let mut budget = PeerErrorBudget::new();
let now = Instant::now();
let peer = addr(1);
for i in 0..PEER_ERROR_BURST {
assert!(spend(&mut budget, &peer, now), "burst signal {i} refused");
}
assert!(!spend(&mut budget, &peer, now));
}
#[test]
fn an_exhausted_budget_refills_at_the_sustained_rate() {
let mut budget = PeerErrorBudget::new();
let start = Instant::now();
let peer = addr(1);
for _ in 0..PEER_ERROR_BURST {
assert!(spend(&mut budget, &peer, start));
}
assert!(!spend(&mut budget, &peer, start));
// One second of refill buys exactly the sustained rate back.
let later = start + Duration::from_secs(1);
for i in 0..PEER_ERROR_RATE_PER_SEC {
assert!(
spend(&mut budget, &peer, later),
"refilled signal {i} refused"
);
}
assert!(!spend(&mut budget, &peer, later));
}
#[test]
fn one_peer_exhausting_its_budget_does_not_silence_another() {
let mut budget = PeerErrorBudget::new();
let now = Instant::now();
let noisy = addr(1);
let quiet = addr(2);
for _ in 0..PEER_ERROR_BURST {
assert!(spend(&mut budget, &noisy, now));
}
assert!(!spend(&mut budget, &noisy, now));
assert!(spend(&mut budget, &quiet, now));
}
#[test]
fn peeking_does_not_spend_a_token() {
let mut budget = PeerErrorBudget::with_rate(1, 1);
let now = Instant::now();
let peer = addr(1);
assert!(budget.has_token(&peer, now));
assert!(budget.has_token(&peer, now));
budget.commit(&peer, now);
assert!(!budget.has_token(&peer, now));
}
#[test]
fn full_buckets_are_dropped_by_the_sweep() {
let mut budget = PeerErrorBudget::new();
let start = Instant::now();
for i in 0..50u8 {
assert!(spend(&mut budget, &addr(i), start));
}
assert_eq!(budget.len(), 50);
// Long enough for every bucket to have refilled to full.
let later = start + SWEEP_INTERVAL + Duration::from_secs(1);
assert!(spend(&mut budget, &addr(200), later));
assert_eq!(budget.len(), 1);
}
}
+163 -6
View File
@@ -2,11 +2,56 @@
//!
//! Prevents routing error floods (CoordsRequired / PathBroken) by
//! rate-limiting error signals per destination address at transit nodes.
//!
//! The destination address is chosen by whoever sent the datagram, so this
//! gate is an aggregate suppressor during a real outage and never a bound on
//! what one sender can induce: a fresh destination is always a first sighting.
//! The bound is `PeerErrorBudget`, keyed on the authenticated link peer and
//! consulted first. What this module owes on top of its interval is that its
//! own map stays bounded and its per-admit cost stays sub-linear whatever the
//! sender does with the key.
use crate::NodeAddr;
use std::collections::HashMap;
use std::time::{Duration, Instant};
/// Maximum number of destinations this limiter remembers at once.
///
/// A hard ceiling on the map an attacker can grow by varying the destination
/// address. Raising it costs one `NodeAddr` plus one `Instant` per entry and
/// buys interval suppression across more simultaneously-unroutable
/// destinations; lowering it makes admission-without-recording (see
/// [`LimitVerdict::AdmitAtCapacity`]) the common case sooner, which weakens
/// the interval gate but never the peer budget.
const MAX_ENTRIES: usize = 4096;
/// Fraction of `max_age` between amortized sweeps.
///
/// The sweep is a full-map `retain`, so running it on every admit made
/// per-packet cost linear in a map the sender sizes. Eight sweeps per entry
/// lifetime keeps expired entries from accumulating without putting the scan
/// on the per-packet path.
const SWEEPS_PER_MAX_AGE: u32 = 8;
/// What the limiter decided about one candidate error signal.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LimitVerdict {
/// Send it; the destination was recorded.
Admit,
/// Send it, but the map was full so the destination was not recorded and
/// the interval will not suppress its successor.
///
/// This gate fails open deliberately. Failing closed would turn a full map
/// into node-wide silence on error signalling, and the map is fullest
/// exactly during partition healing, when many destinations are
/// legitimately unroutable at once and sources most need the signal. The
/// bound on emission is the per-peer budget, not this map.
AdmitAtCapacity,
/// Suppress it; an error for this destination went out within the
/// interval.
Suppress,
}
/// Rate limiter for routing error signals (CoordsRequired / PathBroken).
///
/// Tracks the last time a routing error was sent for each destination
@@ -18,6 +63,12 @@ pub struct RoutingErrorRateLimiter {
min_interval: Duration,
/// Maximum age of entries before cleanup.
max_age: Duration,
/// When `cleanup` last ran.
last_sweep: Instant,
/// Sweeps run since construction. Read by the tests that hold the
/// amortization property: a full-map scan per admit is the denial-of-
/// service multiplier this counter exists to catch coming back.
sweeps: u64,
}
impl RoutingErrorRateLimiter {
@@ -29,6 +80,8 @@ impl RoutingErrorRateLimiter {
last_sent: HashMap::new(),
min_interval: Duration::from_millis(100),
max_age: Duration::from_secs(10),
last_sweep: Instant::now(),
sweeps: 0,
}
}
@@ -38,6 +91,8 @@ impl RoutingErrorRateLimiter {
last_sent: HashMap::new(),
min_interval,
max_age: Duration::from_secs(10),
last_sweep: Instant::now(),
sweeps: 0,
}
}
@@ -47,29 +102,57 @@ impl RoutingErrorRateLimiter {
/// this destination, or if this is the first error. Updates internal
/// state when returning true.
pub fn should_send(&mut self, dest_addr: &NodeAddr) -> bool {
let now = Instant::now();
self.check(dest_addr, Instant::now()) != LimitVerdict::Suppress
}
/// Decide about one error signal at an explicit `now`, reporting whether
/// the destination could be recorded.
///
/// Callers that distinguish the at-capacity admission use this; callers
/// that only need a yes or no use [`Self::should_send`].
pub fn check(&mut self, dest_addr: &NodeAddr, now: Instant) -> LimitVerdict {
if let Some(&last) = self.last_sent.get(dest_addr)
&& now.duration_since(last) < self.min_interval
&& now.saturating_duration_since(last) < self.min_interval
{
return false;
return LimitVerdict::Suppress;
}
if self.last_sent.len() >= MAX_ENTRIES && !self.last_sent.contains_key(dest_addr) {
self.maybe_cleanup(now);
if self.last_sent.len() >= MAX_ENTRIES {
return LimitVerdict::AdmitAtCapacity;
}
}
self.last_sent.insert(*dest_addr, now);
self.cleanup(now);
true
self.maybe_cleanup(now);
LimitVerdict::Admit
}
/// Run the sweep if one is due.
fn maybe_cleanup(&mut self, now: Instant) {
if now.saturating_duration_since(self.last_sweep) >= self.max_age / SWEEPS_PER_MAX_AGE {
self.cleanup(now);
}
}
/// Remove entries older than max_age.
fn cleanup(&mut self, now: Instant) {
self.last_sweep = now;
self.sweeps += 1;
self.last_sent
.retain(|_, &mut last| now.duration_since(last) < self.max_age);
.retain(|_, &mut last| now.saturating_duration_since(last) < self.max_age);
}
#[cfg(test)]
pub fn len(&self) -> usize {
self.last_sent.len()
}
#[cfg(test)]
pub fn sweeps(&self) -> u64 {
self.sweeps
}
}
impl Default for RoutingErrorRateLimiter {
@@ -89,6 +172,15 @@ mod tests {
NodeAddr::from_bytes(bytes)
}
/// A distinct destination address per index, standing for the fresh
/// `dest_addr` a flooding sender puts on every datagram.
fn minted_addr(val: u32) -> NodeAddr {
let mut bytes = [0u8; 16];
bytes[..4].copy_from_slice(&val.to_le_bytes());
bytes[15] = 0xff;
NodeAddr::from_bytes(bytes)
}
#[test]
fn test_first_send_allowed() {
let mut limiter = RoutingErrorRateLimiter::new();
@@ -144,6 +236,71 @@ mod tests {
assert_eq!(limiter.len(), 1);
}
#[test]
fn the_map_stays_bounded_when_a_sender_mints_distinct_destination_keys() {
let mut limiter = RoutingErrorRateLimiter::new();
let now = Instant::now();
for i in 0..100_000u32 {
limiter.check(&minted_addr(i), now);
}
assert!(
limiter.len() <= MAX_ENTRIES,
"limiter held {} entries, above the {MAX_ENTRIES} ceiling",
limiter.len()
);
}
#[test]
fn an_admission_at_capacity_still_sends_rather_than_going_silent() {
let mut limiter = RoutingErrorRateLimiter::new();
let now = Instant::now();
for i in 0..MAX_ENTRIES as u32 {
assert_eq!(limiter.check(&minted_addr(i), now), LimitVerdict::Admit);
}
// The map is full and nothing in it is old enough to evict, so the
// next distinct destination cannot be recorded. It must still be sent.
assert_eq!(
limiter.check(&minted_addr(MAX_ENTRIES as u32), now),
LimitVerdict::AdmitAtCapacity
);
}
#[test]
fn the_map_scan_does_not_run_once_per_admitted_destination() {
let mut limiter = RoutingErrorRateLimiter::new();
let now = Instant::now();
let before = limiter.sweeps();
for i in 0..1_000u32 {
limiter.check(&minted_addr(i), now);
}
assert_eq!(
limiter.sweeps() - before,
0,
"the full-map scan ran inside a single sweep interval"
);
}
#[test]
fn the_map_scan_still_runs_once_a_sweep_interval_has_passed() {
let mut limiter = RoutingErrorRateLimiter::new();
let start = Instant::now();
limiter.check(&minted_addr(0), start);
let before = limiter.sweeps();
let later = start + Duration::from_secs(11);
limiter.check(&minted_addr(1), later);
assert_eq!(limiter.sweeps() - before, 1);
// The first destination aged out, so the sweep did its job.
assert_eq!(limiter.len(), 1);
}
#[test]
fn test_with_interval_custom_rate() {
let mut limiter = RoutingErrorRateLimiter::with_interval(Duration::from_millis(500));
+3
View File
@@ -421,6 +421,9 @@ pub struct ErrorSignalStatsSnapshot {
pub unbound_broken: u64,
pub unbound_mtu: u64,
pub unbound_forged: u64,
pub emit_over_peer_budget: u64,
pub emit_over_dest_interval: u64,
pub emit_limiter_at_capacity: u64,
}
#[derive(Clone, Debug, Default, Serialize)]
+88
View File
@@ -5,6 +5,7 @@
//! multi-hop forwarding through live node topologies.
use super::*;
use crate::node::peer_error_budget::PEER_ERROR_BURST;
use crate::node::session_wire::{FSP_FLAG_CP, build_fsp_header};
use crate::protocol::{SessionAck, SessionDatagram, SessionSetup, encode_coords};
use crate::tree::TreeCoordinate;
@@ -1111,3 +1112,90 @@ fn test_sample_transport_congestion() {
node.sample_transport_congestion();
assert!(!node.transport_drops[&tid].dropping);
}
// --- Emission bounds on induced routing errors ---
/// A distinct destination per index, standing for the fresh `dest_addr` a
/// flooding sender puts on every datagram to escape the per-destination gate.
fn minted_dest(val: u32) -> NodeAddr {
let mut bytes = [0u8; 16];
bytes[..4].copy_from_slice(&val.to_le_bytes());
bytes[15] = 0xfe;
NodeAddr::from_bytes(bytes)
}
/// Feed one transit datagram whose destination this node cannot route.
async fn inject_unroutable(node: &mut Node, from: &NodeAddr, src: NodeAddr, dest: NodeAddr) {
let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(8);
let encoded = dg.encode();
node.handle_session_datagram(from, &encoded[1..], false)
.await;
}
#[tokio::test]
async fn one_link_peer_cannot_induce_unbounded_errors_by_varying_the_destination() {
let mut node = make_node();
let attacker = make_node_addr(0xAA);
let overshoot = 10u32;
for i in 0..PEER_ERROR_BURST + overshoot {
// Fresh destination and fresh spoofed source per packet: neither
// address-keyed gate sees a repeat.
inject_unroutable(
&mut node,
&attacker,
minted_dest(i + 1_000_000),
minted_dest(i),
)
.await;
}
let errors = &node.metrics().errors;
assert_eq!(
errors.emit_over_dest_interval.get(),
0,
"the per-destination gate cannot bound a sender that varies the destination"
);
assert_eq!(
errors.emit_over_peer_budget.get(),
u64::from(overshoot),
"everything past the link peer's burst must be refused"
);
}
#[tokio::test]
async fn a_destination_suppressed_error_does_not_spend_the_link_peer_budget() {
let mut node = make_node();
let peer = make_node_addr(0xAA);
let src = make_node_addr(0x01);
let dest = make_node_addr(0x02);
let injected = PEER_ERROR_BURST * 4;
for _ in 0..injected {
inject_unroutable(&mut node, &peer, src, dest).await;
}
let errors = &node.metrics().errors;
assert_eq!(
errors.emit_over_peer_budget.get(),
0,
"an outage on one destination must not spend the peer's budget for the others"
);
assert_eq!(
errors.emit_over_dest_interval.get(),
u64::from(injected - 1),
"only the first error for a destination goes out within the interval"
);
}
#[tokio::test]
async fn a_single_unroutable_datagram_still_produces_its_error() {
let mut node = make_node();
let peer = make_node_addr(0xAA);
inject_unroutable(&mut node, &peer, make_node_addr(0x01), make_node_addr(0x02)).await;
let errors = &node.metrics().errors;
assert_eq!(errors.emit_over_peer_budget.get(), 0);
assert_eq!(errors.emit_over_dest_interval.get(), 0);
}