From b8635bbc7a10d1b3dff9b804d7a15eb64f98ef56 Mon Sep 17 00:00:00 2001 From: Sameer Kankute Date: Wed, 6 May 2026 05:19:20 +0530 Subject: [PATCH] feat(realtime): OpenAI Realtime GA support and beta compatibility (#27110) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(realtime): OpenAI Realtime GA support and beta compatibility - Normalize beta-style session.update to GA for upstream OpenAI; optional GA→beta event translation when client sends OpenAI-Beta: realtime=v1 - Default upstream WebSocket without OpenAI-Beta; forward header when client opts in - Extend OpenAI realtime types for GA event names and conversation item shapes - Relax LiteLLMRealtimeStreamLoggingObject.results to List[Any] for GA events - Update proxy client_secrets fallback to omit beta header; dashboard RealtimePlayground - Add unit tests for remap, translation, and beta header helper Co-authored-by: Cursor * fix results * fix greptile * Fix mypy issues * Remove unused class constants _GA_TEXT_DELTA_TYPES and _GA_AUDIO_DELTA_TYPES These frozensets were defined as class-level constants in realtime_streaming.py but never referenced anywhere in the codebase. Removing dead code. Co-authored-by: Sameer Kankute * fix(realtime): use GA-shaped session.update in guardrail injections The guardrail VAD injection code sent a beta-style session.update with a flat turn_detection field: {"session": {"turn_detection": {"create_response": false}}} When the upstream OpenAI backend operates in GA mode (no OpenAI-Beta header forwarded), it requires the nested GA shape: {"session": {"type": "realtime", "audio": {"input": {"turn_detection": {"create_response": false}}}}} The _remap_beta_session_to_ga helper was only applied to client- originated session.update messages in client_ack_messages. Internally- generated session.updates (sent via _send_to_backend) in two paths: - _handle_raw_backend_message (raw/no provider_config path, line 518) - backend_to_client_send_messages provider_config path (line 481) bypassed the remap, so GA upstreams ignored or rejected them, breaking audio transcription guardrails for all non-beta clients. Fix: add _make_disable_auto_response_message() helper that always emits the correct GA-shaped session.update, and replace both injection sites with it. Update existing tests to assert the GA nested shape instead of the old flat beta shape, and add a new unit test for the helper itself. Co-authored-by: Sameer Kankute * Log realtime session type * Fix beta realtime session payloads * Fix realtime audio format remapping edge case * Fix Azure realtime beta session shape --------- Co-authored-by: Cursor Co-authored-by: Sameer Kankute Co-authored-by: mateo-berri <277851410+mateo-berri@users.noreply.github.com> --- .../litellm_core_utils/realtime_streaming.py | 308 ++++++++++++++++-- litellm/llms/azure/realtime/handler.py | 6 +- litellm/llms/custom_httpx/llm_http_handler.py | 1 - litellm/llms/openai/realtime/handler.py | 46 ++- litellm/llms/xai/realtime/handler.py | 7 +- litellm/types/llms/openai.py | 83 ++++- .../test_realtime_streaming.py | 192 ++++++++++- .../realtime/test_azure_realtime_handler.py | 8 + .../realtime/test_openai_realtime_handler.py | 63 +++- .../playground/chat_ui/RealtimePlayground.tsx | 26 +- 10 files changed, 669 insertions(+), 71 deletions(-) diff --git a/litellm/litellm_core_utils/realtime_streaming.py b/litellm/litellm_core_utils/realtime_streaming.py index 4493a58f78b..c4528ff74e3 100644 --- a/litellm/litellm_core_utils/realtime_streaming.py +++ b/litellm/litellm_core_utils/realtime_streaming.py @@ -31,6 +31,8 @@ DefaultLoggedRealTimeEventTypes = [ "session.created", "response.create", "response.done", + "conversation.item.added", # GA + "conversation.item.done", # GA ] @@ -44,6 +46,7 @@ class RealTimeStreaming: model: str = "", user_api_key_dict: Optional[Any] = None, request_data: Optional[Dict] = None, + backend_uses_beta_protocol: Optional[bool] = None, ): self.websocket = websocket self.backend_ws = backend_ws @@ -54,6 +57,14 @@ class RealTimeStreaming: self.session_tools: List[Dict] = [] self.tool_calls: List[Dict] = [] + # Detect whether the client is explicitly opting into the beta protocol. + self._client_wants_beta = self._detect_beta_header(websocket) + self._backend_uses_beta_protocol = ( + self._client_wants_beta + if backend_uses_beta_protocol is None + else backend_uses_beta_protocol + ) + _logged_real_time_event_types = litellm.logged_real_time_event_types if _logged_real_time_event_types is None: @@ -76,6 +87,27 @@ class RealTimeStreaming: # response.create can be rewritten to include the failure context. self._pending_guardrail_message: Optional[str] = None + _SESSION_EVENT_TYPES = frozenset(["session.created", "session.updated"]) + _AUDIO_FORMAT_MAP: Dict[str, Dict[str, Any]] = { + "pcm16": {"type": "audio/pcm", "rate": 24000}, + "g711_ulaw": {"type": "audio/G711-ulaw", "rate": 8000}, + "g711_alaw": {"type": "audio/G711-alaw", "rate": 8000}, + } + # GA name → beta name (when client WebSocket includes OpenAI-Beta: realtime=v1) + _GA_TO_BETA_EVENT_TYPES: Dict[str, str] = { + "conversation.item.added": "conversation.item.created", + "response.output_text.delta": "response.text.delta", + "response.output_audio.delta": "response.audio.delta", + "response.output_audio_transcript.delta": "response.audio_transcript.delta", + "response.output_text.done": "response.text.done", + "response.output_audio.done": "response.audio.done", + "response.output_audio_transcript.done": "response.audio_transcript.done", + } + _GA_TO_BETA_CONTENT_TYPES: Dict[str, str] = { + "output_text": "text", + "output_audio": "audio", + } + def _should_store_message( self, message_obj: Union[dict, OpenAIRealtimeEvents], @@ -92,24 +124,27 @@ class RealTimeStreaming: if isinstance(message, bytes): message = message.decode("utf-8") if isinstance(message, dict): - message_obj = message + # TypedDict union members do not narrow to plain dict for mypy. + message_obj: Dict[str, Any] = cast(Dict[str, Any], message) else: - message_obj = json.loads(message) + message_obj = cast(Dict[str, Any], json.loads(cast(str, message))) self._collect_tool_calls_from_response_done(cast(dict, message_obj)) try: - if ( - not isinstance(message, dict) - or message_obj.get("type") == "session.created" - or message_obj.get("type") == "session.updated" - ): - message_obj = OpenAIRealtimeStreamSessionEvents(**message_obj) # type: ignore - elif not isinstance(message, dict): - message_obj = OpenAIRealtimeStreamResponseBaseObject(**message_obj) # type: ignore + event_type = message_obj.get("type", "") + if event_type in self._SESSION_EVENT_TYPES: + typed_obj = OpenAIRealtimeStreamSessionEvents(**message_obj) # type: ignore + else: + # Use the base object as a safe catch-all for all other event types + # (both beta and GA), so unknown/new event names never raise here. + typed_obj = OpenAIRealtimeStreamResponseBaseObject(**message_obj) # type: ignore except Exception as e: verbose_logger.debug(f"Error parsing message for logging: {e}") - raise e - if self._should_store_message(message_obj): - self.messages.append(message_obj) + # Don't re-raise — a parse failure must not drop or delay the message + if self._should_store_message(message_obj): + self.messages.append(message_obj) # type: ignore[arg-type] + return + if self._should_store_message(typed_obj): + self.messages.append(typed_obj) def _collect_user_input_from_client_event(self, message: Union[str, dict]) -> None: """Extract user text content from client WebSocket events for spend logging.""" @@ -147,6 +182,8 @@ class RealTimeStreaming: tools = session.get("tools") if tools and isinstance(tools, list): self.session_tools = tools + # GA: session.type is required; log it for traceability but no action needed + verbose_logger.debug(f"Realtime session.type: {session.get('type')}") except (json.JSONDecodeError, AttributeError, TypeError): pass @@ -228,6 +265,23 @@ class RealTimeStreaming: else: await self.backend_ws.send(message) # type: ignore[union-attr, attr-defined] + def _make_disable_auto_response_message(self) -> str: + """Return a session.update that disables VAD auto-response.""" + if self._backend_uses_beta_protocol: + session: Dict[str, Any] = { + "turn_detection": {"create_response": False}, + } + else: + session = { + "type": "realtime", + "audio": { + "input": { + "turn_detection": {"create_response": False}, + } + }, + } + return json.dumps({"type": "session.update", "session": session}) + def _has_realtime_guardrails(self) -> bool: """Return True if any callback is registered for realtime guardrail event types.""" from litellm.integrations.custom_guardrail import CustomGuardrail @@ -435,14 +489,7 @@ class RealTimeStreaming: ): self.store_message(event_str) await self.websocket.send_text(event_str) - await self._send_to_backend( - json.dumps( - { - "type": "session.update", - "session": {"turn_detection": {"create_response": False}}, - } - ) - ) + await self._send_to_backend(self._make_disable_auto_response_message()) continue ## GUARDRAIL: run on transcription events in provider_config path too if ( @@ -484,14 +531,7 @@ class RealTimeStreaming: ): self.store_message(raw_response) await self.websocket.send_text(raw_response) - await self._send_to_backend( - json.dumps( - { - "type": "session.update", - "session": {"turn_detection": {"create_response": False}}, - } - ) - ) + await self._send_to_backend(self._make_disable_auto_response_message()) return True if ( @@ -542,7 +582,20 @@ class RealTimeStreaming: continue ## LOGGING self.store_message(raw_response) - await self.websocket.send_text(raw_response) + + # If the client opted into beta protocol, translate GA event + # names/shapes back to the beta equivalents before forwarding. + if self._client_wants_beta: + try: + event_dict = json.loads(raw_response) + translated = self._translate_event_to_beta(event_dict) + if translated is None: + continue # drop GA-only events (e.g. conversation.item.done) + await self.websocket.send_text(json.dumps(translated)) + except Exception: + await self.websocket.send_text(raw_response) + else: + await self.websocket.send_text(raw_response) except websockets.exceptions.ConnectionClosed as e: # type: ignore verbose_logger.exception( @@ -553,6 +606,183 @@ class RealTimeStreaming: finally: await self.log_messages() + @staticmethod + def _detect_beta_header(websocket: Any) -> bool: + """Return True if the client sent 'OpenAI-Beta: realtime=v1'. + + Checks the raw ASGI scope headers so it works for both FastAPI WebSocket + objects and any test doubles that expose a .scope dict. + """ + try: + headers = websocket.scope.get("headers", []) + for name, value in headers: + if isinstance(name, bytes): + name = name.decode("latin-1") + if isinstance(value, bytes): + value = value.decode("latin-1") + if name.lower() == "openai-beta" and "realtime=v1" in value.lower(): + return True + except Exception: + pass + return False + + @staticmethod + def _remap_beta_session_to_ga(session: dict) -> dict: + """ + Convert a beta-style session.update payload to the GA nested schema. + + Beta → GA field mappings + ───────────────────────────────────────────────────────────────────── + session.type (inject "realtime" if absent) + session.modalities → session.output_modalities + session.voice → session.audio.output.voice + session.input_audio_format → session.audio.input.format (with type/rate) + session.output_audio_format → session.audio.output.format (with type/rate) + session.turn_detection → session.audio.input.turn_detection + session.input_audio_transcription → session.audio.input.transcription + ───────────────────────────────────────────────────────────────────── + Fields not in the mapping (instructions, tools, etc.) are passed through. + GA clients that already use the nested shape are unaffected. + """ + # Work on a shallow copy so we don't mutate the caller's dict + session = dict(session) + + # 1. Ensure session.type is present + if "type" not in session: + session["type"] = "realtime" + + # 2. Rename modalities → output_modalities and normalise combinations. + # Beta allowed ["audio", "text"] together; GA only supports ["audio"] or + # ["text"] as single-element lists. When both are present we prefer + # ["audio"] because audio mode already delivers transcripts via events. + if "modalities" in session: + mods = session.pop("modalities") + if "output_modalities" not in session: + mods_set = {m.lower() for m in (mods or [])} + if "audio" in mods_set: + session["output_modalities"] = ["audio"] + elif "text" in mods_set: + session["output_modalities"] = ["text"] + + # 3-7. Lift flat audio fields into the nested audio object + audio: Dict[str, Any] = {} + inp: Dict[str, Any] = {} + out: Dict[str, Any] = {} + + # voice → audio.output.voice + if "voice" in session: + out["voice"] = session.pop("voice") + + # input_audio_format → audio.input.format + if "input_audio_format" in session: + raw = session.pop("input_audio_format") + inp["format"] = ( + RealTimeStreaming._AUDIO_FORMAT_MAP.get(raw, raw) + if isinstance(raw, str) + else raw + ) + + # output_audio_format → audio.output.format + if "output_audio_format" in session: + raw = session.pop("output_audio_format") + out["format"] = ( + RealTimeStreaming._AUDIO_FORMAT_MAP.get(raw, raw) + if isinstance(raw, str) + else raw + ) + + # turn_detection → audio.input.turn_detection + if "turn_detection" in session: + inp["turn_detection"] = session.pop("turn_detection") + + # input_audio_transcription → audio.input.transcription + if "input_audio_transcription" in session: + inp["transcription"] = session.pop("input_audio_transcription") + + if inp: + audio["input"] = inp + if out: + audio["output"] = out + + if audio: + # Merge with any existing GA-style `audio` block the client already set, + # letting the remapped values take precedence within each sub-key. + existing = session.get("audio") or {} + for sub_key, sub_val in audio.items(): + if ( + sub_key in existing + and isinstance(existing[sub_key], dict) + and isinstance(sub_val, dict) + ): + existing[sub_key] = {**existing[sub_key], **sub_val} + else: + existing[sub_key] = sub_val + session["audio"] = existing + + return session + + @staticmethod + def _translate_event_to_beta(event: dict) -> Optional[dict]: + """Translate a single GA event dict to its beta equivalent. + + Returns None if the event should be dropped entirely (e.g. the GA-only + conversation.item.done has no beta counterpart). + Returns the (possibly mutated copy of the) event otherwise. + """ + event_type = event.get("type", "") + + # conversation.item.done has no beta equivalent — the client already + # received conversation.item.created (translated from .added). + if event_type == "conversation.item.done": + return None + + # Shallow-copy so we don't mutate the stored message + translated = dict(event) + + # Rename the type field + if event_type in RealTimeStreaming._GA_TO_BETA_EVENT_TYPES: + translated["type"] = RealTimeStreaming._GA_TO_BETA_EVENT_TYPES[event_type] + + # Fix content block types inside items (response.done output list, + # conversation.item.created item content, etc.) + if "item" in translated and isinstance(translated["item"], dict): + translated["item"] = RealTimeStreaming._translate_item_content_types( + dict(translated["item"]) + ) + if "response" in translated and isinstance(translated["response"], dict): + resp = dict(translated["response"]) + if "output" in resp and isinstance(resp["output"], list): + resp["output"] = [ + ( + RealTimeStreaming._translate_item_content_types(dict(o)) + if isinstance(o, dict) + else o + ) + for o in resp["output"] + ] + translated["response"] = resp + + return translated + + @staticmethod + def _translate_item_content_types(item: dict) -> dict: + """Replace GA content type names with beta names inside a single item.""" + if "content" not in item or not isinstance(item["content"], list): + return item + new_content = [] + for block in item["content"]: + if ( + isinstance(block, dict) + and block.get("type") in RealTimeStreaming._GA_TO_BETA_CONTENT_TYPES + ): + block = dict(block) + block["type"] = RealTimeStreaming._GA_TO_BETA_CONTENT_TYPES[ + block["type"] + ] + new_content.append(block) + item["content"] = new_content + return item + async def client_ack_messages(self): try: while True: @@ -594,6 +824,19 @@ class RealTimeStreaming: self._pending_guardrail_message = None continue + # GA compatibility: remap beta-style session fields only when + # the upstream is in GA mode. Beta upstreams expect the flat + # session shape unchanged. + if ( + msg_type == "session.update" + and not self._backend_uses_beta_protocol + ): + session = msg_obj.get("session", {}) + if isinstance(session, dict): + session = self._remap_beta_session_to_ga(session) + msg_obj["session"] = session + message = json.dumps(msg_obj) + except (json.JSONDecodeError, AttributeError): pass @@ -627,3 +870,8 @@ class RealTimeStreaming: await forward_task except asyncio.CancelledError: pass + + +def client_sent_openai_beta_realtime_header(websocket: Any) -> bool: + """True when the client WebSocket includes ``OpenAI-Beta: realtime=v1``.""" + return RealTimeStreaming._detect_beta_header(websocket) diff --git a/litellm/llms/azure/realtime/handler.py b/litellm/llms/azure/realtime/handler.py index 1f3428f2ca5..1f3357fd788 100644 --- a/litellm/llms/azure/realtime/handler.py +++ b/litellm/llms/azure/realtime/handler.py @@ -89,9 +89,10 @@ class AzureOpenAIRealtime(AzureChatCompletion): if api_base is None: raise ValueError("api_base is required for Azure OpenAI calls") - if api_version is None and ( + backend_uses_beta_protocol = ( realtime_protocol is None or realtime_protocol.upper() not in ("GA", "V1") - ): + ) + if api_version is None and backend_uses_beta_protocol: raise ValueError("api_version is required for Azure OpenAI calls") url = self._construct_url( @@ -114,6 +115,7 @@ class AzureOpenAIRealtime(AzureChatCompletion): logging_obj, user_api_key_dict=user_api_key_dict, request_data={"litellm_metadata": litellm_metadata or {}}, + backend_uses_beta_protocol=backend_uses_beta_protocol, ) await realtime_streaming.bidirectional_forward() diff --git a/litellm/llms/custom_httpx/llm_http_handler.py b/litellm/llms/custom_httpx/llm_http_handler.py index 2ffc7acbfb1..fa1253d9005 100644 --- a/litellm/llms/custom_httpx/llm_http_handler.py +++ b/litellm/llms/custom_httpx/llm_http_handler.py @@ -5255,7 +5255,6 @@ class BaseLLMHTTPHandler: headers = { "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", - "OpenAI-Beta": "realtime=v1", } if extra_headers: diff --git a/litellm/llms/openai/realtime/handler.py b/litellm/llms/openai/realtime/handler.py index 5ca0a3186f7..f34dae2df09 100644 --- a/litellm/llms/openai/realtime/handler.py +++ b/litellm/llms/openai/realtime/handler.py @@ -6,12 +6,15 @@ This requires websockets, and is currently only supported on LiteLLM Proxy. from typing import Any, Optional, cast -from litellm._logging import _redact_string +from litellm._logging import _redact_string, verbose_logger from litellm.constants import REALTIME_WEBSOCKET_MAX_MESSAGE_SIZE_BYTES from litellm.types.realtime import RealtimeQueryParams from ....litellm_core_utils.litellm_logging import Logging as LiteLLMLogging -from ....litellm_core_utils.realtime_streaming import RealTimeStreaming +from ....litellm_core_utils.realtime_streaming import ( + RealTimeStreaming, + client_sent_openai_beta_realtime_header, +) from ....llms.custom_httpx.http_handler import get_shared_realtime_ssl_context from ..openai import OpenAIChatCompletion @@ -33,21 +36,24 @@ class OpenAIRealtime(OpenAIChatCompletion): """ return "https://api.openai.com/" - def _get_additional_headers(self, api_key: str) -> dict: + def _get_additional_headers( + self, + api_key: str, + *, + openai_beta_realtime: bool = False, + ) -> dict: """ - Get additional headers beyond Authorization. - Override this in subclasses to customize headers (e.g., remove OpenAI-Beta). + Headers for the upstream OpenAI Realtime WebSocket. - Args: - api_key: API key for authentication - - Returns: - Dictionary of additional headers + When the client sent ``OpenAI-Beta: realtime=v1`` on the proxy WebSocket, + ``openai_beta_realtime`` is True and the same header is forwarded upstream + so the legacy beta API is used. GA clients omit that header on the client + connection and must send GA-shaped ``session.update`` payloads. """ - return { - "Authorization": f"Bearer {api_key}", - "OpenAI-Beta": "realtime=v1", - } + headers: dict = {"Authorization": f"Bearer {api_key}"} + if openai_beta_realtime: + headers["OpenAI-Beta"] = "realtime=v1" + return headers def _get_ssl_config(self, url: str) -> Any: """ @@ -120,8 +126,16 @@ class OpenAIRealtime(OpenAIChatCompletion): # Get provider-specific SSL configuration ssl_config = self._get_ssl_config(url) - # Get provider-specific headers - headers = self._get_additional_headers(api_key) + openai_beta_realtime = client_sent_openai_beta_realtime_header(websocket) + if not openai_beta_realtime: + verbose_logger.debug( + "OpenAI Realtime: connecting with GA protocol (no OpenAI-Beta header). " + "If your client expects beta event names, add 'OpenAI-Beta: realtime=v1' " + "to the WebSocket headers sent to the LiteLLM proxy." + ) + headers = self._get_additional_headers( + api_key, openai_beta_realtime=openai_beta_realtime + ) # Log a masked request preview consistent with other endpoints. logging_obj.pre_call( diff --git a/litellm/llms/xai/realtime/handler.py b/litellm/llms/xai/realtime/handler.py index 805cce5a264..eab19f4a6c8 100644 --- a/litellm/llms/xai/realtime/handler.py +++ b/litellm/llms/xai/realtime/handler.py @@ -28,7 +28,12 @@ class XAIRealtime(OpenAIRealtime): """xAI uses a different API base URL.""" return XAI_API_BASE - def _get_additional_headers(self, api_key: str) -> dict: + def _get_additional_headers( + self, + api_key: str, + *, + openai_beta_realtime: bool = False, + ) -> dict: """ xAI does NOT require the OpenAI-Beta header. Only send Authorization header. diff --git a/litellm/types/llms/openai.py b/litellm/types/llms/openai.py index 986ec39f3bb..abe58199dfd 100644 --- a/litellm/types/llms/openai.py +++ b/litellm/types/llms/openai.py @@ -1878,14 +1878,24 @@ class OpenAIRealtimeStreamSessionEvents(TypedDict): class OpenAIRealtimeStreamResponseOutputItemContent(TypedDict, total=False): audio: str - """Base64-encoded audio bytes, used for 'input_audio' content types""" + """Base64-encoded audio bytes, used for 'input_audio' / 'audio' / 'output_audio' content types""" id: str """The ID of the previous conversation item for reference""" text: str - """The text content, used for 'input_text' and 'text' content types""" + """The text content, used for 'input_text' / 'text' / 'output_text' content types""" transcript: str - """The transcript content, used for 'input_audio' content types""" - type: Literal["input_audio", "input_text", "text", "item_reference", "audio"] + """The transcript content, used for 'input_audio' / 'audio' content types""" + type: Literal[ + "input_audio", + "input_text", + # Beta assistant content types + "text", + "audio", + "item_reference", + # GA assistant content types (aligns with Responses API) + "output_text", + "output_audio", + ] """The type of content""" @@ -1945,23 +1955,46 @@ class OpenAIRealtimeConversationCreated(TypedDict, total=False): class OpenAIRealtimeConversationItemCreated(TypedDict, total=False): + """Beta: single event emitted when a conversation item is created.""" + type: Required[Literal["conversation.item.created"]] item: OpenAIRealtimeStreamResponseOutputItem event_id: str - previous_item_id: str + previous_item_id: Optional[str] # None when this is the first item + + +class OpenAIRealtimeConversationItemAdded(TypedDict, total=False): + """GA: emitted immediately when a conversation item is added (replaces .created).""" + + type: Required[Literal["conversation.item.added"]] + item: OpenAIRealtimeStreamResponseOutputItem + event_id: str + previous_item_id: Optional[str] # None when this is the first item + + +class OpenAIRealtimeConversationItemDone(TypedDict, total=False): + """GA: emitted when a conversation item is fully complete (e.g. transcription done).""" + + type: Required[Literal["conversation.item.done"]] + item: OpenAIRealtimeStreamResponseOutputItem + event_id: str + previous_item_id: Optional[str] # None when this is the first item class OpenAIRealtimeResponseContentPart(TypedDict, total=False): audio: str - """Base64-encoded audio bytes, if type is 'audio'""" + """Base64-encoded audio bytes, if type is 'audio' or 'output_audio'""" text: str - """The text content, if type is 'text'""" + """The text content, if type is 'text' or 'output_text'""" transcript: str - """The transcript content, if type is 'audio'""" + """The transcript content, if type is 'audio' or 'output_audio'""" - type: Literal["audio", "text"] + type: Union[ + Literal["audio", "text"], # beta + Literal["output_audio", "output_text"], # GA + ] """The type of content""" @@ -1982,7 +2015,14 @@ class OpenAIRealtimeResponseDelta(TypedDict): item_id: str output_index: int response_id: str - type: Union[Literal["response.text.delta"], Literal["response.audio.delta"]] + type: Union[ + Literal["response.text.delta"], + Literal["response.audio.delta"], + # GA renamed events + Literal["response.output_text.delta"], + Literal["response.output_audio.delta"], + Literal["response.output_audio_transcript.delta"], + ] class OpenAIRealtimeResponseTextDone(TypedDict): @@ -1992,7 +2032,10 @@ class OpenAIRealtimeResponseTextDone(TypedDict): output_index: int response_id: str text: str - type: Literal["response.text.done"] + type: Union[ + Literal["response.text.done"], + Literal["response.output_text.done"], # GA rename + ] class OpenAIRealtimeResponseAudioDone(TypedDict): @@ -2001,7 +2044,11 @@ class OpenAIRealtimeResponseAudioDone(TypedDict): item_id: str output_index: int response_id: str - type: Literal["response.audio.done"] + type: Union[ + Literal["response.audio.done"], + Literal["response.output_audio.done"], # GA rename + Literal["response.output_audio_transcript.done"], # GA rename + ] class OpenAIRealtimeContentPartDone(TypedDict): @@ -2046,10 +2093,18 @@ class OpenAIRealtimeDoneEvent(TypedDict): class OpenAIRealtimeEventTypes(Enum): SESSION_CREATED = "session.created" + # Beta delta event names RESPONSE_TEXT_DELTA = "response.text.delta" RESPONSE_AUDIO_DELTA = "response.audio.delta" RESPONSE_TEXT_DONE = "response.text.done" RESPONSE_AUDIO_DONE = "response.audio.done" + # GA renamed delta event names + RESPONSE_OUTPUT_TEXT_DELTA = "response.output_text.delta" + RESPONSE_OUTPUT_AUDIO_DELTA = "response.output_audio.delta" + RESPONSE_OUTPUT_AUDIO_TRANSCRIPT_DELTA = "response.output_audio_transcript.delta" + RESPONSE_OUTPUT_TEXT_DONE = "response.output_text.done" + RESPONSE_OUTPUT_AUDIO_DONE = "response.output_audio.done" + RESPONSE_OUTPUT_AUDIO_TRANSCRIPT_DONE = "response.output_audio_transcript.done" RESPONSE_DONE = "response.done" RESPONSE_OUTPUT_ITEM_ADDED = "response.output_item.added" RESPONSE_CONTENT_PART_ADDED = "response.content_part.added" @@ -2060,7 +2115,11 @@ OpenAIRealtimeEvents = Union[ OpenAIRealtimeStreamSessionEvents, OpenAIRealtimeStreamResponseOutputItemAdded, OpenAIRealtimeResponseContentPartAdded, + # Beta conversation item event OpenAIRealtimeConversationItemCreated, + # GA conversation item events + OpenAIRealtimeConversationItemAdded, + OpenAIRealtimeConversationItemDone, OpenAIRealtimeConversationCreated, OpenAIRealtimeResponseDelta, OpenAIRealtimeResponseTextDone, diff --git a/tests/test_litellm/litellm_core_utils/test_realtime_streaming.py b/tests/test_litellm/litellm_core_utils/test_realtime_streaming.py index 8d842fefb7b..2b238b0cdf7 100644 --- a/tests/test_litellm/litellm_core_utils/test_realtime_streaming.py +++ b/tests/test_litellm/litellm_core_utils/test_realtime_streaming.py @@ -13,7 +13,10 @@ sys.path.insert( ) # Adds the parent directory to the system path from litellm.integrations.custom_guardrail import CustomGuardrail -from litellm.litellm_core_utils.realtime_streaming import RealTimeStreaming +from litellm.litellm_core_utils.realtime_streaming import ( + RealTimeStreaming, + client_sent_openai_beta_realtime_header, +) from litellm.types.guardrails import GuardrailEventHooks from litellm.types.llms.openai import ( OpenAIRealtimeStreamResponseBaseObject, @@ -80,6 +83,175 @@ def test_realtime_streaming_store_message(): assert len(streaming.messages) == 2 # Should not store the new message +def test_remap_beta_session_to_ga_normalizes_modalities_and_audio(): + out = RealTimeStreaming._remap_beta_session_to_ga( + {"modalities": ["audio", "text"], "voice": "alloy"} + ) + assert out["type"] == "realtime" + assert out["output_modalities"] == ["audio"] + assert out["audio"]["output"]["voice"] == "alloy" + + +def test_remap_beta_session_to_ga_preserves_ga_audio_format_dicts(): + input_format = {"type": "audio/pcm", "rate": 24000} + output_format = {"type": "audio/G711-ulaw", "rate": 8000} + + out = RealTimeStreaming._remap_beta_session_to_ga( + { + "input_audio_format": input_format, + "output_audio_format": output_format, + } + ) + + assert out["audio"]["input"]["format"] == input_format + assert out["audio"]["output"]["format"] == output_format + + +def test_make_disable_auto_response_message_produces_ga_shape(): + """_make_disable_auto_response_message must produce a GA-shaped session.update. + + The GA Realtime API requires: + - session.type = "realtime" + - turn_detection nested at session.audio.input.turn_detection + The old beta-style flat ``session.turn_detection`` is rejected by GA upstreams. + """ + websocket = MagicMock() + backend_ws = MagicMock() + logging_obj = MagicMock() + streaming = RealTimeStreaming(websocket, backend_ws, logging_obj) + + raw = streaming._make_disable_auto_response_message() + msg = json.loads(raw) + + assert msg["type"] == "session.update" + session = msg["session"] + assert ( + session.get("type") == "realtime" + ), "GA session.update must include session.type='realtime'" + # turn_detection must NOT be at the flat beta location + assert ( + "turn_detection" not in session + ), "turn_detection must not be at the top-level session (beta shape); use audio.input" + # turn_detection must be nested under audio.input + assert session["audio"]["input"]["turn_detection"]["create_response"] is False + + +def test_make_disable_auto_response_message_produces_beta_shape_for_beta_clients(): + websocket = MagicMock() + websocket.scope = {"headers": [(b"openai-beta", b"realtime=v1")]} + backend_ws = MagicMock() + logging_obj = MagicMock() + streaming = RealTimeStreaming(websocket, backend_ws, logging_obj) + + raw = streaming._make_disable_auto_response_message() + msg = json.loads(raw) + + assert msg["type"] == "session.update" + session = msg["session"] + assert session == {"turn_detection": {"create_response": False}} + + +@pytest.mark.asyncio +async def test_client_ack_messages_keeps_beta_session_shape_for_beta_clients(): + client_ws = MagicMock() + client_ws.scope = {"headers": [(b"openai-beta", b"realtime=v1")]} + session_update = json.dumps( + { + "type": "session.update", + "session": { + "modalities": ["audio", "text"], + "voice": "alloy", + "turn_detection": {"create_response": False}, + }, + } + ) + client_ws.receive_text = AsyncMock( + side_effect=[ + session_update, + Exception("connection closed"), + ] + ) + backend_ws = MagicMock() + backend_ws.send = AsyncMock() + logging_obj = MagicMock() + logging_obj.pre_call = MagicMock() + streaming = RealTimeStreaming(client_ws, backend_ws, logging_obj) + + await streaming.client_ack_messages() + + sent_to_backend = json.loads(backend_ws.send.call_args_list[0].args[0]) + session = sent_to_backend["session"] + assert session["modalities"] == ["audio", "text"] + assert session["voice"] == "alloy" + assert session["turn_detection"] == {"create_response": False} + assert "type" not in session + assert "output_modalities" not in session + assert "audio" not in session + + +@pytest.mark.asyncio +async def test_client_ack_messages_keeps_beta_session_shape_for_beta_backend(): + client_ws = MagicMock() + session_update = json.dumps( + { + "type": "session.update", + "session": { + "modalities": ["audio", "text"], + "voice": "alloy", + "turn_detection": {"create_response": False}, + }, + } + ) + client_ws.receive_text = AsyncMock( + side_effect=[ + session_update, + Exception("connection closed"), + ] + ) + backend_ws = MagicMock() + backend_ws.send = AsyncMock() + logging_obj = MagicMock() + logging_obj.pre_call = MagicMock() + streaming = RealTimeStreaming( + client_ws, backend_ws, logging_obj, backend_uses_beta_protocol=True + ) + + await streaming.client_ack_messages() + + sent_to_backend = json.loads(backend_ws.send.call_args_list[0].args[0]) + session = sent_to_backend["session"] + assert session["modalities"] == ["audio", "text"] + assert session["voice"] == "alloy" + assert session["turn_detection"] == {"create_response": False} + assert "type" not in session + assert "output_modalities" not in session + assert "audio" not in session + + +def test_translate_event_to_beta_renames_delta_types(): + ev = RealTimeStreaming._translate_event_to_beta( + {"type": "response.output_audio.delta", "delta": "abc", "event_id": "e1"} + ) + assert ev is not None + assert ev["type"] == "response.audio.delta" + + +def test_translate_event_to_beta_drops_conversation_item_done(): + assert ( + RealTimeStreaming._translate_event_to_beta({"type": "conversation.item.done"}) + is None + ) + + +def test_client_sent_openai_beta_realtime_header_detects_header(): + ws = MagicMock() + ws.scope = {"headers": [(b"openai-beta", b"realtime=v1")]} + assert client_sent_openai_beta_realtime_header(ws) is True + empty = MagicMock() + empty.scope = {"headers": []} + assert client_sent_openai_beta_realtime_header(empty) is False + + def test_collect_user_input_from_text_conversation_item(): """ Test that conversation.item.create with input_text content is collected as user input. @@ -761,7 +933,14 @@ async def test_realtime_session_created_injects_session_update_for_audio_guardra assert ( len(session_updates) == 1 ), f"Expected one session.update injected to backend, got: {sent_to_backend}" - assert session_updates[0]["session"]["turn_detection"]["create_response"] is False + # GA shape: turn_detection must be nested under audio.input, not at top-level session + injected_session = session_updates[0]["session"] + assert ( + injected_session["type"] == "realtime" + ), "GA session.update must include session.type='realtime'" + assert ( + injected_session["audio"]["input"]["turn_detection"]["create_response"] is False + ), "GA session.update must nest turn_detection under audio.input" litellm.callbacks = [] # cleanup @@ -818,7 +997,14 @@ async def test_realtime_session_created_injects_session_update_for_pre_call_guar assert ( len(session_updates) == 1 ), f"pre_call guardrail should inject session.update to gate audio responses, got: {sent_to_backend}" - assert session_updates[0]["session"]["turn_detection"]["create_response"] is False + # GA shape: turn_detection must be nested under audio.input, not at top-level session + injected_session = session_updates[0]["session"] + assert ( + injected_session["type"] == "realtime" + ), "GA session.update must include session.type='realtime'" + assert ( + injected_session["audio"]["input"]["turn_detection"]["create_response"] is False + ), "GA session.update must nest turn_detection under audio.input" litellm.callbacks = [] # cleanup diff --git a/tests/test_litellm/llms/azure/realtime/test_azure_realtime_handler.py b/tests/test_litellm/llms/azure/realtime/test_azure_realtime_handler.py index 42108e46b59..41d301c5d5f 100644 --- a/tests/test_litellm/llms/azure/realtime/test_azure_realtime_handler.py +++ b/tests/test_litellm/llms/azure/realtime/test_azure_realtime_handler.py @@ -247,6 +247,10 @@ async def test_async_realtime_uses_ga_protocol_end_to_end(): assert "model=gpt-4o-realtime-preview" in called_url assert "api-version" not in called_url assert "deployment" not in called_url + assert ( + mock_realtime_streaming.call_args.kwargs["backend_uses_beta_protocol"] + is False + ) @pytest.mark.asyncio @@ -419,3 +423,7 @@ async def test_async_realtime_default_maintains_backwards_compatibility(): called_url = mock_ws_connect.call_args[0][0] assert "/openai/realtime?" in called_url assert "/openai/v1/realtime" not in called_url + assert ( + mock_realtime_streaming.call_args.kwargs["backend_uses_beta_protocol"] + is True + ) diff --git a/tests/test_litellm/llms/openai/realtime/test_openai_realtime_handler.py b/tests/test_litellm/llms/openai/realtime/test_openai_realtime_handler.py index 4f5764e3d6e..e9798f45dce 100644 --- a/tests/test_litellm/llms/openai/realtime/test_openai_realtime_handler.py +++ b/tests/test_litellm/llms/openai/realtime/test_openai_realtime_handler.py @@ -213,12 +213,12 @@ async def test_async_realtime_url_contains_model(): assert called_url.startswith("wss://api.openai.com/v1/realtime?") assert f"model={model}" in called_url - # Verify proper headers were set + # Verify proper headers were set (GA default: no OpenAI-Beta unless client sent it) called_kwargs = mock_ws_connect.call_args[1] assert "additional_headers" in called_kwargs additional_headers = called_kwargs["additional_headers"] assert additional_headers["Authorization"] == f"Bearer {api_key}" - assert additional_headers["OpenAI-Beta"] == "realtime=v1" + assert "OpenAI-Beta" not in additional_headers # Verify SSL is configured (should be an SSLContext or True, not None or False) assert called_kwargs["ssl"] is not None assert called_kwargs["ssl"] is not False @@ -227,6 +227,65 @@ async def test_async_realtime_url_contains_model(): mock_streaming_instance.bidirectional_forward.assert_awaited_once() +@pytest.mark.asyncio +async def test_async_realtime_forwards_openai_beta_header_when_client_sends_it(): + """Upstream WS gets OpenAI-Beta: realtime=v1 only when the client WebSocket included it.""" + from litellm.llms.openai.realtime.handler import OpenAIRealtime + from litellm.types.realtime import RealtimeQueryParams + + handler = OpenAIRealtime() + api_base = "https://api.openai.com/" + api_key = "test-key" + model = "gpt-4o-mini-realtime-preview" + query_params: RealtimeQueryParams = {"model": model} + + dummy_websocket = MagicMock() + dummy_websocket.scope = { + "headers": [ + (b"openai-beta", b"realtime=v1"), + ] + } + dummy_logging_obj = MagicMock() + mock_backend_ws = AsyncMock() + + class DummyAsyncContextManager: + def __init__(self, value): + self.value = value + + async def __aenter__(self): + return self.value + + async def __aexit__(self, exc_type, exc, tb): + return None + + with ( + patch( + "websockets.connect", return_value=DummyAsyncContextManager(mock_backend_ws) + ) as mock_ws_connect, + patch( + "litellm.llms.openai.realtime.handler.RealTimeStreaming" + ) as mock_realtime_streaming, + ): + mock_streaming_instance = MagicMock() + mock_realtime_streaming.return_value = mock_streaming_instance + mock_streaming_instance.bidirectional_forward = AsyncMock() + + await handler.async_realtime( + model=model, + websocket=dummy_websocket, + logging_obj=dummy_logging_obj, + api_base=api_base, + api_key=api_key, + query_params=query_params, + ) + + mock_ws_connect.assert_called_once() + called_kwargs = mock_ws_connect.call_args[1] + additional_headers = called_kwargs["additional_headers"] + assert additional_headers["Authorization"] == f"Bearer {api_key}" + assert additional_headers["OpenAI-Beta"] == "realtime=v1" + + @pytest.mark.asyncio async def test_async_realtime_uses_max_size_parameter(): """ diff --git a/ui/litellm-dashboard/src/components/playground/chat_ui/RealtimePlayground.tsx b/ui/litellm-dashboard/src/components/playground/chat_ui/RealtimePlayground.tsx index a3dd864b894..682ea62cbe9 100644 --- a/ui/litellm-dashboard/src/components/playground/chat_ui/RealtimePlayground.tsx +++ b/ui/litellm-dashboard/src/components/playground/chat_ui/RealtimePlayground.tsx @@ -134,10 +134,12 @@ const RealtimePlayground: React.FC = ({ const type = data.type; if (type === "session.created") { + // GA: session.type is required ("realtime" | "transcription") ws.send( JSON.stringify({ type: "session.update", session: { + type: "realtime", modalities: ["text", "audio"], voice: selectedVoice, input_audio_format: "pcm16", @@ -149,24 +151,36 @@ const RealtimePlayground: React.FC = ({ ); } else if (type === "session.updated") { // session configured - } else if (type === "response.audio.delta") { + } else if ( + // GA: response.output_audio.delta | beta: response.audio.delta + type === "response.output_audio.delta" || type === "response.audio.delta" + ) { if (data.delta) playAudioChunk(data.delta); - } else if (type === "response.audio_transcript.delta" || type === "response.text.delta") { + } else if ( + // GA: response.output_text.delta / response.output_audio_transcript.delta + // beta: response.text.delta / response.audio_transcript.delta + type === "response.output_text.delta" || + type === "response.output_audio_transcript.delta" || + type === "response.audio_transcript.delta" || + type === "response.text.delta" + ) { if (data.delta) appendAssistantText(data.delta); } else if ( type === "conversation.item.input_audio_transcription.completed" ) { if (data.transcript) addMessage("user", data.transcript); } else if (type === "response.done") { - // Ensure we have the full text if deltas were missed + // Ensure we have the full text if deltas were missed. + // Accept both beta (type=text/audio) and GA (type=output_text/output_audio) content. setMessages((prev) => { const last = prev[prev.length - 1]; if (last && last.role === "assistant" && last.content) return prev; - // No assistant message yet — extract from response.done const output = data.response?.output || []; const texts: string[] = []; for (const item of output) { for (const c of item.content || []) { + // beta: c.text (type=text), c.transcript (type=audio) + // GA: c.text (type=output_text), c.transcript (type=output_audio) const t = c.text || c.transcript; if (t) texts.push(t); } @@ -219,10 +233,12 @@ const RealtimePlayground: React.FC = ({ if (!wsRef.current || wsRef.current.readyState !== WebSocket.OPEN) return; // Switch to server VAD mode for voice input + // GA: session.type is required wsRef.current.send( JSON.stringify({ type: "session.update", session: { + type: "realtime", modalities: ["text", "audio"], voice: selectedVoice, input_audio_format: "pcm16", @@ -304,10 +320,12 @@ const RealtimePlayground: React.FC = ({ if (!wsRef.current || wsRef.current.readyState !== WebSocket.OPEN) return; if (configureSessionRef.current) return; configureSessionRef.current = true; + // GA: session.type is required wsRef.current.send( JSON.stringify({ type: "session.update", session: { + type: "realtime", modalities: ["text", "audio"], voice: selectedVoice, input_audio_format: "pcm16",