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
« 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"""
5import os
6from collections.abc import Callable
7from datetime import datetime
8from typing import TYPE_CHECKING, Final, Literal
10from fastapi import HTTPException
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
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
46def _get_priority_settings() -> "PriorityReservationSettings":
47 """
48 Get the priority reservation settings, guaranteed to be non-None.
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
58 return PriorityReservationSettings()
59 return settings
62def _is_latin1_encodable(value: object) -> bool:
63 return all(ord(char) < 256 for char in str(value))
66class _PROXY_DynamicRateLimitHandlerV3(CustomLogger):
67 """
68 Saturation-aware priority-based rate limiter using v3 infrastructure.
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)
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 """
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)
94 def update_variables(self, llm_router: Router):
95 self.llm_router = llm_router
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
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.
108 Uses DualCache with configurable TTL for local cache storage.
109 TTL is configurable via litellm.priority_reservation_settings.saturation_check_cache_ttl
111 Args:
112 counter_key: The cache key for the saturation counter
114 Returns:
115 Counter value as string, or None if not found
116 """
117 local_cache_ttl: Final = self._get_saturation_check_cache_ttl()
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 )
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
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.
145 Checks team metadata first (takes precedence), then falls back to key metadata.
147 Args:
148 user_api_key_dict: User authentication info
150 Returns:
151 Priority string if found, None otherwise
152 """
153 priority: str | None = None
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)
159 # Fall back to key metadata
160 if priority is None:
161 priority = user_api_key_dict.metadata.get("priority", None)
163 return priority
165 def _normalize_priority_weights(self, model_info: ModelGroupInfo) -> dict[str, float]:
166 """
167 Normalize priority weights if they sum to > 1.0
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 {}
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)
180 total_weight: Final = sum(weights.values())
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
187 return weights
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.
199 For explicit priorities: returns specific allocation and unique pool key
200 For default priority: returns default allocation and shared pool key
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)
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 )
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"
229 return priority_weight, priority_key
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.
239 Reuses v3 limiter's Redis-based tracking (works across multiple instances).
240 Reads counters WITHOUT incrementing them.
242 Returns:
243 float: Saturation ratio (0.0 = empty, 1.0 = at capacity, >1.0 = over)
244 """
245 try:
246 max_saturation = 0.0
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 )
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)
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)
265 verbose_proxy_logger.debug(
266 f"Model {model} RPM: {current_requests}/{model_group_info.rpm} ({rpm_saturation:.1%})"
267 )
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 )
277 counter_value = await self._get_saturation_value_from_cache(counter_key=counter_key)
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)
284 verbose_proxy_logger.debug(
285 f"Model {model} TPM: {current_tokens}/{model_group_info.tpm} ({tpm_saturation:.1%})"
286 )
288 verbose_proxy_logger.debug(f"Model {model} overall saturation: {max_saturation:.1%}")
290 return max_saturation
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
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.
306 Uses normalized weights to handle over-allocation scenarios.
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]] = []
313 if litellm.priority_reservation is None:
314 return descriptors
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
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 )
330 rate_limit_config: Final[RateLimitDescriptorRateLimitObject] = {}
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
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
341 if rate_limit_config:
342 rate_limit_config["window_size"] = self.v3_limiter.window_size
344 descriptors.append(
345 RateLimitDescriptor(
346 key="priority_model",
347 value=priority_key,
348 rate_limit=rate_limit_config,
349 )
350 )
352 return descriptors
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.
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)
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 )
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.
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)
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
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)
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
415 Raises:
416 HTTPException: If any limit is exceeded
417 """
418 import json
420 saturation_threshold: Final = _get_priority_settings().saturation_threshold
421 should_enforce_priority: Final = saturation >= saturation_threshold
423 # Build ALL descriptors upfront
424 descriptors_to_check: Final[list[RateLimitDescriptor]] = []
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)
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)
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)
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 )
465 verbose_proxy_logger.debug(
466 "Atomic check+increment response: %s", json.dumps(atomic_response, indent=2, default=list)
467 )
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 )
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 )
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
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.
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
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)
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)
600 Prevents bugs where:
601 - Model counter increments but priority check fails → model over-capacity
602 - Priority counter increments but not enforced → inaccurate metrics
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
610 Returns:
611 None if request is allowed, otherwise raises HTTPException
612 """
613 if "model" not in data:
614 return None
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)
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
626 try:
627 # STEP 1: Check current saturation level
628 saturation: Final = await self._check_model_saturation(model, model_group_info)
630 saturation_threshold: Final = _get_priority_settings().saturation_threshold
632 verbose_proxy_logger.debug(
633 f"[Dynamic Rate Limiter] Model={model}, Saturation={saturation:.1%}, "
634 f"Threshold={saturation_threshold:.1%}, Priority={priority}"
635 )
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 )
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
657 return None
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 )
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"
682 return response
684 except Exception as e:
685 verbose_proxy_logger.exception("Error in dynamic rate limiter v3 post-call hook: %s", e)
686 return response
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.
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
705 try:
706 verbose_proxy_logger.debug("INSIDE dynamic rate limiter ASYNC SUCCESS LOGGING")
708 litellm_parent_otel_span: Final = _get_parent_otel_span_from_kwargs(kwargs)
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 {}
714 # Get model and priority
715 model_group: Final = get_model_group_from_litellm_kwargs(kwargs)
716 if not model_group:
717 return
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")
724 # Get total tokens from response
725 total_tokens = 0
726 rate_limit_type: Final = self.v3_limiter.get_rate_limit_type()
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
738 if total_tokens == 0:
739 return
741 # Create pipeline operations for token increments
742 pipeline_operations: Final[list[RedisPipelineIncrementOperation]] = []
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 )
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 )
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"
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 )
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 )
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 )
801 except Exception as e:
802 verbose_proxy_logger.exception("Error in dynamic rate limiter success event: %s", e)