diff --git a/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py b/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py index b13e5719b73..d1850eaeab2 100644 --- a/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py +++ b/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py @@ -1301,7 +1301,7 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] return healthy_deployments resolved_request_kwargs: Final = request_kwargs or _EMPTY_MAPPING - await self._release_stale_hop_reservations(resolved_request_kwargs) + stale_request_keys: Final = await self._release_stale_hop_reservations(resolved_request_kwargs) metadata_variable_name: Final = get_metadata_variable_name_from_kwargs(resolved_request_kwargs) team_id: Final = _extract_team_id(resolved_request_kwargs, metadata_variable_name) # Built from the full routing-group membership, not `healthy_deployments` @@ -1372,7 +1372,15 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] partition.internal_usage_cache, key, configured_limit.entry.limit, - 1.0, + # A "requests" key matching one already charged by a + # superseded earlier hop of this same request (see + # _release_stale_hop_reservations) renews that same + # charge at zero net cost instead of adding a second + # unit on top of it -- folded into this same + # all-or-nothing batch so a hop that goes on to fail + # a *different* check here never commits a refund + # with nothing to replace it. + 0.0 if configured_limit.unit == "requests" and key in stale_request_keys else 1.0, self._ttl_for(configured_limit), ) for partition, (configured_limit, _tag_value, key) in zip(atomic_partitions, atomic_checks) @@ -1572,7 +1580,7 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] "model_based_tag_rate_limits_hook: failed to release concurrency slot %s: %s", key, e ) - async def _release_stale_hop_reservations(self, request_kwargs: Mapping[str, object]) -> None: + async def _release_stale_hop_reservations(self, request_kwargs: Mapping[str, object]) -> frozenset[str]: """ A concurrency reservation still queued when a *new* hop's admission runs can only belong to an earlier hop of this same request that @@ -1594,23 +1602,26 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] closes that residual case instead, via the cache mirror `_PENDING_RESERVATIONS_CACHE_KEY_PREFIX` documents. - The identical invariant -- a queued entry still present here can only - belong to an already-failed earlier hop -- holds for a "requests" - atomic increment too, so this also refunds any stale entry queued - under `_PENDING_REQUEST_INCREMENTS_FIELD`; see that field's own - docstring for why, unlike concurrency, a hop that goes on to succeed - (or is the chain's own final failure) is deliberately never refunded. + A "requests" atomic increment queued under + `_PENDING_REQUEST_INCREMENTS_FIELD` is never refunded here, even + though the identical staleness invariant holds for it too: an + unconditional refund followed by this hop's own admission is not one + atomic operation, so a hop that goes on to fail a *different* check + (a read-only limit, or another entry in the same atomic batch) would + leave the refund committed with nothing to replace it, undercounting + a logical request that genuinely made an earlier, real attempt. The + returned keys let `async_filter_deployments` fold the swap into its + own atomic batch instead -- see its own comment for how. """ logging_obj: Final = request_kwargs.get("litellm_logging_obj") model_call_details: Final = getattr(logging_obj, "model_call_details", None) if not isinstance(model_call_details, dict): - return + return frozenset() release_keys: Final = await self._pop_pending_concurrency_keys(model_call_details) if release_keys: await self._release_keys(release_keys) stale_request_increments: Final = _pop_reservations(model_call_details, _PENDING_REQUEST_INCREMENTS_FIELD) - if stale_request_increments: - await self._release_keys(stale_request_increments) + return frozenset(key for key, _partition_key in stale_request_increments) async def _pop_pending_concurrency_keys( self, kwargs: Mapping[str, object] diff --git a/tests/test_litellm/proxy/hooks/test_model_based_tag_rate_limits_hook.py b/tests/test_litellm/proxy/hooks/test_model_based_tag_rate_limits_hook.py index dffbe79ddc2..9a2c986c248 100644 --- a/tests/test_litellm/proxy/hooks/test_model_based_tag_rate_limits_hook.py +++ b/tests/test_litellm/proxy/hooks/test_model_based_tag_rate_limits_hook.py @@ -3006,6 +3006,95 @@ async def test_successful_hops_own_request_increment_is_not_refunded(time_contro ) +@pytest.mark.asyncio +async def test_a_hops_own_rejection_on_a_different_check_does_not_undercount_the_prior_hops_request_charge( + time_controller, +): + """ + Regression test for Cursor Bugbot's follow-up finding on this exact fix: + an earlier version refunded the prior hop's "requests" charge + unconditionally at the top of the next hop's admission, before knowing + whether that next hop would itself be admitted. If the next hop then + failed a *different* check (here, concurrency) before ever reaching its + own requests renewal, the refund had already committed with nothing to + replace it -- a logical request that genuinely made one real attempt + (hop 1) would end up charged zero, letting a caller bypass the requests + cap simply by having a later hop collide with someone else's + concurrency slot. + + Fixed by folding the renewal into the same all-or-nothing atomic batch + as every other check on that hop: a "requests" key matching an earlier + hop's charge renews at zero net cost instead of being refunded first, + so a batch-wide rollback (concurrency's own rejection here) refunds that + zero-cost renewal -- a genuine no-op -- leaving hop 1's real charge + exactly as it was. + """ + # Two independent tag identities: "requests" is scoped to end_user_id + # (private to our own request, never shared with the unrelated + # contender below), "concurrency" is scoped to a separate shared_pool + # tag that both our request and the unrelated contender carry, so they + # compete for the same slot without also colliding on the requests cap. + limiter = _make_limiter(time_controller) + router = litellm.Router( + model_list=[ + _deployment( + "grp", + "dep-1", + { + "request_limits": { + "limits": [{"name": "per_period", "tag_id": "end_user_id", "limit": 1, "period_seconds": 300}] + }, + "concurrency_limits": { + "limits": [{"name": "inflight", "tag_id": "shared_pool", "limit": 1, "period_seconds": 300}] + }, + }, + ) + ] + ) + limiter.update_variables(llm_router=router) + healthy = router.model_list + + # Hop 1 of our request admits (claiming both the requests unit and the + # only concurrency slot), then fails for real -- its concurrency + # reservation is released the normal way, but its requests charge is + # left queued as this-hop's-own-charge, not refunded. + request_kwargs, kwargs = _call_context(["end_user_id:u1", "shared_pool:pool-a"]) + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs=request_kwargs + ) + await limiter.async_log_failure_event(kwargs=kwargs, response_obj=None, start_time=0, end_time=0) + + # A second, unrelated request -- no end_user_id tag at all, so it never + # touches the requests bucket -- now claims the concurrency slot our + # hop 1 just released, and holds it. + other_request_kwargs, other_kwargs = _call_context(["shared_pool:pool-a"]) + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs=other_request_kwargs + ) + + # Hop 2 of our original request: its own "requests" renewal would + # trivially succeed alone (net zero cost), but the concurrency slot is + # now held by the unrelated request above, so the whole atomic batch + # must reject -- and must NOT leave hop 1's requests charge refunded. + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs=request_kwargs + ) + + # The unrelated request finishes, freeing the concurrency slot again. + await limiter.async_log_success_event(kwargs=other_kwargs, response_obj=None, start_time=0, end_time=0) + + # A fresh probe against the same end_user_id tag, with concurrency now + # free, must still be rejected by the requests cap: hop 1's real attempt + # already spent the only unit for this period, and it must not have + # been silently erased by hop 2's unrelated, different-check rejection. + probe_request_kwargs, _probe_kwargs = _call_context(["end_user_id:u1"]) + with pytest.raises(ProxyRateLimitError): + await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs=probe_request_kwargs + ) + + @pytest.mark.asyncio async def test_own_rejection_does_not_release_a_live_reservation(time_controller): """