From 9e07e685f3fbd13684e3092419e1616bd2f342a6 Mon Sep 17 00:00:00 2001 From: Deepanshu Date: Mon, 17 Aug 2026 11:02:39 -0400 Subject: [PATCH] feat(rate-limiting): make the in-memory cache size configurable per tag entry Adds TagRateLimitEntry.max_in_memory_cache_size so a single high-cardinality entry can get its own dedicated in-memory cache partition instead of sharing the hook's single default one, keeping the knob alongside key_ttl_seconds and the rest of that entry's config rather than only as a proxy-wide setting. Partitions are keyed by each entry's full signature, not the override value alone, so two unrelated entries that happen to pick the same size don't get merged. tokens/dollars accounting and concurrency-slot release are both partition-aware too: each partition owns its own v3 handler, and a concurrency reservation is released against the exact partition it was incremented on so it can't leak onto the default partition instead. Also fixes a bug this surfaced: _configured_limit_for_signature rebuilt each entry from a 5-field dedup signature that excluded key_ttl_seconds and max_in_memory_cache_size, silently resetting both to their defaults for every entry reachable through the real indexing path used by async_filter_deployments and async_log_success_event. --- litellm/proxy/hooks/tag_rate_limiter.py | 257 ++++++++++----- litellm/types/router.py | 15 + .../proxy/hooks/test_tag_rate_limiter.py | 312 +++++++++++++++++- ui/litellm-dashboard/src/lib/http/schema.d.ts | 2 + 4 files changed, 499 insertions(+), 87 deletions(-) diff --git a/litellm/proxy/hooks/tag_rate_limiter.py b/litellm/proxy/hooks/tag_rate_limiter.py index 15ed2ae4f1e..8b5425a137e 100644 --- a/litellm/proxy/hooks/tag_rate_limiter.py +++ b/litellm/proxy/hooks/tag_rate_limiter.py @@ -203,29 +203,22 @@ def _entries_for_unit(deployment: Mapping[str, object], unit: _LimitUnit) -> tup def _configured_limit_for_signature( unit: _LimitUnit, - signature: tuple[str, str, float, int, bool], + entry: TagRateLimitEntry, declaring_ids: Sequence[str], is_chain_wide: bool, ) -> _ConfiguredLimit | None: - tag_id, name, limit, period_seconds, scope_by_key_hash = signature if unit == "concurrency" and not is_chain_wide: verbose_proxy_logger.warning( "tag_rate_limiter: concurrency_limits entry %r (tag_id=%s) is not declared identically by every " "deployment sharing this model_name; per-deployment-scoped concurrency limits are not supported " "and this entry is being skipped entirely.", - name, - tag_id, + entry.name, + entry.tag_id, ) return None return _ConfiguredLimit( unit=unit, - entry=TagRateLimitEntry( - name=name, - tag_id=tag_id, - limit=limit, - period_seconds=period_seconds, - scope_by_key_hash=scope_by_key_hash, - ), + entry=entry, deployment_scope=None if is_chain_wide else tuple(sorted(declaring_ids)), ) @@ -264,6 +257,14 @@ def _build_group_limits(deployments: Sequence[Mapping[str, object]], unit: _Limi # would scramble that first-seen order, so this stays a plain # accumulator instead. declaring_ids_by_signature: Final = {} # mutable-ok: first-seen order here decides which limit's error raises first (see comment above); sorting to use groupby would scramble it + # The dedup signature is deliberately narrower than the full entry: two + # deployments agreeing on (tag_id, name, limit, period_seconds, + # scope_by_key_hash) share one bucket even if they set key_ttl_seconds or + # max_in_memory_cache_size differently. Whichever deployment's entry is + # seen first for a given signature supplies those fields for the whole + # group -- an arbitrary but deterministic tie-break, consistent with the + # first-seen-order precedent already established above. + representative_entry_by_signature: Final[dict[tuple[str, str, float, int, bool], TagRateLimitEntry]] = {} # mutable-ok: see comment above for deployment in deployments: dep_id = _deployment_id(deployment) if dep_id is None: @@ -272,6 +273,7 @@ def _build_group_limits(deployments: Sequence[Mapping[str, object]], unit: _Limi signature = (entry.tag_id, entry.name, entry.limit, entry.period_seconds, entry.scope_by_key_hash) ids_for_signature = declaring_ids_by_signature.setdefault(signature, []) # mutable-ok: see comment above ids_for_signature.append(dep_id) + representative_entry_by_signature.setdefault(signature, entry) # mutable-ok: see comment above distinct_signature_count_by_name: Final[Mapping[tuple[str, str], int]] = MappingProxyType( { @@ -291,7 +293,7 @@ def _build_group_limits(deployments: Sequence[Mapping[str, object]], unit: _Limi if ( configured_limit := _configured_limit_for_signature( unit, - signature, + representative_entry_by_signature[signature], declaring_ids, is_chain_wide=( distinct_signature_count_by_name[(signature[0], signature[1])] == 1 @@ -446,19 +448,24 @@ _CONCURRENCY_MIN_SAFETY_TTL_SECONDS: Final = 3600 # Concurrency reservation keys accumulated for the current logical request, -# not yet released. Held via a ContextVar bound to a mutable holder object -# (not an immutable tuple rebound with `.set()`) because `asyncio.create_task` -# only copies which *object* a ContextVar is bound to, not a snapshot of that -# object's contents: a `.set()` performed inside a task forked off this -# context mutates only that task's own binding, invisible to the parent task -# that continues on to a fallback hop. Mutating a shared holder in place is +# not yet released, paired with the cache-size override (from +# TagRateLimitEntry.max_in_memory_cache_size) each reservation was +# incremented under: releasing a reservation must decrement the exact same +# cache partition it was incremented on, or the release silently no-ops on +# the wrong (default) partition and the reservation leaks forever. Held via +# a ContextVar bound to a mutable holder object (not an immutable tuple +# rebound with `.set()`) because `asyncio.create_task` only copies which +# *object* a ContextVar is bound to, not a snapshot of that object's +# contents: a `.set()` performed inside a task forked off this context +# mutates only that task's own binding, invisible to the parent task that +# continues on to a fallback hop. Mutating a shared holder in place is # visible from every task forked after the holder was first created, # regardless of which task performs the mutation. class _PendingConcurrencyKeys: __slots__ = ("keys",) def __init__(self) -> None: - self.keys: list[str] = [] # mutable-ok: shared across asyncio.create_task forks by design; see class docstring + self.keys: list[tuple[str, "_PartitionKey"]] = [] # mutable-ok: shared across asyncio.create_task forks by design; see class docstring _pending_concurrency_keys: Final[contextvars.ContextVar[_PendingConcurrencyKeys | None]] = contextvars.ContextVar( @@ -639,6 +646,29 @@ def _resolve_max_in_memory_cache_size() -> int | None: return None +# None => this entry shares the hook's single default cache partition +# (matching every entry's behavior before this override existed). Otherwise +# a value-stable signature -- not the override int alone -- so two different +# entries that happen to choose the identical max_in_memory_cache_size don't +# get merged into one shared partition; the same entry (same config content) +# always resolves to the same signature across index rebuilds, which is what +# keeps _PROXY_TagRateLimiter._partitions from leaking a fresh partition +# every time _TagRateLimitIndex rebuilds and reconstructs `_ConfiguredLimit`s. +_PartitionKey: TypeAlias = tuple[str, str, float, int, bool, int] | None + + +def _partition_key(entry: TagRateLimitEntry) -> _PartitionKey: + if entry.max_in_memory_cache_size is None: + return None + return (entry.tag_id, entry.name, entry.limit, entry.period_seconds, entry.scope_by_key_hash, entry.max_in_memory_cache_size) + + +@dataclass(frozen=True, slots=True) +class _CachePartition: + internal_usage_cache: InternalUsageCache + v3: _PROXY_MaxParallelRequestsHandler_v3 + + class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only referenced via the deferred import in litellm_logging.py's callback resolver; basedpyright doesn't trace that usage CustomLogger ): @@ -654,24 +684,28 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer # flooding this hook's own caller-controlled tag buckets past that # ceiling could evict an unrelated, authentication-bound counter and # exceed a limit nothing here configured. The real Redis connection - # (if any) is still shared, so cross-instance correctness is unaffected. - # - # This cache's own 200-item default is still shared across every - # distinct tag value this hook sees. Deployments rate-limiting on a - # high-cardinality tag_id (e.g. per end user) without Redis can raise - # `litellm_settings.tag_rate_limiter_max_in_memory_cache_size` so - # active buckets aren't evicted before their period elapses. - isolated_dual_cache: Final = DualCache( - in_memory_cache=InMemoryCache(max_size_in_memory=_resolve_max_in_memory_cache_size()), - redis_cache=internal_usage_cache.redis_cache, - ) - self.internal_usage_cache = InternalUsageCache(dual_cache=isolated_dual_cache) - self._v3 = _PROXY_MaxParallelRequestsHandler_v3(self.internal_usage_cache, time_provider=time_provider) + # (if any) is still shared across every partition (see _build_partition), + # so cross-instance correctness is unaffected regardless of partitioning. + self._redis_cache: Final = internal_usage_cache.redis_cache self._time_provider = time_provider or datetime.now + # Every distinct _PartitionKey gets its own dedicated partition + # (in-memory cache + its own v3 handler), lazily built and memoized + # here -- see _partition_for. None (the key every entry uses unless + # it sets its own max_in_memory_cache_size) is this hook's single + # default partition, sized by + # litellm.tag_rate_limiter_max_in_memory_cache_size (200 if that's + # also unset), matching today's behavior for every entry that doesn't + # opt into its own partition. + self._partitions: dict[_PartitionKey, _CachePartition] = {} # mutable-ok: lazily memoized per distinct partition key, guarded by _partitions_lock; see _partition_for + self._partitions_lock = asyncio.Lock() + default_partition: Final = self._build_partition(_resolve_max_in_memory_cache_size()) + self._partitions[None] = default_partition + self.internal_usage_cache = default_partition.internal_usage_cache + self._v3 = default_partition.v3 self._index = _TagRateLimitIndex(time_provider=self._time_provider) self._lock = asyncio.Lock() self.llm_router: Router | None = None - redis_cache: Final = self.internal_usage_cache.dual_cache.redis_cache + redis_cache: Final = self._redis_cache self._check_and_incr_script = ( redis_cache.async_register_script(TAG_RL_CHECK_AND_INCR_SCRIPT) if redis_cache is not None else None ) @@ -682,7 +716,33 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer def update_variables(self, llm_router: Router) -> None: self.llm_router = llm_router - async def _check_and_increment_one(self, key: str, limit: float, increment: float, ttl: int) -> tuple[bool, float]: + def _build_partition(self, cache_size_override: int | None) -> _CachePartition: + dual_cache: Final = DualCache( + in_memory_cache=InMemoryCache(max_size_in_memory=cache_size_override), + redis_cache=self._redis_cache, + ) + cache: Final = InternalUsageCache(dual_cache=dual_cache) + return _CachePartition( + internal_usage_cache=cache, + v3=_PROXY_MaxParallelRequestsHandler_v3(cache, time_provider=self._time_provider), + ) + + async def _partition_for(self, partition_key: _PartitionKey) -> _CachePartition: + existing: Final = self._partitions.get(partition_key) + if existing is not None: + return existing + async with self._partitions_lock: + existing_after_lock: Final = self._partitions.get(partition_key) + if existing_after_lock is not None: + return existing_after_lock + cache_size_override: Final = partition_key[-1] if partition_key is not None else None + built: Final = self._build_partition(cache_size_override) + self._partitions[partition_key] = built # mutable-ok: lazily memoized per distinct partition key, guarded by _partitions_lock above + return built + + async def _check_and_increment_one( + self, cache: InternalUsageCache, key: str, limit: float, increment: float, ttl: int + ) -> tuple[bool, float]: """Single-key atomic check-and-increment. Always one key per Lua call -- see TAG_RL_CHECK_AND_INCR_SCRIPT's module docstring for why.""" if self._check_and_incr_script is not None: @@ -690,40 +750,33 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer return bool(raw[0]), float(raw[1]) async with self._lock: - current_value: Final = await self.internal_usage_cache.async_get_cache( - key=key, litellm_parent_otel_span=None - ) + current_value: Final = await cache.async_get_cache(key=key, litellm_parent_otel_span=None) current: Final = float(current_value) if current_value is not None else 0.0 if current + increment > limit: return False, current new_value: Final = current + increment - await self.internal_usage_cache.async_set_cache( - key=key, value=new_value, ttl=ttl, litellm_parent_otel_span=None - ) + await cache.async_set_cache(key=key, value=new_value, ttl=ttl, litellm_parent_otel_span=None) return True, new_value - async def _decrement_floor_zero(self, key: str, delta: float) -> None: + async def _decrement_floor_zero(self, cache: InternalUsageCache, key: str, delta: float) -> None: if self._decr_floor_zero_script is not None: await self._decr_floor_zero_script(keys=(key,), args=(delta,)) return async with self._lock: - current_value: Final = await self.internal_usage_cache.async_get_cache( - key=key, litellm_parent_otel_span=None - ) + current_value: Final = await cache.async_get_cache(key=key, litellm_parent_otel_span=None) current: Final = float(current_value) if current_value is not None else 0.0 - await self.internal_usage_cache.async_set_cache( - key=key, value=max(0.0, current + delta), litellm_parent_otel_span=None - ) + await cache.async_set_cache(key=key, value=max(0.0, current + delta), litellm_parent_otel_span=None) async def _atomic_check_and_increment( self, - checks: Sequence[tuple[str, float, float, int]], + checks: Sequence[tuple[InternalUsageCache, str, float, float, int]], ) -> tuple[int | None, tuple[float, ...]]: """ - All-or-nothing across every (key, limit, increment, ttl) in `checks`: - if any would exceed its limit, none are incremented -- a single hop's - requests-unit and concurrency-unit checks must commit together or not - at all. Each key is checked/incremented in its own single-key Lua + All-or-nothing across every (cache, key, limit, increment, ttl) in + `checks`: if any would exceed its limit, none are incremented -- a + single hop's requests-unit and concurrency-unit checks must commit + together or not at all, even when they span more than one cache + partition. Each key is checked/incremented in its own single-key Lua call (cluster-safe by construction); all-or-nothing across the batch is enforced here by refunding every earlier admission the moment a later key is rejected, not by a single multi-key script call. @@ -748,15 +801,15 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer # accumulated so far in favor of refunding and returning early, so # this can't be expressed as a one-shot comprehension. admitted_values: Final = [] # mutable-ok: sequential async accumulator, discardable on early rejection; see comment above - for index, (key, limit, increment, ttl) in enumerate(checks): - admitted, value = await self._check_and_increment_one(key, limit, increment, ttl) + for index, (cache, key, limit, increment, ttl) in enumerate(checks): + admitted, value = await self._check_and_increment_one(cache, key, limit, increment, ttl) if admitted: admitted_values.append(value) # mutable-ok: see accumulator comment above continue for refund_index in range(index): - refund_key, _limit, refund_increment, _ttl = checks[refund_index] + refund_cache, refund_key, _limit, refund_increment, _ttl = checks[refund_index] try: - await self._decrement_floor_zero(refund_key, -refund_increment) + await self._decrement_floor_zero(refund_cache, refund_key, -refund_increment) except Exception as e: # noqa: BLE001 - one failed refund must not block refunding the rest verbose_proxy_logger.warning("tag_rate_limiter: failed to refund %s on rollback: %s", refund_key, e) return index, (value,) @@ -820,21 +873,27 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer self._raise_if_over_limit(read_only_checks, current_values, model) if atomic_checks: + atomic_partitions_list: Final = [] # mutable-ok: sequential async lookups, one per atomic_checks entry (a genexpr can't `await` here); zipped with atomic_checks immediately below + for configured_limit, _tag_value, _key in atomic_checks: + atomic_partitions_list.append(await self._partition_for(_partition_key(configured_limit.entry))) # mutable-ok: see comment above + atomic_partitions: Final = tuple(atomic_partitions_list) failing_index, values = await self._atomic_check_and_increment( tuple( - (key, configured_limit.entry.limit, 1.0, self._ttl_for(configured_limit)) - for configured_limit, _tag_value, key in atomic_checks + (partition.internal_usage_cache, key, configured_limit.entry.limit, 1.0, self._ttl_for(configured_limit)) + for partition, (configured_limit, _tag_value, key) in zip(atomic_partitions, atomic_checks) ) ) if failing_index is not None: configured_limit, tag_value, _key = atomic_checks[failing_index] self._raise_over_limit(configured_limit, tag_value, model, current=values[0]) - concurrency_keys: Final = tuple( - key for configured_limit, _tag_value, key in atomic_checks if configured_limit.unit == "concurrency" + concurrency_reservations: Final = tuple( + (key, _partition_key(configured_limit.entry)) + for configured_limit, _tag_value, key in atomic_checks + if configured_limit.unit == "concurrency" ) - if concurrency_keys: - _pending_concurrency_holder().keys.extend(concurrency_keys) + if concurrency_reservations: + _pending_concurrency_holder().keys.extend(concurrency_reservations) return healthy_deployments @@ -860,13 +919,33 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer ) -> tuple[float | None, ...]: if not read_only_checks: return () - keys: Final = tuple(key for _cfg, _tag_value, key in read_only_checks) - current_values: Final = await self.internal_usage_cache.async_batch_get_cache( - keys=list(keys), # mutable-ok: async_batch_get_cache requires a real list; converted only at this boundary - parent_otel_span=parent_otel_span, - local_only=False, - ) - return tuple(current_values) if current_values is not None else tuple(None for _ in keys) + + # Grouped by cache partition (one batched read per partition), then + # reassembled back into read_only_checks's original order: a hop can + # mix entries from more than one partition (e.g. a default-cache + # dollar_limits entry alongside a dedicated-partition request_limits + # entry), and _raise_if_over_limit below zips this result positionally + # against read_only_checks, so order must be preserved exactly. + indices_by_partition: Final[dict[_PartitionKey, list[int]]] = {} # mutable-ok: groups positions sharing a cache partition; reassembled into original order below + for index, (configured_limit, _tag_value, _key) in enumerate(read_only_checks): + indices_by_partition.setdefault(_partition_key(configured_limit.entry), []).append(index) # mutable-ok: see comment above + + values_by_index: Final[dict[int, float | None]] = {} # mutable-ok: see comment above + for partition_key, indices in indices_by_partition.items(): + # not `Final`: rebound each loop iteration, which basedpyright's + # LIT010/Final-in-loop check forbids + partition = await self._partition_for(partition_key) + keys = [read_only_checks[i][2] for i in indices] # mutable-ok: async_batch_get_cache requires a real list; converted only at this boundary + current_values = await partition.internal_usage_cache.async_batch_get_cache( + keys=keys, + parent_otel_span=parent_otel_span, + local_only=False, + ) + resolved = current_values if current_values is not None else [None] * len(keys) + for i, value in zip(indices, resolved): + values_by_index[i] = value # mutable-ok: see comment above + + return tuple(values_by_index[i] for i in range(len(read_only_checks))) def _raise_if_over_limit( self, @@ -913,7 +992,7 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer llm_provider="litellm_proxy", ) - async def _release_keys(self, keys: Sequence[str]) -> None: + async def _release_keys(self, reservations: Sequence[tuple[str, _PartitionKey]]) -> None: """ Release each key by one slot. This does not verify the completing request still owns a live reservation (no per-request slot identity @@ -923,15 +1002,20 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer held. Flooring at 0 (TAG_RL_DECR_FLOOR_ZERO_SCRIPT) bounds the damage to under-counting (briefly under-enforcing the limit) rather than a negative counter, which would admit unlimited requests. + + Each reservation is released against the exact cache partition + (`_partition_for(partition_key)`) its increment used -- see + `_PendingConcurrencyKeys`'s docstring for why this must match. """ - for key in keys: + for key, partition_key in reservations: try: - await self._decrement_floor_zero(key, -1.0) + partition = await self._partition_for(partition_key) # not Final: rebound each loop iteration + await self._decrement_floor_zero(partition.internal_usage_cache, key, -1.0) except Exception as e: # noqa: BLE001 - releasing a slot must never raise into the caller's request path verbose_proxy_logger.warning("tag_rate_limiter: failed to release concurrency slot %s: %s", key, e) @staticmethod - def _pop_pending_concurrency_keys() -> tuple[str, ...]: + def _pop_pending_concurrency_keys() -> tuple[tuple[str, _PartitionKey], ...]: # Snapshot then remove only those exact keys, never a blanket clear: # a sibling hop can still be live and appending to the same shared # holder concurrently (see the holder's own comment above), so @@ -1014,8 +1098,8 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer } ) - operations: Final = tuple( - operation + operation_by_limit: Final = tuple( + (configured_limit, operation) for configured_limit in configured if ( operation := _increment_operation_for_limit( @@ -1025,12 +1109,23 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer is not None ) - if not operations: + if not operation_by_limit: return - asyncio.create_task( - self._v3.async_increment_tokens_with_ttl_preservation( - pipeline_operations=operations, - parent_otel_span=_get_parent_otel_span_from_kwargs(kwargs), + # Grouped by cache partition: a hop's tokens/dollars entries can span + # more than one partition, and each partition owns its own v3 + # handler (see _build_partition), so each group's operations are + # pipelined through that partition's own handler. + operations_by_partition: Final[dict[_PartitionKey, list[RedisPipelineIncrementOperation]]] = {} # mutable-ok: groups operations by cache partition before dispatching each group's pipeline call + for configured_limit, operation in operation_by_limit: + operations_by_partition.setdefault(_partition_key(configured_limit.entry), []).append(operation) # mutable-ok: see comment above + + parent_otel_span: Final = _get_parent_otel_span_from_kwargs(kwargs) + for partition_key, group_operations in operations_by_partition.items(): + partition = await self._partition_for(partition_key) # not Final: rebound each loop iteration + asyncio.create_task( + partition.v3.async_increment_tokens_with_ttl_preservation( + pipeline_operations=tuple(group_operations), + parent_otel_span=parent_otel_span, + ) ) - ) diff --git a/litellm/types/router.py b/litellm/types/router.py index f99f7890f09..d8b2e457706 100644 --- a/litellm/types/router.py +++ b/litellm/types/router.py @@ -150,6 +150,15 @@ class TagRateLimitEntry(BaseModel): # this lets an operator shed them sooner without shortening # period_seconds itself. key_ttl_seconds: int | None = None + # Overrides the size of the dedicated in-memory cache partition this + # entry's own keys live in, when Redis isn't configured (or as a local + # fast-path cache when it is). Unset means this entry shares the hook's + # single default partition, sized by + # litellm.tag_rate_limiter_max_in_memory_cache_size (200 if that's also + # unset). A high-cardinality tag_id can churn past that shared cap and + # evict another entry's active counters; setting this gives the entry + # its own dedicated partition instead. + max_in_memory_cache_size: int | None = None model_config = ConfigDict(protected_namespaces=()) @@ -165,6 +174,12 @@ class TagRateLimitEntry(BaseModel): raise ValueError("key_ttl_seconds must be a positive integer when set") return self + @model_validator(mode="after") + def _validate_max_in_memory_cache_size(self) -> "TagRateLimitEntry": + if self.max_in_memory_cache_size is not None and self.max_in_memory_cache_size <= 0: + raise ValueError("max_in_memory_cache_size must be a positive integer when set") + return self + class TagRateLimitGroup(BaseModel): limits: tuple[TagRateLimitEntry, ...] = () diff --git a/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py b/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py index bd7cf85e92a..d8df1fc88f0 100644 --- a/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py +++ b/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py @@ -1474,6 +1474,41 @@ def test_build_limits_index_is_also_keyed_by_team_public_model_name(): assert by_alias[0].team_scope == "team-1" +def test_build_limits_index_preserves_key_ttl_seconds_and_max_in_memory_cache_size(): + """ + Regression test: _configured_limit_for_signature used to reconstruct a + fresh TagRateLimitEntry from a 5-field dedup signature that didn't + include key_ttl_seconds or max_in_memory_cache_size, silently resetting + both to None for every entry that went through the real indexing path + (which is every entry reachable from async_filter_deployments / + async_log_success_event) -- only entries built directly in a test, never + through _build_limits_index, kept their configured values. + """ + deployment = _deployment( + "grp", + "dep-1", + { + "request_limits": { + "limits": [ + { + "name": "user_cap", + "tag_id": "end_user_id", + "limit": 5, + "period_seconds": 60, + "key_ttl_seconds": 120, + "max_in_memory_cache_size": 500, + } + ] + } + }, + ) + index = _build_limits_index([deployment]) + configured = index.resolve("grp", team_id=None) + assert len(configured) == 1 + assert configured[0].entry.key_ttl_seconds == 120 + assert configured[0].entry.max_in_memory_cache_size == 500 + + def test_build_limits_index_keeps_different_teams_same_alias_separate(): """ `team_public_model_name` is only unique per team: Router itself lets two @@ -1845,7 +1880,7 @@ async def test_cross_unit_refund_leaves_no_phantom_increment_in_memory(time_cont async def test_release_floors_at_zero_instead_of_going_negative(time_controller): limiter = _make_limiter(time_controller) key = "{tag_rl:test:concurrency:floor:chain:u1}:inflight" - await limiter._decrement_floor_zero(key, -1.0) + await limiter._decrement_floor_zero(limiter.internal_usage_cache, key, -1.0) value = await limiter.internal_usage_cache.async_get_cache(key=key, litellm_parent_otel_span=None) assert (float(value) if value is not None else 0.0) == 0.0 @@ -1869,18 +1904,18 @@ async def test_refund_failure_on_one_key_does_not_block_others_or_raise(time_con rejecting_key = "{tag_rl:test:refund-fail:c}:requests" class _FlakyLimiter(_PROXY_TagRateLimiter): - async def _decrement_floor_zero(self, key: str, delta: float) -> None: + async def _decrement_floor_zero(self, cache, key: str, delta: float) -> None: if key == failing_key: raise RuntimeError("simulated transient redis failure") - await super()._decrement_floor_zero(key, delta) + await super()._decrement_floor_zero(cache, key, delta) flaky = _FlakyLimiter(internal_usage_cache=DualCache(), time_provider=time_controller.now) failing_index, values = await flaky._atomic_check_and_increment( [ - (failing_key, 10.0, 1.0, 60), - (other_key, 10.0, 1.0, 60), - (rejecting_key, 0.0, 1.0, 60), + (flaky.internal_usage_cache, failing_key, 10.0, 1.0, 60), + (flaky.internal_usage_cache, other_key, 10.0, 1.0, 60), + (flaky.internal_usage_cache, rejecting_key, 0.0, 1.0, 60), ] ) @@ -2343,3 +2378,268 @@ def test_ttl_for_concurrency_never_drops_below_the_safety_floor_even_with_a_lowe def test_tag_rate_limit_entry_rejects_non_positive_key_ttl_seconds(): with pytest.raises(ValueError): TagRateLimitEntry(name="per_minute", limit=1, period_seconds=60, key_ttl_seconds=0) + + +# --------------------------------------------------------------------------- +# per-tag max_in_memory_cache_size override -- dedicated cache partitions +# --------------------------------------------------------------------------- + + +def test_tag_rate_limit_entry_rejects_non_positive_max_in_memory_cache_size(): + with pytest.raises(ValueError): + TagRateLimitEntry(name="per_minute", limit=1, period_seconds=60, max_in_memory_cache_size=0) + + +def _two_request_limit_router(team_limit: int, user_limit: int, user_cache_size: int | None) -> "litellm.Router": + return litellm.Router( + model_list=[ + _deployment( + "grp", + "dep-1", + { + "request_limits": { + "limits": [ + {"name": "team_cap", "tag_id": "team_id", "limit": team_limit, "period_seconds": 60}, + { + "name": "user_cap", + "tag_id": "end_user_id", + "limit": user_limit, + "period_seconds": 60, + "max_in_memory_cache_size": user_cache_size, + }, + ] + } + }, + ) + ] + ) + + +@pytest.mark.asyncio +async def test_max_in_memory_cache_size_override_isolates_a_flood_on_that_entry_from_a_default_partition_entry( + time_controller, +): + """ + An entry with its own max_in_memory_cache_size gets a dedicated cache + partition. Flooding that entry's own high-cardinality tag values must + never evict a *different* entry's bucket that was never given an + override and still lives on the hook's single default partition. + """ + limiter = _make_limiter(time_controller) + router = _two_request_limit_router(team_limit=1, user_limit=1000, user_cache_size=5) + limiter.update_variables(llm_router=router) + healthy = router.model_list + + # team_cap's bucket (default partition) is created and admitted once. + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs={"metadata": {"tags": ["team_id:t1"]}} + ) + + # Flood user_cap's own dedicated partition (cap=5) past its own capacity + # many times over -- this must stay fully confined to user_cap's own + # partition and never touch team_cap's default-partition bucket. + for i in range(250): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": [f"end_user_id:flood-{i}"]}}, + ) + + # team_cap's bucket must still be at its limit (1) -- a second team_id:t1 + # request is rejected. If it had been evicted by user_cap's flood, this + # would instead admit (a fresh, zeroed counter). + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": ["team_id:t1"]}}, + ) + + +@pytest.mark.asyncio +async def test_two_entries_sharing_the_identical_max_in_memory_cache_size_still_get_separate_partitions( + time_controller, +): + """ + Partitions are keyed by the entry's full signature, not the override + value alone: two unrelated entries that happen to choose the identical + max_in_memory_cache_size must not be merged into one shared cache, or + flooding one would evict the other's bucket exactly like the bug this + override exists to fix. + """ + limiter = _make_limiter(time_controller) + # Both team_cap and user_cap set the identical max_in_memory_cache_size (5). + router = litellm.Router( + model_list=[ + _deployment( + "grp", + "dep-1", + { + "request_limits": { + "limits": [ + { + "name": "team_cap", + "tag_id": "team_id", + "limit": 1, + "period_seconds": 60, + "max_in_memory_cache_size": 5, + }, + { + "name": "user_cap", + "tag_id": "end_user_id", + "limit": 1000, + "period_seconds": 60, + "max_in_memory_cache_size": 5, + }, + ] + } + }, + ) + ] + ) + limiter.update_variables(llm_router=router) + healthy = router.model_list + + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs={"metadata": {"tags": ["team_id:t1"]}} + ) + + for i in range(250): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": [f"end_user_id:flood-{i}"]}}, + ) + + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": ["team_id:t1"]}}, + ) + + +@pytest.mark.asyncio +async def test_concurrency_slot_with_a_cache_size_override_is_released_against_the_same_partition(time_controller): + """ + A concurrency reservation on an entry with its own max_in_memory_cache_size + must be released against that same dedicated partition. If the release + path fell back to the default partition instead, it would silently no-op + (nothing to decrement there) and the reservation would leak forever. + """ + router = litellm.Router( + model_list=[ + _deployment( + "grp", + "dep-1", + { + "concurrency_limits": { + "limits": [ + { + "name": "inflight", + "tag_id": "end_user_id", + "limit": 1, + "period_seconds": 300, + "max_in_memory_cache_size": 10, + } + ] + } + }, + ) + ] + ) + limiter = _make_limiter(time_controller) + limiter.update_variables(llm_router=router) + healthy = router.model_list + + kwargs = {"metadata": {"tags": ["end_user_id:u1"]}} + await limiter.async_filter_deployments(model="grp", healthy_deployments=healthy, messages=None, request_kwargs=kwargs) + + # At capacity: a second concurrent reservation for the same tag is rejected. + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": ["end_user_id:u1"]}}, + ) + + # The first request completes -- its slot is released against the + # overridden partition -- freeing capacity again. + kwargs["standard_logging_object"] = { + "model_group": "grp", + "model_id": "dep-1", + "total_tokens": 0, + "response_cost": 0, + } + await limiter.async_log_success_event(kwargs=kwargs, response_obj=None, start_time=0, end_time=0) + await asyncio.sleep(0) + + result = await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": ["end_user_id:u1"]}}, + ) + assert result == healthy + + +@pytest.mark.asyncio +async def test_token_accounting_with_a_cache_size_override_lands_on_that_entrys_own_partition(time_controller): + """ + tokens/dollars increments go through a per-partition v3 handler (grouped + in async_log_success_event), not always the default one -- an entry with + its own max_in_memory_cache_size must have its usage actually accounted, + not silently dropped or misrouted to the default partition's handler. + """ + router = litellm.Router( + model_list=[ + _deployment( + "grp", + "dep-1", + { + "token_limits": { + "limits": [ + { + "name": "daily", + "tag_id": "end_user_id", + "limit": 100, + "period_seconds": 86400, + "max_in_memory_cache_size": 10, + } + ] + } + }, + ) + ] + ) + limiter = _make_limiter(time_controller) + limiter.update_variables(llm_router=router) + healthy = router.model_list + + kwargs = { + "metadata": {"tags": ["end_user_id:u1"]}, + "standard_logging_object": { + "model_group": "grp", + "model_id": "dep-1", + "total_tokens": 150, + "response_cost": 0, + }, + } + await limiter.async_log_success_event(kwargs=kwargs, response_obj=None, start_time=0, end_time=0) + await asyncio.sleep(0) + + # 150 tokens already used, over the limit of 100 -- the next admission + # check must reject. If the increment had been silently dropped (never + # reaching the overridden partition), this would incorrectly admit. + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", + healthy_deployments=healthy, + messages=None, + request_kwargs={"metadata": {"tags": ["end_user_id:u1"]}}, + ) diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index a707f516c4d..1f5bffbf00b 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -35116,6 +35116,8 @@ export interface components { key_ttl_seconds?: number | null; /** Limit */ limit: number; + /** Max In Memory Cache Size */ + max_in_memory_cache_size?: number | null; /** Name */ name: string; /** Period Seconds */