mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
* fix(proxy): bump health-check max_tokens default to 16 for GPT-5 compatibility (#30708) OpenAI GPT-5 models require max_completion_tokens >= 16. Health checks were using 5 (proxy/health_check.py) and 10 (health_check_helpers.py), causing failures on GPT-5 models. Fixes #23836 * fix: increase health check max_tokens from 5 to 16 (#23836) (#26610) GPT-5 models enforce a minimum of 16 for max_output_tokens. The current default of 5 still causes health checks to fail for these models. Bump the non-wildcard default to 16 — the smallest value that satisfies all known provider minimums while keeping health checks lightweight. Also tightens the wildcard test assertion from a weak disjunctive check to strict key-absence. Co-authored-by: Sameer Kankute <sameer@berri.ai> * fix: ensure checks show gemini-3-flash-preview supports responseJsonS… (#30696) * fix: ensure checks show gemini-3-flash-preview supports responseJsonSchema. * fix: remove async keyword from test. * fix: make Bedrock Mantle Responses routing data-driven per model (#30700) * Make Bedrock Mantle Responses routing data-driven per model Route Bedrock Mantle models to the native Responses API based on each model's price-map capability signal instead of a hardcoded model-name heuristic, and derive the OpenAI-compatible base path segment per model. Responses dispatch now selects the native config when the model advertises responses support (/v1/responses in supported_endpoints, or mode=responses), both overridable via register_model and proxy model_info. This enables native Responses for gpt-oss-120b/20b and the gemma-4 family while keeping chat-only models (gpt-oss safeguard, nvidia, mistral, ...) on the existing chat-completions emulation. Capability is per-model, so gpt-oss-120b routes natively while gpt-oss-safeguard-120b does not despite sharing the gpt-oss substring. The wire path is a separate concern, driven by the existing use_openai_responses_path flag rather than a model-name match: gpt-5.x and gemma-4-* on /openai/v1, everything else (incl. gpt-oss) on /v1. The chat config now derives its base from the same flag, fixing gemma-4 chat-completions requests that previously went to /v1 instead of /openai/v1. Cost maps: add supported_endpoints to the gpt-oss entries (responses for the non-safeguard variants, chat-only for safeguard) and supported_endpoints + use_openai_responses_path to all three gemma-4 entries. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Address review: move capability helper into bedrock_mantle package Move the Responses capability check out of utils.py into litellm/llms/bedrock_mantle/common_utils.py as mantle_supports_responses, alongside its companion wire-path helper mantle_base_segment. Both are now pure functions of (model, model_cost): the price-map mode/supported_endpoints read replaces the get_model_info call, so the rules are unit-testable without patching global state and the Bedrock Mantle package is self-contained. Use str | None instead of Optional[str] on the new signatures to satisfy the ruff UP045 strict-rule gate. Add direct unit tests for both helpers. Fix test_register_model_restore_undoes_existing_key_overwrite: gpt-oss-120b now legitimately supports Responses, so it can no longer be the "None after restore" vehicle; use the chat-only safeguard variant, which isolates the register/restore effect from the model's own capability. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: Sameer Kankute <sameer@berri.ai> * fix(proxy): fail fast on non-PostgreSQL DATABASE_URL instead of hanging on startup (#30366) * fix(proxy): fail fast on non-PostgreSQL DATABASE_URL instead of hanging on startup LiteLLM's Prisma datasource is pinned to provider = 'postgresql', so a sqlite:// or mysql:// DATABASE_URL can never connect. Today that surfaces as an opaque startup stall where the port never binds, and a separate 'DB not connected' 500 on /key/generate when no DATABASE_URL is set at all leaves operators guessing what to configure. Validate the DATABASE_URL / DIRECT_URL scheme in run_server before any Prisma call and exit with an actionable message naming the unsupported scheme. Also reword CommonProxyErrors.db_not_connected_error to tell the operator to set DATABASE_URL to a postgresql:// connection string. Add regression tests covering postgres acceptance and sqlite/mysql/mssql rejection. * fix: resolve CI failures and proxy DB URL typing issue * fix(dashscope): treat an explicit 0.0 tier cost as a real price, not missing (#30653) The tiered cost calculator resolved a tier's per-token cost with `tier.get(cost_key) or tier.get(fallback_cost_key, 0)`. Because `or` short-circuits on any falsy value, a tier that legitimately prices a component at 0.0 (e.g. a free-cache-read tier with cache_read_input_token_cost: 0.0, or a free-reasoning tier) is treated as missing and silently billed at the full fallback rate (input_cost_per_token / output_cost_per_token). The flat-pricing path in the same module already handles this correctly with an `is None` guard. Resolve tier costs through a small helper that mirrors it, so 0.0 is honored at both the in-range and overflow sites. No shipped model currently has a 0.0 tier cost, so this is a latent defect; the fix makes the tiered path consistent with the flat path and prevents over-charging the first time such a tier appears. Adds unit tests covering the in-range and overflow paths, and drops an unused import flagged by ruff in the touched test file. * feat(proxy): show session-aggregate cost and duration in request logs (#25708) (#30507) * fix(anthropic): don't leak tool 'type' into OpenAI function parameters schema (#30618) In the messages->chat/completions bridge, translate_anthropic_tools_to_openai merged every non-mapped tool key into the function parameters dict. The Anthropic tool 'type' (e.g. 'custom') thus overwrote parameters.type ('object' -> 'custom'), and providers reject it ('custom' is not a valid JSON-Schema type). Exclude 'type' from the passthrough. Fixes #30557. * fix(proxy): stop IAM-refresh engine restart from cascading reconnects (#29176) (#30183) An RDS IAM token refresh recreates the Prisma client, which SIGKILLs the running query-engine and spawns a new one. That planned kill was indistinguishable from a crash, and three reconnect paths used two uncoordinated locks, so a single refresh triggered a cascade of engine kill/respawn cycles: 1. `_safe_refresh_token` (holds `_reconnection_lock`) -> recreate -> kill old engine, spawn new one. 2. The engine-death watcher sees that kill, assumes a crash, and calls `attempt_db_reconnect(force=True)` (a different lock, `_db_reconnect_lock`) -> recreate again -> kills the fresh engine. 3. In-flight queries failing during the swap are classified as transport errors and trigger their own `attempt_db_reconnect` -> recreate again. Fix coordinates planned restarts across the wrapper and the watcher: - PrismaWrapper records the old engine PID in `_expected_engine_deaths` before killing it; all four watcher death-detectors (waitpid thread, pidfd, already-dead probe, os.kill poll) consume that PID and skip the reconnect instead of treating it as a crash. - `recreate_prisma_client` now serializes through `_reconnection_lock` and bumps a monotonic `_engine_generation`. Callers pass `expected_generation` as an optimistic-lock token, so racing/cascading recreates collapse into a single restart (losers no-op). This closes the two-lock gap. - The direct reconnect path probes the writer with SELECT 1 before recreating; a healthy connection (e.g. engine already replaced by a refresh) skips the recreate entirely. - `_safe_refresh_token` coalesces: it skips when the current token still has more than the refresh buffer of runway, so stacked triggers (proactive loop + __getattr__ fallback) don't each restart the engine. An `on_engine_replaced` hook re-arms the watcher on the new PID. RoutingPrismaWrapper forwards `expected_generation` and skips recreating the reader when the writer recreate was skipped. * feat(bedrock): support file content retrieval for batch output files (#30595) Implements transform_file_content_request and transform_file_content_response in BedrockFilesConfig so GET /v1/files/{id}/content works for Bedrock batch files. The request transform resolves the file id (direct s3:// URI or base64 unified id) to its S3 object, validates bucket and key prefix against the server-configured bucket, and SigV4-signs an S3 GetObject using the same credential and region resolution as the existing upload path. The credential and region params are validated into a typed model at the boundary, so the only untyped values left are the botocore signing primitives. Also fixes the proxy managed-files path: CredentialLiteLLMParams now carries s3_bucket_name (previously dropped when building deployment credentials) and the managed-files hook passes the deployment credential snapshot when routing afile_content, so unified-id content retrieval works with per-model bucket config instead of only the AWS_S3_BUCKET_NAME env var. Preserves managed-file access control: the proxy file-content endpoint now rejects raw cloud-storage ids (s3://, gs://), which would otherwise skip the owner/team check that only runs for unified ids and let a caller read another tenant's batch output by its object key. Managed outputs are reachable only through their unified file id. The afile_content "not found" error now reports the caller's unified id rather than the resolved internal S3 URI. Fixes #16186, #15563 * fix(oci): make Cohere {{trace}} judges work (tool param types + agentic tool-calling continuation) (#30646) * fix(oci): map Cohere tool array/object params to lowercase builtins OCI's Cohere backend returns HTTP 500 on a tool parameter typed as a bare "List", which is what OCI_JSON_TO_PYTHON_TYPES produced for JSON-schema arrays. MLflow {{trace}} judges trip this: their tools (get_root_span, get_span) take an attributes_to_fetch array. The lowercase builtins list/dict are accepted; only the bare "List" 500s ("Dict" happens to be tolerated, but both are lowercased for consistency). Verified live against us-chicago-1 (cohere.command-a-03-2025 and command-latest). Adds a unit regression on the transformed parameterDefinitions plus a gated integration test exercising an array-param tool end to end. * fix(oci): make Cohere agentic tool-calling continuation work Two bugs broke the OCI Cohere tool-calling loop that MLflow {{trace}} judges drive once a tool has been executed and its result is fed back. Request side: litellm pulled the last user message into the top-level `message` and emitted the tool result as a TOOL entry in chatHistory. OCI rejects that ("cannot specify message if the last entry in chat history contains tool results"), and an empty message alone is rejected too ("message must be at least 1 token long or tool results must be specified"). OCI carries the current turn's results in a dedicated top-level `toolResults` field. The Cohere transform now sends an empty message, keeps the user turn in chatHistory, and puts the results in `toolResults`, matching the langchain-oracle reference. Tool results are no longer represented as chatHistory entries. Response side: tool-grounded answers come back with citations carrying `documentIds` (camelCase) and no `document_ids`, which made the required `CohereCitation.document_ids` field fail validation and sink the whole response parse. Those citations are never surfaced, so the field (and CohereSearchQuery's generation_id) is now optional. Verified live against us-chicago-1 (cohere.command-a-03-2025 and command-latest), single and multi-round tool loops. Adds unit regressions on the transformed request shape and on citation parsing, plus gated integration tests for the continuation. * feat: integrate Repelloai Argus guardrail (#30673) * feat(guardrails): add RepelloAI Argus guardrail integration (#1) * feat(guardrails): add RepelloAI Argus guardrail integration Add a new guardrail hook backed by RepelloAI Argus, with dashboard-managed asset policies enforced via an asset_id and X-API-Key auth. * fix(guardrails): harden RepelloAI Argus guardrail - scan streaming responses on output (was bypassing the guardrail) - log blocked verdicts as guardrail_intervened instead of success - treat auth/config errors (401/403/404/422) as misconfiguration that always blocks, not a fail-open-able unreachable error - default unreachable_fallback to fail_closed and read it directly; block on unknown/malformed verdicts so an API change can't silently disable enforcement - type unreachable_fallback as a Literal, drop the duplicate config model, expose unreachable_fallback in the config schema, and stop leaking the raw provider response / exception strings to the client * fix(guardrails): address RepelloAI Argus review feedback - support ARGUS_API_KEY (with REPELLOAI_API_KEY fallback) - make asset_id required in the config model - normalize unreachable_fallback so only fail_open opens; block on 400 misconfig - correct the shared unreachable_fallback field description * docs(guardrails): add RepelloAI Argus docs page and dashboard listing - add docs page covering config, env vars, modes, verdicts, failure semantics - list RepelloAI Argus in the Guardrail Garden with provider/logo mappings - add a regression test for the provider logo and display-name resolution * fix(guardrails): keep RepelloAI asset_id optional in config model A required asset_id leaked onto the shared LitellmParams (which inherits RepelloAIGuardrailConfigModel), breaking validation for every other guardrail. Keep it optional like sibling models; the guardrail __init__ still raises when asset_id is missing, which is the real enforcement. * Add comment for last user turn scanning * feat(guardrails): harden repelloai scanning * feat(guardrails): expand repelloai scanning to include tool definitions Add extraction of tool definitions and tool call arguments to the RepelloAI guardrail scanning. Improves detection coverage by including function schemas and parameters in the prompt sent to the guardrail service. Also captures detailed error responses in logs and adds guardrail header to streaming responses. * refactor(guardrails): fix and harden repelloai schema text extraction - Fix duplicate text in _iter_schema_text: previously all dict values were re-queued onto the stack even after scalar/list keys were already extracted explicitly, causing names/descriptions to appear twice in the scanned prompt - Extract schema key frozensets to module-level constants so they are not reconstructed on every call - Change _iter_schema_text from @classmethod to @staticmethod (cls unused) - Narrow _call_analyze stage param from str to Literal["prompt", "response"] - Add HttpxResponse type annotation to _raise_for_config_error - Add LLMResponseTypes annotation to async_post_call_success_hook response param * fix(guardrails): resolve pyright type errors in repelloai guardrail - Narrow async_handler.post return from Response|None to Response with explicit None guard before calling raise_for_status/json - Fix list comprehension returning str|None by switching to explicit loop with isinstance guard so pyright tracks the narrowing - Cast model_dump() result to Dict since hasattr does not narrow object type in pyright * fix(guardrails/repello): include Responses API instructions field in prompt scan The /v1/responses top-level `instructions` field was not included in _extract_prompt_text, allowing a caller to bypass guardrail policy checks by putting blocked content in `instructions` while keeping `input` benign. * feat: add api_key to config model and read prompt from data dict * fix(guardrails/repello): plug input_text and tool-call response bypass gaps Responses API input content parts with type 'input_text' were silently dropped by build_inspection_messages (which only handles type='text'), allowing callers to send blocked content via that path without triggering the pre-call scan. Fix: add _extract_input_text_parts to RepelloAIGuardrail and call it when walking the Responses API input messages. Post-call scanning skipped responses whose choices contained only tool_calls or function_call (message.content=None), letting models put blocked output in function arguments undetected. Fix: _extract_chat_completion_text now calls _extract_tool_call_args_from_message on each choice message. Also replace typing.Dict/List with builtin dict/list to clear TID251 strict ruff violations introduced by this file. * fix(guardrails/repello): scan Responses API function_call output arguments Output items with type 'function_call' in a /v1/responses response were skipped by _extract_responses_api_text; only 'message' items were walked. A model could return blocked content in function_call.arguments undetected. Now extract arguments from function_call output items before scanning. * refactor(guardrails/repello): clean up typing and remove lint-any workarounds - Replace Optional[X]/Union[X,Y] with X|None/X|Y union syntax throughout - Use dict[str, object] instead of bare dict in all signatures - Remove **kwargs from __init__; declare guardrail_name, event_hook, default_on explicitly - Replace getattr(litellm_params, ...) with direct attribute access now that LitellmParams inherits RepelloAIGuardrailConfigModel - Add _event_hook_from_mode() to convert str|list[str]|Mode to typed GuardrailEventHooks - Use TypeAdapter.validate_json() instead of response.json() + manual dict construction - Add _is_object_dict/_is_object_list TypeGuard helpers to narrow object types without Any - Remove cast() workarounds and typed intermediate variables that existed only for the now-removed lint-any CI check - Drop _AddLiteLLMCallback Protocol; budget has sufficient slack for the one reportUnknownMemberType - Fix GuardrailConfigModel missing type arg: GuardrailConfigModel[BaseModel] * fix(guardrails/repello): suppress LIT007 on TypeGuard helpers and add streaming scan-skip warning - Add guard-ok suppressions to _is_object_dict and _is_object_list to satisfy the LIT007 hard-zero budget gate - Emit verbose_proxy_logger.warning when the streaming hook finds no inspectable text after assembly, matching observability of pre/post hooks * refactor: modifications for lint check * feat: add Pinstripes as an OpenAI-compatible provider (#30567) * feat: add Pinstripes as an OpenAI-compatible provider Pinstripes (https://pinstripes.io) is an OpenAI-compatible inference provider serving open-source models (GLM-4.5-Air, Qwen3, DeepSeek, etc.) with per-token pricing and no subscriptions. Changes: - `litellm/llms/openai_like/providers.json`: register pinstripes with base_url, api_key_env, and max_completion_tokens→max_tokens mapping - `litellm/types/utils.py`: add `PINSTRIPES = "pinstripes"` to LlmProviders - `litellm/constants.py`: add to openai_compatible_providers and openai_compatible_endpoints lists - `litellm/litellm_core_utils/get_llm_provider_logic.py`: auto-detect provider when api_base is "https://pinstripes.io/v1" - `provider_endpoints_support.json`: document supported endpoints - `tests/`: 7 unit tests covering provider registration, resolution, URL auto-detection, api_base override, and Router config Usage: import litellm response = litellm.completion( model="pinstripes/ps/glm-4.5-air", messages=[{"role": "user", "content": "Hello"}], api_key=os.environ["PINSTRIPES_API_KEY"], ) Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(pinstripes): resolve Greptile P1 review comments - Add api_base_env: PINSTRIPES_API_BASE to providers.json so env var override works - Set responses: false in provider_endpoints_support.json — not actually wired up - Remove docs/my-website/docs/providers/pinstripes.md — belongs in litellm-docs repo Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(pinstripes): add api_base_env and correct responses capability - Add api_base_env: PINSTRIPES_API_BASE to providers.json - Set responses: false in provider_endpoints_support.json Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(pinstripes): wire up Responses API — add supported_endpoints Adds supported_endpoints: ["/v1/chat/completions", "/v1/responses"] so JSONProviderRegistry.supports_responses_api returns true correctly, matching what provider_endpoints_support.json advertises. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * feat(pinstripes): enable embeddings endpoint Pinstripes serves nomic-embed-text-v1.5 and bge-m3 via /v1/embeddings. Add /v1/embeddings to supported_endpoints and set embeddings: true. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(pinstripes): use 4-space indentation in model_prices_and_context_window.json Matches the file's existing convention. Flagged by Greptile review. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(pinstripes): set a2a: false — A2A protocol not implemented All comparable JSON-configured providers (tensormesh, parasail, empiriolabs, libertai, neosantara) have a2a: false. Pinstripes does not implement the Google A2A protocol, so this should be false to match. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> --------- Co-authored-by: inference_provider <max@redactedlab.com> Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(rag): attach existing OpenAI file ids (#30628) * fix(rag): attach existing OpenAI file ids * chore: use modern typing in rag ingest fix * chore: retrigger ci * fix(anthropic-messages): apply cache_control_injection_points on /v1/messages path (#30341) cache_control_injection_points was only consumed by the chat/completions prompt-management hook; on the native Anthropic /v1/messages path it was forwarded unused, so deployment-level cache injection was silently dropped (cache_creation_input_tokens stayed 0 for Anthropic-native clients). Add AnthropicCacheControlHook.apply_to_anthropic_messages_request to inject cache_control at block level for system / tools / message locations (the only forms /v1/messages accepts), wire it into the native anthropic_messages handler, and pop the param so it does not leak upstream as an unknown field. A {location: message, role: system} config is redirected to the top-level system prompt so the same YAML works on both endpoints. Injection respects Anthropic's 4-block cache_control limit shared across system, tools, and messages: client-supplied markers count toward the cap and are never overwritten, a slot is reserved per Bedrock tool_config point, and injection stops once the budget is exhausted. Locations this path cannot represent (tool_config) are forwarded downstream instead of being silently consumed, mirroring get_chat_completion_prompt's remaining_points pass-through. Built on litellm_internal_staging. Refs BerriAI/litellm#30293 * fix(proxy): release budget reservation when a request is cancelled mid-flight (#30522) * fix(proxy): release budget reservation on cancel when no chunk was delivered The pre-call budget reservation increments the cross-pod spend counter by a request's worst-case cost, then reconciles it on success (cost callback) or error (failure hook). A client disconnect or timeout cancels the request and surfaces as CancelledError / GeneratorExit, which neither path catches, so the reservation leaks. Under a retry storm the leaked holds accumulate, pin the counter above real spend, and return spurious 429 "Budget has been exceeded" to keys whose spend is far below budget; the counter only recovers when its TTL lapses, so the failure is intermittent and self-healing. Release the reservation in async_streaming_data_generator (which the Anthropic and Google SSE generators delegate to) on the (CancelledError, GeneratorExit) path, alongside the existing max_parallel_requests release. release_budget_ reservation_on_cancel runs under asyncio.shield so it completes despite the in-progress cancellation, is guarded by the reservation's finalized flag, and swallows a failing release so it cannot replace the in-flight cancellation. The refund is gated on whether a chunk reached the client. The flag is set immediately before the yield, after the slow-path hook await: an async generator suspends at the yield, so a GeneratorExit on disconnect after a delivered chunk sees it True (keep the hold), while a cancellation during the slow-path await leaves it False (refund, nothing sent). A non-streaming cancellation delivers nothing and a completed non-streaming response is reconciled by the success callback, so neither needs a release here. Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * fix(proxy): reconcile a cancelled reservation to input cost, not zero A streaming request cancelled before the first chunk previously reconciled its reservation to zero and finalized it. But by the time the generator is consuming the response the provider call was already dispatched, so the input tokens were billed even though no chunk reached the client, and the success/failure cost callbacks are skipped on cancellation. Refunding to zero let a caller send an expensive request and abort pre-token to dodge the input charge. Compute the request's input-token cost at reservation time and reconcile the cancelled reservation to it instead of zero. The worst-case output portion of the reservation is still released (so a legitimate mid-flight cancellation no longer pins the counter and 429s the key), while the input the provider already processed is charged. --------- Co-authored-by: Bytechoreographer <Bytechoreographer@users.noreply.github.com> Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * fix(caching): encode object name in GCS cache GET path (#30378) GCS cache reads always missed when gcs_path was set. The GET methods interpolated the object name directly into the URL path, while the GCS JSON API requires it to be URL-encoded (a "/" must be sent as %2F). With gcs_path configured the object name is "<prefix>/<sha256>", so the raw slash produced a malformed object path and GCS returned 404. httpx does not raise on 4xx, so the status_code == 200 check fell through and get/async_get returned None, silently missing on every read. Without gcs_path the key has no slash, which is why this went unnoticed. Wrap the object name with urllib.parse.quote(..., safe="") in get_cache and async_get_cache. Apply the same encoding to the name= query parameter in set_cache and async_set_cache so the key written matches the key read back. Adds regression tests asserting the GET path and SET query are encoded (%2F) when gcs_path is set, for both sync and async paths; these fail on the unpatched code. Fixes #30377 * chore: add soniox stt-async-v5 model (#30672) * fix(proxy): include model group aliases in v1 model info (#30626) * Include model group aliases in v1 model info * Fix model info alias implementation * removed extra blank line * chore: rerun CI * fix(lint): remove redundant noqa directive in proxy_cli.py * fix: address greptile review - restore bedrock_mantle auth symbols, guard OCI empty message list, validate DIRECT_URL scheme * Revert "fix: address greptile review - restore bedrock_mantle auth symbols, guard OCI empty message list, validate DIRECT_URL scheme" This reverts commit52c7a07777. * Revert "fix(anthropic-messages): apply cache_control_injection_points on /v1/messages path (#30341)" This reverts commitc9e8a177bd. * Revert "fix(proxy): stop IAM-refresh engine restart from cascading reconnects (#29176) (#30183)" This reverts commit85828da695. * fix(proxy): stop IAM-refresh engine restart from cascading reconnects (#29176) (#30183) An RDS IAM token refresh recreates the Prisma client, which SIGKILLs the running query-engine and spawns a new one. That planned kill was indistinguishable from a crash, and three reconnect paths used two uncoordinated locks, so a single refresh triggered a cascade of engine kill/respawn cycles: 1. `_safe_refresh_token` (holds `_reconnection_lock`) -> recreate -> kill old engine, spawn new one. 2. The engine-death watcher sees that kill, assumes a crash, and calls `attempt_db_reconnect(force=True)` (a different lock, `_db_reconnect_lock`) -> recreate again -> kills the fresh engine. 3. In-flight queries failing during the swap are classified as transport errors and trigger their own `attempt_db_reconnect` -> recreate again. Fix coordinates planned restarts across the wrapper and the watcher: - PrismaWrapper records the old engine PID in `_expected_engine_deaths` before killing it; all four watcher death-detectors (waitpid thread, pidfd, already-dead probe, os.kill poll) consume that PID and skip the reconnect instead of treating it as a crash. - `recreate_prisma_client` now serializes through `_reconnection_lock` and bumps a monotonic `_engine_generation`. Callers pass `expected_generation` as an optimistic-lock token, so racing/cascading recreates collapse into a single restart (losers no-op). This closes the two-lock gap. - The direct reconnect path probes the writer with SELECT 1 before recreating; a healthy connection (e.g. engine already replaced by a refresh) skips the recreate entirely. - `_safe_refresh_token` coalesces: it skips when the current token still has more than the refresh buffer of runway, so stacked triggers (proactive loop + __getattr__ fallback) don't each restart the engine. An `on_engine_replaced` hook re-arms the watcher on the new PID. RoutingPrismaWrapper forwards `expected_generation` and skips recreating the reader when the writer recreate was skipped. * fix(lint): modernize type annotations in IAM-refresh prisma client files (UP006/UP045) * Revert "feat(proxy): show session-aggregate cost and duration in request logs (#25708) (#30507)" This reverts commitf530b2237c. * Revert "fix(dashscope): treat an explicit 0.0 tier cost as a real price, not missing (#30653)" This reverts commit4f58bd0df5. * Revert "fix(oci): make Cohere {{trace}} judges work (tool param types + agentic tool-calling continuation) (#30646)" This reverts commit50f34e0b15. * Revert "fix(proxy): fail fast on non-PostgreSQL DATABASE_URL instead of hanging on startup (#30366)" This reverts commit0544eed6ea. * fix(bedrock_mantle): restore BedrockMantleAuthMixin and constants removed by routing rewrite * fix(key management): restore exact /key/list user_id & key_alias matching by default (#30593) Before substring search was added (commit33bd570d5e), /key/list matched user_id and key_alias exactly. That change made admin-authenticated calls substring-match by default, breaking the prior contract: a caller passing an exact user_id as an access filter (e.g. an integration scoping to one user with an admin key) then received other users' keys -- user_id="alice" also returned "alice2", "alice-test", etc. This is a cross-user key disclosure. Make substring matching opt-in via a new admin-only substring_matching=true query param; default to exact, restoring the prior behavior. The dashboard search box (keyListCall) passes the flag so partial search still works. Non-admins remain exact and scoped to their own keys. Updates the proxy-behavior key_alias test to opt in and adds an exact-by-default guard; adds list_keys unit coverage for the opt-in gate. --------- Co-authored-by: perseus <51974392+tcconnally@users.noreply.github.com> Co-authored-by: Hannah Smith <64043506+hannahmadison@users.noreply.github.com> Co-authored-by: Charlie Patterson <Pattersoncharlesl@gmail.com> Co-authored-by: Matthew Lapointe <mlapointe@alpha-sense.com> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: KRISH SONI <67964054+krishvsoni@users.noreply.github.com> Co-authored-by: Yash Raj Pandey <55940078+devYRPauli@users.noreply.github.com> Co-authored-by: Nitish Agarwal <1592163+nitishagar@users.noreply.github.com> Co-authored-by: hcl <chenglunhu@gmail.com> Co-authored-by: tushar8408 <32977767+tushar8408@users.noreply.github.com> Co-authored-by: AD Mohanraj <admohanraj@gmail.com> Co-authored-by: Fede Kamelhar <federico.kamelhar@oracle.com> Co-authored-by: Lavish Bansal <lavish.bansal619@gmail.com> Co-authored-by: max-amos <gruffulom@gmail.com> Co-authored-by: inference_provider <max@redactedlab.com> Co-authored-by: NK <93352237+Nithish-Yenaganti@users.noreply.github.com> Co-authored-by: 安妮的心动录 <74543653+anneheartrecord@users.noreply.github.com> Co-authored-by: Rick <26716961+Bytechoreographer@users.noreply.github.com> Co-authored-by: Bytechoreographer <Bytechoreographer@users.noreply.github.com> Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> Co-authored-by: Burak Ömür <burak.omur.1998@gmail.com> Co-authored-by: Dan Lemon <daniel.lemon@amazee.io> Co-authored-by: Vanika Dangi <166420943+vanika02@users.noreply.github.com> Co-authored-by: Jay Gowdy <130084966+jgowdy-godaddy@users.noreply.github.com>
2029 lines
68 KiB
Python
2029 lines
68 KiB
Python
import asyncio
|
||
from datetime import datetime, timedelta, timezone
|
||
from unittest.mock import AsyncMock, MagicMock, patch
|
||
|
||
import pytest
|
||
|
||
import litellm
|
||
from litellm.caching.dual_cache import DualCache
|
||
from litellm.proxy._types import (
|
||
LiteLLM_BudgetTable,
|
||
LiteLLM_EndUserTable,
|
||
LiteLLM_OrganizationTable,
|
||
LiteLLM_TagTable,
|
||
LiteLLM_TeamMembership,
|
||
LiteLLM_TeamTable,
|
||
LiteLLM_UserTable,
|
||
UserAPIKeyAuth,
|
||
)
|
||
from litellm.proxy.common_request_processing import ProxyBaseLLMRequestProcessing
|
||
from litellm.proxy.spend_tracking.budget_reservation import (
|
||
estimate_request_max_cost,
|
||
get_budget_window_start,
|
||
invalidate_budget_reservation_counters,
|
||
release_budget_reservation,
|
||
release_budget_reservation_on_cancel,
|
||
reserve_budget_for_request,
|
||
)
|
||
from litellm.proxy.utils import ProxyLogging
|
||
|
||
|
||
@pytest.fixture()
|
||
def spend_counter_state():
|
||
import litellm.proxy.proxy_server as ps
|
||
|
||
original_counter_cache = ps.spend_counter_cache
|
||
original_key_cache = ps.user_api_key_cache
|
||
original_prisma_client = ps.prisma_client
|
||
|
||
counter_cache = DualCache()
|
||
key_cache = DualCache()
|
||
ps.spend_counter_cache = counter_cache
|
||
ps.user_api_key_cache = key_cache
|
||
ps.prisma_client = None
|
||
|
||
try:
|
||
yield counter_cache, key_cache
|
||
finally:
|
||
ps.spend_counter_cache = original_counter_cache
|
||
ps.user_api_key_cache = original_key_cache
|
||
ps.prisma_client = original_prisma_client
|
||
|
||
|
||
def _request_body() -> dict:
|
||
return {
|
||
"model": "gpt-4o-mini",
|
||
"messages": [{"role": "user", "content": "hello"}],
|
||
"max_tokens": 10,
|
||
}
|
||
|
||
|
||
def test_should_not_serialize_budget_reservation_on_user_api_key_auth():
|
||
auth = UserAPIKeyAuth(
|
||
token="key-budget-runtime-state",
|
||
budget_reservation={
|
||
"reserved_cost": 0.5,
|
||
"entries": [{"counter_key": "spend:key:key-budget-runtime-state"}],
|
||
},
|
||
)
|
||
|
||
assert "budget_reservation" not in auth.model_dump()
|
||
assert "budget_reservation" not in auth.model_dump(exclude_none=True)
|
||
assert "budget_reservation" not in auth.model_dump_json()
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_shrink_second_key_reservation_to_remaining_budget(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-race",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.6,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert reservation is not None
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:key:key-budget-race")
|
||
== 0.6
|
||
)
|
||
|
||
second_reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert second_reservation is not None
|
||
assert second_reservation["reserved_cost"] == pytest.approx(0.4)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
with pytest.raises(litellm.BudgetExceededError):
|
||
await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
await release_budget_reservation(second_reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-race"
|
||
) == pytest.approx(0.6)
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_shrink_second_end_user_reservation_to_remaining_budget(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-end-user",
|
||
end_user_id="end-user-budget-race",
|
||
)
|
||
end_user_object = LiteLLM_EndUserTable(
|
||
user_id="end-user-budget-race",
|
||
blocked=False,
|
||
spend=0.0,
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=1.0),
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.6,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
end_user_object=end_user_object,
|
||
)
|
||
assert reservation is not None
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:end-user-budget-race"
|
||
) == pytest.approx(0.6)
|
||
|
||
second_reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
end_user_object=end_user_object,
|
||
)
|
||
assert second_reservation is not None
|
||
assert second_reservation["reserved_cost"] == pytest.approx(0.4)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:end-user-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
with pytest.raises(litellm.BudgetExceededError):
|
||
await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
end_user_object=end_user_object,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:end-user-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
await release_budget_reservation(second_reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:end-user-budget-race"
|
||
) == pytest.approx(0.6)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token=None,
|
||
team_id=None,
|
||
user_id=None,
|
||
response_cost=0.2,
|
||
budget_reservation=reservation,
|
||
end_user_id="end-user-budget-race",
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:end-user-budget-race"
|
||
) == pytest.approx(0.2)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_shrink_second_tag_reservation_to_remaining_budget(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(token="key-budget-tag")
|
||
request_body = _request_body()
|
||
request_body["metadata"] = {
|
||
"tags": ["tag-budget-race", "tag-without-budget", "tag-budget-race"]
|
||
}
|
||
await key_cache.async_set_cache(
|
||
key="tag:tag-budget-race",
|
||
value=LiteLLM_TagTable(
|
||
tag_name="tag-budget-race",
|
||
spend=0.0,
|
||
budget_id="tag-budget-id",
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=1.0),
|
||
).model_dump(),
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="tag:tag-without-budget",
|
||
value=LiteLLM_TagTable(
|
||
tag_name="tag-without-budget",
|
||
spend=0.0,
|
||
).model_dump(),
|
||
)
|
||
prisma_client = MagicMock()
|
||
prisma_client.db.litellm_tagtable.find_many = AsyncMock(return_value=[])
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.6,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=prisma_client,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert reservation is not None
|
||
assert reservation["entries"] == [
|
||
{
|
||
"counter_key": "spend:tag:tag-budget-race",
|
||
"entity_type": "Tag",
|
||
"entity_id": "tag-budget-race",
|
||
"reserved_cost": 0.6,
|
||
"applied_adjustment": 0.0,
|
||
}
|
||
]
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:tag-budget-race"
|
||
) == pytest.approx(0.6)
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:tag:tag-without-budget")
|
||
is None
|
||
)
|
||
|
||
second_reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=prisma_client,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert second_reservation is not None
|
||
assert second_reservation["reserved_cost"] == pytest.approx(0.4)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:tag-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
with pytest.raises(litellm.BudgetExceededError):
|
||
await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=prisma_client,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:tag-budget-race"
|
||
) == pytest.approx(1.0)
|
||
|
||
await release_budget_reservation(second_reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:tag-budget-race"
|
||
) == pytest.approx(0.6)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token=None,
|
||
team_id=None,
|
||
user_id=None,
|
||
response_cost=0.2,
|
||
budget_reservation=reservation,
|
||
tags=["tag-budget-race"],
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:tag-budget-race"
|
||
) == pytest.approx(0.2)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_seed_and_update_end_user_and_tag_counters_without_reservation(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
await key_cache.async_set_cache(
|
||
key="end_user_id:customer-1",
|
||
value=LiteLLM_EndUserTable(
|
||
user_id="customer-1",
|
||
blocked=False,
|
||
spend=4.0,
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=10.0),
|
||
).model_dump(),
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="tag:paid-tag",
|
||
value=LiteLLM_TagTable(
|
||
tag_name="paid-tag",
|
||
spend=7.0,
|
||
).model_dump(),
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="tag:other-tag",
|
||
value=LiteLLM_TagTable(
|
||
tag_name="other-tag",
|
||
spend=2.0,
|
||
).model_dump(),
|
||
)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token=None,
|
||
team_id=None,
|
||
user_id=None,
|
||
response_cost=0.50,
|
||
end_user_id="customer-1",
|
||
tags=["paid-tag", "paid-tag", "other-tag", ""],
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:end_user:customer-1"
|
||
) == pytest.approx(4.50)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:paid-tag"
|
||
) == pytest.approx(7.50)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:tag:other-tag"
|
||
) == pytest.approx(2.50)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reserve_team_member_and_org_budget_counters(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-shared",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
user_id="user-budget-shared",
|
||
team_id="team-budget-shared",
|
||
org_id="org-budget-shared",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-budget-shared",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
user_object = LiteLLM_UserTable(
|
||
user_id="user-budget-shared",
|
||
spend=0.0,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="team_membership:user-budget-shared:team-budget-shared",
|
||
value=LiteLLM_TeamMembership(
|
||
user_id="user-budget-shared",
|
||
team_id="team-budget-shared",
|
||
spend=0.1,
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=1.0),
|
||
).model_dump(),
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="org_id:org-budget-shared:with_budget",
|
||
value=LiteLLM_OrganizationTable(
|
||
organization_id="org-budget-shared",
|
||
organization_alias="shared-org",
|
||
budget_id="org-budget-id",
|
||
spend=0.1,
|
||
models=[],
|
||
created_by="test",
|
||
updated_by="test",
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=1.0),
|
||
).model_dump(),
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.3,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=user_object,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:team_member:user-budget-shared:team-budget-shared"
|
||
) == pytest.approx(0.4)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:org:org-budget-shared"
|
||
) == pytest.approx(0.4)
|
||
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_seed_org_counter_from_with_budget_cache(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
await key_cache.async_set_cache(
|
||
key="org_id:org-counter-with-budget:with_budget",
|
||
value=LiteLLM_OrganizationTable(
|
||
organization_id="org-counter-with-budget",
|
||
organization_alias="shared-org",
|
||
budget_id="org-budget-id",
|
||
spend=2.0,
|
||
models=[],
|
||
created_by="test",
|
||
updated_by="test",
|
||
litellm_budget_table=LiteLLM_BudgetTable(max_budget=10.0),
|
||
).model_dump(),
|
||
)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token=None,
|
||
team_id=None,
|
||
user_id=None,
|
||
org_id="org-counter-with-budget",
|
||
response_cost=0.25,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:org:org-counter-with-budget"
|
||
) == pytest.approx(2.25)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_seed_org_counter_from_plain_org_cache(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
await key_cache.async_set_cache(
|
||
key="org_id:org-counter-plain",
|
||
value=LiteLLM_OrganizationTable(
|
||
organization_id="org-counter-plain",
|
||
organization_alias="shared-org",
|
||
budget_id="org-budget-id",
|
||
spend=2.0,
|
||
models=[],
|
||
created_by="test",
|
||
updated_by="test",
|
||
).model_dump(),
|
||
)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token=None,
|
||
team_id=None,
|
||
user_id=None,
|
||
org_id="org-counter-plain",
|
||
response_cost=0.25,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:org:org-counter-plain"
|
||
) == pytest.approx(2.25)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_cap_known_estimate_to_remaining_budget(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-known-estimate-cap",
|
||
spend=0.9,
|
||
max_budget=1.0,
|
||
)
|
||
counter_cache.in_memory_cache.set_cache(
|
||
key="spend:key:key-budget-known-estimate-cap",
|
||
value=0.9,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.6,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert reservation["reserved_cost"] == pytest.approx(0.1)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-known-estimate-cap"
|
||
) == pytest.approx(1.0)
|
||
|
||
await release_budget_reservation(reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-known-estimate-cap"
|
||
) == pytest.approx(0.9)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_clamp_reservation_to_default_when_output_cap_missing(
|
||
spend_counter_state,
|
||
):
|
||
"""When max_tokens is not specified, _estimate_output_tokens falls back to
|
||
DEFAULT_MAX_OUTPUT_TOKENS_FALLBACK (16K), clamped by the model's
|
||
max_output_tokens. Reservation must be a bounded per-request amount
|
||
(mirroring parallel_request_limiter_v3's DEFAULT_MAX_TOKENS_ESTIMATE),
|
||
not the entire remaining headroom."""
|
||
from litellm.proxy.spend_tracking.budget_reservation import (
|
||
DEFAULT_MAX_OUTPUT_TOKENS_FALLBACK,
|
||
)
|
||
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-uncapped",
|
||
spend=0.2,
|
||
max_budget=10000.0,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="key-budget-uncapped",
|
||
value=valid_token,
|
||
)
|
||
request_body = _request_body()
|
||
request_body.pop("max_tokens")
|
||
|
||
output_cost_per_token = 1e-5 # roughly Opus 4.5/4.7 output rate
|
||
expected_cost = DEFAULT_MAX_OUTPUT_TOKENS_FALLBACK * output_cost_per_token
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"input_cost_per_token": 0.0,
|
||
"output_cost_per_token": output_cost_per_token,
|
||
"max_output_tokens": 200000, # well above the 16K fallback
|
||
},
|
||
):
|
||
estimated = estimate_request_max_cost(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
)
|
||
assert estimated == pytest.approx(expected_cost)
|
||
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert reservation["reserved_cost"] == pytest.approx(expected_cost)
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_clamp_reservation_to_model_ceiling_when_caller_overrequests(
|
||
spend_counter_state,
|
||
):
|
||
"""An adversarial caller sending max_tokens=999_999_999 must not be able
|
||
to inflate the per-request reservation up to the entire remaining team
|
||
headroom. _estimate_output_tokens clamps the explicit value at the
|
||
model's max_output_tokens — the model can only physically emit that
|
||
many tokens anyway, so anything more is both wasteful and a DoS surface."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-overrequest",
|
||
spend=0.0,
|
||
max_budget=10000.0,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="key-budget-overrequest",
|
||
value=valid_token,
|
||
)
|
||
|
||
request_body = _request_body()
|
||
request_body["max_tokens"] = 999_999_999
|
||
|
||
output_cost_per_token = 1e-5
|
||
model_ceiling = 128_000
|
||
expected_cost = model_ceiling * output_cost_per_token
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"input_cost_per_token": 0.0,
|
||
"output_cost_per_token": output_cost_per_token,
|
||
"max_output_tokens": model_ceiling,
|
||
},
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert reservation["reserved_cost"] == pytest.approx(expected_cost)
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reserve_image_generation_cost_per_image(
|
||
spend_counter_state,
|
||
):
|
||
"""Image-generation requests reserve `n × per-image cost` so concurrent
|
||
requests against a depleted budget cannot all bypass the admission gate.
|
||
The OpenAI ``dall-e-3`` entry exposes the per-image price as
|
||
``input_cost_per_image`` (a naming quirk), while other providers use
|
||
``output_cost_per_image`` — both must be honored."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-image-gen",
|
||
spend=0.0,
|
||
max_budget=10.0,
|
||
)
|
||
await key_cache.async_set_cache(key="key-image-gen", value=valid_token)
|
||
|
||
request_body = {"model": "dall-e-3", "prompt": "a cat", "n": 3}
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"mode": "image_generation",
|
||
"input_cost_per_image": 0.04,
|
||
},
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/v1/images/generations",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert reservation["reserved_cost"] == pytest.approx(0.12) # 3 × $0.04
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reject_concurrent_image_request_against_depleted_budget(
|
||
spend_counter_state,
|
||
):
|
||
"""Greptile P1 regression: with image-gen reservation in place, a second
|
||
concurrent image request against a budget already pinned at the cap by
|
||
the first reservation must raise BudgetExceededError instead of
|
||
silently reaching the provider."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-image-deplete",
|
||
spend=0.0,
|
||
team_id="team-image-deplete",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-image-deplete",
|
||
max_budget=0.04,
|
||
spend=0.0,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key=f"team_id:{team_object.team_id}",
|
||
value=team_object,
|
||
)
|
||
|
||
request_body = {"model": "dall-e-3", "prompt": "a cat"}
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"mode": "image_generation",
|
||
"input_cost_per_image": 0.04,
|
||
},
|
||
):
|
||
first = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/v1/images/generations",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert first is not None
|
||
|
||
with pytest.raises(litellm.BudgetExceededError):
|
||
await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/v1/images/generations",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
await release_budget_reservation(first)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_skip_reservation_for_per_pixel_image_model(
|
||
spend_counter_state,
|
||
):
|
||
"""DALL-E 2-style per-pixel pricing depends on the requested ``size``,
|
||
which we don't decode here. Fall through to read-time enforcement
|
||
rather than guess."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-image-per-pixel",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
await key_cache.async_set_cache(key="key-image-per-pixel", value=valid_token)
|
||
|
||
request_body = {"model": "dall-e-2", "prompt": "a cat", "size": "256x256"}
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"mode": "image_generation",
|
||
"input_cost_per_pixel": 2.4414e-07,
|
||
"output_cost_per_pixel": 0.0,
|
||
},
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/v1/images/generations",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is None
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_use_token_pricing_for_chat_model_with_image_cost_field(
|
||
spend_counter_state,
|
||
):
|
||
"""Several chat and embedding models carry ``input_cost_per_image`` /
|
||
``output_cost_per_image`` to price multimodal vision *input*, not image
|
||
generation (e.g. gemini-3.1-pro-preview, azure/gpt-realtime-*,
|
||
amazon.titan-embed-image-v1). _estimate_image_generation_cost must gate
|
||
on ``mode`` so these models still go through the token-priced path —
|
||
otherwise a long chat reserves a fraction of a cent instead of the true
|
||
token cost."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-multimodal-chat",
|
||
spend=0.0,
|
||
max_budget=10.0,
|
||
)
|
||
await key_cache.async_set_cache(key="key-multimodal-chat", value=valid_token)
|
||
|
||
# Roughly the gemini-3.1-pro-preview shape: chat-mode model that
|
||
# carries an output_cost_per_image alongside token pricing.
|
||
output_cost_per_token = 1.2e-5
|
||
request_body = {
|
||
"model": "gemini-3.1-pro-preview",
|
||
"messages": [{"role": "user", "content": "hello"}],
|
||
"max_tokens": 1000,
|
||
}
|
||
expected_cost = 1000 * output_cost_per_token # token-priced path, not 1 × $0.00012
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"mode": "chat",
|
||
"input_cost_per_token": 2e-6,
|
||
"output_cost_per_token": output_cost_per_token,
|
||
"output_cost_per_image": 0.00012,
|
||
"max_output_tokens": 64000,
|
||
},
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
# Token-priced path: reservation ≈ output_tokens × output_cost_per_token,
|
||
# plus a small input-token contribution. Must NOT collapse to the
|
||
# per-image price ($0.00012) which would indicate the image-gen branch
|
||
# incorrectly fired for this chat model.
|
||
assert reservation["reserved_cost"] == pytest.approx(expected_cost, rel=0.05)
|
||
assert reservation["reserved_cost"] > 0.001 # well above per-image price
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reserve_image_edit_cost_per_image(
|
||
spend_counter_state,
|
||
):
|
||
"""``image_edit`` models (Flux Kontext, Stability inpaint/outpaint, etc.)
|
||
bill per generated image just like ``image_generation`` and must get
|
||
the same atomic per-image reservation."""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-image-edit",
|
||
spend=0.0,
|
||
max_budget=10.0,
|
||
)
|
||
await key_cache.async_set_cache(key="key-image-edit", value=valid_token)
|
||
|
||
request_body = {"model": "stability/inpaint", "prompt": "a cat", "n": 2}
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"mode": "image_edit",
|
||
"output_cost_per_image": 0.05,
|
||
},
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/v1/images/edits",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert reservation["reserved_cost"] == pytest.approx(0.10) # 2 × $0.05
|
||
await release_budget_reservation(reservation)
|
||
|
||
|
||
def test_should_start_window_without_reset_at_at_duration_boundary():
|
||
before = datetime.now(timezone.utc) - timedelta(hours=1)
|
||
|
||
window_start = get_budget_window_start({"budget_duration": "1h"})
|
||
|
||
after = datetime.now(timezone.utc) - timedelta(hours=1)
|
||
assert window_start is not None
|
||
assert before <= window_start <= after
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_skip_budget_window_with_unparseable_duration(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-malformed-window",
|
||
spend=0.9,
|
||
max_budget=10.0,
|
||
budget_limits=[
|
||
{
|
||
"budget_duration": "not-a-duration",
|
||
"max_budget": 1.0,
|
||
}
|
||
],
|
||
)
|
||
counter_cache.in_memory_cache.set_cache(
|
||
key="spend:key:key-budget-malformed-window",
|
||
value=0.9,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.2,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is not None
|
||
assert [entry["counter_key"] for entry in reservation["entries"]] == [
|
||
"spend:key:key-budget-malformed-window"
|
||
]
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-malformed-window"
|
||
) == pytest.approx(1.1)
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-malformed-window:window:not-a-duration"
|
||
)
|
||
is None
|
||
)
|
||
|
||
await release_budget_reservation(reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-malformed-window"
|
||
) == pytest.approx(0.9)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_skip_window_reservation_when_db_baseline_unavailable(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-window-db-unavailable",
|
||
budget_limits=[
|
||
{
|
||
"budget_duration": "1h",
|
||
"max_budget": 1.0,
|
||
}
|
||
],
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.5,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is None
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-window-db-unavailable:window:1h"
|
||
)
|
||
is None
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_skip_reservation_when_counter_increment_fails(
|
||
spend_counter_state,
|
||
monkeypatch,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-reserve-unavailable",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
async def fail_increment_cache(*args, **kwargs):
|
||
raise RuntimeError("counter unavailable")
|
||
|
||
monkeypatch.setattr(counter_cache, "async_increment_cache", fail_increment_cache)
|
||
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.5,
|
||
),
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.verbose_proxy_logger.warning"
|
||
) as mock_warning,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is None
|
||
assert mock_warning.call_count >= 1
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-reserve-unavailable"
|
||
)
|
||
is None
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_skip_reservation_when_counter_initialization_fails(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-reserve-init-unavailable",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.5,
|
||
),
|
||
patch(
|
||
"litellm.proxy.proxy_server._ensure_spend_counter_initialized",
|
||
side_effect=RuntimeError("redis unavailable"),
|
||
),
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.verbose_proxy_logger.warning"
|
||
) as mock_warning,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is None
|
||
assert mock_warning.call_count >= 1
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-reserve-init-unavailable"
|
||
)
|
||
is None
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_release_tracked_entry_when_reservation_fails_after_increment(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-reserve-after-increment-failure",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
import litellm.proxy.proxy_server as ps
|
||
|
||
original_increment_counter = ps._increment_spend_counter_cache
|
||
first_increment = True
|
||
|
||
async def fail_after_increment(counter_key: str, increment: float):
|
||
nonlocal first_increment
|
||
if first_increment:
|
||
first_increment = False
|
||
await counter_cache.async_increment_cache(key=counter_key, value=increment)
|
||
raise RuntimeError("lost increment response")
|
||
return await original_increment_counter(
|
||
counter_key=counter_key,
|
||
increment=increment,
|
||
)
|
||
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.5,
|
||
),
|
||
patch(
|
||
"litellm.proxy.proxy_server._increment_spend_counter_cache",
|
||
side_effect=fail_after_increment,
|
||
),
|
||
patch(
|
||
"litellm.proxy.proxy_server._invalidate_spend_counter",
|
||
side_effect=RuntimeError("invalidate unavailable"),
|
||
),
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert reservation is None
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-reserve-after-increment-failure"
|
||
) == pytest.approx(0.0)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reconcile_reserved_counter_to_actual_spend(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-reconcile",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.6,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
from litellm.proxy.proxy_server import increment_spend_counters
|
||
|
||
await increment_spend_counters(
|
||
token="key-budget-reconcile",
|
||
team_id="team-without-budget",
|
||
user_id=None,
|
||
response_cost=0.2,
|
||
budget_reservation=reservation,
|
||
)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-reconcile"
|
||
) == pytest.approx(0.2)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:team:team-without-budget"
|
||
) == pytest.approx(0.2)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_release_reservation_on_failure(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-release",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.4,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
await release_budget_reservation(reservation)
|
||
await release_budget_reservation(reservation)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-release"
|
||
) == pytest.approx(0.0)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_retry_partial_release_without_double_decrement(
|
||
spend_counter_state,
|
||
monkeypatch,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-partial-release",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
team_id="team-budget-partial-release",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-budget-partial-release",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.4,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
original_increment_cache = counter_cache.async_increment_cache
|
||
fail_next_team_release = True
|
||
|
||
async def flaky_increment_cache(key, value, *args, **kwargs):
|
||
nonlocal fail_next_team_release
|
||
if (
|
||
key == "spend:team:team-budget-partial-release"
|
||
and value < 0
|
||
and fail_next_team_release
|
||
):
|
||
fail_next_team_release = False
|
||
raise RuntimeError("simulated counter failure")
|
||
return await original_increment_cache(key=key, value=value, *args, **kwargs)
|
||
|
||
monkeypatch.setattr(counter_cache, "async_increment_cache", flaky_increment_cache)
|
||
|
||
with pytest.raises(RuntimeError):
|
||
await release_budget_reservation(reservation)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-partial-release"
|
||
) == pytest.approx(0.0)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:team:team-budget-partial-release"
|
||
) == pytest.approx(0.4)
|
||
|
||
await release_budget_reservation(reservation)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-partial-release"
|
||
) == pytest.approx(0.0)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:team:team-budget-partial-release"
|
||
) == pytest.approx(0.0)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_preserve_budget_error_and_continue_partial_cleanup(
|
||
spend_counter_state,
|
||
monkeypatch,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-cleanup-failure",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
team_id="team-budget-cleanup-failure",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-budget-cleanup-failure",
|
||
spend=0.3,
|
||
max_budget=0.3,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key="team_id:team-budget-cleanup-failure",
|
||
value=team_object,
|
||
)
|
||
|
||
original_increment_cache = counter_cache.async_increment_cache
|
||
fail_key_cleanup = True
|
||
|
||
async def flaky_increment_cache(key, value, *args, **kwargs):
|
||
nonlocal fail_key_cleanup
|
||
if key == "spend:key:key-budget-cleanup-failure" and value < 0:
|
||
if fail_key_cleanup:
|
||
fail_key_cleanup = False
|
||
raise RuntimeError("simulated cleanup failure")
|
||
return await original_increment_cache(key=key, value=value, *args, **kwargs)
|
||
|
||
monkeypatch.setattr(counter_cache, "async_increment_cache", flaky_increment_cache)
|
||
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.4,
|
||
),
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.verbose_proxy_logger.exception"
|
||
) as mock_log_exception,
|
||
):
|
||
with pytest.raises(litellm.BudgetExceededError):
|
||
await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-cleanup-failure"
|
||
)
|
||
is None
|
||
)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:team:team-budget-cleanup-failure"
|
||
) == pytest.approx(0.3)
|
||
mock_log_exception.assert_called()
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_missing_counter_reseeds_from_db_instead_of_failing(
|
||
spend_counter_state,
|
||
):
|
||
"""A reconcile/release that finds the counter missing must NOT delete it and
|
||
raise (the old fail-open that left budgets unenforced after a Redis reload).
|
||
It reseeds from the authoritative DB; with no DB it leaves the counter
|
||
untouched and finalizes."""
|
||
counter_cache, _ = spend_counter_state
|
||
reservation = {
|
||
"reserved_cost": 0.4,
|
||
"entries": [
|
||
{
|
||
"counter_key": "spend:key:key-budget-missing-release",
|
||
"reserved_cost": 0.4,
|
||
"applied_adjustment": 0.0,
|
||
}
|
||
],
|
||
"finalized": False,
|
||
}
|
||
|
||
# must not raise
|
||
await release_budget_reservation(reservation)
|
||
|
||
# counter not driven negative / not corrupted; left absent (no DB to reseed)
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-missing-release"
|
||
)
|
||
is None
|
||
)
|
||
assert reservation["finalized"] is True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_underflow_counter_reseeds_from_db(spend_counter_state):
|
||
"""When the release delta would drive the counter negative (counter was
|
||
reset/reseeded mid-flight), reseed from the authoritative DB rather than
|
||
deleting and failing open."""
|
||
import litellm.proxy.proxy_server as ps
|
||
|
||
counter_cache, _ = spend_counter_state
|
||
await counter_cache.async_increment_cache(
|
||
key="spend:key:key-budget-underflow-release",
|
||
value=0.1,
|
||
)
|
||
reservation = {
|
||
"reserved_cost": 0.4,
|
||
"entries": [
|
||
{
|
||
"counter_key": "spend:key:key-budget-underflow-release",
|
||
"reserved_cost": 0.4,
|
||
"applied_adjustment": 0.0,
|
||
}
|
||
],
|
||
"finalized": False,
|
||
}
|
||
|
||
with patch.object(ps.SpendCounterReseed, "from_db", AsyncMock(return_value=0.25)):
|
||
await release_budget_reservation(reservation)
|
||
|
||
# counter reseeded up to the authoritative DB value, not deleted or negated
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-underflow-release"
|
||
) == pytest.approx(0.25)
|
||
assert reservation["finalized"] is True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_non_numeric_counter_reseeds_from_db(spend_counter_state):
|
||
"""A non-numeric counter value (corrupt/stale) during release is recovered by
|
||
reseeding from the DB, not by deleting the counter and raising."""
|
||
import litellm.proxy.proxy_server as ps
|
||
|
||
counter_cache, _ = spend_counter_state
|
||
counter_cache.in_memory_cache.set_cache(
|
||
key="spend:key:key-budget-nonnumeric-release",
|
||
value="stale",
|
||
)
|
||
reservation = {
|
||
"reserved_cost": 0.4,
|
||
"entries": [
|
||
{
|
||
"counter_key": "spend:key:key-budget-nonnumeric-release",
|
||
"reserved_cost": 0.4,
|
||
"applied_adjustment": 0.0,
|
||
}
|
||
],
|
||
"finalized": False,
|
||
}
|
||
|
||
with patch.object(ps.SpendCounterReseed, "from_db", AsyncMock(return_value=0.5)):
|
||
await release_budget_reservation(reservation)
|
||
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-budget-nonnumeric-release"
|
||
) == pytest.approx(0.5)
|
||
assert reservation["finalized"] is True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_invalidate_reserved_counters_after_persisted_spend_failure(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, _ = spend_counter_state
|
||
await counter_cache.async_increment_cache(
|
||
key="spend:key:key-budget-invalidate",
|
||
value=0.4,
|
||
)
|
||
await counter_cache.async_increment_cache(
|
||
key="spend:team:team-budget-invalidate",
|
||
value=0.4,
|
||
)
|
||
|
||
await invalidate_budget_reservation_counters(
|
||
{
|
||
"reserved_cost": 0.4,
|
||
"entries": [
|
||
{"counter_key": "spend:key:key-budget-invalidate"},
|
||
{"counter_key": "spend:team:team-budget-invalidate"},
|
||
],
|
||
}
|
||
)
|
||
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:key:key-budget-invalidate")
|
||
is None
|
||
)
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:team:team-budget-invalidate")
|
||
is None
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_reserve_all_budgeted_counters(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-budget-all",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
team_id="team-budget-all",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-budget-all",
|
||
spend=0.0,
|
||
max_budget=1.0,
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=0.3,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:key:key-budget-all") == 0.3
|
||
)
|
||
assert (
|
||
counter_cache.in_memory_cache.get_cache(key="spend:team:team-budget-all") == 0.3
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_should_not_block_concurrent_team_request_when_first_request_lacks_max_tokens(
|
||
spend_counter_state,
|
||
):
|
||
"""
|
||
Regression test: a team-bound request with no max_tokens must not pin the
|
||
team's spend counter at max_budget for the duration of the request.
|
||
|
||
Repro of the integration-test team being falsely budget-blocked at the
|
||
$2000 cap while DB spend is $0.144: the first request without max_tokens
|
||
used to reserve the entire remaining headroom, leaving any subsequent
|
||
request stuck behind a counter sitting at the cap until the success
|
||
callback finished reconciling.
|
||
"""
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-team-integration-tests",
|
||
spend=0.0,
|
||
team_id="team-integration-tests",
|
||
)
|
||
team_object = LiteLLM_TeamTable(
|
||
team_id="team-integration-tests",
|
||
max_budget=2000.0,
|
||
spend=0.144,
|
||
)
|
||
await key_cache.async_set_cache(
|
||
key=f"team_id:{team_object.team_id}",
|
||
value=team_object,
|
||
)
|
||
|
||
request_body = _request_body()
|
||
request_body.pop("max_tokens")
|
||
|
||
# Realistic Opus 4.7 output pricing — the 16K fallback × $25/M ≈ $0.40
|
||
# reservation per request, leaving ~5000 admittable concurrent requests
|
||
# against a $2000 team budget.
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation._get_model_cost_info",
|
||
return_value={
|
||
"input_cost_per_token": 5e-6,
|
||
"output_cost_per_token": 2.5e-5,
|
||
"max_output_tokens": 128000,
|
||
},
|
||
):
|
||
first_reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
|
||
# The team counter must not be pinned at max_budget while the first
|
||
# request is in flight, otherwise concurrent requests false-positive.
|
||
team_counter_after_first = (
|
||
counter_cache.in_memory_cache.get_cache(
|
||
key=f"spend:team:{team_object.team_id}"
|
||
)
|
||
or 0.0
|
||
)
|
||
assert team_counter_after_first < team_object.max_budget, (
|
||
f"Team counter sat at {team_counter_after_first} after one uncapped "
|
||
f"reservation against a {team_object.max_budget} budget — concurrent "
|
||
"requests will be falsely blocked."
|
||
)
|
||
|
||
# Second request — same shape — must succeed without raising.
|
||
second_reservation = await reserve_budget_for_request(
|
||
request_body=request_body,
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=team_object,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert second_reservation is not None
|
||
|
||
if first_reservation is not None:
|
||
await release_budget_reservation(first_reservation)
|
||
if second_reservation is not None:
|
||
await release_budget_reservation(second_reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_budget_reservation_on_cancel_gives_back_counter(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-cancel-give-back", spend=0.0, max_budget=10.0
|
||
)
|
||
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=3.0,
|
||
),
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_input_cost",
|
||
return_value=0.5,
|
||
),
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert reservation is not None
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-give-back"
|
||
) == pytest.approx(3.0)
|
||
|
||
await release_budget_reservation_on_cancel(reservation)
|
||
|
||
# the provider already received the input, so the reservation is reconciled
|
||
# to the input cost (0.5), not refunded to zero; the worst-case output
|
||
# reservation (3.0 -> 0.5) is released
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-give-back"
|
||
) == pytest.approx(0.5)
|
||
assert reservation["finalized"] is True
|
||
|
||
# idempotent: a second cancel reconcile must not change the counter again
|
||
await release_budget_reservation_on_cancel(reservation)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-give-back"
|
||
) == pytest.approx(0.5)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_budget_reservation_on_cancel_noop_when_finalized(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token = UserAPIKeyAuth(
|
||
token="key-cancel-finalized", spend=0.0, max_budget=10.0
|
||
)
|
||
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=3.0,
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert reservation is not None
|
||
reservation["finalized"] = True
|
||
|
||
await release_budget_reservation_on_cancel(reservation)
|
||
|
||
# already reconciled by the success/failure path -> must stay untouched
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-finalized"
|
||
) == pytest.approx(3.0)
|
||
|
||
|
||
async def _reserve_for_stream(counter_cache, key_cache, proxy_logging_obj, token: str):
|
||
valid_token = UserAPIKeyAuth(token=token, spend=0.0, max_budget=10.0)
|
||
with (
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_max_cost",
|
||
return_value=2.0,
|
||
),
|
||
patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.estimate_request_input_cost",
|
||
return_value=0.5,
|
||
),
|
||
):
|
||
reservation = await reserve_budget_for_request(
|
||
request_body=_request_body(),
|
||
route="/chat/completions",
|
||
llm_router=None,
|
||
valid_token=valid_token,
|
||
team_object=None,
|
||
user_object=None,
|
||
prisma_client=None,
|
||
user_api_key_cache=key_cache,
|
||
proxy_logging_obj=proxy_logging_obj,
|
||
)
|
||
assert reservation is not None
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key=f"spend:key:{token}"
|
||
) == pytest.approx(2.0)
|
||
valid_token.budget_reservation = reservation
|
||
return valid_token, reservation
|
||
|
||
|
||
def _drive_streaming_cancel(valid_token, iterator_hook):
|
||
streaming_logging_obj = MagicMock()
|
||
streaming_logging_obj.async_post_call_streaming_iterator_hook = iterator_hook
|
||
return ProxyBaseLLMRequestProcessing.async_streaming_data_generator(
|
||
response=MagicMock(),
|
||
user_api_key_dict=valid_token,
|
||
request_data=_request_body(),
|
||
proxy_logging_obj=streaming_logging_obj,
|
||
serialize_chunk=lambda chunk: chunk,
|
||
serialize_error=lambda exc: str(exc),
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_streaming_cancel_before_any_chunk_reconciles_to_input_cost(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token, reservation = await _reserve_for_stream(
|
||
counter_cache, key_cache, proxy_logging_obj, "key-cancel-no-chunk"
|
||
)
|
||
|
||
# Client disconnects before the upstream produced any output.
|
||
async def cancel_before_chunk(user_api_key_dict, response, request_data):
|
||
if False:
|
||
yield "" # make this an async generator
|
||
raise asyncio.CancelledError()
|
||
|
||
generator = _drive_streaming_cancel(valid_token, cancel_before_chunk)
|
||
received = []
|
||
with pytest.raises(asyncio.CancelledError):
|
||
async for chunk in generator:
|
||
received.append(chunk)
|
||
|
||
assert received == []
|
||
# no chunk delivered, but the provider already received the input, so the
|
||
# reservation is reconciled to the input cost (0.5), not refunded to zero
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-no-chunk"
|
||
) == pytest.approx(0.5)
|
||
assert reservation["finalized"] is True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_streaming_cancel_after_chunk_keeps_reservation(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token, reservation = await _reserve_for_stream(
|
||
counter_cache, key_cache, proxy_logging_obj, "key-cancel-after-chunk"
|
||
)
|
||
|
||
# Client consumes a chunk, then disconnects. Cancellation logs no cost, so
|
||
# refunding here would let the caller read partial output for free.
|
||
async def cancel_after_chunk(user_api_key_dict, response, request_data):
|
||
yield "data: chunk\n\n"
|
||
raise asyncio.CancelledError()
|
||
|
||
generator = _drive_streaming_cancel(valid_token, cancel_after_chunk)
|
||
received = []
|
||
with pytest.raises(asyncio.CancelledError):
|
||
async for chunk in generator:
|
||
received.append(chunk)
|
||
|
||
assert received == ["data: chunk\n\n"]
|
||
# a consumed stream must NOT be refunded
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-after-chunk"
|
||
) == pytest.approx(2.0)
|
||
assert reservation.get("finalized") is not True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_release_budget_reservation_on_cancel_swallows_release_errors():
|
||
# If the release itself fails (e.g. Redis unavailable) it must not escape
|
||
# the helper: doing so would replace the in-flight CancelledError /
|
||
# GeneratorExit at the call site and disrupt the disconnect teardown.
|
||
reservation = {
|
||
"reserved_cost": 3.0,
|
||
"entries": [{"counter_key": "spend:key:key-cancel-error"}],
|
||
"finalized": False,
|
||
"input_cost": 0.5,
|
||
}
|
||
with patch(
|
||
"litellm.proxy.spend_tracking.budget_reservation.reconcile_budget_reservation",
|
||
new=AsyncMock(side_effect=RuntimeError("redis down")),
|
||
):
|
||
# must return without raising
|
||
await release_budget_reservation_on_cancel(reservation)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_streaming_cancel_in_slow_path_before_yield_refunds(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token, reservation = await _reserve_for_stream(
|
||
counter_cache, key_cache, proxy_logging_obj, "key-cancel-slowpath"
|
||
)
|
||
|
||
async def one_chunk(user_api_key_dict, response, request_data):
|
||
yield "data: chunk\n\n"
|
||
|
||
streaming_logging_obj = MagicMock()
|
||
streaming_logging_obj.async_post_call_streaming_iterator_hook = one_chunk
|
||
# On the slow path the per-chunk hook is awaited before the chunk is yielded
|
||
# to the client; cancel there. Nothing has reached the client yet.
|
||
streaming_logging_obj.async_post_call_streaming_hook = AsyncMock(
|
||
side_effect=asyncio.CancelledError()
|
||
)
|
||
|
||
generator = ProxyBaseLLMRequestProcessing.async_streaming_data_generator(
|
||
response=MagicMock(),
|
||
user_api_key_dict=valid_token,
|
||
request_data=_request_body(),
|
||
proxy_logging_obj=streaming_logging_obj,
|
||
serialize_chunk=lambda chunk: chunk,
|
||
serialize_error=lambda exc: str(exc),
|
||
)
|
||
|
||
received = []
|
||
# include_cost_in_streaming_usage forces fast_path off, so the hook above runs
|
||
with patch.object(litellm, "include_cost_in_streaming_usage", True, create=True):
|
||
with pytest.raises(asyncio.CancelledError):
|
||
async for chunk in generator:
|
||
received.append(chunk)
|
||
|
||
assert received == []
|
||
# cancellation happened before any chunk reached the client, but the
|
||
# provider already received the input -> reconcile to the input cost (0.5)
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-cancel-slowpath"
|
||
) == pytest.approx(0.5)
|
||
assert reservation["finalized"] is True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_streaming_disconnect_after_consuming_chunk_keeps_reservation(
|
||
spend_counter_state,
|
||
):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token, reservation = await _reserve_for_stream(
|
||
counter_cache, key_cache, proxy_logging_obj, "key-disconnect-after-chunk"
|
||
)
|
||
|
||
async def two_chunks(user_api_key_dict, response, request_data):
|
||
yield "data: a\n\n"
|
||
yield "data: b\n\n"
|
||
|
||
generator = _drive_streaming_cancel(valid_token, two_chunks)
|
||
|
||
# Client consumes one chunk, then disconnects. aclose() raises GeneratorExit
|
||
# at the suspended yield, after the chunk already reached the client.
|
||
first = await generator.__anext__()
|
||
assert first == "data: a\n\n"
|
||
await generator.aclose()
|
||
|
||
# output was delivered, so the reservation must NOT be refunded
|
||
assert counter_cache.in_memory_cache.get_cache(
|
||
key="spend:key:key-disconnect-after-chunk"
|
||
) == pytest.approx(2.0)
|
||
assert reservation.get("finalized") is not True
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_streaming_slow_path_processes_and_yields_chunk(spend_counter_state):
|
||
counter_cache, key_cache = spend_counter_state
|
||
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
|
||
valid_token, _ = await _reserve_for_stream(
|
||
counter_cache, key_cache, proxy_logging_obj, "key-slowpath-ok"
|
||
)
|
||
|
||
async def one_chunk(user_api_key_dict, response, request_data):
|
||
yield {"content": "hi"}
|
||
|
||
streaming_logging_obj = MagicMock()
|
||
streaming_logging_obj.async_post_call_streaming_iterator_hook = one_chunk
|
||
streaming_logging_obj.async_post_call_streaming_hook = AsyncMock(
|
||
side_effect=lambda **kwargs: kwargs["response"]
|
||
)
|
||
|
||
generator = ProxyBaseLLMRequestProcessing.async_streaming_data_generator(
|
||
response=MagicMock(),
|
||
user_api_key_dict=valid_token,
|
||
request_data=_request_body(),
|
||
proxy_logging_obj=streaming_logging_obj,
|
||
serialize_chunk=lambda chunk: chunk,
|
||
serialize_error=lambda exc: str(exc),
|
||
)
|
||
|
||
received = []
|
||
# include_cost_in_streaming_usage forces the slow path so the per-chunk hook,
|
||
# content accumulation, and cost-injection branch all run to a successful yield
|
||
with patch.object(litellm, "include_cost_in_streaming_usage", True, create=True):
|
||
async for chunk in generator:
|
||
received.append(chunk)
|
||
|
||
assert received == [{"content": "hi"}]
|
||
streaming_logging_obj.async_post_call_streaming_hook.assert_awaited_once()
|