mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix(realtime): release the budget reservation on a failed session and scrub relayed close details
A refused or failed /v1/realtime session never ran the success cost callback or a failure hook, so its pre-call budget reservation stayed open and kept the key/team/user spend counters pinned above real spend, 429ing later requests on the same key until the counter's TTL expired. The endpoint now reconciles the reservation in a finally, reusing a shared release_or_invalidate_budget_reservation helper that mirrors the success/failure paths (release to zero, else invalidate the reserved counters and finalize). The relayed upstream close message and reason also go through the proxy's client-facing redaction, so a credential, internal hostname, private IP, or server path echoed by the upstream never reaches the client verbatim.
This commit is contained in:
parent
74613f9bd4
commit
af3ddb477a
5 changed files with 137 additions and 10 deletions
|
|
@ -9,7 +9,7 @@ from typing import TYPE_CHECKING, Any, Final, NoReturn, Protocol, TypedDict, cas
|
|||
from typing_extensions import ReadOnly
|
||||
|
||||
import litellm
|
||||
from litellm._logging import _redact_string, verbose_logger
|
||||
from litellm._logging import redact_internal_details_from_client_message, verbose_logger
|
||||
from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER
|
||||
from litellm.llms.base_llm.realtime.transformation import BaseRealtimeConfig
|
||||
from litellm.types.llms.openai import (
|
||||
|
|
@ -1567,8 +1567,8 @@ class RealTimeStreaming:
|
|||
await asyncio.gather(forward_task, client_task, return_exceptions=True)
|
||||
|
||||
async def _close_client(self, close: BackendClose) -> None:
|
||||
redacted_message: Final = _redact_string(close.message)
|
||||
redacted_reason: Final = _redact_string(close.reason)
|
||||
redacted_message: Final = redact_internal_details_from_client_message(close.message)
|
||||
redacted_reason: Final = redact_internal_details_from_client_message(close.reason)
|
||||
try:
|
||||
if close.code != 1000:
|
||||
await self.websocket.send_text(realtime_error_event(redacted_message, error_type="server_error"))
|
||||
|
|
|
|||
|
|
@ -11453,6 +11453,16 @@ def _realtime_query_params_template(model: str | None, intent: str | None) -> tu
|
|||
return tuple(params)
|
||||
|
||||
|
||||
async def _release_realtime_budget_reservation(user_api_key_dict: UserAPIKeyAuth) -> None:
|
||||
from litellm.proxy.spend_tracking.budget_reservation import (
|
||||
release_or_invalidate_budget_reservation,
|
||||
)
|
||||
|
||||
await release_or_invalidate_budget_reservation(
|
||||
budget_reservation=user_api_key_dict.budget_reservation,
|
||||
)
|
||||
|
||||
|
||||
@app.websocket("/openai/v1/realtime")
|
||||
@app.websocket("/v1/realtime")
|
||||
@app.websocket("/realtime")
|
||||
|
|
@ -11592,6 +11602,8 @@ async def realtime_websocket_endpoint(
|
|||
)
|
||||
except Exception: # noqa: BLE001 # the lower layer may have closed the socket already; closing twice is not an error
|
||||
verbose_proxy_logger.debug("Could not close realtime client websocket; it is already gone")
|
||||
finally:
|
||||
await _release_realtime_budget_reservation(user_api_key_dict)
|
||||
|
||||
|
||||
######################################################################
|
||||
|
|
|
|||
|
|
@ -373,6 +373,31 @@ async def invalidate_budget_reservation_counters(
|
|||
await _invalidate_spend_counter(counter_key=counter_key)
|
||||
|
||||
|
||||
async def release_or_invalidate_budget_reservation(
|
||||
budget_reservation: dict | None, # mutable-ok: stamps finalized on the caller's shared reservation dict
|
||||
) -> None:
|
||||
"""Reconcile a still-open reservation on a terminal path that settles no cost.
|
||||
|
||||
A failed or upstream-refused request never runs the success cost callback, so
|
||||
its pre-call reservation stays open and keeps the spend counter pinned above
|
||||
real spend until the counter's TTL expires, 429ing later requests on the same
|
||||
key. Release it to zero; if the release itself fails (e.g. the counter store is
|
||||
unreachable) drop the reserved counters directly and mark the reservation
|
||||
finalized so nothing reprocesses it. Idempotent: the finalized guard makes a
|
||||
second call a no-op once success or failure handling already reconciled.
|
||||
"""
|
||||
if budget_reservation is None or budget_reservation.get("finalized") is True:
|
||||
return
|
||||
try:
|
||||
await release_budget_reservation(budget_reservation=budget_reservation)
|
||||
except Exception: # noqa: BLE001 # a cleanup failure must not pin the counter; drop it directly instead
|
||||
verbose_proxy_logger.exception("Failed to release budget reservation; invalidating counters")
|
||||
try:
|
||||
await invalidate_budget_reservation_counters(budget_reservation=budget_reservation)
|
||||
finally:
|
||||
budget_reservation["finalized"] = True
|
||||
|
||||
|
||||
async def _get_budget_counters(
|
||||
request_body: dict,
|
||||
valid_token: UserAPIKeyAuth,
|
||||
|
|
|
|||
|
|
@ -3229,21 +3229,28 @@ async def test_bidirectional_forward_relays_upstream_policy_close_to_client():
|
|||
client_ws.close.assert_awaited_once_with(code=1008, reason=_UPSTREAM_REFUSAL)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"leaked_detail",
|
||||
(
|
||||
pytest.param("sk-live-abcdef0123456789abcdef0123", id="credential"),
|
||||
pytest.param("vertex-int.svc.cluster.local", id="internal-hostname"),
|
||||
pytest.param("/etc/litellm/service-account.json", id="filesystem-path"),
|
||||
),
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_upstream_close_reason_with_a_secret_is_redacted_before_reaching_the_client():
|
||||
"""LIT-6973: the relayed close mirrors the handshake path and scrubs credential
|
||||
patterns, so an upstream error echoing a token never reaches the client verbatim."""
|
||||
secret: Final = "sk-live-abcdef0123456789abcdef0123"
|
||||
async def test_upstream_close_details_are_scrubbed_before_reaching_the_client(leaked_detail: str):
|
||||
"""LIT-6973: the relayed close goes through the proxy's client-facing redaction, so an upstream
|
||||
error echoing a credential, an internal host, or a server path never reaches the client verbatim."""
|
||||
client_ws: Final = _client_ws_that_never_sends()
|
||||
upstream_close: Final = ConnectionClosed(Close(1008, f"auth failed for {secret}"), None)
|
||||
upstream_close: Final = ConnectionClosed(Close(1008, f"upstream rejected: {leaked_detail}"), None)
|
||||
session: Final = _relay_session(client_ws, _backend_ws_closing_with(upstream_close))
|
||||
|
||||
await session.run()
|
||||
|
||||
(error_event,) = _error_events_sent_to(client_ws)
|
||||
assert secret not in error_event["error"]["message"]
|
||||
assert leaked_detail not in error_event["error"]["message"]
|
||||
relayed_reason: Final = client_ws.close.await_args.kwargs["reason"]
|
||||
assert secret not in relayed_reason
|
||||
assert leaked_detail not in relayed_reason
|
||||
assert "REDACTED" in relayed_reason
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -9521,6 +9521,89 @@ def test_realtime_websocket_route_aliases_registered():
|
|||
)
|
||||
|
||||
|
||||
def _lit6973_fake_realtime_ws() -> MagicMock:
|
||||
ws = MagicMock()
|
||||
ws.headers = {}
|
||||
ws.scope = {"headers": [], "type": "websocket"}
|
||||
ws.url = "ws://testserver/v1/realtime"
|
||||
ws.accept = AsyncMock()
|
||||
ws.send_text = AsyncMock()
|
||||
ws.close = AsyncMock()
|
||||
return ws
|
||||
|
||||
|
||||
async def _lit6973_drive_refused_realtime_session(reservation: dict) -> None:
|
||||
"""Drive realtime_websocket_endpoint through a session the upstream refused.
|
||||
|
||||
route_request resolves normally because the relay handles the refusal
|
||||
internally (sends the error event, closes the client), so neither the
|
||||
success cost callback nor a failure hook runs on _ProxyDBLogger. The
|
||||
endpoint itself must reconcile the pre-call budget reservation, so the
|
||||
real release runs (entries is empty, so it touches no counter store) and
|
||||
the caller asserts on the observable reservation state afterwards."""
|
||||
from litellm.proxy import proxy_server as ps
|
||||
|
||||
user_api_key_dict: Final = UserAPIKeyAuth(api_key="sk-test", token="hashed-token")
|
||||
user_api_key_dict.budget_reservation = reservation
|
||||
|
||||
completed: Final = asyncio.get_running_loop().create_future()
|
||||
completed.set_result(None)
|
||||
|
||||
pre_call: Final = AsyncMock(return_value=({"model": "vertex_ai/gemini-live-2.5-flash"}, MagicMock()))
|
||||
can_call = patch.object(ps, "can_key_call_resolved_model", new=AsyncMock()) # test-quality-ok: no HTTP boundary; fakes in-process auth to reach the finally under test
|
||||
pre = patch.object(ps.ProxyBaseLLMRequestProcessing, "common_processing_pre_call_logic", new=pre_call) # test-quality-ok: fakes phase-1 wiring; assertion checks observable reservation state
|
||||
route = patch.object(ps, "route_request", new=AsyncMock(return_value=completed)) # test-quality-ok: fakes the relay that already handled the refusal so the session returns normally
|
||||
with can_call, pre, route:
|
||||
await ps.realtime_websocket_endpoint(
|
||||
websocket=_lit6973_fake_realtime_ws(),
|
||||
model="vertex_ai/gemini-live-2.5-flash",
|
||||
intent=None,
|
||||
guardrails=None,
|
||||
user_api_key_dict=user_api_key_dict,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_refused_realtime_session_releases_the_budget_reservation():
|
||||
"""LIT-6973: reclassifying a refused realtime session as a failure removed the
|
||||
success-path reservation release, so the pre-call reservation stayed open and
|
||||
pinned the key/team/user spend counters, locking the key after a couple of
|
||||
refusals. The endpoint must reconcile it: the reservation ends up finalized."""
|
||||
reservation: Final = {"reserved_cost": 0.55, "input_cost": 0.0, "finalized": False, "entries": []}
|
||||
|
||||
await _lit6973_drive_refused_realtime_session(reservation)
|
||||
|
||||
assert reservation["finalized"] is True
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_release_or_invalidate_falls_back_to_invalidating_the_counters():
|
||||
"""If releasing the reservation itself fails (e.g. the counter store is down),
|
||||
the reserved counters must be invalidated directly so the estimate does not
|
||||
stay pinned, and the reservation is finalized so nothing reprocesses it."""
|
||||
from litellm.proxy import proxy_server as ps
|
||||
from litellm.proxy.spend_tracking import budget_reservation as br
|
||||
|
||||
reservation: Final = {
|
||||
"reserved_cost": 0.55,
|
||||
"input_cost": 0.0,
|
||||
"finalized": False,
|
||||
"entries": [{"counter_key": "spend:key:hashed-token"}],
|
||||
}
|
||||
invalidated: Final[list[str]] = []
|
||||
|
||||
async def _record(counter_key: str) -> None:
|
||||
invalidated.append(counter_key)
|
||||
|
||||
failing_release = patch.object(br, "release_budget_reservation", new=AsyncMock(side_effect=RuntimeError("counter store down"))) # test-quality-ok: forces the failure branch; assertion observes which counter key got invalidated
|
||||
sink = patch.object(ps, "_invalidate_spend_counter", new=_record) # test-quality-ok: fakes the counter-store sink so the invalidated key is observable
|
||||
with failing_release, sink:
|
||||
await br.release_or_invalidate_budget_reservation(budget_reservation=reservation)
|
||||
|
||||
assert invalidated == ["spend:key:hashed-token"]
|
||||
assert reservation["finalized"] is True
|
||||
|
||||
|
||||
class TestTransformRequestBannedParams:
|
||||
"""
|
||||
/utils/transform_request applies the same banned-param check as LLM endpoints.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue