diff --git a/litellm/proxy/hooks/parallel_request_limiter_v3.py b/litellm/proxy/hooks/parallel_request_limiter_v3.py index d10e8aabdd0..9fd40115e10 100644 --- a/litellm/proxy/hooks/parallel_request_limiter_v3.py +++ b/litellm/proxy/hooks/parallel_request_limiter_v3.py @@ -69,7 +69,7 @@ from litellm.router_utils.add_retry_fallback_headers import ( response_has_hidden_params, ) from litellm.router_utils.common_utils import resolve_model_group_alias -from litellm.router_utils.ptu_shares import PTUTeamCeiling, model_group_deployments, team_ptu_ceiling +from litellm.router_utils.ptu_shares import PTUTeamCeiling, team_ptu_ceiling from litellm.types.caching import RedisPipelineIncrementOperation from litellm.types.llms.openai import BaseLiteLLMOpenAIResponseObject, ResponseAPIUsage from litellm.types.utils import ( @@ -134,7 +134,7 @@ def _resolve_ptu_team_ceiling_via_proxy_router(team_id: str, model_group: str) - if llm_router is None or not is_ptu_cost_attribution_enabled(): return None - return team_ptu_ceiling(model_group_deployments(llm_router.get_model_list() or (), model_group), team_id) + return team_ptu_ceiling(llm_router.get_model_list() or (), team_id, model_group) FAIL_CLOSED_RATE_LIMIT_ENFORCEMENT_SETTING: Final = "fail_closed_rate_limit_enforcement" @@ -3301,7 +3301,7 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): descriptors.append( RateLimitDescriptor( key=PTU_TEAM_DESCRIPTOR_KEY, - value=f"{user_api_key_dict.team_id}:{model.group}", + value=f"{user_api_key_dict.team_id}:{ceiling.model_group}", rate_limit=RateLimitDescriptorRateLimitObject( requests_per_unit=None, tokens_per_unit=ceiling.tpm_limit, window_size=self.window_size ), @@ -4930,12 +4930,15 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): team_id: Final = standard_logging_metadata.get("user_api_key_team_id") if reconcile_model is None or not isinstance(team_id, str) or not team_id: return () - scope: Final = (PTU_TEAM_DESCRIPTOR_KEY, f"{team_id}:{reconcile_model.group}") ceiling: Final = ( reserved_ceiling if reserved_ceiling is not None else self._ptu_team_ceiling_resolver(team_id, reconcile_model.group) ) + scope: Final = ( + PTU_TEAM_DESCRIPTOR_KEY, + f"{team_id}:{ceiling.model_group if ceiling is not None else reconcile_model.group}", + ) if ceiling is None and scope not in reserved_scopes: return () return self._build_reservation_aware_tpm_ops( diff --git a/litellm/router_utils/ptu_shares.py b/litellm/router_utils/ptu_shares.py index 11b61676ce5..91ccc585ba7 100644 --- a/litellm/router_utils/ptu_shares.py +++ b/litellm/router_utils/ptu_shares.py @@ -18,6 +18,7 @@ _DeploymentT = TypeVar("_DeploymentT", bound=Mapping[str, object]) @dataclass(frozen=True, slots=True) class PTUTeamCeiling: + model_group: str tpm_limit: int output_to_input_ratio: float cached_input_ratio: float @@ -55,17 +56,24 @@ def filter_ptu_shared_deployments( return PTUShareFilterResult(deployments=kept, withheld=len(kept) < len(checks)) -def team_ptu_ceiling(deployments: Sequence[Mapping[str, object]], team_id: str) -> PTUTeamCeiling | None: - """The per-minute normalized-token ceiling ``team_id``'s shares across ``deployments`` add - up to, else None when the team holds no share on a deployment with a known sizing row. +def team_ptu_ceiling( + deployments: Sequence[Mapping[str, object]], team_id: str, requested_model: str +) -> PTUTeamCeiling | None: + """The per-minute normalized-token ceiling ``team_id``'s shares on the group serving + ``requested_model`` add up to, else None when the team holds no share on a deployment with + a known sizing row. + + A request naming one deployment by its id or provider model counts against that deployment's + group, so every name the router serves it under shares one ceiling. Two shared deployments of different models in one group are weighted by the larger output and cached-input ratios, which over-counts those tokens on the cheaper one rather than under-counting them on the dearer one. """ + model_group: Final = _model_group_of(deployments, requested_model) priced: Final = tuple( (shares[team_id], capacity) - for deployment in deployments + for deployment in model_group_deployments(deployments, model_group) if (shares := _deployment_shares(deployment)) is not None and team_id in shares and (capacity := deployment_ptu_capacity(deployment)) is not None @@ -73,6 +81,7 @@ def team_ptu_ceiling(deployments: Sequence[Mapping[str, object]], team_id: str) if not priced: return None return PTUTeamCeiling( + model_group=model_group, tpm_limit=sum(capacity.input_tpm_for(share) for share, capacity in priced), output_to_input_ratio=max(capacity.output_to_input_ratio for _, capacity in priced), cached_input_ratio=max(capacity.cached_input_ratio for _, capacity in priced), @@ -94,6 +103,36 @@ def model_group_deployments(deployments: Sequence[_DeploymentT], model_group: st ) +def _names_deployment(deployment: Mapping[str, object], name: str) -> bool: + model_info: Final = deployment.get("model_info") + litellm_params: Final = deployment.get("litellm_params") + return (isinstance(model_info, Mapping) and model_info.get("id") == name) or ( + isinstance(litellm_params, Mapping) and litellm_params.get("model") == name + ) + + +def _deployment_model_group(deployment: Mapping[str, object]) -> str | None: + model_info: Final = deployment.get("model_info") + public_name: Final = model_info.get("team_public_model_name") if isinstance(model_info, Mapping) else None + if isinstance(public_name, str): + return public_name + model_name: Final = deployment.get("model_name") + return model_name if isinstance(model_name, str) else None + + +def _model_group_of(deployments: Sequence[Mapping[str, object]], requested_model: str) -> str: + """The group ``requested_model`` routes to: itself when it names a group, else the group of + the deployment it names by id or by provider model, the way the router falls back to them.""" + if model_group_deployments(deployments, requested_model): + return requested_model + named: Final = tuple(deployment for deployment in deployments if _names_deployment(deployment, requested_model)) + shared_first: Final = sorted(named, key=lambda deployment: _deployment_shares(deployment) is None) + return next( + (group for deployment in shared_first if (group := _deployment_model_group(deployment)) is not None), + requested_model, + ) + + def model_group_ptu_capacity(deployments: Sequence[Mapping[str, object]]) -> PTUCapacity | None: """The sizing row of the group's first reserved deployment, single-team or shared, so a team's tokens on that group convert to PTU-hours.""" diff --git a/tests/test_litellm/proxy/hooks/test_parallel_request_limiter_v3.py b/tests/test_litellm/proxy/hooks/test_parallel_request_limiter_v3.py index dfabeceb94c..0b3e54a8358 100644 --- a/tests/test_litellm/proxy/hooks/test_parallel_request_limiter_v3.py +++ b/tests/test_litellm/proxy/hooks/test_parallel_request_limiter_v3.py @@ -7576,7 +7576,9 @@ def _ptu_ceiling_for(team_id: str, model_group: str, tpm_limit: int, ratio: floa calls.append((requested_team, requested_group)) if (requested_team, requested_group) != (team_id, model_group): return None - return PTUTeamCeiling(tpm_limit=tpm_limit, output_to_input_ratio=ratio, cached_input_ratio=cached_ratio) + return PTUTeamCeiling( + model_group=model_group, tpm_limit=tpm_limit, output_to_input_ratio=ratio, cached_input_ratio=cached_ratio + ) return resolve, calls @@ -7760,6 +7762,37 @@ async def test_the_proxy_router_turns_a_teams_share_into_its_ceiling_when_attrib assert "model_per_team_ptu" in str(exc.value.detail) +@pytest.mark.asyncio +@pytest.mark.parametrize("deployment_name", ["shared-ptu", "azure/gpt-4.1"]) +async def test_naming_the_shared_deployment_directly_draws_on_the_same_ceiling_as_its_group( + monkeypatch, deployment_name +): + """The router also serves a deployment named by its id or its provider model, so a team that + spent its share by group name cannot keep going under the deployment's other names.""" + monkeypatch.setenv(PTU_COST_ATTRIBUTION_ENV_VAR, "true") + cache = DualCache() + handler = _PROXY_MaxParallelRequestsHandler(internal_usage_cache=InternalUsageCache(cache)) + key = UserAPIKeyAuth(api_key=hash_token("sk-ptu"), team_id="t") + + with patch("litellm.proxy.proxy_server.llm_router", _shared_ptu_router("test-model")): + await handler.async_pre_call_hook( + user_api_key_dict=key, cache=cache, data=_two_thirds_of_a_ptu_minute(), call_type="acompletion" + ) + with pytest.raises(HTTPException) as exc: + await handler.async_pre_call_hook( + user_api_key_dict=key, + cache=cache, + data={**_two_thirds_of_a_ptu_minute(), "model": deployment_name}, + call_type="acompletion", + ) + + assert exc.value.status_code == 429 + assert "model_per_team_ptu" in str(exc.value.detail) + ptu_keys = [cache_key for cache_key in cache.in_memory_cache.cache_dict if "model_per_team_ptu" in cache_key] + assert handler.create_rate_limit_keys("model_per_team_ptu", "t:test-model", "tokens") in ptu_keys + assert all(cache_key.startswith("{model_per_team_ptu:t:test-model}") for cache_key in ptu_keys) + + @pytest.mark.asyncio async def test_the_proxy_router_sets_no_ceiling_while_attribution_is_off(monkeypatch): monkeypatch.delenv(PTU_COST_ATTRIBUTION_ENV_VAR, raising=False) @@ -7869,7 +7902,9 @@ async def test_a_reservation_is_settled_even_after_the_teams_share_is_gone(): """The share can be removed between admission and completion; the reserved tokens still come off the counter instead of standing in the window.""" ceiling: dict[str, PTUTeamCeiling | None] = { - "current": PTUTeamCeiling(tpm_limit=2000, output_to_input_ratio=4.0, cached_input_ratio=0.0) + "current": PTUTeamCeiling( + model_group="test-model", tpm_limit=2000, output_to_input_ratio=4.0, cached_input_ratio=0.0 + ) } cache = DualCache() handler = _PROXY_MaxParallelRequestsHandler( diff --git a/tests/unit/router_utils/test_ptu_shares.py b/tests/unit/router_utils/test_ptu_shares.py index 2c54c3a1c8d..0012b5628c2 100644 --- a/tests/unit/router_utils/test_ptu_shares.py +++ b/tests/unit/router_utils/test_ptu_shares.py @@ -77,8 +77,9 @@ def test_a_single_team_deployment_and_a_malformed_share_map_are_not_filtered_her def test_a_share_converts_to_the_models_input_tpm_per_ptu(): - ceiling: Final = team_ptu_ceiling([_shared()], "team-a") + ceiling: Final = team_ptu_ceiling([_shared()], "team-a", "gpt-4.1-ptu") assert ceiling == PTUTeamCeiling( + model_group="gpt-4.1-ptu", tpm_limit=30 * _GPT41.input_tpm_per_ptu, output_to_input_ratio=_GPT41.output_to_input_ratio, cached_input_ratio=_GPT41.cached_input_ratio, @@ -87,7 +88,7 @@ def test_a_share_converts_to_the_models_input_tpm_per_ptu(): def test_shares_across_deployments_add_up_and_the_larger_output_ratio_wins(): gpt4o: Final = _shared(model="azure/gpt-4o", shares={"team-a": 10}, deployment_id="shared-4o") - ceiling: Final = team_ptu_ceiling([_shared(), gpt4o, _OPEN], "team-a") + ceiling: Final = team_ptu_ceiling([_shared(), gpt4o, _OPEN], "team-a", "gpt-4.1-ptu") assert ceiling is not None assert ceiling.tpm_limit == 30 * _GPT41.input_tpm_per_ptu + 10 * _GPT4O.input_tpm_per_ptu assert ceiling.output_to_input_ratio == max(_GPT41.output_to_input_ratio, _GPT4O.output_to_input_ratio) @@ -97,7 +98,7 @@ def test_the_larger_cached_input_ratio_wins_across_deployments(): """A team sharing two models is weighted by the one that charges more for cache reads, whichever order the deployments come in.""" gpt6sol: Final = _shared(model="azure/gpt-6-sol", shares={"team-a": 10}, deployment_id="shared-6") - ceiling: Final = team_ptu_ceiling([_shared(), gpt6sol], "team-a") + ceiling: Final = team_ptu_ceiling([_shared(), gpt6sol], "team-a", "gpt-4.1-ptu") assert ceiling is not None assert _GPT41.cached_input_ratio < _GPT6SOL.cached_input_ratio assert ceiling.cached_input_ratio == _GPT6SOL.cached_input_ratio @@ -120,9 +121,45 @@ def test_a_group_is_served_by_name_or_by_a_team_scoped_deployments_public_name() def test_no_share_or_no_sizing_row_sets_no_ceiling(): - assert team_ptu_ceiling([_shared()], "team-c") is None - assert team_ptu_ceiling([_shared(model="azure/unknown-deployment")], "team-a") is None - assert team_ptu_ceiling([_single_team(), _OPEN], "team-a") is None + assert team_ptu_ceiling([_shared()], "team-c", "gpt-4.1-ptu") is None + assert team_ptu_ceiling([_shared(model="azure/unknown-deployment")], "team-a", "gpt-4.1-ptu") is None + assert team_ptu_ceiling([_single_team(), _OPEN], "team-a", "gpt-4.1-ptu") is None + + +def test_naming_a_shared_deployment_by_id_or_provider_model_draws_on_its_groups_ceiling(): + """The router serves a deployment id or a provider model string when no group has that + name, so those names share the group's ceiling instead of bypassing it.""" + payg: Final = {"model_name": "gpt-4.1-payg", "litellm_params": {"model": "azure/gpt-4.1"}, "model_info": {"id": "payg"}} + deployments: Final = [payg, _shared(), _OPEN] + by_group: Final = team_ptu_ceiling(deployments, "team-a", "gpt-4.1-ptu") + assert by_group is not None + assert by_group.model_group == "gpt-4.1-ptu" + assert team_ptu_ceiling(deployments, "team-a", "shared") == by_group + assert team_ptu_ceiling(deployments, "team-a", "azure/gpt-4.1") == by_group + assert team_ptu_ceiling(deployments, "team-a", "payg") is None + assert team_ptu_ceiling(deployments, "team-a", "missing") is None + + +def test_a_group_name_wins_over_a_deployment_id_it_collides_with(): + """The router routes a name that is both a group and a deployment id to the group.""" + colliding: Final = { + "model_name": "shared", + "litellm_params": {"model": "azure/gpt-4o"}, + "model_info": {"id": "colliding"}, + } + assert team_ptu_ceiling([_shared(), colliding], "team-a", "shared") is None + + +def test_a_team_scoped_deployment_named_by_id_draws_on_its_public_groups_ceiling(): + team_scoped: Final = { + "model_name": "gpt-4.1-ptu-3f9c1b", + "litellm_params": {"model": "azure/gpt-4.1"}, + "model_info": {**_shared()["model_info"], "id": "team-scoped", "team_public_model_name": "gpt-4.1-ptu"}, + } + ceiling: Final = team_ptu_ceiling([team_scoped], "team-a", "team-scoped") + assert ceiling is not None + assert ceiling.model_group == "gpt-4.1-ptu" + assert ceiling == team_ptu_ceiling([team_scoped], "team-a", "gpt-4.1-ptu") def test_a_groups_capacity_comes_from_its_first_reserved_deployment_with_a_row():