diff --git a/CHANGELOG.md b/CHANGELOG.md index f09f5eb3..a1bf8244 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/src/bin/fipstop/ui/routing.rs b/src/bin/fipstop/ui/routing.rs index 819541e4..015234d6 100644 --- a/src/bin/fipstop/ui/routing.rs +++ b/src/bin/fipstop/ui/routing.rs @@ -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("")); diff --git a/src/bin/fipstop/ui/snapshots.rs b/src/bin/fipstop/ui/snapshots.rs index b108cc66..a91611b4 100644 --- a/src/bin/fipstop/ui/snapshots.rs +++ b/src/bin/fipstop/ui/snapshots.rs @@ -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); }); diff --git a/src/control/snapshots/show_routing.json b/src/control/snapshots/show_routing.json index 9741572e..1e5328e2 100644 --- a/src/control/snapshots/show_routing.json +++ b/src/control/snapshots/show_routing.json @@ -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, diff --git a/src/node/handlers/forwarding.rs b/src/node/handlers/forwarding.rs index c0b4dd42..3124231f 100644 --- a/src/node/handlers/forwarding.rs +++ b/src/node/handlers/forwarding.rs @@ -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; } diff --git a/src/node/metrics.rs b/src/node/metrics.rs index e81457b7..6dfc0bb3 100644 --- a/src/node/metrics.rs +++ b/src/node/metrics.rs @@ -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(), } } } diff --git a/src/node/mod.rs b/src/node/mod.rs index d88fcb6b..3a52a44f 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -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), ), diff --git a/src/node/peer_error_budget.rs b/src/node/peer_error_budget.rs new file mode 100644 index 00000000..785337bc --- /dev/null +++ b/src/node/peer_error_budget.rs @@ -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, + /// 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); + } +} diff --git a/src/node/routing_error_rate_limit.rs b/src/node/routing_error_rate_limit.rs index 32cdab48..be8fff32 100644 --- a/src/node/routing_error_rate_limit.rs +++ b/src/node/routing_error_rate_limit.rs @@ -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)); diff --git a/src/node/stats.rs b/src/node/stats.rs index acf5f7c9..d92b838e 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -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)] diff --git a/src/node/tests/forwarding.rs b/src/node/tests/forwarding.rs index 4ba0db08..c71bd9e7 100644 --- a/src/node/tests/forwarding.rs +++ b/src/node/tests/forwarding.rs @@ -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); +}