pub struct BandwidthRefiller {
bucket: Arc<AtomicTokenBucket>,
rx: UnboundedReceiver<(u64, Arc<RefillWaiter>)>,
head: Option<(u64, Arc<RefillWaiter>)>,
held: u64,
}Expand description
A bandwidth refiller is in charge of refilling the associated
super::BandwidthPool and processing any pending RefillWaiter that were enqueued
by the pool.
There is exactly one refiller per pool as it owns the receiving end of the request channel. The channel is a FIFO of waiters.
Dropping it closes the pool which makes queued and new acquirers fail with
super::BwPoolError::PoolClosed.
Fields§
§bucket: Arc<AtomicTokenBucket>The shared token bucket which comes from the super::BandwidthPool.
rx: UnboundedReceiver<(u64, Arc<RefillWaiter>)>Receiving end of the request channel
head: Option<(u64, Arc<RefillWaiter>)>A single request we have taken off the channel to inspect but cannot yet fund. If only mpsc channels had an “is_empty()”.
This is populated by the Self::wait function
held: u64Tokens that have been drained out of the fast-path bucket or handed in via
Self::refill_and_serve but not yet distributed.
If anything is left at the end of the refill loop, it is published back in the main pool fast path.
Implementations§
Source§impl BandwidthRefiller
impl BandwidthRefiller
Sourcepub(super) fn new(
bucket: Arc<AtomicTokenBucket>,
rx: UnboundedReceiver<(u64, Arc<RefillWaiter>)>,
) -> Self
pub(super) fn new( bucket: Arc<AtomicTokenBucket>, rx: UnboundedReceiver<(u64, Arc<RefillWaiter>)>, ) -> Self
Constructor.
Sourcepub async fn run<SP: SleepProvider>(self, sleep: SP, rate: u64)
pub async fn run<SP: SleepProvider>(self, sleep: SP, rate: u64)
Start the refiller main loop. This should be run in its own task as it is blocking until the bandwidth pool closes.
Using the given SleepProvider, we drive the refill with it along side a
TokenBucket that is built at the start with the given config.
This waits on new request that comes in when the fast-path is depleted that is a request waiting for a refill. Once a request is received, a refill is triggered and then the loop sleeps until the needed deficit is available.
As an example, if the queue has a request for 10 tokens but only 5 are available in the pool after an immediate refill, we will sleep the exact time it takes to get another 5 tokens. Then, it goes on to the next request and so on until the queue is empty.
Keen observer will notice that once a request is received, the refiller will have to empty the entire queue before the fast path could even see 1 token added back by a refill. That is because each iteration of the loop sleeps the exact amount of time to fulfill the pending request.
Returns when the pool is closed or if the config rate is zero or if the requested amount of token is above the pool capacity.
Sourcepub async fn wait(&mut self) -> bool
pub async fn wait(&mut self) -> bool
Wait until at least one bandwidth request is queued.
Returns true once a request is received or false if the pool has been closed
meaning the tx end is closed.
Returns true immediately if a request is already held as the head.
This should only be used as a “doorbell” that is indicating someone is at the
door with a request rather than waiting for the next request. One should use
Self::serve for that.
Sourcepub fn refill_and_serve(&mut self, tokens: u64) -> Option<u64>
pub fn refill_and_serve(&mut self, tokens: u64) -> Option<u64>
Add tokens to the pool and then serve all pending requests if any.
The very first thing that this function does is drain the pool’s fast path tokens in order to reserve its current balance for the waiting queue.
Any surplus left will be put back into the pool’s fast path.
Returns None if no one is left waiting indicating the pool is now idle and
Self::wait can be safely used to get notified of a new request.
Returns Some(deficit) if an acquirer is still waiting where deficit is how
many more tokens are needed before it can be served. The caller can use this to
decide how long to wait before the next refill.
Sourcefn serve(&mut self, capacity: u64)
fn serve(&mut self, capacity: u64)
Serve pending requests with the token we are holding.
The given capacity is essentially the maximum we can give a single request.
If we have a token deficit, the head is updated with the latest request that we can’t serve which indicates the caller we are in deficit.
Sourcefn publish_held(&mut self)
fn publish_held(&mut self)
Publish any held surplus back to the fast-path.
Trait Implementations§
Source§impl Debug for BandwidthRefiller
impl Debug for BandwidthRefiller
Auto Trait Implementations§
impl !RefUnwindSafe for BandwidthRefiller
impl !UnwindSafe for BandwidthRefiller
impl Freeze for BandwidthRefiller
impl Send for BandwidthRefiller
impl Sync for BandwidthRefiller
impl Unpin for BandwidthRefiller
impl UnsafeUnpin for BandwidthRefiller
Blanket Implementations§
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more