use std::net::SocketAddr;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use dashmap::DashMap;
use freenet_stdlib::prelude::ContractInstanceId;
use tokio::time::Instant;
use super::Ring;
use super::resync_rate_limit::{BucketOutcome, TokenBucketLimiter};
use crate::util::time_source::TimeSource;
pub(crate) const MIN_UPDATE_INTERVAL: Duration = Duration::from_millis(100);
pub(crate) const MIN_BROADCAST_INTERVAL: Duration = Duration::from_millis(20);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum UpdateClass {
Request,
Broadcast,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct PairStamps {
request: Option<Instant>,
broadcast: Option<Instant>,
latest: Instant,
}
impl UpdateClass {
fn index(self) -> usize {
match self {
UpdateClass::Request => 0,
UpdateClass::Broadcast => 1,
}
}
fn as_str(self) -> &'static str {
match self {
UpdateClass::Request => "request",
UpdateClass::Broadcast => "broadcast",
}
}
}
impl PairStamps {
fn new_at(class: UpdateClass, now: Instant) -> Self {
let mut stamps = Self {
request: None,
broadcast: None,
latest: now,
};
stamps.set(class, now);
stamps
}
fn get(&self, class: UpdateClass) -> Option<Instant> {
match class {
UpdateClass::Request => self.request,
UpdateClass::Broadcast => self.broadcast,
}
}
fn set(&mut self, class: UpdateClass, now: Instant) {
match class {
UpdateClass::Request => self.request = Some(now),
UpdateClass::Broadcast => self.broadcast = Some(now),
}
self.latest = self.latest.max(now);
}
fn latest(&self) -> Instant {
self.latest
}
}
pub(crate) const CLEANUP_AGE: Duration = Duration::from_secs(5 * 60);
pub(crate) const MAX_TRACKED_PAIRS: usize = 16_384;
const EVICTION_BATCH_DIVISOR: usize = 64;
const EVICTION_LOG_INTERVAL: Duration = Duration::from_secs(60);
const MAX_ADMISSION_ATTEMPTS: usize = 6;
const NEW_PAIR_BURST: f64 = 1024.0;
const NEW_PAIR_REFILL_INTERVAL: Duration = Duration::from_millis(5);
const SENDER_TRACKING_HEADROOM: usize = 8;
const MIN_TRACKED_SENDERS: usize = 64;
const NEW_PAIR_LOG_INTERVAL: Duration = EVICTION_LOG_INTERVAL;
const REJECTED_LOG_INTERVAL: Duration = EVICTION_LOG_INTERVAL;
#[cfg_attr(test, derive(Debug, PartialEq, Eq))]
enum EvictionOutcome {
Reserved,
Retry,
CapIsZero,
}
fn log_due(slot: &Mutex<Option<Instant>>, now: Instant, interval: Duration) -> bool {
let mut last = match slot.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
let due = match *last {
Some(prev) => now.saturating_duration_since(prev) >= interval,
None => true,
};
if due {
*last = Some(now);
}
due
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RateLimitDecision {
Allowed,
Rejected {
elapsed: Duration,
min_interval: Duration,
},
CapacityExceeded,
SenderNewPairBudget,
}
impl RateLimitDecision {
pub fn is_allowed(self) -> bool {
matches!(self, RateLimitDecision::Allowed)
}
}
pub(crate) struct UpdateRateLimiter {
last_accepted: DashMap<(SocketAddr, ContractInstanceId), PairStamps>,
size: AtomicUsize,
min_interval: Duration,
broadcast_interval: Duration,
max_tracked_pairs: usize,
time_source: Arc<dyn TimeSource + Send + Sync>,
accepted_total: AtomicU64,
rejected_total: AtomicU64,
capacity_rejected_total: AtomicU64,
capacity_evicted_total: AtomicU64,
last_eviction_log: Mutex<Option<Instant>>,
eviction_lock: Mutex<()>,
new_pair_budget: TokenBucketLimiter<SocketAddr>,
new_pair_budget_rejected_total: AtomicU64,
new_pair_budget_untracked_total: AtomicU64,
last_new_pair_log: Mutex<Option<Instant>>,
last_rejected_log: [Mutex<Option<Instant>>; 2],
}
impl UpdateRateLimiter {
pub fn new(time_source: Arc<dyn TimeSource + Send + Sync>, max_connections: usize) -> Self {
Self::with_new_pair_budget(
time_source,
MIN_UPDATE_INTERVAL,
MIN_BROADCAST_INTERVAL,
MAX_TRACKED_PAIRS,
NEW_PAIR_BURST,
max_connections
.saturating_mul(SENDER_TRACKING_HEADROOM)
.max(MIN_TRACKED_SENDERS),
)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn with_config(
time_source: Arc<dyn TimeSource + Send + Sync>,
min_interval: Duration,
max_tracked_pairs: usize,
) -> Self {
Self::with_class_intervals(time_source, min_interval, min_interval, max_tracked_pairs)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn with_class_intervals(
time_source: Arc<dyn TimeSource + Send + Sync>,
min_interval: Duration,
broadcast_interval: Duration,
max_tracked_pairs: usize,
) -> Self {
Self::with_new_pair_budget(
time_source,
min_interval,
broadcast_interval,
max_tracked_pairs,
NEW_PAIR_BURST,
Ring::DEFAULT_MAX_CONNECTIONS * SENDER_TRACKING_HEADROOM,
)
}
pub fn with_new_pair_budget(
time_source: Arc<dyn TimeSource + Send + Sync>,
min_interval: Duration,
broadcast_interval: Duration,
max_tracked_pairs: usize,
new_pair_burst: f64,
max_tracked_senders: usize,
) -> Self {
Self {
new_pair_budget: TokenBucketLimiter::new(
time_source.clone(),
new_pair_burst,
NEW_PAIR_REFILL_INTERVAL,
max_tracked_senders,
),
new_pair_budget_rejected_total: AtomicU64::new(0),
new_pair_budget_untracked_total: AtomicU64::new(0),
last_new_pair_log: Mutex::new(None),
last_rejected_log: [Mutex::new(None), Mutex::new(None)],
last_accepted: DashMap::new(),
size: AtomicUsize::new(0),
min_interval,
broadcast_interval,
max_tracked_pairs,
time_source,
accepted_total: AtomicU64::new(0),
rejected_total: AtomicU64::new(0),
capacity_rejected_total: AtomicU64::new(0),
capacity_evicted_total: AtomicU64::new(0),
last_eviction_log: Mutex::new(None),
eviction_lock: Mutex::new(()),
}
}
pub fn check_and_record(
&self,
sender: SocketAddr,
contract: ContractInstanceId,
class: UpdateClass,
) -> RateLimitDecision {
let now = self.time_source.now();
let min_interval = self.interval_for(class);
let key = (sender, contract);
use dashmap::mapref::entry::Entry;
let mut budget_spent = false;
let mut reserved = false;
for attempt in 0..MAX_ADMISSION_ATTEMPTS {
match self.last_accepted.entry(key) {
Entry::Occupied(mut entry) => {
if reserved {
self.size.fetch_sub(1, Ordering::Relaxed);
}
let Some(last) = entry.get().get(class) else {
entry.get_mut().set(class, now);
self.accepted_total.fetch_add(1, Ordering::Relaxed);
return RateLimitDecision::Allowed;
};
let elapsed = now.saturating_duration_since(last);
if elapsed < min_interval {
self.rejected_total.fetch_add(1, Ordering::Relaxed);
drop(entry);
self.log_rejected(now, sender, contract, class);
return RateLimitDecision::Rejected {
elapsed,
min_interval,
};
}
entry.get_mut().set(class, now);
self.accepted_total.fetch_add(1, Ordering::Relaxed);
return RateLimitDecision::Allowed;
}
Entry::Vacant(entry) => {
if !reserved {
if !budget_spent {
match self.new_pair_budget.check_and_record_detailed(sender) {
BucketOutcome::Allowed => budget_spent = true,
BucketOutcome::RateLimited => {
drop(entry);
self.new_pair_budget_rejected_total
.fetch_add(1, Ordering::Relaxed);
self.log_new_pair_budget(now, sender);
return RateLimitDecision::SenderNewPairBudget;
}
BucketOutcome::Untracked => {
self.new_pair_budget_untracked_total
.fetch_add(1, Ordering::Relaxed);
budget_spent = true;
}
}
}
let prev = self.size.fetch_add(1, Ordering::Relaxed);
if prev >= self.max_tracked_pairs {
self.size.fetch_sub(1, Ordering::Relaxed);
drop(entry);
if attempt + 1 == MAX_ADMISSION_ATTEMPTS {
break;
}
match self.evict_oldest(now) {
EvictionOutcome::Reserved => {
reserved = true;
continue;
}
EvictionOutcome::Retry => continue,
EvictionOutcome::CapIsZero => {
self.capacity_rejected_total.fetch_add(1, Ordering::Relaxed);
return RateLimitDecision::CapacityExceeded;
}
}
}
}
entry.insert(PairStamps::new_at(class, now));
self.accepted_total.fetch_add(1, Ordering::Relaxed);
return RateLimitDecision::Allowed;
}
}
}
if reserved {
self.size.fetch_sub(1, Ordering::Relaxed);
}
self.capacity_rejected_total.fetch_add(1, Ordering::Relaxed);
RateLimitDecision::CapacityExceeded
}
fn interval_for(&self, class: UpdateClass) -> Duration {
match class {
UpdateClass::Request => self.min_interval,
UpdateClass::Broadcast => self.broadcast_interval,
}
}
fn remove_victims<'a, I>(&self, victims: I) -> usize
where
I: IntoIterator<Item = &'a (SocketAddr, ContractInstanceId)>,
{
victims
.into_iter()
.filter(|victim| self.last_accepted.remove(*victim).is_some())
.count()
}
fn evict_oldest(&self, now: Instant) -> EvictionOutcome {
if self.max_tracked_pairs == 0 {
return EvictionOutcome::CapIsZero;
}
if self.size.load(Ordering::Relaxed) < self.max_tracked_pairs {
return EvictionOutcome::Retry;
}
let evicting = match self.eviction_lock.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
};
if self.size.load(Ordering::Relaxed) < self.max_tracked_pairs {
return EvictionOutcome::Retry;
}
let batch = (self.max_tracked_pairs / EVICTION_BATCH_DIVISOR).max(1);
let mut entries: Vec<(Instant, (SocketAddr, ContractInstanceId))> = self
.last_accepted
.iter()
.map(|e| (e.value().latest(), *e.key()))
.collect();
if entries.is_empty() {
return EvictionOutcome::Retry;
}
let batch = batch.min(entries.len());
entries.select_nth_unstable_by_key(batch - 1, |(stamped, _)| *stamped);
let removed = self.remove_victims(entries[..batch].iter().map(|(_, k)| k));
if removed > 0 {
self.size.fetch_sub(removed - 1, Ordering::Relaxed);
self.capacity_evicted_total
.fetch_add(removed as u64, Ordering::Relaxed);
drop(evicting);
self.log_eviction(now, removed);
return EvictionOutcome::Reserved;
}
EvictionOutcome::Retry
}
fn log_eviction(&self, now: Instant, removed: usize) {
if !log_due(&self.last_eviction_log, now, EVICTION_LOG_INTERVAL) {
return;
}
tracing::info!(
evicted = removed,
evicted_total = self.capacity_evicted_total.load(Ordering::Relaxed),
reserved_plus_tracked = self.size.load(Ordering::Relaxed),
max_tracked_pairs = self.max_tracked_pairs,
"UPDATE rate limiter at capacity: evicted least-recently-used \
(sender, contract) pairs to admit new ones. Expected on a node \
relaying for many peers and contracts; an evicted pair's next \
UPDATE is treated as new. Throttled to one line per minute."
);
}
fn log_new_pair_budget(&self, now: Instant, sender: SocketAddr) {
if !log_due(&self.last_new_pair_log, now, NEW_PAIR_LOG_INTERVAL) {
return;
}
tracing::info!(
%sender,
dropped_total = self.new_pair_budget_rejected_total.load(Ordering::Relaxed),
burst = NEW_PAIR_BURST,
refill_interval_ms = NEW_PAIR_REFILL_INTERVAL.as_millis() as u64,
"UPDATE rate limiter: peer is presenting (sender, contract) pairs this \
node is not tracking faster than its budget allows, so those UPDATEs \
are being dropped. Its traffic for tracked pairs is unaffected. If this \
peer is not churning contract ids, the budget is set too low for this \
node's working set. Throttled to one line per minute."
);
}
fn log_rejected(
&self,
now: Instant,
sender: SocketAddr,
contract: ContractInstanceId,
class: UpdateClass,
) {
if !log_due(
&self.last_rejected_log[class.index()],
now,
REJECTED_LOG_INTERVAL,
) {
return;
}
tracing::info!(
%sender,
%contract,
class = class.as_str(),
rejected_total = self.rejected_total.load(Ordering::Relaxed),
min_interval_ms = self.interval_for(class).as_millis() as u64,
"UPDATE rate limiter: dropping UPDATEs from a (sender, contract) pair \
that is exceeding its class's minimum inter-update interval. \
class=broadcast means co-host fan-out for this contract is \
outrunning the per-pair budget, and each dropped broadcast is \
repaired by a throttled ResyncRequest (#5510). class=request means \
a peer's ROUTED CLIENT WRITES are being refused, which is the \
rarer and more alarming of the two; its originator observes the \
failure and retries. Throttled to one line per minute PER CLASS."
);
}
pub fn cleanup(&self) {
let now = self.time_source.now();
let cutoff = match now.checked_sub(CLEANUP_AGE) {
Some(t) => t,
None => return, };
self.last_accepted.retain(|_, stamps| {
let keep = stamps.latest() >= cutoff;
if !keep {
self.size.fetch_sub(1, Ordering::Relaxed);
}
keep
});
self.new_pair_budget.cleanup();
}
pub fn accepted_total(&self) -> u64 {
self.accepted_total.load(Ordering::Relaxed)
}
pub fn rejected_total(&self) -> u64 {
self.rejected_total.load(Ordering::Relaxed)
}
pub fn capacity_rejected_total(&self) -> u64 {
self.capacity_rejected_total.load(Ordering::Relaxed)
}
pub fn capacity_evicted_total(&self) -> u64 {
self.capacity_evicted_total.load(Ordering::Relaxed)
}
pub fn new_pair_budget_rejected_total(&self) -> u64 {
self.new_pair_budget_rejected_total.load(Ordering::Relaxed)
}
pub fn new_pair_budget_untracked_total(&self) -> u64 {
self.new_pair_budget_untracked_total.load(Ordering::Relaxed)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn tracked_senders(&self) -> usize {
self.new_pair_budget.len()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn len(&self) -> usize {
self.last_accepted.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::util::time_source::SharedMockTimeSource;
fn mk_sender(byte: u8) -> SocketAddr {
SocketAddr::from(([10, 0, 0, byte], 30000 + byte as u16))
}
fn mk_contract(byte: u8) -> ContractInstanceId {
ContractInstanceId::new([byte; 32])
}
fn mk_limiter() -> (UpdateRateLimiter, SharedMockTimeSource) {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::new(Arc::new(ts.clone()), Ring::DEFAULT_MAX_CONNECTIONS);
(limiter, ts)
}
trait Advance {
fn advance(&self, d: Duration);
}
impl Advance for SharedMockTimeSource {
fn advance(&self, d: Duration) {
self.advance_time(d);
}
}
#[test]
fn first_update_for_pair_is_allowed() {
let (l, _ts) = mk_limiter();
let d = l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
assert_eq!(d, RateLimitDecision::Allowed);
assert_eq!(l.accepted_total(), 1);
assert_eq!(l.rejected_total(), 0);
}
#[test]
fn second_update_within_min_interval_is_rejected() {
let (l, ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(10));
let d = l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
assert!(
matches!(d, RateLimitDecision::Rejected { .. }),
"second UPDATE 10ms after first must be rejected, got {d:?}"
);
assert_eq!(l.accepted_total(), 1);
assert_eq!(l.rejected_total(), 1);
}
#[test]
fn broadcast_class_is_charged_against_its_own_interval() {
let (l, ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Broadcast)
.is_allowed()
);
assert!(
l.check_and_record(mk_sender(2), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(30));
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Broadcast)
.is_allowed(),
"a broadcast 30ms after the last one must be ALLOWED — co-host fan-out \
carries the contract's AGGREGATE commit rate, not one writer's cadence \
(#5510)"
);
let d = l.check_and_record(mk_sender(2), mk_contract(1), UpdateClass::Request);
assert!(
matches!(d, RateLimitDecision::Rejected { .. }),
"a routed request 30ms after the last one must still be REJECTED — the \
write-cadence budget is unchanged by #5510, got {d:?}"
);
}
#[test]
fn broadcast_traffic_does_not_starve_routed_requests_on_the_same_pair() {
let (l, ts) = mk_limiter();
let (sender, contract) = (mk_sender(1), mk_contract(1));
assert!(
l.check_and_record(sender, contract, UpdateClass::Request)
.is_allowed(),
"precondition: the pair's first request is allowed and stamps the \
request class"
);
for i in 0..8 {
ts.advance(Duration::from_millis(25));
assert!(
l.check_and_record(sender, contract, UpdateClass::Broadcast)
.is_allowed(),
"broadcast {i} at 25ms spacing must be allowed"
);
}
assert!(
l.check_and_record(sender, contract, UpdateClass::Request)
.is_allowed(),
"a routed request must be judged against the REQUEST stamp alone — \
sharing one stamp with fan-out would starve client writes on every \
busy contract (#5510)"
);
}
#[test]
fn first_message_of_a_class_on_an_existing_pair_is_allowed() {
let (l, _ts) = mk_limiter();
let (sender, contract) = (mk_sender(1), mk_contract(1));
assert!(
l.check_and_record(sender, contract, UpdateClass::Broadcast)
.is_allowed(),
"precondition: the broadcast creates the entry"
);
assert!(
l.check_and_record(sender, contract, UpdateClass::Request)
.is_allowed(),
"the pair's FIRST request must be allowed even with zero elapsed time — \
it has no request stamp to be measured against"
);
assert!(
!l.check_and_record(sender, contract, UpdateClass::Request)
.is_allowed(),
"the second request must be refused — the first must have stamped"
);
}
#[test]
fn capacity_eviction_orders_on_the_most_recent_stamp_of_either_class() {
let ts = SharedMockTimeSource::new();
let l = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 2);
let (busy, quiet, newcomer) = (mk_contract(1), mk_contract(2), mk_contract(3));
let sender = mk_sender(1);
assert!(
l.check_and_record(sender, busy, UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_secs(10));
assert!(
l.check_and_record(sender, quiet, UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_secs(10));
assert!(
l.check_and_record(sender, busy, UpdateClass::Broadcast)
.is_allowed()
);
ts.advance(Duration::from_secs(10));
assert!(
l.check_and_record(sender, newcomer, UpdateClass::Request)
.is_allowed()
);
assert_eq!(l.len(), 2, "the cap must still hold after the admission");
assert!(
l.last_accepted.contains_key(&(sender, busy)),
"the pair kept live by BROADCAST traffic must have SURVIVED — ordering \
eviction on the request stamp alone would evict it, resetting its \
broadcast window and handing the sender a free message (#5510)"
);
assert!(
!l.last_accepted.contains_key(&(sender, quiet)),
"the genuinely least-recently-used pair must be the one evicted"
);
}
#[test]
fn ttl_sweep_reads_the_most_recent_stamp_of_either_class() {
let (l, ts) = mk_limiter();
let (sender, contract) = (mk_sender(1), mk_contract(1));
assert!(
l.check_and_record(sender, contract, UpdateClass::Request)
.is_allowed()
);
ts.advance(CLEANUP_AGE + Duration::from_secs(1));
assert!(
l.check_and_record(sender, contract, UpdateClass::Broadcast)
.is_allowed()
);
l.cleanup();
assert_eq!(
l.len(),
1,
"an entry whose broadcast stamp is fresh must survive the sweep even \
when its request stamp is older than CLEANUP_AGE"
);
ts.advance(CLEANUP_AGE + Duration::from_secs(1));
l.cleanup();
assert_eq!(l.len(), 0, "an entry stale in both classes must be swept");
}
#[test]
fn every_broadcast_opcode_is_charged_to_the_broadcast_class() {
const NODE_SRC: &str = include_str!("../node.rs");
let start = NODE_SRC.find("let update_class = match op {").expect(
"the UPDATE dispatch must classify the message before the rate-limit \
check; if the classifier moved or was renamed, update this pin rather \
than deleting it",
);
let body = &NODE_SRC[start..];
let end = body
.find("\n };")
.expect("could not find the end of the update_class match");
let body: String = body[..end].chars().filter(|c| !c.is_whitespace()).collect();
let class_at = |opcode: &str| -> &'static str {
let at = body.find(opcode).unwrap_or_else(|| {
panic!(
"opcode `{opcode}` is not classified in node.rs — an \
unclassified UPDATE opcode is a budget bypass"
)
});
let next_request = body[at..].find("UpdateClass::Request");
let next_broadcast = body[at..].find("UpdateClass::Broadcast");
match (next_request, next_broadcast) {
(Some(r), Some(b)) => {
if r < b {
"Request"
} else {
"Broadcast"
}
}
(Some(_), None) => "Request",
(None, Some(_)) => "Broadcast",
(None, None) => panic!(
"opcode `{opcode}` is followed by no UpdateClass at all — the \
classifier is not a match over the two classes any more"
),
}
};
for opcode in [
"UpdateMsg::BroadcastTo{..}",
"UpdateMsg::BroadcastToV2{..}",
"UpdateMsg::BroadcastToStreaming{..}",
"UpdateMsg::BroadcastToStreamingV2{..}",
] {
assert_eq!(
class_at(opcode),
"Broadcast",
"broadcast opcode `{opcode}` must be charged to the Broadcast \
class — the four broadcast opcodes MUST share one budget so \
switching between them gains nothing (#5510)"
);
}
for opcode in [
"UpdateMsg::RequestUpdate{..}",
"UpdateMsg::RequestUpdateStreaming{..}",
] {
assert_eq!(
class_at(opcode),
"Request",
"request opcode `{opcode}` must be charged to the Request class — \
a routed client write must not gain the wider fan-out budget \
(#5510)"
);
}
assert!(
!body.contains("_=>"),
"the classifier must stay an exhaustive match with NO catch-all: a new \
UPDATE wire variant must fail to COMPILE rather than silently inherit \
whichever class the wildcard names"
);
let variants = body.matches("UpdateMsg::").count();
assert_eq!(
variants, 6,
"the classifier must name all six UpdateMsg variants explicitly, found \
{variants}. Fewer means a variant is being matched by something other \
than its own name — which is a catch-all however it is spelled"
);
let arms = body.matches("=>").count();
assert_eq!(
arms, 2,
"expected exactly two arms (Request and Broadcast), found {arms}; a \
third arm is a catch-all or an unintended reclassification"
);
for (i, _) in body.match_indices("=>") {
let preceding = &body[..i];
assert!(
preceding.ends_with("}"),
"every classifier arm must match a named variant with a `{{..}}` \
pattern, but one arm's pattern ends with {:?} — a bare identifier \
binds EVERY variant (`other => ...`) and a guard leaves the \
remainder unmatched, either of which lets a new wire variant \
inherit a class silently",
preceding.chars().rev().take(12).collect::<String>()
);
}
}
#[test]
fn update_after_min_interval_is_allowed() {
let (l, ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(200));
let d = l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
assert_eq!(d, RateLimitDecision::Allowed);
assert_eq!(l.accepted_total(), 2);
assert_eq!(l.rejected_total(), 0);
}
#[test]
fn different_senders_same_contract_independent() {
let (l, ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(1));
assert!(
l.check_and_record(mk_sender(2), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
let d = l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
assert!(matches!(d, RateLimitDecision::Rejected { .. }));
}
#[test]
fn same_sender_different_contracts_independent() {
let (l, _ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
assert!(
l.check_and_record(mk_sender(1), mk_contract(2), UpdateClass::Request)
.is_allowed()
);
assert_eq!(l.accepted_total(), 2);
}
#[test]
fn rejected_attempts_do_not_extend_window() {
let (l, ts) = mk_limiter();
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
for _ in 0..9 {
ts.advance(Duration::from_millis(10));
assert!(
!l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
}
ts.advance(Duration::from_millis(5));
assert!(
!l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(10));
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed(),
"after 105ms+ from original accept, next attempt MUST be allowed — \
rejected attempts must not have moved the window forward"
);
}
#[test]
fn cleanup_removes_stale_entries() {
let (l, ts) = mk_limiter();
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
l.check_and_record(mk_sender(2), mk_contract(2), UpdateClass::Request);
assert_eq!(l.len(), 2);
ts.advance(CLEANUP_AGE + Duration::from_secs(1));
l.cleanup();
assert_eq!(l.len(), 0, "all stale entries must be cleared");
}
#[test]
fn cleanup_preserves_fresh_entries() {
let (l, ts) = mk_limiter();
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request);
ts.advance(CLEANUP_AGE / 2);
l.cleanup();
assert_eq!(l.len(), 1, "fresh entry must be preserved");
}
#[test]
fn counters_track_accepts_and_rejects() {
let (l, ts) = mk_limiter();
for i in 0..5 {
ts.advance(MIN_UPDATE_INTERVAL + Duration::from_millis(1));
assert!(
l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed(),
"iter {i}"
);
}
for _ in 0..3 {
assert!(
!l.check_and_record(mk_sender(1), mk_contract(1), UpdateClass::Request)
.is_allowed()
);
}
assert_eq!(l.accepted_total(), 5);
assert_eq!(l.rejected_total(), 3);
}
#[test]
fn may21_flood_pattern_is_throttled() {
let (l, ts) = mk_limiter();
let sender = mk_sender(1);
let contract = mk_contract(1);
for _ in 0..1000 {
l.check_and_record(sender, contract, UpdateClass::Request);
ts.advance(Duration::from_millis(1));
}
let accepted = l.accepted_total();
let rejected = l.rejected_total();
assert!(
(9..=12).contains(&accepted),
"expected ~10 admits over 1s of flooding, got {accepted}"
);
assert_eq!(accepted + rejected, 1000);
assert!(
rejected as f64 / 1000.0 > 0.95,
"expected >95% rejection rate, got {}",
rejected as f64 / 1000.0
);
}
#[test]
fn at_capacity_new_pair_is_admitted_and_the_map_stays_bounded() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
8, );
for i in 0..8 {
let d = limiter.check_and_record(
mk_sender(i + 1),
mk_contract(i + 1),
UpdateClass::Request,
);
assert_eq!(d, RateLimitDecision::Allowed, "pair {i} should be allowed");
ts.advance(Duration::from_millis(1));
}
assert_eq!(limiter.len(), 8);
let d = limiter.check_and_record(mk_sender(99), mk_contract(99), UpdateClass::Request);
assert_eq!(
d,
RateLimitDecision::Allowed,
"a new pair at capacity must be admitted, not starved (#4981)"
);
assert_eq!(
limiter.capacity_rejected_total(),
0,
"admission by eviction must not count as a capacity rejection"
);
assert_eq!(limiter.capacity_evicted_total(), 1);
assert_eq!(
limiter.len(),
8,
"the cap is still a hard bound: one in, one out"
);
ts.advance(MIN_UPDATE_INTERVAL + Duration::from_millis(1));
let d = limiter.check_and_record(mk_sender(8), mk_contract(8), UpdateClass::Request);
assert_eq!(
d,
RateLimitDecision::Allowed,
"existing pair must keep working at cap"
);
}
#[test]
fn eviction_counts_actual_removals_not_the_selected_batch() {
const CAP: usize = 8;
let ts = SharedMockTimeSource::new();
let limiter =
UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, CAP);
for i in 0..CAP {
assert_eq!(
limiter.check_and_record(
mk_sender(i as u8 + 1),
mk_contract(i as u8 + 1),
UpdateClass::Request
),
RateLimitDecision::Allowed
);
ts.advance(Duration::from_millis(1));
}
assert_eq!(limiter.len(), CAP);
assert_eq!(limiter.size.load(Ordering::Relaxed), CAP);
let present_a = (mk_sender(1), mk_contract(1));
let present_b = (mk_sender(2), mk_contract(2));
let already_gone = (mk_sender(3), mk_contract(3));
assert!(limiter.last_accepted.remove(&already_gone).is_some());
let victims = [present_a, already_gone, present_b];
let removed = limiter.remove_victims(victims.iter());
assert_eq!(
removed,
2,
"two of the three victims were still present, so the pass removed \
two; returning the selected batch size ({}) instead over-decrements \
`size` and lets the map grow past the cap",
victims.len()
);
assert_eq!(
limiter.len(),
CAP - 3,
"all three victims are gone from the map either way"
);
limiter.size.fetch_sub(removed - 1, Ordering::Relaxed);
assert!(
limiter.size.load(Ordering::Relaxed) >= limiter.len(),
"`size` must never fall below the map's true length: size={} len={}",
limiter.size.load(Ordering::Relaxed),
limiter.len()
);
}
#[test]
fn a_throttled_sender_cannot_evict_other_peers_entries() {
const CAP: usize = 8;
const BURST: usize = 8;
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
CAP,
BURST as f64,
Ring::DEFAULT_MAX_CONNECTIONS * SENDER_TRACKING_HEADROOM,
);
let incumbent = mk_sender(200);
for i in 0..CAP {
assert_eq!(
limiter.check_and_record(incumbent, mk_contract(i as u8 + 1), UpdateClass::Request),
RateLimitDecision::Allowed
);
}
assert_eq!(limiter.len(), CAP);
let evicted_before = limiter.capacity_evicted_total();
let attacker = mk_sender(1);
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(
attacker,
mk_contract(100 + i as u8),
UpdateClass::Request
),
RateLimitDecision::Allowed
);
}
let evicted_at_budget_end = limiter.capacity_evicted_total();
assert!(
evicted_at_budget_end > evicted_before,
"fixture check: at this cap an admission must actually evict, \
otherwise the assertion below cannot discriminate"
);
for i in BURST..(BURST + 32) {
assert_eq!(
limiter.check_and_record(
attacker,
mk_contract(100 + i as u8),
UpdateClass::Request
),
RateLimitDecision::SenderNewPairBudget
);
}
assert_eq!(
limiter.capacity_evicted_total(),
evicted_at_budget_end,
"a sender past its budget must not evict anything: the budget is \
charged before the reserve/evict block, so being throttled costs \
other peers nothing. 32 refused messages evicted {} entries.",
limiter.capacity_evicted_total() - evicted_at_budget_end
);
assert_eq!(limiter.len(), CAP, "the map stays exactly at the cap");
}
#[test]
fn log_due_fires_once_per_interval() {
const INTERVAL: Duration = Duration::from_secs(60);
let slot = Mutex::new(None);
let t0 = Instant::now();
assert!(
log_due(&slot, t0, INTERVAL),
"the first call must be due — an unfired throttle that starts \
closed would suppress the signal entirely"
);
assert!(
!log_due(&slot, t0, INTERVAL),
"an immediate second call must be throttled"
);
assert!(
!log_due(&slot, t0 + INTERVAL - Duration::from_millis(1), INTERVAL),
"just inside the interval is still throttled"
);
assert!(
log_due(&slot, t0 + INTERVAL, INTERVAL),
"at the interval the line is due again"
);
}
#[test]
fn a_sender_churning_fresh_contract_ids_is_cut_off_after_its_burst() {
const BURST: usize = 8;
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
MAX_TRACKED_PAIRS,
BURST as f64,
Ring::DEFAULT_MAX_CONNECTIONS * SENDER_TRACKING_HEADROOM,
);
let attacker = mk_sender(1);
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(attacker, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::Allowed,
"fresh id {i} is within the burst"
);
}
for i in BURST..(BURST + 32) {
assert_eq!(
limiter.check_and_record(attacker, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::SenderNewPairBudget,
"fresh id {i} is past the burst and must be refused"
);
}
assert_eq!(limiter.new_pair_budget_rejected_total(), 32);
assert_eq!(
limiter.len(),
BURST,
"a throttled sender must not have grown the tracking map"
);
assert_eq!(
limiter.capacity_evicted_total(),
0,
"a throttled sender must not have evicted anything — the budget is \
charged BEFORE the capacity path so churn cannot push other peers out"
);
ts.advance(NEW_PAIR_REFILL_INTERVAL);
assert_eq!(
limiter.check_and_record(attacker, mk_contract(200), UpdateClass::Request),
RateLimitDecision::Allowed,
"one refill interval must buy exactly one more fresh pair"
);
assert_eq!(
limiter.check_and_record(attacker, mk_contract(201), UpdateClass::Request),
RateLimitDecision::SenderNewPairBudget,
"and only one"
);
}
#[test]
fn the_new_pair_budget_never_throttles_established_pairs() {
const BURST: usize = 4;
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
MAX_TRACKED_PAIRS,
BURST as f64,
Ring::DEFAULT_MAX_CONNECTIONS * SENDER_TRACKING_HEADROOM,
);
let peer = mk_sender(1);
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(peer, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::Allowed
);
}
assert_eq!(
limiter.check_and_record(peer, mk_contract(99), UpdateClass::Request),
RateLimitDecision::SenderNewPairBudget
);
for round in 0..200 {
ts.advance(MIN_UPDATE_INTERVAL + Duration::from_millis(1));
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(peer, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::Allowed,
"round {round}, established pair {i} must be unaffected by the \
new-pair budget"
);
}
}
assert_eq!(
limiter.new_pair_budget_rejected_total(),
1,
"only the one fresh pair was refused; established traffic never \
touches the budget"
);
assert_eq!(limiter.accepted_total(), (BURST + BURST * 200) as u64);
}
#[test]
fn under_saturation_re_admitted_pairs_are_not_throttled_at_the_sustained_rate() {
const CAP: usize = 64;
const CONTRACTS: u8 = 128;
const CYCLES: usize = 20;
let ts = SharedMockTimeSource::new();
let limiter =
UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, CAP);
let peer = mk_sender(1);
for cycle in 0..CYCLES {
for c in 0..CONTRACTS {
assert_eq!(
limiter.check_and_record(peer, mk_contract(c), UpdateClass::Request),
RateLimitDecision::Allowed,
"cycle {cycle}, contract {c}: a peer relaying at the sustained rate \
must not be throttled, even though eviction keeps making its pairs \
look new"
);
ts.advance(NEW_PAIR_REFILL_INTERVAL);
}
}
assert!(
limiter.capacity_evicted_total() > 0,
"fixture must actually saturate, or it proves nothing about re-admission"
);
assert_eq!(
limiter.new_pair_budget_rejected_total(),
0,
"nothing may be dropped at the sustained rate"
);
let attempts = NEW_PAIR_BURST as usize * 2;
let mut refused = 0;
for c in 0..attempts {
let mut id = [0u8; 32];
id[..8].copy_from_slice(&(c as u64).to_be_bytes());
id[8] = 0xEE; if limiter.check_and_record(peer, ContractInstanceId::new(id), UpdateClass::Request)
== RateLimitDecision::SenderNewPairBudget
{
refused += 1;
}
}
assert!(
refused > 0,
"a peer presenting unfamiliar pairs with no time passing must eventually \
be refused, or the budget bounds nothing"
);
assert!(
refused < attempts,
"and the burst must let a substantial run through before that: \
{refused} of {attempts} refused"
);
}
#[test]
fn the_new_pair_budget_is_scoped_per_sender() {
const BURST: usize = 2;
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
MAX_TRACKED_PAIRS,
BURST as f64,
Ring::DEFAULT_MAX_CONNECTIONS * SENDER_TRACKING_HEADROOM,
);
let noisy = mk_sender(1);
let quiet = mk_sender(2);
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(noisy, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::Allowed
);
}
assert_eq!(
limiter.check_and_record(noisy, mk_contract(50), UpdateClass::Request),
RateLimitDecision::SenderNewPairBudget
);
for i in 0..BURST {
assert_eq!(
limiter.check_and_record(quiet, mk_contract(i as u8), UpdateClass::Request),
RateLimitDecision::Allowed,
"a second sender has its own budget, untouched by the first"
);
}
assert_eq!(limiter.tracked_senders(), 2);
}
#[test]
fn the_new_pair_budget_map_is_bounded_and_fails_open_when_full() {
const MAX_SENDERS: usize = 4;
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
MAX_TRACKED_PAIRS,
NEW_PAIR_BURST,
MAX_SENDERS,
);
for i in 0..(MAX_SENDERS * 4) {
let d =
limiter.check_and_record(mk_sender(i as u8), mk_contract(1), UpdateClass::Request);
assert_eq!(
d,
RateLimitDecision::Allowed,
"sender {i}: a full sender map must fail OPEN — refusing here is \
#4981 one map over, and this sender has spent nothing"
);
assert!(
limiter.tracked_senders() <= MAX_SENDERS,
"sender {i}: the budget map must never exceed its cap"
);
}
assert_eq!(limiter.tracked_senders(), MAX_SENDERS);
assert_eq!(
limiter.new_pair_budget_untracked_total(),
(MAX_SENDERS * 3) as u64,
"every admission past the cap must be counted as untracked"
);
assert_eq!(
limiter.new_pair_budget_rejected_total(),
0,
"a full map is a sizing accident, not evidence about any sender"
);
}
#[test]
fn evict_oldest_has_exactly_one_terminal_exit() {
const THIS_FILE: &str = include_str!("update_rate_limit.rs");
let start = THIS_FILE
.find(concat!("fn evict_", "oldest("))
.expect("evict_oldest must exist; if renamed, update this pin");
let tail = &THIS_FILE[start..];
let end = tail[1..]
.find("\n fn ")
.map(|i| i + 1)
.unwrap_or(tail.len());
let body = &tail[..end];
let terminal = concat!("EvictionOutcome::", "CapIsZero");
assert_eq!(
body.matches(terminal).count(),
1,
"evict_oldest must have exactly ONE terminal exit. Every other way out \
means slots are, or are about to be, free — including an empty map, \
which `cleanup` can produce mid-scan — and dropping there loses a \
legitimate UPDATE without spending the retry budget (#4981, #4997)."
);
let guard_pos = body
.find("self.max_tracked_pairs == 0")
.expect("the terminal exit must be guarded by the cap being zero");
let terminal_pos = body.find(terminal).expect("checked above");
assert!(
guard_pos < terminal_pos,
"the one terminal exit must be the zero-cap guard, not a condition \
inferred from the map's contents"
);
}
#[test]
fn an_empty_map_at_a_nonzero_cap_retries_rather_than_dropping() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 8);
limiter.size.store(8, Ordering::Relaxed);
assert_eq!(limiter.len(), 0);
assert!(
matches!(limiter.evict_oldest(ts.now()), EvictionOutcome::Retry),
"an empty map at a non-zero cap means slots are about to exist, not that \
the cap is zero"
);
let zero = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 0);
assert!(matches!(
zero.evict_oldest(ts.now()),
EvictionOutcome::CapIsZero
));
}
#[test]
fn cleanup_leaves_the_size_counter_equal_to_the_map_length() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 8);
for i in 1..=8u8 {
assert!(
limiter
.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request)
.is_allowed()
);
}
ts.advance(CLEANUP_AGE + Duration::from_secs(1));
for i in 1..=4u8 {
assert!(
limiter
.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request)
.is_allowed()
);
}
limiter.cleanup();
assert_eq!(limiter.len(), 4);
assert_eq!(
limiter.size.load(Ordering::Relaxed),
limiter.len(),
"size must equal the map length after a sweep"
);
}
#[test]
fn a_zero_cap_refuses_every_pair_without_spinning() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 0);
for i in 1..=4u8 {
assert_eq!(
limiter.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request),
RateLimitDecision::CapacityExceeded,
"with a zero cap there is no slot for pair {i} and nothing to evict"
);
}
assert_eq!(limiter.len(), 0);
assert_eq!(limiter.capacity_rejected_total(), 4);
assert_eq!(
limiter.capacity_evicted_total(),
0,
"an empty map has nothing to evict"
);
assert_eq!(limiter.accepted_total(), 0);
}
const ORDER_TEST_CAP: usize = 128;
const ORDER_TEST_EVICTED: usize = 6;
fn mk_pair(i: usize) -> (SocketAddr, ContractInstanceId) {
let sender = SocketAddr::from(([10, 1, (i >> 8) as u8, (i & 0xff) as u8], 30000));
(sender, ContractInstanceId::new([0xAB; 32]))
}
fn admit_until_evicted(
limiter: &UpdateRateLimiter,
ts: &SharedMockTimeSource,
from: usize,
) -> usize {
let mut i = from;
while limiter.capacity_evicted_total() < ORDER_TEST_EVICTED as u64 {
let (s, c) = mk_pair(i);
assert_eq!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Allowed,
"newcomer {i} must be admitted"
);
i += 1;
ts.advance(Duration::from_micros(1));
assert!(
i < from + ORDER_TEST_CAP,
"newcomers should have driven {ORDER_TEST_EVICTED} evictions long before this"
);
}
assert_eq!(
limiter.capacity_evicted_total(),
ORDER_TEST_EVICTED as u64,
"batches of 2 must land exactly on {ORDER_TEST_EVICTED}"
);
i
}
#[test]
fn at_capacity_evicts_the_oldest_pairs_not_arbitrary_ones() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
ORDER_TEST_CAP,
);
assert_eq!(
ORDER_TEST_CAP / EVICTION_BATCH_DIVISOR,
2,
"fixture assumes a 2-entry batch"
);
let start = ts.now();
for i in 0..ORDER_TEST_CAP {
let (s, c) = mk_pair(i);
assert_eq!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Allowed,
"fill {i}"
);
ts.advance(Duration::from_micros(1));
}
assert_eq!(limiter.len(), ORDER_TEST_CAP);
admit_until_evicted(&limiter, &ts, ORDER_TEST_CAP);
assert!(
ts.now().saturating_duration_since(start) < MIN_UPDATE_INTERVAL,
"the fixture must stay inside min_interval, or a surviving pair and an \
evicted one both answer Allowed and the assertions below are vacuous"
);
for i in ORDER_TEST_EVICTED..ORDER_TEST_CAP {
let (s, c) = mk_pair(i);
assert!(
matches!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Rejected { .. }
),
"pair {i} is newer than the {ORDER_TEST_EVICTED} oldest and must have survived"
);
}
for i in 0..ORDER_TEST_EVICTED {
let (s, c) = mk_pair(i);
assert_eq!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Allowed,
"pair {i} is among the {ORDER_TEST_EVICTED} oldest and must have been evicted"
);
}
}
#[test]
fn a_refreshed_pair_moves_to_the_back_of_the_eviction_order() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
ORDER_TEST_CAP,
);
for i in 0..ORDER_TEST_CAP {
let (s, c) = mk_pair(i);
assert_eq!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Allowed,
"fill {i}"
);
ts.advance(Duration::from_micros(1));
}
ts.advance(MIN_UPDATE_INTERVAL);
let victims = ORDER_TEST_EVICTED..(2 * ORDER_TEST_EVICTED);
let refreshed = ts.now();
for i in (0..ORDER_TEST_CAP).filter(|i| !victims.contains(i)) {
let (s, c) = mk_pair(i);
assert_eq!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Allowed,
"refresh {i} must be accepted a full min_interval after the fill"
);
}
assert_eq!(
limiter.capacity_evicted_total(),
0,
"restamping tracked pairs must not evict anything — it never takes the \
capacity path"
);
admit_until_evicted(&limiter, &ts, ORDER_TEST_CAP);
assert!(
ts.now().saturating_duration_since(refreshed) < MIN_UPDATE_INTERVAL,
"every refreshed pair must still be inside min_interval, or the \
assertions below are vacuous"
);
for i in (0..ORDER_TEST_CAP).filter(|i| !victims.contains(i)) {
let (s, c) = mk_pair(i);
assert!(
matches!(
limiter.check_and_record(s, c, UpdateClass::Request),
RateLimitDecision::Rejected { .. }
),
"pair {i} was refreshed, so it must be at the back of the eviction \
order and must have survived"
);
}
}
#[test]
fn busy_pairs_cannot_hold_slots_against_newcomers() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, 8);
let start = ts.now();
for i in 1..=8u8 {
assert!(
limiter
.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(1));
}
for round in 0..20u8 {
ts.advance(MIN_UPDATE_INTERVAL + Duration::from_millis(1));
for i in 1..=8u8 {
limiter.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request);
}
let newcomer = limiter.check_and_record(
mk_sender(100 + round),
mk_contract(200),
UpdateClass::Request,
);
assert_eq!(
newcomer,
RateLimitDecision::Allowed,
"round {round}: a newcomer must not be starved by busy incumbents"
);
assert!(
limiter.len() <= 8,
"round {round}: the cap must still bound the map"
);
}
assert_eq!(
limiter.capacity_rejected_total(),
0,
"no UPDATE should have been dropped for capacity"
);
assert!(
ts.now().saturating_duration_since(start) < CLEANUP_AGE,
"sanity: the fixture must stay inside CLEANUP_AGE, so the TTL sweep \
cannot be what admitted the newcomers"
);
}
#[test]
fn eviction_is_batched_at_larger_caps() {
let ts = SharedMockTimeSource::new();
let cap = 128;
let limiter =
UpdateRateLimiter::with_config(Arc::new(ts.clone()), MIN_UPDATE_INTERVAL, cap);
let expected_batch = cap / EVICTION_BATCH_DIVISOR;
assert_eq!(expected_batch, 2, "fixture assumes a 2-entry batch");
for i in 0..cap {
let sender = SocketAddr::from(([10, 1, (i / 256) as u8, (i % 256) as u8], 30000));
assert!(
limiter
.check_and_record(sender, mk_contract(1), UpdateClass::Request)
.is_allowed()
);
ts.advance(Duration::from_millis(1));
}
assert_eq!(limiter.len(), cap);
assert!(
limiter
.check_and_record(mk_sender(99), mk_contract(99), UpdateClass::Request)
.is_allowed()
);
assert_eq!(
limiter.capacity_evicted_total(),
expected_batch as u64,
"one admission at capacity must evict a whole batch"
);
assert_eq!(
limiter.len(),
cap - expected_batch + 1,
"the batch leaves headroom, so the next admissions skip the scan"
);
assert!(
limiter
.check_and_record(mk_sender(98), mk_contract(98), UpdateClass::Request)
.is_allowed()
);
assert_eq!(
limiter.capacity_evicted_total(),
expected_batch as u64,
"an admission with headroom must not trigger another eviction"
);
}
#[test]
fn concurrent_check_and_record_admits_one_per_window() {
use std::sync::{Arc as StdArc, Barrier};
use std::thread;
let ts = SharedMockTimeSource::new();
let limiter = StdArc::new(UpdateRateLimiter::new(
Arc::new(ts.clone()),
Ring::DEFAULT_MAX_CONNECTIONS,
));
let sender = mk_sender(1);
let contract = mk_contract(1);
const THREADS: usize = 16;
let barrier = StdArc::new(Barrier::new(THREADS));
let mut handles = Vec::with_capacity(THREADS);
for _ in 0..THREADS {
let l = limiter.clone();
let b = barrier.clone();
handles.push(thread::spawn(move || {
b.wait();
l.check_and_record(sender, contract, UpdateClass::Request)
}));
}
let mut allowed = 0;
let mut rejected = 0;
for h in handles {
match h.join().unwrap() {
RateLimitDecision::Allowed => allowed += 1,
RateLimitDecision::Rejected { .. } => rejected += 1,
RateLimitDecision::CapacityExceeded => panic!("unexpected cap"),
RateLimitDecision::SenderNewPairBudget => {
panic!("one pair cannot exhaust a new-pair budget")
}
}
}
assert_eq!(
allowed, 1,
"exactly ONE concurrent caller must be admitted per window; \
got {allowed} admits, {rejected} rejects"
);
assert_eq!(rejected, THREADS - 1);
assert_eq!(limiter.accepted_total(), 1);
assert_eq!(limiter.rejected_total(), (THREADS - 1) as u64);
}
#[test]
fn update_dispatch_gates_all_four_wire_variants() {
const NODE_SRC: &str = include_str!("../node.rs");
let block_start = NODE_SRC
.find("NetMessageV1::Update(ref op) =>")
.expect("could not locate UPDATE dispatch block in node.rs");
let tail = &NODE_SRC[block_start + 1..];
let block_len = tail
.find("\n NetMessageV1::")
.or_else(|| tail.find("\n NetMessageV1::"))
.unwrap_or(tail.len());
let block = &NODE_SRC[block_start..block_start + 1 + block_len];
let rate_limit_pos = block
.find("update_rate_limiter")
.expect("update_rate_limiter not invoked in UPDATE dispatch block");
let first_spawn_pos = block
.find("start_relay_request_update(")
.expect("start_relay_request_update spawn not found in block");
assert!(
rate_limit_pos < first_spawn_pos,
"rate limit gate (offset {rate_limit_pos}) must appear BEFORE \
the first relay spawn (offset {first_spawn_pos}) so rejected \
messages don't pay the spawn cost"
);
for variant in [
"UpdateMsg::RequestUpdate {",
"UpdateMsg::BroadcastTo {",
"UpdateMsg::RequestUpdateStreaming {",
"UpdateMsg::BroadcastToStreaming {",
"UpdateMsg::BroadcastToV2 {",
"UpdateMsg::BroadcastToStreamingV2 {",
] {
assert!(
block.contains(variant),
"UPDATE dispatch block missing wire variant: `{variant}`. \
If a new UPDATE wire variant was added, gate it through \
the rate limiter and update this list. If a variant was \
removed, update this list."
);
}
for spawn in [
"start_relay_request_update(",
"start_relay_broadcast_to(",
"start_relay_request_update_streaming(",
"start_relay_broadcast_to_streaming(",
] {
let count = block.matches(spawn).count();
assert!(
count >= 1,
"UPDATE dispatch block does not invoke `{spawn}` — the \
corresponding wire variant is not actually gated."
);
}
}
#[test]
fn concurrent_distinct_keys_do_not_overshoot_cap() {
use std::sync::{Arc as StdArc, Barrier};
use std::thread;
const CAP: usize = 8;
const THREADS: usize = 64;
let ts = SharedMockTimeSource::new();
let limiter = StdArc::new(UpdateRateLimiter::with_config(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
CAP,
));
let barrier = StdArc::new(Barrier::new(THREADS));
let mut handles = Vec::with_capacity(THREADS);
for i in 0..THREADS {
let l = limiter.clone();
let b = barrier.clone();
handles.push(thread::spawn(move || {
b.wait();
l.check_and_record(
mk_sender((i + 1) as u8),
mk_contract((i + 1) as u8),
UpdateClass::Request,
)
}));
}
let mut allowed = 0;
let mut cap_rejected = 0;
let mut rate_rejected = 0;
for h in handles {
match h.join().unwrap() {
RateLimitDecision::Allowed => allowed += 1,
RateLimitDecision::CapacityExceeded => cap_rejected += 1,
RateLimitDecision::Rejected { .. } => rate_rejected += 1,
RateLimitDecision::SenderNewPairBudget => {
panic!("each thread uses a distinct sender, so no budget can be spent")
}
}
}
assert!(
limiter.len() <= CAP,
"strict cap: map size must never exceed CAP after a 64-thread \
concurrent flood of distinct keys, got {}",
limiter.len()
);
assert!(
allowed >= CAP,
"at least CAP admissions expected under flood, got {allowed}"
);
assert_eq!(allowed + cap_rejected + rate_rejected, THREADS);
assert_eq!(rate_rejected, 0);
assert_eq!(limiter.capacity_rejected_total(), cap_rejected as u64);
assert_eq!(
limiter.len() as u64,
allowed as u64 - limiter.capacity_evicted_total(),
"inserts minus evictions must account for every tracked entry: \
len={} allowed={allowed} evicted={}",
limiter.len(),
limiter.capacity_evicted_total()
);
assert!(
limiter.capacity_evicted_total() > 0,
"64 threads against a cap of {CAP} must have driven at least one \
eviction; zero means the at-capacity path refused instead of \
evicting, which is the #4981 regression"
);
assert_eq!(
limiter.size.load(Ordering::Relaxed),
limiter.len(),
"after all callers finish, the reservation counter must have \
settled back onto the map's true length"
);
}
#[test]
fn capacity_signals_are_logged_above_debug_level() {
fn strip_ws(s: &str) -> String {
s.chars()
.filter(|c| !c.is_whitespace() && *c != '\\')
.collect()
}
fn fn_body<'a>(src: &'a str, sig: &str, what: &str) -> &'a str {
let start = src.find(sig).unwrap_or_else(|| {
panic!(
"the {what} log site's enclosing fn (`{sig}`) must exist; \
if it was renamed, update this pin rather than deleting it"
)
});
let rest = &src[start + sig.len()..];
let end = rest.find("\n fn ").unwrap_or(rest.len());
&rest[..end]
}
fn level_of(haystack: &str, marker: &str, what: &str) -> String {
let pos = haystack.find(marker).unwrap_or_else(|| {
panic!("the {what} log line must exist; if it was renamed, update this pin rather than deleting it")
});
let start = haystack[..pos]
.rfind(&strip_ws("tracing::"))
.unwrap_or_else(|| {
panic!(
"no `tracing::<level>!` precedes the {what} log line inside its own \
function. If the site now uses an imported `info!`/`debug!` without \
the `tracing::` path, this pin can no longer read its level — restore \
the qualified form rather than widening the search, because widening \
it is exactly what lets a silent downgrade pass."
)
});
let tail = &haystack[start..pos];
tail[..tail.find('!').unwrap_or(tail.len())].to_string()
}
let this_file_raw = include_str!("update_rate_limit.rs");
let this_file = strip_ws(fn_body(
this_file_raw,
"fn log_eviction(&self, now: Instant, removed: usize) {",
"eviction",
));
assert_eq!(
level_of(
&this_file,
&strip_ws(concat!("UPDATE rate limiter at ", "capacity: evicted")),
"eviction",
),
strip_ws("tracing::info"),
"the eviction log must be emitted at info! or higher; debug! is compiled out of \
release builds by release_max_level_info (#4981)"
);
let new_pair_fn = strip_ws(fn_body(
this_file_raw,
"fn log_new_pair_budget(&self, now: Instant, sender: SocketAddr) {",
"new-pair-budget",
));
assert_eq!(
level_of(
&new_pair_fn,
&strip_ws(concat!(
"UPDATE rate limiter: peer is presenting ",
"(sender, contract) pairs this node is not tracking"
)),
"new-pair-budget",
),
strip_ws("tracing::info"),
"the new-pair-budget log must be emitted at info! or higher; debug! is compiled out \
of release builds by release_max_level_info (#4981)"
);
let rejected_fn = strip_ws(fn_body(
this_file_raw,
" fn log_rejected(",
"per-pair-rejection",
));
assert_eq!(
level_of(
&rejected_fn,
&strip_ws(concat!(
"UPDATE rate limiter: dropping UPDATEs from a ",
"(sender, contract) pair"
)),
"per-pair-rejection",
),
strip_ws("tracing::info"),
"the per-pair rejection log must be emitted at info! or higher; debug! is compiled \
out of release builds by release_max_level_info, which is how the #5510 broadcast \
drops left no evidence on a production node"
);
let node_rs = strip_ws(include_str!("../node.rs"));
assert_eq!(
level_of(
&node_rs,
&strip_ws(concat!("update_dispatch_", "capacity_dropped")),
"capacity-drop",
),
strip_ws("tracing::info"),
"dropping an UPDATE for capacity must be logged at info! or higher, so the drop is \
greppable on a production node (#4981)"
);
}
#[test]
fn cleanup_decrements_size_counter() {
let ts = SharedMockTimeSource::new();
let limiter = UpdateRateLimiter::with_config(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
4, );
for i in 0..4 {
assert_eq!(
limiter.check_and_record(
mk_sender(i + 1),
mk_contract(i + 1),
UpdateClass::Request
),
RateLimitDecision::Allowed
);
}
assert_eq!(
limiter.check_and_record(mk_sender(5), mk_contract(5), UpdateClass::Request),
RateLimitDecision::Allowed
);
assert_eq!(limiter.len(), 4, "eviction keeps the map at the cap");
ts.advance(CLEANUP_AGE + Duration::from_secs(1));
limiter.cleanup();
assert_eq!(limiter.len(), 0);
for i in 10..14 {
assert_eq!(
limiter.check_and_record(mk_sender(i), mk_contract(i), UpdateClass::Request),
RateLimitDecision::Allowed,
"after cleanup, new pair (sender={i}) should be admitted"
);
}
}
#[test]
fn concurrent_admission_at_capacity_almost_never_drops() {
use std::sync::{Arc as StdArc, Barrier};
use std::thread;
const CAP: usize = 1024;
const THREADS: usize = 16;
const PER_THREAD: usize = 400;
const ADMISSIONS: usize = THREADS * PER_THREAD;
const MAX_DROPS_PER_THOUSAND: usize = 2;
let ts = SharedMockTimeSource::new();
let limiter = StdArc::new(UpdateRateLimiter::with_new_pair_budget(
Arc::new(ts.clone()),
MIN_UPDATE_INTERVAL,
MIN_UPDATE_INTERVAL,
CAP,
f64::from(u32::MAX),
THREADS,
));
let barrier = StdArc::new(Barrier::new(THREADS));
let mut handles = Vec::with_capacity(THREADS);
for t in 0..THREADS {
let l = limiter.clone();
let b = barrier.clone();
handles.push(thread::spawn(move || {
b.wait();
for i in 0..PER_THREAD {
let sender = mk_sender(t as u8);
let mut id = [0u8; 32];
id[..8].copy_from_slice(&(i as u64).to_be_bytes());
l.check_and_record(sender, ContractInstanceId::new(id), UpdateClass::Request);
}
}));
}
for h in handles {
h.join().unwrap();
}
assert!(
limiter.capacity_evicted_total() > 0,
"fixture must actually saturate the cap, or it pins nothing"
);
let drops = limiter.capacity_rejected_total();
assert!(
drops as usize * 1000 <= ADMISSIONS * MAX_DROPS_PER_THOUSAND,
"at most {MAX_DROPS_PER_THOUSAND} drop per 1000 admissions, got {drops} \
in {ADMISSIONS}. An eviction that removed nothing means a concurrent \
caller freed the slots, so the admission must retry, not drop."
);
assert!(limiter.len() <= CAP);
}
}