Skip to main content

BandwidthRefiller

Struct BandwidthRefiller 

Source
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: u64

Tokens 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

Source

pub(super) fn new( bucket: Arc<AtomicTokenBucket>, rx: UnboundedReceiver<(u64, Arc<RefillWaiter>)>, ) -> Self

Constructor.

Source

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.

Source

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.

Source

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.

Source

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.

Source

fn publish_held(&mut self)

Publish any held surplus back to the fast-path.

Trait Implementations§

Source§

impl Debug for BandwidthRefiller

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Drop for BandwidthRefiller

Source§

fn drop(&mut self)

Wake all queued waiters on teardown so they wake up and get to realize the pool is closed on their next poll.

Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts 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
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more