slack_morphism/ratectl/
throttling_counter.rs1use std::ops::Add;
2use std::time::{Duration, Instant};
3
4#[derive(Clone, Debug)]
5pub struct ThrottlingCounter {
6 capacity: i64,
7 max_capacity: usize,
8 last_arrived: Instant,
9 last_updated: Instant,
10 rate_limit_in_millis: u64,
11 delay: Duration,
12}
13
14impl ThrottlingCounter {
15 pub fn new(max_capacity: usize, rate_limit_in_millis: u64) -> Self {
16 Self {
17 capacity: max_capacity as i64,
18 max_capacity,
19 last_arrived: Instant::now(),
20 last_updated: Instant::now(),
21 rate_limit_in_millis,
22 delay: Duration::from_millis(0),
23 }
24 }
25
26 pub fn update(&self, now: Instant) -> Self {
27 let time_elapsed_millis = now
28 .checked_duration_since(self.last_arrived)
29 .unwrap_or_else(|| Duration::from_millis(0))
30 .as_millis() as u64;
31
32 let (arrived, new_last_arrived) = {
33 if time_elapsed_millis >= self.rate_limit_in_millis {
34 let arrived_in_time = time_elapsed_millis / self.rate_limit_in_millis;
35 let new_last_updated = self.last_arrived.add(Duration::from_millis(
36 arrived_in_time * self.rate_limit_in_millis,
37 ));
38 (arrived_in_time as usize, new_last_updated)
39 } else {
40 (0usize, self.last_arrived)
41 }
42 };
43
44 let new_available_capacity =
45 std::cmp::min(self.capacity + arrived as i64, self.max_capacity as i64);
46
47 if new_available_capacity > 0 {
48 Self {
49 capacity: new_available_capacity - 1,
50 last_arrived: new_last_arrived,
51 last_updated: now,
52 delay: Duration::from_millis(0),
53 ..self.clone()
54 }
55 } else {
56 let updated_capacity = new_available_capacity - 1;
57
58 let base_delay_in_millis = now
59 .checked_duration_since(self.last_updated)
60 .map_or(0u64, |d| d.as_millis() as u64);
61
62 let delay_penalty = (self.rate_limit_in_millis as f64 * self.capacity.abs() as f64
63 / self.max_capacity as f64) as u64;
64
65 let delay_in_millis = self
66 .rate_limit_in_millis
67 .saturating_sub(base_delay_in_millis);
68 let delay_with_penalty = Duration::from_millis(delay_in_millis + delay_penalty);
69
70 Self {
71 capacity: updated_capacity,
72 last_arrived: new_last_arrived,
73 last_updated: now,
74 delay: delay_with_penalty,
75 ..self.clone()
76 }
77 }
78 }
79
80 pub fn delay(&self) -> &Duration {
81 &self.delay
82 }
83}
84
85#[test]
86fn check_decreased() {
87 use crate::ratectl::*;
88 let rate_limit = SlackApiRateControlLimit::new(15, Duration::from_secs(60));
89 let rate_limit_in_ms = rate_limit.to_rate_limit_in_ms();
90 let rate_limit_capacity = rate_limit.to_rate_limit_capacity();
91
92 let now = Instant::now();
93 let counter = ThrottlingCounter::new(rate_limit_capacity, rate_limit_in_ms);
94 let updated_counter = counter.update(now.add(Duration::from_millis(rate_limit_in_ms - 1)));
95
96 assert_eq!(updated_counter.last_arrived, counter.last_arrived);
97 assert_eq!(updated_counter.delay, Duration::from_millis(0));
98 assert_eq!(updated_counter.capacity, counter.capacity - 1);
99}
100
101#[test]
102fn check_max_available() {
103 use crate::ratectl::*;
104 let rate_limit = SlackApiRateControlLimit::new(15, Duration::from_secs(60));
105 let rate_limit_in_ms = rate_limit.to_rate_limit_in_ms();
106 let rate_limit_capacity = rate_limit.to_rate_limit_capacity();
107
108 let now = Instant::now();
109 let counter = ThrottlingCounter::new(rate_limit_capacity, rate_limit_in_ms);
110 let updated_counter = counter.update(now.add(Duration::from_millis(rate_limit_in_ms + 1)));
111
112 assert_eq!(updated_counter.delay, Duration::from_millis(0));
113 assert_eq!(updated_counter.capacity, (counter.max_capacity - 1) as i64);
114}
115
116#[test]
117fn check_delay() {
118 use crate::ratectl::*;
119 let counter =
120 SlackApiRateControlLimit::new(15, Duration::from_secs(60)).to_throttling_counter();
121
122 let now = Instant::now();
123
124 let updated_counter =
125 (0..counter.capacity + 1).fold(counter.clone(), |result, _| result.update(now));
126
127 assert_eq!(
128 updated_counter.delay,
129 Duration::from_millis(counter.rate_limit_in_millis)
130 );
131}