From 791d38d87bae40a78e351f381c329d59649fa368 Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Tue, 22 Sep 2026 23:24:08 +0800 Subject: [PATCH 1/7] fix(proxy): enforce project ITPM/OTPM on context-management summaries Summary subrequests skipped project model_itpm_limit/model_otpm_limit on both the read-only gate and post-call charge path. Include those descriptors with estimate headroom checks, and charge unreserved summary call ids from logging metadata like combined TPM. Signed-off-by: sinksilk <785976238@qq.com> --- .../context_management/editors/compact.py | 53 +++++- .../hooks/parallel_request_limiter_v3.py | 90 ++++++++- tests/unit/proxy/hooks/test_tpm_concurrent.py | 173 ++++++++++++++++++ 3 files changed, 310 insertions(+), 6 deletions(-) diff --git a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py index 62826865894..59d85528d9e 100644 --- a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py +++ b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py @@ -521,6 +521,9 @@ def _without_parallel_request_gauge(descriptor: "RateLimitDescriptor") -> "RateL async def _check_summary_model_rate_limit( user_api_key_auth: Optional["UserAPIKeyAuth"], summary_model: str, + *, + estimated_input_tokens: int = 1, + estimated_output_tokens: int = 1, ) -> bool: """Return True when the caller is within their configured RPM/TPM limits for ``summary_model``. @@ -530,14 +533,22 @@ async def _check_summary_model_rate_limit( user RPM or TPM could still drive an extra summary-model completion per allowed ``/v1/messages`` request. This mirrors the read side of ``_PROXY_MaxParallelRequestsHandler_v3.async_pre_call_hook`` for the - summary model: it builds the same descriptor set and runs the check in - ``read_only`` mode so no counter is reserved or incremented — the summary - call's actual usage is still charged exactly once by the limiter's - post-call success hook (via the propagated ``litellm_metadata``). + summary model: it builds the same descriptor set (including project + ITPM/OTPM) and runs the check in ``read_only`` mode so no counter is + reserved or incremented — the summary call's actual usage is still + charged exactly once by the limiter's post-call success hook (via the + propagated ``litellm_metadata``). ``max_parallel_requests`` gauges are left out of the check: the summary call runs inside the caller's already admitted request, whose own slot would otherwise count against it. + Project ITPM/OTPM are reservation-style quotas (pre-call reserves an + estimate). A read-only ``OVER_LIMIT`` only fires once the counter is + already at the cap, so this gate also compares ``limit_remaining`` to + ``estimated_input_tokens`` / ``estimated_output_tokens`` for those + descriptors — matching how an ordinary request would be refused when the + next reservation cannot fit. + Returns True (allow) outside the proxy, when the active limiter does not expose the read-only descriptor check (legacy limiter), or when the descriptor set cannot be built — the deny signals are a definitive @@ -566,6 +577,9 @@ async def _check_summary_model_rate_limit( add_project_descriptor: Final[_AddModelRateLimitDescriptor | None] = getattr( limiter, "_add_project_model_rate_limit_descriptor_from_metadata", None ) + add_project_io_descriptor: Final[_AddModelRateLimitDescriptor | None] = getattr( + limiter, "add_project_io_token_rate_limit_descriptors_from_metadata", None + ) create_org_descriptors: Final[_CreateOrgRateLimitDescriptors | None] = getattr( limiter, "create_organization_rate_limit_descriptor", None ) @@ -599,6 +613,15 @@ async def _check_summary_model_rate_limit( requested_model=summary_model, descriptors=base_descriptors, ) + # Project ITPM/OTPM are reserved (not merely read) on the main pre-call + # path, so the summary gate must add those descriptors explicitly — + # otherwise an exhausted project IO quota still allows compaction. + if add_project_io_descriptor is not None: + add_project_io_descriptor( + user_api_key_dict=user_api_key_auth, + requested_model=summary_model, + descriptors=base_descriptors, + ) descriptors: Final = _without_parallel_request_gauges( (*base_descriptors, *create_org_descriptors(user_api_key_auth, summary_model)) ) @@ -624,7 +647,25 @@ async def _check_summary_model_rate_limit( e, ) return True - return response.get("overall_code") != "OVER_LIMIT" + if response.get("overall_code") == "OVER_LIMIT": + return False + + # Reservation-style project IO quotas: deny when the estimated summary + # cannot fit in remaining headroom (ordinary traffic fails the same way). + input_estimate: Final = max(1, estimated_input_tokens) + output_estimate: Final = max(1, estimated_output_tokens) + for status in response.get("statuses") or (): + if not isinstance(status, Mapping): + continue + descriptor_key = status.get("descriptor_key") + remaining = status.get("limit_remaining") + if not isinstance(remaining, int): + continue + if descriptor_key == "model_per_project_itpm" and remaining < input_estimate: + return False + if descriptor_key == "model_per_project_otpm" and remaining < output_estimate: + return False + return True def _find_latest_compaction_index( @@ -1293,6 +1334,8 @@ async def apply_compact_20260112( if not await _check_summary_model_rate_limit( user_api_key_auth=user_api_key_auth, summary_model=summary_model, + estimated_input_tokens=current_tokens, + estimated_output_tokens=_read_summary_max_tokens_setting(), ): verbose_logger.warning( "compact_20260112: caller over rate limit for summary_model=%s; skipping summary call", diff --git a/litellm/proxy/hooks/parallel_request_limiter_v3.py b/litellm/proxy/hooks/parallel_request_limiter_v3.py index b33cea5742d..2f3b10b07f6 100644 --- a/litellm/proxy/hooks/parallel_request_limiter_v3.py +++ b/litellm/proxy/hooks/parallel_request_limiter_v3.py @@ -4637,6 +4637,90 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): return 0, 0, False + def _collect_project_io_scope_targets( + self, + standard_logging_metadata: Mapping[str, Any], + model_group: str | None, + ) -> list[tuple[str, str]]: + """Rebuild project ITPM/OTPM scopes from logging metadata. + + Combined TPM already charges ``model_per_project`` from metadata when + no reservation owns the call. Summary/compaction subrequests carry a + distinct ``litellm_call_id`` and never own the parent stash, so their + IO quotas must use the same metadata rebuild. + """ + user_api_key_project_id: Final = standard_logging_metadata.get("user_api_key_project_id") + if not user_api_key_project_id or not model_group: + return [] + descriptor_value: Final = f"{user_api_key_project_id}:{model_group}" + return [ # mutable-ok: caller may filter ITPM vs OTPM scopes + (PROJECT_ITPM_DESCRIPTOR_KEY, descriptor_value), + (PROJECT_OTPM_DESCRIPTOR_KEY, descriptor_value), + ] + + def _build_unreserved_project_io_token_ops( + self, + kwargs: dict[str, Any], + response_obj: object, + ) -> Sequence[RedisPipelineIncrementOperation]: + """Charge full actual ITPM/OTPM when no pre-call reservation owns this call. + + Summary subrequests never claim the parent stash (``owner_litellm_call_id`` + pins it), so without this path their input/output tokens never hit the + project IO counters even though combined TPM still charges them. + """ + from litellm.proxy.common_utils.callback_utils import ( + get_model_group_from_litellm_kwargs, + ) + + standard_logging_object: Final = kwargs.get("standard_logging_object") or {} + if not isinstance(standard_logging_object, dict): + return () + standard_logging_metadata: Final = standard_logging_object.get("metadata") or {} + if not isinstance(standard_logging_metadata, Mapping): + return () + + model_group: Final = get_model_group_from_litellm_kwargs(kwargs) or ( + standard_logging_object.get("model_group") + if isinstance(standard_logging_object.get("model_group"), str) + else None + ) + targets: Final = self._collect_project_io_scope_targets( + standard_logging_metadata=standard_logging_metadata, + model_group=model_group if isinstance(model_group, str) else None, + ) + if not targets: + return () + + response_usage: Final = self._resolve_io_token_reconcile_usage(response_obj) + combined_usage: Final = self._resolve_io_token_reconcile_usage(kwargs.get("combined_usage_object")) + aggregate_total: Final = self._aggregate_only_total_tokens( + self._response_usage(response_obj) + ) or self._aggregate_only_total_tokens(self._response_usage(kwargs.get("combined_usage_object"))) + if not response_usage[2] and not combined_usage[2] and aggregate_total <= 0: + return () + resolved_usage: Final = ( + response_usage + if response_usage[2] + else combined_usage + if combined_usage[2] + else (aggregate_total, aggregate_total, True) + ) + billable_input, completion_tokens, _ = resolved_usage + itpm_targets: Final = [t for t in targets if t[0] == PROJECT_ITPM_DESCRIPTOR_KEY] + otpm_targets: Final = [t for t in targets if t[0] == PROJECT_OTPM_DESCRIPTOR_KEY] + return self._build_reservation_aware_tpm_ops( + targets=itpm_targets, + reserved_scopes=frozenset(), + actual_tokens=billable_input, + reserved_tokens=0, + ) + self._build_reservation_aware_tpm_ops( + targets=otpm_targets, + reserved_scopes=frozenset(), + actual_tokens=completion_tokens, + reserved_tokens=0, + ) + def _build_io_token_reservation_ops( self, kwargs: object, @@ -4649,12 +4733,16 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): are stored in the same ":tokens" cache bucket as combined TPM, just under distinct scope keys, so the reservation-aware increment math is identical; only the usage fields being reconciled against differ. + + When this call id does not own the request stash (summary/compaction + subrequest), falls through to the unreserved metadata rebuild so + project IO quotas still receive the summary's actual usage. """ if not isinstance(kwargs, dict): return () stash: Final = get_request_stash_for_call(_call_id_from_callback_kwargs(kwargs)) if stash is None: - return () + return self._build_unreserved_project_io_token_ops(kwargs, response_obj) itpm_reserved: Final = stash.itpm_reserved_tokens otpm_reserved: Final = stash.otpm_reserved_tokens diff --git a/tests/unit/proxy/hooks/test_tpm_concurrent.py b/tests/unit/proxy/hooks/test_tpm_concurrent.py index 42c1f489bdd..91aad818f52 100644 --- a/tests/unit/proxy/hooks/test_tpm_concurrent.py +++ b/tests/unit/proxy/hooks/test_tpm_concurrent.py @@ -3708,5 +3708,178 @@ async def test_the_project_itpm_reservation_counts_the_request_off_the_event_loo assert_loop_stayed_free(took, lags) +@pytest.mark.asyncio +async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): + """Regression for #41395: context-management summary subrequests must be + gated on project ITPM/OTPM and charged against those buckets afterwards. + + The summary call carries its own litellm_call_id, so it never owns the + parent stash. Without the unreserved IO charge path, combined TPM still + increments while model_per_project_itpm/otpm stay untouched. + + Project IO quotas are reservation-style: after ordinary traffic is refused + there may still be residual headroom, so the summary gate passes the + estimated summary size (as apply_compact does) to compare against + limit_remaining. + """ + from litellm.llms.anthropic.experimental_pass_through.context_management.editors.compact import ( + _check_summary_model_rate_limit, + ) + from litellm.proxy import proxy_server + + handler, _cache = rate_limiter + proxy_server.proxy_logging_obj.max_parallel_request_limiter = handler + + model = "gpt-4o-mini" + project = "proj-summary-io" + itpm_limit = 2000 + otpm_limit = 10**6 + + def make_auth(**extra_project_metadata) -> UserAPIKeyAuth: + return UserAPIKeyAuth( + api_key="sk-proj-key", + project_id=project, + project_metadata={ + "model_itpm_limit": {model: itpm_limit}, + "model_otpm_limit": {model: otpm_limit}, + **extra_project_metadata, + }, + ) + + def request_data() -> dict: + return { + "model": model, + "messages": [{"role": "user", "content": "x " * 300}], + "litellm_call_id": "parent-call-id", + } + + async def drive_until_refused(make_auth_fn) -> tuple[int, str | None]: + allowed = 0 + for _ in range(30): + try: + await handler.async_pre_call_hook( + user_api_key_dict=make_auth_fn(), + cache=DualCache(), + data=request_data(), + call_type="completion", + ) + allowed += 1 + except Exception as e: + return allowed, str(e) + return allowed, None + + allowed, refusal = await drive_until_refused(make_auth) + assert allowed >= 1 + assert refusal is not None + assert "model_per_project_itpm" in refusal + + # Same residual headroom that refused the next ordinary reservation must + # refuse a summary whose estimated input cannot fit. + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=make_auth(), + summary_model=model, + estimated_input_tokens=200, + estimated_output_tokens=1, + ) + is False + ) + + # CONTROL: RPM exhaustion still denies via the same gate. + rpm_handler = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) + proxy_server.proxy_logging_obj.max_parallel_request_limiter = rpm_handler + allowed_rpm, refusal_rpm = 0, None + for _ in range(30): + try: + await rpm_handler.async_pre_call_hook( + user_api_key_dict=make_auth(model_rpm_limit={model: 4}), + cache=DualCache(), + data=request_data(), + call_type="completion", + ) + allowed_rpm += 1 + except Exception as e: + refusal_rpm = str(e) + break + assert allowed_rpm == 4 + assert refusal_rpm is not None + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=make_auth(model_rpm_limit={model: 4}), + summary_model=model, + ) + is False + ) + + # Post-call: summary call id must charge ITPM/OTPM like combined TPM. + charging = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) + proxy_server.proxy_logging_obj.max_parallel_request_limiter = charging + + async def in_one_request_context(): + await charging.async_pre_call_hook( + user_api_key_dict=make_auth(), + cache=DualCache(), + data=request_data(), + call_type="completion", + ) + response = ModelResponse( + usage=Usage(prompt_tokens=5000, completion_tokens=5000, total_tokens=10000) + ) + metadata = { + "user_api_key_project_id": project, + "user_api_key_hash": "sk-proj-key", + "model_group": model, + } + + def kwargs_for(call_id: str) -> dict: + return { + "litellm_call_id": call_id, + "model": model, + "litellm_params": {"metadata": metadata}, + "standard_logging_object": {"metadata": metadata, "model_group": model}, + } + + parent_ops = list( + charging._build_io_token_reservation_ops( + kwargs=kwargs_for("parent-call-id"), + response_obj=response, + ) + ) + summary_ops = list( + charging._build_io_token_reservation_ops( + kwargs=kwargs_for("summary-call-id"), + response_obj=response, + ) + ) + summary_tpm = charging._build_success_event_pipeline_operations( + kwargs=kwargs_for("summary-call-id"), + response_obj=response, + rate_limit_type=charging.get_rate_limit_type(), + ) + return parent_ops, summary_ops, summary_tpm + + parent_ops, summary_ops, summary_tpm = await asyncio.create_task( + in_one_request_context() + ) + assert parent_ops, "parent call should reconcile reserved ITPM/OTPM" + assert summary_ops, "summary call must charge project ITPM/OTPM without owning the stash" + summary_keys = {op["key"] for op in summary_ops} + assert any("model_per_project_itpm" in key for key in summary_keys) + assert any("model_per_project_otpm" in key for key in summary_keys) + assert any( + "model_per_project:" in op["key"] and op["increment_value"] == 10000 + for op in summary_tpm + ) + # Unreserved summary path charges full actual usage (no reservation delta). + assert any( + "model_per_project_itpm" in op["key"] and op["increment_value"] == 5000 + for op in summary_ops + ) + assert any( + "model_per_project_otpm" in op["key"] and op["increment_value"] == 5000 + for op in summary_ops + ) + + if __name__ == "__main__": pytest.main([__file__, "-v", "-s"]) From a77f23396137f6718eeffe7f9d8a56ef41bf103e Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Tue, 22 Sep 2026 23:32:20 +0800 Subject: [PATCH 2/7] fix(proxy): clear LIT002 in unreserved project IO token helpers Use immutable tuples and drop or-{} fallbacks so type_discipline_gate no longer blames the summary ITPM/OTPM path. Signed-off-by: sinksilk <785976238@qq.com> --- .../proxy/hooks/parallel_request_limiter_v3.py | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/litellm/proxy/hooks/parallel_request_limiter_v3.py b/litellm/proxy/hooks/parallel_request_limiter_v3.py index 2f3b10b07f6..1e08403083a 100644 --- a/litellm/proxy/hooks/parallel_request_limiter_v3.py +++ b/litellm/proxy/hooks/parallel_request_limiter_v3.py @@ -4641,7 +4641,7 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): self, standard_logging_metadata: Mapping[str, Any], model_group: str | None, - ) -> list[tuple[str, str]]: + ) -> Sequence[tuple[str, str]]: """Rebuild project ITPM/OTPM scopes from logging metadata. Combined TPM already charges ``model_per_project`` from metadata when @@ -4651,12 +4651,12 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): """ user_api_key_project_id: Final = standard_logging_metadata.get("user_api_key_project_id") if not user_api_key_project_id or not model_group: - return [] + return () descriptor_value: Final = f"{user_api_key_project_id}:{model_group}" - return [ # mutable-ok: caller may filter ITPM vs OTPM scopes + return ( (PROJECT_ITPM_DESCRIPTOR_KEY, descriptor_value), (PROJECT_OTPM_DESCRIPTOR_KEY, descriptor_value), - ] + ) def _build_unreserved_project_io_token_ops( self, @@ -4673,10 +4673,10 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): get_model_group_from_litellm_kwargs, ) - standard_logging_object: Final = kwargs.get("standard_logging_object") or {} + standard_logging_object: Final = kwargs.get("standard_logging_object") if not isinstance(standard_logging_object, dict): return () - standard_logging_metadata: Final = standard_logging_object.get("metadata") or {} + standard_logging_metadata: Final = standard_logging_object.get("metadata") if not isinstance(standard_logging_metadata, Mapping): return () @@ -4707,8 +4707,8 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): else (aggregate_total, aggregate_total, True) ) billable_input, completion_tokens, _ = resolved_usage - itpm_targets: Final = [t for t in targets if t[0] == PROJECT_ITPM_DESCRIPTOR_KEY] - otpm_targets: Final = [t for t in targets if t[0] == PROJECT_OTPM_DESCRIPTOR_KEY] + itpm_targets: Final = tuple(t for t in targets if t[0] == PROJECT_ITPM_DESCRIPTOR_KEY) + otpm_targets: Final = tuple(t for t in targets if t[0] == PROJECT_OTPM_DESCRIPTOR_KEY) return self._build_reservation_aware_tpm_ops( targets=itpm_targets, reserved_scopes=frozenset(), From 723ca75b8aac0f2d2f9978f21195186d2ffbc505 Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Fri, 25 Sep 2026 01:09:56 +0800 Subject: [PATCH 3/7] fix(proxy): seed windows for unreserved summary ITPM/OTPM charges Unreserved summary charges now open the TPM window with the increment so a later reservation cannot wipe them. Admission estimates use the built summary messages and summary_model Signed-off-by: sinksilk <785976238@qq.com> --- .../context_management/editors/compact.py | 31 +- .../hooks/parallel_request_limiter_v3.py | 101 +++++-- tests/unit/proxy/hooks/test_tpm_concurrent.py | 265 ++++++++---------- 3 files changed, 234 insertions(+), 163 deletions(-) diff --git a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py index 59d85528d9e..4be645b8f72 100644 --- a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py +++ b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py @@ -1037,6 +1037,25 @@ def _build_summary_messages( return summary_messages +async def _estimate_summary_input_tokens( + *, + summary_model: str, + summary_messages: Sequence[Mapping[str, object]], + fallback_tokens: int, +) -> int: + try: + return await asyncify(litellm.token_counter)( + model=summary_model, + messages=list(summary_messages), + ) + except Exception as e: + verbose_logger.warning( + "compact_20260112: summary token estimate failed; falling back to parent current_tokens: %s", + e, + ) + return fallback_tokens + + def _is_user_message(msg: object) -> bool: return isinstance(msg, dict) and msg.get("role") == "user" @@ -1331,10 +1350,18 @@ async def apply_compact_20260112( applied_edits=[applied], ) + prompt: Final = _build_summary_prompt(edit_spec, tools) + summary_messages: Final = _build_summary_messages(effective_messages, prompt, system=augmented_system) + estimated_summary_input: Final = await _estimate_summary_input_tokens( + summary_model=summary_model, + summary_messages=summary_messages, + fallback_tokens=current_tokens, + ) + if not await _check_summary_model_rate_limit( user_api_key_auth=user_api_key_auth, summary_model=summary_model, - estimated_input_tokens=current_tokens, + estimated_input_tokens=estimated_summary_input, estimated_output_tokens=_read_summary_max_tokens_setting(), ): verbose_logger.warning( @@ -1348,8 +1375,6 @@ async def apply_compact_20260112( applied_edits=[applied], ) - prompt: Final = _build_summary_prompt(edit_spec, tools) - summary_messages: Final = _build_summary_messages(effective_messages, prompt, system=augmented_system) propagated_metadata: Final = _propagate_metadata(litellm_metadata) allowed_model_region: Final = getattr(user_api_key_auth, "allowed_model_region", None) diff --git a/litellm/proxy/hooks/parallel_request_limiter_v3.py b/litellm/proxy/hooks/parallel_request_limiter_v3.py index 1e08403083a..f8ad7fd2f7e 100644 --- a/litellm/proxy/hooks/parallel_request_limiter_v3.py +++ b/litellm/proxy/hooks/parallel_request_limiter_v3.py @@ -554,6 +554,7 @@ class ReservationAwareIncrementOperation(RedisPipelineIncrementOperation): window_key: NotRequired[str] expected_window_start: NotRequired[str] reservation_backend: NotRequired[Literal["redis", "local"]] + seed_window_if_absent: NotRequired[bool] class RateLimitResponseWithDescriptors(TypedDict): @@ -4500,18 +4501,16 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): parent_otel_span: Span | None = None, ) -> None: for operation in pipeline_operations: - if operation.get("window_key") is None or operation.get("expected_window_start") is None: - await self.internal_usage_cache.async_increment_cache( - key=operation["key"], - value=operation["increment_value"], - litellm_parent_otel_span=parent_otel_span, - ttl=operation["ttl"], - ) + await self._apply_one_reservation_aware_token_increment( + operation=operation, + parent_otel_span=parent_otel_span, + ) local_guarded_operations: Final = tuple( operation for operation in pipeline_operations if operation.get("window_key") is not None and operation.get("expected_window_start") is not None + and not operation.get("seed_window_if_absent") and operation.get("reservation_backend") == "local" ) redis_guarded_operations: Final = tuple( @@ -4519,6 +4518,7 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): for operation in pipeline_operations if operation.get("window_key") is not None and operation.get("expected_window_start") is not None + and not operation.get("seed_window_if_absent") and operation.get("reservation_backend") != "local" ) if local_guarded_operations: @@ -4532,6 +4532,57 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): parent_otel_span=parent_otel_span, ) + async def _apply_one_reservation_aware_token_increment( + self, + *, + operation: ReservationAwareIncrementOperation, + parent_otel_span: Span | None, + ) -> None: + if operation.get("seed_window_if_absent"): + window_key: Final = operation.get("window_key") + if window_key is not None: + await self._seed_rate_limit_window_if_absent( + window_key=window_key, + ttl=operation["ttl"], + parent_otel_span=parent_otel_span, + ) + await self.internal_usage_cache.async_increment_cache( + key=operation["key"], + value=operation["increment_value"], + litellm_parent_otel_span=parent_otel_span, + ttl=operation["ttl"], + ) + return + if operation.get("window_key") is None or operation.get("expected_window_start") is None: + await self.internal_usage_cache.async_increment_cache( + key=operation["key"], + value=operation["increment_value"], + litellm_parent_otel_span=parent_otel_span, + ttl=operation["ttl"], + ) + + async def _seed_rate_limit_window_if_absent( + self, + *, + window_key: str, + ttl: int | None, + parent_otel_span: Span | None = None, + ) -> None: + """Open a TPM window around an unreserved charge so a later reservation does not wipe it.""" + active_window: Final = await self.internal_usage_cache.async_get_cache( + key=window_key, + litellm_parent_otel_span=parent_otel_span, + ) + if active_window is not None: + return + window_ttl: Final = ttl if ttl is not None else self.window_size + await self.internal_usage_cache.async_set_cache( + key=window_key, + value=str(int(self._get_current_time().timestamp())), + ttl=window_ttl, + litellm_parent_otel_span=parent_otel_span, + ) + def get_rate_limit_type(self) -> Literal["output", "input", "total"]: from litellm.proxy.proxy_server import general_settings @@ -4639,7 +4690,7 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): def _collect_project_io_scope_targets( self, - standard_logging_metadata: Mapping[str, Any], + standard_logging_metadata: Mapping[str, object], model_group: str | None, ) -> Sequence[tuple[str, str]]: """Rebuild project ITPM/OTPM scopes from logging metadata. @@ -4658,16 +4709,38 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): (PROJECT_OTPM_DESCRIPTOR_KEY, descriptor_value), ) + def _build_unreserved_scoped_token_ops( + self, + targets: Sequence[tuple[str, str]], + actual_tokens: int, + ) -> tuple[ReservationAwareIncrementOperation, ...]: + if actual_tokens == 0: + return () + return tuple( + ReservationAwareIncrementOperation( + key=self.create_rate_limit_keys(scope_key, scope_value, "tokens"), + increment_value=actual_tokens, + ttl=self.window_size, + window_key=f"{{{scope_key}:{scope_value}}}:window", + seed_window_if_absent=True, + ) + for scope_key, scope_value in targets + ) + def _build_unreserved_project_io_token_ops( self, - kwargs: dict[str, Any], + kwargs: Mapping[str, object], response_obj: object, - ) -> Sequence[RedisPipelineIncrementOperation]: + ) -> tuple[ReservationAwareIncrementOperation, ...]: """Charge full actual ITPM/OTPM when no pre-call reservation owns this call. Summary subrequests never claim the parent stash (``owner_litellm_call_id`` pins it), so without this path their input/output tokens never hit the project IO counters even though combined TPM still charges them. + + Unreserved charges also seed the TPM window key. A plain counter increment + without a window is wiped when the next ordinary reservation treats a + missing window as expired and resets sibling counters. """ from litellm.proxy.common_utils.callback_utils import ( get_model_group_from_litellm_kwargs, @@ -4709,16 +4782,12 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): billable_input, completion_tokens, _ = resolved_usage itpm_targets: Final = tuple(t for t in targets if t[0] == PROJECT_ITPM_DESCRIPTOR_KEY) otpm_targets: Final = tuple(t for t in targets if t[0] == PROJECT_OTPM_DESCRIPTOR_KEY) - return self._build_reservation_aware_tpm_ops( + return self._build_unreserved_scoped_token_ops( targets=itpm_targets, - reserved_scopes=frozenset(), actual_tokens=billable_input, - reserved_tokens=0, - ) + self._build_reservation_aware_tpm_ops( + ) + self._build_unreserved_scoped_token_ops( targets=otpm_targets, - reserved_scopes=frozenset(), actual_tokens=completion_tokens, - reserved_tokens=0, ) def _build_io_token_reservation_ops( diff --git a/tests/unit/proxy/hooks/test_tpm_concurrent.py b/tests/unit/proxy/hooks/test_tpm_concurrent.py index 91aad818f52..83f00211895 100644 --- a/tests/unit/proxy/hooks/test_tpm_concurrent.py +++ b/tests/unit/proxy/hooks/test_tpm_concurrent.py @@ -3710,175 +3710,152 @@ async def test_the_project_itpm_reservation_counts_the_request_off_the_event_loo @pytest.mark.asyncio async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): - """Regression for #41395: context-management summary subrequests must be - gated on project ITPM/OTPM and charged against those buckets afterwards. + """Regression for #41395: summary subrequests gate and charge project ITPM/OTPM.""" + from typing import Final - The summary call carries its own litellm_call_id, so it never owns the - parent stash. Without the unreserved IO charge path, combined TPM still - increments while model_per_project_itpm/otpm stay untouched. - - Project IO quotas are reservation-style: after ordinary traffic is refused - there may still be residual headroom, so the summary gate passes the - estimated summary size (as apply_compact does) to compare against - limit_remaining. - """ from litellm.llms.anthropic.experimental_pass_through.context_management.editors.compact import ( _check_summary_model_rate_limit, ) from litellm.proxy import proxy_server handler, _cache = rate_limiter + previous_limiter: Final = getattr(proxy_server.proxy_logging_obj, "max_parallel_request_limiter", None) proxy_server.proxy_logging_obj.max_parallel_request_limiter = handler + try: + model: Final = "gpt-4o-mini" + project: Final = "proj-summary-io" - model = "gpt-4o-mini" - project = "proj-summary-io" - itpm_limit = 2000 - otpm_limit = 10**6 - - def make_auth(**extra_project_metadata) -> UserAPIKeyAuth: - return UserAPIKeyAuth( - api_key="sk-proj-key", - project_id=project, - project_metadata={ - "model_itpm_limit": {model: itpm_limit}, - "model_otpm_limit": {model: otpm_limit}, - **extra_project_metadata, - }, - ) - - def request_data() -> dict: - return { - "model": model, - "messages": [{"role": "user", "content": "x " * 300}], - "litellm_call_id": "parent-call-id", - } - - async def drive_until_refused(make_auth_fn) -> tuple[int, str | None]: - allowed = 0 - for _ in range(30): - try: - await handler.async_pre_call_hook( - user_api_key_dict=make_auth_fn(), - cache=DualCache(), - data=request_data(), - call_type="completion", - ) - allowed += 1 - except Exception as e: - return allowed, str(e) - return allowed, None - - allowed, refusal = await drive_until_refused(make_auth) - assert allowed >= 1 - assert refusal is not None - assert "model_per_project_itpm" in refusal - - # Same residual headroom that refused the next ordinary reservation must - # refuse a summary whose estimated input cannot fit. - assert ( - await _check_summary_model_rate_limit( - user_api_key_auth=make_auth(), - summary_model=model, - estimated_input_tokens=200, - estimated_output_tokens=1, - ) - is False - ) - - # CONTROL: RPM exhaustion still denies via the same gate. - rpm_handler = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) - proxy_server.proxy_logging_obj.max_parallel_request_limiter = rpm_handler - allowed_rpm, refusal_rpm = 0, None - for _ in range(30): - try: - await rpm_handler.async_pre_call_hook( - user_api_key_dict=make_auth(model_rpm_limit={model: 4}), - cache=DualCache(), - data=request_data(), - call_type="completion", + def make_auth(**extra_project_metadata) -> UserAPIKeyAuth: + return UserAPIKeyAuth( + api_key="sk-proj-key", + project_id=project, + project_metadata={ + "model_itpm_limit": {model: 2000}, + "model_otpm_limit": {model: 10**6}, + **extra_project_metadata, + }, ) - allowed_rpm += 1 - except Exception as e: - refusal_rpm = str(e) - break - assert allowed_rpm == 4 - assert refusal_rpm is not None - assert ( - await _check_summary_model_rate_limit( - user_api_key_auth=make_auth(model_rpm_limit={model: 4}), - summary_model=model, - ) - is False - ) - # Post-call: summary call id must charge ITPM/OTPM like combined TPM. - charging = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) - proxy_server.proxy_logging_obj.max_parallel_request_limiter = charging + def request_data() -> dict[str, object]: + return { + "model": model, + "messages": [{"role": "user", "content": "x " * 300}], + "litellm_call_id": "parent-call-id", + } - async def in_one_request_context(): - await charging.async_pre_call_hook( - user_api_key_dict=make_auth(), - cache=DualCache(), - data=request_data(), - call_type="completion", + async def drive_until_refused( + limiter: RateLimitHandler, auth: UserAPIKeyAuth + ) -> tuple[int, str | None]: + successes: Final[list[bool]] = [] + for _ in range(30): + try: + await limiter.async_pre_call_hook( + user_api_key_dict=auth, + cache=DualCache(), + data=request_data(), + call_type="completion", + ) + successes.append(True) + except Exception as e: + return len(successes), str(e) + return len(successes), None + + allowed, refusal = await drive_until_refused(handler, make_auth()) + assert allowed >= 1 + assert refusal is not None + assert "model_per_project_itpm" in refusal + + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=make_auth(), + summary_model=model, + estimated_input_tokens=200, + estimated_output_tokens=1, + ) + is False ) - response = ModelResponse( - usage=Usage(prompt_tokens=5000, completion_tokens=5000, total_tokens=10000) + + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=make_auth( + model_itpm_limit={model: 10**6}, + model_otpm_limit={model: 5}, + ), + summary_model=model, + estimated_input_tokens=1, + estimated_output_tokens=20, + ) + is False ) - metadata = { + + rpm_handler: Final = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) + proxy_server.proxy_logging_obj.max_parallel_request_limiter = rpm_handler + allowed_rpm, refusal_rpm = await drive_until_refused( + rpm_handler, make_auth(model_rpm_limit={model: 4}) + ) + assert allowed_rpm == 4 + assert refusal_rpm is not None + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=make_auth(model_rpm_limit={model: 4}), + summary_model=model, + ) + is False + ) + + charging: Final = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) + proxy_server.proxy_logging_obj.max_parallel_request_limiter = charging + summary_response: Final = ModelResponse( + usage=Usage(prompt_tokens=60, completion_tokens=40, total_tokens=100) + ) + metadata: Final = { "user_api_key_project_id": project, "user_api_key_hash": "sk-proj-key", "model_group": model, } - - def kwargs_for(call_id: str) -> dict: - return { - "litellm_call_id": call_id, + await charging.async_log_success_event( + kwargs={ + "litellm_call_id": "summary-call-id", "model": model, "litellm_params": {"metadata": metadata}, "standard_logging_object": {"metadata": metadata, "model_group": model}, - } + }, + response_obj=summary_response, + start_time=datetime.now(), + end_time=datetime.now(), + ) - parent_ops = list( - charging._build_io_token_reservation_ops( - kwargs=kwargs_for("parent-call-id"), - response_obj=response, - ) + itpm_key: Final = charging.create_rate_limit_keys( + PROJECT_ITPM_DESCRIPTOR_KEY, f"{project}:{model}", "tokens" ) - summary_ops = list( - charging._build_io_token_reservation_ops( - kwargs=kwargs_for("summary-call-id"), - response_obj=response, - ) + otpm_key: Final = charging.create_rate_limit_keys( + PROJECT_OTPM_DESCRIPTOR_KEY, f"{project}:{model}", "tokens" ) - summary_tpm = charging._build_success_event_pipeline_operations( - kwargs=kwargs_for("summary-call-id"), - response_obj=response, - rate_limit_type=charging.get_rate_limit_type(), - ) - return parent_ops, summary_ops, summary_tpm + itpm_window: Final = f"{{{PROJECT_ITPM_DESCRIPTOR_KEY}:{project}:{model}}}:window" + otpm_window: Final = f"{{{PROJECT_OTPM_DESCRIPTOR_KEY}:{project}:{model}}}:window" + dual: Final = charging.internal_usage_cache.dual_cache + assert int(await dual.async_get_cache(key=itpm_key) or 0) == 60 + assert int(await dual.async_get_cache(key=otpm_key) or 0) == 40 + assert await dual.async_get_cache(key=itpm_window) is not None + assert await dual.async_get_cache(key=otpm_window) is not None - parent_ops, summary_ops, summary_tpm = await asyncio.create_task( - in_one_request_context() - ) - assert parent_ops, "parent call should reconcile reserved ITPM/OTPM" - assert summary_ops, "summary call must charge project ITPM/OTPM without owning the stash" - summary_keys = {op["key"] for op in summary_ops} - assert any("model_per_project_itpm" in key for key in summary_keys) - assert any("model_per_project_otpm" in key for key in summary_keys) - assert any( - "model_per_project:" in op["key"] and op["increment_value"] == 10000 - for op in summary_tpm - ) - # Unreserved summary path charges full actual usage (no reservation delta). - assert any( - "model_per_project_itpm" in op["key"] and op["increment_value"] == 5000 - for op in summary_ops - ) - assert any( - "model_per_project_otpm" in op["key"] and op["increment_value"] == 5000 - for op in summary_ops - ) + await charging.async_pre_call_hook( + user_api_key_dict=make_auth( + model_itpm_limit={model: 10**6}, + model_otpm_limit={model: 10**6}, + ), + cache=DualCache(), + data={ + "model": model, + "messages": [{"role": "user", "content": "hi"}], + "max_tokens": 10, + "litellm_call_id": "follow-up-call", + }, + call_type="completion", + ) + assert int(await dual.async_get_cache(key=itpm_key) or 0) >= 60 + finally: + proxy_server.proxy_logging_obj.max_parallel_request_limiter = previous_limiter if __name__ == "__main__": From daf87b4064405c30acbdcb0c6dc1996b795be4f9 Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Tue, 29 Sep 2026 23:24:32 +0800 Subject: [PATCH 4/7] fix(proxy): retarget summary ITPM regression after pass_through rename Import compact from anthropic.pass_through and install the limiter via proxy_hook_mapping so get_proxy_hook resolves it in the summary gate. Signed-off-by: sinksilk <785976238@qq.com> --- tests/unit/proxy/hooks/test_tpm_concurrent.py | 26 ++++++++++++++++--- 1 file changed, 22 insertions(+), 4 deletions(-) diff --git a/tests/unit/proxy/hooks/test_tpm_concurrent.py b/tests/unit/proxy/hooks/test_tpm_concurrent.py index 83f00211895..354e86bee81 100644 --- a/tests/unit/proxy/hooks/test_tpm_concurrent.py +++ b/tests/unit/proxy/hooks/test_tpm_concurrent.py @@ -3713,14 +3713,24 @@ async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): """Regression for #41395: summary subrequests gate and charge project ITPM/OTPM.""" from typing import Final - from litellm.llms.anthropic.experimental_pass_through.context_management.editors.compact import ( + from litellm.llms.anthropic.pass_through.context_management.editors.compact import ( _check_summary_model_rate_limit, ) from litellm.proxy import proxy_server handler, _cache = rate_limiter previous_limiter: Final = getattr(proxy_server.proxy_logging_obj, "max_parallel_request_limiter", None) - proxy_server.proxy_logging_obj.max_parallel_request_limiter = handler + previous_hook: Final = proxy_server.proxy_logging_obj.proxy_hook_mapping.get( + "parallel_request_limiter" + ) + + def install_limiter(limiter: RateLimitHandler) -> None: + # Summary gate resolves the limiter via get_proxy_hook(), not the + # legacy max_parallel_request_limiter attribute alone. + proxy_server.proxy_logging_obj.max_parallel_request_limiter = limiter + proxy_server.proxy_logging_obj.proxy_hook_mapping["parallel_request_limiter"] = limiter + + install_limiter(handler) try: model: Final = "gpt-4o-mini" project: Final = "proj-summary-io" @@ -3789,7 +3799,7 @@ async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): ) rpm_handler: Final = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) - proxy_server.proxy_logging_obj.max_parallel_request_limiter = rpm_handler + install_limiter(rpm_handler) allowed_rpm, refusal_rpm = await drive_until_refused( rpm_handler, make_auth(model_rpm_limit={model: 4}) ) @@ -3804,7 +3814,7 @@ async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): ) charging: Final = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) - proxy_server.proxy_logging_obj.max_parallel_request_limiter = charging + install_limiter(charging) summary_response: Final = ModelResponse( usage=Usage(prompt_tokens=60, completion_tokens=40, total_tokens=100) ) @@ -3856,6 +3866,14 @@ async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): assert int(await dual.async_get_cache(key=itpm_key) or 0) >= 60 finally: proxy_server.proxy_logging_obj.max_parallel_request_limiter = previous_limiter + if previous_hook is None: + proxy_server.proxy_logging_obj.proxy_hook_mapping.pop( + "parallel_request_limiter", None + ) + else: + proxy_server.proxy_logging_obj.proxy_hook_mapping[ + "parallel_request_limiter" + ] = previous_hook if __name__ == "__main__": From 4d6f12578de91a14e5e2377807d5459cea145603 Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Thu, 1 Oct 2026 22:15:12 +0800 Subject: [PATCH 5/7] fix(proxy): count summary tokens without building a new list Signed-off-by: sinksilk <785976238@qq.com> --- .../pass_through/context_management/editors/compact.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py index 4be645b8f72..81245b4e8b3 100644 --- a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py +++ b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py @@ -1046,7 +1046,7 @@ async def _estimate_summary_input_tokens( try: return await asyncify(litellm.token_counter)( model=summary_model, - messages=list(summary_messages), + messages=summary_messages, ) except Exception as e: verbose_logger.warning( From 25ca54ca2b30d2582f06f5121fdb83cd7793d10a Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Fri, 2 Oct 2026 00:52:05 +0800 Subject: [PATCH 6/7] fix(proxy): type summary quota helpers without unknown arguments Signed-off-by: sinksilk <785976238@qq.com> --- .../context_management/editors/compact.py | 16 ++++++- .../hooks/parallel_request_limiter_v3.py | 42 +++++++++++++------ 2 files changed, 45 insertions(+), 13 deletions(-) diff --git a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py index 81245b4e8b3..69089d20cba 100644 --- a/litellm/llms/anthropic/pass_through/context_management/editors/compact.py +++ b/litellm/llms/anthropic/pass_through/context_management/editors/compact.py @@ -1037,6 +1037,20 @@ def _build_summary_messages( return summary_messages +def _count_summary_message_tokens( + model: str, + messages: Sequence[Mapping[str, object]], +) -> int: + """Count summary-call tokens through a fully annotated wrapper. + + ``litellm.token_counter``'s own signature is partially unknown (bare + ``Sequence`` / ``dict`` parameters). Passing that function into + ``asyncify`` is a new ``reportUnknownArgumentType``. This wrapper's + signature is fully known, so the asyncify boundary stays typed. + """ + return litellm.token_counter(model=model, messages=messages) + + async def _estimate_summary_input_tokens( *, summary_model: str, @@ -1044,7 +1058,7 @@ async def _estimate_summary_input_tokens( fallback_tokens: int, ) -> int: try: - return await asyncify(litellm.token_counter)( + return await asyncify(_count_summary_message_tokens)( model=summary_model, messages=summary_messages, ) diff --git a/litellm/proxy/hooks/parallel_request_limiter_v3.py b/litellm/proxy/hooks/parallel_request_limiter_v3.py index f8ad7fd2f7e..3c389d2fd09 100644 --- a/litellm/proxy/hooks/parallel_request_limiter_v3.py +++ b/litellm/proxy/hooks/parallel_request_limiter_v3.py @@ -24,6 +24,7 @@ from typing import ( Protocol, TypeAlias, TypedDict, + cast, ) from fastapi import HTTPException @@ -740,6 +741,19 @@ def _call_id_from_callback_kwargs(kwargs: object) -> str | None: return call_id if isinstance(call_id, str) else None +def _as_str_object_dict(value: object) -> dict[str, object] | None: # mutable-ok: model-group helper requires a dict + """Return a ``dict[str, object]`` view of an untyped callback payload. + + ``isinstance(..., dict)`` narrows to ``dict[Unknown, Unknown]``. Keep an + ``object`` alias from before that narrowing and cast that alias, so the + cast argument stays a known type. + """ + raw: Final[object] = value + if not isinstance(value, dict): + return None + return cast("dict[str, object]", raw) # cast-ok: success-callback payload is an untyped dict + + def _parse_output_cap_value(raw_value: object) -> int | None: if isinstance(raw_value, bool) or not isinstance(raw_value, (int, float, str)): return None @@ -4729,7 +4743,7 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): def _build_unreserved_project_io_token_ops( self, - kwargs: Mapping[str, object], + kwargs: object, response_obj: object, ) -> tuple[ReservationAwareIncrementOperation, ...]: """Charge full actual ITPM/OTPM when no pre-call reservation owns this call. @@ -4746,17 +4760,19 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): get_model_group_from_litellm_kwargs, ) - standard_logging_object: Final = kwargs.get("standard_logging_object") - if not isinstance(standard_logging_object, dict): + callback_kwargs: Final = _as_str_object_dict(kwargs) + if callback_kwargs is None: return () - standard_logging_metadata: Final = standard_logging_object.get("metadata") - if not isinstance(standard_logging_metadata, Mapping): + logging_map: Final = _as_str_object_dict(callback_kwargs.get("standard_logging_object")) + if logging_map is None: + return () + standard_logging_metadata: Final = _as_str_object_dict(logging_map.get("metadata")) + if standard_logging_metadata is None: return () - model_group: Final = get_model_group_from_litellm_kwargs(kwargs) or ( - standard_logging_object.get("model_group") - if isinstance(standard_logging_object.get("model_group"), str) - else None + logged_group: Final = logging_map.get("model_group") + model_group: Final = get_model_group_from_litellm_kwargs(callback_kwargs) or ( + logged_group if isinstance(logged_group, str) else None ) targets: Final = self._collect_project_io_scope_targets( standard_logging_metadata=standard_logging_metadata, @@ -4766,10 +4782,11 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): return () response_usage: Final = self._resolve_io_token_reconcile_usage(response_obj) - combined_usage: Final = self._resolve_io_token_reconcile_usage(kwargs.get("combined_usage_object")) + combined_usage_object: Final = callback_kwargs.get("combined_usage_object") + combined_usage: Final = self._resolve_io_token_reconcile_usage(combined_usage_object) aggregate_total: Final = self._aggregate_only_total_tokens( self._response_usage(response_obj) - ) or self._aggregate_only_total_tokens(self._response_usage(kwargs.get("combined_usage_object"))) + ) or self._aggregate_only_total_tokens(self._response_usage(combined_usage_object)) if not response_usage[2] and not combined_usage[2] and aggregate_total <= 0: return () resolved_usage: Final = ( @@ -4807,11 +4824,12 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger): subrequest), falls through to the unreserved metadata rebuild so project IO quotas still receive the summary's actual usage. """ + callback_kwargs: Final[object] = kwargs if not isinstance(kwargs, dict): return () stash: Final = get_request_stash_for_call(_call_id_from_callback_kwargs(kwargs)) if stash is None: - return self._build_unreserved_project_io_token_ops(kwargs, response_obj) + return self._build_unreserved_project_io_token_ops(callback_kwargs, response_obj) itpm_reserved: Final = stash.itpm_reserved_tokens otpm_reserved: Final = stash.otpm_reserved_tokens From e70b9cb01ff2b932b21056652b2b28185867c805 Mon Sep 17 00:00:00 2001 From: sinksilk <785976238@qq.com> Date: Fri, 2 Oct 2026 01:18:40 +0800 Subject: [PATCH 7/7] test(proxy): cover summary quota fallbacks and bump pypdf to 6.19.0 Signed-off-by: sinksilk <785976238@qq.com> --- tests/unit/proxy/hooks/test_tpm_concurrent.py | 153 ++++++++++++++++++ 1 file changed, 153 insertions(+) diff --git a/tests/unit/proxy/hooks/test_tpm_concurrent.py b/tests/unit/proxy/hooks/test_tpm_concurrent.py index 354e86bee81..f4da0d00223 100644 --- a/tests/unit/proxy/hooks/test_tpm_concurrent.py +++ b/tests/unit/proxy/hooks/test_tpm_concurrent.py @@ -30,6 +30,7 @@ from litellm.proxy.hooks.parallel_request_limiter_v3 import ( _PROXY_MaxParallelRequestsHandler_v3 as RateLimitHandler, ) from litellm.proxy.hooks.parallel_request_limiter_v3 import ( + _as_str_object_dict, _call_id_from_callback_kwargs, _request_stash, get_or_create_request_stash, @@ -3876,5 +3877,157 @@ async def test_summary_subrequest_honors_project_itpm_otpm(rate_limiter): ] = previous_hook +@pytest.mark.asyncio +async def test_unreserved_project_io_covers_empty_and_fallback_usage(rate_limiter): + """Cover the unreserved ITPM/OTPM branches the happy-path summary test skips.""" + from typing import Final + + handler, _cache = rate_limiter + assert _as_str_object_dict("nope") is None + echoed: Final = _as_str_object_dict({"a": 1}) + assert echoed is not None + assert echoed["a"] == 1 + + assert handler._build_unreserved_project_io_token_ops("nope", {}) == () + assert handler._build_unreserved_project_io_token_ops({}, {}) == () + assert ( + handler._build_unreserved_project_io_token_ops( + {"standard_logging_object": "x"}, + {}, + ) + == () + ) + assert ( + handler._build_unreserved_project_io_token_ops( + {"standard_logging_object": {"metadata": "x"}}, + {}, + ) + == () + ) + assert ( + handler._build_unreserved_project_io_token_ops( + { + "standard_logging_object": { + "metadata": {"user_api_key_project_id": "proj"}, + "model_group": 1, + } + }, + {}, + ) + == () + ) + + combined_ops: Final = handler._build_unreserved_project_io_token_ops( + { + "standard_logging_object": { + "metadata": {"user_api_key_project_id": "proj"}, + "model_group": "gpt-4o-mini", + }, + "combined_usage_object": {"prompt_tokens": 5, "completion_tokens": 0}, + }, + {}, + ) + assert len(combined_ops) == 1 + assert combined_ops[0]["increment_value"] == 5 + + aggregate_ops: Final = handler._build_unreserved_project_io_token_ops( + { + "standard_logging_object": { + "metadata": {"user_api_key_project_id": "proj"}, + "model_group": "gpt-4o-mini", + }, + }, + {"total_tokens": 9}, + ) + assert len(aggregate_ops) == 2 + assert {op["increment_value"] for op in aggregate_ops} == {9} + + await handler._seed_rate_limit_window_if_absent(window_key="already-open", ttl=None) + await handler._seed_rate_limit_window_if_absent(window_key="already-open", ttl=None) + await handler._apply_one_reservation_aware_token_increment( + operation={ + "key": "plain-counter", + "increment_value": 3, + "ttl": 60, + }, + parent_otel_span=None, + ) + assert int(await handler.internal_usage_cache.async_get_cache(key="plain-counter", litellm_parent_otel_span=None) or 0) == 3 + + +@pytest.mark.asyncio +async def test_summary_token_estimate_uses_counter_or_falls_back(monkeypatch): + from typing import Final + + from litellm.llms.anthropic.pass_through.context_management.editors.compact import ( + _check_summary_model_rate_limit, + _estimate_summary_input_tokens, + ) + from litellm.proxy import proxy_server + + messages: Final = ({"role": "user", "content": "hi"},) + + def fail_counter(**_kwargs: object) -> int: + raise RuntimeError("counter down") + + monkeypatch.setattr("litellm.token_counter", fail_counter) + assert ( + await _estimate_summary_input_tokens( + summary_model="gpt-4o-mini", + summary_messages=messages, + fallback_tokens=42, + ) + == 42 + ) + + def fixed_counter(**_kwargs: object) -> int: + return 17 + + monkeypatch.setattr("litellm.token_counter", fixed_counter) + assert ( + await _estimate_summary_input_tokens( + summary_model="gpt-4o-mini", + summary_messages=messages, + fallback_tokens=42, + ) + == 17 + ) + + handler: Final = RateLimitHandler(internal_usage_cache=InternalUsageCache(DualCache())) + + async def junk_statuses(self, descriptors, parent_otel_span=None, read_only=False, **_kwargs): + return { + "overall_code": "OK", + "statuses": [ + "not-a-status", + {"descriptor_key": "model_per_project_itpm", "limit_remaining": "lots"}, + {"descriptor_key": "model_per_project_itpm", "limit_remaining": 1000}, + ], + } + + monkeypatch.setattr(RateLimitHandler, "should_rate_limit", junk_statuses) + previous_hook: Final = proxy_server.proxy_logging_obj.proxy_hook_mapping.get("parallel_request_limiter") + proxy_server.proxy_logging_obj.proxy_hook_mapping["parallel_request_limiter"] = handler + try: + assert ( + await _check_summary_model_rate_limit( + user_api_key_auth=UserAPIKeyAuth( + api_key="sk-proj-key", + project_id="proj-summary-io", + project_metadata={"model_itpm_limit": {"gpt-4o-mini": 2000}}, + ), + summary_model="gpt-4o-mini", + estimated_input_tokens=10, + estimated_output_tokens=1, + ) + is True + ) + finally: + if previous_hook is None: + proxy_server.proxy_logging_obj.proxy_hook_mapping.pop("parallel_request_limiter", None) + else: + proxy_server.proxy_logging_obj.proxy_hook_mapping["parallel_request_limiter"] = previous_hook + + if __name__ == "__main__": pytest.main([__file__, "-v", "-s"])