From ac3d86a6cdb5003979f37a4c514a3de8a7f30969 Mon Sep 17 00:00:00 2001 From: Deepanshu Date: Tue, 25 Aug 2026 20:47:24 -0400 Subject: [PATCH] fix(rate-limiting): peek instead of pop for the requests renewal field Cursor Bugbot follow-up finding on the previous commit: _release_stale_hop_reservations popped _PENDING_REQUEST_INCREMENTS_FIELD unconditionally at the top of a new hop's admission, but only re-queued it after that hop's own atomic batch fully succeeded. A hop that failed earlier (a read-only check, or a different entry in the same atomic batch) left the real counter correctly charged (a 0.0-increment renewal rolls back to a genuine no-op) but the bookkeeping field empty, so a later hop's own peek found nothing to renew and charged a fresh unit on top of the one already sitting in the real counter. The field is now peeked, never popped, for "requests" -- concurrency still pops, since concurrency actively releases on every outcome, but a "requests" renewal only ever needs to know whether a key is already charged, never to clear it. Queuing at the end of a hop's own successful batch now skips keys already present, avoiding unbounded duplicate growth across a long retry chain. This matches global_tag_rate_limits_hook's already-correct pattern (an append-only list on the stash, never cleared), which never had this bug in the first place. Regression test reuses the same request_kwargs across three hops -- unlike the existing undercount test, whose final probe uses a fresh, independent context and so can't distinguish "the real counter is correct" from "the bookkeeping a future hop needs is intact" -- confirmed red on the popping version, green on this one. --- .../hooks/model_based_tag_rate_limits_hook.py | 48 ++++++----- .../test_model_based_tag_rate_limits_hook.py | 79 +++++++++++++++++++ 2 files changed, 102 insertions(+), 25 deletions(-) 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 d1850eaeab2..010996e29f9 100644 --- a/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py +++ b/litellm/proxy/hooks/model_based_tag_rate_limits_hook.py @@ -1050,28 +1050,6 @@ def _queue_pending_reservations( pending.extend(reservations) # mutable-ok: see comment above -def _pop_reservations(model_call_details: Mapping[str, object], field: str) -> tuple[tuple[str, "_PartitionKey"], ...]: - """No external cache-mirror interaction -- only - `_PENDING_CONCURRENCY_KEYS_FIELD` needs that (see its docstring); a - "requests" entry here never needs a final-hop release path, so this is - the whole mechanism. Snapshots then removes individual items from the - same list object rather than a blanket pop of `field` itself, matching - `_pop_pending_concurrency_keys`'s own reasoning: a sibling hop sharing - this request's `model_call_details` can still be live and appending - concurrently, so clearing the whole field here could silently strand - that entry instead of it being refunded or left to stand later.""" - pending = model_call_details.get(field) - if not isinstance(pending, list) or not pending: - return () - keys: Final = tuple(pending) - for key in keys: - try: - pending.remove(key) # mutable-ok: see queuing helper above - except ValueError: - pass - return keys - - def _record_admission_time(request_kwargs: Mapping[str, object], now: float) -> None: """Stash this hop's admission timestamp -- see `_ADMISSION_TIME_FIELD`'s docstring for why. Silently a no-op without a real logging object @@ -1405,10 +1383,16 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] concurrency_reservations, ) + # Only genuinely new keys, never one already in stale_request_keys: + # that key's own check just renewed at zero net cost above and is + # still sitting in the field (see _release_stale_hop_reservations' + # own comment on why this is a peek, not a pop) -- appending it + # again here would grow the list with a duplicate entry on every + # hop of a long retry chain without changing what it means. request_increments: Final = tuple( (key, _partition_key(configured_limit.entry)) for configured_limit, _tag_value, key in atomic_checks - if configured_limit.unit == "requests" + if configured_limit.unit == "requests" and key not in stale_request_keys ) if request_increments: _queue_pending_reservations( @@ -1612,6 +1596,18 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] 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. + + Deliberately a peek, not a pop, for that same field: an earlier + version popped it here and only re-queued on a fully successful + atomic batch, so a hop that failed *before* reaching that point (a + read-only check, or a different entry in its own batch) silently + dropped the bookkeeping -- the real counter was untouched (0.0 + renewals roll back to a no-op), but the *next* hop's own peek would + come back empty, no longer recognize the key as already charged, and + charge a fresh unit on top of the one still sitting in the real + counter. Peeking leaves the field exactly as it was for whichever + hop reads it next, regardless of how many hops in between fail + before ever reaching their own successful queuing step. """ logging_obj: Final = request_kwargs.get("litellm_logging_obj") model_call_details: Final = getattr(logging_obj, "model_call_details", None) @@ -1620,8 +1616,10 @@ class _PROXY_ModelBasedTagRateLimitsHook( # pyright: ignore[reportUnusedClass] 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) - return frozenset(key for key, _partition_key in stale_request_increments) + pending_request_increments: Final = model_call_details.get(_PENDING_REQUEST_INCREMENTS_FIELD) + if not isinstance(pending_request_increments, list): + return frozenset() + return frozenset(key for key, _partition_key in pending_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 9a2c986c248..1ad4ec94ff8 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 @@ -3095,6 +3095,85 @@ async def test_a_hops_own_rejection_on_a_different_check_does_not_undercount_the ) +@pytest.mark.asyncio +async def test_a_third_hops_admission_still_recognizes_a_charge_that_survived_a_middle_hops_rejection( + time_controller, +): + """ + Regression test for Cursor Bugbot's follow-up finding on the fix above: + an earlier version had _release_stale_hop_reservations *pop* the pending + "requests" entry, re-queuing it only after this hop's own atomic batch + fully succeeded. A middle hop that failed *before* reaching that point + (exactly the concurrency-rejection scenario the test above covers) left + the real counter correctly charged but the bookkeeping field empty, so + a *third* hop's own peek came back with nothing to renew and charged a + fresh unit on top of the one still sitting in the real counter -- + doubling the charge (or, with a tighter limit, a false 429) despite the + fix that was supposed to prevent exactly that. + + Reuses the same request_kwargs (the same model_call_details) across all + three hops -- unlike the probe above, which deliberately uses a fresh, + independent context and so can't distinguish "the real counter is + correct" from "the bookkeeping that lets a future hop recognize it is + intact"; only a third hop sharing the same chain's own bookkeeping can. + """ + 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 admits (claiming the requests unit and the only concurrency + # slot), then fails for real. + 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) + + # An unrelated request claims the now-free concurrency slot 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 rejects on concurrency (not requests) -- + # its own admission never reaches its own successful queuing step. + 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. Release + # runs as a background task on success (see async_log_success_event's + # own implementation), so let it actually complete before hop 3 checks. + await limiter.async_log_success_event(kwargs=other_kwargs, response_obj=None, start_time=0, end_time=0) + await asyncio.sleep(0) + + # Hop 3 of our original request, same request_kwargs: if the bookkeeping + # survived hop 2's rejection, this renews at zero cost and succeeds. If + # it was lost, this charges a fresh unit on top of hop 1's still-live + # real charge and wrongly rejects. + result = await limiter.async_filter_deployments( + model="grp", healthy_deployments=healthy, messages=None, request_kwargs=request_kwargs + ) + assert result == healthy + + @pytest.mark.asyncio async def test_own_rejection_does_not_release_a_live_reservation(time_controller): """