mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
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.
This commit is contained in:
parent
68405241a0
commit
ac3d86a6cd
2 changed files with 102 additions and 25 deletions
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
"""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue