use super::clock::AdmissionClock;
use super::gate::RateGate;
use super::metrics::{AdmissionMetricsHandle, RefusalCounters};
use super::{
AdmissionBindError, AdmissionDimension, AdmissionDimensionSet, LocalAdmissionGate,
SharedAdmissionGate,
};
use crate::memory_limiter::{LocalReceiverAdmissionState, SharedReceiverAdmissionState};
use otel_arrow_dfe_config::policy::{
RateLimitUnit, RateLimiterDeclarationScope, RateLimiterPolicy,
};
use std::sync::Arc;
use std::sync::atomic::{AtomicU8, Ordering};
const UNBOUND: u8 = 0;
const BOUND: u8 = 1;
#[derive(Debug)]
struct BindingInner {
limiter_name: Arc<str>,
declaration_scope: Option<RateLimiterDeclarationScope>,
policy: RateLimiterPolicy,
state: AtomicU8,
clock: AdmissionClock,
counters: Arc<RefusalCounters>,
}
#[derive(Clone, Debug, Default)]
pub struct AdmissionBinder {
inner: Option<Arc<BindingInner>>,
}
impl AdmissionBinder {
#[must_use]
pub const fn none() -> Self {
Self { inner: None }
}
#[must_use]
pub fn configured(limiter_name: impl Into<Arc<str>>, policy: RateLimiterPolicy) -> Self {
Self::configured_at_scope(limiter_name, None, policy)
}
#[must_use]
pub fn configured_at_scope(
limiter_name: impl Into<Arc<str>>,
declaration_scope: Option<RateLimiterDeclarationScope>,
policy: RateLimiterPolicy,
) -> Self {
Self::with_clock(
limiter_name,
declaration_scope,
policy,
AdmissionClock::system(),
)
}
pub(crate) fn with_clock(
limiter_name: impl Into<Arc<str>>,
declaration_scope: Option<RateLimiterDeclarationScope>,
policy: RateLimiterPolicy,
clock: AdmissionClock,
) -> Self {
Self {
inner: Some(Arc::new(BindingInner {
limiter_name: limiter_name.into(),
declaration_scope,
policy,
state: AtomicU8::new(UNBOUND),
clock,
counters: Arc::new(RefusalCounters::default()),
})),
}
}
#[must_use]
pub fn is_configured(&self) -> bool {
self.inner.is_some()
}
pub fn bind_local(
&self,
dimension: AdmissionDimension,
pressure: LocalReceiverAdmissionState,
) -> Result<Option<LocalAdmissionGate>, AdmissionBindError> {
let Some(inner) = &self.inner else {
return Ok(None);
};
if !inner.validate_dimension(dimension)? {
return Ok(None);
}
inner.claim()?;
Ok(Some(LocalAdmissionGate::new(RateGate::new(
inner.policy,
pressure,
inner.clock.clone(),
Arc::clone(&inner.counters),
))))
}
pub fn bind_shared(
&self,
dimension: AdmissionDimension,
pressure: SharedReceiverAdmissionState,
) -> Result<Option<SharedAdmissionGate>, AdmissionBindError> {
let Some(inner) = &self.inner else {
return Ok(None);
};
if !inner.validate_dimension(dimension)? {
return Ok(None);
}
inner.claim()?;
Ok(Some(SharedAdmissionGate::new(RateGate::new(
inner.policy,
pressure,
inner.clock.clone(),
Arc::clone(&inner.counters),
))))
}
pub fn validate_factory_consumption(
&self,
component_urn: &str,
) -> Result<(), AdmissionBindError> {
let Some(inner) = &self.inner else {
return Ok(());
};
if inner.state.load(Ordering::Acquire) != BOUND {
return Err(AdmissionBindError::ExplicitBindingNotConsumed {
limiter: inner.limiter_name.to_string(),
component: component_urn.to_owned(),
});
}
Ok(())
}
#[must_use]
pub fn was_bound(&self) -> bool {
self.inner
.as_ref()
.is_some_and(|inner| inner.state.load(Ordering::Acquire) == BOUND)
}
#[must_use]
pub fn limiter_name(&self) -> Option<&str> {
self.inner.as_ref().map(|inner| inner.limiter_name.as_ref())
}
#[must_use]
pub fn declaration_scope(&self) -> Option<&RateLimiterDeclarationScope> {
self.inner
.as_ref()
.and_then(|inner| inner.declaration_scope.as_ref())
}
pub(crate) fn metrics_handle(
&self,
pipeline_ctx: &crate::PipelineContext,
) -> Option<AdmissionMetricsHandle> {
let inner = self.inner.as_ref()?;
(inner.state.load(Ordering::Acquire) == BOUND).then(|| {
AdmissionMetricsHandle::new(
pipeline_ctx,
AdmissionDimension::from(inner.policy.unit),
Arc::clone(&inner.counters),
)
})
}
}
impl BindingInner {
fn validate_dimension(
&self,
requested: AdmissionDimension,
) -> Result<bool, AdmissionBindError> {
let supported = AdmissionDimension::from(self.policy.unit);
if requested == supported {
return Ok(true);
}
Err(AdmissionBindError::UnsupportedDimension {
requested,
supported: AdmissionDimensionSet::single(supported),
})
}
fn claim(&self) -> Result<(), AdmissionBindError> {
self.state
.compare_exchange(UNBOUND, BOUND, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| AdmissionBindError::AlreadyBound)
}
}
impl From<RateLimitUnit> for AdmissionDimension {
fn from(unit: RateLimitUnit) -> Self {
match unit {
RateLimitUnit::RequestBytes => Self::Bytes,
RateLimitUnit::Messages => Self::Messages,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::admission::{AdmissionContext, AdmissionDecision};
use crate::memory_limiter::{MemoryPressureState, SharedReceiverAdmissionState};
use otel_arrow_dfe_config::policy::{
RateLimitAggregation, RateLimitEnforcement, RateLimitPressure, TokenBucketPolicy,
};
use std::time::Duration;
fn policy(unit: RateLimitUnit) -> RateLimiterPolicy {
RateLimiterPolicy {
enforcement: RateLimitEnforcement::Enforce,
aggregation: RateLimitAggregation::ReceiverInstance,
unit,
pressure: RateLimitPressure::Soft,
token_bucket: TokenBucketPolicy {
allow: 1,
interval: Duration::from_secs(1),
burst: Some(1),
},
}
}
#[test]
fn binding_is_one_shot_across_context_clones() {
let binder = AdmissionBinder::configured_at_scope(
"ingress",
Some(RateLimiterDeclarationScope::Engine),
policy(RateLimitUnit::RequestBytes),
);
let clone = binder.clone();
let pressure =
SharedReceiverAdmissionState::from_process_state(&MemoryPressureState::default());
let gate = binder
.bind_shared(AdmissionDimension::Bytes, pressure.clone())
.expect("first bind")
.expect("configured gate");
assert!(matches!(
clone.bind_shared(AdmissionDimension::Bytes, pressure),
Err(AdmissionBindError::AlreadyBound)
));
assert_eq!(
gate.admit(1, AdmissionContext::EMPTY),
AdmissionDecision::Admit
);
assert_eq!(
clone.declaration_scope(),
Some(&RateLimiterDeclarationScope::Engine)
);
}
#[test]
fn shared_gate_clones_charge_one_bucket() {
let binder = AdmissionBinder::configured("ingress", policy(RateLimitUnit::RequestBytes));
let process_state = MemoryPressureState::default();
process_state.set_level_for_tests(crate::memory_limiter::MemoryPressureLevel::Soft);
let pressure = SharedReceiverAdmissionState::from_process_state(&process_state);
pressure.apply(process_state.current_update(1));
let grpc_gate = binder
.bind_shared(AdmissionDimension::Bytes, pressure)
.expect("bind shared gate")
.expect("configured gate");
let http_gate = grpc_gate.clone();
assert_eq!(
grpc_gate.admit(1, AdmissionContext::EMPTY),
AdmissionDecision::Admit
);
assert!(matches!(
http_gate.admit(1, AdmissionContext::EMPTY),
AdmissionDecision::Throttle { .. }
));
}
#[test]
fn dimension_is_validated_before_binding_is_consumed() {
let binder = AdmissionBinder::configured("ingress", policy(RateLimitUnit::Messages));
let pressure =
SharedReceiverAdmissionState::from_process_state(&MemoryPressureState::default());
assert!(matches!(
binder.bind_shared(AdmissionDimension::Bytes, pressure.clone()),
Err(AdmissionBindError::UnsupportedDimension { .. })
));
assert!(
binder
.bind_shared(AdmissionDimension::Messages, pressure)
.expect("binding remains available")
.is_some()
);
}
#[test]
fn unconsumed_binding_is_an_error() {
let binder = AdmissionBinder::configured("ingress", policy(RateLimitUnit::Messages));
assert!(matches!(
binder.validate_factory_consumption("urn:otel:receiver:host_metrics"),
Err(AdmissionBindError::ExplicitBindingNotConsumed { .. })
));
}
}