Skip to content

core

Token Bucket Rate Limiter Implementation

AsyncTokenBucket

Asynchronous Leaky Bucket Rate Limiter

Parameters:

Name Type Description Default
bucket_config BucketConfig | None

Configuration for the token bucket with the max capacity and time period in seconds

None
max_concurrent int | None

Maximum number of concurrent requests allowed to acquire capacity

None
Note

This implementation is synchronous and supports bursts up to the capacity within the specified time period

Source code in limitor/token_bucket/core.py
class AsyncTokenBucket:
    """Asynchronous Leaky Bucket Rate Limiter

    Args:
        bucket_config: Configuration for the token bucket with the max capacity and time period in seconds
        max_concurrent: Maximum number of concurrent requests allowed to acquire capacity

    Note:
        This implementation is synchronous and supports bursts up to the capacity within the specified time period
    """

    def __init__(
        self,
        bucket_config: BucketConfig | None = None,
        max_concurrent: int | None = None,
    ):
        config = bucket_config or BucketConfig()
        self.capacity = config.capacity
        self.seconds = config.seconds

        self.fill_rate = self.capacity / self.seconds

        self._bucket_level = self.capacity
        self._last_fill = time.monotonic()

        self.max_concurrent = max_concurrent
        self._lock = asyncio.Lock()
        self._semaphore: asyncio.Semaphore | AbstractAsyncContextManager[Any] | None = (
            None
        )

    def _fill(self) -> None:
        """Fill the bucket based on the elapsed time since the last fill"""
        now = time.monotonic()
        elapsed = now - self._last_fill
        self._bucket_level = min(
            self.capacity, self._bucket_level + elapsed * self.fill_rate
        )
        self._last_fill = now

    def capacity_info(self, amount: float = 1) -> Capacity:
        """Get the current capacity information of the token bucket

        Args:
            amount: The amount of capacity to check for, defaults to 1

        Returns:
            A named tuple indicating if the bucket has enough capacity and how much more is needed
        """
        self._fill()
        # we need at least `amount` tokens to proceed
        needed = amount - self._bucket_level
        return Capacity(has_capacity=needed <= 0, needed_capacity=needed)

    async def _acquire_logic(self, amount: float = 1) -> None:
        """Core logic for acquiring capacity from the token bucket.

        Args:
            amount: The amount of capacity to check for, defaults to 1

        Notes:
            Adding a lock here ensures that the acquire logic is atomic, but it also means that the
                requests are going to be done in the order they were received  i.e. not out-of-order like
                most async programs.
            The benefit is that with multiple concurrent requests, we can ensure that the bucket level
                is updated correctly and that we don't have multiple requests trying to update the bucket level
                at the same time, which could lead to an inconsistent state i.e. a race condition.
        """
        async with (
            self._lock
        ):  # ensures atomicity given we can have multiple concurrent requests
            capacity_info = self.capacity_info(amount=amount)
            while not capacity_info.has_capacity:
                needed = capacity_info.needed_capacity
                # amount we need to wait to leak (either part or the entire capacity)
                # needed is guaranteed to be positive here, so we can use it directly
                wait_time = needed / self.fill_rate
                if wait_time > 0:
                    await asyncio.sleep(wait_time)

                capacity_info = self.capacity_info(amount=amount)

            self._bucket_level -= amount

    async def _semaphore_acquire(self, amount: float = 1) -> None:
        """Acquire capacity using a semaphore to limit concurrency.

        Args:
            amount: The amount of capacity to acquire, defaults to 1
        """
        if self._semaphore is None:
            self._semaphore = (
                asyncio.Semaphore(self.max_concurrent)
                if self.max_concurrent
                else nullcontext()
            )

        async with self._semaphore:
            await self._acquire_logic(amount)

    async def acquire(self, amount: float = 1, timeout: float | None = None) -> None:
        """Acquire capacity from the token bucket, waiting asynchronously until allowed.

        Supports timeouts and cancellations.

        Args:
            amount: The amount of capacity to acquire, defaults to 1
            timeout: Optional timeout in seconds for the acquire operation

        Raises:
            TimeoutError: If the acquire operation times out after the specified timeout period
        """
        validate_amount(self, amount=amount)

        if timeout is not None:
            try:
                await asyncio.wait_for(self._semaphore_acquire(amount), timeout=timeout)
            except TimeoutError as error:
                raise TimeoutError(
                    f"Acquire timed out after {timeout} seconds for amount={amount}"
                ) from error
        else:
            await self._semaphore_acquire(amount)

    @asynccontextmanager
    async def acquire_ctx(
        self, amount: float = 1, timeout: float | None = None
    ) -> AsyncGenerator[AsyncTokenBucket]:
        """Async context manager that acquires a specific amount of capacity.

        Args:
            amount: The amount of capacity to acquire, defaults to 1
            timeout: Optional timeout in seconds for the acquire operation

        Yields:
            AsyncGenerator[AsyncTokenBucket]: The AsyncTokenBucket instance
        """
        await self.acquire(amount=amount, timeout=timeout)
        try:
            yield self
        finally:
            pass

    async def reconcile(self, actual: float, estimated: float) -> None:
        """Reconcile the actual tokens/capacity used against the estimated amount.

        Args:
            actual: The actual amount of tokens/capacity used
            estimated: The estimated amount of tokens/capacity previously acquired
        """
        diff = actual - estimated
        if diff > 0:
            await self.acquire(amount=diff)
        elif diff < 0:
            async with self._lock:
                self._fill()
                self._bucket_level = min(self.capacity, self._bucket_level - diff)

    async def __aenter__(self) -> AsyncTokenBucket:
        """Enter the context manager, acquiring resources if necessary"""
        await self.acquire()
        return self

    async def __aexit__(
        self,
        exc_type: type[BaseException] | None,
        exc_val: BaseException | None,
        exc_tb: TracebackType | None,
    ) -> None:
        """Exit the context manager, releasing any resources if necessary"""
        return None

__aenter__() async

Enter the context manager, acquiring resources if necessary

Source code in limitor/token_bucket/core.py
async def __aenter__(self) -> AsyncTokenBucket:
    """Enter the context manager, acquiring resources if necessary"""
    await self.acquire()
    return self

__aexit__(exc_type, exc_val, exc_tb) async

Exit the context manager, releasing any resources if necessary

Source code in limitor/token_bucket/core.py
async def __aexit__(
    self,
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: TracebackType | None,
) -> None:
    """Exit the context manager, releasing any resources if necessary"""
    return None

acquire(amount=1, timeout=None) async

Acquire capacity from the token bucket, waiting asynchronously until allowed.

Supports timeouts and cancellations.

Parameters:

Name Type Description Default
amount float

The amount of capacity to acquire, defaults to 1

1
timeout float | None

Optional timeout in seconds for the acquire operation

None

Raises:

Type Description
TimeoutError

If the acquire operation times out after the specified timeout period

Source code in limitor/token_bucket/core.py
async def acquire(self, amount: float = 1, timeout: float | None = None) -> None:
    """Acquire capacity from the token bucket, waiting asynchronously until allowed.

    Supports timeouts and cancellations.

    Args:
        amount: The amount of capacity to acquire, defaults to 1
        timeout: Optional timeout in seconds for the acquire operation

    Raises:
        TimeoutError: If the acquire operation times out after the specified timeout period
    """
    validate_amount(self, amount=amount)

    if timeout is not None:
        try:
            await asyncio.wait_for(self._semaphore_acquire(amount), timeout=timeout)
        except TimeoutError as error:
            raise TimeoutError(
                f"Acquire timed out after {timeout} seconds for amount={amount}"
            ) from error
    else:
        await self._semaphore_acquire(amount)

acquire_ctx(amount=1, timeout=None) async

Async context manager that acquires a specific amount of capacity.

Parameters:

Name Type Description Default
amount float

The amount of capacity to acquire, defaults to 1

1
timeout float | None

Optional timeout in seconds for the acquire operation

None

Yields:

Type Description
AsyncGenerator[AsyncTokenBucket]

AsyncGenerator[AsyncTokenBucket]: The AsyncTokenBucket instance

Source code in limitor/token_bucket/core.py
@asynccontextmanager
async def acquire_ctx(
    self, amount: float = 1, timeout: float | None = None
) -> AsyncGenerator[AsyncTokenBucket]:
    """Async context manager that acquires a specific amount of capacity.

    Args:
        amount: The amount of capacity to acquire, defaults to 1
        timeout: Optional timeout in seconds for the acquire operation

    Yields:
        AsyncGenerator[AsyncTokenBucket]: The AsyncTokenBucket instance
    """
    await self.acquire(amount=amount, timeout=timeout)
    try:
        yield self
    finally:
        pass

capacity_info(amount=1)

Get the current capacity information of the token bucket

Parameters:

Name Type Description Default
amount float

The amount of capacity to check for, defaults to 1

1

Returns:

Type Description
Capacity

A named tuple indicating if the bucket has enough capacity and how much more is needed

Source code in limitor/token_bucket/core.py
def capacity_info(self, amount: float = 1) -> Capacity:
    """Get the current capacity information of the token bucket

    Args:
        amount: The amount of capacity to check for, defaults to 1

    Returns:
        A named tuple indicating if the bucket has enough capacity and how much more is needed
    """
    self._fill()
    # we need at least `amount` tokens to proceed
    needed = amount - self._bucket_level
    return Capacity(has_capacity=needed <= 0, needed_capacity=needed)

reconcile(actual, estimated) async

Reconcile the actual tokens/capacity used against the estimated amount.

Parameters:

Name Type Description Default
actual float

The actual amount of tokens/capacity used

required
estimated float

The estimated amount of tokens/capacity previously acquired

required
Source code in limitor/token_bucket/core.py
async def reconcile(self, actual: float, estimated: float) -> None:
    """Reconcile the actual tokens/capacity used against the estimated amount.

    Args:
        actual: The actual amount of tokens/capacity used
        estimated: The estimated amount of tokens/capacity previously acquired
    """
    diff = actual - estimated
    if diff > 0:
        await self.acquire(amount=diff)
    elif diff < 0:
        async with self._lock:
            self._fill()
            self._bucket_level = min(self.capacity, self._bucket_level - diff)

SyncTokenBucket

Token Bucket Rate Limiter

Parameters:

Name Type Description Default
bucket_config BucketConfig | None

Configuration for the token bucket with the max capacity and time period in seconds

None
Note

This implementation is synchronous and supports bursts up to the capacity within the specified time period

Source code in limitor/token_bucket/core.py
class SyncTokenBucket:
    """Token Bucket Rate Limiter

    Args:
        bucket_config: Configuration for the token bucket with the max capacity and time period in seconds

    Note:
        This implementation is synchronous and supports bursts up to the capacity within the specified time period
    """

    def __init__(self, bucket_config: BucketConfig | None = None):
        # import config and set attributes
        config = bucket_config or BucketConfig()
        self.capacity = config.capacity
        self.seconds = config.seconds

        self.fill_rate = self.capacity / self.seconds  # units per second

        self._bucket_level = self.capacity  # current volume of tokens in the bucket
        self._last_fill = time.monotonic()  # last refill time

        # thread-safe
        self._lock = threading.Lock()

    def _fill(self) -> None:
        """Fill the bucket based on the elapsed time since the last fill"""
        now = time.monotonic()
        elapsed = now - self._last_fill
        self._bucket_level = min(
            self.capacity, self._bucket_level + elapsed * self.fill_rate
        )
        self._last_fill = now

    def capacity_info(self, amount: float = 1) -> Capacity:
        """Get the current capacity information of the token bucket

        Args:
            amount: The amount of capacity to check for, defaults to 1

        Returns:
            A named tuple indicating if the bucket has enough capacity and how much more is needed
        """
        self._fill()
        # we need at least `amount` tokens to proceed
        needed = amount - self._bucket_level
        return Capacity(has_capacity=needed <= 0, needed_capacity=needed)

    def acquire(self, amount: float = 1) -> None:
        """Acquire capacity from the token bucket, blocking until enough capacity is available.

        This method will block and sleep until the requested amount can be acquired
        without exceeding the bucket's capacity, simulating rate limiting.

        Args:
            amount: The amount of capacity to acquire, defaults to 1

        Notes:
            The while loop is just to make sure nothing funny happens while waiting
        """
        validate_amount(self, amount=amount)

        with self._lock:
            capacity_info = self.capacity_info(amount=amount)
            while not capacity_info.has_capacity:
                needed = capacity_info.needed_capacity
                # amount we need to wait to leak
                # needed is guaranteed to be positive here, so we can use it directly
                wait_time = needed / self.fill_rate
                if wait_time > 0:
                    time.sleep(wait_time)

                capacity_info = self.capacity_info(amount=amount)

            self._bucket_level -= amount

    @contextmanager
    def acquire_ctx(self, amount: float = 1) -> Generator[SyncTokenBucket]:
        """Context manager that acquires a specific amount of capacity.

        Args:
            amount: The amount of capacity to acquire, defaults to 1

        Yields:
            Generator[SyncTokenBucket]: The SyncTokenBucket instance
        """
        self.acquire(amount=amount)
        try:
            yield self
        finally:
            pass

    def reconcile(self, actual: float, estimated: float) -> None:
        """Reconcile the actual tokens/capacity used against the estimated amount.

        Args:
            actual: The actual amount of tokens/capacity used
            estimated: The estimated amount of tokens/capacity previously acquired
        """
        diff = actual - estimated
        if diff > 0:
            self.acquire(amount=diff)
        elif diff < 0:
            with self._lock:
                self._fill()
                self._bucket_level = min(self.capacity, self._bucket_level - diff)

    def __enter__(self) -> SyncTokenBucket:
        """Enter the context manager, acquiring resources if necessary"""
        self.acquire()
        return self

    def __exit__(
        self,
        exc_type: type[BaseException] | None,
        exc_val: BaseException | None,
        exc_tb: TracebackType | None,
    ) -> None:
        """Exit the context manager, releasing any resources if necessary"""
        return None

__enter__()

Enter the context manager, acquiring resources if necessary

Source code in limitor/token_bucket/core.py
def __enter__(self) -> SyncTokenBucket:
    """Enter the context manager, acquiring resources if necessary"""
    self.acquire()
    return self

__exit__(exc_type, exc_val, exc_tb)

Exit the context manager, releasing any resources if necessary

Source code in limitor/token_bucket/core.py
def __exit__(
    self,
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: TracebackType | None,
) -> None:
    """Exit the context manager, releasing any resources if necessary"""
    return None

acquire(amount=1)

Acquire capacity from the token bucket, blocking until enough capacity is available.

This method will block and sleep until the requested amount can be acquired without exceeding the bucket's capacity, simulating rate limiting.

Parameters:

Name Type Description Default
amount float

The amount of capacity to acquire, defaults to 1

1
Notes

The while loop is just to make sure nothing funny happens while waiting

Source code in limitor/token_bucket/core.py
def acquire(self, amount: float = 1) -> None:
    """Acquire capacity from the token bucket, blocking until enough capacity is available.

    This method will block and sleep until the requested amount can be acquired
    without exceeding the bucket's capacity, simulating rate limiting.

    Args:
        amount: The amount of capacity to acquire, defaults to 1

    Notes:
        The while loop is just to make sure nothing funny happens while waiting
    """
    validate_amount(self, amount=amount)

    with self._lock:
        capacity_info = self.capacity_info(amount=amount)
        while not capacity_info.has_capacity:
            needed = capacity_info.needed_capacity
            # amount we need to wait to leak
            # needed is guaranteed to be positive here, so we can use it directly
            wait_time = needed / self.fill_rate
            if wait_time > 0:
                time.sleep(wait_time)

            capacity_info = self.capacity_info(amount=amount)

        self._bucket_level -= amount

acquire_ctx(amount=1)

Context manager that acquires a specific amount of capacity.

Parameters:

Name Type Description Default
amount float

The amount of capacity to acquire, defaults to 1

1

Yields:

Type Description
Generator[SyncTokenBucket]

Generator[SyncTokenBucket]: The SyncTokenBucket instance

Source code in limitor/token_bucket/core.py
@contextmanager
def acquire_ctx(self, amount: float = 1) -> Generator[SyncTokenBucket]:
    """Context manager that acquires a specific amount of capacity.

    Args:
        amount: The amount of capacity to acquire, defaults to 1

    Yields:
        Generator[SyncTokenBucket]: The SyncTokenBucket instance
    """
    self.acquire(amount=amount)
    try:
        yield self
    finally:
        pass

capacity_info(amount=1)

Get the current capacity information of the token bucket

Parameters:

Name Type Description Default
amount float

The amount of capacity to check for, defaults to 1

1

Returns:

Type Description
Capacity

A named tuple indicating if the bucket has enough capacity and how much more is needed

Source code in limitor/token_bucket/core.py
def capacity_info(self, amount: float = 1) -> Capacity:
    """Get the current capacity information of the token bucket

    Args:
        amount: The amount of capacity to check for, defaults to 1

    Returns:
        A named tuple indicating if the bucket has enough capacity and how much more is needed
    """
    self._fill()
    # we need at least `amount` tokens to proceed
    needed = amount - self._bucket_level
    return Capacity(has_capacity=needed <= 0, needed_capacity=needed)

reconcile(actual, estimated)

Reconcile the actual tokens/capacity used against the estimated amount.

Parameters:

Name Type Description Default
actual float

The actual amount of tokens/capacity used

required
estimated float

The estimated amount of tokens/capacity previously acquired

required
Source code in limitor/token_bucket/core.py
def reconcile(self, actual: float, estimated: float) -> None:
    """Reconcile the actual tokens/capacity used against the estimated amount.

    Args:
        actual: The actual amount of tokens/capacity used
        estimated: The estimated amount of tokens/capacity previously acquired
    """
    diff = actual - estimated
    if diff > 0:
        self.acquire(amount=diff)
    elif diff < 0:
        with self._lock:
            self._fill()
            self._bucket_level = min(self.capacity, self._bucket_level - diff)