Skip to content

latent.rate_limiter.wrapper

Provider-agnostic rate-limited wrappers around litellm completion APIs.

Two entry points:

  • :func:rate_limited_completion — async, non-streaming. Slot held for the duration of the await.
  • :func:rate_limited_completion_stream — async generator. Slot held for the entire iteration so open streams still count against the cap.

Both detect litellm.RateLimitError (already normalized across providers — Bedrock ThrottlingException, OpenAI 429, Anthropic 429, Vertex quota errors all surface as the same class), notify the backend, and re-raise so the caller's retry policy (Prefect @task retries, by default) can take over.

The backend doesn't know which LLM provider it's coordinating for; the per-provider concurrency cap comes from :func:latent.rate_limiter.profiles.get_profile.

Functions

is_litellm_available

is_litellm_available() -> bool

True when litellm is importable and the wrappers can dispatch calls.

Use this at lazy-import call sites that need to silently degrade or raise a custom message when litellm isn't installed — checking import latent.rate_limiter alone won't tell you (the package imports successfully even without litellm; the failure surfaces only at call time as ImportError("litellm is required for ...")).

rate_limited_completion

rate_limited_completion(kwargs: Any = {}) -> Any

Non-streaming async rate-limited wrapper around litellm.acompletion.

Drop-in replacement at call sites that do::

response = await rate_limited_completion(model=..., messages=..., ...)

rate_limited_completion_stream

rate_limited_completion_stream(kwargs: Any = {}) -> AsyncIterator[Any]

Streaming async rate-limited wrapper. Yields chunks; slot held for the entire iteration so open streams still count against the per-provider cap.

Drop-in replacement at call sites that do::

response = await acompletion(stream=True, ...)
async for chunk in response:
    ...

becomes::

async for chunk in rate_limited_completion_stream(stream=True, ...):
    ...