From 533ae89c324aef46724039bbeecf884bfa07287b Mon Sep 17 00:00:00 2001 From: Deepanshu Date: Thu, 20 Aug 2026 12:39:02 -0400 Subject: [PATCH] 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. --- litellm/proxy/hooks/tag_rate_limiter.py | 49 +++++- .../proxy/hooks/test_tag_rate_limiter.py | 164 ++++++++++++++++++ 2 files changed, 207 insertions(+), 6 deletions(-) diff --git a/litellm/proxy/hooks/tag_rate_limiter.py b/litellm/proxy/hooks/tag_rate_limiter.py index b38f981e819..2e41ba6271a 100644 --- a/litellm/proxy/hooks/tag_rate_limiter.py +++ b/litellm/proxy/hooks/tag_rate_limiter.py @@ -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 diff --git a/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py b/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py index eb8fd892519..1e869ab71f9 100644 --- a/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py +++ b/tests/test_litellm/proxy/hooks/test_tag_rate_limiter.py @@ -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): """