Skip to main content

AtomicTokenBucket

Struct AtomicTokenBucket 

Source
pub(super) struct AtomicTokenBucket {
    available: AtomicU64,
    capacity: u64,
}
Expand description

The atomic token bucket minus the clock component.

This is lock-free and the owner needs to refill it with a specific number of tokens it wants to be distributed across many actors.

Fields§

§available: AtomicU64

Every access to these counters is with Ordering::Relaxed as they don’t protect any outside data and so we only care about the atomicity action on the counter. The current token count.

§capacity: u64

The maximum number of tokens the bucket may hold a.k.a the burst.

This is atomic because it can be set during runtime. For instance, a config option reconfigure of a bandwidth rate.

Implementations§

Source§

impl AtomicTokenBucket

Source

pub(super) fn new(capacity: u64) -> Self

A new bucket capped at capacity tokens.

It starts full.

Source

pub(super) fn claim(&self, tokens: u64) -> Option<u64>

Claim all tokens from the bucket or nothing if not enough available.

The request is first clamped to the capacity because the bucket can never hold more than its burst.

Returns Some(claimed) if the bucket held at least the clamped request which indicates that they are now granted to the caller. Note that claimed can be lower than tokens hence why the caller needs to use that value and not what it asked for.

Returns None otherwise, nothing is taken.

Source

pub(super) fn refill(&self, tokens: u64)

Add tokens to the bucket which is capped at the capacity.

We use a CAS loop (Compare-And-Swap) because we need to cap the refill to the internal capacity atomically. A fetch_add + fetch_sub is not possible in order to refill atomically because of this race:

  capacity = 100, available = 90:
    refill: fetch_add(50)  -> available = 140 (above capacity)
    claim:  fetch_sub(140) -> available = 0   (illegal claim, above capacity)
    refill: fetch_sub(40)  -> available underflows (adjust overshoot too late)

Any tokens that overshoot the capacity are forfeited.

Finally, the CAS loop is considered ok because the refill is not in the fast path and should work most of the time on the first iteration.

Source

pub(super) fn drain(&self) -> u64

Take every token out of the bucket and return the amount.

Source

pub(super) fn capacity(&self) -> u64

Return the capacity.

Trait Implementations§

Source§

impl Debug for AtomicTokenBucket

Source§

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

Formats the value using the given formatter. 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