Skip to main content

tor_async_utils/bw_pool/
bucket.rs

1//! An atomic and shareable token bucket BUT with one caveat, it is without the clock component
2//! that is the refill is not taking into account any rate or time tracking.
3//!
4//! This [`AtomicTokenBucket`] is used within a [`super::BandwidthPool`] to keep track of the
5//! available bandwidth.
6//!
7//! It is refilled through an entity, the [`super::BandwidthRefiller`], which is owned by an
8//! independent task built to refill this bucket at the appropriate time.
9//!
10//! See the [`super`] documentation for more information on how these objects interact with each
11//! other.
12
13// TODO MSRV 1.95: Remove this, and use try_update instead of fetch_update.
14#![allow(deprecated)]
15
16use std::sync::atomic::{AtomicU64, Ordering};
17
18/// The atomic token bucket minus the clock component.
19///
20/// This is lock-free and the owner needs to refill it with a specific number of tokens it wants to
21/// be distributed across many actors.
22#[derive(Debug)]
23pub(super) struct AtomicTokenBucket {
24    /// Every access to these counters is with [`Ordering::Relaxed`] as they don't
25    /// protect any outside data and so we only care about the atomicity action on the
26    /// counter.
27
28    /// The current token count.
29    available: AtomicU64,
30    /// The maximum number of tokens the bucket may hold a.k.a the burst.
31    ///
32    /// This is atomic because it can be set during runtime. For instance, a config
33    /// option reconfigure of a bandwidth rate.
34    capacity: u64,
35}
36
37impl AtomicTokenBucket {
38    /// A new bucket capped at `capacity` tokens.
39    ///
40    /// It starts full.
41    pub(super) fn new(capacity: u64) -> Self {
42        AtomicTokenBucket {
43            available: AtomicU64::new(capacity),
44            capacity,
45        }
46    }
47
48    /// Claim all `tokens` from the bucket or nothing if not enough available.
49    ///
50    /// The request is first clamped to the capacity because the bucket can never hold
51    /// more than its burst.
52    ///
53    /// Returns `Some(claimed)` if the bucket held at least the clamped request which indicates
54    /// that they are now granted to the caller. Note that `claimed` can be lower than `tokens`
55    /// hence why the caller needs to use that value and not what it asked for.
56    ///
57    /// Returns `None` otherwise, nothing is taken.
58    #[must_use]
59    pub(super) fn claim(&self, tokens: u64) -> Option<u64> {
60        let tokens = tokens.min(self.capacity);
61
62        // NOTE: fetch_update() is deprecated in 1.99.0 and replaced by try_update()
63        // starting in 1.95.0. Our current MSRV is 1.89.0.
64        //
65        // Relaxed: only the atomic subtract matters, the counter gates no other memory.
66        self.available
67            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |cur| {
68                cur.checked_sub(tokens)
69            })
70            .is_ok()
71            .then_some(tokens)
72    }
73
74    /// Add `tokens` to the bucket which is capped at the capacity.
75    ///
76    /// We use a CAS loop (Compare-And-Swap) because we need to cap the refill to the internal
77    /// capacity atomically. A fetch_add + fetch_sub is not possible in order to refill atomically
78    /// because of this race:
79    ///
80    /// ```text
81    ///   capacity = 100, available = 90:
82    ///     refill: fetch_add(50)  -> available = 140 (above capacity)
83    ///     claim:  fetch_sub(140) -> available = 0   (illegal claim, above capacity)
84    ///     refill: fetch_sub(40)  -> available underflows (adjust overshoot too late)
85    /// ```
86    ///
87    /// Any tokens that overshoot the capacity are forfeited.
88    ///
89    /// Finally, the CAS loop is considered ok because the refill is not in the fast path and
90    /// should work most of the time on the first iteration.
91    pub(super) fn refill(&self, tokens: u64) {
92        if tokens == 0 {
93            return;
94        }
95
96        // NOTE: fetch_update() is deprecated in 1.99.0 and replaced by try_update()
97        // starting in 1.95.0. Our current MSRV is 1.89.0.
98        //
99        // Relaxed: only the atomic add matters, the counter gates no other memory.
100        let _ = self
101            .available
102            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |cur| {
103                Some(cur.saturating_add(tokens).min(self.capacity))
104            });
105    }
106
107    /// Take every token out of the bucket and return the amount.
108    pub(super) fn drain(&self) -> u64 {
109        // Relaxed: only the atomic exchange matters, the counter gates no other memory.
110        self.available.swap(0, Ordering::Relaxed)
111    }
112
113    /// Return the capacity.
114    pub(super) fn capacity(&self) -> u64 {
115        self.capacity
116    }
117
118    /// Return the available balance.
119    #[cfg(test)]
120    pub(super) fn available(&self) -> u64 {
121        // Relaxed: we only need the current value of the counter.
122        self.available.load(Ordering::Relaxed)
123    }
124}
125
126#[cfg(test)]
127mod test {
128    // @@ begin test lint list maintained by maint/add_warning @@
129    #![allow(clippy::bool_assert_comparison)]
130    #![allow(clippy::clone_on_copy)]
131    #![allow(clippy::dbg_macro)]
132    #![allow(clippy::mixed_attributes_style)]
133    #![allow(clippy::print_stderr)]
134    #![allow(clippy::print_stdout)]
135    #![allow(clippy::single_char_pattern)]
136    #![allow(clippy::unwrap_used)]
137    #![allow(clippy::unchecked_time_subtraction)]
138    #![allow(clippy::useless_vec)]
139    #![allow(clippy::needless_pass_by_value)]
140    #![allow(clippy::string_slice)] // See arti#2571
141    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
142
143    use super::*;
144
145    #[test]
146    fn claim_and_refill() {
147        let b = AtomicTokenBucket::new(100);
148        assert_eq!(b.available(), 100);
149        assert_eq!(b.capacity(), 100);
150
151        assert_eq!(b.claim(60), Some(60));
152        // Only 40 left. All or nothing.
153        assert_eq!(b.claim(50), None);
154        assert_eq!(b.available(), 40);
155        assert_eq!(b.claim(40), Some(40));
156        assert_eq!(b.claim(1), None);
157
158        // Refill is capped at capacity.
159        b.refill(1000);
160        assert_eq!(b.available(), 100);
161    }
162
163    #[test]
164    fn claim_above_capacity() {
165        let b = AtomicTokenBucket::new(100);
166
167        // A claim above the burst is clamped to it and not refused.
168        assert_eq!(b.claim(1000), Some(100));
169        assert_eq!(b.available(), 0);
170
171        // Even clamped, it is all or nothing.
172        b.refill(50);
173        assert_eq!(b.claim(1000), None);
174        assert_eq!(b.available(), 50);
175    }
176
177    #[test]
178    fn drain() {
179        let b = AtomicTokenBucket::new(100);
180        assert_eq!(b.claim(90), Some(90)); // 10 tokens left
181
182        assert_eq!(b.drain(), 10);
183        assert_eq!(b.available(), 0);
184
185        // Draining an empty bucket takes nothing.
186        assert_eq!(b.drain(), 0);
187    }
188}