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 theawait. - :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¶
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¶
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¶
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, ...):
...