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

1# What is this? 

2## Allocates dynamic tpm/rpm quota for a project based on current traffic 

3## Tracks num active projects per minute 

4 

5import asyncio 

6import os 

7from collections.abc import Callable 

8from datetime import datetime 

9from typing import Final 

10 

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 

26 

27 

28class DynamicRateLimiterCache: 

29 """ 

30 Thin wrapper on DualCache for this file. 

31 

32 Track number of active projects calling a model. 

33 """ 

34 

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 

39 

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 

49 

50 async def async_set_cache_sadd(self, model: str, value: list): 

51 """ 

52 Add value to set. 

53 

54 Parameters: 

55 - model: str, the name of the model group 

56 - value: str, the team id 

57 

58 Returns: 

59 - None 

60 

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") 

67 

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 

75 

76 

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) 

81 

82 def update_variables(self, llm_router: Router): 

83 self.llm_router = llm_router 

84 

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 

90 

91 Params: 

92 - model: str, the name of the model in the router model_list 

93 - priority: Optional[str], the priority for the request. 

94 

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) 

106 

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) 

122 

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 

135 

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 

141 

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 

147 

148 available_tpm: int | None = None 

149 

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) 

155 

156 if available_tpm is not None and available_tpm < 0: 

157 available_tpm = 0 

158 

159 available_rpm: int | None = None 

160 

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) 

166 

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 

181 

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 

236 

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 } 

260 

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