Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/hooks/dynamic_rate_limiter_v3.py: 12%

236 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:01 +0000

1""" 

2Dynamic rate limiter v3 - Saturation-aware priority-based rate limiting 

3""" 

4 

5import os 

6from collections.abc import Callable 

7from datetime import datetime 

8from typing import TYPE_CHECKING, Final, Literal 

9 

10from fastapi import HTTPException 

11 

12import litellm 

13from litellm import ModelResponse, Router 

14from litellm._logging import verbose_proxy_logger 

15from litellm.caching.caching import DualCache 

16from litellm.integrations.custom_logger import CustomLogger 

17from litellm.proxy._types import UserAPIKeyAuth 

18from litellm.proxy.common_utils.proxy_rate_limit_error import ( 

19 ProxyRateLimitError, 

20 map_v3_rate_limit_type, 

21) 

22from litellm.proxy.hooks.parallel_request_limiter_v3 import ( 

23 RateLimitDescriptor, 

24 RateLimitDescriptorRateLimitObject, 

25 RateLimitResponse, 

26 _PROXY_MaxParallelRequestsHandler_v3, 

27 claim_request_stash_for_data, 

28 get_or_create_request_stash, 

29) 

30from litellm.proxy.hooks.rate_limiter_utils import ( 

31 convert_priority_to_percent, 

32 resolve_llm_provider_for_rate_limit, 

33) 

34from litellm.proxy.utils import InternalUsageCache 

35from litellm.router_utils.add_retry_fallback_headers import ( 

36 ensure_response_additional_headers, 

37 response_has_hidden_params, 

38) 

39from litellm.types.router import ModelGroupInfo 

40from litellm.types.utils import CallTypesLiteral 

41 

42if TYPE_CHECKING: 42 ↛ 43line 42 didn't jump to line 43 because the condition on line 42 was never true

43 from litellm.types.utils import PriorityReservationSettings 

44 

45 

46def _get_priority_settings() -> "PriorityReservationSettings": 

47 """ 

48 Get the priority reservation settings, guaranteed to be non-None. 

49 

50 The settings are lazy-loaded in litellm.__init__ and always return an instance. 

51 This helper provides proper type narrowing for mypy. 

52 """ 

53 settings: Final = litellm.priority_reservation_settings 

54 if settings is None: 

55 # This should never happen due to lazy loading, but satisfy mypy 

56 from litellm.types.utils import PriorityReservationSettings 

57 

58 return PriorityReservationSettings() 

59 return settings 

60 

61 

62def _is_latin1_encodable(value: object) -> bool: 

63 return all(ord(char) < 256 for char in str(value)) 

64 

65 

66class _PROXY_DynamicRateLimitHandlerV3(CustomLogger): 

67 """ 

68 Saturation-aware priority-based rate limiter using v3 infrastructure. 

69 

70 Key features: 

71 1. Model capacity ALWAYS enforced at 100% (prevents over-allocation) 

72 2. Priority usage tracked from first request (accurate accounting) 

73 3. Priority limits only enforced when saturated >= threshold 

74 4. Three-phase checking prevents partial counter increments 

75 5. Reuses v3 limiter's Redis-based tracking (multi-instance safe) 

76 

77 How it works: 

78 - Phase 1: Read-only check of ALL limits (no increments) 

79 - Phase 2: Decide enforcement based on saturation 

80 - Phase 3: Increment counters only if request allowed 

81 - When under-saturated: priorities can borrow unused capacity (generous) 

82 - When saturated: strict priority-based limits enforced (fair) 

83 - Uses v3 limiter's atomic Lua scripts for race-free increments 

84 """ 

85 

86 def __init__( 

87 self, 

88 internal_usage_cache: DualCache, 

89 time_provider: Callable[[], datetime] | None = None, 

90 ): 

91 self.internal_usage_cache = InternalUsageCache(dual_cache=internal_usage_cache) 

92 self.v3_limiter = _PROXY_MaxParallelRequestsHandler_v3(self.internal_usage_cache, time_provider=time_provider) 

93 

94 def update_variables(self, llm_router: Router): 

95 self.llm_router = llm_router 

96 

97 def _get_saturation_check_cache_ttl(self) -> int: 

98 """Get the configurable TTL for local cache when reading saturation values.""" 

99 return _get_priority_settings().saturation_check_cache_ttl 

100 

101 async def _get_saturation_value_from_cache( 

102 self, 

103 counter_key: str, 

104 ) -> str | None: 

105 """ 

106 Get saturation value with configurable local cache TTL. 

107 

108 Uses DualCache with configurable TTL for local cache storage. 

109 TTL is configurable via litellm.priority_reservation_settings.saturation_check_cache_ttl 

110 

111 Args: 

112 counter_key: The cache key for the saturation counter 

113 

114 Returns: 

115 Counter value as string, or None if not found 

116 """ 

117 local_cache_ttl: Final = self._get_saturation_check_cache_ttl() 

118 

119 return await self.internal_usage_cache.async_get_cache( 

120 key=counter_key, 

121 litellm_parent_otel_span=None, 

122 local_only=False, 

123 ttl=local_cache_ttl, 

124 ) 

125 

126 def _get_priority_weight(self, priority: str | None, model_info: ModelGroupInfo | None = None) -> float: 

127 """Get the weight for a given priority from litellm.priority_reservation""" 

128 weight: float = _get_priority_settings().default_priority 

129 if litellm.priority_reservation is None or priority not in litellm.priority_reservation: 

130 verbose_proxy_logger.debug("Priority Reservation not set for the given priority.") 

131 elif priority is not None and litellm.priority_reservation is not None: 

132 if os.getenv("LITELLM_LICENSE", None) is None: 

133 verbose_proxy_logger.error( 

134 "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." 

135 ) 

136 else: 

137 value: Final = litellm.priority_reservation[priority] 

138 weight = convert_priority_to_percent(value, model_info) 

139 return weight 

140 

141 def _get_priority_from_user_api_key_dict(self, user_api_key_dict: UserAPIKeyAuth) -> str | None: 

142 """ 

143 Get priority from user_api_key_dict. 

144 

145 Checks team metadata first (takes precedence), then falls back to key metadata. 

146 

147 Args: 

148 user_api_key_dict: User authentication info 

149 

150 Returns: 

151 Priority string if found, None otherwise 

152 """ 

153 priority: str | None = None 

154 

155 # Check team metadata first (takes precedence) 

156 if user_api_key_dict.team_metadata is not None: 

157 priority = user_api_key_dict.team_metadata.get("priority", None) 

158 

159 # Fall back to key metadata 

160 if priority is None: 

161 priority = user_api_key_dict.metadata.get("priority", None) 

162 

163 return priority 

164 

165 def _normalize_priority_weights(self, model_info: ModelGroupInfo) -> dict[str, float]: 

166 """ 

167 Normalize priority weights if they sum to > 1.0 

168 

169 Handles over-allocation: {key_a: 0.60, key_b: 0.80} -> {key_a: 0.43, key_b: 0.57} 

170 Converts absolute rpm/tpm values to percentages based on model capacity. 

171 """ 

172 if litellm.priority_reservation is None: 

173 return {} 

174 

175 # Convert all values to percentages first 

176 weights: Final[dict[str, float]] = {} 

177 for k, v in litellm.priority_reservation.items(): 

178 weights[k] = convert_priority_to_percent(v, model_info) 

179 

180 total_weight: Final = sum(weights.values()) 

181 

182 if total_weight > 1.0: 

183 normalized: Final = {k: v / total_weight for k, v in weights.items()} 

184 verbose_proxy_logger.debug("Normalized over-allocated priorities: %s -> %s", weights, normalized) 

185 return normalized 

186 

187 return weights 

188 

189 def _get_priority_allocation( 

190 self, 

191 model: str, 

192 priority: str | None, 

193 normalized_weights: dict[str, float], 

194 model_info: ModelGroupInfo | None = None, 

195 ) -> tuple[float, str]: 

196 """ 

197 Get priority weight and pool key for a given priority. 

198 

199 For explicit priorities: returns specific allocation and unique pool key 

200 For default priority: returns default allocation and shared pool key 

201 

202 Args: 

203 model: Model name 

204 priority: Priority level (None for default) 

205 normalized_weights: Pre-computed normalized weights 

206 model_info: Model configuration (optional, for fallback conversion) 

207 

208 Returns: 

209 tuple: (priority_weight, priority_key) 

210 """ 

211 # Check if this key has an explicit priority in litellm.priority_reservation 

212 has_explicit_priority: Final = ( 

213 priority is not None 

214 and litellm.priority_reservation is not None 

215 and priority in litellm.priority_reservation 

216 ) 

217 

218 if has_explicit_priority and priority is not None: 

219 # Explicit priority: get its specific allocation 

220 priority_weight = normalized_weights.get(priority, self._get_priority_weight(priority, model_info)) 

221 # Use unique key per priority level 

222 priority_key = f"{model}:{priority}" 

223 else: 

224 # No explicit priority: share the default_priority pool with ALL other default keys 

225 priority_weight = _get_priority_settings().default_priority 

226 # Use shared key for all default-priority requests 

227 priority_key = f"{model}:default_pool" 

228 

229 return priority_weight, priority_key 

230 

231 async def _check_model_saturation( 

232 self, 

233 model: str, 

234 model_group_info: ModelGroupInfo, 

235 ) -> float: 

236 """ 

237 Check current saturation by directly querying v3 limiter's cache keys. 

238 

239 Reuses v3 limiter's Redis-based tracking (works across multiple instances). 

240 Reads counters WITHOUT incrementing them. 

241 

242 Returns: 

243 float: Saturation ratio (0.0 = empty, 1.0 = at capacity, >1.0 = over) 

244 """ 

245 try: 

246 max_saturation = 0.0 

247 

248 # Query RPM saturation - always read from Redis for multi-node consistency 

249 if model_group_info.rpm is not None and model_group_info.rpm > 0: 

250 # Use v3 limiter's key format: {key:value}:rate_limit_type 

251 counter_key = self.v3_limiter.create_rate_limit_keys( 

252 key="model_saturation_check", 

253 value=model, 

254 rate_limit_type="requests", 

255 ) 

256 

257 # Query Redis directly for current counter value (skip local cache for consistency) 

258 counter_value = await self._get_saturation_value_from_cache(counter_key=counter_key) 

259 

260 if counter_value is not None: 

261 current_requests: Final = int(counter_value) 

262 rpm_saturation: Final = current_requests / model_group_info.rpm 

263 max_saturation = max(max_saturation, rpm_saturation) 

264 

265 verbose_proxy_logger.debug( 

266 f"Model {model} RPM: {current_requests}/{model_group_info.rpm} ({rpm_saturation:.1%})" 

267 ) 

268 

269 # Query TPM saturation 

270 if model_group_info.tpm is not None and model_group_info.tpm > 0: 

271 counter_key = self.v3_limiter.create_rate_limit_keys( 

272 key="model_saturation_check", 

273 value=model, 

274 rate_limit_type="tokens", 

275 ) 

276 

277 counter_value = await self._get_saturation_value_from_cache(counter_key=counter_key) 

278 

279 if counter_value is not None: 

280 current_tokens: Final = float(counter_value) 

281 tpm_saturation: Final = current_tokens / model_group_info.tpm 

282 max_saturation = max(max_saturation, tpm_saturation) 

283 

284 verbose_proxy_logger.debug( 

285 f"Model {model} TPM: {current_tokens}/{model_group_info.tpm} ({tpm_saturation:.1%})" 

286 ) 

287 

288 verbose_proxy_logger.debug(f"Model {model} overall saturation: {max_saturation:.1%}") 

289 

290 return max_saturation 

291 

292 except Exception as e: 

293 verbose_proxy_logger.error("Error checking saturation for %s: %s", model, e) 

294 # Fail open: assume not saturated on error 

295 return 0.0 

296 

297 def _create_priority_based_descriptors( 

298 self, 

299 model: str, 

300 user_api_key_dict: UserAPIKeyAuth, 

301 priority: str | None, 

302 ) -> list[RateLimitDescriptor]: 

303 """ 

304 Create rate limit descriptors with normalized priority weights. 

305 

306 Uses normalized weights to handle over-allocation scenarios. 

307 

308 For explicit priorities: each priority gets its own pool (e.g., prod gets 75%) 

309 For default priority: ALL keys without explicit priority share ONE pool (e.g., all share 25%) 

310 """ 

311 descriptors: Final[list[RateLimitDescriptor]] = [] 

312 

313 if litellm.priority_reservation is None: 

314 return descriptors 

315 

316 # Get model group info 

317 model_group_info: Final[ModelGroupInfo | None] = self.llm_router.get_model_group_info(model_group=model) 

318 if model_group_info is None: 

319 return descriptors 

320 

321 # Get normalized priority weight and pool key 

322 normalized_weights: Final = self._normalize_priority_weights(model_group_info) 

323 priority_weight, priority_key = self._get_priority_allocation( 

324 model=model, 

325 priority=priority, 

326 normalized_weights=normalized_weights, 

327 model_info=model_group_info, 

328 ) 

329 

330 rate_limit_config: Final[RateLimitDescriptorRateLimitObject] = {} 

331 

332 # Apply priority weight to model limits 

333 if model_group_info.tpm is not None: 

334 reserved_tpm: Final = int(model_group_info.tpm * priority_weight) 

335 rate_limit_config["tokens_per_unit"] = reserved_tpm 

336 

337 if model_group_info.rpm is not None: 

338 reserved_rpm: Final = int(model_group_info.rpm * priority_weight) 

339 rate_limit_config["requests_per_unit"] = reserved_rpm 

340 

341 if rate_limit_config: 

342 rate_limit_config["window_size"] = self.v3_limiter.window_size 

343 

344 descriptors.append( 

345 RateLimitDescriptor( 

346 key="priority_model", 

347 value=priority_key, 

348 rate_limit=rate_limit_config, 

349 ) 

350 ) 

351 

352 return descriptors 

353 

354 def _create_model_tracking_descriptor( 

355 self, 

356 model: str, 

357 model_group_info: ModelGroupInfo, 

358 high_limit_multiplier: int = 1, 

359 ) -> RateLimitDescriptor: 

360 """ 

361 Create a descriptor for tracking model-wide usage. 

362 

363 Args: 

364 model: Model name 

365 model_group_info: Model configuration with RPM/TPM limits 

366 high_limit_multiplier: Multiplier for limits (use >1 for tracking-only) 

367 

368 Returns: 

369 Rate limit descriptor for model-wide tracking 

370 """ 

371 return RateLimitDescriptor( 

372 key="model_saturation_check", 

373 value=model, 

374 rate_limit={ 

375 "requests_per_unit": (model_group_info.rpm * high_limit_multiplier if model_group_info.rpm else None), 

376 "tokens_per_unit": (model_group_info.tpm * high_limit_multiplier if model_group_info.tpm else None), 

377 "window_size": self.v3_limiter.window_size, 

378 }, 

379 ) 

380 

381 async def _check_rate_limits( 

382 self, 

383 model: str, 

384 model_group_info: ModelGroupInfo, 

385 user_api_key_dict: UserAPIKeyAuth, 

386 priority: str | None, 

387 saturation: float, 

388 ) -> None: 

389 """ 

390 Check rate limits using THREE-PHASE approach to prevent partial increments. 

391 

392 Phase 1: Read-only check of ALL limits (no increments) 

393 Phase 2: Decide which limits to enforce based on saturation 

394 Phase 3: Increment ALL counters atomically (model + priority) 

395 

396 This prevents the bug where: 

397 - Model counter increments in stage 1 

398 - Priority check fails in stage 2 

399 - Request blocked but model counter already incremented 

400 

401 Key behaviors: 

402 - All checks performed first (read-only) 

403 - Only increment counters if request will be allowed 

404 - Model capacity: Always enforced at 100% 

405 - Priority limits: Only enforced when saturated >= threshold 

406 - Both counters tracked from first request (accurate accounting) 

407 

408 Args: 

409 model: Model name 

410 model_group_info: Model configuration 

411 user_api_key_dict: User authentication info 

412 priority: User's priority level 

413 saturation: Current saturation level 

414 

415 Raises: 

416 HTTPException: If any limit is exceeded 

417 """ 

418 import json 

419 

420 saturation_threshold: Final = _get_priority_settings().saturation_threshold 

421 should_enforce_priority: Final = saturation >= saturation_threshold 

422 

423 # Build ALL descriptors upfront 

424 descriptors_to_check: Final[list[RateLimitDescriptor]] = [] 

425 

426 # Model-wide descriptor (always enforce) 

427 model_wide_descriptor: Final = self._create_model_tracking_descriptor( 

428 model=model, 

429 model_group_info=model_group_info, 

430 high_limit_multiplier=1, 

431 ) 

432 descriptors_to_check.append(model_wide_descriptor) 

433 

434 # Priority descriptors (always track, conditionally enforce) 

435 priority_descriptors: Final = self._create_priority_based_descriptors( 

436 model=model, 

437 user_api_key_dict=user_api_key_dict, 

438 priority=priority, 

439 ) 

440 if priority_descriptors: 

441 descriptors_to_check.extend(priority_descriptors) 

442 

443 # Atomic check-and-increment for the ENFORCED descriptor set: 

444 # - model_saturation_check is always enforced 

445 # - priority_model is enforced only when saturation crosses threshold 

446 # 

447 # Backed by a Redis Lua script (multi-process atomic) with an 

448 # asyncio.Lock + in-memory fallback for single-process deployments. 

449 # All-or-nothing: if any enforced descriptor would exceed its limit, 

450 # no counter is modified and the response carries "OVER_LIMIT". 

451 enforced_descriptors: Final[list[RateLimitDescriptor]] = [model_wide_descriptor] 

452 if priority_descriptors and should_enforce_priority: 

453 enforced_descriptors.extend(priority_descriptors) 

454 

455 per_request_increment: Final[dict[Literal["requests", "tokens"], int]] = { 

456 "requests": 1, 

457 "tokens": 0, 

458 } 

459 atomic_response: Final = await self.v3_limiter.atomic_check_and_increment_by_n( 

460 descriptors=enforced_descriptors, 

461 increments=[per_request_increment for _ in enforced_descriptors], 

462 parent_otel_span=user_api_key_dict.parent_otel_span, 

463 ) 

464 

465 verbose_proxy_logger.debug( 

466 "Atomic check+increment response: %s", json.dumps(atomic_response, indent=2, default=list) 

467 ) 

468 

469 if atomic_response["overall_code"] == "OVER_LIMIT": 

470 resolved_model, llm_provider = resolve_llm_provider_for_rate_limit(model) 

471 for status in atomic_response["statuses"]: 

472 if status["code"] != "OVER_LIMIT": 

473 continue 

474 descriptor_key = status["descriptor_key"] 

475 if descriptor_key == "model_saturation_check": 

476 raise ProxyRateLimitError( 

477 detail={ 

478 "error": f"Model capacity reached for {model}. " 

479 f"Priority: {priority}, " 

480 f"Rate limit type: {status['rate_limit_type']}, " 

481 f"Model TPM: {model_group_info.tpm if model_group_info.tpm is not None else 'not configured'}, " 

482 f"Model RPM: {model_group_info.rpm if model_group_info.rpm is not None else 'not configured'}, " 

483 f"Remaining: {status['limit_remaining']}" 

484 }, 

485 headers={ 

486 "retry-after": str(self.v3_limiter.window_size), 

487 "rate_limit_type": str(status["rate_limit_type"]), 

488 "x-litellm-priority": priority or "default", 

489 }, 

490 rate_limit_type=map_v3_rate_limit_type(status["rate_limit_type"]), 

491 model=resolved_model, 

492 llm_provider=llm_provider, 

493 ) 

494 if descriptor_key == "priority_model": 

495 verbose_proxy_logger.debug( 

496 f"Enforcing priority limits for {model}, saturation: {saturation:.1%}, priority: {priority}" 

497 ) 

498 raise ProxyRateLimitError( 

499 detail={ 

500 "error": f"Priority-based rate limit exceeded. " 

501 f"Model: {model}, " 

502 f"Priority: {priority}, " 

503 f"Rate limit type: {status['rate_limit_type']}, " 

504 f"Model TPM: {model_group_info.tpm if model_group_info.tpm is not None else 'not configured'}, " 

505 f"Model RPM: {model_group_info.rpm if model_group_info.rpm is not None else 'not configured'}, " 

506 f"Remaining: {status['limit_remaining']}, " 

507 f"Model saturation: {saturation:.1%}" 

508 }, 

509 headers={ 

510 "retry-after": str(self.v3_limiter.window_size), 

511 "rate_limit_type": str(status["rate_limit_type"]), 

512 "x-litellm-priority": priority or "default", 

513 "x-litellm-saturation": f"{saturation:.2%}", 

514 }, 

515 rate_limit_type=map_v3_rate_limit_type(status["rate_limit_type"]), 

516 model=resolved_model, 

517 llm_provider=llm_provider, 

518 ) 

519 

520 # Fail-closed guard: overall_code says OVER_LIMIT but no status 

521 # matched a descriptor key we know how to translate into a 429. 

522 # Refuse the request rather than silently fall through and let an 

523 # over-limit request proceed to the model. Without this, a future 

524 # caller wiring an unfamiliar descriptor into enforced_descriptors 

525 # would silently bypass the rate limit. 

526 offending: Final = next( 

527 (s for s in atomic_response["statuses"] if s["code"] == "OVER_LIMIT"), 

528 None, 

529 ) 

530 verbose_proxy_logger.error( 

531 "Dynamic rate limiter: OVER_LIMIT response with unknown descriptor_key(s) — refusing request. response=%s", 

532 atomic_response, 

533 ) 

534 raise ProxyRateLimitError( 

535 detail={ 

536 "error": "Rate limit exceeded", 

537 "descriptor_key": (offending["descriptor_key"] if offending else "unknown"), 

538 "rate_limit_type": (str(offending["rate_limit_type"]) if offending else "unknown"), 

539 }, 

540 rate_limit_type=map_v3_rate_limit_type(offending["rate_limit_type"] if offending else None), 

541 headers={ 

542 "retry-after": str(self.v3_limiter.window_size), 

543 "x-litellm-priority": priority or "default", 

544 }, 

545 model=resolved_model, 

546 llm_provider=llm_provider, 

547 ) 

548 

549 # If priority is NOT enforced (saturation below threshold) but 

550 # priority_descriptors exist, increment them for tracking only — no 

551 # check, no rollback. This matches the prior tracking semantics. 

552 # 

553 # Using the non-atomic should_rate_limit (instead of 

554 # atomic_check_and_increment_by_n) is intentional here: we don't want 

555 # to enforce the limit, we only want to bump the counter so the 

556 # priority allocation has accurate usage when it later becomes 

557 # enforced. The increment-then-check semantics of should_rate_limit 

558 # are fine because we ignore the OVER_LIMIT response. 

559 if priority_descriptors and not should_enforce_priority: 

560 priority_tracking_response: Final = await self.v3_limiter.should_rate_limit( 

561 descriptors=priority_descriptors, 

562 parent_otel_span=user_api_key_dict.parent_otel_span, 

563 read_only=False, 

564 ) 

565 get_or_create_request_stash().rate_limit_response = RateLimitResponse( 

566 overall_code=atomic_response["overall_code"], 

567 statuses=atomic_response["statuses"] + priority_tracking_response["statuses"], 

568 ) 

569 else: 

570 get_or_create_request_stash().rate_limit_response = atomic_response 

571 

572 async def async_pre_call_hook( 

573 self, 

574 user_api_key_dict: UserAPIKeyAuth, 

575 cache: DualCache, 

576 data: dict, 

577 call_type: CallTypesLiteral, 

578 ) -> Exception | str | dict | None: 

579 """ 

580 Saturation-aware pre-call hook for priority-based rate limiting. 

581 

582 Flow: 

583 1. Check current saturation level 

584 2. THREE-PHASE rate limit check: 

585 - PHASE 1: Read-only check of ALL limits (no increments) 

586 - PHASE 2: Decide which limits to enforce based on saturation 

587 - PHASE 3: Increment ALL counters atomically if request allowed 

588 

589 This three-phase approach ensures: 

590 - Model capacity is NEVER exceeded (always enforced at 100%) 

591 - Priority usage tracked from first request (accurate metrics) 

592 - Counters only increment when request will be allowed (prevents phantom usage) 

593 - When under-saturated: priorities can borrow unused capacity (generous) 

594 - When saturated: fair allocation based on normalized priority weights (strict) 

595 

596 Example with 100 RPM model, 60% priority allocation, 80% threshold: 

597 - Saturation < 80%: Priority can use up to 100 RPM (model limit enforced only) 

598 - Saturation >= 80%: Priority limited to 60 RPM (both limits enforced) 

599 

600 Prevents bugs where: 

601 - Model counter increments but priority check fails → model over-capacity 

602 - Priority counter increments but not enforced → inaccurate metrics 

603 

604 Args: 

605 user_api_key_dict: User authentication and metadata 

606 cache: Dual cache instance 

607 data: Request data containing model name 

608 call_type: Type of API call being made 

609 

610 Returns: 

611 None if request is allowed, otherwise raises HTTPException 

612 """ 

613 if "model" not in data: 

614 return None 

615 

616 claim_request_stash_for_data(data) 

617 model: Final = data["model"] 

618 priority: Final = self._get_priority_from_user_api_key_dict(user_api_key_dict=user_api_key_dict) 

619 

620 # Get model configuration 

621 model_group_info: Final[ModelGroupInfo | None] = self.llm_router.get_model_group_info(model_group=model) 

622 if model_group_info is None: 

623 verbose_proxy_logger.debug("No model group info for %s, allowing request", model) 

624 return None 

625 

626 try: 

627 # STEP 1: Check current saturation level 

628 saturation: Final = await self._check_model_saturation(model, model_group_info) 

629 

630 saturation_threshold: Final = _get_priority_settings().saturation_threshold 

631 

632 verbose_proxy_logger.debug( 

633 f"[Dynamic Rate Limiter] Model={model}, Saturation={saturation:.1%}, " 

634 f"Threshold={saturation_threshold:.1%}, Priority={priority}" 

635 ) 

636 

637 # STEP 2: Check rate limits in THREE phases 

638 # Phase 1: Read-only check of ALL limits (no increments) 

639 # Phase 2: Decide which limits to enforce (based on saturation) 

640 # Phase 3: Increment ALL counters only if request will be allowed 

641 # This prevents partial increments and ensures accurate tracking 

642 await self._check_rate_limits( 

643 model=model, 

644 model_group_info=model_group_info, 

645 user_api_key_dict=user_api_key_dict, 

646 priority=priority, 

647 saturation=saturation, 

648 ) 

649 

650 except HTTPException: 

651 raise 

652 except Exception as e: 

653 verbose_proxy_logger.error("Error in dynamic rate limiter: %s, allowing request", e) 

654 # Fail open on unexpected errors 

655 return None 

656 

657 return None 

658 

659 async def async_post_call_success_hook(self, data: dict, user_api_key_dict: UserAPIKeyAuth, response): 

660 """ 

661 Post-call hook to add rate limit headers to response. 

662 Leverages v3 limiter's post-call hook functionality. 

663 """ 

664 try: 

665 # Call v3 limiter's post-call hook to add standard rate limit headers 

666 await self.v3_limiter.async_post_call_success_hook( 

667 data=data, user_api_key_dict=user_api_key_dict, response=response 

668 ) 

669 

670 if response_has_hidden_params(response): 

671 priority: Final = self._get_priority_from_user_api_key_dict(user_api_key_dict=user_api_key_dict) 

672 additional_headers: Final = ensure_response_additional_headers(response) 

673 priority_header: Final = priority or "default" 

674 if _is_latin1_encodable(priority_header): 

675 additional_headers["x-litellm-priority"] = priority_header 

676 else: 

677 verbose_proxy_logger.debug( 

678 "Skipping x-litellm-priority header: priority %r is not Latin-1 encodable", priority 

679 ) 

680 additional_headers["x-litellm-rate-limiter-version"] = "v3" 

681 

682 return response 

683 

684 except Exception as e: 

685 verbose_proxy_logger.exception("Error in dynamic rate limiter v3 post-call hook: %s", e) 

686 return response 

687 

688 async def async_log_success_event(self, kwargs, response_obj, start_time, end_time): 

689 """ 

690 Update token usage for priority-based rate limiting after successful API calls. 

691 

692 Increments token counters for: 

693 - model_saturation_check: Model-wide token tracking 

694 - priority_model: Priority-specific token tracking 

695 """ 

696 from litellm.litellm_core_utils.core_helpers import ( 

697 _get_parent_otel_span_from_kwargs, 

698 ) 

699 from litellm.proxy.common_utils.callback_utils import ( 

700 get_model_group_from_litellm_kwargs, 

701 ) 

702 from litellm.types.caching import RedisPipelineIncrementOperation 

703 from litellm.types.utils import Usage 

704 

705 try: 

706 verbose_proxy_logger.debug("INSIDE dynamic rate limiter ASYNC SUCCESS LOGGING") 

707 

708 litellm_parent_otel_span: Final = _get_parent_otel_span_from_kwargs(kwargs) 

709 

710 # Get metadata from standard_logging_object 

711 standard_logging_object: Final = kwargs.get("standard_logging_object") or {} 

712 standard_logging_metadata: Final = standard_logging_object.get("metadata") or {} 

713 

714 # Get model and priority 

715 model_group: Final = get_model_group_from_litellm_kwargs(kwargs) 

716 if not model_group: 

717 return 

718 

719 # Get priority from user_api_key_auth_metadata in standard_logging_metadata 

720 # This is where user_api_key_dict.metadata is stored during pre-call 

721 user_api_key_auth_metadata: Final = standard_logging_metadata.get("user_api_key_auth_metadata") or {} 

722 key_priority: Final[str | None] = user_api_key_auth_metadata.get("priority") 

723 

724 # Get total tokens from response 

725 total_tokens = 0 

726 rate_limit_type: Final = self.v3_limiter.get_rate_limit_type() 

727 

728 if isinstance(response_obj, ModelResponse): 

729 _usage: Final = getattr(response_obj, "usage", None) 

730 if _usage and isinstance(_usage, Usage): 

731 if rate_limit_type == "output": 

732 total_tokens = _usage.completion_tokens 

733 elif rate_limit_type == "input": 

734 total_tokens = _usage.prompt_tokens 

735 elif rate_limit_type == "total": 

736 total_tokens = _usage.total_tokens 

737 

738 if total_tokens == 0: 

739 return 

740 

741 # Create pipeline operations for token increments 

742 pipeline_operations: Final[list[RedisPipelineIncrementOperation]] = [] 

743 

744 # Model-wide token tracking (model_saturation_check) 

745 model_token_key: Final = self.v3_limiter.create_rate_limit_keys( 

746 key="model_saturation_check", 

747 value=model_group, 

748 rate_limit_type="tokens", 

749 ) 

750 pipeline_operations.append( 

751 RedisPipelineIncrementOperation( 

752 key=model_token_key, 

753 increment_value=total_tokens, 

754 ttl=self.v3_limiter.window_size, 

755 ) 

756 ) 

757 

758 # Priority-specific token tracking (priority_model) 

759 # Determine priority key (same logic as _get_priority_allocation) 

760 has_explicit_priority: Final = ( 

761 key_priority is not None 

762 and litellm.priority_reservation is not None 

763 and key_priority in litellm.priority_reservation 

764 ) 

765 

766 if has_explicit_priority and key_priority is not None: 

767 priority_key = f"{model_group}:{key_priority}" 

768 else: 

769 priority_key = f"{model_group}:default_pool" 

770 

771 priority_token_key: Final = self.v3_limiter.create_rate_limit_keys( 

772 key="priority_model", 

773 value=priority_key, 

774 rate_limit_type="tokens", 

775 ) 

776 pipeline_operations.append( 

777 RedisPipelineIncrementOperation( 

778 key=priority_token_key, 

779 increment_value=total_tokens, 

780 ttl=self.v3_limiter.window_size, 

781 ) 

782 ) 

783 

784 # Execute token increments with TTL preservation 

785 if pipeline_operations: 

786 await self.v3_limiter.async_increment_tokens_with_ttl_preservation( 

787 pipeline_operations=pipeline_operations, 

788 parent_otel_span=litellm_parent_otel_span, 

789 ) 

790 

791 # Only log 'priority' if it's known safe; otherwise, redact. 

792 SAFE_PRIORITIES: Final = {"low", "medium", "high", "default"} 

793 logged_priority: Final = key_priority if key_priority in SAFE_PRIORITIES else "REDACTED" 

794 verbose_proxy_logger.debug( 

795 "[Dynamic Rate Limiter] Incremented tokens by %s for model=%s, priority=%s", 

796 total_tokens, 

797 model_group, 

798 logged_priority, 

799 ) 

800 

801 except Exception as e: 

802 verbose_proxy_logger.exception("Error in dynamic rate limiter success event: %s", e)