Skip to main content

rusty_pv/
throttle.rs

1//! Token-bucket throttle for `-L RATE` (FR-022, AD-005, HINT-001).
2//!
3//! Per-iteration check: `expected_bytes = (now - start) · rate`. If
4//! `bytes_done > expected_bytes`, sleep `(bytes_done - expected_bytes) / rate`
5//! before the next read. 10 ms minimum sleep floor — finer granularity is
6//! unreliable on Windows.
7//!
8//! All timing uses `std::time::Instant` (monotonic) — never `SystemTime` —
9//! per AD-005, immune to wall-clock jumps.
10
11use std::time::{Duration, Instant};
12
13/// Token-bucket rate-limit throttle.
14#[derive(Debug, Clone)]
15pub struct TokenBucket {
16    rate_bytes_per_sec: u64,
17    start: Instant,
18}
19
20impl TokenBucket {
21    /// Construct a new throttle at `rate_bytes_per_sec`. The throttle's clock
22    /// starts at construction time.
23    #[must_use]
24    pub fn new(rate_bytes_per_sec: u64) -> Self {
25        TokenBucket {
26            rate_bytes_per_sec,
27            start: Instant::now(),
28        }
29    }
30
31    /// Inspect whether the loop is currently over-budget and, if so, return
32    /// the recommended sleep duration to converge to the rate. Returns
33    /// `Duration::ZERO` when the throttle is happy with current progress.
34    #[must_use]
35    pub fn next_sleep(&self, bytes_done: u64) -> Duration {
36        let elapsed = self.start.elapsed().as_secs_f64();
37        let expected = elapsed * self.rate_bytes_per_sec as f64;
38        let over = bytes_done as f64 - expected;
39        if over <= 0.0 {
40            return Duration::ZERO;
41        }
42        let want = over / self.rate_bytes_per_sec as f64;
43        // 10ms minimum so the sleep is meaningful on Windows.
44        Duration::from_secs_f64(want.max(0.010))
45    }
46
47    /// Sleep for the recommended interval if any. Returns `true` if a sleep
48    /// was performed, `false` if no sleep was needed.
49    pub fn maybe_sleep(&self, bytes_done: u64) -> bool {
50        let d = self.next_sleep(bytes_done);
51        if d == Duration::ZERO {
52            return false;
53        }
54        std::thread::sleep(d);
55        true
56    }
57}
58
59#[cfg(test)]
60mod tests {
61    use super::*;
62
63    #[test]
64    fn zero_when_under_budget() {
65        let tb = TokenBucket::new(1_000_000);
66        // Newly constructed: elapsed ≈ 0, expected ≈ 0, bytes_done = 0 → not over.
67        assert_eq!(tb.next_sleep(0), Duration::ZERO);
68    }
69
70    #[test]
71    fn sleep_when_over_budget() {
72        let tb = TokenBucket::new(1_000_000);
73        // Pretend we've done 2 MB in zero elapsed time → over-budget.
74        let d = tb.next_sleep(2_000_000);
75        assert!(d >= Duration::from_millis(10));
76        assert!(d >= Duration::from_millis(500));
77    }
78}