Skip to content

rate_limiter

langroid/language_models/rate_limiter.py

Pro-active rate limiting for OpenAI-compatible chat-completion calls.

This module is standalone: it knows nothing about Langroid agents, and it does not touch the retry-with-exponential-backoff logic in langroid.language_models.utils, which stays as the reactive fallback.

The limiter discovers the account's actual limits instead of asking the user to configure them. OpenAI (and most OpenAI-compatible gateways) report the live budget on every response::

x-ratelimit-limit-requests / x-ratelimit-remaining-requests
x-ratelimit-limit-tokens   / x-ratelimit-remaining-tokens
x-ratelimit-reset-requests / x-ratelimit-reset-tokens

reset is the time until the bucket refills to limit, so the sustained refill rate implied by a single response is::

rate = (limit - remaining) / reset

which is exactly the account's limit for that model, in requests (or tokens) per second. The limiter paces sends at rate * (1 - headroom) so a request is held briefly rather than rejected.

For providers that send no such headers (local servers, Groq/Cerebras, litellm) the limiter falls back to an additive-free AIMD scheme: multiply the inter-send interval up on a 429, decay it down on every success. No prior knowledge of any limit is needed.

Limiters are shared process-wide by key (see get_rate_limiter), so the many cloned agents of a run_batch_tasks job pace against one budget.

Known limitation: until the first response arrives there are no headers to learn from, so an initial burst of concurrent calls is sent unpaced; pacing engages from the first observed response onwards.

RateLimitSnapshot

Bases: BaseModel

One provider-reported view of the remaining rate-limit budget.

has_rate_limit_info property

Did the provider report any rate-limit budget at all?

request_refill_rate property

Implied sustained limit, in requests/second (None if unknowable).

token_refill_rate property

Implied sustained limit, in tokens/second (None if unknowable).

from_headers(headers) classmethod

Build a snapshot from HTTP response headers (case-insensitive).

Source code in langroid/language_models/rate_limiter.py
@classmethod
def from_headers(cls, headers: Mapping[str, Any]) -> "RateLimitSnapshot":
    """Build a snapshot from HTTP response headers (case-insensitive)."""
    lowered = {str(k).lower(): v for k, v in dict(headers).items()}

    def get(name: str) -> Optional[str]:
        raw = lowered.get(name)
        return None if raw is None else str(raw)

    return cls(
        limit_requests=_parse_int(get(LIMIT_REQUESTS_HEADER)),
        remaining_requests=_parse_int(get(REMAINING_REQUESTS_HEADER)),
        reset_requests=parse_reset_duration(get(RESET_REQUESTS_HEADER)),
        limit_tokens=_parse_int(get(LIMIT_TOKENS_HEADER)),
        remaining_tokens=_parse_int(get(REMAINING_TOKENS_HEADER)),
        reset_tokens=parse_reset_duration(get(RESET_TOKENS_HEADER)),
        retry_after=parse_reset_duration(get(RETRY_AFTER_HEADER)),
    )

RateLimitConfig

Bases: BaseSettings

Settings for the pro-active rate limiter.

Every field can be overridden by an env var with the LANGROID_RATE_LIMIT_ prefix, e.g. LANGROID_RATE_LIMIT_ENABLED=1.

RateLimiter(config=None, name='')

Paces sends so provider rate limits are approached, not exceeded.

Thread-safe and asyncio-safe: state is guarded by a threading.Lock held only for the (non-blocking) bookkeeping, while the wait itself happens outside the lock via time.sleep or asyncio.sleep.

Slots are handed out as reservations off a shared monotonic clock, so N concurrent callers queue rather than all retrying at once.

Source code in langroid/language_models/rate_limiter.py
def __init__(
    self, config: Optional[RateLimitConfig] = None, name: str = ""
) -> None:
    self.config = config or RateLimitConfig()
    self.name = name
    self._lock = threading.Lock()
    self._next_send_at = 0.0  # monotonic time of the next free slot
    # Hard cooldown floor: no caller may send before this, including one
    # that is already asleep on an earlier reservation.
    self._gate_until = 0.0
    self._request_rate: Optional[float] = None  # requests/sec, from headers
    self._token_rate: Optional[float] = None  # tokens/sec, from headers
    self._avg_tokens: Optional[float] = None  # EWMA tokens per request
    self._fallback_interval = 0.0  # header-free AIMD interval
    self._seen_headers = False
    self._sends = 0
    self._waits = 0
    self._total_wait = 0.0
    self._rate_limit_errors = 0
    self._capped_waits = 0

acquire()

Block until it is this caller's turn to send. Returns seconds waited.

Source code in langroid/language_models/rate_limiter.py
def acquire(self) -> float:
    """Block until it is this caller's turn to send. Returns seconds waited."""
    waited = 0.0
    for wait in self._wait_steps():
        time.sleep(wait)
        waited += wait
    return waited

acquire_async() async

Async variant of acquire. Returns seconds waited.

Source code in langroid/language_models/rate_limiter.py
async def acquire_async(self) -> float:
    """Async variant of `acquire`. Returns seconds waited."""
    waited = 0.0
    for wait in self._wait_steps():
        await asyncio.sleep(wait)
        waited += wait
    return waited

observe_response(headers=None, tokens_used=None)

Take in what a response reveals about the budget.

Information only: callers may call this more than once per request (the headers arrive with the response, a streaming request's token usage only later), so it must NOT be where the AIMD recovery is applied -- see observe_success, which is called exactly once per request.

Parameters:

Name Type Description Default
headers Optional[Mapping[str, Any]]

Response headers, if the provider/transport exposes them.

None
tokens_used Optional[int]

Total tokens billed for this request, if reported. Used as an EWMA estimate of the per-request token cost, which is what makes token-budget pacing possible. It is an estimate, not a guarantee of provider quota compliance.

None
Source code in langroid/language_models/rate_limiter.py
def observe_response(
    self,
    headers: Optional[Mapping[str, Any]] = None,
    tokens_used: Optional[int] = None,
) -> None:
    """Take in what a response reveals about the budget.

    Information only: callers may call this more than once per request (the
    headers arrive with the response, a streaming request's token usage
    only later), so it must NOT be where the AIMD recovery is applied --
    see `observe_success`, which is called exactly once per request.

    Args:
        headers: Response headers, if the provider/transport exposes them.
        tokens_used: Total tokens billed for this request, if reported.
            Used as an EWMA estimate of the per-request token cost, which
            is what makes token-budget pacing possible. It is an estimate,
            not a guarantee of provider quota compliance.
    """
    snap = RateLimitSnapshot.from_headers(headers) if headers is not None else None
    with self._lock:
        if snap is not None:
            self._apply_snapshot_locked(snap)
        if tokens_used is not None and tokens_used > 0:
            if self._avg_tokens is None:
                self._avg_tokens = float(tokens_used)
            else:
                self._avg_tokens = 0.7 * self._avg_tokens + 0.3 * tokens_used

observe_success(tokens_used=None)

Record that one request was accepted; recover the send rate.

Call this EXACTLY once per request that the provider accepted. It is the only place the header-free AIMD interval decays, so calling it twice for one request would halve the backoff twice over.

Source code in langroid/language_models/rate_limiter.py
def observe_success(self, tokens_used: Optional[int] = None) -> None:
    """Record that one request was accepted; recover the send rate.

    Call this EXACTLY once per request that the provider accepted. It is
    the only place the header-free AIMD interval decays, so calling it
    twice for one request would halve the backoff twice over.
    """
    self.observe_response(tokens_used=tokens_used)
    with self._lock:
        if self._fallback_interval > 0:
            decayed = self._fallback_interval * self.config.recovery_factor
            self._fallback_interval = 0.0 if decayed < 1e-3 else decayed

observe_rate_limit_error(headers=None)

Learn from a 429: back off, and stall until the budget recovers.

Source code in langroid/language_models/rate_limiter.py
def observe_rate_limit_error(
    self, headers: Optional[Mapping[str, Any]] = None
) -> None:
    """Learn from a 429: back off, and stall until the budget recovers."""
    snap = RateLimitSnapshot.from_headers(headers) if headers is not None else None
    with self._lock:
        self._rate_limit_errors += 1
        hold = 0.0
        if snap is not None:
            self._apply_snapshot_locked(snap)
            hold = max(
                snap.retry_after or 0.0,
                self._refill_wait(snap, tokens=False),
                self._refill_wait(snap, tokens=True),
            )
        if hold <= 0:
            # The response says nothing actionable: no headers, too few to
            # derive a rate from, or a 429 whose budgets both read healthy.
            # Multiplicatively reduce the send rate instead, and recover it
            # on each success.
            self._fallback_interval = min(
                self.config.max_interval,
                max(
                    self.config.error_interval,
                    self._fallback_interval * self.config.backoff_factor,
                ),
            )
        self._stall_locked(
            min(max(hold, self._fallback_interval), self.config.max_wait)
        )
        logger.debug(
            "rate limiter %s: 429 observed, interval now %.4fs",
            self.name,
            self._interval_locked(),
        )

stats()

Counters and discovered state, for tests and diagnostics.

sends is the "did it run at all" signal: zero means the limiter was never consulted, which is different from "it was consulted and never needed to wait" (sends > 0, waits == 0).

Source code in langroid/language_models/rate_limiter.py
def stats(self) -> Dict[str, Any]:
    """Counters and discovered state, for tests and diagnostics.

    `sends` is the "did it run at all" signal: zero means the limiter was
    never consulted, which is different from "it was consulted and never
    needed to wait" (`sends > 0, waits == 0`).
    """
    with self._lock:
        return dict(
            name=self.name,
            sends=self._sends,
            waits=self._waits,
            total_wait=self._total_wait,
            capped_waits=self._capped_waits,
            rate_limit_errors=self._rate_limit_errors,
            seen_headers=self._seen_headers,
            request_rate=self._request_rate,
            token_rate=self._token_rate,
            avg_tokens=self._avg_tokens,
            fallback_interval=self._fallback_interval,
            interval=self._interval_locked(),
            gate_remaining=max(0.0, self._gate_until - time.monotonic()),
        )

parse_reset_duration(value)

Parse an OpenAI x-ratelimit-reset-* duration into seconds.

The header is a Go-style duration, e.g. "120ms", "1s", "6m0s", "1h2m3s", "7.66s". A bare number is read as seconds, which is what some OpenAI-compatible gateways send.

Parameters:

Name Type Description Default
value Optional[str]

Raw header value, or None.

required

Returns:

Type Description
Optional[float]

Duration in seconds, or None if value is absent or unparseable.

Source code in langroid/language_models/rate_limiter.py
def parse_reset_duration(value: Optional[str]) -> Optional[float]:
    """Parse an OpenAI `x-ratelimit-reset-*` duration into seconds.

    The header is a Go-style duration, e.g. `"120ms"`, `"1s"`, `"6m0s"`,
    `"1h2m3s"`, `"7.66s"`. A bare number is read as seconds, which is what
    some OpenAI-compatible gateways send.

    Args:
        value: Raw header value, or None.

    Returns:
        Duration in seconds, or None if `value` is absent or unparseable.
    """
    if value is None:
        return None
    text = value.strip().lower()
    if not text:
        return None
    total = 0.0
    matched = False
    i = 0
    n = len(text)
    while i < n:
        start = i
        while i < n and (text[i].isdigit() or text[i] == "."):
            i += 1
        if i == start:
            return None  # unexpected character where a number was expected
        try:
            amount = float(text[start:i])
        except ValueError:
            return None
        unit_start = i
        while i < n and text[i].isalpha():
            i += 1
        unit = text[unit_start:i]
        if unit == "":
            # bare number: seconds
            total += amount
            matched = True
            continue
        if unit not in _UNIT_SECONDS:
            return None
        total += amount * _UNIT_SECONDS[unit]
        matched = True
    return total if matched else None

get_rate_limiter(key, config=None)

Get (creating if needed) the process-wide limiter for key.

Sharing by key is what lets the cloned agents of a batch job pace against one budget. The config of the first caller for a given key wins; later callers get the existing limiter unchanged.

Parameters:

Name Type Description Default
key str

Sharing key, e.g. "https://api.openai.com/v1::gpt-4o-mini".

required
config Optional[RateLimitConfig]

Settings to use if the limiter does not exist yet.

None

Returns:

Type Description
RateLimiter

The shared RateLimiter for key.

Source code in langroid/language_models/rate_limiter.py
def get_rate_limiter(key: str, config: Optional[RateLimitConfig] = None) -> RateLimiter:
    """Get (creating if needed) the process-wide limiter for `key`.

    Sharing by key is what lets the cloned agents of a batch job pace against
    one budget. The config of the *first* caller for a given key wins; later
    callers get the existing limiter unchanged.

    Args:
        key: Sharing key, e.g. `"https://api.openai.com/v1::gpt-4o-mini"`.
        config: Settings to use if the limiter does not exist yet.

    Returns:
        The shared `RateLimiter` for `key`.
    """
    with _registry_lock:
        limiter = _limiters.get(key)
        if limiter is None:
            limiter = RateLimiter(config=config, name=key)
            _limiters[key] = limiter
        return limiter

reset_rate_limiters()

Drop all shared limiters (test/teardown helper).

Source code in langroid/language_models/rate_limiter.py
def reset_rate_limiters() -> None:
    """Drop all shared limiters (test/teardown helper)."""
    with _registry_lock:
        _limiters.clear()

rate_limit_error_headers(exc)

Classify an exception as a rate-limit error and return its headers.

Returns:

Type Description
Optional[Mapping[str, Any]]

None if exc is not a rate-limit (429) error; otherwise the response

Optional[Mapping[str, Any]]

headers, which may be an empty mapping if the provider sent none.

Optional[Mapping[str, Any]]

An empty mapping is therefore meaningfully different from None.

Source code in langroid/language_models/rate_limiter.py
def rate_limit_error_headers(exc: BaseException) -> Optional[Mapping[str, Any]]:
    """Classify an exception as a rate-limit error and return its headers.

    Returns:
        None if `exc` is not a rate-limit (429) error; otherwise the response
        headers, which may be an empty mapping if the provider sent none.
        An empty mapping is therefore meaningfully different from None.
    """
    status = getattr(exc, "status_code", None)
    if status is None:
        status = getattr(getattr(exc, "response", None), "status_code", None)
    is_rate_limit = status == 429
    if not is_rate_limit:
        # litellm and some SDKs signal rate limits by exception class name
        # without exposing a status code.
        is_rate_limit = type(exc).__name__ in (
            "RateLimitError",
            "APIRateLimitError",
        )
    if not is_rate_limit:
        return None
    headers = getattr(getattr(exc, "response", None), "headers", None)
    if headers is None:
        headers = getattr(exc, "headers", None)
    if headers is None:
        return {}
    try:
        return dict(headers)
    except Exception:
        return {}