Skip to main content

tor_async_utils/
bw_pool.rs

1//! A shared token pool with a lock-free fast path and FIFO waiting queue.
2//!
3//! [`BandwidthPool`] is a set of tokens that many tasks can draw from concurrently
4//! representing a bandwidth [`Permit`]. Deciding how many tokens become available and
5//! when is the job of the [`BandwidthRefiller`] which should be run in a task that has a
6//! reference to the pool's token bucket.
7//!
8//! # Pool
9//!
10//! The pool holds the bandwidth token balance (bucket) available for any
11//! [`BandwidthAcquirer`] to attempt to acquire concurrently.
12//!
13//! The [`BandwidthPool`] can be shared between an arbitrary amount of tasks and each
14//! task needs to get a [`BandwidthAcquirer`] from the [`BandwidthPool::new_acquirer`]
15//! method.
16//!
17//! The acquirer is designed to be allocated once and reused throughout the owner object
18//! lifetime, thereby reducing the number of allocations needed at runtime.
19//!
20//! In order to be granted permission to use a certain number of tokens, the task needs
21//! to call [`BandwidthAcquirer::poll_acquire`] with the number of tokens it wants. It
22//! has to be called in the context of a task so it can be woken up once it is granted.
23//!
24//! ## Acquisition Mechanism
25//!
26//! The pool design is that there is a so called fast-path that is meant to allow a
27//! thundering herd to attempt to acquire tokens in an atomic way as long as the pool has
28//! available tokens.
29//!
30//! Once the pool is empty, new acquire requests go into a FIFO queue for fairness and
31//! are served as the pool gets refilled by the [`BandwidthRefiller`] task.
32//!
33//! A [`Permit`] is handed out once the full requested amount has been deducted from the
34//! pool as in available. The permit holds the granted tokens and the holder claims what
35//! it actually uses with [`Permit::claim`] (or [`Permit::claim_all`]). Whatever is left
36//! unclaimed in the permit is refunded to the pool when it is dropped.
37//!
38//! ## Refund
39//!
40//! Refunded tokens always go back to the fast-path pool where they are immediately
41//! available. A newcomer can acquire them while another request is queued but before
42//! the [`BandwidthRefiller`] drains the pool. This favors keeping the fast path busy
43//! over strict fairness between the fast path and the waiting queue.
44//!
45//! ## Teardown
46//!
47//! There is deliberately no cancellation mechanism for a queued request. If a
48//! [`BandwidthAcquirer`] is torn down while its request is queued, the refiller will
49//! eventually fund that request anyway, record the grant, and wake a task that no
50//! longer exists. The granted tokens are simply forfeited.
51//!
52//! This is a considered trade-off and we believe in the context of a Tor relay, losing a
53//! grant is not significant at all in the large picture of available bandwidth. The
54//! trade-off allows us to reduce a lot of complexity.
55
56mod bucket;
57mod refiller;
58
59use futures::channel::mpsc;
60use std::sync::Arc;
61use std::task::{Context, Poll};
62
63use crate::bw_pool::bucket::AtomicTokenBucket;
64use crate::bw_pool::refiller::{BwRequest, RefillWaiter};
65
66// Public export as the outside world needs this.
67pub use crate::bw_pool::refiller::BandwidthRefiller;
68
69/// Error returned by this module.
70#[derive(Clone, Debug, thiserror::Error)]
71#[non_exhaustive]
72pub enum BwPoolError {
73    /// Over permit claim.
74    #[error("Permit's claim() exceeds the granted tokens")]
75    ClaimExceedsGrant,
76    /// The bandwidth pool is closed.
77    #[error("bandwidth pool is closed")]
78    PoolClosed,
79}
80
81/// Proof that the requested number of tokens were granted.
82///
83/// If dropped, any remaining tokens will be refunded hence why there is a reference to
84/// the shared bucket.
85#[derive(Debug)]
86#[must_use = "a Permit are tokens that have already been deducted"]
87pub struct Permit {
88    /// The bucket the unclaimed tokens are refunded to when this permit is dropped.
89    bucket: Arc<AtomicTokenBucket>,
90    /// How many granted tokens this permit still has.
91    ///
92    /// Lowered by the claim methods. What is left is refunded on drop.
93    granted: u64,
94}
95
96impl Permit {
97    /// Construct a permit granting `granted` tokens.
98    ///
99    /// We keep a reference to the bucket for a refund if any.
100    fn new(bucket: Arc<AtomicTokenBucket>, granted: u64) -> Self {
101        Self { bucket, granted }
102    }
103
104    /// Claim `tokens` from this permit.
105    ///
106    /// Returns true if the permit had at least `tokens` that are now claimed. Else,
107    /// return false and nothing is claimed as in the caller is trying to use more
108    /// than it was granted.
109    pub fn claim(&mut self, tokens: u64) -> Result<(), BwPoolError> {
110        match self.granted.checked_sub(tokens) {
111            Some(left) => {
112                self.granted = left;
113                Ok(())
114            }
115            None => Err(BwPoolError::ClaimExceedsGrant),
116        }
117    }
118
119    /// Claim everything this permit holds.
120    ///
121    /// Used in the for a Sink implementation that gets its permit in the poll_ready()
122    /// which is before the start_send().
123    pub fn claim_all(&mut self) {
124        let r = self.claim(self.granted);
125        debug_assert!(r.is_ok());
126    }
127
128    /// How many tokens this permit still has granted a.k.a available.
129    pub fn granted(&self) -> u64 {
130        self.granted
131    }
132}
133
134impl Drop for Permit {
135    fn drop(&mut self) {
136        // Refund into the shared pool. Note that the refund will have a no-op on the
137        // shared counters if the granted value is zero.
138        self.bucket.refill(self.granted);
139    }
140}
141
142/// A reusable bandwidth acquirer that is designed for an async context (poll).
143///
144/// The main entry point is [`Self::poll_acquire`] used within the context of a task
145/// which, if rate limited, the task gets woken up when the bandwidth usage is granted.
146#[derive(Debug)]
147pub struct BandwidthAcquirer {
148    /// The shared waiter. Reused for every acquire.
149    waiter: Arc<RefillWaiter>,
150    /// Whether this acquirer currently has a request enqueued with the refiller.
151    ///
152    /// Prevent multi poll to avoid queuing a duplicate request with a different waker.
153    in_flight: bool,
154    /// The bandwidth pool attached to this acquirer.
155    pool: BandwidthPool,
156}
157
158impl BandwidthAcquirer {
159    /// Create a reusable [`BandwidthAcquirer`] for the given `pool`.
160    ///
161    /// This is the only allocation an acquirer makes. It can only be called from
162    /// [`BandwidthPool::new_acquirer`].
163    ///
164    /// The number of tokens to acquire is chosen per [`Self::poll_acquire`] call.
165    fn new(pool: BandwidthPool) -> Self {
166        Self {
167            waiter: Arc::new(RefillWaiter::new()),
168            in_flight: false,
169            pool,
170        }
171    }
172
173    /// Build the [`Permit`].
174    ///
175    /// This sets the `in-flight` to false and takes the granted tokens out of the waiter,
176    /// resetting it, as we now hand them to the permit.
177    fn grant_permit(&mut self) -> Permit {
178        let granted = self.waiter.take_granted();
179        self.in_flight = false;
180        Permit::new(Arc::clone(&self.pool.bucket), granted)
181    }
182
183    /// Poll to take `tokens` tokens from the pool, async waiting if the pool is dry.
184    ///
185    /// The number of `tokens` is per call so the same acquirer can request a different
186    /// amount each time.
187    ///
188    /// If the request is in flight, the `tokens` argument is ignored and the originally
189    /// requested amount is used. Only when a [`Permit`] is emitted that a new `tokens`
190    /// value can be used.
191    ///
192    /// A request above the pool capacity is clamped to it rather than refused. The
193    /// emitted [`Permit`] then holds less than what was asked and so the caller must use
194    /// [`Permit::granted`] rather than assume it got `tokens`.
195    ///
196    /// Returns a [`BwPoolError::PoolClosed`] error if the [`BandwidthRefiller`] has been
197    /// dropped.
198    pub fn poll_acquire(
199        &mut self,
200        cx: &mut Context<'_>,
201        tokens: u64,
202    ) -> Poll<Result<Permit, BwPoolError>> {
203        if !self.in_flight {
204            // No request in flight. This is the fast path! The thundering herd is allowed,
205            // that is, all tasks race to this and their fairness is sub-contracted to their
206            // task scheduler.
207            //
208            // An under-utilized relay with bandwidth limitation will hit this path most
209            // of the time.
210            //
211            // The fast path is deliberately not gated by queued requests. But there is a
212            // fairness race that we accept:
213            //
214            //      task A: fail the fast path and enqueue
215            //      permit: drop and refund enough tokens
216            //      task B: claim the refund before the refiller sees queued task A
217            //
218            // Task B has jumped task A.
219            //
220            // This is OK because the refiller drains the global bucket before serving
221            // the queue and the refill are never done directly into that bucket.
222            // Instead, it keeps it locally until all enqueued requests have been served.
223            // And so, letting B use the refund(s) keeps the pool active without breaking
224            // the rate/burst limits.
225            //
226            // This prevents a slow enqueuing on one CPU to not block the fast path on
227            // another CPU. This logic also applies for a slow refiller task. We
228            // essentially compromise almost-strict fairness for high throughput.
229            //
230            // For the case of a relay, in normal operations, we expect the write side to
231            // rarely have refunds and so this race would be rare. On the read side,
232            // refunds are much more likely to happen because as a relay we don't know
233            // exactly how many bytes we should expect at every recv().
234            if let Some(permit) = self.pool.try_acquire(tokens) {
235                return Poll::Ready(Ok(permit));
236            }
237
238            // Enqueue the request as it is not in flight and our bw pool is depleted.
239            // This is where the caller gets to wait on bw availability and the refilling
240            // process is triggered.
241            match self.enqueue_request(cx, tokens) {
242                Ok(()) => return Poll::Pending,
243                Err(e) => return Poll::Ready(Err(e)),
244            };
245        }
246
247        // Register our current waker and check if we were granted.
248        if self.waiter.poll_granted(cx.waker()) {
249            return Poll::Ready(Ok(self.grant_permit()));
250        }
251
252        // Check if the refiller is gone. Reason this is done at the end is because we
253        // want to avoid this issue:
254        //
255        //      refiller: grant tokens. waker.wake()
256        //      refiller: drop() as in dropped
257        //      sink:     poll_acquire() is called and notices the grant.
258        //
259        // The race shows that we would miss the grant even if the refiller is gone. The
260        // end result is that we get to at least send these bytes before the whole arti
261        // relay collapses.
262        if self.pool.is_closed() {
263            // The refiller is gone; nobody will ever serve us.
264            self.in_flight = false;
265            return Poll::Ready(Err(BwPoolError::PoolClosed));
266        }
267
268        Poll::Pending
269    }
270
271    /// Enqueue a request for `tokens` tokens that is NOT in flight.
272    ///
273    /// The amount is sent as it was asked. Clamping it to the pool capacity is the
274    /// refiller's job as it is the one holding the tokens by the time this is served.
275    /// The clamping is localized to the refiller for correctness.
276    ///
277    /// Return a [`BwPoolError::PoolClosed`] error if the refiller is gone.
278    fn enqueue_request(&mut self, cx: &mut Context<'_>, tokens: u64) -> Result<(), BwPoolError> {
279        // Prepare the waiter for this new request.
280        self.waiter.prepare(cx.waker());
281        // Send the requested amount along with the waiter. The amount never changes for
282        // the lifetime of the request so it doesn't need to live in the shared waiter.
283        //
284        // Notice, this is the only dynamic allocation in this path because it is sent on
285        // an unbounded MPSC queue.
286        if self
287            .pool
288            .requests
289            .unbounded_send((tokens, Arc::clone(&self.waiter)))
290            .is_err()
291        {
292            // The refiller is gone, the pool is closed.
293            return Err(BwPoolError::PoolClosed);
294        }
295        self.in_flight = true;
296        Ok(())
297    }
298}
299
300/// A shareable bandwidth pool.
301///
302/// Use [`Self::share`] to give it to each task requiring bandwidth limitation.
303///
304/// [`Clone`] is deliberately not implemented so that handing this around can not be
305/// mistaken for creating a second pool with its own bandwidth. There is exactly one pool
306/// per [`Self::new`] call.
307///
308/// Each task needs to hold a [`BandwidthAcquirer`] in order to request bandwidth permits
309/// from this pool.
310#[derive(Debug)]
311pub struct BandwidthPool {
312    /// The shared token bucket the fast path claims from.
313    bucket: Arc<AtomicTokenBucket>,
314    /// Ingress for acquirers that failed the fast path.
315    ///
316    /// Sending a [`BwRequest`] here both enqueues it and wakes the
317    /// [`BandwidthRefiller`]
318    ///
319    /// Unbounded because we never want an acquirer's enqueue to block. the number of
320    /// in-flight requests is bounded by the number of acquirers.
321    requests: mpsc::UnboundedSender<BwRequest>,
322}
323
324impl BandwidthPool {
325    /// Create a new pool that can hold up to `capacity` tokens, and its associated
326    /// refiller [`BandwidthRefiller`] that needs to be run in its own task.
327    ///
328    /// The pool starts full and `capacity` tokens are immediately available to the fast
329    /// path.
330    pub fn new(capacity: u64) -> (BandwidthPool, BandwidthRefiller) {
331        let (tx, rx) = mpsc::unbounded();
332        let bucket = Arc::new(AtomicTokenBucket::new(capacity));
333        let pool = BandwidthPool {
334            bucket: Arc::clone(&bucket),
335            requests: tx,
336        };
337        let refiller = BandwidthRefiller::new(bucket, rx);
338        (pool, refiller)
339    }
340
341    /// Return another [`BandwidthPool`] that is the same as this one.
342    ///
343    /// The tokens are not duplicated, this is so you can share the pool with other
344    /// tasks/objects.
345    pub fn share(&self) -> BandwidthPool {
346        BandwidthPool {
347            bucket: Arc::clone(&self.bucket),
348            requests: self.requests.clone(),
349        }
350    }
351
352    /// Return a new [`BandwidthAcquirer`] associated to this pool.
353    ///
354    /// The number of tokens to acquire is chosen per
355    /// [`BandwidthAcquirer::poll_acquire`] call, so a single acquirer can be reused for
356    /// requests of different sizes.
357    ///
358    /// It is through an acquirer that one can get permission to use bandwidth. See
359    /// [`BandwidthAcquirer::poll_acquire`].
360    pub fn new_acquirer(&self) -> BandwidthAcquirer {
361        BandwidthAcquirer::new(self.share())
362    }
363
364    /// The maximum number of tokens this pool can hold (its burst).
365    pub fn capacity(&self) -> u64 {
366        self.bucket.capacity()
367    }
368
369    /// Return true iff the refiller is gone, meaning this pool is closed.
370    fn is_closed(&self) -> bool {
371        self.requests.is_closed()
372    }
373
374    /// Try to take `tokens` from the pool without waiting.
375    ///
376    /// A request above the pool capacity is clamped to it by the bucket it self so the
377    /// permit can hold less than what was asked. The caller learns what it got with
378    /// [`Permit::granted`].
379    ///
380    /// Returns `Some(permit)` if there were tokens available right now, or `None` if the
381    /// bucket is empty.
382    ///
383    /// This never blocks and never enqueues. It is the fast path that the other
384    /// acquisition methods use first.
385    fn try_acquire(&self, tokens: u64) -> Option<Permit> {
386        let granted = self.bucket.claim(tokens)?;
387        Some(Permit::new(Arc::clone(&self.bucket), granted))
388    }
389
390    /// Unit tests helper: async acquire wrapping a throwaway [`BandwidthAcquirer`].
391    #[cfg(test)]
392    async fn acquire(&self, tokens: u64) -> Result<Permit, BwPoolError> {
393        if let Some(permit) = self.try_acquire(tokens) {
394            return Ok(permit);
395        }
396        let mut acquirer = BandwidthAcquirer::new(self.share());
397        std::future::poll_fn(|cx| acquirer.poll_acquire(cx, tokens)).await
398    }
399
400    /// Unit tests helper: The number of tokens currently available to the fast path.
401    #[cfg(test)]
402    pub(crate) fn available(&self) -> u64 {
403        self.bucket.available()
404    }
405}
406
407#[cfg(test)]
408mod test {
409    // @@ begin test lint list maintained by maint/add_warning @@
410    #![allow(clippy::bool_assert_comparison)]
411    #![allow(clippy::clone_on_copy)]
412    #![allow(clippy::dbg_macro)]
413    #![allow(clippy::mixed_attributes_style)]
414    #![allow(clippy::print_stderr)]
415    #![allow(clippy::print_stdout)]
416    #![allow(clippy::single_char_pattern)]
417    #![allow(clippy::unwrap_used)]
418    #![allow(clippy::unchecked_time_subtraction)]
419    #![allow(clippy::useless_vec)]
420    #![allow(clippy::needless_pass_by_value)]
421    #![allow(clippy::string_slice)] // See arti#2571
422    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
423
424    use super::*;
425
426    use std::future::Future;
427    use std::pin::pin;
428
429    use futures::FutureExt as _;
430    use futures::task::noop_waker_ref;
431
432    /// A `Context` backed by a no-op waker, for deterministic manual polling.
433    ///
434    /// This is a trick so we don't rely on wakeups as grants are committed by a call to
435    /// refill_and_serve() into the waiter's grant flag. So a second poll observes it.
436    fn noop_cx() -> Context<'static> {
437        Context::from_waker(noop_waker_ref())
438    }
439
440    /// Build a [`BandwidthPool`] and drain it so 0 tokens are available.
441    fn drained_pool(cap: u64) -> (BandwidthPool, BandwidthRefiller) {
442        let (pool, refiller) = BandwidthPool::new(cap);
443        let mut permit = pool.try_acquire(cap).expect("a fresh pool starts full");
444        permit.claim_all();
445        (pool, refiller)
446    }
447
448    /// Expect `poll` to be a granted permit and claim all so the drop refunds zero.
449    fn expect_granted(poll: Poll<Result<Permit, BwPoolError>>) {
450        match poll {
451            Poll::Ready(Ok(mut permit)) => permit.claim_all(),
452            other => panic!("expected a granted permit, got {other:?}"),
453        }
454    }
455
456    #[test]
457    fn fast_path() {
458        let (pool, _refiller) = BandwidthPool::new(100);
459        assert_eq!(pool.available(), 100);
460
461        let mut permit = pool.acquire(40).now_or_never().unwrap().unwrap();
462        permit.claim_all();
463        drop(permit);
464        assert_eq!(pool.available(), 60);
465
466        let mut permit = pool.try_acquire(60).unwrap();
467        permit.claim_all();
468        drop(permit);
469        assert_eq!(pool.available(), 0);
470
471        // Nothing left. This should block and the fast path returns nothing.
472        assert!(pool.try_acquire(1).is_none());
473        assert!(pool.acquire(1).now_or_never().is_none());
474    }
475
476    #[test]
477    fn acquirer() {
478        let (pool, mut refiller) = drained_pool(100);
479
480        // Dummy context so we can poll without a wakeup.
481        let mut cx = noop_cx();
482        let mut a = pin!(pool.acquire(30));
483        // No tokens, it should be pending.
484        assert!(a.as_mut().poll(&mut cx).is_pending());
485
486        // Refill to serve it fully.
487        assert_eq!(refiller.refill_and_serve(30), None);
488        expect_granted(a.as_mut().poll(&mut cx));
489    }
490
491    #[test]
492    fn acquirer_reuse() {
493        let (pool, mut refiller) = drained_pool(100);
494
495        // Dummy context so we can poll without a wakeup.
496        let mut cx = noop_cx();
497        let mut acquirer = BandwidthAcquirer::new(pool.share());
498
499        // First blocked acquire through the acquirer.
500        assert!(acquirer.poll_acquire(&mut cx, 30).is_pending());
501        assert_eq!(refiller.refill_and_serve(30), None);
502        expect_granted(acquirer.poll_acquire(&mut cx, 30));
503
504        // The same acquirer is reusable with a different amount. No new allocation.
505        assert!(acquirer.poll_acquire(&mut cx, 50).is_pending());
506        assert_eq!(refiller.refill_and_serve(50), None);
507        expect_granted(acquirer.poll_acquire(&mut cx, 50));
508    }
509
510    #[test]
511    fn refill_newcomer() {
512        let (pool, mut refiller) = drained_pool(100);
513
514        // Dummy context so we can poll without a wakeup.
515        let mut cx = noop_cx();
516        // Enqueue 30 for A then 50 for B.
517        let mut a = pin!(pool.acquire(30));
518        assert!(a.as_mut().poll(&mut cx).is_pending());
519        let mut b = pin!(pool.acquire(50));
520        assert!(b.as_mut().poll(&mut cx).is_pending());
521
522        // Refill of 40: serves A and keeps 10 reserved for B meaning a deficit of 40.
523        assert_eq!(refiller.refill_and_serve(40), Some(40));
524        expect_granted(a.as_mut().poll(&mut cx));
525        assert!(b.as_mut().poll(&mut cx).is_pending());
526
527        // The 10 reserved tokens are NOT visible to the fast path meaning a newcomer
528        // cannot barge ahead of B (deficit of 40 to serve it).
529        assert_eq!(pool.available(), 0);
530        let mut c = pin!(pool.acquire(5));
531        assert!(c.as_mut().poll(&mut cx).is_pending());
532
533        // Second refill of 40 plus held was 10 which is 50 that B needs.
534        assert_eq!(refiller.refill_and_serve(40), Some(5)); // C is the head and 5 is the deficit.
535        expect_granted(b.as_mut().poll(&mut cx));
536        assert!(c.as_mut().poll(&mut cx).is_pending());
537
538        // And finally C.
539        assert_eq!(refiller.refill_and_serve(5), None);
540        expect_granted(c.as_mut().poll(&mut cx));
541    }
542
543    #[test]
544    fn refill_idle() {
545        let (pool, mut refiller) = drained_pool(100);
546
547        // No acquirers. Refill with a large value keeps it cap to the pool capacity.
548        assert_eq!(refiller.refill_and_serve(1000), None);
549        assert_eq!(pool.available(), 100);
550        assert_eq!(pool.capacity(), 100);
551
552        // Fast path is fine.
553        assert!(pool.try_acquire(100).is_some());
554    }
555
556    #[test]
557    fn close_pool() {
558        // Acquire fails once closed and the bucket is dry.
559        let (pool, refiller) = drained_pool(100);
560        drop(refiller);
561        assert!(pool.try_acquire(1).is_none());
562        assert!(matches!(
563            pool.acquire(10).now_or_never(),
564            Some(Err(BwPoolError::PoolClosed)),
565        ));
566
567        // A blocked acquirer is woken with PoolClosed when the refiller drops.
568        let (pool, refiller) = drained_pool(100);
569        // Dummy context so we can poll without a wakeup.
570        let mut cx = noop_cx();
571        let mut a = pin!(pool.acquire(10));
572        assert!(a.as_mut().poll(&mut cx).is_pending());
573        drop(refiller);
574        assert!(matches!(
575            a.as_mut().poll(&mut cx),
576            Poll::Ready(Err(BwPoolError::PoolClosed)),
577        ));
578
579        let (pool, mut refiller) = BandwidthPool::new(0);
580        drop(pool);
581        // No senders left, wait should report closure.
582        assert_eq!(refiller.wait().now_or_never(), Some(false));
583    }
584
585    #[test]
586    fn refiller_wait() {
587        let (pool, mut refiller) = drained_pool(100);
588
589        // Dummy context so we can poll without a wakeup.
590        let mut cx = noop_cx();
591        // Idle: nobody waiting so wait() is pending.
592        {
593            let mut wait_fut = pin!(refiller.wait());
594            assert!(wait_fut.as_mut().poll(&mut cx).is_pending());
595        }
596
597        // A blocked acquire enqueues a request (and wakes the refiller)...
598        let mut a = pin!(pool.acquire(10));
599        assert!(a.as_mut().poll(&mut cx).is_pending());
600
601        // wait() now completes and having taken the request as the head.
602        {
603            let mut wait_fut = pin!(refiller.wait());
604            assert!(matches!(wait_fut.as_mut().poll(&mut cx), Poll::Ready(true)));
605        }
606
607        // And refilling serves the head.
608        assert_eq!(refiller.refill_and_serve(10), None);
609        expect_granted(a.as_mut().poll(&mut cx));
610    }
611
612    #[test]
613    fn permit_claim() {
614        let (pool, _refiller) = BandwidthPool::new(100);
615
616        // Grab 40 and drop. Should be fully refunded.
617        let permit = pool.try_acquire(40).unwrap();
618        assert_eq!(pool.available(), 60);
619        drop(permit);
620        assert_eq!(pool.available(), 100);
621
622        // Claims are cumulative.
623        let mut permit = pool.try_acquire(40).unwrap();
624        // Can't over claim.
625        assert!(matches!(
626            permit.claim(41),
627            Err(BwPoolError::ClaimExceedsGrant),
628        ));
629        assert!(permit.claim(10).is_ok());
630        assert!(permit.claim(20).is_ok());
631        // Over claiming by one.
632        assert!(matches!(
633            permit.claim(11),
634            Err(BwPoolError::ClaimExceedsGrant),
635        ));
636        // 10 remains unclaimed now (claimed: 30, unclaimed: 10, pool: 60)
637        drop(permit);
638        // Refund the 10 left, pool is now 70.
639        assert_eq!(pool.available(), 70);
640
641        // Get what remains.
642        let mut permit = pool.try_acquire(70).unwrap();
643        // A fully claimed permit refunds nothing.
644        permit.claim_all();
645        drop(permit);
646        assert_eq!(pool.available(), 0);
647    }
648
649    #[test]
650    fn refund_newcomer() {
651        let (pool, mut refiller) = BandwidthPool::new(100);
652        let mut cx = noop_cx();
653
654        // Get 60 on the fast path and only use 10 so we keep holding 50.
655        let mut permit_hold = pool.try_acquire(60).unwrap();
656        assert!(permit_hold.claim(10).is_ok());
657        // Drain the rest so the pool is empty.
658        let mut drain = pool.try_acquire(40).unwrap();
659        drain.claim_all();
660        drop(drain);
661        assert_eq!(pool.available(), 0);
662
663        // Task A acquires 50 but gets queued because pool is empty.
664        let mut a = pin!(pool.acquire(50));
665        assert!(a.as_mut().poll(&mut cx).is_pending());
666
667        // Drop permit which refunds 50 to the fast-path balance. A newcomer can claim
668        // the refund even though A is already queued.
669        drop(permit_hold);
670        assert_eq!(pool.available(), 50);
671        let mut newcomer = pool.try_acquire(50).unwrap();
672        newcomer.claim_all();
673        drop(newcomer);
674        assert_eq!(pool.available(), 0);
675
676        // A remains queued and is served by subsequent refills.
677        assert_eq!(refiller.refill_and_serve(10), Some(40));
678        assert!(a.as_mut().poll(&mut cx).is_pending());
679        assert_eq!(refiller.refill_and_serve(40), None);
680        expect_granted(a.as_mut().poll(&mut cx));
681        assert_eq!(pool.available(), 0);
682    }
683}