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}