mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix(rate-limiting): fold a stale requests renewal into the atomic batch
Cursor Bugbot follow-up finding on the previous commit: refunding a prior hop's "requests" charge unconditionally at the top of the next hop's admission, before knowing whether that hop would itself be admitted, could undercount a logical request. If the next hop went on to fail a different check (a read-only limit, or another entry in the same atomic batch), the refund had already committed with nothing to replace it, leaving a genuine, real attempt charged at zero and letting a caller bypass the cap under fallback pressure. _release_stale_hop_reservations no longer refunds a stale "requests" entry; it returns the stale key instead. A "requests" check whose key matches one already charged by a superseded earlier hop now renews at zero net cost as part of the same all-or-nothing atomic batch as every other check on that hop, rather than as a separate refund-then-recharge sequence. If the batch rolls back for any other reason, refunding a zero-cost renewal is a genuine no-op, so the earlier hop's real charge is left exactly as it was. Regression test reproduces the exact scenario: a hop fails for real, an unrelated request claims the now-free concurrency slot, this request's next hop is rejected on concurrency (not requests), and a fresh probe confirms the original requests charge still stands. Live-verified the original multi-hop scenario still works correctly end to end.
This commit is contained in:
parent
ed726bef7b
commit
334de9031b
2 changed files with 112 additions and 12 deletions
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
"""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue