Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/realtime_endpoints/endpoints.py: 33%
234 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1#### Realtime WebRTC Endpoints #####
3import json
4import time
5from typing import TYPE_CHECKING, Any, Final
7import httpx
8from fastapi import APIRouter, Depends, HTTPException, Request, Response
9from fastapi import status as http_status
11from litellm._logging import verbose_proxy_logger
12from litellm.proxy._types import ProxyException, UserAPIKeyAuth
13from litellm.proxy.auth.auth_checks import can_key_call_resolved_model
14from litellm.proxy.auth.user_api_key_auth import user_api_key_auth
15from litellm.proxy.common_utils.encrypt_decrypt_utils import (
16 decrypt_value_helper,
17 encrypt_value_helper,
18)
19from litellm.proxy.common_utils.http_parsing_utils import _read_request_body
20from litellm.proxy.common_utils.openai_error_payload import (
21 error_status_code,
22 openai_error_param,
23 openai_error_type,
24)
25from litellm.types.realtime import (
26 RealtimeClientSecretRequest,
27 RealtimeClientSecretResponse,
28 RealtimeTranscriptionSessionRequest,
29 RealtimeTranscriptionSessionResponse,
30)
32if TYPE_CHECKING: 32 ↛ 33line 32 didn't jump to line 33 because the condition on line 32 was never true
33 from litellm.router import Router
35router: Final = APIRouter()
37_REALTIME_TOKEN_VERSION: Final = "realtime_v1"
38_DEFAULT_REALTIME_MODEL: Final = "gpt-4o-realtime-preview"
39_DEFAULT_TRANSCRIPTION_MODEL: Final = "gpt-realtime-whisper"
40_ALLOWED_SESSION_TYPES: Final = ("realtime", "transcription")
43def _coerce_realtime_session_type(session_type: str | None) -> str:
44 if session_type in _ALLOWED_SESSION_TYPES: 44 ↛ 45line 44 didn't jump to line 45 because the condition on line 44 was never true
45 return session_type
46 return "realtime"
49def _append_model_candidate(candidates: list[str], model: object) -> None:
50 if isinstance(model, str) and model and model not in candidates:
51 candidates.append(model)
54def _transcription_model_candidates_from_session(session: dict) -> list[str]:
55 candidates: Final[list[str]] = []
57 audio: Final = session.get("audio")
58 if isinstance(audio, dict):
59 audio_input: Final = audio.get("input")
60 if isinstance(audio_input, dict):
61 nested_transcription: Final = audio_input.get("transcription")
62 if isinstance(nested_transcription, dict):
63 _append_model_candidate(
64 candidates,
65 nested_transcription.get("model"),
66 )
68 flat_transcription: Final = session.get("input_audio_transcription")
69 if isinstance(flat_transcription, dict):
70 _append_model_candidate(candidates, flat_transcription.get("model"))
72 return candidates
75def _set_transcription_model_on_session(
76 session: dict,
77 model: str,
78 create_if_missing: bool = False,
79) -> None:
80 updated_existing_config = False
82 flat_transcription: Final = session.get("input_audio_transcription")
83 if isinstance(flat_transcription, dict):
84 session["input_audio_transcription"] = {
85 **flat_transcription,
86 "model": model,
87 }
88 updated_existing_config = True
90 audio = session.get("audio")
91 if isinstance(audio, dict):
92 audio_input = audio.get("input")
93 if isinstance(audio_input, dict):
94 nested_transcription: Final = audio_input.get("transcription")
95 if isinstance(nested_transcription, dict):
96 session["audio"] = {
97 **audio,
98 "input": {
99 **audio_input,
100 "transcription": {
101 **nested_transcription,
102 "model": model,
103 },
104 },
105 }
106 updated_existing_config = True
108 if updated_existing_config or not create_if_missing:
109 return
111 audio = audio if isinstance(audio, dict) else {}
112 audio_input = audio.get("input")
113 audio_input = audio_input if isinstance(audio_input, dict) else {}
114 session["audio"] = {
115 **audio,
116 "input": {
117 **audio_input,
118 "transcription": {"model": model},
119 },
120 }
123async def _prepare_client_secret_session(
124 req: RealtimeClientSecretRequest,
125 user_api_key_dict: UserAPIKeyAuth,
126 llm_model_list: list | None,
127 llm_router: "Router | None",
128) -> tuple[str, dict | None, str]:
129 session_type: Final = _coerce_realtime_session_type(req.session.type if req.session else None)
130 session_data: Final[dict | None] = req.session.model_dump(exclude_none=True) if req.session else None
131 if session_data is not None: 131 ↛ 132line 131 didn't jump to line 132 because the condition on line 131 was never true
132 session_data["type"] = session_type
134 session_model: Final = req.session.model if req.session else None
135 model: str = session_model or req.model or _DEFAULT_REALTIME_MODEL
136 if session_type != "transcription": 136 ↛ 145line 136 didn't jump to line 145 because the condition on line 136 was always true
137 await can_key_call_resolved_model(
138 model=model,
139 valid_token=user_api_key_dict,
140 llm_model_list=llm_model_list,
141 llm_router=llm_router,
142 )
143 return model, session_data, session_type
145 transcription_model_candidates: Final = _transcription_model_candidates_from_session(session_data or {})
146 if not transcription_model_candidates:
147 _append_model_candidate(transcription_model_candidates, session_model)
148 _append_model_candidate(transcription_model_candidates, req.model)
149 if not transcription_model_candidates:
150 transcription_model_candidates.append(_DEFAULT_TRANSCRIPTION_MODEL)
152 model = transcription_model_candidates[0]
153 for transcription_model in transcription_model_candidates:
154 await can_key_call_resolved_model(
155 model=transcription_model,
156 valid_token=user_api_key_dict,
157 llm_model_list=llm_model_list,
158 llm_router=llm_router,
159 )
160 if session_data is not None:
161 _set_transcription_model_on_session(
162 session=session_data,
163 model=model,
164 create_if_missing=True,
165 )
166 session_data.pop("model", None)
167 return model, session_data, session_type
170def _encode_realtime_token_payload(
171 ephemeral_key: str,
172 model_id: str,
173 user_id: str | None,
174 team_id: str | None,
175 expires_at: int | None,
176 session_type: str = "realtime",
177) -> str:
178 """
179 Encode metadata with the upstream ephemeral key so /realtime/calls can
180 route without requiring model as a query param.
181 """
182 payload: Final[dict[str, str | int | None]] = {
183 "v": _REALTIME_TOKEN_VERSION,
184 "ephemeral_key": ephemeral_key,
185 "model_id": model_id,
186 "user_id": user_id or "",
187 "team_id": team_id or "",
188 "expires_at": expires_at,
189 "session_type": session_type,
190 }
191 return json.dumps(payload, separators=(",", ":"))
194def _decode_realtime_token_payload(
195 decrypted_value: str,
196) -> dict[str, Any] | None:
197 """
198 Decode realtime token payload; returns None for legacy/raw ephemeral tokens.
199 """
200 try:
201 decoded: Final = json.loads(decrypted_value)
202 except Exception:
203 return None
205 if not isinstance(decoded, dict):
206 return None
207 if decoded.get("v") != _REALTIME_TOKEN_VERSION:
208 return None
209 if not isinstance(decoded.get("ephemeral_key"), str):
210 return None
211 if not isinstance(decoded.get("model_id"), str):
212 return None
213 return decoded
216@router.post(
217 "/v1/realtime/client_secrets",
218 dependencies=[Depends(user_api_key_auth)],
219 tags=["realtime"],
220)
221@router.post(
222 "/realtime/client_secrets",
223 dependencies=[Depends(user_api_key_auth)],
224 tags=["realtime"],
225)
226@router.post(
227 "/openai/v1/realtime/client_secrets",
228 dependencies=[Depends(user_api_key_auth)],
229 tags=["realtime"],
230)
231async def create_realtime_client_secret(
232 request: Request,
233 fastapi_response: Response,
234 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth),
235) -> RealtimeClientSecretResponse:
236 from litellm.proxy.proxy_server import (
237 add_litellm_data_to_request,
238 general_settings,
239 llm_model_list,
240 llm_router,
241 proxy_config,
242 proxy_logging_obj,
243 route_request,
244 user_model,
245 version,
246 )
248 data: dict = {}
249 try:
250 body: Final = await _read_request_body(request=request)
251 req: Final = RealtimeClientSecretRequest(**body)
253 model, session_data, session_type = await _prepare_client_secret_session(
254 req=req,
255 user_api_key_dict=user_api_key_dict,
256 llm_model_list=llm_model_list,
257 llm_router=llm_router,
258 )
260 data = {"model": model}
262 # If session is provided, use it; otherwise create one from model
263 if session_data is not None: 263 ↛ 264line 263 didn't jump to line 264 because the condition on line 263 was never true
264 data["session"] = session_data
265 elif req.model: 265 ↛ 267line 265 didn't jump to line 267 because the condition on line 265 was never true
266 # User provided model at root level, convert to session format
267 data["session"] = {"type": "realtime", "model": model}
269 if req.expires_after: 269 ↛ 270line 269 didn't jump to line 270 because the condition on line 269 was never true
270 data["expires_after"] = req.expires_after.model_dump(exclude_none=True)
272 data = await add_litellm_data_to_request(
273 data=data,
274 request=request,
275 general_settings=general_settings,
276 user_api_key_dict=user_api_key_dict,
277 version=version,
278 proxy_config=proxy_config,
279 )
281 data = await proxy_logging_obj.pre_call_hook(
282 user_api_key_dict=user_api_key_dict,
283 data=data,
284 call_type="acreate_realtime_client_secret",
285 )
287 verbose_proxy_logger.debug("WebRTC: /v1/realtime/client_secrets (model=%s)", model)
289 llm_call: Final = await route_request(
290 data=data,
291 route_type="acreate_realtime_client_secret",
292 llm_router=llm_router,
293 user_model=user_model,
294 )
295 upstream_resp: Final[httpx.Response] = await llm_call
297 except Exception as e:
298 await proxy_logging_obj.post_call_failure_hook(
299 user_api_key_dict=user_api_key_dict,
300 original_exception=e,
301 request_data=data,
302 )
303 verbose_proxy_logger.error(
304 "litellm.proxy.realtime_endpoints.webrtc.create_realtime_client_secret(): Exception - %s",
305 str(e),
306 )
307 if isinstance(e, ProxyException): 307 ↛ 308line 307 didn't jump to line 308 because the condition on line 307 was never true
308 raise e
309 if isinstance(e, HTTPException): 309 ↛ 316line 309 didn't jump to line 316 because the condition on line 309 was always true
310 raise ProxyException(
311 message=getattr(e, "message", str(e)),
312 type=openai_error_type(e, error_status_code(e, http_status.HTTP_400_BAD_REQUEST)),
313 param=openai_error_param(e),
314 code=error_status_code(e, http_status.HTTP_400_BAD_REQUEST),
315 )
316 raise ProxyException(
317 message=getattr(e, "message", str(e)),
318 type=openai_error_type(e, error_status_code(e, 500)),
319 param=openai_error_param(e),
320 code=error_status_code(e, 500),
321 )
323 if upstream_resp.status_code != 200:
324 verbose_proxy_logger.error(
325 "WebRTC client_secrets upstream error %s: %s",
326 upstream_resp.status_code,
327 upstream_resp.text,
328 )
329 return Response(
330 content=upstream_resp.content,
331 status_code=upstream_resp.status_code,
332 media_type="application/json",
333 )
335 upstream_json: Final[dict] = upstream_resp.json()
337 # Encrypt upstream ephemeral key with routing metadata so /realtime/calls
338 # can recover model without requiring query params.
339 raw_value: Final[str] = upstream_json.get("value", "")
340 expires_at: Final = upstream_json.get("expires_at")
341 token_payload: Final = _encode_realtime_token_payload(
342 ephemeral_key=raw_value,
343 model_id=model,
344 user_id=getattr(user_api_key_dict, "user_id", None),
345 team_id=getattr(user_api_key_dict, "team_id", None),
346 expires_at=expires_at if isinstance(expires_at, int) else None,
347 session_type=session_type,
348 )
349 encrypted_token: Final[str] = encrypt_value_helper(token_payload)
350 upstream_json["value"] = encrypted_token
352 session_obj: Final[dict | None] = upstream_json.get("session")
353 if isinstance(session_obj, dict):
354 cs: Final = session_obj.get("client_secret")
355 if isinstance(cs, dict) and "value" in cs:
356 cs["value"] = encrypted_token
357 upstream_json["session"] = session_obj
359 return RealtimeClientSecretResponse(**upstream_json)
362@router.post(
363 "/v1/realtime/calls",
364 tags=["realtime"],
365)
366@router.post(
367 "/realtime/calls",
368 tags=["realtime"],
369)
370@router.post(
371 "/openai/v1/realtime/calls",
372 tags=["realtime"],
373)
374async def proxy_realtime_calls(
375 request: Request,
376 fastapi_response: Response,
377) -> Response:
378 from litellm.proxy.proxy_server import (
379 add_litellm_data_to_request,
380 general_settings,
381 llm_router,
382 proxy_config,
383 proxy_logging_obj,
384 route_request,
385 user_model,
386 version,
387 )
389 # Auth: the Bearer token is the encrypted ephemeral key issued by
390 # /realtime/client_secrets, not a standard proxy API key.
391 auth_header: Final[str | None] = request.headers.get("Authorization")
392 if not auth_header or not auth_header.startswith("Bearer "): 392 ↛ 399line 392 didn't jump to line 399 because the condition on line 392 was always true
393 return Response(
394 content=json.dumps({"error": "Missing or invalid Authorization header"}),
395 status_code=http_status.HTTP_401_UNAUTHORIZED,
396 media_type="application/json",
397 )
399 encrypted_token: Final = auth_header.removeprefix("Bearer ").strip()
400 decrypted_token_value: Final = decrypt_value_helper(
401 value=encrypted_token,
402 key="realtime_calls_auth",
403 )
404 if not decrypted_token_value:
405 return Response(
406 content=json.dumps({"error": "Invalid or expired token"}),
407 status_code=http_status.HTTP_401_UNAUTHORIZED,
408 media_type="application/json",
409 )
411 sdp_body: Final[bytes] = await request.body()
412 decoded_payload: Final = _decode_realtime_token_payload(decrypted_token_value)
413 if decoded_payload is not None:
414 # Check token expiry
415 expires_at: Final = decoded_payload.get("expires_at")
416 if expires_at is not None and isinstance(expires_at, int):
417 if time.time() > expires_at:
418 return Response(
419 content=json.dumps({"error": "Token has expired"}),
420 status_code=http_status.HTTP_401_UNAUTHORIZED,
421 media_type="application/json",
422 )
424 openai_ephemeral_key = decoded_payload.get("ephemeral_key", "")
425 model = decoded_payload.get("model_id") or request.query_params.get("model") or _DEFAULT_REALTIME_MODEL
426 user_id = decoded_payload.get("user_id") or None
427 team_id = decoded_payload.get("team_id") or None
428 session_type = _coerce_realtime_session_type(decoded_payload.get("session_type"))
429 else:
430 # Backward compatibility: older tokens contained only encrypted upstream key.
431 openai_ephemeral_key = decrypted_token_value
432 model = request.query_params.get("model", _DEFAULT_REALTIME_MODEL)
433 user_id = None
434 team_id = None
435 session_type = "realtime"
437 # Build a minimal UserAPIKeyAuth with user/team IDs from the token
438 # so spend tracking and budget enforcement work correctly.
439 minimal_auth: Final = UserAPIKeyAuth(
440 user_id=user_id,
441 team_id=team_id,
442 )
444 data: dict = {}
445 try:
446 session_config: Final = {
447 "type": session_type,
448 }
449 if session_type == "transcription":
450 _set_transcription_model_on_session(
451 session=session_config,
452 model=model,
453 create_if_missing=True,
454 )
455 else:
456 session_config["model"] = model
458 data = {
459 "model": model,
460 "openai_ephemeral_key": openai_ephemeral_key,
461 "sdp_body": sdp_body,
462 "session": session_config,
463 }
465 data = await add_litellm_data_to_request(
466 data=data,
467 request=request,
468 general_settings=general_settings,
469 user_api_key_dict=minimal_auth,
470 version=version,
471 proxy_config=proxy_config,
472 )
474 data = await proxy_logging_obj.pre_call_hook(
475 user_api_key_dict=minimal_auth,
476 data=data,
477 call_type="arealtime_calls",
478 )
480 verbose_proxy_logger.debug("WebRTC: /v1/realtime/calls (model=%s)", model)
482 llm_call: Final = await route_request(
483 data=data,
484 route_type="arealtime_calls",
485 llm_router=llm_router,
486 user_model=user_model,
487 )
488 upstream_resp: Final[httpx.Response] = await llm_call
490 except Exception as e:
491 await proxy_logging_obj.post_call_failure_hook(
492 user_api_key_dict=minimal_auth,
493 original_exception=e,
494 request_data=data,
495 )
496 verbose_proxy_logger.error(
497 "litellm.proxy.realtime_endpoints.webrtc.proxy_realtime_calls(): Exception - %s",
498 str(e),
499 )
500 if isinstance(e, HTTPException):
501 raise ProxyException(
502 message=getattr(e, "message", str(e)),
503 type=openai_error_type(e, error_status_code(e, http_status.HTTP_400_BAD_REQUEST)),
504 param=openai_error_param(e),
505 code=error_status_code(e, http_status.HTTP_400_BAD_REQUEST),
506 )
507 raise ProxyException(
508 message=getattr(e, "message", str(e)),
509 type=openai_error_type(e, error_status_code(e, 500)),
510 param=openai_error_param(e),
511 code=error_status_code(e, 500),
512 )
514 return Response(
515 content=upstream_resp.content,
516 status_code=upstream_resp.status_code,
517 media_type=upstream_resp.headers.get("content-type", "application/sdp"),
518 )
521@router.post(
522 "/v1/realtime/transcription_sessions",
523 dependencies=[Depends(user_api_key_auth)],
524 tags=["realtime"],
525)
526@router.post(
527 "/realtime/transcription_sessions",
528 dependencies=[Depends(user_api_key_auth)],
529 tags=["realtime"],
530)
531@router.post(
532 "/openai/v1/realtime/transcription_sessions",
533 dependencies=[Depends(user_api_key_auth)],
534 tags=["realtime"],
535)
536async def create_realtime_transcription_session(
537 request: Request,
538 fastapi_response: Response,
539 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth),
540) -> RealtimeTranscriptionSessionResponse:
541 """
542 Create an ephemeral Realtime transcription session
543 (POST /v1/realtime/transcription_sessions) for the WebRTC/WebSocket flow.
545 Mirrors the client_secrets route but targets the transcription_sessions
546 endpoint and encrypts the ephemeral key returned under `client_secret.value`.
547 """
548 from litellm.proxy.proxy_server import (
549 add_litellm_data_to_request,
550 general_settings,
551 llm_model_list,
552 llm_router,
553 proxy_config,
554 proxy_logging_obj,
555 route_request,
556 user_model,
557 version,
558 )
560 data: dict = {}
561 try:
562 body: Final = await _read_request_body(request=request)
563 req: Final = RealtimeTranscriptionSessionRequest(**body)
565 model: Final[str] = req.resolved_model() or "gpt-realtime-whisper"
566 await can_key_call_resolved_model(
567 model=model,
568 valid_token=user_api_key_dict,
569 llm_model_list=llm_model_list,
570 llm_router=llm_router,
571 )
573 transcription_session: Final = {k: v for k, v in body.items() if k != "model"}
574 data = {"model": model, "transcription_session": transcription_session}
576 data = await add_litellm_data_to_request(
577 data=data,
578 request=request,
579 general_settings=general_settings,
580 user_api_key_dict=user_api_key_dict,
581 version=version,
582 proxy_config=proxy_config,
583 )
585 data = await proxy_logging_obj.pre_call_hook(
586 user_api_key_dict=user_api_key_dict,
587 data=data,
588 call_type="acreate_realtime_transcription_session",
589 )
591 verbose_proxy_logger.debug("Realtime: /v1/realtime/transcription_sessions (model=%s)", model)
593 llm_call: Final = await route_request(
594 data=data,
595 route_type="acreate_realtime_transcription_session",
596 llm_router=llm_router,
597 user_model=user_model,
598 )
599 upstream_resp: Final[httpx.Response] = await llm_call
601 except Exception as e:
602 await proxy_logging_obj.post_call_failure_hook(
603 user_api_key_dict=user_api_key_dict,
604 original_exception=e,
605 request_data=data,
606 )
607 verbose_proxy_logger.error(
608 "litellm.proxy.realtime_endpoints.create_realtime_transcription_session(): Exception - %s",
609 str(e),
610 )
611 if isinstance(e, ProxyException): 611 ↛ 612line 611 didn't jump to line 612 because the condition on line 611 was never true
612 raise e
613 if isinstance(e, HTTPException): 613 ↛ 620line 613 didn't jump to line 620 because the condition on line 613 was always true
614 raise ProxyException(
615 message=getattr(e, "detail", getattr(e, "message", str(e))),
616 type=openai_error_type(e, error_status_code(e, http_status.HTTP_400_BAD_REQUEST)),
617 param=openai_error_param(e),
618 code=error_status_code(e, http_status.HTTP_400_BAD_REQUEST),
619 )
620 raise ProxyException(
621 message=getattr(e, "message", str(e)),
622 type=openai_error_type(e, error_status_code(e, 500)),
623 param=openai_error_param(e),
624 code=error_status_code(e, 500),
625 )
627 if upstream_resp.status_code != 200:
628 verbose_proxy_logger.error(
629 "Realtime transcription_sessions upstream error %s: %s",
630 upstream_resp.status_code,
631 upstream_resp.text,
632 )
633 return Response(
634 content=upstream_resp.content,
635 status_code=upstream_resp.status_code,
636 media_type="application/json",
637 )
639 upstream_json: Final[dict] = upstream_resp.json()
641 # Encrypt the ephemeral key (returned under client_secret.value) with routing
642 # metadata so the follow-up /realtime/calls request can recover the model.
643 client_secret: Final = upstream_json.get("client_secret")
644 if isinstance(client_secret, dict) and "value" in client_secret:
645 raw_value: Final[str] = client_secret.get("value", "")
646 expires_at: Final = client_secret.get("expires_at")
647 token_payload: Final = _encode_realtime_token_payload(
648 ephemeral_key=raw_value,
649 model_id=model,
650 user_id=getattr(user_api_key_dict, "user_id", None),
651 team_id=getattr(user_api_key_dict, "team_id", None),
652 expires_at=expires_at if isinstance(expires_at, int) else None,
653 session_type="transcription",
654 )
655 client_secret["value"] = encrypt_value_helper(token_payload)
656 upstream_json["client_secret"] = client_secret
658 return RealtimeTranscriptionSessionResponse(**upstream_json)