Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/hooks/dynamic_rate_limiter.py: 17%
119 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1# What is this?
2## Allocates dynamic tpm/rpm quota for a project based on current traffic
3## Tracks num active projects per minute
5import asyncio
6import os
7from collections.abc import Callable
8from datetime import datetime
9from typing import Final
11import litellm
12from litellm import ModelResponse, Router
13from litellm._logging import verbose_proxy_logger
14from litellm.caching.caching import DualCache
15from litellm.exceptions import RateLimitType
16from litellm.integrations.custom_logger import CustomLogger
17from litellm.proxy._types import UserAPIKeyAuth
18from litellm.proxy.common_utils.proxy_rate_limit_error import ProxyRateLimitError
19from litellm.proxy.hooks.rate_limiter_utils import (
20 convert_priority_to_percent,
21 resolve_llm_provider_for_rate_limit,
22)
23from litellm.types.router import ModelGroupInfo
24from litellm.types.utils import CallTypesLiteral
25from litellm.utils import get_utc_datetime
28class DynamicRateLimiterCache:
29 """
30 Thin wrapper on DualCache for this file.
32 Track number of active projects calling a model.
33 """
35 def __init__(self, cache: DualCache, time_fn: Callable[[], datetime] = get_utc_datetime) -> None:
36 self.cache = cache
37 self.ttl = 60 # 1 min ttl
38 self.time_fn = time_fn
40 async def async_get_cache(self, model: str) -> int | None:
41 dt: Final = self.time_fn()
42 current_minute: Final = dt.strftime("%H-%M")
43 key_name: Final = f"{current_minute}:{model}"
44 _response: Final = await self.cache.async_get_cache(key=key_name)
45 response: int | None = None
46 if _response is not None:
47 response = len(_response)
48 return response
50 async def async_set_cache_sadd(self, model: str, value: list):
51 """
52 Add value to set.
54 Parameters:
55 - model: str, the name of the model group
56 - value: str, the team id
58 Returns:
59 - None
61 Raises:
62 - Exception, if unable to connect to cache client (if redis caching enabled)
63 """
64 try:
65 dt: Final = self.time_fn()
66 current_minute: Final = dt.strftime("%H-%M")
68 key_name: Final = f"{current_minute}:{model}"
69 await self.cache.async_set_cache_sadd(key=key_name, value=value, ttl=self.ttl)
70 except Exception as e:
71 verbose_proxy_logger.exception(
72 "litellm.proxy.hooks.dynamic_rate_limiter.py::async_set_cache_sadd(): Exception occured - %s", e
73 )
74 raise e
77class _PROXY_DynamicRateLimitHandler(CustomLogger):
78 # Class variables or attributes
79 def __init__(self, internal_usage_cache: DualCache, time_fn: Callable[[], datetime] = get_utc_datetime):
80 self.internal_usage_cache = DynamicRateLimiterCache(cache=internal_usage_cache, time_fn=time_fn)
82 def update_variables(self, llm_router: Router):
83 self.llm_router = llm_router
85 async def check_available_usage(
86 self, model: str, priority: str | None = None
87 ) -> tuple[int | None, int | None, int | None, int | None, int | None]:
88 """
89 For a given model, get its available tpm
91 Params:
92 - model: str, the name of the model in the router model_list
93 - priority: Optional[str], the priority for the request.
95 Returns
96 - Tuple[available_tpm, available_tpm, model_tpm, model_rpm, active_projects]
97 - available_tpm: int or null - always 0 or positive.
98 - available_tpm: int or null - always 0 or positive.
99 - remaining_model_tpm: int or null. If available tpm is int, then this will be too.
100 - remaining_model_rpm: int or null. If available rpm is int, then this will be too.
101 - active_projects: int or null
102 """
103 try:
104 # Get model info first for conversion
105 model_group_info: Final[ModelGroupInfo | None] = self.llm_router.get_model_group_info(model_group=model)
107 weight: float = 1
108 if litellm.priority_reservation is None or priority not in litellm.priority_reservation:
109 verbose_proxy_logger.error(
110 "Priority Reservation not set. priority=%s, but litellm.priority_reservation is %s.",
111 priority,
112 litellm.priority_reservation,
113 )
114 elif priority is not None and litellm.priority_reservation is not None:
115 if os.getenv("LITELLM_LICENSE", None) is None:
116 verbose_proxy_logger.error(
117 "PREMIUM FEATURE: Reserving tpm/rpm by priority is a premium feature. Please add a 'LITELLM_LICENSE' to your .env to enable this.\nGet a license: https://docs.litellm.ai/docs/proxy/enterprise."
118 )
119 else:
120 value: Final = litellm.priority_reservation[priority]
121 weight = convert_priority_to_percent(value, model_group_info)
123 active_projects: Final = await self.internal_usage_cache.async_get_cache(model=model)
124 (
125 current_model_tpm,
126 current_model_rpm,
127 ) = await self.llm_router.get_model_group_usage(model_group=model)
128 total_model_tpm: int | None = None
129 total_model_rpm: int | None = None
130 if model_group_info is not None:
131 if model_group_info.tpm is not None:
132 total_model_tpm = model_group_info.tpm
133 if model_group_info.rpm is not None:
134 total_model_rpm = model_group_info.rpm
136 remaining_model_tpm: int | None = None
137 if total_model_tpm is not None and current_model_tpm is not None:
138 remaining_model_tpm = total_model_tpm - current_model_tpm
139 elif total_model_tpm is not None:
140 remaining_model_tpm = total_model_tpm
142 remaining_model_rpm: int | None = None
143 if total_model_rpm is not None and current_model_rpm is not None:
144 remaining_model_rpm = total_model_rpm - current_model_rpm
145 elif total_model_rpm is not None:
146 remaining_model_rpm = total_model_rpm
148 available_tpm: int | None = None
150 if remaining_model_tpm is not None:
151 if active_projects is not None:
152 available_tpm = int(remaining_model_tpm * weight / active_projects)
153 else:
154 available_tpm = int(remaining_model_tpm * weight)
156 if available_tpm is not None and available_tpm < 0:
157 available_tpm = 0
159 available_rpm: int | None = None
161 if remaining_model_rpm is not None:
162 if active_projects is not None:
163 available_rpm = int(remaining_model_rpm * weight / active_projects)
164 else:
165 available_rpm = int(remaining_model_rpm * weight)
167 if available_rpm is not None and available_rpm < 0:
168 available_rpm = 0
169 return (
170 available_tpm,
171 available_rpm,
172 remaining_model_tpm,
173 remaining_model_rpm,
174 active_projects,
175 )
176 except Exception as e:
177 verbose_proxy_logger.exception(
178 "litellm.proxy.hooks.dynamic_rate_limiter.py::check_available_usage: Exception occurred - %s", e
179 )
180 return None, None, None, None, None
182 async def async_pre_call_hook(
183 self,
184 user_api_key_dict: UserAPIKeyAuth,
185 cache: DualCache,
186 data: dict,
187 call_type: CallTypesLiteral,
188 ) -> (
189 Exception | str | dict | None
190 ): # raise exception if invalid, return a str for the user to receive - if rejected, or return a modified dictionary for passing into litellm
191 """
192 - For a model group
193 - Check if tpm/rpm available
194 - Raise RateLimitError if no tpm/rpm available
195 """
196 if "model" in data:
197 key_priority: Final[str | None] = user_api_key_dict.metadata.get("priority", None)
198 (
199 available_tpm,
200 available_rpm,
201 model_tpm,
202 model_rpm,
203 active_projects,
204 ) = await self.check_available_usage(model=data["model"], priority=key_priority)
205 ### CHECK TPM ###
206 if available_tpm is not None and available_tpm == 0:
207 resolved_model, llm_provider = resolve_llm_provider_for_rate_limit(data.get("model"))
208 raise ProxyRateLimitError(
209 detail={
210 "error": f"Key={user_api_key_dict.api_key} over available TPM={available_tpm}. Model TPM={model_tpm}, Active keys={active_projects}"
211 },
212 rate_limit_type=RateLimitType.TOKENS,
213 model=resolved_model,
214 llm_provider=llm_provider,
215 )
216 ### CHECK RPM ###
217 elif available_rpm is not None and available_rpm == 0:
218 resolved_model, llm_provider = resolve_llm_provider_for_rate_limit(data.get("model"))
219 raise ProxyRateLimitError(
220 detail={
221 "error": f"Key={user_api_key_dict.api_key} over available RPM={available_rpm}. Model RPM={model_rpm}, Active keys={active_projects}"
222 },
223 rate_limit_type=RateLimitType.REQUESTS,
224 model=resolved_model,
225 llm_provider=llm_provider,
226 )
227 elif available_rpm is not None or available_tpm is not None:
228 ## UPDATE CACHE WITH ACTIVE PROJECT
229 asyncio.create_task(
230 self.internal_usage_cache.async_set_cache_sadd( # this is a set
231 model=data["model"],
232 value=[user_api_key_dict.token or "default_key"],
233 )
234 )
235 return None
237 async def async_post_call_success_hook(self, data: dict, user_api_key_dict: UserAPIKeyAuth, response):
238 try:
239 if isinstance(response, ModelResponse):
240 model_info: Final = self.llm_router.get_model_info(id=response._hidden_params["model_id"])
241 assert model_info is not None, "Model info for model with id={} is None".format(
242 response._hidden_params["model_id"]
243 )
244 key_priority: Final[str | None] = user_api_key_dict.metadata.get("priority", None)
245 (
246 available_tpm,
247 available_rpm,
248 model_tpm,
249 model_rpm,
250 active_projects,
251 ) = await self.check_available_usage(model=model_info["model_name"], priority=key_priority)
252 response._hidden_params["additional_headers"] = { # Add additional response headers - easier debugging
253 "x-litellm-model_group": model_info["model_name"],
254 "x-ratelimit-remaining-litellm-project-tokens": available_tpm,
255 "x-ratelimit-remaining-litellm-project-requests": available_rpm,
256 "x-ratelimit-remaining-model-tokens": model_tpm,
257 "x-ratelimit-remaining-model-requests": model_rpm,
258 "x-ratelimit-current-active-projects": active_projects,
259 }
261 return response
262 return await super().async_post_call_success_hook(
263 data=data,
264 user_api_key_dict=user_api_key_dict,
265 response=response,
266 )
267 except Exception as e:
268 verbose_proxy_logger.exception(
269 "litellm.proxy.hooks.dynamic_rate_limiter.py::async_post_call_success_hook(): Exception occured - %s", e
270 )
271 return response