CAS Contention in Traffic Bucket bounded

This commit is contained in:
Alexey
2026-09-25 06:32:01 +03:00
parent 08109d53e8
commit 26a574780e
9 changed files with 1308 additions and 159 deletions
+32
View File
@@ -118,6 +118,38 @@ pub(super) fn render(
}
);
let _ = writeln!(
out,
"# HELP telemt_rate_limiter_cas_retry_exhausted_total Traffic limiter operations that exhausted their bounded CAS attempt budget"
);
let _ = writeln!(out, "# TYPE telemt_rate_limiter_cas_retry_exhausted_total counter");
for (scope, direction, reserve, refund) in [
(
"user", "up", limiter_metrics.user_reserve_cas_retry_exhausted_up_total,
limiter_metrics.user_refund_cas_retry_exhausted_up_total,
),
(
"user", "down", limiter_metrics.user_reserve_cas_retry_exhausted_down_total,
limiter_metrics.user_refund_cas_retry_exhausted_down_total,
),
(
"cidr", "up", limiter_metrics.cidr_reserve_cas_retry_exhausted_up_total,
limiter_metrics.cidr_refund_cas_retry_exhausted_up_total,
),
(
"cidr", "down", limiter_metrics.cidr_reserve_cas_retry_exhausted_down_total,
limiter_metrics.cidr_refund_cas_retry_exhausted_down_total,
),
] {
for (operation, value) in [("reserve", reserve), ("refund", refund)] {
let _ = writeln!(
out,
"telemt_rate_limiter_cas_retry_exhausted_total{{scope=\"{scope}\",direction=\"{direction}\",operation=\"{operation}\"}} {}",
if core_enabled { value } else { 0 }
);
}
}
let _ = writeln!(
out,
"# HELP telemt_rate_limiter_active_leases Active relay leases under rate limiting by scope"
+38
View File
@@ -3,10 +3,22 @@ use http_body_util::BodyExt;
use std::net::IpAddr;
use std::time::SystemTime;
use crate::stats::telemetry::TelemetryPolicy;
use crate::tls_front::types::{
CachedTlsData, ParsedServerHello, TlsBehaviorProfile, TlsCertPayload, TlsProfileSource,
};
const CAS_CONTENTION_SERIES: [(&str, &str, &str, u64); 8] = [
("user", "up", "reserve", 1),
("user", "down", "reserve", 2),
("user", "up", "refund", 3),
("user", "down", "refund", 4),
("cidr", "up", "reserve", 5),
("cidr", "down", "reserve", 6),
("cidr", "up", "refund", 7),
("cidr", "down", "refund", 8),
];
fn test_web_publication() -> crate::web::control::WebRuntimePublication {
let control = crate::web::control::WebRuntimeControl::new();
control.subscribe().borrow().clone()
@@ -18,6 +30,9 @@ async fn test_render_metrics_format() {
let shared_state = ProxySharedState::new();
let tracker = UserIpTracker::new();
let mut config = ProxyConfig::default();
shared_state
.traffic_limiter
.set_cas_contention_metrics_for_test([1, 2, 3, 4, 5, 6, 7, 8]);
config
.access
.user_max_unique_ips
@@ -156,6 +171,17 @@ async fn test_render_metrics_format() {
assert!(output.contains("telemt_ip_tracker_users{scope=\"active\"} 1"));
assert!(output.contains("telemt_ip_tracker_entries{scope=\"active\"} 1"));
assert!(output.contains("telemt_ip_tracker_cleanup_queue_len 0"));
for (scope, direction, operation, value) in CAS_CONTENTION_SERIES {
assert!(output.contains(&format!(
"telemt_rate_limiter_cas_retry_exhausted_total{{scope=\"{scope}\",direction=\"{direction}\",operation=\"{operation}\"}} {value}"
)));
}
assert_eq!(
output
.matches("telemt_rate_limiter_cas_retry_exhausted_total{")
.count(),
8
);
}
#[tokio::test]
@@ -310,6 +336,13 @@ async fn process_tls_budget_metrics_survive_a_generation_without_tls_cache() {
async fn test_render_empty_stats() {
let stats = Stats::new();
let shared_state = ProxySharedState::new();
stats.apply_telemetry_policy(TelemetryPolicy {
core_enabled: false,
..TelemetryPolicy::default()
});
shared_state
.traffic_limiter
.set_cas_contention_metrics_for_test([1, 2, 3, 4, 5, 6, 7, 8]);
let tracker = UserIpTracker::new();
let config = ProxyConfig::default();
let output = render_metrics(
@@ -330,6 +363,11 @@ async fn test_render_empty_stats() {
assert!(output.contains("telemt_auth_budget_exhausted_total 0"));
assert!(output.contains("telemt_user_unique_ips_current{user="));
assert!(output.contains("telemt_user_unique_ips_recent_window{user="));
for (scope, direction, operation, _) in CAS_CONTENTION_SERIES {
assert!(output.contains(&format!(
"telemt_rate_limiter_cas_retry_exhausted_total{{scope=\"{scope}\",direction=\"{direction}\",operation=\"{operation}\"}} 0"
)));
}
}
#[tokio::test]
+96 -2
View File
@@ -33,45 +33,91 @@ const REGISTRY_SHARDS: usize = 64;
const FAIR_EPOCH_MS: u64 = 20;
const MAX_BORROW_CHUNK_BYTES: u64 = 32 * 1024;
const CLEANUP_INTERVAL_SECS: u64 = 60;
const RESERVE_CAS_ATTEMPT_LIMIT: usize = 16;
const REFUND_CAS_ATTEMPT_LIMIT: usize = 8;
const PACKED_USAGE_BITS: u32 = 28;
const PACKED_USAGE_MASK: u64 = (1u64 << PACKED_USAGE_BITS) - 1;
const PACKED_EPOCH_MAX: u64 = (1u64 << (u64::BITS - PACKED_USAGE_BITS)) - 1;
/// Traffic direction used by rate-limit accounting.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RateDirection {
/// Traffic received from a client.
Up,
/// Traffic sent to a client.
Down,
}
/// Result of an immediate traffic-budget request.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TrafficConsumeResult {
/// Number of bytes granted by all applicable limits.
pub granted: u64,
/// Whether the per-user limit prevented a grant.
pub blocked_user: bool,
/// Whether the per-CIDR limit prevented a grant.
pub blocked_cidr: bool,
}
/// Process-wide traffic-limiter counters and gauges.
#[derive(Debug, Clone, Copy)]
pub struct TrafficLimiterMetricsSnapshot {
/// Per-user upload throttle events.
pub user_throttle_up_total: u64,
/// Per-user download throttle events.
pub user_throttle_down_total: u64,
/// Per-CIDR upload throttle events.
pub cidr_throttle_up_total: u64,
/// Per-CIDR download throttle events.
pub cidr_throttle_down_total: u64,
/// Per-user accumulated upload wait time in milliseconds.
pub user_wait_up_ms_total: u64,
/// Per-user accumulated download wait time in milliseconds.
pub user_wait_down_ms_total: u64,
/// Per-CIDR accumulated upload wait time in milliseconds.
pub cidr_wait_up_ms_total: u64,
/// Per-CIDR accumulated download wait time in milliseconds.
pub cidr_wait_down_ms_total: u64,
/// Per-user upload reservations that exhausted their CAS budget.
pub user_reserve_cas_retry_exhausted_up_total: u64,
/// Per-user download reservations that exhausted their CAS budget.
pub user_reserve_cas_retry_exhausted_down_total: u64,
/// Per-user upload refunds that exhausted their CAS budget.
pub user_refund_cas_retry_exhausted_up_total: u64,
/// Per-user download refunds that exhausted their CAS budget.
pub user_refund_cas_retry_exhausted_down_total: u64,
/// Per-CIDR upload reservations that exhausted their CAS budget.
pub cidr_reserve_cas_retry_exhausted_up_total: u64,
/// Per-CIDR download reservations that exhausted their CAS budget.
pub cidr_reserve_cas_retry_exhausted_down_total: u64,
/// Per-CIDR upload refunds that exhausted their CAS budget.
pub cidr_refund_cas_retry_exhausted_up_total: u64,
/// Per-CIDR download refunds that exhausted their CAS budget.
pub cidr_refund_cas_retry_exhausted_down_total: u64,
/// Active leases with a per-user rate limit.
pub user_active_leases: u64,
/// Active leases with a per-CIDR rate limit.
pub cidr_active_leases: u64,
/// Configured per-user rate-limit entries.
pub user_policy_entries: u64,
/// Configured per-CIDR rate-limit entries.
pub cidr_policy_entries: u64,
}
#[derive(Default)]
struct CasContentionMetrics {
reserve_exhausted_total: AtomicU64,
refund_exhausted_total: AtomicU64,
}
#[derive(Default)]
struct ScopeMetrics {
throttle_up_total: AtomicU64,
throttle_down_total: AtomicU64,
wait_up_ms_total: AtomicU64,
wait_down_ms_total: AtomicU64,
contention_up: Arc<CasContentionMetrics>,
contention_down: Arc<CasContentionMetrics>,
active_leases: AtomicU64,
policy_entries: AtomicU64,
}
@@ -86,6 +132,54 @@ struct AtomicRatePair {
#[derive(Default)]
struct DirectionBucket {
state: AtomicU64,
contention: Arc<CasContentionMetrics>,
#[cfg(test)]
forced_reserve_failures: AtomicU64,
#[cfg(test)]
forced_refund_failures: AtomicU64,
#[cfg(test)]
reserve_cas_attempts: AtomicU64,
#[cfg(test)]
refund_cas_attempts: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BucketReserveError {
StaleEpoch,
Contended,
RefundContended,
ReserveAndRefundContended,
FairShareContended,
}
impl BucketReserveError {
fn exhausted_reserve_budget(self) -> bool {
matches!(
self,
Self::Contended | Self::ReserveAndRefundContended
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BucketRefundOutcome {
Complete,
Contended,
}
struct ReserveCasBudget {
remaining: usize,
}
struct BucketReservation {
granted: u64,
debit: Option<DirectionDebit>,
}
struct CidrReservation {
granted: u64,
aggregate_debit: Option<DirectionDebit>,
user_debit: Option<DirectionDebit>,
}
struct UserBucket {
@@ -95,13 +189,11 @@ struct UserBucket {
active_leases: AtomicU64,
}
#[derive(Default)]
struct CidrDirectionBucket {
used: Arc<DirectionBucket>,
active_users: Arc<DirectionBucket>,
}
#[derive(Default)]
struct CidrUserDirectionState {
used: Arc<DirectionBucket>,
}
@@ -165,6 +257,7 @@ struct TrafficLeaseBinding {
cidr_user_share: Option<Arc<CidrUserShare>>,
}
/// A live traffic-limiter binding for one authenticated client session.
pub struct TrafficLease {
limiter: Arc<TrafficLimiter>,
user: String,
@@ -173,6 +266,7 @@ pub struct TrafficLease {
refresh: ParkingMutex<()>,
}
/// Owns rate-limit policy, shared buckets, and limiter telemetry.
pub struct TrafficLimiter {
policy: ArcSwap<PolicySnapshot>,
policy_update: ParkingMutex<()>,
+305 -119
View File
@@ -1,27 +1,22 @@
use super::*;
impl ScopeMetrics {
pub(super) fn throttle(&self, direction: RateDirection) {
match direction {
RateDirection::Up => {
self.throttle_up_total.fetch_add(1, Ordering::Relaxed);
}
RateDirection::Down => {
self.throttle_down_total.fetch_add(1, Ordering::Relaxed);
}
impl ReserveCasBudget {
pub(super) fn new() -> Self {
Self {
remaining: RESERVE_CAS_ATTEMPT_LIMIT,
}
}
pub(super) fn wait_ms(&self, direction: RateDirection, wait_ms: u64) {
match direction {
RateDirection::Up => {
self.wait_up_ms_total.fetch_add(wait_ms, Ordering::Relaxed);
}
RateDirection::Down => {
self.wait_down_ms_total
.fetch_add(wait_ms, Ordering::Relaxed);
}
fn take(&mut self) -> bool {
if self.remaining == 0 {
return false;
}
self.remaining -= 1;
true
}
pub(super) fn is_exhausted(&self) -> bool {
self.remaining == 0
}
}
@@ -51,6 +46,49 @@ impl AtomicRatePair {
}
impl DirectionBucket {
fn new(contention: Arc<CasContentionMetrics>) -> Self {
Self {
state: AtomicU64::new(0),
contention,
#[cfg(test)]
forced_reserve_failures: AtomicU64::new(0),
#[cfg(test)]
forced_refund_failures: AtomicU64::new(0),
#[cfg(test)]
reserve_cas_attempts: AtomicU64::new(0),
#[cfg(test)]
refund_cas_attempts: AtomicU64::new(0),
}
}
#[inline(always)]
fn compare_exchange_reserve(&self, current: u64, next: u64) -> Result<u64, u64> {
#[cfg(test)]
if self.should_force_reserve_failure() {
return Err(current);
}
self.state.compare_exchange(
current,
next,
Ordering::Relaxed,
Ordering::Relaxed,
)
}
#[inline(always)]
fn compare_exchange_refund(&self, current: u64, next: u64) -> Result<u64, u64> {
#[cfg(test)]
if self.should_force_refund_failure() {
return Err(current);
}
self.state.compare_exchange(
current,
next,
Ordering::Relaxed,
Ordering::Relaxed,
)
}
fn unpack(state: u64) -> (u64, u64) {
(state >> PACKED_USAGE_BITS, state & PACKED_USAGE_MASK)
}
@@ -62,6 +100,7 @@ impl DirectionBucket {
Some((epoch << PACKED_USAGE_BITS) | used)
}
#[cfg(test)]
pub(super) fn used_at(&self, epoch: u64) -> Option<u64> {
if epoch > PACKED_EPOCH_MAX {
return None;
@@ -70,22 +109,35 @@ impl DirectionBucket {
(current_epoch == epoch).then_some(used)
}
fn used_in_epoch(&self, epoch: u64) -> Result<u64, BucketReserveError> {
let (observed_epoch, used) = Self::unpack(self.state.load(Ordering::Relaxed));
if observed_epoch != epoch {
return Err(BucketReserveError::StaleEpoch);
}
Ok(used)
}
pub(super) fn try_reserve_at(
self: &Arc<Self>,
epoch: u64,
cap: u64,
requested: u64,
) -> Option<DirectionDebit> {
if requested == 0 || cap == 0 || epoch > PACKED_EPOCH_MAX {
return None;
budget: &mut ReserveCasBudget,
) -> Result<Option<DirectionDebit>, BucketReserveError> {
if requested == 0 || cap == 0 {
return Ok(None);
}
if epoch > PACKED_EPOCH_MAX {
return Err(BucketReserveError::StaleEpoch);
}
let cap = cap.min(PACKED_USAGE_MASK);
let mut observed = self.state.load(Ordering::Relaxed);
// The transaction-owned budget bounds retries across every participating bucket.
loop {
let (observed_epoch, observed_used) = Self::unpack(observed);
if observed_epoch > epoch {
return None;
return Err(BucketReserveError::StaleEpoch);
}
let used = if observed_epoch == epoch {
observed_used
@@ -93,71 +145,80 @@ impl DirectionBucket {
0
};
if used >= cap {
return None;
return Ok(None);
}
let remaining = cap - used;
let grant = requested.min(remaining);
if grant == 0 {
return None;
return Ok(None);
}
let next = Self::pack(epoch, used + grant)?;
match self.state.compare_exchange_weak(
observed,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
let Some(next) = Self::pack(epoch, used + grant) else {
return Err(BucketReserveError::StaleEpoch);
};
if !budget.take() {
return Err(BucketReserveError::Contended);
}
match self.compare_exchange_reserve(observed, next) {
Ok(_) => {
return Some(DirectionDebit {
return Ok(Some(DirectionDebit {
bucket: Arc::clone(self),
epoch,
refundable: grant,
});
}));
}
Err(actual) => observed = actual,
}
}
}
fn refund_at(&self, epoch: u64, bytes: u64) {
fn refund_at(&self, epoch: u64, bytes: u64) -> BucketRefundOutcome {
if bytes == 0 || epoch > PACKED_EPOCH_MAX {
return;
return BucketRefundOutcome::Complete;
}
let mut observed = self.state.load(Ordering::Relaxed);
loop {
for _ in 0..REFUND_CAS_ATTEMPT_LIMIT {
let (observed_epoch, used) = Self::unpack(observed);
if observed_epoch != epoch || used == 0 {
return;
return BucketRefundOutcome::Complete;
}
let next = Self::pack(epoch, used.saturating_sub(bytes)).unwrap_or(observed);
match self.state.compare_exchange_weak(
observed,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => return,
match self.compare_exchange_refund(observed, next) {
Ok(_) => return BucketRefundOutcome::Complete,
Err(actual) => observed = actual,
}
}
let (observed_epoch, used) = Self::unpack(observed);
if observed_epoch != epoch || used == 0 {
return BucketRefundOutcome::Complete;
}
// Retain the charge after bounded contention so accounting cannot under-enforce.
self.contention
.refund_exhausted_total
.fetch_add(1, Ordering::Relaxed);
BucketRefundOutcome::Contended
}
}
impl DirectionDebit {
fn granted(&self) -> u64 {
self.refundable
}
pub(super) fn shrink_to(&mut self, retained: u64) {
pub(super) fn shrink_to(&mut self, retained: u64) -> BucketRefundOutcome {
let retained = retained.min(self.refundable);
self.bucket
let outcome = self
.bucket
.refund_at(self.epoch, self.refundable - retained);
// Never retry the attempted delta in Drop; failed refunds remain charged.
self.refundable = retained;
outcome
}
fn refund_all(&mut self) -> BucketRefundOutcome {
self.shrink_to(0)
}
pub(super) fn settle(&mut self, committed: u64) {
self.shrink_to(committed);
let _ = self.shrink_to(committed);
self.refundable = 0;
}
@@ -170,16 +231,21 @@ impl DirectionDebit {
impl Drop for DirectionDebit {
fn drop(&mut self) {
self.bucket.refund_at(self.epoch, self.refundable);
let _ = self.bucket.refund_at(self.epoch, self.refundable);
}
}
impl UserBucket {
pub(super) fn new(revision: u64, limits: RateLimitBps) -> Self {
pub(super) fn new(
revision: u64,
limits: RateLimitBps,
up_contention: Arc<CasContentionMetrics>,
down_contention: Arc<CasContentionMetrics>,
) -> Self {
Self {
rates: AtomicRatePair::new(revision, limits),
up: Arc::new(DirectionBucket::default()),
down: Arc::new(DirectionBucket::default()),
up: Arc::new(DirectionBucket::new(up_contention)),
down: Arc::new(DirectionBucket::new(down_contention)),
active_leases: AtomicU64::new(0),
}
}
@@ -191,107 +257,195 @@ impl UserBucket {
pub(super) fn try_reserve(
&self,
direction: RateDirection,
epoch: u64,
requested: u64,
) -> (u64, Option<DirectionDebit>) {
budget: &mut ReserveCasBudget,
) -> Result<BucketReservation, BucketReserveError> {
let cap_bps = self.rates.get(direction);
if cap_bps == 0 {
return (requested, None);
return Ok(BucketReservation {
granted: requested,
debit: None,
});
}
let cap = bytes_per_epoch(cap_bps);
let debit = match direction {
RateDirection::Up => self.up.try_reserve_at(current_epoch(), cap, requested),
RateDirection::Down => self.down.try_reserve_at(current_epoch(), cap, requested),
RateDirection::Up => self.up.try_reserve_at(epoch, cap, requested, budget)?,
RateDirection::Down => self.down.try_reserve_at(epoch, cap, requested, budget)?,
};
let granted = debit.as_ref().map(DirectionDebit::granted).unwrap_or(0);
(granted, debit)
Ok(BucketReservation { granted, debit })
}
}
impl CidrDirectionBucket {
fn new(contention: Arc<CasContentionMetrics>) -> Self {
Self {
used: Arc::new(DirectionBucket::new(Arc::clone(&contention))),
active_users: Arc::new(DirectionBucket::new(contention)),
}
}
pub(super) fn try_reserve(
&self,
user_state: &CidrUserDirectionState,
epoch: u64,
cap_epoch: u64,
requested: u64,
) -> (u64, Option<DirectionDebit>, Option<DirectionDebit>) {
budget: &mut ReserveCasBudget,
) -> Result<CidrReservation, BucketReserveError> {
if requested == 0 || cap_epoch == 0 {
return (0, None, None);
return Ok(CidrReservation {
granted: 0,
aggregate_debit: None,
user_debit: None,
});
}
let epoch = current_epoch();
if !user_state.ensure_active(epoch, &self.active_users) {
return (0, None, None);
if !user_state.ensure_active(epoch, &self.active_users, budget)? {
return Ok(CidrReservation {
granted: 0,
aggregate_debit: None,
user_debit: None,
});
}
let Some(active_users) = self.active_users.used_at(epoch) else {
return (0, None, None);
};
let active_users = self.active_users.used_in_epoch(epoch)?;
let active_users = active_users.max(1);
let fair_share = cap_epoch.saturating_div(active_users).max(1);
loop {
let Some(user_used) = user_state.used.used_at(epoch) else {
return (0, None, None);
};
let guaranteed_remaining = fair_share.saturating_sub(user_used);
let (user_cap, desired) = if guaranteed_remaining > 0 {
(fair_share, requested.min(guaranteed_remaining))
let user_used = user_state.used.used_in_epoch(epoch)?;
let guaranteed_remaining = fair_share.saturating_sub(user_used);
let (user_cap, desired) = if guaranteed_remaining > 0 {
(fair_share, requested.min(guaranteed_remaining))
} else {
(PACKED_USAGE_MASK, requested.min(MAX_BORROW_CHUNK_BYTES))
};
let mut user_debit = user_state
.used
.try_reserve_at(epoch, user_cap, desired, budget)?;
// A competing reservation can consume the guaranteed share between the snapshot and CAS.
if user_debit.is_none() && guaranteed_remaining > 0 {
let refreshed_used = user_state.used.used_in_epoch(epoch)?;
if refreshed_used >= fair_share {
user_debit = user_state.used.try_reserve_at(
epoch,
PACKED_USAGE_MASK,
requested.min(MAX_BORROW_CHUNK_BYTES),
budget,
)?;
} else {
(PACKED_USAGE_MASK, requested.min(MAX_BORROW_CHUNK_BYTES))
};
let Some(mut user_debit) = user_state.used.try_reserve_at(epoch, user_cap, desired)
else {
if guaranteed_remaining > 0 {
continue;
}
return (0, None, None);
};
let user_granted = user_debit.granted();
let Some(aggregate_debit) = self.used.try_reserve_at(epoch, cap_epoch, user_granted)
else {
return (0, None, None);
};
let granted = aggregate_debit.granted();
if granted < user_granted {
user_debit.shrink_to(granted);
user_debit = user_state.used.try_reserve_at(
epoch,
fair_share,
requested.min(fair_share - refreshed_used),
budget,
)?;
}
if user_debit.is_none() {
return Err(BucketReserveError::FairShareContended);
}
return (granted, Some(aggregate_debit), Some(user_debit));
}
let Some(mut user_debit) = user_debit else {
return Ok(CidrReservation {
granted: 0,
aggregate_debit: None,
user_debit: None,
});
};
let user_granted = user_debit.granted();
let Some(aggregate_debit) = self
.used
.try_reserve_at(epoch, cap_epoch, user_granted, budget)?
else {
return Ok(CidrReservation {
granted: 0,
aggregate_debit: None,
user_debit: None,
});
};
let granted = aggregate_debit.granted();
if granted < user_granted {
let _ = user_debit.shrink_to(granted);
}
Ok(CidrReservation {
granted,
aggregate_debit: Some(aggregate_debit),
user_debit: Some(user_debit),
})
}
}
#[cfg(test)]
impl Default for CidrDirectionBucket {
fn default() -> Self {
Self::new(Arc::new(CasContentionMetrics::default()))
}
}
impl CidrUserDirectionState {
pub(super) fn ensure_active(&self, epoch: u64, active_users: &Arc<DirectionBucket>) -> bool {
fn new(contention: Arc<CasContentionMetrics>) -> Self {
Self {
used: Arc::new(DirectionBucket::new(contention)),
}
}
pub(super) fn ensure_active(
&self,
epoch: u64,
active_users: &Arc<DirectionBucket>,
budget: &mut ReserveCasBudget,
) -> Result<bool, BucketReserveError> {
if epoch > PACKED_EPOCH_MAX {
return false;
return Err(BucketReserveError::StaleEpoch);
}
let mut observed = self.used.state.load(Ordering::Relaxed);
// Every retry spends shared reserve budget before publishing the user epoch.
loop {
let (observed_epoch, _) = DirectionBucket::unpack(observed);
if observed_epoch == epoch {
return true;
return Ok(true);
}
if observed_epoch > epoch {
return false;
return Err(BucketReserveError::StaleEpoch);
}
let Some(mut active_debit) = active_users.try_reserve_at(epoch, PACKED_USAGE_MASK, 1)
let Some(mut active_debit) = active_users.try_reserve_at(
epoch,
PACKED_USAGE_MASK,
1,
budget,
)?
else {
return false;
return Ok(false);
};
let Some(next) = DirectionBucket::pack(epoch, 0) else {
return false;
return Err(BucketReserveError::StaleEpoch);
};
match self.used.state.compare_exchange(
observed,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
if !budget.take() {
if active_debit.refund_all() == BucketRefundOutcome::Contended {
return Err(BucketReserveError::ReserveAndRefundContended);
}
return Err(BucketReserveError::Contended);
}
match self.used.compare_exchange_reserve(observed, next) {
Ok(_) => {
active_debit.commit_all();
return true;
return Ok(true);
}
Err(actual) => {
drop(active_debit);
let refund = active_debit.refund_all();
let (actual_epoch, _) = DirectionBucket::unpack(actual);
if refund == BucketRefundOutcome::Contended {
return Err(if budget.is_exhausted() {
BucketReserveError::ReserveAndRefundContended
} else {
BucketReserveError::RefundContended
});
}
if actual_epoch == epoch {
return Ok(true);
}
observed = actual;
}
}
@@ -299,22 +453,37 @@ impl CidrUserDirectionState {
}
}
#[cfg(test)]
impl Default for CidrUserDirectionState {
fn default() -> Self {
Self::new(Arc::new(CasContentionMetrics::default()))
}
}
impl CidrUserShare {
pub(super) fn new() -> Self {
pub(super) fn new(
up_contention: Arc<CasContentionMetrics>,
down_contention: Arc<CasContentionMetrics>,
) -> Self {
Self {
active_conns: AtomicU64::new(0),
up: CidrUserDirectionState::default(),
down: CidrUserDirectionState::default(),
up: CidrUserDirectionState::new(up_contention),
down: CidrUserDirectionState::new(down_contention),
}
}
}
impl CidrBucket {
pub(super) fn new(revision: u64, limits: RateLimitBps) -> Self {
pub(super) fn new(
revision: u64,
limits: RateLimitBps,
up_contention: Arc<CasContentionMetrics>,
down_contention: Arc<CasContentionMetrics>,
) -> Self {
Self {
rates: AtomicRatePair::new(revision, limits),
up: CidrDirectionBucket::default(),
down: CidrDirectionBucket::default(),
up: CidrDirectionBucket::new(up_contention),
down: CidrDirectionBucket::new(down_contention),
users: ShardedRegistry::new(REGISTRY_SHARDS),
active_leases: AtomicU64::new(0),
}
@@ -325,10 +494,15 @@ impl CidrBucket {
}
pub(super) fn acquire_user_share(&self, user: &str) -> Arc<CidrUserShare> {
self.users
.get_or_insert_with(user, CidrUserShare::new, |share| {
let up_contention = Arc::clone(&self.up.used.contention);
let down_contention = Arc::clone(&self.down.used.contention);
self.users.get_or_insert_with(
user,
|| CidrUserShare::new(up_contention, down_contention),
|share| {
share.active_conns.fetch_add(1, Ordering::Relaxed);
})
},
)
}
pub(super) fn release_user_share(&self, user: &str, share: &Arc<CidrUserShare>) {
@@ -344,16 +518,28 @@ impl CidrBucket {
&self,
direction: RateDirection,
share: &CidrUserShare,
epoch: u64,
requested: u64,
) -> (u64, Option<DirectionDebit>, Option<DirectionDebit>) {
budget: &mut ReserveCasBudget,
) -> Result<CidrReservation, BucketReserveError> {
let cap_bps = self.rates.get(direction);
if cap_bps == 0 {
return (requested, None, None);
return Ok(CidrReservation {
granted: requested,
aggregate_debit: None,
user_debit: None,
});
}
let cap_epoch = bytes_per_epoch(cap_bps);
match direction {
RateDirection::Up => self.up.try_reserve(&share.up, cap_epoch, requested),
RateDirection::Down => self.down.try_reserve(&share.down, cap_epoch, requested),
RateDirection::Up => {
self.up
.try_reserve(&share.up, epoch, cap_epoch, requested, budget)
}
RateDirection::Down => {
self.down
.try_reserve(&share.down, epoch, cap_epoch, requested, budget)
}
}
}
+1
View File
@@ -2,6 +2,7 @@ use std::sync::OnceLock;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use super::*;
/// Returns the delay until the next traffic-limiter refill boundary.
pub fn next_refill_delay() -> Duration {
let start = limiter_epoch_start();
let elapsed_ms = start.elapsed().as_millis() as u64;
+79 -12
View File
@@ -42,13 +42,53 @@ impl TrafficLease {
cidr_user: None,
};
}
if binding.user_bucket.is_none() && binding.cidr_bucket.is_none() {
return TrafficReservation {
result: TrafficConsumeResult {
granted: requested,
blocked_user: false,
blocked_cidr: false,
},
_binding: binding,
user: None,
cidr: None,
cidr_user: None,
};
}
let epoch = current_epoch();
let mut budget = ReserveCasBudget::new();
let mut granted = requested;
let mut user_debit = None;
if let Some(user_bucket) = binding.user_bucket.as_ref() {
let (user_granted, debit) = user_bucket.try_reserve(direction, granted);
user_debit = debit;
if user_granted == 0 {
let user_reservation = match user_bucket.try_reserve(
direction,
epoch,
granted,
&mut budget,
) {
Ok(reservation) => reservation,
Err(error) => {
if error.exhausted_reserve_budget() {
self.limiter
.user_scope
.reserve_cas_retry_exhausted(direction);
}
return TrafficReservation {
result: TrafficConsumeResult {
granted: 0,
blocked_user: false,
blocked_cidr: false,
},
_binding: binding,
user: None,
cidr: None,
cidr_user: None,
};
}
};
user_debit = user_reservation.debit;
if user_reservation.granted == 0 {
self.limiter.observe_throttle(direction, true, false);
return TrafficReservation {
result: TrafficConsumeResult {
@@ -62,7 +102,7 @@ impl TrafficLease {
cidr_user: None,
};
}
granted = user_granted;
granted = user_reservation.granted;
}
let mut cidr_debit = None;
@@ -70,16 +110,41 @@ impl TrafficLease {
if let (Some(cidr_bucket), Some(cidr_user_share)) =
(binding.cidr_bucket.as_ref(), binding.cidr_user_share.as_ref())
{
let (cidr_granted, aggregate_debit, share_debit) =
cidr_bucket.try_reserve_for_user(direction, cidr_user_share, granted);
cidr_debit = aggregate_debit;
cidr_user_debit = share_debit;
if cidr_granted < granted
let cidr_reservation = match cidr_bucket.try_reserve_for_user(
direction,
cidr_user_share,
epoch,
granted,
&mut budget,
) {
Ok(reservation) => reservation,
Err(error) => {
if error.exhausted_reserve_budget() {
self.limiter
.cidr_scope
.reserve_cas_retry_exhausted(direction);
}
return TrafficReservation {
result: TrafficConsumeResult {
granted: 0,
blocked_user: false,
blocked_cidr: false,
},
_binding: binding,
user: user_debit,
cidr: None,
cidr_user: None,
};
}
};
cidr_debit = cidr_reservation.aggregate_debit;
cidr_user_debit = cidr_reservation.user_debit;
if cidr_reservation.granted < granted
&& let Some(debit) = user_debit.as_mut()
{
debit.shrink_to(cidr_granted);
let _ = debit.shrink_to(cidr_reservation.granted);
}
if cidr_granted == 0 {
if cidr_reservation.granted == 0 {
self.limiter.observe_throttle(direction, false, true);
return TrafficReservation {
result: TrafficConsumeResult {
@@ -93,7 +158,7 @@ impl TrafficLease {
cidr_user: cidr_user_debit,
};
}
granted = cidr_granted;
granted = cidr_reservation.granted;
}
TrafficReservation {
@@ -109,6 +174,7 @@ impl TrafficLease {
}
}
/// Immediately consumes the budget granted by all applicable limits.
pub fn try_consume(&self, direction: RateDirection, requested: u64) -> TrafficConsumeResult {
let reservation = self.try_reserve(direction, requested);
let result = reservation.result();
@@ -116,6 +182,7 @@ impl TrafficLease {
result
}
/// Records caller-observed limiter wait time for the blocking scopes.
pub fn observe_wait_ms(
&self,
direction: RateDirection,
+130 -2
View File
@@ -1,7 +1,51 @@
use crate::config::CidrRateLimitKey;
use super::*;
impl ScopeMetrics {
pub(super) fn throttle(&self, direction: RateDirection) {
match direction {
RateDirection::Up => {
self.throttle_up_total.fetch_add(1, Ordering::Relaxed);
}
RateDirection::Down => {
self.throttle_down_total.fetch_add(1, Ordering::Relaxed);
}
}
}
pub(super) fn wait_ms(&self, direction: RateDirection, wait_ms: u64) {
match direction {
RateDirection::Up => {
self.wait_up_ms_total.fetch_add(wait_ms, Ordering::Relaxed);
}
RateDirection::Down => {
self.wait_down_ms_total
.fetch_add(wait_ms, Ordering::Relaxed);
}
}
}
pub(super) fn contention(&self, direction: RateDirection) -> Arc<CasContentionMetrics> {
match direction {
RateDirection::Up => Arc::clone(&self.contention_up),
RateDirection::Down => Arc::clone(&self.contention_down),
}
}
pub(super) fn reserve_cas_retry_exhausted(&self, direction: RateDirection) {
let metrics = match direction {
RateDirection::Up => self.contention_up.as_ref(),
RateDirection::Down => self.contention_down.as_ref(),
};
metrics
.reserve_exhausted_total
.fetch_add(1, Ordering::Relaxed);
}
}
impl TrafficLimiter {
/// Creates an empty limiter with no active rate-limit policy.
pub fn new() -> Arc<Self> {
Arc::new(Self {
policy: ArcSwap::from_pointee(PolicySnapshot::default()),
@@ -15,6 +59,7 @@ impl TrafficLimiter {
}
#[cfg(test)]
/// Replaces the in-memory policy for isolated limiter tests.
pub fn apply_policy(
&self,
user_limits: HashMap<String, RateLimitBps>,
@@ -125,6 +170,7 @@ impl TrafficLimiter {
true
}
/// Creates a lease that follows policy revisions for one client identity.
pub fn acquire_lease(
self: &Arc<Self>,
user: &str,
@@ -151,7 +197,14 @@ impl TrafficLimiter {
if let Some(limit) = policy.user_limits.get(user).copied() {
let bucket = self.user_buckets.get_or_insert_with(
user,
|| UserBucket::new(policy.revision, limit),
|| {
UserBucket::new(
policy.revision,
limit,
self.user_scope.contention(RateDirection::Up),
self.user_scope.contention(RateDirection::Down),
)
},
|bucket| {
bucket.active_leases.fetch_add(1, Ordering::Relaxed);
},
@@ -173,7 +226,14 @@ impl TrafficLimiter {
};
let bucket = self.cidr_buckets.get_or_insert_with(
key,
|| CidrBucket::new(policy.revision, limits),
|| {
CidrBucket::new(
policy.revision,
limits,
self.cidr_scope.contention(RateDirection::Up),
self.cidr_scope.contention(RateDirection::Down),
)
},
|bucket| {
bucket.active_leases.fetch_add(1, Ordering::Relaxed);
},
@@ -198,6 +258,7 @@ impl TrafficLimiter {
})
}
/// Captures limiter telemetry without locking bucket registries.
pub fn metrics_snapshot(&self) -> TrafficLimiterMetricsSnapshot {
TrafficLimiterMetricsSnapshot {
user_throttle_up_total: self.user_scope.throttle_up_total.load(Ordering::Relaxed),
@@ -208,6 +269,46 @@ impl TrafficLimiter {
user_wait_down_ms_total: self.user_scope.wait_down_ms_total.load(Ordering::Relaxed),
cidr_wait_up_ms_total: self.cidr_scope.wait_up_ms_total.load(Ordering::Relaxed),
cidr_wait_down_ms_total: self.cidr_scope.wait_down_ms_total.load(Ordering::Relaxed),
user_reserve_cas_retry_exhausted_up_total: self
.user_scope
.contention_up
.reserve_exhausted_total
.load(Ordering::Relaxed),
user_reserve_cas_retry_exhausted_down_total: self
.user_scope
.contention_down
.reserve_exhausted_total
.load(Ordering::Relaxed),
user_refund_cas_retry_exhausted_up_total: self
.user_scope
.contention_up
.refund_exhausted_total
.load(Ordering::Relaxed),
user_refund_cas_retry_exhausted_down_total: self
.user_scope
.contention_down
.refund_exhausted_total
.load(Ordering::Relaxed),
cidr_reserve_cas_retry_exhausted_up_total: self
.cidr_scope
.contention_up
.reserve_exhausted_total
.load(Ordering::Relaxed),
cidr_reserve_cas_retry_exhausted_down_total: self
.cidr_scope
.contention_down
.reserve_exhausted_total
.load(Ordering::Relaxed),
cidr_refund_cas_retry_exhausted_up_total: self
.cidr_scope
.contention_up
.refund_exhausted_total
.load(Ordering::Relaxed),
cidr_refund_cas_retry_exhausted_down_total: self
.cidr_scope
.contention_down
.refund_exhausted_total
.load(Ordering::Relaxed),
user_active_leases: self.user_scope.active_leases.load(Ordering::Relaxed),
cidr_active_leases: self.cidr_scope.active_leases.load(Ordering::Relaxed),
user_policy_entries: self.user_scope.policy_entries.load(Ordering::Relaxed),
@@ -215,6 +316,33 @@ impl TrafficLimiter {
}
}
/// Sets fixed contention counters for renderer mapping tests.
#[cfg(test)]
pub(crate) fn set_cas_contention_metrics_for_test(&self, values: [u64; 8]) {
let [
user_reserve_up,
user_reserve_down,
user_refund_up,
user_refund_down,
cidr_reserve_up,
cidr_reserve_down,
cidr_refund_up,
cidr_refund_down,
] = values;
for (counter, value) in [
(&self.user_scope.contention_up.reserve_exhausted_total, user_reserve_up),
(&self.user_scope.contention_down.reserve_exhausted_total, user_reserve_down),
(&self.user_scope.contention_up.refund_exhausted_total, user_refund_up),
(&self.user_scope.contention_down.refund_exhausted_total, user_refund_down),
(&self.cidr_scope.contention_up.reserve_exhausted_total, cidr_reserve_up),
(&self.cidr_scope.contention_down.reserve_exhausted_total, cidr_reserve_down),
(&self.cidr_scope.contention_up.refund_exhausted_total, cidr_refund_up),
(&self.cidr_scope.contention_down.refund_exhausted_total, cidr_refund_down),
] {
counter.store(value, Ordering::Relaxed);
}
}
pub(super) fn observe_throttle(
&self,
direction: RateDirection,
+126 -24
View File
@@ -1,10 +1,78 @@
use super::*;
use crate::config::CidrRateLimitKey;
mod bucket_contention;
impl DirectionBucket {
pub(crate) fn should_force_reserve_failure(&self) -> bool {
self.reserve_cas_attempts.fetch_add(1, Ordering::Relaxed);
let remaining = self.forced_reserve_failures.load(Ordering::Relaxed);
remaining > 0
&& self
.forced_reserve_failures
.compare_exchange(
remaining,
remaining - 1,
Ordering::Relaxed,
Ordering::Relaxed,
)
.is_ok()
}
pub(crate) fn should_force_refund_failure(&self) -> bool {
self.refund_cas_attempts.fetch_add(1, Ordering::Relaxed);
let remaining = self.forced_refund_failures.load(Ordering::Relaxed);
remaining > 0
&& self
.forced_refund_failures
.compare_exchange(
remaining,
remaining - 1,
Ordering::Relaxed,
Ordering::Relaxed,
)
.is_ok()
}
pub(crate) fn force_reserve_failures(&self, failures: usize) {
self.reserve_cas_attempts.store(0, Ordering::Relaxed);
self.forced_reserve_failures
.store(failures as u64, Ordering::Relaxed);
}
pub(crate) fn force_refund_failures(&self, failures: usize) {
self.refund_cas_attempts.store(0, Ordering::Relaxed);
self.forced_refund_failures
.store(failures as u64, Ordering::Relaxed);
}
pub(crate) fn reserve_cas_attempts(&self) -> u64 {
self.reserve_cas_attempts.load(Ordering::Relaxed)
}
pub(crate) fn refund_cas_attempts(&self) -> u64 {
self.refund_cas_attempts.load(Ordering::Relaxed)
}
}
fn rate(up_bps: u64, down_bps: u64) -> RateLimitBps {
RateLimitBps { up_bps, down_bps }
}
fn reserve_at(
bucket: &Arc<DirectionBucket>,
epoch: u64,
cap: u64,
requested: u64,
) -> Result<Option<DirectionDebit>, BucketReserveError> {
bucket.try_reserve_at(
epoch,
cap,
requested,
&mut ReserveCasBudget::new(),
)
}
#[test]
fn stale_runtime_cannot_overwrite_newer_rate_policy() {
let limiter = TrafficLimiter::new();
@@ -151,12 +219,20 @@ fn auto_cidr_bucket_key_canonicalizes_network_address() {
#[test]
fn refund_from_an_old_epoch_does_not_reduce_the_current_epoch() {
let bucket = Arc::new(DirectionBucket::default());
let old_debit = bucket.try_reserve_at(7, 100, 80).unwrap();
let current_debit = bucket.try_reserve_at(8, 100, 60).unwrap();
let old_debit = reserve_at(&bucket, 7, 100, 80).unwrap().unwrap();
let current_debit = reserve_at(&bucket, 8, 100, 60).unwrap().unwrap();
drop(old_debit);
assert_eq!(bucket.used_at(8), Some(60));
assert_eq!(bucket.refund_cas_attempts(), 0);
assert_eq!(
bucket
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
0
);
drop(current_debit);
}
@@ -172,8 +248,8 @@ fn concurrent_rollover_cannot_publish_multiple_epoch_budgets() {
let barrier = Arc::clone(&barrier);
threads.push(std::thread::spawn(move || {
barrier.wait();
bucket
.try_reserve_at(9, 100, 100)
reserve_at(&bucket, 9, 100, 100)
.unwrap()
.map(|mut debit| debit.commit_all())
.unwrap_or(0)
}));
@@ -202,8 +278,8 @@ fn scheduler_pressure_never_exceeds_a_packed_epoch_budget() {
let mut grants = Vec::with_capacity(EPOCHS);
for epoch in 1..=EPOCHS as u64 {
barrier.wait();
let granted = bucket
.try_reserve_at(epoch, 100, 100)
let granted = reserve_at(&bucket, epoch, 100, 100)
.unwrap()
.map(|mut debit| debit.commit_all())
.unwrap_or(0);
grants.push(granted);
@@ -228,7 +304,12 @@ fn scheduler_pressure_never_exceeds_a_packed_epoch_budget() {
#[test]
fn stale_policy_revision_cannot_restore_an_old_rate() {
let bucket = UserBucket::new(2, rate(2_000, 3_000));
let bucket = UserBucket::new(
2,
rate(2_000, 3_000),
Arc::new(CasContentionMetrics::default()),
Arc::new(CasContentionMetrics::default()),
);
bucket.set_rates(3, rate(4_000, 5_000));
bucket.set_rates(2, rate(6_000, 7_000));
@@ -240,16 +321,15 @@ fn stale_policy_revision_cannot_restore_an_old_rate() {
#[test]
fn dropped_debit_refunds_only_its_packed_epoch() {
let bucket = Arc::new(DirectionBucket::default());
let debit = bucket.try_reserve_at(11, 100, 80).unwrap();
let debit = reserve_at(&bucket, 11, 100, 80).unwrap().unwrap();
drop(debit);
assert_eq!(bucket.used_at(11), Some(0));
assert!(
bucket
.try_reserve_at(PACKED_EPOCH_MAX + 1, 100, 1)
.is_none()
);
assert!(matches!(
reserve_at(&bucket, PACKED_EPOCH_MAX + 1, 100, 1),
Err(BucketReserveError::StaleEpoch)
));
}
#[test]
@@ -266,14 +346,34 @@ fn concurrent_first_use_counts_one_active_cidr_user() {
let barrier = Arc::clone(&barrier);
threads.push(std::thread::spawn(move || {
barrier.wait();
assert!(user.ensure_active(13, &bucket.active_users));
user.ensure_active(
13,
&bucket.active_users,
&mut ReserveCasBudget::new(),
)
}));
}
for thread in threads {
thread.join().unwrap();
}
let results: Vec<_> = threads
.into_iter()
.map(|thread| thread.join().unwrap())
.collect();
assert_eq!(bucket.active_users.used_at(13), Some(1));
assert!(results.iter().any(|result| *result == Ok(true)));
assert!(results.iter().all(|result| matches!(
result,
Ok(true)
| Err(BucketReserveError::Contended)
| Err(BucketReserveError::RefundContended)
| Err(BucketReserveError::ReserveAndRefundContended)
)));
let leaked = bucket
.active_users
.contention
.refund_exhausted_total
.load(Ordering::Relaxed);
assert_eq!(bucket.active_users.used_at(13), Some(1 + leaked));
assert!(1 + leaked <= CONTENDERS as u64);
assert_eq!(user.used.used_at(13), Some(0));
}
#[test]
@@ -300,6 +400,8 @@ fn dropped_traffic_reservation_refunds_user_and_cidr_debits() {
let reservation = lease.try_reserve(RateDirection::Down, 800);
assert_eq!(reservation.result().granted, 800);
let epoch = reservation.user.as_ref().unwrap().epoch;
assert_eq!(reservation.cidr.as_ref().unwrap().epoch, epoch);
assert_eq!(reservation.cidr_user.as_ref().unwrap().epoch, epoch);
drop(reservation);
let binding = lease.binding.load_full();
@@ -347,12 +449,12 @@ fn active_lease_observes_policy_removal() {
let lease = limiter
.acquire_lease("alice", "203.0.113.7".parse().unwrap())
.unwrap();
assert_eq!(lease.try_consume(RateDirection::Up, 1).granted, 1);
assert_eq!(lease.try_consume(RateDirection::Up, 1).granted, 0);
assert!(lease.binding.load().user_bucket.is_some());
limiter.apply_policy(HashMap::new(), HashMap::new());
assert_eq!(lease.try_consume(RateDirection::Up, 1).granted, 1);
assert_eq!(lease.try_consume(RateDirection::Up, 2).granted, 2);
assert!(lease.binding.load().user_bucket.is_none());
}
#[test]
@@ -361,14 +463,14 @@ fn active_lease_observes_policy_addition() {
let lease = limiter
.acquire_lease("alice", "203.0.113.7".parse().unwrap())
.unwrap();
assert_eq!(lease.try_consume(RateDirection::Up, 2).granted, 2);
assert!(lease.binding.load().user_bucket.is_none());
let mut user_limits = HashMap::new();
user_limits.insert("alice".to_string(), rate(1, 0));
limiter.apply_policy(user_limits, HashMap::new());
assert_eq!(lease.try_consume(RateDirection::Up, 1).granted, 1);
assert_eq!(lease.try_consume(RateDirection::Up, 1).granted, 0);
assert_eq!(lease.try_consume(RateDirection::Up, 2).granted, 1);
assert!(lease.binding.load().user_bucket.is_some());
}
#[test]
@@ -0,0 +1,501 @@
use super::*;
use proptest::prelude::*;
#[test]
fn reserve_stops_after_the_attempt_limit() {
let bucket = Arc::new(DirectionBucket::default());
bucket.force_reserve_failures(RESERVE_CAS_ATTEMPT_LIMIT);
let mut budget = ReserveCasBudget::new();
let reservation = bucket.try_reserve_at(1, 100, 1, &mut budget);
assert!(matches!(
reservation,
Err(BucketReserveError::Contended)
));
assert_eq!(
bucket.reserve_cas_attempts(),
RESERVE_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(bucket.used_at(1), None);
}
#[test]
fn reserve_succeeds_on_the_last_allowed_attempt() {
let bucket = Arc::new(DirectionBucket::default());
bucket.force_reserve_failures(RESERVE_CAS_ATTEMPT_LIMIT - 1);
let mut budget = ReserveCasBudget::new();
let mut debit = bucket
.try_reserve_at(1, 100, 80, &mut budget)
.unwrap()
.unwrap();
assert_eq!(debit.commit_all(), 80);
assert_eq!(bucket.reserve_cas_attempts(), RESERVE_CAS_ATTEMPT_LIMIT as u64);
assert!(budget.is_exhausted());
assert_eq!(bucket.used_at(1), Some(80));
}
#[test]
fn refund_stops_after_the_attempt_limit() {
let bucket = Arc::new(DirectionBucket::default());
let debit = reserve_at(&bucket, 1, 100, 80).unwrap().unwrap();
bucket.force_refund_failures(REFUND_CAS_ATTEMPT_LIMIT);
drop(debit);
assert_eq!(
bucket.refund_cas_attempts(),
REFUND_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(bucket.used_at(1), Some(80));
assert_eq!(
bucket
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
1
);
let mut tail = reserve_at(&bucket, 1, 100, 100).unwrap().unwrap();
assert_eq!(tail.commit_all(), 20);
assert!(reserve_at(&bucket, 1, 100, 1).unwrap().is_none());
assert_eq!(bucket.used_at(1), Some(100));
let mut next = reserve_at(&bucket, 2, 100, 100).unwrap().unwrap();
assert_eq!(next.commit_all(), 100);
}
#[test]
fn refund_succeeds_on_the_last_allowed_attempt() {
let bucket = Arc::new(DirectionBucket::default());
let debit = reserve_at(&bucket, 1, 100, 80).unwrap().unwrap();
bucket.force_refund_failures(REFUND_CAS_ATTEMPT_LIMIT - 1);
drop(debit);
assert_eq!(
bucket.refund_cas_attempts(),
REFUND_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(bucket.used_at(1), Some(0));
assert_eq!(
bucket
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
0
);
}
#[test]
fn partial_refund_exhaustion_retains_the_full_charge() {
let bucket = Arc::new(DirectionBucket::default());
let mut debit = reserve_at(&bucket, 1, 100, 80).unwrap().unwrap();
bucket.force_refund_failures(REFUND_CAS_ATTEMPT_LIMIT);
debit.settle(30);
drop(debit);
assert_eq!(bucket.used_at(1), Some(80));
assert_eq!(
bucket.refund_cas_attempts(),
REFUND_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(
bucket
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
1
);
}
#[test]
fn lease_contention_is_not_reported_as_throttling() {
let limiter = TrafficLimiter::new();
let mut user_limits = HashMap::new();
user_limits.insert("alice".to_string(), rate(400_000, 400_000));
limiter.apply_policy(user_limits, HashMap::new());
let lease = limiter
.acquire_lease("alice", "203.0.113.7".parse().unwrap())
.unwrap();
let bucket = Arc::clone(
&lease
.binding
.load_full()
.user_bucket
.as_ref()
.unwrap()
.down,
);
bucket.force_reserve_failures(RESERVE_CAS_ATTEMPT_LIMIT);
let result = lease.try_consume(RateDirection::Down, 1);
let metrics = limiter.metrics_snapshot();
assert_eq!(result.granted, 0);
assert!(!result.blocked_user);
assert!(!result.blocked_cidr);
assert_eq!(metrics.user_reserve_cas_retry_exhausted_down_total, 1);
assert_eq!(metrics.user_throttle_down_total, 0);
}
#[test]
fn cidr_contention_rolls_back_provisional_user_debits() {
let limiter = TrafficLimiter::new();
let mut user_limits = HashMap::new();
user_limits.insert("alice".to_string(), rate(400_000, 400_000));
let mut cidr_limits = HashMap::new();
cidr_limits.insert(
CidrRateLimitKey::Network("203.0.113.0/24".parse().unwrap()),
rate(400_000, 400_000),
);
limiter.apply_policy(user_limits, cidr_limits);
let lease = limiter
.acquire_lease("alice", "203.0.113.7".parse().unwrap())
.unwrap();
let binding = lease.binding.load_full();
let cidr_bucket = binding.cidr_bucket.as_ref().unwrap();
cidr_bucket
.down
.used
.force_reserve_failures(RESERVE_CAS_ATTEMPT_LIMIT);
let reservation = lease.try_reserve(RateDirection::Down, 800);
let result = reservation.result();
let epoch = reservation.user.as_ref().unwrap().epoch;
let user_attempts = binding
.user_bucket
.as_ref()
.unwrap()
.down
.reserve_cas_attempts();
let active_user_attempts = cidr_bucket.down.active_users.reserve_cas_attempts();
let cidr_user_attempts = binding
.cidr_user_share
.as_ref()
.unwrap()
.down
.used
.reserve_cas_attempts();
let aggregate_attempts = cidr_bucket.down.used.reserve_cas_attempts();
drop(reservation);
assert_eq!(result.granted, 0);
assert!(!result.blocked_user);
assert!(!result.blocked_cidr);
assert_eq!(
binding
.user_bucket
.as_ref()
.unwrap()
.down
.used_at(epoch),
Some(0)
);
assert_eq!(cidr_bucket.down.used.used_at(epoch), None);
assert_eq!(
user_attempts + active_user_attempts + cidr_user_attempts + aggregate_attempts,
RESERVE_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(
binding
.cidr_user_share
.as_ref()
.unwrap()
.down
.used
.used_at(epoch),
Some(0)
);
let metrics = limiter.metrics_snapshot();
assert_eq!(metrics.cidr_reserve_cas_retry_exhausted_down_total, 1);
assert_eq!(metrics.cidr_throttle_down_total, 0);
}
#[test]
fn contention_snapshot_preserves_scope_direction_and_operation() {
let limiter = TrafficLimiter::new();
limiter
.user_scope
.contention_up
.reserve_exhausted_total
.store(1, Ordering::Relaxed);
limiter
.user_scope
.contention_down
.reserve_exhausted_total
.store(2, Ordering::Relaxed);
limiter
.user_scope
.contention_up
.refund_exhausted_total
.store(3, Ordering::Relaxed);
limiter
.user_scope
.contention_down
.refund_exhausted_total
.store(4, Ordering::Relaxed);
limiter
.cidr_scope
.contention_up
.reserve_exhausted_total
.store(5, Ordering::Relaxed);
limiter
.cidr_scope
.contention_down
.reserve_exhausted_total
.store(6, Ordering::Relaxed);
limiter
.cidr_scope
.contention_up
.refund_exhausted_total
.store(7, Ordering::Relaxed);
limiter
.cidr_scope
.contention_down
.refund_exhausted_total
.store(8, Ordering::Relaxed);
let snapshot = limiter.metrics_snapshot();
assert_eq!(snapshot.user_reserve_cas_retry_exhausted_up_total, 1);
assert_eq!(snapshot.user_reserve_cas_retry_exhausted_down_total, 2);
assert_eq!(snapshot.user_refund_cas_retry_exhausted_up_total, 3);
assert_eq!(snapshot.user_refund_cas_retry_exhausted_down_total, 4);
assert_eq!(snapshot.cidr_reserve_cas_retry_exhausted_up_total, 5);
assert_eq!(snapshot.cidr_reserve_cas_retry_exhausted_down_total, 6);
assert_eq!(snapshot.cidr_refund_cas_retry_exhausted_up_total, 7);
assert_eq!(snapshot.cidr_refund_cas_retry_exhausted_down_total, 8);
}
#[test]
fn cidr_activation_consumes_one_shared_attempt_budget() {
let bucket = CidrDirectionBucket::default();
let user = CidrUserDirectionState::default();
user.used
.force_reserve_failures(RESERVE_CAS_ATTEMPT_LIMIT);
let mut budget = ReserveCasBudget::new();
let activation = user.ensure_active(13, &bucket.active_users, &mut budget);
assert_eq!(activation, Err(BucketReserveError::Contended));
assert_eq!(
user.used.reserve_cas_attempts() + bucket.active_users.reserve_cas_attempts(),
RESERVE_CAS_ATTEMPT_LIMIT as u64
);
assert_eq!(bucket.active_users.used_at(13), Some(0));
assert_eq!(user.used.used_at(13), None);
assert_eq!(
bucket
.active_users
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
0
);
}
#[test]
fn cidr_activation_reports_refund_contention_without_reserve_exhaustion() {
let bucket = CidrDirectionBucket::default();
let user = CidrUserDirectionState::default();
user.used.force_reserve_failures(1);
bucket
.active_users
.force_refund_failures(REFUND_CAS_ATTEMPT_LIMIT);
let mut budget = ReserveCasBudget::new();
let activation = user.ensure_active(13, &bucket.active_users, &mut budget);
assert_eq!(activation, Err(BucketReserveError::RefundContended));
assert!(!budget.is_exhausted());
assert!(!activation.unwrap_err().exhausted_reserve_budget());
assert_eq!(bucket.active_users.used_at(13), Some(1));
assert_eq!(user.used.used_at(13), None);
assert_eq!(
bucket
.active_users
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
1
);
}
#[test]
fn cidr_activation_preserves_combined_contention_cause() {
let bucket = CidrDirectionBucket::default();
let user = CidrUserDirectionState::default();
bucket
.active_users
.force_refund_failures(REFUND_CAS_ATTEMPT_LIMIT);
let mut budget = ReserveCasBudget { remaining: 1 };
let activation = user.ensure_active(13, &bucket.active_users, &mut budget);
assert_eq!(
activation,
Err(BucketReserveError::ReserveAndRefundContended)
);
assert!(budget.is_exhausted());
assert!(activation.unwrap_err().exhausted_reserve_budget());
assert_eq!(bucket.active_users.used_at(13), Some(1));
assert_eq!(user.used.used_at(13), None);
assert_eq!(
bucket
.active_users
.contention
.refund_exhausted_total
.load(Ordering::Relaxed),
1
);
}
#[test]
fn cidr_first_grants_preserve_the_current_soft_fair_share() {
let bucket = CidrDirectionBucket::default();
let first = CidrUserDirectionState::default();
let second = CidrUserDirectionState::default();
assert_eq!(
first.ensure_active(
17,
&bucket.active_users,
&mut ReserveCasBudget::new(),
),
Ok(true)
);
assert_eq!(
second.ensure_active(
17,
&bucket.active_users,
&mut ReserveCasBudget::new(),
),
Ok(true)
);
let mut first_reservation = bucket
.try_reserve(&first, 17, 100, 100, &mut ReserveCasBudget::new())
.unwrap();
let mut second_reservation = bucket
.try_reserve(&second, 17, 100, 100, &mut ReserveCasBudget::new())
.unwrap();
assert_eq!(first_reservation.granted, 50);
assert_eq!(second_reservation.granted, 50);
assert_eq!(bucket.active_users.used_at(17), Some(2));
first_reservation
.aggregate_debit
.as_mut()
.unwrap()
.commit_all();
first_reservation
.user_debit
.as_mut()
.unwrap()
.commit_all();
second_reservation
.aggregate_debit
.as_mut()
.unwrap()
.commit_all();
second_reservation
.user_debit
.as_mut()
.unwrap()
.commit_all();
assert_eq!(bucket.used.used_at(17), Some(100));
assert_eq!(first.used.used_at(17), Some(50));
assert_eq!(second.used.used_at(17), Some(50));
}
#[test]
fn concurrent_commit_and_refund_preserve_fail_closed_accounting() {
const WORKERS: usize = 32;
const OPERATIONS_PER_WORKER: usize = 1_024;
const EPOCH: u64 = 19;
const CAP: u64 = (WORKERS * OPERATIONS_PER_WORKER) as u64;
let bucket = Arc::new(DirectionBucket::default());
let committed = Arc::new(AtomicU64::new(0));
let barrier = Arc::new(std::sync::Barrier::new(WORKERS));
let mut threads = Vec::with_capacity(WORKERS);
for worker in 0..WORKERS {
let bucket = Arc::clone(&bucket);
let committed = Arc::clone(&committed);
let barrier = Arc::clone(&barrier);
threads.push(std::thread::spawn(move || {
barrier.wait();
for operation in 0..OPERATIONS_PER_WORKER {
match reserve_at(&bucket, EPOCH, CAP, 1) {
Ok(Some(mut debit)) if (worker + operation) % 2 == 0 => {
committed.fetch_add(debit.commit_all(), Ordering::Relaxed);
}
Ok(Some(debit)) => drop(debit),
Err(BucketReserveError::Contended) => {}
Ok(None) => panic!("capacity exhausted before every operation ran"),
Err(error) => panic!("unexpected reservation error: {error:?}"),
}
}
}));
}
for thread in threads {
thread.join().unwrap();
}
let committed = committed.load(Ordering::Relaxed);
let leaked = bucket
.contention
.refund_exhausted_total
.load(Ordering::Relaxed);
let used = bucket.used_at(EPOCH).unwrap();
assert_eq!(used, committed + leaked);
assert!(used <= CAP);
}
proptest! {
#[test]
fn sequential_debit_lifecycle_matches_the_epoch_model(
operations in prop::collection::vec((0u8..4, 1u64..128), 1..128),
) {
const CAP: u64 = 512;
let bucket = Arc::new(DirectionBucket::default());
let mut epoch = 1u64;
let mut model_used = 0u64;
for (operation, requested) in operations {
if operation == 3 {
epoch += 1;
model_used = 0;
let debit = reserve_at(&bucket, epoch, CAP, 1).unwrap().unwrap();
drop(debit);
prop_assert_eq!(bucket.used_at(epoch), Some(0));
continue;
}
let reservation = reserve_at(&bucket, epoch, CAP, requested).unwrap();
let expected_grant = requested.min(CAP.saturating_sub(model_used));
if expected_grant == 0 {
prop_assert!(reservation.is_none());
continue;
}
let mut debit = reservation.unwrap();
prop_assert_eq!(debit.refundable, expected_grant);
match operation {
0 => {
model_used += debit.commit_all();
}
1 => {
let committed = expected_grant / 2;
debit.settle(committed);
model_used += committed;
}
2 => drop(debit),
_ => unreachable!(),
}
prop_assert_eq!(bucket.used_at(epoch), Some(model_used));
}
}
}