Evict from the lookup dedup cache instead of refusing the request

The discovery dedup cache is also the reverse-path table for responses in
flight, and at its 4096-entry bound it dropped the arriving request. That
drop sat ahead of both the check for whether the request names this node
and the forwarding path, so one link peer emitting fresh request_ids could
stop the node answering lookups for itself and stop it carrying anyone
else's, for as long as it kept the cache full.

Make room instead of refusing. An arrival-order index partitioned by the
link peer the request came from says who pays: a peer over its own share
loses its oldest entry, and at global capacity the peer holding the most
entries loses its oldest. A light peer's reverse path is therefore never
taken to admit a heavy one, and extra identities buy a flooder
proportionally less. A share is the cache divided by the current link-peer
count with a floor of 64, so it tracks the peer count rather than being
pinned to a number a many-peer node outgrows. The loosening this accepts
is that an evicted request_id arriving again inside the window is
forwarded a second time rather than recognised as a duplicate; the
per-target forward limiter and TTL already bound that.

Meter answering lookups for ourselves per link peer, in the same change,
because the cache filling up was the only thing bounding it. The response
proof is signed over the requester's request_id, so every request
addressed to this node costs a fresh Schnorr signature that cannot be
cached or served twice. A token bucket of 256 signatures refilling at 32
per second per link peer absorbs the legitimate burst that follows a
topology change, when many correspondents re-look-up at once through the
few links that lead here, while capping what one neighbour can make the
node sign. A refused request keeps its dedup entry, and retries carry
fresh request_ids, so a refusal cannot suppress the retry.

Evictions count as req_dedup_evicted and signing refusals as
req_sign_rate_limited, both in show routing, show metrics and the fipstop
routing pane. The old req_dedup_cache_full counter stays in place, frozen
at zero, so a dashboard carried across versions does not lose the series.
This commit is contained in:
Johnathan Corgan
2026-08-23 11:46:53 +01:00
parent cbe35f1cac
commit 65321617ae
10 changed files with 516 additions and 19 deletions
+32
View File
@@ -1077,6 +1077,38 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
because a request is flooded to every qualifying tree peer and the
duplicate replies land there once the first has been accepted.
- A flooded discovery dedup cache no longer makes a node unresolvable. The
cache is both the duplicate filter and the reverse-path table for lookup
responses, and at its 4096-entry bound it dropped the arriving request.
That drop sat ahead of both the check for whether the request names this
node and the forwarding path, so one link peer emitting fresh request_ids
could stop the node answering lookups for itself and stop it carrying
anyone else's, for as long as it kept the cache full. The cache now makes
room instead of refusing: over a peer's own share it drops that peer's
oldest entry, and at global capacity it drops the oldest entry of whichever
peer holds the most, so a light peer's reverse path is never taken to admit
a heavy one and extra identities buy a flooder proportionally less. A
peer's share is the cache divided by the current link-peer count, with a
floor of 64. The loosening this accepts is that an evicted request_id
arriving again inside the window is forwarded a second time rather than
recognised as a duplicate, which the per-target forward limiter and TTL
already bound. Evictions are counted as `req_dedup_evicted`; the old
`req_dedup_cache_full` counter stays in place, frozen at zero, so a
dashboard carried across versions does not lose the series.
- Answering a lookup for ourselves is now metered per link peer. The response
proof is signed over the requester's `request_id`, so every request
addressed to this node costs a fresh Schnorr signature that cannot be
cached or served twice, and until now the only thing bounding that rate was
the dedup cache filling up, which is the defect above. A token bucket per
link peer, 256 signatures of burst refilling at 32 per second, absorbs the
legitimate burst that follows a topology change, when many correspondents
re-look-up at once through the few links that lead here, while capping what
one neighbour can make the node sign. Refusals are counted as
`req_sign_rate_limited` and visible in `show routing`, `show metrics` and
the fipstop routing pane. A refused request keeps its dedup entry, and
retries carry fresh request_ids, so a refusal cannot suppress the retry.
#### Admission / peer caps
- The Ethernet transport's discovery buffer is now bounded and no longer costs
+2
View File
@@ -192,6 +192,8 @@ fn draw_routing_stats(
("Bloom Miss", disc("req_bloom_miss")),
("Backoff Suppressed", disc("req_backoff_suppressed")),
("Fwd Rate Limited", disc("req_forward_rate_limited")),
("Sign Rate Limited", disc("req_sign_rate_limited")),
("Dedup Evicted", disc("req_dedup_evicted")),
("TTL Exhausted", disc("req_ttl_exhausted")),
("Decode Error", disc("req_decode_error")),
],
+2
View File
@@ -12,6 +12,7 @@
"req_bloom_miss": 0,
"req_decode_error": 0,
"req_dedup_cache_full": 0,
"req_dedup_evicted": 0,
"req_deduplicated": 0,
"req_duplicate": 0,
"req_fallback_forwarded": 0,
@@ -20,6 +21,7 @@
"req_initiated": 0,
"req_no_tree_peer": 0,
"req_received": 0,
"req_sign_rate_limited": 0,
"req_target_is_us": 0,
"req_ttl_exhausted": 0,
"resp_accepted": 0,
+138
View File
@@ -12,6 +12,11 @@
//! - **`DiscoveryForwardRateLimiter`** (transit-side): Per-target minimum
//! interval for forwarded requests. Defense-in-depth against misbehaving
//! nodes generating fresh request_ids at high rate.
//!
//! - **`LookupSignRateLimiter`** (target-side): Per-link-peer token bucket
//! on answering lookups for ourselves. Every such answer costs a fresh
//! Schnorr signature, because the proof is bound to the requester's
//! `request_id` and so cannot be cached or reused.
use crate::NodeAddr;
use std::collections::HashMap;
@@ -222,6 +227,113 @@ impl Default for DiscoveryForwardRateLimiter {
}
}
// ============================================================================
// Target-side: Lookup Signing Budget
// ============================================================================
/// Signatures one link peer may buy in a burst before the refill paces it.
///
/// Sized for the case that actually produces a burst: a topology change
/// flushes correspondents' coordinate caches and they all look this node up
/// at once, through whichever few link peers lead here, each retrying on the
/// `node.discovery.attempt_timeouts_secs` ladder. Lowering this makes a
/// genuinely popular node intermittently unresolvable, which is the same
/// symptom as the flood it defends against; raising it raises the worst-case
/// signing burst one neighbour can force.
const DEFAULT_SIGN_BURST: f64 = 256.0;
/// Sustained signatures per second per link peer.
///
/// At the default eight or so link peers this caps the node near 256
/// signatures per second in the sustained case. The real cost of one
/// `Identity::sign` on this codebase has not been measured, so this number
/// is a bound rather than a tuned value; it is the one line to change if a
/// measurement says otherwise.
const DEFAULT_SIGN_RATE: f64 = 32.0;
/// Maximum age of an idle bucket before cleanup.
const SIGN_MAX_AGE: Duration = Duration::from_secs(300);
/// Token bucket per link peer for lookups this node answers about itself.
///
/// A min-interval limiter is the wrong shape here: a popular node receives
/// legitimate bursts of lookups for itself through the few link peers that
/// lead to it, and a min interval refuses all but the first of each burst.
/// A bucket absorbs the burst and paces the sustained rate.
pub struct LookupSignRateLimiter {
buckets: HashMap<NodeAddr, SignBucket>,
burst: f64,
rate: f64,
}
struct SignBucket {
/// Tokens remaining, at most `burst`.
tokens: f64,
/// When `tokens` was last refilled.
updated: Instant,
}
impl LookupSignRateLimiter {
/// Create with default burst and refill rate.
pub fn new() -> Self {
Self::with_params(DEFAULT_SIGN_BURST, DEFAULT_SIGN_RATE)
}
/// Create with a custom burst and refill rate.
pub fn with_params(burst: f64, rate: f64) -> Self {
Self {
buckets: HashMap::new(),
burst,
rate,
}
}
/// Spend one token for `from`, or report that its budget is exhausted.
///
/// Returns true when the signature may be produced. A zero burst is
/// read as "unlimited" rather than "refuse everything", so a
/// misconfiguration cannot make this node unresolvable.
pub fn should_sign(&mut self, from: &NodeAddr) -> bool {
if self.burst <= 0.0 {
return true;
}
let now = Instant::now();
let burst = self.burst;
let rate = self.rate;
let bucket = self.buckets.entry(*from).or_insert(SignBucket {
tokens: burst,
updated: now,
});
let elapsed = now.duration_since(bucket.updated).as_secs_f64();
bucket.tokens = (bucket.tokens + elapsed * rate).min(burst);
bucket.updated = now;
if bucket.tokens < 1.0 {
return false;
}
bucket.tokens -= 1.0;
self.cleanup(now);
true
}
/// Drop buckets untouched for longer than [`SIGN_MAX_AGE`]; a full
/// bucket carries no state worth keeping.
fn cleanup(&mut self, now: Instant) {
self.buckets
.retain(|_, b| now.duration_since(b.updated) < SIGN_MAX_AGE);
}
#[cfg(test)]
pub fn len(&self) -> usize {
self.buckets.len()
}
}
impl Default for LookupSignRateLimiter {
fn default() -> Self {
Self::new()
}
}
// ============================================================================
// Tests
// ============================================================================
@@ -373,4 +485,30 @@ mod tests {
limiter.cleanup(Instant::now());
assert_eq!(limiter.len(), 1);
}
#[test]
fn test_sign_budget_is_spent_per_peer_and_does_not_touch_another_peer() {
let mut limiter = LookupSignRateLimiter::with_params(4.0, 0.0);
for _ in 0..4 {
assert!(limiter.should_sign(&addr(1)));
}
assert!(
!limiter.should_sign(&addr(1)),
"the burst is the whole budget when nothing refills it"
);
assert!(
limiter.should_sign(&addr(2)),
"one peer spending its budget must not spend another's"
);
assert_eq!(limiter.len(), 2);
}
#[test]
fn test_sign_budget_of_zero_burst_is_read_as_unlimited() {
let mut limiter = LookupSignRateLimiter::with_params(0.0, 0.0);
for _ in 0..1000 {
assert!(limiter.should_sign(&addr(1)));
}
assert_eq!(limiter.len(), 0, "unlimited keeps no per-peer state");
}
}
+98 -17
View File
@@ -12,7 +12,19 @@ use crate::transport::{TransportAddr, TransportId};
use crate::{NodeAddr, PeerIdentity};
use tracing::{debug, info, trace, warn};
const MAX_RECENT_DISCOVERY_REQUESTS: usize = 4096;
/// Cap on the discovery request dedup cache, which is also the reverse-path
/// table for responses in flight.
pub(in crate::node) const MAX_RECENT_DISCOVERY_REQUESTS: usize = 4096;
/// Floor under one link peer's share of the dedup cache.
///
/// A peer's share is the cache divided by the current link-peer count, and
/// this is what stops that share collapsing to nothing on a node with very
/// many links. It is a cap and not a reservation: shares can sum past the
/// cache size, in which case the peer holding the most entries pays for the
/// next admission. Raising it lets one busy neighbour hold more of the
/// cache; lowering it clips a genuine transit burst.
pub(in crate::node) const MIN_RECENT_PER_PEER: usize = 64;
impl Node {
/// Handle an incoming LookupRequest from a peer.
@@ -56,26 +68,44 @@ impl Node {
return;
}
if self.recent_requests.len() >= MAX_RECENT_DISCOVERY_REQUESTS {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqDedupCacheFull);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
recent_requests = self.recent_requests.len(),
max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS,
"Discovery request dedup cache full, dropping LookupRequest"
);
return;
}
// A full cache evicts rather than refuses. Refusing meant one peer
// could fill the cache with fresh request_ids and stop this node
// answering lookups for itself and forwarding anyone else's, which
// is a denial of the service the cache exists to protect. The
// eviction is charged to the peer that filled the cache: over its
// own share it pays for itself, and at global capacity the peer
// holding the most entries pays, so extra identities buy a flooder
// proportionally less and a light peer's reverse path survives.
self.make_room_for_request(from);
// Record for reverse-path forwarding and dedup
self.recent_requests
.insert(request.request_id, RecentRequest::new(*from, now_ms));
self.recent_by_peer
.entry(*from)
.or_default()
.push_back(request.request_id);
// Are we the target?
if request.target == *self.node_addr() {
// Answering costs a fresh Schnorr signature every time: the
// proof is bound to the requester's request_id, so it cannot be
// cached or served twice. Meter that per link peer, or a
// neighbour generating request_ids sets this node's signing
// rate. The dedup entry above stays regardless, so a refused
// request still occupies its id and a retry, which carries a
// fresh id, is unaffected.
if !self.discovery_sign_limiter.should_sign(from) {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqSignRateLimited);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
"Lookup signing budget spent for this peer, not answering"
);
return;
}
self.metrics().discovery.req_target_is_us.inc();
debug!(
request_id = request.request_id,
@@ -715,10 +745,61 @@ impl Node {
}
/// Remove expired entries from the recent_requests cache.
fn purge_expired_requests(&mut self, current_time_ms: u64) {
pub(in crate::node) fn purge_expired_requests(&mut self, current_time_ms: u64) {
let expiry_ms = self.config().node.discovery.recent_expiry_secs * 1000;
self.recent_requests
.retain(|_, entry| !entry.is_expired(current_time_ms, expiry_ms));
let recent = &mut self.recent_requests;
recent.retain(|_, entry| !entry.is_expired(current_time_ms, expiry_ms));
self.recent_by_peer.retain(|_, ids| {
ids.retain(|id| recent.contains_key(id));
!ids.is_empty()
});
}
/// Evict from the dedup cache if admitting one more request would put
/// this peer over its share, or the cache over its capacity.
///
/// The share is the cache divided by the current link-peer count, with
/// [`MIN_RECENT_PER_PEER`] as a floor, so it tracks the peer count
/// instead of being pinned to a number that a many-peer node outgrows.
fn make_room_for_request(&mut self, from: &NodeAddr) {
let share =
(MAX_RECENT_DISCOVERY_REQUESTS / self.peers.len().max(1)).max(MIN_RECENT_PER_PEER);
let over_share = self
.recent_by_peer
.get(from)
.is_some_and(|ids| ids.len() >= share);
let victim = if over_share {
Some(*from)
} else if self.recent_requests.len() >= MAX_RECENT_DISCOVERY_REQUESTS {
// Never take from a peer under its share: charge the fattest.
self.recent_by_peer
.iter()
.max_by_key(|(_, ids)| ids.len())
.map(|(peer, _)| *peer)
} else {
return;
};
let Some(victim) = victim else { return };
let Some(ids) = self.recent_by_peer.get_mut(&victim) else {
return;
};
let Some(evicted) = ids.pop_front() else {
return;
};
if ids.is_empty() {
self.recent_by_peer.remove(&victim);
}
self.recent_requests.remove(&evicted);
self.metrics().discovery.req_dedup_evicted.inc();
debug!(
request_id = evicted,
evicted_from = %self.peer_display_name(&victim),
admitting = %self.peer_display_name(from),
share = share,
"Discovery dedup cache full, evicting the oldest entry to make room"
);
}
/// Min-fold our outgoing-link MTU into a LookupResponse's `path_mtu`.
+5
View File
@@ -266,6 +266,8 @@ pub struct DiscoveryMetrics {
pub req_decode_error: Counter,
pub req_duplicate: Counter,
pub req_dedup_cache_full: Counter,
pub req_dedup_evicted: Counter,
pub req_sign_rate_limited: Counter,
pub req_target_is_us: Counter,
pub req_forwarded: Counter,
pub req_ttl_exhausted: Counter,
@@ -296,6 +298,7 @@ impl DiscoveryMetrics {
DiscoveryReject::ReqDecodeError => self.req_decode_error.inc(),
DiscoveryReject::ReqDuplicate => self.req_duplicate.inc(),
DiscoveryReject::ReqDedupCacheFull => self.req_dedup_cache_full.inc(),
DiscoveryReject::ReqSignRateLimited => self.req_sign_rate_limited.inc(),
DiscoveryReject::ReqTtlExhausted => self.req_ttl_exhausted.inc(),
DiscoveryReject::RespDecodeError => self.resp_decode_error.inc(),
DiscoveryReject::RespIdentityMiss => self.resp_identity_miss.inc(),
@@ -312,6 +315,8 @@ impl DiscoveryMetrics {
req_decode_error: self.req_decode_error.get(),
req_duplicate: self.req_duplicate.get(),
req_dedup_cache_full: self.req_dedup_cache_full.get(),
req_dedup_evicted: self.req_dedup_evicted.get(),
req_sign_rate_limited: self.req_sign_rate_limited.get(),
req_target_is_us: self.req_target_is_us.get(),
req_forwarded: self.req_forwarded.get(),
req_ttl_exhausted: self.req_ttl_exhausted.get(),
+24 -2
View File
@@ -32,7 +32,9 @@ mod tests;
mod tree;
pub(crate) mod wire;
use self::discovery_rate_limit::{DiscoveryBackoff, DiscoveryForwardRateLimiter};
use self::discovery_rate_limit::{
DiscoveryBackoff, DiscoveryForwardRateLimiter, LookupSignRateLimiter,
};
use self::peer_error_budget::PeerErrorBudget;
use self::rate_limit::{HandshakeRateLimiter, SessionSetupRateLimiter};
use self::reloadable::Reloadable;
@@ -69,7 +71,7 @@ use crate::upper::tun::{TunError, TunOutboundRx, TunState, TunTx};
use crate::utils::index::IndexAllocator;
use crate::{Config, ConfigError, Identity, IdentityError, NodeAddr, PeerIdentity};
use rand::Rng;
use std::collections::{HashMap, HashSet, VecDeque};
use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
use std::fmt;
use std::sync::Arc;
use std::thread::JoinHandle;
@@ -367,6 +369,13 @@ pub struct Node {
/// Recent discovery requests (dedup + reverse-path forwarding).
/// Maps request_id → RecentRequest.
recent_requests: HashMap<u64, RecentRequest>,
/// Arrival-order index over `recent_requests`, partitioned by the link
/// peer each request arrived from. The cache is full-then-evict rather
/// than full-then-refuse, and this is what lets an eviction be charged
/// to the peer that filled the cache instead of to whoever happens to
/// be oldest. Timestamps are nondecreasing across inserts, so each
/// deque is in arrival order and the front is the oldest.
recent_by_peer: BTreeMap<NodeAddr, VecDeque<u64>>,
/// Per-destination path MTU lookup, keyed by FipsAddress (mirrors
/// `coord_cache.entries[*].path_mtu`). Sync read-only access from
/// the TUN reader/writer threads at TCP MSS clamp time so the
@@ -516,6 +525,8 @@ pub struct Node {
discovery_backoff: DiscoveryBackoff,
/// Rate limiter for forwarded discovery requests (transit-side).
discovery_forward_limiter: DiscoveryForwardRateLimiter,
/// Signing budget for lookups we answer about ourselves (target-side).
discovery_sign_limiter: LookupSignRateLimiter,
// === Pending Transport Connects ===
/// Links waiting for transport-level connection establishment before
@@ -761,6 +772,7 @@ impl Node {
bloom_state,
coord_cache,
recent_requests: HashMap::new(),
recent_by_peer: BTreeMap::new(),
transports: HashMap::new(),
transport_drops: HashMap::new(),
links: HashMap::new(),
@@ -813,6 +825,7 @@ impl Node {
discovery_forward_limiter: DiscoveryForwardRateLimiter::with_interval(
std::time::Duration::from_secs(forward_min_interval_secs),
),
discovery_sign_limiter: LookupSignRateLimiter::new(),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
nostr_discovery: None,
@@ -922,6 +935,7 @@ impl Node {
bloom_state,
coord_cache,
recent_requests: HashMap::new(),
recent_by_peer: BTreeMap::new(),
transports: HashMap::new(),
transport_drops: HashMap::new(),
links: HashMap::new(),
@@ -972,6 +986,7 @@ impl Node {
),
discovery_backoff: DiscoveryBackoff::new(),
discovery_forward_limiter: DiscoveryForwardRateLimiter::new(),
discovery_sign_limiter: LookupSignRateLimiter::new(),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
nostr_discovery: None,
@@ -2547,6 +2562,13 @@ impl Node {
// === End-to-End Sessions ===
/// Get a session by remote NodeAddr.
/// Set the per-link-peer lookup signing budget (for tests).
#[cfg(test)]
pub(crate) fn set_discovery_sign_budget(&mut self, burst: f64, rate: f64) {
self.discovery_sign_limiter =
discovery_rate_limit::LookupSignRateLimiter::with_params(burst, rate);
}
/// Disable the discovery forward rate limiter (for tests).
#[cfg(test)]
pub(crate) fn disable_discovery_forward_rate_limit(&mut self) {
+13
View File
@@ -120,7 +120,19 @@ pub enum DiscoveryReject {
/// Request dedup cache (`recent_requests`) is at capacity, so the
/// `LookupRequest` is dropped without being forwarded. Tracked via
/// [`DiscoveryStats::req_dedup_cache_full`](crate::node::stats::DiscoveryStats).
///
/// Frozen at zero: a full cache now evicts its oldest entry and admits
/// the request, counted as
/// [`DiscoveryStats::req_dedup_evicted`](crate::node::stats::DiscoveryStats).
/// The variant and its counter stay so an operator reading a dashboard
/// across versions does not find the series missing.
ReqDedupCacheFull,
/// This node is the lookup target, but the link peer the request
/// arrived from has spent its signing budget. Answering costs a fresh
/// Schnorr signature per request, so the budget bounds what one
/// neighbour can make this node sign. Tracked via
/// [`DiscoveryStats::req_sign_rate_limited`](crate::node::stats::DiscoveryStats).
ReqSignRateLimited,
/// Request arrived with TTL=0 — no more forwarding hops allowed.
/// Tracked via
/// [`DiscoveryStats::req_ttl_exhausted`](crate::node::stats::DiscoveryStats).
@@ -376,6 +388,7 @@ mod tests {
DiscoveryReject::ReqDuplicate,
DiscoveryReject::ReqDedupCacheFull,
DiscoveryReject::ReqTtlExhausted,
DiscoveryReject::ReqSignRateLimited,
DiscoveryReject::RespDecodeError,
DiscoveryReject::RespIdentityMiss,
DiscoveryReject::RespProofFailed,
+2
View File
@@ -319,6 +319,8 @@ pub struct DiscoveryStatsSnapshot {
pub req_decode_error: u64,
pub req_duplicate: u64,
pub req_dedup_cache_full: u64,
pub req_dedup_evicted: u64,
pub req_sign_rate_limited: u64,
pub req_target_is_us: u64,
pub req_forwarded: u64,
pub req_ttl_exhausted: u64,
+200
View File
@@ -614,6 +614,206 @@ async fn test_recent_request_expiry() {
assert!(node.recent_requests.contains_key(&789));
}
// ============================================================================
// Unit Tests — dedup cache capacity policy
// ============================================================================
use crate::node::handlers::discovery::{MAX_RECENT_DISCOVERY_REQUESTS, MIN_RECENT_PER_PEER};
/// Encode a LookupRequest for `target` carrying `request_id`, ready for
/// `handle_lookup_request` (which is handed the payload without the
/// msg_type byte).
fn lookup_request_payload(request_id: u64, target: &crate::NodeAddr) -> Vec<u8> {
let origin = make_node_addr(0xCC);
let coords = TreeCoordinate::from_addrs(vec![origin, make_node_addr(0)]).unwrap();
LookupRequest::new(request_id, *target, origin, coords, 5, 0).encode()[1..].to_vec()
}
/// Deliver `count` distinct requests from `from`, ids starting at `first_id`.
async fn flood_requests(node: &mut Node, from: &crate::NodeAddr, first_id: u64, count: u64) {
let target = make_node_addr(0xBB);
for i in 0..count {
let payload = lookup_request_payload(first_id + i, &target);
node.handle_lookup_request(from, &payload).await;
}
}
/// Register `count` peers so the per-peer share of the dedup cache is the
/// floor rather than the whole cache, and return their addresses.
fn register_peers(node: &mut Node, count: usize) -> Vec<crate::NodeAddr> {
(0..count)
.map(|i| {
let identity = Identity::generate();
let addr = *identity.node_addr();
let peer_identity = crate::PeerIdentity::from_pubkey(identity.pubkey());
node.peers.insert(
addr,
ActivePeer::new(peer_identity, LinkId::new(i as u64), 0),
);
addr
})
.collect()
}
#[tokio::test]
async fn test_a_full_dedup_cache_admits_the_new_request_by_evicting_the_oldest() {
// A full cache used to drop the arriving request, which let one peer
// spend 4096 fresh request_ids and stop the node forwarding anyone
// else's lookups until the entries aged out.
let mut node = make_node();
let from = make_node_addr(0xAA);
flood_requests(&mut node, &from, 1, MAX_RECENT_DISCOVERY_REQUESTS as u64).await;
assert_eq!(
node.recent_requests.len(),
MAX_RECENT_DISCOVERY_REQUESTS,
"precondition: the cache is full, or the rest observes nothing"
);
let payload = lookup_request_payload(u64::MAX, &make_node_addr(0xBB));
node.handle_lookup_request(&from, &payload).await;
assert!(
node.recent_requests.contains_key(&u64::MAX),
"the arriving request must be recorded, so its response can be routed back"
);
assert!(
!node.recent_requests.contains_key(&1),
"room is made by dropping the oldest entry"
);
assert_eq!(
node.recent_requests.len(),
MAX_RECENT_DISCOVERY_REQUESTS,
"the cache stays at its bound"
);
assert_eq!(node.metrics().discovery.req_dedup_evicted.get(), 1);
assert_eq!(
node.metrics().discovery.req_dedup_cache_full.get(),
0,
"the cache-full drop is gone, and its counter stays frozen at zero"
);
}
#[tokio::test]
async fn test_a_flooding_peer_evicts_only_its_own_dedup_entries() {
// The whole point of partitioning the cache by link peer: one peer
// filling its share must not cost another peer the reverse path its own
// lookup depends on.
let mut node = make_node();
let peers = register_peers(&mut node, 64);
let flooder = peers[0];
let light = peers[1];
let payload = lookup_request_payload(7, &make_node_addr(0xBB));
node.handle_lookup_request(&light, &payload).await;
// One over the share, so the flooder pays for its own admission.
flood_requests(&mut node, &flooder, 1000, MIN_RECENT_PER_PEER as u64 + 1).await;
assert!(
node.recent_requests.contains_key(&7),
"a light peer's reverse-path entry must survive a neighbour's flood"
);
assert!(
!node.recent_requests.contains_key(&1000),
"the flooder's own oldest entry is what pays for its newest"
);
assert!(
node.recent_requests
.contains_key(&(1000 + MIN_RECENT_PER_PEER as u64)),
"and its newest is admitted rather than dropped"
);
}
#[tokio::test]
async fn test_a_node_whose_dedup_cache_is_flooded_still_answers_a_lookup_for_itself() {
// The availability claim. Filling the cache used to make the node
// unresolvable, because the cache-full drop sat ahead of the check for
// whether the request names us.
let mut node = make_node();
let flooder = make_node_addr(0xAA);
let other = make_node_addr(0xAB);
flood_requests(&mut node, &flooder, 1, MAX_RECENT_DISCOVERY_REQUESTS as u64).await;
let my_addr = *node.node_addr();
let payload = lookup_request_payload(u64::MAX, &my_addr);
node.handle_lookup_request(&other, &payload).await;
assert_eq!(
node.metrics().discovery.req_target_is_us.get(),
1,
"a flooded cache must not stop the node answering lookups for itself"
);
}
#[tokio::test]
async fn test_the_dedup_index_stays_level_with_the_cache_across_insert_duplicate_and_purge() {
// Two containers where there was one, so the desync is the maintenance
// risk. Everything the eviction policy decides reads the index, so an
// index that has drifted evicts the wrong entry or none at all.
let mut node = make_node();
let peers = register_peers(&mut node, 64);
flood_requests(&mut node, &peers[0], 1, 70).await;
flood_requests(&mut node, &peers[1], 500, 5).await;
// Duplicates, which must not be indexed twice.
flood_requests(&mut node, &peers[1], 500, 5).await;
let indexed: usize = node.recent_by_peer.values().map(|ids| ids.len()).sum();
assert_eq!(
indexed,
node.recent_requests.len(),
"every cached request is indexed exactly once"
);
// Age everything out and purge through the ordinary request path.
let expiry_ms = node.config().node.discovery.recent_expiry_secs * 1000;
let future = Node::now_ms() + expiry_ms + 1;
node.purge_expired_requests(future);
assert!(
node.recent_requests.is_empty(),
"precondition: the purge removed everything"
);
assert!(
node.recent_by_peer.is_empty(),
"the index must not keep entries the cache no longer holds"
);
}
#[tokio::test]
async fn test_answering_lookups_for_ourselves_stops_at_the_per_peer_signing_budget() {
// Each answer costs a fresh Schnorr signature, because the proof is
// bound to the requester's request_id and cannot be reused. Without a
// budget, one neighbour sets this node's signing rate.
let mut node = make_node();
node.set_discovery_sign_budget(3.0, 0.0);
let from = make_node_addr(0xAA);
let other = make_node_addr(0xAB);
let my_addr = *node.node_addr();
for id in 0..4u64 {
let payload = lookup_request_payload(id, &my_addr);
node.handle_lookup_request(&from, &payload).await;
}
assert_eq!(
node.metrics().discovery.req_target_is_us.get(),
3,
"the burst is answered and the fourth request is not"
);
assert_eq!(node.metrics().discovery.req_sign_rate_limited.get(), 1);
let payload = lookup_request_payload(100, &my_addr);
node.handle_lookup_request(&other, &payload).await;
assert_eq!(
node.metrics().discovery.req_target_is_us.get(),
4,
"one peer spending its budget must not make the node unresolvable through another"
);
}
// ============================================================================
// Integration Tests — Multi-Node Forwarding
// ============================================================================