use super::bucket::{BucketOutcome, RateBucket};
use super::clock::AdmissionClock;
use super::metrics::{AdmissionRefusal, RefusalCounters};
use super::{AdmissionContext, AdmissionDecision};
use crate::memory_limiter::{
LocalReceiverAdmissionState, MemoryPressureLevel, SharedReceiverAdmissionState,
};
use otel_arrow_dfe_config::policy::{RateLimitEnforcement, RateLimitPressure, RateLimiterPolicy};
use std::rc::Rc;
use std::sync::Arc;
pub(crate) trait PressureSource {
fn level(&self) -> MemoryPressureLevel;
}
impl PressureSource for SharedReceiverAdmissionState {
#[inline]
fn level(&self) -> MemoryPressureLevel {
SharedReceiverAdmissionState::level(self)
}
}
impl PressureSource for LocalReceiverAdmissionState {
#[inline]
fn level(&self) -> MemoryPressureLevel {
LocalReceiverAdmissionState::level(self)
}
}
const NANOS_PER_SECOND: u64 = 1_000_000_000;
fn retry_after_secs(retry_after_nanos: u64) -> u32 {
let seconds = retry_after_nanos.div_ceil(NANOS_PER_SECOND).max(1);
u32::try_from(seconds).unwrap_or(u32::MAX)
}
#[derive(Debug)]
pub(crate) struct RateGate<P> {
policy: RateLimiterPolicy,
pressure: P,
bucket: RateBucket,
counters: Arc<RefusalCounters>,
}
impl<P: PressureSource> RateGate<P> {
pub(crate) fn new(
policy: RateLimiterPolicy,
pressure: P,
clock: AdmissionClock,
counters: Arc<RefusalCounters>,
) -> Self {
Self {
bucket: RateBucket::with_clock(&policy, clock),
policy,
pressure,
counters,
}
}
#[inline]
fn pressure_active(&self) -> bool {
let level = P::level(&self.pressure);
match self.policy.pressure {
RateLimitPressure::Soft => {
matches!(level, MemoryPressureLevel::Soft | MemoryPressureLevel::Hard)
}
}
}
#[inline]
fn decide(&self, units: u64) -> AdmissionDecision {
let pressure_active = self.pressure_active();
let enforce = pressure_active && self.policy.enforcement == RateLimitEnforcement::Enforce;
let outcome = if enforce {
self.bucket.check_units(units)
} else {
self.bucket.observe_units(units)
};
match outcome {
BucketOutcome::WithinLimit => AdmissionDecision::Admit,
_ if !pressure_active => AdmissionDecision::Admit,
BucketOutcome::OverLimit { .. } | BucketOutcome::Oversized
if self.policy.enforcement == RateLimitEnforcement::ObserveOnly =>
{
self.counters.record(AdmissionRefusal::WouldThrottle);
AdmissionDecision::WouldThrottle
}
BucketOutcome::OverLimit { retry_after_nanos } => {
self.counters.record(AdmissionRefusal::Throttle);
AdmissionDecision::Throttle {
retry_after_secs: retry_after_secs(retry_after_nanos),
}
}
BucketOutcome::Oversized => {
self.counters.record(AdmissionRefusal::Oversized);
AdmissionDecision::Oversized
}
}
}
#[inline]
fn saturated(&self) -> bool {
if !self.pressure_active() || self.policy.enforcement != RateLimitEnforcement::Enforce {
return false;
}
self.bucket.is_exhausted()
}
fn record_instance_saturation_refusal(&self) {
self.counters.record(AdmissionRefusal::Throttle);
}
}
#[derive(Clone, Debug)]
pub struct LocalAdmissionGate {
inner: Rc<RateGate<LocalReceiverAdmissionState>>,
}
impl LocalAdmissionGate {
pub(crate) fn new(inner: RateGate<LocalReceiverAdmissionState>) -> Self {
Self {
inner: Rc::new(inner),
}
}
#[inline]
#[must_use]
pub fn admit(&self, units: u64, _context: AdmissionContext<'_>) -> AdmissionDecision {
self.inner.decide(units)
}
#[inline]
#[must_use]
pub fn is_instance_saturated(&self) -> bool {
self.inner.saturated()
}
}
#[derive(Clone, Debug)]
pub struct SharedAdmissionGate {
inner: Arc<RateGate<SharedReceiverAdmissionState>>,
}
impl SharedAdmissionGate {
pub(crate) fn new(inner: RateGate<SharedReceiverAdmissionState>) -> Self {
Self {
inner: Arc::new(inner),
}
}
#[inline]
#[must_use]
pub fn admit(&self, units: u64, _context: AdmissionContext<'_>) -> AdmissionDecision {
self.inner.decide(units)
}
#[inline]
#[must_use]
pub fn is_instance_saturated(&self) -> bool {
self.inner.saturated()
}
#[must_use]
pub fn refuse_if_instance_saturated(&self) -> bool {
if !self.inner.saturated() {
return false;
}
self.inner.record_instance_saturation_refusal();
true
}
pub fn record_probed_instance_saturation_refusal(&self) {
self.inner.record_instance_saturation_refusal()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::admission::clock::ManualClock;
use crate::memory_limiter::{MemoryPressureChanged, MemoryPressureState};
use otel_arrow_dfe_config::policy::{RateLimitAggregation, RateLimitUnit, TokenBucketPolicy};
use std::time::Duration;
const RETRY_AFTER_SECS: u32 = 7;
fn policy(enforcement: RateLimitEnforcement) -> RateLimiterPolicy {
RateLimiterPolicy {
enforcement,
aggregation: RateLimitAggregation::ReceiverInstance,
unit: RateLimitUnit::RequestBytes,
pressure: RateLimitPressure::Soft,
token_bucket: TokenBucketPolicy {
allow: 10,
interval: Duration::from_secs(1),
burst: Some(10),
},
}
}
struct Harness<P> {
gate: RateGate<P>,
pressure: P,
clock: Arc<ManualClock>,
counters: Arc<RefusalCounters>,
}
impl<P: PressureSource + Clone> Harness<P> {
fn new(enforcement: RateLimitEnforcement, pressure: P) -> Self {
let clock = Arc::new(ManualClock::new(0));
let counters = Arc::new(RefusalCounters::default());
let gate = RateGate::new(
policy(enforcement),
pressure.clone(),
AdmissionClock::Manual(Arc::clone(&clock)),
Arc::clone(&counters),
);
Self {
gate,
pressure,
clock,
counters,
}
}
}
fn shared_harness(enforcement: RateLimitEnforcement) -> Harness<SharedReceiverAdmissionState> {
Harness::new(
enforcement,
SharedReceiverAdmissionState::from_process_state(&MemoryPressureState::default()),
)
}
fn local_harness(enforcement: RateLimitEnforcement) -> Harness<LocalReceiverAdmissionState> {
Harness::new(
enforcement,
LocalReceiverAdmissionState::from_process_state(&MemoryPressureState::default()),
)
}
fn raise_pressure<P: PressureSource + ApplyPressure>(pressure: &P, level: MemoryPressureLevel) {
pressure.apply_change(MemoryPressureChanged {
generation: 1,
level,
retry_after_secs: RETRY_AFTER_SECS,
usage_bytes: 1,
});
}
trait ApplyPressure {
fn apply_change(&self, change: MemoryPressureChanged);
}
impl ApplyPressure for SharedReceiverAdmissionState {
fn apply_change(&self, change: MemoryPressureChanged) {
self.apply(change);
}
}
impl ApplyPressure for LocalReceiverAdmissionState {
fn apply_change(&self, change: MemoryPressureChanged) {
self.apply(change);
}
}
fn admit_shared(
gate: &RateGate<SharedReceiverAdmissionState>,
units: u64,
) -> AdmissionDecision {
gate.decide(units)
}
fn admit_local(gate: &RateGate<LocalReceiverAdmissionState>, units: u64) -> AdmissionDecision {
gate.decide(units)
}
#[test]
fn over_limit_traffic_is_admitted_without_pressure() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
for _ in 0..100 {
assert_eq!(admit_shared(&harness.gate, 1), AdmissionDecision::Admit);
}
assert!(!harness.gate.saturated());
assert_eq!(harness.counters.drain(), [0, 0, 0]);
}
#[test]
fn enforcement_starts_from_observed_history() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
for _ in 0..100 {
let _ = admit_shared(&harness.gate, 1);
}
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
assert!(matches!(
admit_shared(&harness.gate, 1),
AdmissionDecision::Throttle { .. }
));
}
#[test]
fn throttling_reports_retry_hint_and_recovers() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
for _ in 0..10 {
assert_eq!(admit_shared(&harness.gate, 1), AdmissionDecision::Admit);
}
let decision = admit_shared(&harness.gate, 1);
assert_eq!(
decision,
AdmissionDecision::Throttle {
retry_after_secs: 1
}
);
assert!(harness.gate.saturated());
harness.clock.advance(100_000_000);
assert_eq!(admit_shared(&harness.gate, 1), AdmissionDecision::Admit);
}
#[test]
fn observe_only_reports_without_refusing() {
let harness = shared_harness(RateLimitEnforcement::ObserveOnly);
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
for _ in 0..10 {
assert_eq!(admit_shared(&harness.gate, 1), AdmissionDecision::Admit);
}
assert_eq!(
admit_shared(&harness.gate, 1),
AdmissionDecision::WouldThrottle
);
assert!(!harness.gate.saturated());
assert_eq!(harness.counters.drain(), [1, 0, 0]);
}
#[test]
fn oversized_request_is_distinguished_from_throttling() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
assert_eq!(
admit_shared(&harness.gate, 11),
AdmissionDecision::Oversized
);
assert_eq!(harness.counters.drain(), [0, 0, 1]);
}
#[test]
fn only_refusals_are_staged_for_telemetry() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
for _ in 0..10 {
let _ = admit_shared(&harness.gate, 1);
}
assert_eq!(harness.counters.drain(), [0, 0, 0]);
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
let _ = admit_shared(&harness.gate, 1);
let _ = admit_shared(&harness.gate, 11);
assert_eq!(harness.counters.drain(), [0, 1, 1]);
}
#[test]
fn probed_saturation_refusal_is_staged_once() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
for _ in 0..10 {
let _ = admit_shared(&harness.gate, 1);
}
assert!(harness.gate.saturated());
harness.gate.record_instance_saturation_refusal();
assert_eq!(harness.counters.drain(), [0, 1, 0]);
}
#[test]
fn each_bind_gets_independent_capacity() {
let first = shared_harness(RateLimitEnforcement::Enforce);
let second = shared_harness(RateLimitEnforcement::Enforce);
raise_pressure(&first.pressure, MemoryPressureLevel::Soft);
raise_pressure(&second.pressure, MemoryPressureLevel::Soft);
for _ in 0..10 {
let _ = admit_shared(&first.gate, 1);
}
assert!(matches!(
admit_shared(&first.gate, 1),
AdmissionDecision::Throttle { .. }
));
assert_eq!(admit_shared(&second.gate, 1), AdmissionDecision::Admit);
}
#[test]
fn hard_pressure_activates_a_soft_threshold_policy() {
let harness = shared_harness(RateLimitEnforcement::Enforce);
raise_pressure(&harness.pressure, MemoryPressureLevel::Hard);
for _ in 0..10 {
let _ = admit_shared(&harness.gate, 1);
}
assert!(matches!(
admit_shared(&harness.gate, 1),
AdmissionDecision::Throttle { .. }
));
}
#[test]
fn local_gate_matches_shared_gate_semantics() {
let harness = local_harness(RateLimitEnforcement::Enforce);
for _ in 0..100 {
assert_eq!(admit_local(&harness.gate, 1), AdmissionDecision::Admit);
}
raise_pressure(&harness.pressure, MemoryPressureLevel::Soft);
assert_eq!(
admit_local(&harness.gate, 1),
AdmissionDecision::Throttle {
retry_after_secs: 2
}
);
assert!(harness.gate.saturated());
assert_eq!(admit_local(&harness.gate, 11), AdmissionDecision::Oversized);
harness.clock.advance(2_000_000_000);
assert_eq!(admit_local(&harness.gate, 1), AdmissionDecision::Admit);
assert_eq!(harness.counters.drain(), [0, 1, 1]);
}
}