fix(realtime): stop revalidating realtime events at the logging boundary (#31054)

Realtime websocket sessions emit events outside the OpenAIRealtimeEvents
union (e.g. rate_limits.updated, response.function_call_arguments.delta,
surfaced when logged_real_time_event_types="*"). Building
LiteLLMRealtimeStreamLoggingObject revalidated every stored event against
the 16-member union, producing thousands of ValidationErrors per session
(12,670 for a ~281-event session). That synchronous work blocked the
asyncio event loop, degrading realtime time-to-first-audio and dial latency
and tripping readiness probes, and the raised error discarded the session
usage so no cost was tracked.

Type results as SkipValidation[OpenAIRealtimeStreamList] and serialize the
events verbatim, so already-formed event dicts are not revalidated. The
flood drops from 12,670 errors to 0 and the combined usage survives to the
cost calculator.

Resolves LIT-3919
Resolves LIT-3920
This commit is contained in:
Yassin Kortam 2026-06-22 20:43:03 -07:00 • committed by GitHub
parent 6f6aec2930
commit a9de75b1f7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 81 additions and 1 deletions

View file

@ -37,6 +37,8 @@ from pydantic import (
ConfigDict,
Field,
PrivateAttr,
SkipValidation,
field_serializer,
field_validator,
)
from typing_extensions import Required, TypedDict
@ -3610,10 +3612,20 @@ class LiteLLMBatch(Batch):
class LiteLLMRealtimeStreamLoggingObject(LiteLLMPydanticObjectBase):
results: OpenAIRealtimeStreamList
# Events are already well-formed provider dicts. Validating them against the
# OpenAIRealtimeEvents union makes Pydantic try every member per event, which
# floods thousands of ValidationErrors for events outside the union (e.g.
# rate_limits.updated), blocks the event loop, and discards the session usage.
results: SkipValidation[OpenAIRealtimeStreamList]
usage: Usage
_hidden_params: dict = {}
@field_serializer("results")
def _serialize_results(
self, results: OpenAIRealtimeStreamList
) -> List[Dict[str, Any]]:
return [dict(event) for event in results]
def __contains__(self, key):
# Define custom behavior for the 'in' operator
return hasattr(self, key)

View file

@ -505,6 +505,74 @@ def test_realtime_logging_object_allows_null_transcript_in_conversation_item_add
assert logging_result.results[0]["item"]["content"][0]["transcript"] is None
def test_realtime_logging_object_does_not_validate_unknown_event_types():
"""
A realtime session emits events outside the OpenAIRealtimeEvents union (e.g.
rate_limits.updated, response.function_call_arguments.delta). Building the
logging object must not revalidate every event against the union; doing so
produces thousands of Pydantic ValidationErrors per session, blocks the event
loop, and the raised error discards the session's usage. The events must
survive verbatim, the combined usage must be preserved, and serialization
must stay clean.
"""
import warnings
results: OpenAIRealtimeStreamList = [
{"type": "session.created", "event_id": "ev0", "session": {"id": "s"}},
]
for i in range(50):
results += [
{
"type": "rate_limits.updated",
"event_id": f"rl{i}",
"rate_limits": [{"name": "requests", "limit": 1000, "remaining": 900}],
},
{
"type": "response.function_call_arguments.delta",
"event_id": f"fc{i}",
"delta": "{}",
},
{
"type": "response.done",
"event_id": f"rd{i}",
"response": {
"usage": {
"input_tokens": 4,
"output_tokens": 6,
"total_tokens": 10,
}
},
},
]
usage = RealtimeAPITokenUsageProcessor.collect_and_combine_usage_from_realtime_stream_results(
results=results
)
# On unfixed code this raises pydantic ValidationError instead of returning.
logging_result = RealtimeAPITokenUsageProcessor.create_logging_realtime_object(
usage=usage,
results=results,
)
assert logging_result.usage.total_tokens == 500
assert len(logging_result.results) == len(results)
unknown_types = {
r["type"]
for r in logging_result.results
if r["type"]
in ("rate_limits.updated", "response.function_call_arguments.delta")
}
assert unknown_types == {
"rate_limits.updated",
"response.function_call_arguments.delta",
}
with warnings.catch_warnings():
warnings.simplefilter("error")
dumped = logging_result.model_dump()
assert len(dumped["results"]) == len(results)
def test_realtime_transcription_duration_cost(monkeypatch):
"""
gpt-realtime-whisper transcription sessions are billed by input audio duration