use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use nmbrs_metrics::controls::ControlApplier;
use crate::RateLimiter;
use crate::spec::RateSpec;
pub struct RateLimiterApplier {
limiter: Arc<RateLimiter>,
}
impl RateLimiterApplier {
pub fn new(limiter: Arc<RateLimiter>) -> Self {
Self { limiter }
}
}
impl ControlApplier<RateSpec> for RateLimiterApplier {
fn apply(
&self,
value: RateSpec,
) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
let limiter = self.limiter.clone();
Box::pin(async move { limiter.reconfigure(value) })
}
}
#[cfg(test)]
mod tests {
use super::*;
use nmbrs_metrics::controls::{ControlBuilder, ControlOrigin, SetError};
use std::time::{Duration, Instant};
#[tokio::test]
async fn applier_reconfigures_limiter_through_control() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(100.0)));
let control: nmbrs_metrics::controls::Control<RateSpec> =
ControlBuilder::new("rate", RateSpec::new(100.0))
.validator(|spec| {
if spec.ops_per_sec <= 0.0 {
Err(format!("rate must be > 0, got {}", spec.ops_per_sec))
} else {
Ok(())
}
})
.build();
control.register_applier(RateLimiterApplier::new(limiter.clone()));
assert_eq!(limiter.rate(), 100.0);
control
.set(RateSpec::new(5_000.0), ControlOrigin::Test)
.await
.expect("reconfigure should succeed");
assert_eq!(limiter.rate(), 5_000.0);
drop(control);
}
#[tokio::test]
async fn validator_rejection_leaves_rate_unchanged() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(250.0)));
let control: nmbrs_metrics::controls::Control<RateSpec> =
ControlBuilder::new("rate", RateSpec::new(250.0))
.validator(|spec| {
if spec.ops_per_sec > 0.0 && spec.ops_per_sec <= 10_000.0 {
Ok(())
} else {
Err("rate out of range (0, 10_000]".into())
}
})
.build();
control.register_applier(RateLimiterApplier::new(limiter.clone()));
match control
.set(RateSpec::new(100_000.0), ControlOrigin::Test)
.await
{
Err(SetError::ValidationFailed(_)) => {}
other => panic!("expected validation rejection, got {other:?}"),
}
assert_eq!(limiter.rate(), 250.0);
}
#[tokio::test]
async fn live_reconfigure_changes_throughput_observed_by_acquire() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(200.0)));
let control: nmbrs_metrics::controls::Control<RateSpec> =
ControlBuilder::new("rate", RateSpec::new(200.0)).build();
control.register_applier(RateLimiterApplier::new(limiter.clone()));
tokio::time::sleep(Duration::from_millis(30)).await;
let slow_start = Instant::now();
for _ in 0..10 {
limiter.acquire().await;
}
let slow_elapsed = slow_start.elapsed();
control
.set(RateSpec::new(100_000.0), ControlOrigin::Test)
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(30)).await;
let fast_start = Instant::now();
for _ in 0..10 {
limiter.acquire().await;
}
let fast_elapsed = fast_start.elapsed();
assert!(
fast_elapsed < slow_elapsed,
"after live reconfigure the limiter should drain faster: \
slow={slow_elapsed:?} fast={fast_elapsed:?}",
);
}
#[tokio::test]
async fn reconfigure_preserves_in_flight_backlog_field() {
let limiter = Arc::new(RateLimiter::start(RateSpec::new(10.0))); let control: nmbrs_metrics::controls::Control<RateSpec> =
ControlBuilder::new("rate", RateSpec::new(10.0)).build();
control.register_applier(RateLimiterApplier::new(limiter.clone()));
tokio::time::sleep(Duration::from_millis(60)).await;
let backlog_before = limiter.wait_time_nanos();
control
.set(RateSpec::new(1_000.0), ControlOrigin::Test)
.await
.unwrap();
let _backlog_after = limiter.wait_time_nanos();
assert!(limiter.rate() == 1_000.0);
let _ = backlog_before;
}
}