fix(rate-limiting): dedup routing-group unions and pin down a fire-and-forget release task

Two real findings from Bugbot:

Routing groups charge every member (High): resolve_any unions limits
across every distinct model_name reachable from a routing-group hop's
healthy_deployments, then async_filter_deployments atomically checks and
increments every returned entry for that one hop -- even though only one
member deployment ends up serving. Members declaring an identical
signature and scope now dedup to one shared entry, so the hop reserves
capacity once, not once per member. Members that genuinely disagree on
the limit stay separate, unchanged from before this fix: resolving that
ambiguity needs knowing which deployment gets picked, which isn't known
yet at this admission-time hook.

Success release can leak slots (Medium): async_log_success_event fires
its concurrency release via a bare, unreferenced asyncio.create_task to
keep the hot success path from waiting on a Redis round trip. Per
asyncio.create_task's own docs, the event loop only holds a weak
reference to a task with no other referrer, and by the time this one
would run its keys are already popped out of model_call_details, so a
collected task's release is unrecoverable, not just delayed. Added
_BACKGROUND_RELEASE_TASKS, a module-level set holding a strong reference
for exactly as long as each release task is pending, discarding it via
the task's own done-callback once it completes.
This commit is contained in:
Deepanshu 2026-08-20 12:39:02 -04:00
parent f70e48872f
commit 533ae89c32
2 changed files with 207 additions and 6 deletions

View file

@ -75,6 +75,19 @@ _UNIT_TO_RATE_LIMIT_TYPE: Final[Mapping[_LimitUnit, RateLimitType]] = MappingPro
# every one of these call sites just to immediately call `.get()` on it.
_EMPTY_MAPPING: Final[Mapping[str, object]] = MappingProxyType({})
# `asyncio.create_task`'s own docs: "Save a reference to the result of this
# function, to avoid a task disappearing mid-execution. The event loop only
# keeps weak references to tasks. A task that isn't referenced elsewhere may
# get garbage collected at any time, even before it's done." The success path
# below deliberately fires-and-forgets its release (unlike the failure/
# disconnect paths, which await it directly) to keep the hot success-response
# path from waiting on a Redis round trip; by the time that background task
# would run, its keys have already been popped out of model_call_details, so
# a collected task's release is unrecoverable, not just delayed. Holding a
# strong reference here until the task's own completion callback discards it
# is the standard fix.
_BACKGROUND_RELEASE_TASKS: Final[set["asyncio.Task[None]"]] = set() # mutable-ok: see comment above
# Single-key atomic check-and-increment. Deliberately one key per script call
# (never a batch of differently-hash-tagged keys in one call): every tag_rl
# key carries its own self-contained {..} hash tag so unrelated buckets never
@ -357,15 +370,37 @@ class _LimitsIndex:
stamped with the `model_name` it actually came from (`resolved_group`)
so hashing stays namespaced per underlying group even when the
candidates span more than one `model_name`.
Members declaring the identical signature and scope are deduped to
one shared entry: only one deployment in the group ends up actually
serving a given hop, but every member's own `model_name` is resolved
independently above, so an undeduped union would check and charge
every member's bucket for that one hop -- request/concurrency
capacity a caller never actually used, and a false 429 for a sibling
member that was never over its own limit. Divergent configs across
model_names (different limit/period/scope for the same tag_id+name)
are left as separate entries, same as before this dedup: resolving
that ambiguity needs knowing which deployment will be picked, which
isn't known yet at this admission-time hook.
"""
direct: Final = self.resolve(model, team_id)
if direct:
return direct
return tuple(
replace(limit, resolved_group=name)
for name in frozenset(candidate_model_names)
for limit in self.by_model_name.get(name, ())
)
deduped: Final[dict[tuple[object, ...], _ConfiguredLimit]] = {} # mutable-ok: see docstring above
for name in frozenset(candidate_model_names):
for limit in self.by_model_name.get(name, ()):
key = (
limit.unit,
limit.entry.tag_id,
limit.entry.name,
limit.entry.limit,
limit.entry.period_seconds,
limit.entry.scope_by_key_hash,
limit.deployment_scope,
limit.team_scope,
)
deduped.setdefault(key, replace(limit, resolved_group=name)) # mutable-ok: see docstring above
return tuple(deduped.values())
def _team_alias_key(deployment: Mapping[str, object]) -> tuple[str, str] | None:
@ -1178,7 +1213,9 @@ class _PROXY_TagRateLimiter( # pyright: ignore[reportUnusedClass] # only refer
async def async_log_success_event(self, kwargs, response_obj, start_time, end_time) -> None:
release_keys: Final = self._pop_pending_concurrency_keys(kwargs)
if release_keys:
asyncio.create_task(self._release_keys(release_keys))
release_task: Final = asyncio.create_task(self._release_keys(release_keys))
_BACKGROUND_RELEASE_TASKS.add(release_task) # mutable-ok: see _BACKGROUND_RELEASE_TASKS's own docstring
release_task.add_done_callback(_BACKGROUND_RELEASE_TASKS.discard)
if self.llm_router is None:
return

View file

@ -24,6 +24,7 @@ from litellm.proxy.hooks.tag_rate_limiter import (
_extract_key_hash,
_extract_team_id,
_fixed_length_identity,
_BACKGROUND_RELEASE_TASKS,
_inflight_key,
_partition_key,
_PENDING_CONCURRENCY_KEYS_FIELD,
@ -485,6 +486,79 @@ async def test_filter_deployments_routing_group_does_not_collide_across_differen
assert result == [healthy[1]]
def test_resolve_any_dedups_identical_signature_across_member_model_names():
"""
A real routing-group hop presents every member simultaneously (Router
resolves the group name to its full member list for one filtering pass,
then picks exactly one afterwards), and resolve_any is called once for
that one hop with every member's model_name as a candidate. Two members
declaring the identical concurrency signature must resolve to one shared
entry for that hop, not two: `async_filter_deployments` checks and
atomically increments every entry `resolve_any` returns as belonging to
this one hop, so two entries here means the hop reserves capacity twice
(once per member) even though only one deployment will actually serve --
over-charging the caller's own usage and risking a false 429 against a
sibling member that was never over its own limit.
Admission-level round-trip tests can't distinguish this from "two
separate entries with identical limits, always incremented together":
every hop that presents the same member set moves both buckets in
lockstep regardless of whether they're actually one shared entry or two,
so the dedup can only be verified directly at this level.
"""
concurrency_limits = {
"concurrency_limits": {
"limits": [{"name": "inflight", "tag_id": "end_user_id", "limit": 1, "period_seconds": 300}]
}
}
index = _build_limits_index(
[
_deployment("backend-a", "dep-a", concurrency_limits),
_deployment("backend-b", "dep-b", concurrency_limits),
]
)
resolved = index.resolve_any("my-group", team_id=None, candidate_model_names=("backend-a", "backend-b"))
assert len(resolved) == 1
assert resolved[0].unit == "concurrency"
assert resolved[0].entry.limit == 1
def test_resolve_any_keeps_divergent_signatures_across_member_model_names_separate():
"""
Companion to the dedup test above: members that genuinely disagree on
the limit for the same tag_id+name must not be silently collapsed --
which of two different limits would even apply isn't knowable at this
admission-time hook, before a specific deployment is picked, so both
stay as their own entries (today's pre-existing behavior for a
divergent config, left unchanged by the identical-signature dedup).
"""
index = _build_limits_index(
[
_deployment(
"backend-a",
"dep-a",
{
"concurrency_limits": {
"limits": [{"name": "inflight", "tag_id": "end_user_id", "limit": 1, "period_seconds": 300}]
}
},
),
_deployment(
"backend-b",
"dep-b",
{
"concurrency_limits": {
"limits": [{"name": "inflight", "tag_id": "end_user_id", "limit": 2, "period_seconds": 300}]
}
},
),
]
)
resolved = index.resolve_any("my-group", team_id=None, candidate_model_names=("backend-a", "backend-b"))
assert len(resolved) == 2
assert {c.entry.limit for c in resolved} == {1, 2}
@pytest.mark.asyncio
async def test_filter_deployments_per_entry_fail_open_when_tag_absent(time_controller):
"""
@ -1041,6 +1115,96 @@ async def test_concurrency_slot_released_on_success_frees_capacity(time_controll
assert result == healthy
@pytest.mark.asyncio
async def test_background_release_tasks_registry_holds_a_reference_until_done():
"""
async_log_success_event fires its release via a bare asyncio.create_task
(unlike failure/disconnect, which await it directly) to keep the hot
success path from waiting on a Redis round trip. asyncio.create_task's
own docs warn the event loop only holds a *weak* reference to a task, so
one with no other referrer can be garbage collected before it runs --
and by the time it would run here, its keys are already popped out of
model_call_details, so a collected task's release is unrecoverable, not
merely delayed. _BACKGROUND_RELEASE_TASKS exists to hold a strong
reference for exactly as long as the task is pending, then release it via
the task's own done-callback -- exercised directly here (an Event gate
gives a deterministic pending window; going through the real
async_log_success_event doesn't, since its own further awaits let a fast
in-memory release resolve before a test could ever observe it pending).
"""
assert len(_BACKGROUND_RELEASE_TASKS) == 0
gate = asyncio.Event()
async def _pending_release():
await gate.wait()
task = asyncio.create_task(_pending_release())
_BACKGROUND_RELEASE_TASKS.add(task)
task.add_done_callback(_BACKGROUND_RELEASE_TASKS.discard)
assert task in _BACKGROUND_RELEASE_TASKS
gate.set()
await task
# The done-callback removes it -- the registry doesn't grow unbounded
# across requests.
assert task not in _BACKGROUND_RELEASE_TASKS
assert len(_BACKGROUND_RELEASE_TASKS) == 0
@pytest.mark.asyncio
async def test_success_event_release_is_wired_through_the_background_registry(time_controller):
"""
End-to-end check that async_log_success_event's fire-and-forget release
is genuinely wired through _BACKGROUND_RELEASE_TASKS, not a bare
unreferenced asyncio.create_task -- the registry must be empty again once
the (fast, in-memory) release has had a chance to run, and the release
itself must have actually happened.
"""
limiter = _make_limiter(time_controller)
router = _concurrency_router(limit=1)
limiter.update_variables(llm_router=router)
healthy = router.model_list
request_kwargs, kwargs = _call_context(["end_user_id:u1"])
await limiter.async_filter_deployments(
model="grp", healthy_deployments=healthy, messages=None, request_kwargs=request_kwargs
)
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)
# The registry was actually populated: proves the release ran through
# _BACKGROUND_RELEASE_TASKS, not a bare unreferenced asyncio.create_task
# (which would never touch this set at all, and an "empty at the end"
# check alone can't tell the two apart -- an empty registry throughout
# would satisfy that just as well as one that filled and drained).
assert len(_BACKGROUND_RELEASE_TASKS) == 1
# Two ticks: one for the release task itself to finish (it may already be
# done by the time async_log_success_event returns, given that method's
# own further awaits), and one for its done-callback -- scheduled via
# call_soon when the task completes -- to actually run and discard it.
await asyncio.sleep(0)
await asyncio.sleep(0)
assert len(_BACKGROUND_RELEASE_TASKS) == 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_concurrency_slot_released_on_disconnect_frees_capacity(time_controller):
"""