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}