diff --git a/src/metrics/render/traffic.rs b/src/metrics/render/traffic.rs index 29636f6..057a0d9 100644 --- a/src/metrics/render/traffic.rs +++ b/src/metrics/render/traffic.rs @@ -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" diff --git a/src/metrics/tests.rs b/src/metrics/tests.rs index 9d9d830..2b57ed6 100644 --- a/src/metrics/tests.rs +++ b/src/metrics/tests.rs @@ -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] diff --git a/src/proxy/traffic_limiter.rs b/src/proxy/traffic_limiter.rs index 51e2b6e..766c734 100644 --- a/src/proxy/traffic_limiter.rs +++ b/src/proxy/traffic_limiter.rs @@ -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, + contention_down: Arc, active_leases: AtomicU64, policy_entries: AtomicU64, } @@ -86,6 +132,54 @@ struct AtomicRatePair { #[derive(Default)] struct DirectionBucket { state: AtomicU64, + contention: Arc, + #[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, +} + +struct CidrReservation { + granted: u64, + aggregate_debit: Option, + user_debit: Option, } struct UserBucket { @@ -95,13 +189,11 @@ struct UserBucket { active_leases: AtomicU64, } -#[derive(Default)] struct CidrDirectionBucket { used: Arc, active_users: Arc, } -#[derive(Default)] struct CidrUserDirectionState { used: Arc, } @@ -165,6 +257,7 @@ struct TrafficLeaseBinding { cidr_user_share: Option>, } +/// A live traffic-limiter binding for one authenticated client session. pub struct TrafficLease { limiter: Arc, 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, policy_update: ParkingMutex<()>, diff --git a/src/proxy/traffic_limiter/buckets.rs b/src/proxy/traffic_limiter/buckets.rs index 8213151..7755283 100644 --- a/src/proxy/traffic_limiter/buckets.rs +++ b/src/proxy/traffic_limiter/buckets.rs @@ -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) -> 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 { + #[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 { + #[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 { 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 { + 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, epoch: u64, cap: u64, requested: u64, - ) -> Option { - if requested == 0 || cap == 0 || epoch > PACKED_EPOCH_MAX { - return None; + budget: &mut ReserveCasBudget, + ) -> Result, 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, + down_contention: Arc, + ) -> 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) { + budget: &mut ReserveCasBudget, + ) -> Result { 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) -> 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, Option) { + budget: &mut ReserveCasBudget, + ) -> Result { 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) -> bool { + fn new(contention: Arc) -> Self { + Self { + used: Arc::new(DirectionBucket::new(contention)), + } + } + + pub(super) fn ensure_active( + &self, + epoch: u64, + active_users: &Arc, + budget: &mut ReserveCasBudget, + ) -> Result { 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, + down_contention: Arc, + ) -> 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, + down_contention: Arc, + ) -> 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 { - 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) { @@ -344,16 +518,28 @@ impl CidrBucket { &self, direction: RateDirection, share: &CidrUserShare, + epoch: u64, requested: u64, - ) -> (u64, Option, Option) { + budget: &mut ReserveCasBudget, + ) -> Result { 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) + } } } diff --git a/src/proxy/traffic_limiter/helpers.rs b/src/proxy/traffic_limiter/helpers.rs index a25eada..c869d85 100644 --- a/src/proxy/traffic_limiter/helpers.rs +++ b/src/proxy/traffic_limiter/helpers.rs @@ -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; diff --git a/src/proxy/traffic_limiter/lease.rs b/src/proxy/traffic_limiter/lease.rs index 02ea7f0..01ac375 100644 --- a/src/proxy/traffic_limiter/lease.rs +++ b/src/proxy/traffic_limiter/lease.rs @@ -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, diff --git a/src/proxy/traffic_limiter/limiter.rs b/src/proxy/traffic_limiter/limiter.rs index 6c55b89..e805791 100644 --- a/src/proxy/traffic_limiter/limiter.rs +++ b/src/proxy/traffic_limiter/limiter.rs @@ -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 { + 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 { 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, @@ -125,6 +170,7 @@ impl TrafficLimiter { true } + /// Creates a lease that follows policy revisions for one client identity. pub fn acquire_lease( self: &Arc, 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, diff --git a/src/proxy/traffic_limiter/tests.rs b/src/proxy/traffic_limiter/tests.rs index 975dff8..847c639 100644 --- a/src/proxy/traffic_limiter/tests.rs +++ b/src/proxy/traffic_limiter/tests.rs @@ -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, + epoch: u64, + cap: u64, + requested: u64, +) -> Result, 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] diff --git a/src/proxy/traffic_limiter/tests/bucket_contention.rs b/src/proxy/traffic_limiter/tests/bucket_contention.rs new file mode 100644 index 0000000..dd96a36 --- /dev/null +++ b/src/proxy/traffic_limiter/tests/bucket_contention.rs @@ -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)); + } + } +}