1use std::future::Future;
16use std::pin::Pin;
17use std::sync::Arc;
18
19use nmbrs_metrics::controls::ControlApplier;
20
21use crate::RateLimiter;
22use crate::spec::RateSpec;
23
24pub struct RateLimiterApplier {
29 limiter: Arc<RateLimiter>,
30}
31
32impl RateLimiterApplier {
33 pub fn new(limiter: Arc<RateLimiter>) -> Self {
34 Self { limiter }
35 }
36}
37
38impl ControlApplier<RateSpec> for RateLimiterApplier {
39 fn apply(
40 &self,
41 value: RateSpec,
42 ) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
43 let limiter = self.limiter.clone();
44 Box::pin(async move { limiter.reconfigure(value) })
45 }
46}
47
48#[cfg(test)]
49mod tests {
50 use super::*;
51 use nmbrs_metrics::controls::{ControlBuilder, ControlOrigin, SetError};
52 use std::time::{Duration, Instant};
53
54 #[tokio::test]
55 async fn applier_reconfigures_limiter_through_control() {
56 let limiter = Arc::new(RateLimiter::start(RateSpec::new(100.0)));
57 let control: nmbrs_metrics::controls::Control<RateSpec> =
58 ControlBuilder::new("rate", RateSpec::new(100.0))
59 .validator(|spec| {
60 if spec.ops_per_sec <= 0.0 {
61 Err(format!("rate must be > 0, got {}", spec.ops_per_sec))
62 } else {
63 Ok(())
64 }
65 })
66 .build();
67 control.register_applier(RateLimiterApplier::new(limiter.clone()));
68
69 assert_eq!(limiter.rate(), 100.0);
70
71 control
72 .set(RateSpec::new(5_000.0), ControlOrigin::Test)
73 .await
74 .expect("reconfigure should succeed");
75
76 assert_eq!(limiter.rate(), 5_000.0);
77
78 drop(control);
82 }
84
85 #[tokio::test]
86 async fn validator_rejection_leaves_rate_unchanged() {
87 let limiter = Arc::new(RateLimiter::start(RateSpec::new(250.0)));
88 let control: nmbrs_metrics::controls::Control<RateSpec> =
89 ControlBuilder::new("rate", RateSpec::new(250.0))
90 .validator(|spec| {
91 if spec.ops_per_sec > 0.0 && spec.ops_per_sec <= 10_000.0 {
92 Ok(())
93 } else {
94 Err("rate out of range (0, 10_000]".into())
95 }
96 })
97 .build();
98 control.register_applier(RateLimiterApplier::new(limiter.clone()));
99
100 match control
101 .set(RateSpec::new(100_000.0), ControlOrigin::Test)
102 .await
103 {
104 Err(SetError::ValidationFailed(_)) => {}
105 other => panic!("expected validation rejection, got {other:?}"),
106 }
107 assert_eq!(limiter.rate(), 250.0);
108 }
109
110 #[tokio::test]
111 async fn live_reconfigure_changes_throughput_observed_by_acquire() {
112 let limiter = Arc::new(RateLimiter::start(RateSpec::new(200.0)));
116 let control: nmbrs_metrics::controls::Control<RateSpec> =
117 ControlBuilder::new("rate", RateSpec::new(200.0)).build();
118 control.register_applier(RateLimiterApplier::new(limiter.clone()));
119
120 tokio::time::sleep(Duration::from_millis(30)).await;
121
122 let slow_start = Instant::now();
123 for _ in 0..10 {
124 limiter.acquire().await;
125 }
126 let slow_elapsed = slow_start.elapsed();
127
128 control
129 .set(RateSpec::new(100_000.0), ControlOrigin::Test)
130 .await
131 .unwrap();
132
133 tokio::time::sleep(Duration::from_millis(30)).await;
135
136 let fast_start = Instant::now();
137 for _ in 0..10 {
138 limiter.acquire().await;
139 }
140 let fast_elapsed = fast_start.elapsed();
141
142 assert!(
143 fast_elapsed < slow_elapsed,
144 "after live reconfigure the limiter should drain faster: \
145 slow={slow_elapsed:?} fast={fast_elapsed:?}",
146 );
147 }
148
149 #[tokio::test]
150 async fn reconfigure_preserves_in_flight_backlog_field() {
151 let limiter = Arc::new(RateLimiter::start(RateSpec::new(10.0))); let control: nmbrs_metrics::controls::Control<RateSpec> =
156 ControlBuilder::new("rate", RateSpec::new(10.0)).build();
157 control.register_applier(RateLimiterApplier::new(limiter.clone()));
158
159 tokio::time::sleep(Duration::from_millis(60)).await;
166 let backlog_before = limiter.wait_time_nanos();
167 control
168 .set(RateSpec::new(1_000.0), ControlOrigin::Test)
169 .await
170 .unwrap();
171 let _backlog_after = limiter.wait_time_nanos();
174 assert!(limiter.rate() == 1_000.0);
175 let _ = backlog_before;
176 }
177}