Skip to main content

nmbrs_rate/
applier.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! `Control<RateSpec>` → [`RateLimiter`] bridge.
5//!
6//! See SRD 23 (Dynamic Controls). A workload that wants its rate
7//! limiter to be live-reconfigurable declares a `Control<RateSpec>`
8//! on its owning component and registers a [`RateLimiterApplier`]
9//! against that control. When any writer (CLI, TUI, GK, API)
10//! calls `set` on the control, the applier calls
11//! [`RateLimiter::reconfigure`] and the new rate takes effect on
12//! the next `acquire`/refill cycle without restarting the refill
13//! task or dropping in-flight backlog.
14
15use 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
24/// A [`ControlApplier<RateSpec>`] that reconfigures a
25/// [`RateLimiter`] in place. Clone the limiter `Arc` before
26/// registering — the applier holds its own handle so the limiter
27/// outlives any single writer.
28pub 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        // Cleanup: stop the limiter by dropping our Arc; the
79        // original limiter goes away when the last reference
80        // does. Drop order: applier → control → limiter.
81        drop(control);
82        // limiter is the only remaining ref-holder; let it drop.
83    }
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        // Start low, pull a small batch, reconfigure to a much
113        // higher rate, pull another batch. The second batch should
114        // drain materially faster than the first.
115        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        // Give the refill loop one cycle to see the new config.
134        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        // The waiting pool is kept across a reconfigure — a writer
152        // that raises the rate shouldn't lose track of backlog
153        // that was already owed under the old rate.
154        let limiter = Arc::new(RateLimiter::start(RateSpec::new(10.0))); // very slow
155        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        // Let the refill task build up some backlog under the slow
160        // rate. 150ms at 10 ops/s → refill issues ~1 tick before
161        // the overflow logic kicks in. Hard to force backlog
162        // deterministically without a heavy acquire loop; assert
163        // the API at minimum: after reconfigure, wait_time_nanos
164        // still reads without panicking and the new rate is live.
165        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        // Read the backlog through the new (unit may differ) —
172        // the call is valid, which is the protocol guarantee.
173        let _backlog_after = limiter.wait_time_nanos();
174        assert!(limiter.rate() == 1_000.0);
175        let _ = backlog_before;
176    }
177}