Module livekit.plugins.openai.realtime

Classes

class GPTLiveDelegation (id: str, pending_transcript: str)
Expand source code
@dataclass
class GPTLiveDelegation:
    """Work the model handed to the application, under client delegation.

    It carries no task text: the ask is whatever the conversation says, which the agent's chat
    context holds. Answer it with :meth:`GPTLiveSession.append_commentary`, passing this id.
    """

    id: str
    pending_transcript: str
    """The caller's current turn, not yet in the chat context when the model delegates."""

Work the model handed to the application, under client delegation.

It carries no task text: the ask is whatever the conversation says, which the agent's chat context holds. Answer it with :meth:GPTLiveSession.append_commentary(), passing this id.

Instance variables

var id : str
var pending_transcript : str

The caller's current turn, not yet in the chat context when the model delegates.

class GPTLiveModel (*,
model: str = 'gpt-live-1',
voice: GPTLiveVoices | str | dict[str, Any] = 'marin',
delegation: types.DelegationTarget = 'responses',
responses_options: NotGivenOr[ResponsesDelegationOptions] = NOT_GIVEN,
service_tier: NotGivenOr[types.ServiceTier] = NOT_GIVEN,
api_key: str | None = None,
base_url: NotGivenOr[str] = NOT_GIVEN,
http_session: aiohttp.ClientSession | None = None,
max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
conn_options: APIConnectOptions = APIConnectOptions(max_retry=3, retry_interval=2.0, timeout=10.0),
azure_deployment: str | None = None,
entra_token: str | None = None)
Expand source code
class GPTLiveModel(llm.DuplexModel):
    """OpenAI GPT-Live full-duplex voice model, ready to pass to ``AgentSession(llm=)``."""

    def __init__(
        self,
        *,
        model: str = DEFAULT_MODEL,
        voice: GPTLiveVoices | str | dict[str, Any] = DEFAULT_VOICE,
        delegation: types.DelegationTarget = "responses",
        responses_options: NotGivenOr[ResponsesDelegationOptions] = NOT_GIVEN,
        service_tier: NotGivenOr[types.ServiceTier] = NOT_GIVEN,
        api_key: str | None = None,
        base_url: NotGivenOr[str] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
        azure_deployment: str | None = None,
        entra_token: str | None = None,
    ) -> None:
        """
        Args:
            model: GPT-Live voice model slug. Ignored on Azure, where ``azure_deployment`` names
                the model.
            voice: Output voice: a name from :data:`GPTLiveVoices`, another supported name, or
                ``{"id": "voice_..."}`` for an authorized custom voice. Defaults to ``marin``.
                Immutable after the session starts.
            delegation: Where delegated work goes, fixed for the life of the session.
                ``responses`` runs it on a backend model, so ``@function_tool`` works as usual;
                ``client`` hands it to the application as a ``delegation_created`` event, which
                no framework tool can answer.
            responses_options: The backend Responses model, under ``delegation="responses"``.
            service_tier: Processing tier the session runs under, sent as the
                ``OpenAI-Service-Tier`` header on the connection, for example ``ultrafast``.
                Unset leaves the header off, and the account's default applies.
            api_key: OpenAI API key. Falls back to ``OPENAI_API_KEY``, or to
                ``AZURE_OPENAI_API_KEY`` on Azure unless ``entra_token`` is given.
            base_url: HTTP base url of the OpenAI API. On Azure, the resource endpoint, falling
                back to ``AZURE_OPENAI_ENDPOINT``.
            http_session: Optional shared HTTP session.
            max_session_duration: Seconds before the connection is recycled.
            conn_options: Retry/backoff and connection settings.
            azure_deployment: Azure OpenAI deployment of the voice model. Giving it or
                ``entra_token`` selects Azure; prefer :meth:`with_azure`.
            entra_token: Microsoft Entra ID token for Azure, instead of ``api_key``.
        """
        super().__init__(
            capabilities=llm.DuplexCapabilities(
                user_transcription=True,
                # the model continues on its own once every tool result reaches the backend
                auto_tool_reply_generation=True,
                mutable_chat_context=False,
                mutable_instructions=False,
                # tools live on the backend model, and a client delegation has none
                mutable_tools=delegation == "responses",
            )
        )
        responses = (
            responses_options if is_given(responses_options) else ResponsesDelegationOptions()
        )
        is_azure = azure_deployment is not None or entra_token is not None
        if not is_azure:
            api_key = api_key or os.environ.get("OPENAI_API_KEY")
            if api_key is None:
                raise ValueError(
                    "The api_key client option must be set either by passing api_key "
                    "to the client or by setting the OPENAI_API_KEY environment variable"
                )
            resolved_base_url = (
                base_url if is_given(base_url) else os.getenv("OPENAI_BASE_URL", OPENAI_BASE_URL)
            )
        else:
            if not azure_deployment:
                raise ValueError("Azure needs azure_deployment, the voice model's deployment name")
            model = azure_deployment

            if api_key and entra_token:
                raise ValueError("api_key and entra_token are mutually exclusive")
            if entra_token is None:
                api_key = api_key or os.environ.get("AZURE_OPENAI_API_KEY")
                if not api_key:
                    raise ValueError(
                        "Missing Azure credentials. Pass api_key or entra_token, "
                        "or set the AZURE_OPENAI_API_KEY environment variable"
                    )

            endpoint = base_url if is_given(base_url) else os.environ.get("AZURE_OPENAI_ENDPOINT")
            if not endpoint:
                raise ValueError(
                    "Missing Azure endpoint. Pass azure_endpoint or base_url, or set the "
                    "AZURE_OPENAI_ENDPOINT environment variable"
                )
            resolved_base_url = endpoint

            # the service resolves the backend model as a deployment of this resource, where the
            # OpenAI default names nothing: every delegated turn would fail while the voice
            # model keeps promising an answer, so ask for the deployment up front
            if delegation == "responses" and not responses.get("model"):
                raise ValueError(
                    "Azure responses delegation needs responses_options['model'], the name of a "
                    "Responses deployment in the same resource; or use delegation='client'"
                )

        self._opts = _LiveOptions(
            model=model,
            voice=voice,
            delegation=delegation,
            responses=responses,
            service_tier=service_tier if is_given(service_tier) else None,
            api_key=api_key,
            base_url=resolved_base_url,
            conn_options=conn_options,
            max_session_duration=max_session_duration if is_given(max_session_duration) else None,
            is_azure=is_azure,
            entra_token=entra_token,
        )
        self._http_session = http_session
        self._http_session_owned = False
        self._provider_label = "OpenAI Live API"

    @classmethod
    def with_azure(
        cls,
        *,
        azure_deployment: str,
        azure_endpoint: str | None = None,
        api_key: str | None = None,
        entra_token: str | None = None,
        base_url: str | None = None,
        voice: GPTLiveVoices | str | dict[str, Any] = DEFAULT_VOICE,
        delegation: types.DelegationTarget = "responses",
        responses_options: NotGivenOr[ResponsesDelegationOptions] = NOT_GIVEN,
        service_tier: NotGivenOr[types.ServiceTier] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
    ) -> GPTLiveModel:
        """Create a GPTLiveModel served by Azure OpenAI.

        Args:
            azure_deployment: Deployment name of the GPT-Live voice model.
            azure_endpoint: Resource endpoint, such as ``https://<resource>.openai.azure.com``.
                Falls back to ``AZURE_OPENAI_ENDPOINT``.
            api_key: Azure API key. Falls back to ``AZURE_OPENAI_API_KEY`` unless
                ``entra_token`` is given.
            entra_token: Microsoft Entra ID token, instead of ``api_key``.
            base_url: Explicit base URL, for example a gateway in front of the resource.
                Mutually exclusive with ``azure_endpoint``.
            voice: Output voice, as in :class:`GPTLiveModel`.
            delegation: Where delegated work goes, as in :class:`GPTLiveModel`.
            responses_options: The backend Responses model under ``delegation="responses"``.
                Its ``model`` is required, and names a deployment in the same resource.
            service_tier: Processing tier, as in :class:`GPTLiveModel`.
            http_session: Optional shared HTTP session.
            max_session_duration: Seconds before the connection is recycled.
            conn_options: Retry/backoff and connection settings.

        Returns:
            GPTLiveModel: A model that connects to the Azure resource.

        Raises:
            ValueError: If the deployment, credentials, or endpoint are missing, if both
                ``api_key`` and ``entra_token`` or both ``base_url`` and ``azure_endpoint`` are
                given, or if responses delegation has no backend deployment.

        Example:
            ```python
            from livekit.plugins.openai.realtime import GPTLiveModel

            model = GPTLiveModel.with_azure(
                azure_deployment="gpt-live-1",
                azure_endpoint="https://<resource>.openai.azure.com",
                api_key="<api-key>",
                responses_options={"model": "<responses-deployment>"},
            )
            ```
        """
        if base_url is not None and azure_endpoint is not None:
            raise ValueError("base_url and azure_endpoint are mutually exclusive")
        endpoint = base_url if base_url is not None else azure_endpoint

        return cls(
            voice=voice,
            delegation=delegation,
            responses_options=responses_options,
            service_tier=service_tier,
            api_key=api_key,
            base_url=endpoint if endpoint is not None else NOT_GIVEN,
            http_session=http_session,
            max_session_duration=max_session_duration,
            conn_options=conn_options,
            azure_deployment=azure_deployment,
            entra_token=entra_token,
        )

    @property
    def model(self) -> str:
        return self._opts.model

    @property
    def provider(self) -> str:
        return urlparse(self._opts.base_url).netloc

    def _ensure_http_session(self) -> aiohttp.ClientSession:
        if not self._http_session:
            try:
                self._http_session = utils.http_context.http_session()
            except RuntimeError:
                self._http_session = aiohttp.ClientSession()
                self._http_session_owned = True
        return self._http_session

    def audio_gate(self) -> llm.AudioGate:
        return llm.FixedGate(_SILENCE_RMS, min_silence_duration=_MIN_SILENCE_DURATION)

    def session(self) -> GPTLiveSession:
        return GPTLiveSession(self)

    async def aclose(self) -> None:
        if self._http_session_owned and self._http_session:
            await self._http_session.close()

OpenAI GPT-Live full-duplex voice model, ready to pass to AgentSession(llm=).

Args

model
GPT-Live voice model slug. Ignored on Azure, where azure_deployment names the model.
voice
Output voice: a name from :data:GPTLiveVoices, another supported name, or {"id": "voice_..."} for an authorized custom voice. Defaults to marin. Immutable after the session starts.
delegation
Where delegated work goes, fixed for the life of the session. responses runs it on a backend model, so @function_tool works as usual; client hands it to the application as a delegation_created event, which no framework tool can answer.
responses_options
The backend Responses model, under delegation="responses".
service_tier
Processing tier the session runs under, sent as the OpenAI-Service-Tier header on the connection, for example ultrafast. Unset leaves the header off, and the account's default applies.
api_key
OpenAI API key. Falls back to OPENAI_API_KEY, or to AZURE_OPENAI_API_KEY on Azure unless entra_token is given.
base_url
HTTP base url of the OpenAI API. On Azure, the resource endpoint, falling back to AZURE_OPENAI_ENDPOINT.
http_session
Optional shared HTTP session.
max_session_duration
Seconds before the connection is recycled.
conn_options
Retry/backoff and connection settings.
azure_deployment
Azure OpenAI deployment of the voice model. Giving it or entra_token selects Azure; prefer :meth:with_azure.
entra_token
Microsoft Entra ID token for Azure, instead of api_key.

Ancestors

  • livekit.agents.llm.duplex.DuplexModel
  • abc.ABC

Static methods

def with_azure(*,
azure_deployment: str,
azure_endpoint: str | None = None,
api_key: str | None = None,
entra_token: str | None = None,
base_url: str | None = None,
voice: GPTLiveVoices | str | dict[str, Any] = 'marin',
delegation: types.DelegationTarget = 'responses',
responses_options: NotGivenOr[ResponsesDelegationOptions] = NOT_GIVEN,
service_tier: NotGivenOr[types.ServiceTier] = NOT_GIVEN,
http_session: aiohttp.ClientSession | None = None,
max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
conn_options: APIConnectOptions = APIConnectOptions(max_retry=3, retry_interval=2.0, timeout=10.0)) ‑> livekit.plugins.openai.realtime.gpt_live_model.GPTLiveModel

Create a GPTLiveModel served by Azure OpenAI.

Args

azure_deployment
Deployment name of the GPT-Live voice model.
azure_endpoint
Resource endpoint, such as https://<resource>.openai.azure.com. Falls back to AZURE_OPENAI_ENDPOINT.
api_key
Azure API key. Falls back to AZURE_OPENAI_API_KEY unless entra_token is given.
entra_token
Microsoft Entra ID token, instead of api_key.
base_url
Explicit base URL, for example a gateway in front of the resource. Mutually exclusive with azure_endpoint.
voice
Output voice, as in :class:GPTLiveModel.
delegation
Where delegated work goes, as in :class:GPTLiveModel.
responses_options
The backend Responses model under delegation="responses". Its model is required, and names a deployment in the same resource.
service_tier
Processing tier, as in :class:GPTLiveModel.
http_session
Optional shared HTTP session.
max_session_duration
Seconds before the connection is recycled.
conn_options
Retry/backoff and connection settings.

Returns

GPTLiveModel
A model that connects to the Azure resource.

Raises

ValueError
If the deployment, credentials, or endpoint are missing, if both api_key and entra_token or both base_url and azure_endpoint are given, or if responses delegation has no backend deployment.

Example

from livekit.plugins.openai.realtime import GPTLiveModel

model = GPTLiveModel.with_azure(
    azure_deployment="gpt-live-1",
    azure_endpoint="https://<resource>.openai.azure.com",
    api_key="<api-key>",
    responses_options={"model": "<responses-deployment>"},
)

Instance variables

prop model : str
Expand source code
@property
def model(self) -> str:
    return self._opts.model
prop provider : str
Expand source code
@property
def provider(self) -> str:
    return urlparse(self._opts.base_url).netloc

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    if self._http_session_owned and self._http_session:
        await self._http_session.close()
def audio_gate(self) ‑> livekit.agents.llm.duplex_adapter.AudioGate
Expand source code
def audio_gate(self) -> llm.AudioGate:
    return llm.FixedGate(_SILENCE_RMS, min_silence_duration=_MIN_SILENCE_DURATION)

The gate for this model's output, or None to let the adapter infer one.

def session(self) ‑> livekit.plugins.openai.realtime.gpt_live_model.GPTLiveSession
Expand source code
def session(self) -> GPTLiveSession:
    return GPTLiveSession(self)

Open a session; the adapter configures it with _update_session before use.

class GPTLiveSession (duplex_model: GPTLiveModel)
Expand source code
class GPTLiveSession(
    llm.DuplexSession[
        Literal["openai_server_event_received", "openai_client_event_queued", "delegation_created"]
    ]
):
    """A session for the OpenAI GPT-Live API (WebSocket), reached with ``Agent.duplex_session``.

    Exposes three extra events:
    - openai_server_event_received: raw server events
    - openai_client_event_queued: raw client events sent to the server
    - delegation_created: a :class:`GPTLiveDelegation`, under client delegation

    The ``append_*`` methods queue context and return without waiting. Their ``*.appended``
    events arrive at the estimated context-injection end; they do not mean speech has finished.
    """

    def __init__(self, duplex_model: GPTLiveModel) -> None:
        super().__init__(duplex_model)
        self._live_model = duplex_model
        self._opts = replace(duplex_model._opts, responses=duplex_model._opts.responses.copy())
        self._tools = llm.ToolContext.empty()
        # the agent's instructions, set by _update_session before session.start and immutable after
        self._instructions: str | None = None
        self._msg_ch = utils.aio.Chan[types.ClientEvent | dict[str, Any]]()
        self._audio_ch = utils.aio.Chan[llm.DuplexAudioFrame]()
        self._input_resampler: rtc.AudioResampler | None = None

        # session.start opens a connection and carries the config that is immutable after it
        self._session_start_sent = False
        self._session_started_fut: asyncio.Future[None] = asyncio.Future()
        self._session_closed_fut: asyncio.Future[None] = asyncio.Future()
        self._session_id: str | None = None
        self._num_retries = 0
        # session usage is reported cumulatively; kept to emit per-event deltas
        self._usage_total = types.Usage()

        # everything the model has been told, both speakers' words included, to reseed a
        # reconnect; the framework's chat context is the adapter's, not this
        self._history = llm.ChatContext.empty()
        self._speech: dict[Role, _Speech] = {}

        # delegation_id -> call_ids of its response that has not completed yet; on completion the
        # call_ids move to the open set, and stay there until each output is sent to the backend.
        # response_pending is set once an output is sent and cleared by the response.create
        self._backend_running_responses: dict[str | None, set[str]] = {}
        self._backend_open_calls: set[str] = set()
        self._backend_response_pending = False

        # the newest history item the last ask was about, so an ask never repeats one
        self._asked_item_id: str | None = None

        self._bstream = utils.audio.AudioByteStream(
            SAMPLE_RATE, NUM_CHANNELS, samples_per_channel=SAMPLE_RATE // 10
        )

        self._main_atask = asyncio.create_task(self._main_task(), name="GPTLiveSession._main")

    # outbound

    def send_event(self, event: types.ClientEvent | dict[str, Any]) -> None:
        with contextlib.suppress(utils.aio.channel.ChanClosed):
            self._msg_ch.send_nowait(event)

    def _build_delegation(self) -> types.Delegation:
        if self._opts.delegation == "client":
            return types.Delegation(type="client")
        opts = self._opts.responses
        return types.Delegation(
            type="responses",
            responses=types.ResponsesConfig(
                model=opts.get("model", DEFAULT_BACKEND_MODEL),
                instructions=opts.get("instructions"),
                tools=_build_delegation_tools(self._tools.flatten()) or None,
                tool_choice=_to_tool_choice(opts["tool_choice"]) if "tool_choice" in opts else None,
                parallel_tool_calls=opts.get("parallel_tool_calls"),
                reasoning=opts.get("reasoning"),
                text=opts.get("text"),
                service_tier=opts.get("service_tier"),
                max_output_tokens=opts.get("max_output_tokens"),
            ),
        )

    def _session_start_event(self) -> types.SessionStartEvent:
        """The whole configuration, composed fresh for each connection."""
        # the conversation so far is startup history, newest first until the cap is reached
        items: list[types.InputItem] = []
        dropped = 0
        for item in reversed(self._history.items):
            if (rendered := _render_item(item)) is None:
                continue
            role, text = rendered
            if len(items) >= _MAX_INPUT_ITEMS:
                dropped += 1
                continue
            part = (
                types.OutputTextPart(text=text)
                if role == "assistant"
                else types.InputTextPart(text=text)
            )
            items.append(types.InputItem(role=role, content=[part]))
        if dropped:
            logger.warning(
                "gpt-live startup history exceeds what a session accepts; dropping the oldest",
                extra={"dropped": dropped, "kept": len(items)},
            )
        items.reverse()

        return types.SessionStartEvent(
            event_id=utils.shortuuid("session_start_"),
            session=types.SessionConfig(
                model=self._opts.model,
                instructions=self._instructions,
                input=items or None,
                audio=types.AudioConfig(
                    format=types.AudioFormat(type="audio/pcm", rate=SAMPLE_RATE),
                    output=types.AudioOutput(voice=self._opts.voice),
                ),
                delegation=self._build_delegation(),
            ),
        )

    def _send_delegation_update(self, responses: types.ResponsesConfig) -> None:
        """A sparse session.update: only the backend settings can change once started."""
        if self._opts.delegation != "responses" or not self._session_start_sent:
            return
        self.send_event(
            types.SessionUpdateEvent(
                event_id=utils.shortuuid("delegation_update_"),
                session=types.SessionUpdateConfig(
                    delegation=types.Delegation(type="responses", responses=responses)
                ),
            )
        )

    # connection loop

    @utils.log_exceptions(logger=logger)
    async def _main_task(self) -> None:
        max_retries = self._opts.conn_options.max_retry
        reconnecting = False

        try:
            while not self._msg_ch.closed:
                try:
                    ws_conn = await self._create_ws_conn()
                    if reconnecting:
                        self._reset_for_reconnect()
                        self.emit("session_reconnected", llm.RealtimeSessionReconnectedEvent())
                    try:
                        await self._run_ws(ws_conn)
                    finally:
                        # what arrives now is history for the next connection, not an append
                        self._session_start_sent = False
                except APIError as e:
                    if max_retries == 0 or not e.retryable:
                        self._emit_error(e, recoverable=False)
                        raise
                    elif self._num_retries == max_retries:
                        self._emit_error(e, recoverable=False)
                        raise APIConnectionError(
                            f"{self._live_model._provider_label} connection failed after "
                            f"{self._num_retries} attempts",
                        ) from e
                    else:
                        self._emit_error(e, recoverable=True)
                        interval = self._opts.conn_options._interval_for_retry(self._num_retries)
                        logger.warning(
                            f"{self._live_model._provider_label} connection failed, "
                            f"retrying in {interval}s",
                            exc_info=e,
                        )
                        await asyncio.sleep(interval)
                    self._num_retries += 1
                except Exception as e:
                    logger.error("gpt-live session failed", extra={"error_type": type(e).__name__})
                    error = APIConnectionError("GPT-Live session failed", retryable=False)
                    self._emit_error(error, recoverable=False)
                    raise error from None
                reconnecting = True
        finally:
            self._audio_ch.close()

    def _reset_for_reconnect(self) -> None:
        # a new connection is a new session, reseeded from the history; the rest of what the
        # dropped one was carrying never arrives
        self._bstream.clear()
        self._input_resampler = None
        self._session_started_fut = asyncio.Future()
        self._session_closed_fut = asyncio.Future()
        self._end_speech("user")
        self._speech.clear()
        self._backend_running_responses.clear()
        self._backend_open_calls.clear()
        self._backend_response_pending = False
        self._usage_total = types.Usage()
        self._session_id = None

    async def _create_ws_conn(self) -> aiohttp.ClientWebSocketResponse:
        headers = {"User-Agent": "LiveKit Agents"}
        if not self._opts.is_azure:
            headers["Authorization"] = f"Bearer {self._opts.api_key}"
        elif self._opts.entra_token:
            headers["Authorization"] = f"Bearer {self._opts.entra_token}"
        elif self._opts.api_key:
            # Azure answers a bearer api key with a redirect to ?api-key=, which aiohttp does not
            # follow for a websocket, so the key goes in its own header
            headers["api-key"] = self._opts.api_key
        if self._opts.service_tier:
            headers["OpenAI-Service-Tier"] = self._opts.service_tier
        url = _live_sessions_url(self._opts.base_url, is_azure=self._opts.is_azure)
        if lk_oai_debug:
            logger.debug("connecting to GPT-Live API", extra={"lk.pii.url": url})

        t0 = time.perf_counter()
        try:
            ws = await asyncio.wait_for(
                self._live_model._ensure_http_session().ws_connect(url=url, headers=headers),
                self._opts.conn_options.timeout,
            )
            self._report_connection_acquired(time.perf_counter() - t0)
            return ws
        except (aiohttp.ClientError, asyncio.TimeoutError):
            raise APIConnectionError(
                f"{self._live_model._provider_label} connection error"
            ) from None

    async def _run_ws(self, ws_conn: aiohttp.ClientWebSocketResponse) -> None:
        closing = False

        async def _close_ws() -> None:
            nonlocal closing
            if not closing:
                closing = True
                if self._session_start_sent:
                    with contextlib.suppress(Exception):
                        await self._ws_send(ws_conn, types.SessionCloseEvent())
            if self._session_started_fut.done() and not self._session_started_fut.cancelled():
                with contextlib.suppress(asyncio.TimeoutError):
                    await asyncio.wait_for(
                        asyncio.shield(self._session_closed_fut), _SESSION_CLOSE_TIMEOUT
                    )
            await ws_conn.close()

        async def _send_task() -> None:
            nonlocal closing
            # instructions, voice and history are immutable once the session starts
            await self._configured.wait()
            if self._closing:
                # closed before the configuration landed; there is no session to start or drain
                closing = True
                await ws_conn.close()
                return
            start = self._session_start_event()
            self._session_start_sent = True
            await self._ws_send(ws_conn, start)

            async for msg in self._msg_ch:
                # the protocol asks for session.started before any audio or command goes out
                if not self._session_started_fut.done():
                    await self._session_started_fut
                await self._ws_send(ws_conn, msg)

            await _close_ws()

        async def _recv_task() -> None:
            while True:
                try:
                    msg = await ws_conn.receive()
                except (aiohttp.ClientError, ConnectionError, asyncio.TimeoutError):
                    raise APIConnectionError("GPT-Live receive failed") from None
                if msg.type in (
                    aiohttp.WSMsgType.CLOSED,
                    aiohttp.WSMsgType.CLOSE,
                    aiohttp.WSMsgType.CLOSING,
                ):
                    if closing or self._session_closed_fut.done():
                        return
                    raise APIConnectionError(
                        f"{self._live_model._provider_label} connection closed unexpectedly"
                    )
                if msg.type != aiohttp.WSMsgType.TEXT:
                    continue

                event = json.loads(msg.data)
                self.emit("openai_server_event_received", event)
                # while closing only the shutdown events matter; the framework has torn down
                if self._closing and event.get("type") not in _CLOSING_EVENTS:
                    continue
                try:
                    self._handle_event(event)
                except Exception as e:
                    if isinstance(e, APIError) and not e.retryable:
                        raise
                    logger.warning(
                        "failed to handle gpt-live event",
                        extra={"type": event.get("type"), "error_type": type(e).__name__},
                    )

        send_task = asyncio.create_task(_send_task(), name="_send_task")
        tasks = [asyncio.create_task(_recv_task(), name="_recv_task"), send_task]
        wait_reconnect_task: asyncio.Task | None = None
        if self._opts.max_session_duration is not None:
            wait_reconnect_task = asyncio.create_task(
                asyncio.sleep(self._opts.max_session_duration), name="_timeout_task"
            )
            tasks.append(wait_reconnect_task)
        try:
            done, _ = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
            for task in done:
                if task != wait_reconnect_task:
                    task.result()
            if wait_reconnect_task in done:
                await utils.aio.cancel_and_wait(send_task)
                await _close_ws()
        finally:
            await utils.aio.cancel_and_wait(*tasks)
            await ws_conn.close()

    async def _ws_send(
        self, ws_conn: aiohttp.ClientWebSocketResponse, event: types.ClientEvent | dict[str, Any]
    ) -> None:
        raw = event if isinstance(event, dict) else event.model_dump(exclude_none=True)
        self.emit("openai_client_event_queued", raw)
        if lk_oai_debug and raw.get("type") != "session.input_audio.append":
            logger.debug("gpt-live client event", extra={"lk.pii.event": raw})
        try:
            await ws_conn.send_str(json.dumps(raw))
        except (aiohttp.ClientError, ConnectionError, asyncio.TimeoutError):
            raise APIConnectionError("GPT-Live send failed") from None

    # inbound events

    def _handle_event(self, event: dict[str, Any]) -> None:
        etype = event.get("type", "")
        if lk_oai_debug and etype != "session.output_audio.delta":
            logger.debug("gpt-live server event", extra={"lk.pii.event": event})

        if etype == "session.started":
            self._handle_session_started(types.SessionStartedEvent.construct(**event))
        elif etype == "session.output_audio.delta":
            self._handle_output_audio_delta(types.OutputAudioDeltaEvent.construct(**event))
        elif etype == "session.output_transcript.delta":
            self._handle_transcript_delta(
                "assistant", types.TranscriptDeltaEvent.construct(**event)
            )
        elif etype == "session.input_transcript.delta":
            self._handle_transcript_delta("user", types.TranscriptDeltaEvent.construct(**event))
        elif etype == "session.delegation.created":
            self._handle_delegation_created(types.SessionDelegationCreatedEvent.construct(**event))
        elif etype == "response.event":
            self._handle_response_event(types.ResponseEventEnvelope.construct(**event))
        elif etype == "session.usage.updated":
            self._handle_session_usage_updated(types.SessionUsageUpdatedEvent.construct(**event))
        elif etype == "session.closed":
            self._handle_session_closed(types.SessionClosedEvent.construct(**event))
        elif etype == "error":
            self._handle_error(types.ErrorEvent.construct(**event).error)
        elif etype in (
            "session.updated",
            "session.input_audio.muted",
            "session.input_audio.unmuted",
            "session.instructions.appended",
            "session.thinking.appended",
            "session.commentary.appended",
        ):
            # Context append receipts arrive at the estimated injection end, not speech end.
            # Nothing waits on these acknowledgments.
            logger.debug(
                "gpt-live acknowledged a command",
                extra={"type": etype, "client_event_id": event.get("client_event_id")},
            )
        elif lk_oai_debug:
            logger.debug("unhandled gpt-live event", extra={"lk.pii.type": etype})

    def _handle_session_started(self, event: types.SessionStartedEvent) -> None:
        self._session_id = event.session.id or self._session_id
        self._num_retries = 0
        if not self._session_started_fut.done():
            self._session_started_fut.set_result(None)

    def _handle_output_audio_delta(self, event: types.OutputAudioDeltaEvent) -> None:
        # every frame is published, silence included: the framework decides what plays
        if not (data := base64.b64decode(event.delta or "")):
            return
        frame = rtc.AudioFrame(
            data=data,
            sample_rate=SAMPLE_RATE,
            num_channels=NUM_CHANNELS,
            samples_per_channel=len(data) // 2,
        )
        with contextlib.suppress(utils.aio.channel.ChanClosed):
            self._audio_ch.send_nowait(llm.DuplexAudioFrame(frame=frame))

    def _handle_transcript_delta(self, role: Role, event: types.TranscriptDeltaEvent) -> None:
        if not event.delta:
            return
        speech = self._speech.get(role)
        # a pause on the model's clock ends the message even when its fragments arrived together
        if (
            speech is not None
            and speech.end_ms is not None
            and event.start_ms is not None
            and event.start_ms - speech.end_ms > _MIN_SILENCE_MS
        ):
            self._end_speech(role)
            speech = None
        if speech is None:
            speech = self._speech[role] = _Speech(message_id=utils.shortuuid("speech_"))
            self._history.insert(
                llm.ChatMessage(
                    id=speech.message_id,
                    role=role,
                    content=[""],
                    # spoken, not typed: what tells transcribed messages from the app's
                    transcript_confidence=1.0 if role == "user" else None,
                )
            )
            if role == "user":
                self.emit("input_speech_started", llm.InputSpeechStartedEvent())
        speech.text += event.delta
        speech.quiet_ms = 0
        if event.end_ms is not None:
            speech.end_ms = max(speech.end_ms or 0, event.end_ms)
        if isinstance(message := self._history.get_by_id(speech.message_id), llm.ChatMessage):
            message.content[0] = speech.text

        if role == "user":
            self.emit(
                "input_audio_transcription_completed",
                llm.InputTranscriptionCompleted(
                    item_id=speech.message_id,
                    transcript=speech.text,
                    is_final=False,
                    # the model answers over the caller, so the turn is stamped where it began
                    turn_started_at=speech.started_at,
                ),
            )
        else:
            self.emit(
                "transcript_delta",
                llm.DuplexOutputTranscriptDelta(
                    text=event.delta, start_ms=event.start_ms, end_ms=event.end_ms
                ),
            )

    def _end_speech(self, role: Role) -> None:
        """Close a speaker's message; for the user this is the end of their turn."""
        if (speech := self._speech.pop(role, None)) is None or role != "user":
            return
        # the final transcript goes out first, so nothing waits for one after the stop
        self.emit(
            "input_audio_transcription_completed",
            llm.InputTranscriptionCompleted(
                item_id=speech.message_id,
                transcript=speech.text,
                is_final=True,
                turn_started_at=speech.started_at,
            ),
        )
        self.emit(
            "input_speech_stopped", llm.InputSpeechStoppedEvent(user_transcription_enabled=False)
        )

    def _handle_delegation_created(self, event: types.SessionDelegationCreatedEvent) -> None:
        delegation = event.delegation
        if not delegation.id:
            logger.warning("gpt-live delegation has no id; nothing can answer it")
        elif delegation.target == "client":
            speech = self._speech.get("user")
            self.emit(
                "delegation_created",
                GPTLiveDelegation(
                    id=delegation.id, pending_transcript=speech.text if speech else ""
                ),
            )

    def _handle_response_event(self, envelope: types.ResponseEventEnvelope) -> None:
        # the inner event carries no response id, so a delegation's responses are followed in
        # sequence: a continuation is the next response.created under the same delegation, and a
        # response the application started itself has a null delegation
        event = envelope.event
        response = event.response
        d_id = envelope.delegation_id

        if event.type == "response.created":
            self._backend_running_responses[d_id] = set()

        elif event.type == "response.output_item.done":
            # only the completed item carries the name, call id and arguments together
            item = event.item
            if item is None or item.type != "function_call":
                return
            if item.status != "completed":
                logger.debug(
                    "gpt-live ignoring incomplete function call",
                    extra={"function": item.name, "status": item.status},
                )
                return
            if not item.call_id or not item.name or item.arguments is None:
                logger.warning(
                    "gpt-live dropping function call with missing fields",
                    extra={"call_id": item.call_id, "function": item.name},
                )
                return
            if (calls := self._backend_running_responses.get(d_id)) is None:
                logger.warning(
                    "gpt-live function call outside a known response",
                    extra={"call_id": item.call_id, "function": item.name, "delegation_id": d_id},
                )
                calls = self._backend_open_calls
            if item.call_id in calls:
                return
            calls.add(item.call_id)

            fnc_call = llm.FunctionCall(
                id=item.id or utils.shortuuid("fc_"),
                call_id=item.call_id,
                name=item.name,
                arguments=item.arguments,
            )
            self._history.insert(fnc_call)
            self.emit("function_call", fnc_call)

        elif event.type == "response.completed":
            if response is not None and (usage := response.usage) is not None:
                # the voice model is billed by duration; these tokens are the backend's and are
                # reported under its name
                self.emit(
                    "metrics_collected",
                    LLMMetrics(
                        label=self._live_model.label,
                        request_id=response.id or "",
                        timestamp=time.time(),
                        duration=0,
                        ttft=-1,
                        cancelled=False,
                        prompt_tokens=usage.input_tokens,
                        prompt_cached_tokens=usage.input_tokens_details.cached_tokens,
                        cache_creation_tokens=usage.input_tokens_details.cache_write_tokens,
                        completion_tokens=usage.output_tokens,
                        reasoning_tokens=usage.output_tokens_details.reasoning_tokens,
                        total_tokens=usage.total_tokens,
                        tokens_per_second=0,
                        metadata=Metadata(
                            model_name=response.model
                            or self._opts.responses.get("model", DEFAULT_BACKEND_MODEL),
                            model_provider=self._live_model.provider,
                        ),
                    ),
                )
            if (calls := self._backend_running_responses.pop(d_id, None)) is not None:
                self._backend_open_calls |= calls
            self._maybe_continue_response()

        elif event.type in ("response.failed", "response.incomplete"):
            logger.warning(
                "gpt-live backend response did not complete",
                extra={
                    "type": event.type,
                    "delegation_id": d_id,
                    "lk.pii.error": response.error if response else None,
                    "lk.pii.incomplete_details": response.incomplete_details if response else None,
                },
            )
            # the service discards a failed response's calls: an output for one is refused, so
            # it goes to the voice model as context instead
            self._backend_running_responses.pop(d_id, None)
            self._maybe_continue_response()

    def _maybe_continue_response(self) -> None:
        # one response.create continues the chain, and only once nothing is still asking and every
        # call in the conversation has its answer: the service refuses a partial batch
        if (
            self._backend_running_responses
            or self._backend_open_calls
            or not self._backend_response_pending
        ):
            return
        self._backend_response_pending = False
        self.send_event(types.ResponseCreateEvent(event_id=utils.shortuuid("response_create_")))

    # metrics and errors

    def _handle_session_usage_updated(self, event: types.SessionUsageUpdatedEvent) -> None:
        if event.context_window is not None and event.context_window.usage_ratio is not None:
            logger.debug(
                "gpt-live context window utilization",
                extra={"usage_ratio": event.context_window.usage_ratio},
            )
        self._handle_usage(event.usage)

    def _handle_session_closed(self, event: types.SessionClosedEvent) -> None:
        logger.debug(
            "gpt-live session closed",
            extra={"reason": event.reason, "session_id": self._session_id},
        )
        self._handle_usage(event.usage)
        if not self._session_closed_fut.done():
            self._session_closed_fut.set_result(None)

    def _handle_usage(self, usage: types.Usage) -> None:
        # reported cumulatively for the whole session, so only the delta goes to the collectors
        previous, self._usage_total = self._usage_total, usage
        self.emit(
            "metrics_collected",
            RealtimeModelMetrics(
                timestamp=time.time(),
                request_id=self._session_id or "",
                ttft=-1,
                duration=0,
                session_duration=max(0.0, usage.seconds - previous.seconds),
                cancelled=False,
                label=self._live_model.label,
                input_tokens=0,
                output_tokens=0,
                total_tokens=0,
                tokens_per_second=0,
                input_token_details=RealtimeModelMetrics.InputTokenDetails(),
                output_token_details=RealtimeModelMetrics.OutputTokenDetails(),
                metadata=Metadata(
                    model_name=self._live_model.model, model_provider=self._live_model.provider
                ),
            ),
        )

    def _handle_error(self, error: types.ErrorBody) -> None:
        logger.error(
            "gpt-live returned an error",
            extra={"lk.pii.error": error.model_dump(exclude_none=True)},
        )
        recoverable = (error.code or error.type or "") not in _FATAL_ERROR_CODES
        api_error = APIError(
            message="GPT-Live returned an error",
            retryable=recoverable,
        )
        if not recoverable:
            raise api_error
        self._emit_error(api_error, recoverable=True)

    def _emit_error(self, error: Exception, recoverable: bool) -> None:
        self.emit(
            "error",
            llm.RealtimeModelError(
                timestamp=time.time(),
                label=self._live_model.label,
                error=error,
                recoverable=recoverable,
            ),
        )

    # DuplexSession interface

    @property
    def session_id(self) -> str | None:
        """The service's id for the current connection."""
        return self._session_id

    @property
    def audio_stream(self) -> AsyncIterable[llm.DuplexAudioFrame]:
        return self._audio_ch

    @property
    def tools(self) -> llm.ToolContext:
        return self._tools.copy()

    def push_audio(self, frame: rtc.AudioFrame) -> None:
        # the caller's turn ends on their own audio: this much pushed since their last fragment
        if (speech := self._speech.get("user")) is not None:
            speech.quiet_ms += round(frame.duration * 1000)
            if speech.quiet_ms >= _MIN_SILENCE_MS:
                self._end_speech("user")

        if self._input_resampler and frame.sample_rate != self._input_resampler._input_rate:
            self._input_resampler = None
        if self._input_resampler is None and (
            frame.sample_rate != SAMPLE_RATE or frame.num_channels != NUM_CHANNELS
        ):
            self._input_resampler = rtc.AudioResampler(
                input_rate=frame.sample_rate, output_rate=SAMPLE_RATE, num_channels=NUM_CHANNELS
            )
        frames = self._input_resampler.push(frame) if self._input_resampler else [frame]
        for f in frames:
            for nf in self._bstream.write(f.data.tobytes()):
                self.send_event(
                    types.InputAudioAppendEvent(audio=base64.b64encode(nf.data).decode("utf-8"))
                )

    def append_instructions(self, text: str, *, delegation_id: str | None = None) -> None:
        """Add a standing rule to the model's instructions, capped at 500 tokens."""
        self._append(types.InstructionsAppendEvent, text, delegation_id)

    def append_thinking(self, text: str, *, delegation_id: str | None = None) -> None:
        """Give the model something to know without saying it, capped at 500 tokens."""
        self._append(types.ThinkingAppendEvent, text, delegation_id)

    def append_commentary(self, text: str, *, delegation_id: str | None = None) -> None:
        """Give the model something to say once, in its own words, capped at 500 tokens.

        Under client delegation this answers a :class:`GPTLiveDelegation`; repeated calls with the
        same ``delegation_id`` continue that work.
        """
        self._append(types.CommentaryAppendEvent, text, delegation_id)

    def _append(
        self,
        event_cls: type[types.InstructionsAppendEvent]
        | type[types.ThinkingAppendEvent]
        | type[types.CommentaryAppendEvent],
        text: str,
        delegation_id: str | None,
    ) -> None:
        self.send_event(
            event_cls(
                event_id=utils.shortuuid("append_"), delegation_id=delegation_id, content=text
            )
        )

    def mute_input(self) -> None:
        """Replace microphone input with silence; the model keeps generating and speaking."""
        self.send_event(types.InputAudioMuteEvent(event_id=utils.shortuuid("mute_")))

    def unmute_input(self) -> None:
        self.send_event(types.InputAudioUnmuteEvent(event_id=utils.shortuuid("unmute_")))

    async def aclose(self) -> None:
        await super().aclose()
        if not self._session_started_fut.done():
            self._session_started_fut.cancel()
        self._msg_ch.close()
        with contextlib.suppress(asyncio.CancelledError):
            await self._main_atask

    # framework hooks

    async def _update_instructions(self, instructions: str) -> None:
        if self._session_start_sent and instructions != self._instructions:
            raise llm.RealtimeError(
                "gpt-live voice instructions are immutable after session start; use "
                "append_instructions for a standing rule"
            )
        self._instructions = instructions

    async def _update_tools(self, tools: list[llm.Tool]) -> None:
        self._tools = llm.ToolContext(tools)
        if self._opts.delegation == "client":
            if tools:
                # dropping them silently leaves an agent whose tools simply never run
                raise llm.RealtimeError(
                    "gpt-live client delegation has no tool channel, so the model can never call "
                    f"{sorted(tool.id for tool in self._tools.flatten())}. Leave the agent's tools "
                    "empty and answer delegation_created with append_commentary, or pass "
                    'delegation="responses" to run tools on the backend model.'
                )
            return
        self._send_delegation_update(
            types.ResponsesConfig(tools=_build_delegation_tools(self._tools.flatten()))
        )

    async def _append_items(self, items: list[llm.ChatItem]) -> None:
        self._history.insert(items)
        if not self._session_start_sent:
            return  # startup history, rendered into session.start

        # a system or developer message is a standing rule for the voice model, a tool result
        # answering a call the backend delegated goes back on the backend's channel, and everything
        # else is context for the voice model, as one append
        backend_outputs: list[llm.FunctionCallOutput] = []
        lines: list[str] = []
        for item in items:
            if isinstance(item, llm.ChatMessage) and item.role in ("system", "developer"):
                if text := item.text_content:
                    self.append_instructions(text)
            elif isinstance(item, llm.FunctionCallOutput) and (
                item.call_id in self._backend_open_calls
                or any(item.call_id in c for c in self._backend_running_responses.values())
            ):
                backend_outputs.append(item)
            elif (rendered := _render_item(item)) is not None:
                lines.append("{}: {}".format(*rendered))

        if lines:
            self.append_thinking("\n".join(lines))

        for output in backend_outputs:
            self.send_event(
                types.ResponseItemCreateEvent(
                    event_id=utils.shortuuid("tool_output_"),
                    item=FunctionCallOutput(
                        type="function_call_output", call_id=output.call_id, output=output.output
                    ),
                )
            )
            self._backend_open_calls.discard(output.call_id)
            for calls in self._backend_running_responses.values():
                calls.discard(output.call_id)
        if backend_outputs:
            if silenced := [o.name or o.call_id for o in backend_outputs if not o.reply_required]:
                logger.warning(
                    "a tool result wants no reply, but GPT Live will answer it anyway: the "
                    "backend has no way to close a call without a spoken continuation, and an "
                    "unanswered call holds every later tool call. Continuing regardless.",
                    extra={"functions": silenced},
                )
            self._backend_response_pending = True
            self._maybe_continue_response()

        # TODO: under client delegation, answer a GPTLiveDelegation handled as a tool call with
        # append_commentary(output, delegation_id=...) here; nothing reaches the model for it yet
        # A manual call to append_commentary() is the only way to answer a GPTLiveDelegation for now

    def _generate_reply(
        self,
        *,
        instructions: NotGivenOr[str] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice] = NOT_GIVEN,
        tools: NotGivenOr[list[llm.Tool]] = NOT_GIVEN,
    ) -> None:
        if is_given(instructions):
            self.append_commentary(f"{_ASK_INSTRUCTED}\n\n{instructions}")
            return
        # a typed message still the newest thing said rides in the ask, once; after speech has
        # moved the conversation on, the ask points at the context instead
        newest = self._history.items[-1] if self._history.items else None
        typed = (
            newest.text_content
            if isinstance(newest, llm.ChatMessage)
            and newest.role == "user"
            and newest.transcript_confidence is None
            and newest.id != self._asked_item_id
            else None
        )
        self._asked_item_id = newest.id if newest is not None else None
        self.append_commentary(f"{_ASK_TYPED}\n\n{typed}" if typed else _ASK_BARE)

    def _update_options(
        self, *, tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN
    ) -> None:
        if is_given(tool_choice):
            self._opts.responses["tool_choice"] = tool_choice
            self._send_delegation_update(
                types.ResponsesConfig(tool_choice=_to_tool_choice(tool_choice))
            )

A session for the OpenAI GPT-Live API (WebSocket), reached with Agent.duplex_session.

Exposes three extra events: - openai_server_event_received: raw server events - openai_client_event_queued: raw client events sent to the server - delegation_created: a :class:GPTLiveDelegation, under client delegation

The append_* methods queue context and return without waiting. Their *.appended events arrive at the estimated context-injection end; they do not mean speech has finished.

Ancestors

  • livekit.agents.llm.duplex.DuplexSession
  • abc.ABC
  • EventEmitter
  • typing.Generic

Instance variables

prop audio_stream : AsyncIterable[llm.DuplexAudioFrame]
Expand source code
@property
def audio_stream(self) -> AsyncIterable[llm.DuplexAudioFrame]:
    return self._audio_ch

The model's output audio, its own silence included; it may stop between bursts.

prop session_id : str | None
Expand source code
@property
def session_id(self) -> str | None:
    """The service's id for the current connection."""
    return self._session_id

The service's id for the current connection.

prop tools : llm.ToolContext
Expand source code
@property
def tools(self) -> llm.ToolContext:
    return self._tools.copy()

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    await super().aclose()
    if not self._session_started_fut.done():
        self._session_started_fut.cancel()
    self._msg_ch.close()
    with contextlib.suppress(asyncio.CancelledError):
        await self._main_atask

Close the session; an override calls super().aclose() first.

It releases a model that waits on _configured, which then reads _closing to see that the configuration was abandoned rather than applied.

def append_commentary(self, text: str, *, delegation_id: str | None = None) ‑> None
Expand source code
def append_commentary(self, text: str, *, delegation_id: str | None = None) -> None:
    """Give the model something to say once, in its own words, capped at 500 tokens.

    Under client delegation this answers a :class:`GPTLiveDelegation`; repeated calls with the
    same ``delegation_id`` continue that work.
    """
    self._append(types.CommentaryAppendEvent, text, delegation_id)

Give the model something to say once, in its own words, capped at 500 tokens.

Under client delegation this answers a :class:GPTLiveDelegation; repeated calls with the same delegation_id continue that work.

def append_instructions(self, text: str, *, delegation_id: str | None = None) ‑> None
Expand source code
def append_instructions(self, text: str, *, delegation_id: str | None = None) -> None:
    """Add a standing rule to the model's instructions, capped at 500 tokens."""
    self._append(types.InstructionsAppendEvent, text, delegation_id)

Add a standing rule to the model's instructions, capped at 500 tokens.

def append_thinking(self, text: str, *, delegation_id: str | None = None) ‑> None
Expand source code
def append_thinking(self, text: str, *, delegation_id: str | None = None) -> None:
    """Give the model something to know without saying it, capped at 500 tokens."""
    self._append(types.ThinkingAppendEvent, text, delegation_id)

Give the model something to know without saying it, capped at 500 tokens.

def mute_input(self) ‑> None
Expand source code
def mute_input(self) -> None:
    """Replace microphone input with silence; the model keeps generating and speaking."""
    self.send_event(types.InputAudioMuteEvent(event_id=utils.shortuuid("mute_")))

Replace microphone input with silence; the model keeps generating and speaking.

def push_audio(self, frame: rtc.AudioFrame) ‑> None
Expand source code
def push_audio(self, frame: rtc.AudioFrame) -> None:
    # the caller's turn ends on their own audio: this much pushed since their last fragment
    if (speech := self._speech.get("user")) is not None:
        speech.quiet_ms += round(frame.duration * 1000)
        if speech.quiet_ms >= _MIN_SILENCE_MS:
            self._end_speech("user")

    if self._input_resampler and frame.sample_rate != self._input_resampler._input_rate:
        self._input_resampler = None
    if self._input_resampler is None and (
        frame.sample_rate != SAMPLE_RATE or frame.num_channels != NUM_CHANNELS
    ):
        self._input_resampler = rtc.AudioResampler(
            input_rate=frame.sample_rate, output_rate=SAMPLE_RATE, num_channels=NUM_CHANNELS
        )
    frames = self._input_resampler.push(frame) if self._input_resampler else [frame]
    for f in frames:
        for nf in self._bstream.write(f.data.tobytes()):
            self.send_event(
                types.InputAudioAppendEvent(audio=base64.b64encode(nf.data).decode("utf-8"))
            )
def send_event(self, event: types.ClientEvent | dict[str, Any]) ‑> None
Expand source code
def send_event(self, event: types.ClientEvent | dict[str, Any]) -> None:
    with contextlib.suppress(utils.aio.channel.ChanClosed):
        self._msg_ch.send_nowait(event)
def unmute_input(self) ‑> None
Expand source code
def unmute_input(self) -> None:
    self.send_event(types.InputAudioUnmuteEvent(event_id=utils.shortuuid("unmute_")))

Inherited members

class InferenceRealtimeModel (model: str,
*,
provider: str | None = None,
base_url: str | None = None,
api_key: str | None = None,
api_secret: str | None = None,
inference_class: InferenceClass | None = None,
voice: NotGivenOr[str] = NOT_GIVEN,
modalities: "NotGivenOr[list[Literal['text', 'audio']]]" = NOT_GIVEN,
input_audio_transcription: NotGivenOr[AudioTranscription | None] = NOT_GIVEN,
input_audio_noise_reduction: NotGivenOr[NoiseReductionType | NoiseReduction | None] = NOT_GIVEN,
turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | None] = NOT_GIVEN,
tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
speed: NotGivenOr[float] = NOT_GIVEN,
tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
http_session: aiohttp.ClientSession | None = None,
max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
conn_options: APIConnectOptions = APIConnectOptions(max_retry=3, retry_interval=2.0, timeout=10.0))
Expand source code
class RealtimeModel(_RealtimeModel):
    """OpenAI-compatible realtime model authenticated through LiveKit Inference."""

    def __init__(
        self,
        model: str,
        *,
        provider: str | None = None,
        base_url: str | None = None,
        api_key: str | None = None,
        api_secret: str | None = None,
        inference_class: InferenceClass | None = None,
        voice: NotGivenOr[str] = NOT_GIVEN,
        modalities: NotGivenOr[list[Literal["text", "audio"]]] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[AudioTranscription | None] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReductionType | NoiseReduction | None
        ] = NOT_GIVEN,
        turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | None] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
    ) -> None:
        if "/" not in model:
            raise ValueError("model must be provider-prefixed, for example 'openai/gpt-realtime'")

        resolved_api_key, resolved_api_secret = resolve_credentials(api_key, api_secret)

        is_xai = model.startswith("xai/")
        resolved_voice = voice if is_given(voice) else "eve" if is_xai else DEFAULT_VOICE
        resolved_transcription = input_audio_transcription
        resolved_turn_detection = turn_detection
        # Preserve whether xAI's server VAD is a default the framework may disable.
        can_disable_turn_detection = not is_given(turn_detection)
        if is_xai:
            if not is_given(resolved_transcription):
                resolved_transcription = _XAI_DEFAULT_INPUT_AUDIO_TRANSCRIPTION
            if not is_given(resolved_turn_detection):
                resolved_turn_detection = _XAI_DEFAULT_TURN_DETECTION

        super().__init__(
            model=model,
            voice=resolved_voice,
            modalities=modalities,
            input_audio_transcription=resolved_transcription,
            input_audio_noise_reduction=input_audio_noise_reduction,
            turn_detection=resolved_turn_detection,
            tool_choice=tool_choice,
            speed=speed,
            tracing=tracing,
            truncation=truncation,
            reasoning=reasoning,
            api_key="livekit-inference",
            base_url=base_url or get_default_inference_url(),
            http_session=http_session,
            max_session_duration=max_session_duration,
            conn_options=conn_options,
        )
        # LiveKit Inference always uses the OpenAI-compatible protocol; ambient Azure
        # settings must not change its URL or session wire format.
        self._opts.is_azure = False
        self._opts.api_version = None
        if is_xai:
            self._capabilities.can_disable_turn_detection = can_disable_turn_detection
        self._inference_opts = _InferenceOptions(
            provider=provider,
            api_key=resolved_api_key,
            api_secret=resolved_api_secret,
            inference_class=inference_class,
        )
        self._provider_label = "LiveKit Inference Realtime"

    @classmethod
    def from_model_string(cls, model: str) -> RealtimeModel:
        """Create a RealtimeModel instance from a model string"""
        return cls(model)

    @property
    def provider(self) -> str:
        return "livekit"

    def session(self, *, turn_detection_disabled: bool = False) -> RealtimeSession:
        sess = RealtimeSession(self, turn_detection_disabled=turn_detection_disabled)
        self._sessions.add(sess)
        return sess

OpenAI-compatible realtime model authenticated through LiveKit Inference.

Initialize a Realtime model client for OpenAI or Azure OpenAI.

Args

model : str
Realtime model name, e.g., "gpt-realtime".
voice : str
Voice used for audio responses. Defaults to "marin".
modalities (list[Literal["text", "audio"]] | NotGiven): Modalities to enable. Defaults to ["text", "audio"] if not provided.
tool_choice : llm.ToolChoice | None | NotGiven
Tool selection policy for responses.
base_url : str | NotGiven
HTTP base URL of the OpenAI/Azure API. If not provided, uses OPENAI_BASE_URL for OpenAI; for Azure, constructed from AZURE_OPENAI_ENDPOINT.
input_audio_transcription : AudioTranscription | None | NotGiven
Options for transcribing input audio.
input_audio_noise_reduction : NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None | NotGiven
Input audio noise reduction settings.
turn_detection : RealtimeAudioInputTurnDetection | None | NotGiven
Server-side turn-detection options.
speed : float | NotGiven
Audio playback speed multiplier.
tracing : Tracing | None | NotGiven
Tracing configuration for OpenAI Realtime.
truncation : RealtimeTruncation | None | NotGiven
Truncation configuration for OpenAI Realtime.
reasoning : RealtimeReasoning | None | NotGiven
Reasoning config for reasoning-capable models (e.g. gpt-realtime-2), e.g. RealtimeReasoning(effort="low").
api_key : str | None
OpenAI API key. If None and not using Azure, read from OPENAI_API_KEY.
http_session : aiohttp.ClientSession | None
Optional shared HTTP session.
azure_deployment : str | None
Azure deployment name. Presence of any Azure-specific option enables Azure mode.
entra_token : str | None
Azure Entra token auth (alternative to api_key).
max_session_duration : float | None | NotGiven
Seconds before recycling the connection.
conn_options : APIConnectOptions
Retry/backoff and connection settings.
temperature : float | NotGiven
Deprecated; ignored by Realtime v1.

Raises

ValueError
If OPENAI_API_KEY is missing in non-Azure mode, or if Azure endpoint cannot be determined when in Azure mode.

Examples

Basic OpenAI usage:

from livekit.plugins.openai.realtime import RealtimeModel
from openai.types import realtime

model = RealtimeModel(
    voice="marin",
    modalities=["audio"],
    input_audio_transcription=realtime.AudioTranscription(
        model="gpt-4o-transcribe",
    ),
    input_audio_noise_reduction="near_field",
    turn_detection=realtime.realtime_audio_input_turn_detection.SemanticVad(
        type="semantic_vad",
        create_response=True,
        eagerness="auto",
        interrupt_response=True,
    ),
)
session = AgentSession(llm=model)

Ancestors

  • livekit.agents.llm._realtime.openai.RealtimeModel
  • livekit.agents.llm.realtime.RealtimeModel

Static methods

def from_model_string(model: str) ‑> RealtimeModel

Create a RealtimeModel instance from a model string

Instance variables

prop provider : str
Expand source code
@property
def provider(self) -> str:
    return "livekit"

Methods

def session(self, *, turn_detection_disabled: bool = False) ‑> RealtimeSession
Expand source code
def session(self, *, turn_detection_disabled: bool = False) -> RealtimeSession:
    sess = RealtimeSession(self, turn_detection_disabled=turn_detection_disabled)
    self._sessions.add(sess)
    return sess

Create a new session, optionally with server-side turn detection disabled.

turn_detection_disabled is honored only by plugins reporting can_disable_turn_detection; the model itself is left unchanged and reusable.

class RealtimeModel (*,
model: str = 'gpt-realtime',
voice: str = 'marin',
modalities: "NotGivenOr[list[Literal['text', 'audio']]]" = NOT_GIVEN,
tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
base_url: NotGivenOr[str] = NOT_GIVEN,
input_audio_transcription: NotGivenOr[AudioTranscription | InputAudioTranscription | None] = NOT_GIVEN,
input_audio_noise_reduction: NotGivenOr[NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None] = NOT_GIVEN,
turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | TurnDetection | None] = NOT_GIVEN,
speed: NotGivenOr[float] = NOT_GIVEN,
tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
api_key: str | None = None,
http_session: aiohttp.ClientSession | None = None,
azure_deployment: str | None = None,
entra_token: str | None = None,
max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
conn_options: APIConnectOptions = APIConnectOptions(max_retry=3, retry_interval=2.0, timeout=10.0),
temperature: NotGivenOr[float] = NOT_GIVEN,
**kwargs: Any)
Expand source code
class RealtimeModel(llm.RealtimeModel):
    @overload
    def __init__(
        self,
        *,
        model: RealtimeModels | str = "gpt-realtime",
        voice: str = DEFAULT_VOICE,
        modalities: NotGivenOr[list[Literal["text", "audio"]]] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[
            AudioTranscription | InputAudioTranscription | None
        ] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None
        ] = NOT_GIVEN,
        turn_detection: NotGivenOr[
            RealtimeAudioInputTurnDetection | TurnDetection | None
        ] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        api_key: str | None = None,
        base_url: NotGivenOr[str] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
        temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
    ) -> None: ...

    @overload
    def __init__(
        self,
        *,
        azure_deployment: str | None = None,
        entra_token: str | None = None,
        api_key: str | None = None,
        api_version: str | None = None,
        base_url: NotGivenOr[str] = NOT_GIVEN,
        voice: str = DEFAULT_VOICE,
        modalities: NotGivenOr[list[Literal["text", "audio"]]] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[
            AudioTranscription | InputAudioTranscription | None
        ] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None
        ] = NOT_GIVEN,
        turn_detection: NotGivenOr[
            RealtimeAudioInputTurnDetection | TurnDetection | None
        ] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
        temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
    ) -> None: ...

    def __init__(
        self,
        *,
        model: str = "gpt-realtime",
        voice: str = DEFAULT_VOICE,
        modalities: NotGivenOr[list[Literal["text", "audio"]]] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        base_url: NotGivenOr[str] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[
            AudioTranscription | InputAudioTranscription | None
        ] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None
        ] = NOT_GIVEN,
        turn_detection: NotGivenOr[
            RealtimeAudioInputTurnDetection | TurnDetection | None
        ] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        api_key: str | None = None,
        http_session: aiohttp.ClientSession | None = None,
        azure_deployment: str | None = None,
        entra_token: str | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
        temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
        **kwargs: Any,
    ) -> None:
        """
        Initialize a Realtime model client for OpenAI or Azure OpenAI.

        Args:
            model (str): Realtime model name, e.g., "gpt-realtime".
            voice (str): Voice used for audio responses. Defaults to "marin".
            modalities (list[Literal["text", "audio"]] | NotGiven): Modalities to enable. Defaults to ["text", "audio"] if not provided.
            tool_choice (llm.ToolChoice | None | NotGiven): Tool selection policy for responses.
            base_url (str | NotGiven): HTTP base URL of the OpenAI/Azure API. If not provided, uses OPENAI_BASE_URL for OpenAI; for Azure, constructed from AZURE_OPENAI_ENDPOINT.
            input_audio_transcription (AudioTranscription | None | NotGiven): Options for transcribing input audio.
            input_audio_noise_reduction (NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None | NotGiven): Input audio noise reduction settings.
            turn_detection (RealtimeAudioInputTurnDetection | None | NotGiven): Server-side turn-detection options.
            speed (float | NotGiven): Audio playback speed multiplier.
            tracing (Tracing | None | NotGiven): Tracing configuration for OpenAI Realtime.
            truncation (RealtimeTruncation | None | NotGiven): Truncation configuration for OpenAI Realtime.
            reasoning (RealtimeReasoning | None | NotGiven): Reasoning config for reasoning-capable models (e.g. ``gpt-realtime-2``), e.g. ``RealtimeReasoning(effort="low")``.
            api_key (str | None): OpenAI API key. If None and not using Azure, read from OPENAI_API_KEY.
            http_session (aiohttp.ClientSession | None): Optional shared HTTP session.
            azure_deployment (str | None): Azure deployment name. Presence of any Azure-specific option enables Azure mode.
            entra_token (str | None): Azure Entra token auth (alternative to api_key).
            max_session_duration (float | None | NotGiven): Seconds before recycling the connection.
            conn_options (APIConnectOptions): Retry/backoff and connection settings.
            temperature (float | NotGiven): Deprecated; ignored by Realtime v1.

        Raises:
            ValueError: If OPENAI_API_KEY is missing in non-Azure mode, or if Azure endpoint cannot be determined when in Azure mode.

        Examples:
            Basic OpenAI usage:

            ```python
            from livekit.plugins.openai.realtime import RealtimeModel
            from openai.types import realtime

            model = RealtimeModel(
                voice="marin",
                modalities=["audio"],
                input_audio_transcription=realtime.AudioTranscription(
                    model="gpt-4o-transcribe",
                ),
                input_audio_noise_reduction="near_field",
                turn_detection=realtime.realtime_audio_input_turn_detection.SemanticVad(
                    type="semantic_vad",
                    create_response=True,
                    eagerness="auto",
                    interrupt_response=True,
                ),
            )
            session = AgentSession(llm=model)
            ```
        """
        api_version: str | None = kwargs.get("api_version") or os.getenv("OPENAI_API_VERSION")
        if kwargs.get("api_version"):
            logger.warning(
                "The `api_version` parameter is deprecated and will be removed on April 30, 2026."
            )
        elif os.getenv("OPENAI_API_VERSION"):
            logger.warning(
                "The OPENAI_API_VERSION environment variable is deprecated and will be removed "
                "on April 30, 2026."
            )

        modalities = modalities if is_given(modalities) else ["text", "audio"]
        resolved_turn_detection = to_turn_detection(turn_detection)
        _warn_on_half_disabled_turn_taking(resolved_turn_detection)
        super().__init__(
            capabilities=llm.RealtimeCapabilities(
                message_truncation=True,
                turn_detection=_server_turn_taking_enabled(resolved_turn_detection),
                can_disable_turn_detection=not is_given(turn_detection),
                user_transcription=input_audio_transcription is not None,
                auto_tool_reply_generation=False,
                audio_output="audio" in modalities,
                manual_function_calls=True,
                mutable_chat_context=True,
                mutable_instructions=True,
                mutable_tools=True,
                per_response_tool_choice=True,
            )
        )
        if type(self) is RealtimeModel:
            # Preserve the pre-move metrics label for direct OpenAI sessions.
            self._label = "livekit.plugins.openai.realtime.realtime_model.RealtimeModel"

        is_azure = (
            api_version is not None or entra_token is not None or azure_deployment is not None
        )

        api_key = api_key or os.environ.get("OPENAI_API_KEY")
        if api_key is None and not is_azure:
            raise ValueError(
                "The api_key client option must be set either by passing api_key "
                "to the client or by setting the OPENAI_API_KEY environment variable"
            )

        if is_given(base_url):
            base_url_val = base_url
        else:
            if is_azure:
                azure_endpoint = os.getenv("AZURE_OPENAI_ENDPOINT")
                if azure_endpoint is None:
                    raise ValueError(
                        "Missing Azure endpoint. Please pass base_url "
                        "or set AZURE_OPENAI_ENDPOINT environment variable."
                    )
                base_url_val = f"{azure_endpoint.rstrip('/')}/openai"
            else:
                base_url_val = os.getenv("OPENAI_BASE_URL", OPENAI_BASE_URL)

        self._opts = _RealtimeOptions(
            model=model,
            voice=voice,
            tool_choice=tool_choice or None,
            modalities=modalities,
            input_audio_transcription=to_audio_transcription(input_audio_transcription),
            input_audio_noise_reduction=to_noise_reduction(input_audio_noise_reduction),
            turn_detection=resolved_turn_detection,
            api_key=api_key,
            base_url=base_url_val,
            is_azure=is_azure,
            azure_deployment=azure_deployment,
            entra_token=entra_token,
            api_version=api_version,
            max_response_output_tokens=DEFAULT_MAX_RESPONSE_OUTPUT_TOKENS,  # type: ignore
            speed=speed if is_given(speed) else 1.0,
            tracing=tracing if is_given(tracing) else None,
            truncation=truncation if is_given(truncation) else None,
            reasoning=reasoning if is_given(reasoning) else None,
            max_session_duration=max_session_duration
            if is_given(max_session_duration)
            else DEFAULT_MAX_SESSION_DURATION,
            conn_options=conn_options,
        )
        self._http_session = http_session
        self._http_session_owned = False
        self._sessions = weakref.WeakSet[RealtimeSession]()
        self._provider_label = "OpenAI Realtime API"

    @property
    def model(self) -> str:
        return self._opts.model

    @property
    def provider(self) -> str:
        from urllib.parse import urlparse

        return urlparse(self._opts.base_url).netloc

    @classmethod
    def with_azure(
        cls,
        *,
        azure_deployment: str,
        azure_endpoint: str | None = None,
        api_key: str | None = None,
        entra_token: str | None = None,
        base_url: str | None = None,
        voice: str = DEFAULT_VOICE,
        modalities: NotGivenOr[list[Literal["text", "audio"]]] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[
            AudioTranscription | InputAudioTranscription | None
        ] = NOT_GIVEN,
        input_audio_noise_reduction: NoiseReductionType | InputAudioNoiseReduction | None = None,
        turn_detection: NotGivenOr[
            RealtimeAudioInputTurnDetection | TurnDetection | None
        ] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        http_session: aiohttp.ClientSession | None = None,
        max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
        temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
        **kwargs: Any,
    ) -> RealtimeModel:
        """
        Create a RealtimeModel configured for Azure OpenAI.

        Args:
            azure_deployment (str): Azure OpenAI deployment name.
            azure_endpoint (str | None): Azure endpoint URL; if None, taken from AZURE_OPENAI_ENDPOINT.
            api_key (str | None): Azure API key; if None, taken from AZURE_OPENAI_API_KEY. Omit if using `entra_token`.
            entra_token (str | None): Azure Entra token for AAD auth. Provide instead of `api_key`.
            base_url (str | None): Explicit base URL. Mutually exclusive with `azure_endpoint`. If provided, used as-is.
            voice (str): Voice used for audio responses.
            modalities (list[Literal["text", "audio"]] | NotGiven): Modalities to enable. Defaults to ["text", "audio"] if not provided.
            input_audio_transcription (AudioTranscription | InputAudioTranscription | None | NotGiven): Transcription options; defaults to Azure-optimized values when not provided.
            input_audio_noise_reduction (NoiseReductionType | InputAudioNoiseReduction | None): Input noise reduction settings. Defaults to None.
            turn_detection (RealtimeAudioInputTurnDetection | TurnDetection | None | NotGiven): Server-side VAD; defaults to Azure-optimized values when not provided.
            speed (float | NotGiven): Audio playback speed multiplier.
            tracing (Tracing | None | NotGiven): Tracing configuration for OpenAI Realtime.
            reasoning (RealtimeReasoning | None | NotGiven): Reasoning config for reasoning-capable models, e.g. ``RealtimeReasoning(effort="low")``.
            http_session (aiohttp.ClientSession | None): Optional shared HTTP session.
            max_session_duration (float | None | NotGiven): Seconds before recycling the connection.
            temperature (float | NotGiven): Deprecated; ignored by Realtime v1.

        Returns:
            RealtimeModel: Configured client for Azure OpenAI Realtime.

        Raises:
            ValueError: If credentials are missing, Azure endpoint cannot be determined, or both `base_url` and `azure_endpoint` are provided.

        Examples:
            Azure usage with api-version 2024-10-01-preview:

            ```python
            from livekit.plugins.openai.realtime import RealtimeModel
            from openai.types.beta import realtime

            model = openai.realtime.RealtimeModel.with_azure(
                azure_deployment="gpt-realtime",
                azure_endpoint="https://yourendpoint.azure.com",
                api_version="2024-10-01-preview",
                api_key="your-api-key",
                modalities=["text", "audio"],
                input_audio_transcription=realtime.session.InputAudioTranscription(
                    model="gpt-4o-transcribe",
                ),
                input_audio_noise_reduction=realtime.session.InputAudioNoiseReduction(
                    type="near_field",
                ),
                turn_detection=realtime.session.TurnDetection(
                    type="semantic_vad",
                    create_response=True,
                    eagerness="auto",
                    interrupt_response=True,
                ),
            )
            ```

            Azure usage with api-version 2025-08-28:
            ```python
            from livekit.plugins.openai.realtime import RealtimeModel
            from openai.types import realtime

            model = RealtimeModel(
                azure_deployment="gpt-realtime",
                azure_endpoint="https://yourendpoint.azure.com",
                api_version="2024-10-01-preview",
                api_key="your-api-key",
                input_audio_transcription=realtime.AudioTranscription(
                    model="gpt-4o-transcribe",
                ),
                input_audio_noise_reduction="near_field",
                turn_detection=realtime.realtime_audio_input_turn_detection.SemanticVad(
                    type="semantic_vad",
                    create_response=True,
                    eagerness="auto",
                    interrupt_response=True,
                ),
            )
            ```
        """
        if kwargs.get("api_version"):
            logger.warning(
                "The `api_version` parameter in `with_azure` is deprecated and will be removed "
                "on April 30, 2026."
            )
        elif os.getenv("OPENAI_API_VERSION"):
            logger.warning(
                "The OPENAI_API_VERSION environment variable is deprecated and will be removed "
                "on April 30, 2026."
            )

        api_key = api_key or os.getenv("AZURE_OPENAI_API_KEY")
        if api_key is None and entra_token is None:
            raise ValueError(
                "Missing credentials. Please pass one of `api_key`, `entra_token`, "
                "or the `AZURE_OPENAI_API_KEY` environment variable."
            )

        api_version: str | None = kwargs.get("api_version") or os.getenv("OPENAI_API_VERSION")

        if base_url is None:
            azure_endpoint = azure_endpoint or os.getenv("AZURE_OPENAI_ENDPOINT")
            if azure_endpoint is None:
                raise ValueError(
                    "Missing Azure endpoint. Please pass the `azure_endpoint` "
                    "parameter or set the `AZURE_OPENAI_ENDPOINT` environment variable."
                )

            base_url = f"{azure_endpoint.rstrip('/')}/openai"
        elif azure_endpoint is not None:
            raise ValueError("base_url and azure_endpoint are mutually exclusive")

        if not is_given(input_audio_transcription):
            input_audio_transcription = AZURE_DEFAULT_INPUT_AUDIO_TRANSCRIPTION

        # capture intent before applying the azure default, so the framework can still
        # auto-disable server-side turn detection when the user didn't configure it
        can_disable_turn_detection = not is_given(turn_detection)
        if not is_given(turn_detection):
            turn_detection = AZURE_DEFAULT_TURN_DETECTION

        model = RealtimeModel(
            voice=voice,
            modalities=modalities,
            input_audio_transcription=input_audio_transcription,
            input_audio_noise_reduction=input_audio_noise_reduction,
            turn_detection=turn_detection,
            speed=speed,
            tracing=tracing,
            reasoning=reasoning,
            api_key=api_key,
            http_session=http_session,
            azure_deployment=azure_deployment,
            api_version=api_version,
            entra_token=entra_token,
            base_url=base_url,
            max_session_duration=max_session_duration,
        )
        model._capabilities.can_disable_turn_detection = can_disable_turn_detection
        return model

    def update_options(
        self,
        *,
        voice: NotGivenOr[str] = NOT_GIVEN,
        turn_detection: NotGivenOr[
            RealtimeAudioInputTurnDetection | TurnDetection | None
        ] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[
            InputAudioTranscription | AudioTranscription | None
        ] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReduction | NoiseReductionType | InputAudioNoiseReduction | None
        ] = NOT_GIVEN,
        max_response_output_tokens: NotGivenOr[int | Literal["inf"] | None] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
        temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
    ) -> None:
        if is_given(voice):
            self._opts.voice = voice

        if is_given(turn_detection):
            # a derived capability has to follow the option it is derived from
            self._opts.turn_detection = to_turn_detection(turn_detection)
            self._capabilities.turn_detection = _server_turn_taking_enabled(
                self._opts.turn_detection
            )
            # only the model warns: it re-runs the update on every session it owns
            _warn_on_half_disabled_turn_taking(self._opts.turn_detection)

        if is_given(tool_choice):
            self._opts.tool_choice = tool_choice

        if is_given(input_audio_transcription):
            self._opts.input_audio_transcription = to_audio_transcription(input_audio_transcription)
            self._capabilities.user_transcription = self._opts.input_audio_transcription is not None

        if is_given(input_audio_noise_reduction):
            self._opts.input_audio_noise_reduction = to_noise_reduction(input_audio_noise_reduction)

        if is_given(max_response_output_tokens):
            self._opts.max_response_output_tokens = max_response_output_tokens

        if is_given(speed):
            self._opts.speed = speed

        if is_given(tracing):
            self._opts.tracing = tracing

        if is_given(truncation):
            self._opts.truncation = truncation

        if is_given(reasoning):
            self._opts.reasoning = reasoning

        for sess in self._sessions:
            sess.update_options(
                voice=voice,
                # only propagate when the caller set it, so a session that opted out of
                # server-side turn detection isn't force-synced back on by an unrelated update
                turn_detection=self._opts.turn_detection if is_given(turn_detection) else NOT_GIVEN,
                tool_choice=tool_choice,
                input_audio_transcription=self._opts.input_audio_transcription,
                input_audio_noise_reduction=self._opts.input_audio_noise_reduction,
                max_response_output_tokens=max_response_output_tokens,
                speed=speed,
                tracing=tracing,
                truncation=truncation,
                reasoning=reasoning,
            )

    def _ensure_http_session(self) -> aiohttp.ClientSession:
        if not self._http_session:
            try:
                self._http_session = utils.http_context.http_session()
            except RuntimeError:
                self._http_session = aiohttp.ClientSession()
                self._http_session_owned = True

        return self._http_session

    def session(self, *, turn_detection_disabled: bool = False) -> RealtimeSession:
        sess = RealtimeSession(self, turn_detection_disabled=turn_detection_disabled)
        self._sessions.add(sess)
        return sess

    async def aclose(self) -> None:
        if self._http_session_owned and self._http_session:
            await self._http_session.close()

Initialize a Realtime model client for OpenAI or Azure OpenAI.

Args

model : str
Realtime model name, e.g., "gpt-realtime".
voice : str
Voice used for audio responses. Defaults to "marin".
modalities (list[Literal["text", "audio"]] | NotGiven): Modalities to enable. Defaults to ["text", "audio"] if not provided.
tool_choice : llm.ToolChoice | None | NotGiven
Tool selection policy for responses.
base_url : str | NotGiven
HTTP base URL of the OpenAI/Azure API. If not provided, uses OPENAI_BASE_URL for OpenAI; for Azure, constructed from AZURE_OPENAI_ENDPOINT.
input_audio_transcription : AudioTranscription | None | NotGiven
Options for transcribing input audio.
input_audio_noise_reduction : NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None | NotGiven
Input audio noise reduction settings.
turn_detection : RealtimeAudioInputTurnDetection | None | NotGiven
Server-side turn-detection options.
speed : float | NotGiven
Audio playback speed multiplier.
tracing : Tracing | None | NotGiven
Tracing configuration for OpenAI Realtime.
truncation : RealtimeTruncation | None | NotGiven
Truncation configuration for OpenAI Realtime.
reasoning : RealtimeReasoning | None | NotGiven
Reasoning config for reasoning-capable models (e.g. gpt-realtime-2), e.g. RealtimeReasoning(effort="low").
api_key : str | None
OpenAI API key. If None and not using Azure, read from OPENAI_API_KEY.
http_session : aiohttp.ClientSession | None
Optional shared HTTP session.
azure_deployment : str | None
Azure deployment name. Presence of any Azure-specific option enables Azure mode.
entra_token : str | None
Azure Entra token auth (alternative to api_key).
max_session_duration : float | None | NotGiven
Seconds before recycling the connection.
conn_options : APIConnectOptions
Retry/backoff and connection settings.
temperature : float | NotGiven
Deprecated; ignored by Realtime v1.

Raises

ValueError
If OPENAI_API_KEY is missing in non-Azure mode, or if Azure endpoint cannot be determined when in Azure mode.

Examples

Basic OpenAI usage:

from livekit.plugins.openai.realtime import RealtimeModel
from openai.types import realtime

model = RealtimeModel(
    voice="marin",
    modalities=["audio"],
    input_audio_transcription=realtime.AudioTranscription(
        model="gpt-4o-transcribe",
    ),
    input_audio_noise_reduction="near_field",
    turn_detection=realtime.realtime_audio_input_turn_detection.SemanticVad(
        type="semantic_vad",
        create_response=True,
        eagerness="auto",
        interrupt_response=True,
    ),
)
session = AgentSession(llm=model)

Ancestors

  • livekit.agents.llm.realtime.RealtimeModel

Subclasses

Static methods

def with_azure(*,
azure_deployment: str,
azure_endpoint: str | None = None,
api_key: str | None = None,
entra_token: str | None = None,
base_url: str | None = None,
voice: str = 'marin',
modalities: "NotGivenOr[list[Literal['text', 'audio']]]" = NOT_GIVEN,
input_audio_transcription: NotGivenOr[AudioTranscription | InputAudioTranscription | None] = NOT_GIVEN,
input_audio_noise_reduction: NoiseReductionType | InputAudioNoiseReduction | None = None,
turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | TurnDetection | None] = NOT_GIVEN,
speed: NotGivenOr[float] = NOT_GIVEN,
tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
http_session: aiohttp.ClientSession | None = None,
max_session_duration: NotGivenOr[float | None] = NOT_GIVEN,
temperature: NotGivenOr[float] = NOT_GIVEN,
**kwargs: Any) ‑> livekit.agents.llm._realtime.openai.RealtimeModel

Create a RealtimeModel configured for Azure OpenAI.

Args

azure_deployment : str
Azure OpenAI deployment name.
azure_endpoint : str | None
Azure endpoint URL; if None, taken from AZURE_OPENAI_ENDPOINT.
api_key : str | None
Azure API key; if None, taken from AZURE_OPENAI_API_KEY. Omit if using entra_token.
entra_token : str | None
Azure Entra token for AAD auth. Provide instead of api_key.
base_url : str | None
Explicit base URL. Mutually exclusive with azure_endpoint. If provided, used as-is.
voice : str
Voice used for audio responses.
modalities (list[Literal["text", "audio"]] | NotGiven): Modalities to enable. Defaults to ["text", "audio"] if not provided.
input_audio_transcription : AudioTranscription | InputAudioTranscription | None | NotGiven
Transcription options; defaults to Azure-optimized values when not provided.
input_audio_noise_reduction : NoiseReductionType | InputAudioNoiseReduction | None
Input noise reduction settings. Defaults to None.
turn_detection : RealtimeAudioInputTurnDetection | TurnDetection | None | NotGiven
Server-side VAD; defaults to Azure-optimized values when not provided.
speed : float | NotGiven
Audio playback speed multiplier.
tracing : Tracing | None | NotGiven
Tracing configuration for OpenAI Realtime.
reasoning : RealtimeReasoning | None | NotGiven
Reasoning config for reasoning-capable models, e.g. RealtimeReasoning(effort="low").
http_session : aiohttp.ClientSession | None
Optional shared HTTP session.
max_session_duration : float | None | NotGiven
Seconds before recycling the connection.
temperature : float | NotGiven
Deprecated; ignored by Realtime v1.

Returns

RealtimeModel
Configured client for Azure OpenAI Realtime.

Raises

ValueError
If credentials are missing, Azure endpoint cannot be determined, or both base_url and azure_endpoint are provided.

Examples

Azure usage with api-version 2024-10-01-preview:

from livekit.plugins.openai.realtime import RealtimeModel
from openai.types.beta import realtime

model = openai.realtime.RealtimeModel.with_azure(
    azure_deployment="gpt-realtime",
    azure_endpoint="https://yourendpoint.azure.com",
    api_version="2024-10-01-preview",
    api_key="your-api-key",
    modalities=["text", "audio"],
    input_audio_transcription=realtime.session.InputAudioTranscription(
        model="gpt-4o-transcribe",
    ),
    input_audio_noise_reduction=realtime.session.InputAudioNoiseReduction(
        type="near_field",
    ),
    turn_detection=realtime.session.TurnDetection(
        type="semantic_vad",
        create_response=True,
        eagerness="auto",
        interrupt_response=True,
    ),
)

Azure usage with api-version 2025-08-28:

from livekit.plugins.openai.realtime import RealtimeModel
from openai.types import realtime

model = RealtimeModel(
    azure_deployment="gpt-realtime",
    azure_endpoint="https://yourendpoint.azure.com",
    api_version="2024-10-01-preview",
    api_key="your-api-key",
    input_audio_transcription=realtime.AudioTranscription(
        model="gpt-4o-transcribe",
    ),
    input_audio_noise_reduction="near_field",
    turn_detection=realtime.realtime_audio_input_turn_detection.SemanticVad(
        type="semantic_vad",
        create_response=True,
        eagerness="auto",
        interrupt_response=True,
    ),
)

Instance variables

prop model : str
Expand source code
@property
def model(self) -> str:
    return self._opts.model
prop provider : str
Expand source code
@property
def provider(self) -> str:
    from urllib.parse import urlparse

    return urlparse(self._opts.base_url).netloc

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    if self._http_session_owned and self._http_session:
        await self._http_session.close()
def session(self, *, turn_detection_disabled: bool = False) ‑> livekit.agents.llm._realtime.openai.RealtimeSession
Expand source code
def session(self, *, turn_detection_disabled: bool = False) -> RealtimeSession:
    sess = RealtimeSession(self, turn_detection_disabled=turn_detection_disabled)
    self._sessions.add(sess)
    return sess

Create a new session, optionally with server-side turn detection disabled.

turn_detection_disabled is honored only by plugins reporting can_disable_turn_detection; the model itself is left unchanged and reusable.

def update_options(self,
*,
voice: NotGivenOr[str] = NOT_GIVEN,
turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | TurnDetection | None] = NOT_GIVEN,
tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
input_audio_transcription: NotGivenOr[InputAudioTranscription | AudioTranscription | None] = NOT_GIVEN,
input_audio_noise_reduction: NotGivenOr[NoiseReduction | NoiseReductionType | InputAudioNoiseReduction | None] = NOT_GIVEN,
max_response_output_tokens: "NotGivenOr[int | Literal['inf'] | None]" = NOT_GIVEN,
speed: NotGivenOr[float] = NOT_GIVEN,
tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
temperature: NotGivenOr[float] = NOT_GIVEN) ‑> None
Expand source code
def update_options(
    self,
    *,
    voice: NotGivenOr[str] = NOT_GIVEN,
    turn_detection: NotGivenOr[
        RealtimeAudioInputTurnDetection | TurnDetection | None
    ] = NOT_GIVEN,
    tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
    input_audio_transcription: NotGivenOr[
        InputAudioTranscription | AudioTranscription | None
    ] = NOT_GIVEN,
    input_audio_noise_reduction: NotGivenOr[
        NoiseReduction | NoiseReductionType | InputAudioNoiseReduction | None
    ] = NOT_GIVEN,
    max_response_output_tokens: NotGivenOr[int | Literal["inf"] | None] = NOT_GIVEN,
    speed: NotGivenOr[float] = NOT_GIVEN,
    tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
    truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
    reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
    temperature: NotGivenOr[float] = NOT_GIVEN,  # deprecated, unused in v1
) -> None:
    if is_given(voice):
        self._opts.voice = voice

    if is_given(turn_detection):
        # a derived capability has to follow the option it is derived from
        self._opts.turn_detection = to_turn_detection(turn_detection)
        self._capabilities.turn_detection = _server_turn_taking_enabled(
            self._opts.turn_detection
        )
        # only the model warns: it re-runs the update on every session it owns
        _warn_on_half_disabled_turn_taking(self._opts.turn_detection)

    if is_given(tool_choice):
        self._opts.tool_choice = tool_choice

    if is_given(input_audio_transcription):
        self._opts.input_audio_transcription = to_audio_transcription(input_audio_transcription)
        self._capabilities.user_transcription = self._opts.input_audio_transcription is not None

    if is_given(input_audio_noise_reduction):
        self._opts.input_audio_noise_reduction = to_noise_reduction(input_audio_noise_reduction)

    if is_given(max_response_output_tokens):
        self._opts.max_response_output_tokens = max_response_output_tokens

    if is_given(speed):
        self._opts.speed = speed

    if is_given(tracing):
        self._opts.tracing = tracing

    if is_given(truncation):
        self._opts.truncation = truncation

    if is_given(reasoning):
        self._opts.reasoning = reasoning

    for sess in self._sessions:
        sess.update_options(
            voice=voice,
            # only propagate when the caller set it, so a session that opted out of
            # server-side turn detection isn't force-synced back on by an unrelated update
            turn_detection=self._opts.turn_detection if is_given(turn_detection) else NOT_GIVEN,
            tool_choice=tool_choice,
            input_audio_transcription=self._opts.input_audio_transcription,
            input_audio_noise_reduction=self._opts.input_audio_noise_reduction,
            max_response_output_tokens=max_response_output_tokens,
            speed=speed,
            tracing=tracing,
            truncation=truncation,
            reasoning=reasoning,
        )
class RealtimeSession (realtime_model: RealtimeModel,
*,
turn_detection_disabled: bool = False)
Expand source code
class RealtimeSession(
    llm.RealtimeSession[Literal["openai_server_event_received", "openai_client_event_queued"]]
):
    """
    A session for the OpenAI Realtime API.

    This class is used to interact with the OpenAI Realtime API.
    It is responsible for sending events to the OpenAI Realtime API and receiving events from it.

    It exposes two more events:
    - openai_server_event_received: expose the raw server events from the OpenAI Realtime API
    - openai_client_event_queued: expose the raw client events sent to the OpenAI Realtime API
    """

    def __init__(
        self, realtime_model: RealtimeModel, *, turn_detection_disabled: bool = False
    ) -> None:
        super().__init__(realtime_model)
        self._realtime_model: RealtimeModel = realtime_model
        # per-session copy of opts so update_options can diff against session's own state
        self._opts = replace(
            realtime_model._opts,
            turn_detection=None if turn_detection_disabled else realtime_model._opts.turn_detection,
        )
        # this session's own copy: turn detection can be off here and on for the model
        self._capabilities = replace(
            realtime_model.capabilities,
            turn_detection=False
            if turn_detection_disabled
            else realtime_model.capabilities.turn_detection,
        )
        self._tools = llm.ToolContext.empty()
        self._msg_ch = utils.aio.Chan[RealtimeClientEvent | dict[str, Any]]()
        self._input_resampler: rtc.AudioResampler | None = None

        self._instructions: str | None = None
        # set on aclose; trailing server events are ignored while it's set
        self._closing = False
        self._main_atask = asyncio.create_task(self._main_task(), name="RealtimeSession._main_task")
        self.send_event(self._create_session_update_event())

        self._response_created_futures: dict[str, asyncio.Future[llm.GenerationCreatedEvent]] = {}
        self._item_delete_future: dict[str, asyncio.Future] = {}
        self._item_create_future: dict[str, asyncio.Future] = {}

        # future per in-flight chat ctx event, so a rejection settles the one it answers
        self._chat_ctx_event_futures: dict[str, asyncio.Future] = {}

        # generate_reply event_ids cancelled or timed out before response.created arrived; the
        # response is cancelled by id and discarded when it finally arrives
        self._discarded_event_ids: set[str] = set()

        self._reset_input_turn_state()

        self._current_generation: _ResponseGeneration | _DiscardedGeneration | None = None
        self._remote_chat_ctx = llm.remote_chat_context.RemoteChatContext()

        self._update_chat_ctx_lock = asyncio.Lock()
        self._update_fnc_ctx_lock = asyncio.Lock()

        # 100ms chunks
        self._bstream = utils.audio.AudioByteStream(
            SAMPLE_RATE, NUM_CHANNELS, samples_per_channel=SAMPLE_RATE // 10
        )
        self._pushed_duration_s: float = 0  # duration of audio pushed to the OpenAI Realtime API

    def send_event(self, event: RealtimeClientEvent | dict[str, Any]) -> None:
        with contextlib.suppress(utils.aio.channel.ChanClosed):
            self._msg_ch.send_nowait(event)

    def _reset_input_turn_state(self) -> None:
        """Per-turn input state, keyed by item_id and valid only within one connection.

        Every field here must be discarded on reconnect: the server assigns new item ids,
        so a stale entry can never be matched again.
        """
        # accumulates partial input-audio transcripts per (item_id, content_index)
        self._input_transcript_accumulators: dict[str, dict[int, str]] = {}

        # when serverside VAD detected speech onset, per item_id. Correlating through the
        # item keeps each turn paired with its own start; a single "last speech started"
        # value cannot, because a late transcript would consume the next turn's value.
        self._input_speech_started_at: dict[str, float] = {}

    @utils.log_exceptions(logger=logger)
    async def _main_task(self) -> None:
        num_retries: int = 0
        max_retries = self._opts.conn_options.max_retry

        async def _reconnect() -> None:
            logger.debug(
                f"reconnecting to {self._realtime_model._provider_label}",
                extra={"max_session_duration": self._opts.max_session_duration},
            )

            events: list[RealtimeClientEvent | dict[str, Any]] = []

            # options and instructions
            events.append(self._create_session_update_event())

            # tools
            tools = self._tools.flatten()
            if tools:
                events.append(self._create_tools_update_event(tools))

            # chat context. the turn state goes first, since what it settles belongs in the
            # mirror that is replayed below
            self._reset_input_turn_state()
            chat_ctx = self.chat_ctx.copy(
                exclude_function_call=True,
                exclude_instructions=True,
                exclude_empty_message=True,
                exclude_handoff=True,
                exclude_config_update=True,
            )
            old_chat_ctx = self._remote_chat_ctx
            self._remote_chat_ctx = llm.remote_chat_context.RemoteChatContext()
            events.extend(self._create_update_chat_ctx_events(chat_ctx))

            try:
                for ev in events:
                    # certain events could already be in dict format
                    if isinstance(ev, BaseModel):
                        ev = ev.model_dump(
                            by_alias=True, exclude_unset=True, exclude_defaults=False
                        )

                    if self._opts.is_azure and self._opts.api_version:
                        _normalize_azure_client_event(ev)

                    self.emit("openai_client_event_queued", ev)
                    await ws_conn.send_str(json.dumps(ev))
            except Exception as e:
                self._remote_chat_ctx = old_chat_ctx  # restore the old chat context
                raise APIConnectionError(
                    message=(
                        f"Failed to send message to {self._realtime_model._provider_label} during session re-connection"
                    ),
                ) from e

            for fut in self._response_created_futures.values():
                if not fut.done():
                    fut.set_exception(
                        llm.RealtimeError("pending response discarded due to session reconnection")
                    )
            self._response_created_futures.clear()
            self._discarded_event_ids.clear()
            self._close_current_generation("session reconnection")

            logger.debug(f"reconnected to {self._realtime_model._provider_label}")
            self.emit("session_reconnected", llm.RealtimeSessionReconnectedEvent())

        reconnecting = False
        try:
            while not self._msg_ch.closed:
                try:
                    ws_conn = await self._create_ws_conn()
                    if reconnecting:
                        await _reconnect()
                        num_retries = 0  # reset the retry counter
                    await self._run_ws(ws_conn)

                except APIError as e:
                    if max_retries == 0 or not e.retryable:
                        self._emit_error(e, recoverable=False)
                        raise
                    elif num_retries == max_retries:
                        self._emit_error(e, recoverable=False)
                        raise APIConnectionError(
                            f"{self._realtime_model._provider_label} connection failed after {num_retries} attempts",
                        ) from e
                    else:
                        self._emit_error(e, recoverable=True)

                        retry_interval = self._opts.conn_options._interval_for_retry(num_retries)
                        logger.warning(
                            f"{self._realtime_model._provider_label} connection failed, retrying in {retry_interval}s",
                            exc_info=e,
                            extra={"attempt": num_retries, "max_retries": max_retries},
                        )
                        await asyncio.sleep(retry_interval)
                    num_retries += 1

                except Exception as e:
                    self._emit_error(e, recoverable=False)
                    raise

                reconnecting = True
        finally:
            # the session loop has exited (fatal server error, retries exhausted, or
            # close); close any in-progress generation and fail any pending
            # generate_reply futures so consumers don't hang and callers don't wait
            # out their timeout
            self._close_current_generation("session closed")
            for fut in self._response_created_futures.values():
                if not fut.done():
                    fut.set_exception(llm.RealtimeError("realtime session closed"))
            self._response_created_futures.clear()

    async def _create_ws_conn(self) -> aiohttp.ClientWebSocketResponse:
        url, headers = self._create_ws_url_and_headers()

        if lk_oai_debug:
            logger.debug(f"connecting to Realtime API: {url}")

        t0 = time.perf_counter()
        try:
            ws = await asyncio.wait_for(
                self._realtime_model._ensure_http_session().ws_connect(url=url, headers=headers),
                self._opts.conn_options.timeout,
            )
            self._report_connection_acquired(time.perf_counter() - t0)
            return ws
        except aiohttp.ClientError as e:
            raise APIConnectionError(
                f"{self._realtime_model._provider_label} client connection error"
            ) from e
        except asyncio.TimeoutError as e:
            raise APIConnectionError(
                message=f"{self._realtime_model._provider_label} connection timed out",
            ) from e

    def _create_ws_url_and_headers(self) -> tuple[str, dict[str, str]]:
        headers = {"User-Agent": "LiveKit Agents"}
        if self._opts.is_azure:
            if self._opts.entra_token:
                headers["Authorization"] = f"Bearer {self._opts.entra_token}"

            if self._opts.api_key:
                headers["api-key"] = self._opts.api_key
        else:
            headers["Authorization"] = f"Bearer {self._opts.api_key}"

        url = process_base_url(
            self._opts.base_url,
            self._opts.model,
            is_azure=self._opts.is_azure,
            api_version=self._opts.api_version,
            azure_deployment=self._opts.azure_deployment,
        )
        return url, headers

    async def _run_ws(self, ws_conn: aiohttp.ClientWebSocketResponse) -> None:
        closing = False

        @utils.log_exceptions(logger=logger)
        async def _send_task() -> None:
            nonlocal closing
            async for msg in self._msg_ch:
                try:
                    if isinstance(msg, BaseModel):
                        msg = msg.model_dump(
                            by_alias=True, exclude_unset=True, exclude_defaults=False
                        )

                    # Azure uses "text" for assistant content parts, while
                    # the new API uses "output_text" for assistant content.
                    if self._opts.is_azure and self._opts.api_version:
                        _normalize_azure_client_event(msg)

                    self.emit("openai_client_event_queued", msg)
                    await ws_conn.send_str(json.dumps(msg))

                    if lk_oai_debug and msg["type"] != "input_audio_buffer.append":
                        logger.debug(">>>", extra={"lk.pii.event": msg})
                except Exception:
                    logger.exception("failed to send event")

            closing = True
            await ws_conn.close()

        @utils.log_exceptions(logger=logger)
        async def _recv_task() -> None:
            while True:
                msg = await ws_conn.receive()
                if msg.type in (
                    aiohttp.WSMsgType.CLOSED,
                    aiohttp.WSMsgType.CLOSE,
                    aiohttp.WSMsgType.CLOSING,
                ):
                    if closing:  # closing is expected, see _send_task
                        return

                    # this will trigger a reconnection
                    raise APIConnectionError(
                        message=f"{self._realtime_model._provider_label} connection closed unexpectedly"
                    )

                if msg.type != aiohttp.WSMsgType.TEXT:
                    continue

                if self._closing:
                    # draining after aclose; the generation is already discarded
                    continue

                event = json.loads(msg.data)

                # Azure OpenAI uses old-style event names from the beta API.
                # Normalize them to the current OpenAI event names so the rest
                # of the handler code only needs to deal with one set of names.
                if self._opts.is_azure:
                    event_type = event.get("type", "")
                    normalized = _AZURE_EVENT_MAPPING.get(event_type)
                    if normalized is not None:
                        event["type"] = normalized

                # emit the raw json dictionary instead of the BaseModel because different
                # providers can have different event types that are not part of the OpenAI Realtime API  # noqa: E501
                self.emit("openai_server_event_received", event)

                try:
                    if lk_oai_debug:
                        event_copy = event.copy()
                        if event_copy["type"] == "response.output_audio.delta":
                            event_copy = {**event_copy, "delta": "..."}

                        logger.debug("<<<", extra={"lk.pii.event": event_copy})

                    if event["type"] == "input_audio_buffer.speech_started":
                        self._handle_input_audio_buffer_speech_started(
                            InputAudioBufferSpeechStartedEvent.construct(**event)
                        )
                    elif event["type"] == "input_audio_buffer.speech_stopped":
                        self._handle_input_audio_buffer_speech_stopped(
                            InputAudioBufferSpeechStoppedEvent.construct(**event)
                        )
                    elif event["type"] == "response.created":
                        self._handle_response_created(ResponseCreatedEvent.construct(**event))
                    elif event["type"] == "response.output_item.added":
                        self._handle_response_output_item_added(
                            ResponseOutputItemAddedEvent.construct(**event)
                        )
                    elif event["type"] == "response.content_part.added":
                        self._handle_response_content_part_added(
                            ResponseContentPartAddedEvent.construct(**event)
                        )
                    elif event["type"] == "conversation.item.added":
                        self._handle_conversion_item_added(ConversationItemAdded.construct(**event))
                    elif event["type"] == "conversation.item.deleted":
                        self._handle_conversion_item_deleted(
                            ConversationItemDeletedEvent.construct(**event)
                        )
                    elif event["type"] == "conversation.item.input_audio_transcription.delta":
                        self._handle_conversion_item_input_audio_transcription_delta(
                            ConversationItemInputAudioTranscriptionDeltaEvent.construct(**event)
                        )
                    elif event["type"] == "conversation.item.input_audio_transcription.completed":
                        self._handle_conversion_item_input_audio_transcription_completed(
                            ConversationItemInputAudioTranscriptionCompletedEvent.construct(**event)
                        )
                    elif event["type"] == "conversation.item.input_audio_transcription.failed":
                        self._handle_conversion_item_input_audio_transcription_failed(
                            ConversationItemInputAudioTranscriptionFailedEvent.construct(**event)
                        )
                    elif event["type"] == "response.output_text.delta":
                        self._handle_response_text_delta(ResponseTextDeltaEvent.construct(**event))
                    elif event["type"] == "response.output_text.done":
                        self._handle_response_text_done(ResponseTextDoneEvent.construct(**event))
                    elif event["type"] == "response.output_audio_transcript.delta":
                        self._handle_response_audio_transcript_delta(event)
                    elif event["type"] == "response.output_audio.delta":
                        self._handle_response_audio_delta(
                            ResponseAudioDeltaEvent.construct(**event)
                        )
                    elif event["type"] == "response.output_audio.done":
                        self._handle_response_audio_done(ResponseAudioDoneEvent.construct(**event))
                    elif event["type"] == "response.output_item.done":
                        self._handle_response_output_item_done(
                            ResponseOutputItemDoneEvent.construct(**event)
                        )
                    elif event["type"] == "response.done":
                        self._handle_response_done(ResponseDoneEvent.construct(**event))
                    elif event["type"] == "error":
                        self._handle_error(RealtimeErrorEvent.construct(**event))
                    elif lk_oai_debug:
                        logger.debug(
                            f"unhandled event: {event['type']}", extra={"lk.pii.event": event}
                        )
                except Exception as e:
                    # terminal server errors (e.g. insufficient_quota) must break the recv
                    # loop so _main_task stops reconnecting; every other handler failure is
                    # logged and skipped
                    if isinstance(e, APIError) and not e.retryable:
                        raise
                    if event["type"] == "response.output_audio.delta":
                        event["delta"] = event["delta"][:10] + "..."
                    logger.exception("failed to handle event", extra={"lk.pii.event": event})

        tasks = [
            asyncio.create_task(_recv_task(), name="_recv_task"),
            asyncio.create_task(_send_task(), name="_send_task"),
        ]
        wait_reconnect_task: asyncio.Task | None = None
        if self._opts.max_session_duration is not None:
            wait_reconnect_task = asyncio.create_task(
                asyncio.sleep(self._opts.max_session_duration),
                name="_timeout_task",
            )
            tasks.append(wait_reconnect_task)
        try:
            done, _ = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)

            # propagate exceptions from completed tasks
            for task in done:
                if task != wait_reconnect_task:
                    task.result()

            if (
                wait_reconnect_task
                and wait_reconnect_task in done
                and isinstance(self._current_generation, _ResponseGeneration)
            ):
                # wait for the current generation to complete before reconnecting
                await self._current_generation._done_fut
                closing = True

        finally:
            await utils.aio.cancel_and_wait(*tasks)
            await ws_conn.close()

    def _wrap_session_update(
        self, event_id: str, session: RealtimeSessionCreateRequest
    ) -> SessionUpdateEvent | dict[str, Any]:
        """Wrap a session object in the appropriate event type.

        For Azure, converts the new-style session to the old flat format
        and returns a dict (since AzureSessionUpdateEvent is not part of
        the RealtimeClientEvent union).
        """
        if self._opts.is_azure and self._opts.api_version:
            # legacy Azure API: convert to old flat format
            return AzureSessionUpdateEvent(
                type="session.update",
                session=_oai_session_to_azure(session),
                event_id=event_id,
            ).model_dump(by_alias=True, exclude_unset=True, exclude_defaults=False)

        return SessionUpdateEvent(
            type="session.update",
            session=session,
            event_id=event_id,
        )

    def _create_session_update_event(self) -> SessionUpdateEvent | dict[str, Any]:
        audio_format = realtime.realtime_audio_formats.AudioPCM(rate=SAMPLE_RATE, type="audio/pcm")
        # they do not support both text and audio modalities, it'll respond in audio + transcript
        modality = "audio" if "audio" in self._opts.modalities else "text"
        opts = self._opts

        session = RealtimeSessionCreateRequest(
            type="realtime",
            model=opts.model,
            output_modalities=[modality],
            audio=RealtimeAudioConfig(
                input=RealtimeAudioConfigInput(
                    format=audio_format,
                    noise_reduction=opts.input_audio_noise_reduction,
                    transcription=opts.input_audio_transcription,
                    turn_detection=opts.turn_detection,
                ),
                output=RealtimeAudioConfigOutput(
                    format=audio_format,
                    speed=opts.speed,
                    voice=opts.voice,
                ),
            ),
            max_output_tokens=opts.max_response_output_tokens,
            tool_choice=to_oai_tool_choice(opts.tool_choice),
            tracing=opts.tracing,
        )
        if self._instructions is not None:
            session.instructions = self._instructions
        if opts.truncation is not None:
            session.truncation = opts.truncation
        if opts.reasoning is not None:
            session.reasoning = opts.reasoning

        return self._wrap_session_update(
            event_id=utils.shortuuid("session_update_"), session=session
        )

    @property
    def capabilities(self) -> llm.RealtimeCapabilities:
        return self._capabilities

    @property
    def chat_ctx(self) -> llm.ChatContext:
        return self._remote_chat_ctx.to_chat_ctx()

    @property
    def tools(self) -> llm.ToolContext:
        return self._tools.copy()

    def update_options(
        self,
        *,
        tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
        voice: NotGivenOr[str] = NOT_GIVEN,
        turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | None] = NOT_GIVEN,
        max_response_output_tokens: NotGivenOr[int | Literal["inf"] | None] = NOT_GIVEN,
        input_audio_transcription: NotGivenOr[AudioTranscription | None] = NOT_GIVEN,
        input_audio_noise_reduction: NotGivenOr[
            NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None
        ] = NOT_GIVEN,
        speed: NotGivenOr[float] = NOT_GIVEN,
        tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
        truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
        reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
    ) -> None:
        session = RealtimeSessionCreateRequest(type="realtime")
        has_changes = False

        if is_given(tool_choice):
            current_oai = to_oai_tool_choice(self._opts.tool_choice)
            next_oai = to_oai_tool_choice(tool_choice)
            self._opts.tool_choice = tool_choice
            if current_oai != next_oai:
                session.tool_choice = next_oai
                has_changes = True

        if is_given(max_response_output_tokens):
            if self._opts.max_response_output_tokens != max_response_output_tokens:
                session.max_output_tokens = max_response_output_tokens
                has_changes = True
            self._opts.max_response_output_tokens = max_response_output_tokens

        if is_given(tracing):
            if self._opts.tracing != tracing:
                session.tracing = tracing  # type: ignore[assignment]
                has_changes = True
            self._opts.tracing = tracing

        if is_given(truncation):
            if self._opts.truncation != truncation:
                session.truncation = truncation
                has_changes = True
            self._opts.truncation = truncation

        if is_given(reasoning):
            if self._opts.reasoning != reasoning:
                # setting reasoning to None clears it server-side
                session.reasoning = reasoning
                has_changes = True
            self._opts.reasoning = reasoning

        has_audio_config = False
        audio_output = RealtimeAudioConfigOutput()
        audio_input = RealtimeAudioConfigInput()
        audio_config = RealtimeAudioConfig(output=audio_output, input=audio_input)

        if is_given(voice):
            if self._opts.voice != voice:
                audio_output.voice = voice
                has_audio_config = True
            self._opts.voice = voice

        if is_given(turn_detection):
            if self._opts.turn_detection != turn_detection:
                audio_input.turn_detection = turn_detection
                has_audio_config = True
            self._opts.turn_detection = turn_detection
            self._capabilities.turn_detection = _server_turn_taking_enabled(turn_detection)

        if is_given(input_audio_transcription):
            if self._opts.input_audio_transcription != input_audio_transcription:
                audio_input.transcription = input_audio_transcription
                has_audio_config = True
            self._opts.input_audio_transcription = input_audio_transcription
            self._capabilities.user_transcription = input_audio_transcription is not None

        if is_given(input_audio_noise_reduction):
            input_audio_noise_reduction = to_noise_reduction(input_audio_noise_reduction)
            if self._opts.input_audio_noise_reduction != input_audio_noise_reduction:
                audio_input.noise_reduction = input_audio_noise_reduction
                has_audio_config = True
            self._opts.input_audio_noise_reduction = input_audio_noise_reduction

        if is_given(speed):
            if self._opts.speed != speed:
                audio_output.speed = speed
                has_audio_config = True
            self._opts.speed = speed

        if has_audio_config:
            session.audio = audio_config
            has_changes = True

        if has_changes:
            self.send_event(
                self._wrap_session_update(
                    event_id=utils.shortuuid("options_update_"), session=session
                )
            )

    async def update_chat_ctx(self, chat_ctx: llm.ChatContext) -> None:
        async with self._update_chat_ctx_lock:
            chat_ctx = chat_ctx.copy(
                exclude_handoff=True,
                exclude_config_update=True,
            )
            # only remove the instructions but keep other system messages
            remove_instructions(chat_ctx)

            events = self._create_update_chat_ctx_events(chat_ctx)
            futs: list[asyncio.Future[None]] = []
            self._chat_ctx_event_futures = {}

            for ev in events:
                futs.append(f := asyncio.Future[None]())
                if isinstance(ev, ConversationItemDeleteEvent):
                    self._item_delete_future[ev.item_id] = f
                else:
                    assert ev.item.id is not None
                    self._item_create_future[ev.item.id] = f

                # an updated item sends a delete and a create under the same id, so only the
                # event id tells a rejection which of the two it answers
                if ev.event_id:
                    self._chat_ctx_event_futures[ev.event_id] = f
                self.send_event(ev)

            if not futs:
                return
            try:
                results = await asyncio.wait_for(
                    asyncio.gather(*futs, return_exceptions=True), timeout=5.0
                )
            except asyncio.TimeoutError:
                raise llm.RealtimeError("update_chat_ctx timed out.") from None
            finally:
                self._chat_ctx_event_futures = {}
                for ev in events:
                    if isinstance(ev, ConversationItemDeleteEvent):
                        self._item_delete_future.pop(ev.item_id, None)
                    else:
                        assert ev.item.id is not None
                        self._item_create_future.pop(ev.item.id, None)

            # a rejected item is not worth failing the turn over
            if rejected := [str(r) for r in results if isinstance(r, BaseException)]:
                logger.warning(
                    f"{self._realtime_model._provider_label} rejected part of a chat context update",  # noqa: E501
                    extra={"errors": rejected},
                )

    def _create_update_chat_ctx_events(
        self, chat_ctx: llm.ChatContext
    ) -> list[ConversationItemCreateEvent | ConversationItemDeleteEvent]:
        events: list[ConversationItemCreateEvent | ConversationItemDeleteEvent] = []
        remote_ctx = self._remote_chat_ctx.to_chat_ctx()

        # Empty message content can mean either:
        # - a local placeholder that should not be created remotely, or
        # - an existing remote item with non-text content (audio/images) that is not
        #   synced into the agent-side ChatContext.
        # Keep empty messages that already exist remotely so we do not delete them.
        remote_ids = {item.id for item in remote_ctx.items}
        chat_ctx = llm.ChatContext(
            [
                item
                for item in chat_ctx.items
                if item.type != "message" or item.content or item.id in remote_ids
            ]
        )
        diff_ops = llm.utils.compute_chat_ctx_diff(remote_ctx, chat_ctx)

        def _delete_item(msg_id: str) -> None:
            events.append(
                ConversationItemDeleteEvent(
                    type="conversation.item.delete",
                    item_id=msg_id,
                    event_id=utils.shortuuid("chat_ctx_delete_"),
                )
            )

        def _create_item(previous_msg_id: str | None, msg_id: str) -> None:
            chat_item = chat_ctx.get_by_id(msg_id)
            assert chat_item is not None
            events.append(
                ConversationItemCreateEvent(
                    type="conversation.item.create",
                    item=livekit_item_to_openai_item(chat_item),
                    previous_item_id=("root" if previous_msg_id is None else previous_msg_id),
                    event_id=utils.shortuuid("chat_ctx_create_"),
                )
            )

        def _is_content_empty(msg_id: str) -> bool:
            remote_item = remote_ctx.get_by_id(msg_id)
            if remote_item and remote_item.type == "message" and not remote_item.content:
                return True
            return False

        for msg_id in diff_ops.to_remove:
            _delete_item(msg_id)

        for previous_msg_id, msg_id in diff_ops.to_create:
            _create_item(previous_msg_id, msg_id)

        # update the items with the same id but different content
        for previous_msg_id, msg_id in diff_ops.to_update:
            # empty content almost always means the content is not synced down
            # we don't want to recreate these items there
            if _is_content_empty(msg_id):
                continue
            _delete_item(msg_id)
            _create_item(previous_msg_id, msg_id)

        return events

    async def update_tools(self, tools: list[llm.Tool]) -> None:
        async with self._update_fnc_ctx_lock:
            ev = self._create_tools_update_event(tools)
            self.send_event(ev)

            retained_tool_names: set[str] = set()
            for t in ev["session"]["tools"]:
                if name := t.get("name"):
                    retained_tool_names.add(name)
                # TODO(dz): handle MCP tools
            retained_tools = [
                tool
                for tool in tools
                if (
                    isinstance(tool, (llm.FunctionTool, llm.RawFunctionTool))
                    and tool.info.name in retained_tool_names
                )
                or isinstance(tool, llm.ProviderTool)
            ]
            self._tools = llm.ToolContext(retained_tools)

    # this function can be overrided
    def _convert_tools_to_oai(self, tools: list[llm.Tool]) -> list[RealtimeFunctionTool]:
        oai_tools: list[RealtimeFunctionTool] = []

        for tool in tools:
            if isinstance(tool, llm.FunctionTool):
                tool_desc = llm.utils.build_legacy_openai_schema(tool, internally_tagged=True)
            elif isinstance(tool, llm.RawFunctionTool):
                # copy to avoid modifying original
                tool_desc = dict(tool.info.raw_schema)
                tool_desc.pop("meta", None)  # meta is not supported by OpenAI Realtime API
                tool_desc["type"] = "function"  # internally tagged
            elif isinstance(tool, llm.ProviderTool):
                continue  # currently only xAI supports ProviderTools
            else:
                logger.error(
                    f"{self._realtime_model._provider_label} doesn't support this tool type",
                    extra={"tool": tool},
                )
                continue

            try:
                session_tool = RealtimeFunctionTool.model_validate(tool_desc)
                oai_tools.append(session_tool)
            except ValidationError:
                logger.error(
                    f"{self._realtime_model._provider_label} doesn't support this tool",
                    extra={"tool": tool_desc},
                )
                continue

        return oai_tools

    def _create_tools_update_event(self, tools: list[llm.Tool]) -> dict[str, Any]:
        oai_tools = self._convert_tools_to_oai(tools)

        event = self._wrap_session_update(
            event_id=utils.shortuuid("tools_update_"),
            session=RealtimeSessionCreateRequest.model_construct(
                type="realtime",
                model=self._opts.model,
                tools=oai_tools,  # type: ignore
            ),
        )
        if isinstance(event, dict):
            return event
        return event.model_dump(by_alias=True, exclude_unset=True, exclude_defaults=False)

    async def update_instructions(self, instructions: str) -> None:
        self.send_event(
            self._wrap_session_update(
                event_id=utils.shortuuid("instructions_update_"),
                session=RealtimeSessionCreateRequest.model_construct(
                    type="realtime",
                    instructions=instructions,
                ),
            )
        )
        self._instructions = instructions

    def push_audio(self, frame: rtc.AudioFrame) -> None:
        for f in self._resample_audio(frame):
            data = f.data.tobytes()
            for nf in self._bstream.write(data):
                self.send_event(
                    InputAudioBufferAppendEvent(
                        type="input_audio_buffer.append",
                        audio=base64.b64encode(nf.data).decode("utf-8"),
                    )
                )
                self._pushed_duration_s += nf.duration

    def push_video(self, frame: rtc.VideoFrame) -> None:
        message = llm.ChatMessage(
            role="user",
            content=[llm.ImageContent(image=frame)],
        )
        oai_item = livekit_item_to_openai_item(message)
        self.send_event(
            ConversationItemCreateEvent(
                type="conversation.item.create",
                item=oai_item,
                event_id=utils.shortuuid("video_"),
            )
        )

    def commit_audio(self) -> None:
        if self._pushed_duration_s > 0.1:  # OpenAI requires at least 100ms of audio
            self.send_event(InputAudioBufferCommitEvent(type="input_audio_buffer.commit"))
            self._pushed_duration_s = 0

    def clear_audio(self) -> None:
        self.send_event(InputAudioBufferClearEvent(type="input_audio_buffer.clear"))
        self._pushed_duration_s = 0

    def generate_reply(
        self,
        *,
        instructions: NotGivenOr[str] = NOT_GIVEN,
        tool_choice: NotGivenOr[llm.ToolChoice] = NOT_GIVEN,
        tools: NotGivenOr[list[llm.Tool]] = NOT_GIVEN,
    ) -> asyncio.Future[llm.GenerationCreatedEvent]:
        event_id = utils.shortuuid("response_create_")
        fut = asyncio.Future[llm.GenerationCreatedEvent]()
        self._response_created_futures[event_id] = fut

        if is_given(instructions) and self._instructions:
            # in OpenAI realtime, the session-level instructions are completely replaced
            # by the new instructions for this response
            instructions = f"{self._instructions}\n{instructions}"

        params = RealtimeResponseCreateParams(
            instructions=instructions or None,
            metadata={"client_event_id": event_id},
        )
        if is_given(tool_choice):
            params.tool_choice = to_oai_tool_choice(tool_choice)
        if is_given(tools):
            params.tools = self._convert_tools_to_oai(tools)  # type: ignore

        self.send_event(
            ResponseCreateEvent(type="response.create", event_id=event_id, response=params)
        )

        def _on_timeout() -> None:
            self._response_created_futures.pop(event_id, None)
            if fut and not fut.done():
                # discard the response if the server still creates it after the timeout
                self._discarded_event_ids.add(event_id)
                fut.set_exception(llm.RealtimeError("generate_reply timed out."))

        handle = asyncio.get_event_loop().call_later(10.0, _on_timeout)

        def _on_fut_done(f: asyncio.Future[llm.GenerationCreatedEvent]) -> None:
            handle.cancel()
            self._response_created_futures.pop(event_id, None)
            if f.cancelled():
                # response.create was already sent; cancel the response server-side
                self.send_event(ResponseCancelEvent(type="response.cancel"))
                # the cancel above is a no-op if the response isn't created yet; discard it by id
                # when it arrives
                self._discarded_event_ids.add(event_id)

        fut.add_done_callback(_on_fut_done)
        return fut

    @property
    def has_active_generation(self) -> bool:
        return self._current_generation is not None or len(self._response_created_futures) > 0

    def interrupt(self) -> None:
        if not self.has_active_generation:
            return
        self.send_event(ResponseCancelEvent(type="response.cancel"))

    def truncate(
        self,
        *,
        message_id: str,
        modalities: list[Literal["text", "audio"]],
        audio_end_ms: int,
        audio_transcript: NotGivenOr[str] = NOT_GIVEN,
    ) -> None:
        if "audio" in modalities:
            if audio_end_ms > 0:
                self.send_event(
                    ConversationItemTruncateEvent(
                        type="conversation.item.truncate",
                        content_index=0,
                        item_id=message_id,
                        audio_end_ms=audio_end_ms,
                    )
                )
            else:
                self.send_event(
                    ConversationItemDeleteEvent(
                        type="conversation.item.delete",
                        item_id=message_id,
                        event_id=utils.shortuuid("chat_ctx_delete_"),
                    )
                )
        elif utils.is_given(audio_transcript):
            # sync the forwarded text to the remote chat ctx
            chat_ctx = self.chat_ctx.copy(
                exclude_handoff=True,
                exclude_config_update=True,
            )
            if (idx := chat_ctx.index_by_id(message_id)) is not None:
                new_item = copy.copy(chat_ctx.items[idx])
                assert new_item.type == "message"

                new_item.content = [audio_transcript]
                chat_ctx.items[idx] = new_item
                events = self._create_update_chat_ctx_events(chat_ctx)
                for ev in events:
                    self.send_event(ev)

    async def aclose(self) -> None:
        self._closing = True
        self._close_current_generation("session closed")
        self._msg_ch.close()
        await self._main_atask

    def _close_current_generation(self, reason: str | None = None) -> None:
        """Close all channels and resolve _done_fut for the current generation.

        This prevents consumers from hanging indefinitely when a generation is
        interrupted by a reconnection or session close.
        """
        if isinstance(self._current_generation, _DiscardedGeneration):
            self._current_generation = None
            return

        if self._current_generation is None or self._current_generation._done_fut.done():
            return

        for generation in self._current_generation.messages.values():
            generation.text_ch.close()
            generation.audio_ch.close()
            if not generation.modalities.done():
                generation.modalities.set_result(self._opts.modalities)

        self._current_generation.function_ch.close()
        self._current_generation.message_ch.close()

        with contextlib.suppress(asyncio.InvalidStateError):
            self._current_generation._done_fut.set_result(None)
        self._current_generation = None

        if reason:
            logger.warning(f"in-progress generation discarded due to {reason}")

    def _resample_audio(self, frame: rtc.AudioFrame) -> Iterator[rtc.AudioFrame]:
        if self._input_resampler:
            if frame.sample_rate != self._input_resampler._input_rate:
                # input audio changed to a different sample rate
                self._input_resampler = None

        if self._input_resampler is None and (
            frame.sample_rate != SAMPLE_RATE or frame.num_channels != NUM_CHANNELS
        ):
            self._input_resampler = rtc.AudioResampler(
                input_rate=frame.sample_rate,
                output_rate=SAMPLE_RATE,
                num_channels=NUM_CHANNELS,
            )

        if self._input_resampler:
            # TODO(long): flush the resampler when the input source is changed
            yield from self._input_resampler.push(frame)
        else:
            yield frame

    def _handle_input_audio_buffer_speech_started(
        self, event: InputAudioBufferSpeechStartedEvent
    ) -> None:
        if event.item_id:
            self._input_speech_started_at[event.item_id] = time.time()
        self.emit("input_speech_started", llm.InputSpeechStartedEvent())

    def _handle_input_audio_buffer_speech_stopped(
        self, _: InputAudioBufferSpeechStoppedEvent
    ) -> None:
        user_transcription_enabled = self._opts.input_audio_transcription is not None
        self.emit(
            "input_speech_stopped",
            llm.InputSpeechStoppedEvent(user_transcription_enabled=user_transcription_enabled),
        )

    def _handle_response_created(self, event: ResponseCreatedEvent) -> None:
        assert event.response.id is not None, "response.id is None"

        client_event_id: str | None = None
        if isinstance(event.response.metadata, dict):
            client_event_id = event.response.metadata.get("client_event_id")

        if client_event_id and client_event_id in self._discarded_event_ids:
            # interrupted or timed out before the server created it: cancel by id and mark it
            # discarded so its trailing events are skipped, instead of surfacing it
            self._discarded_event_ids.discard(client_event_id)
            self.send_event(
                ResponseCancelEvent(type="response.cancel", response_id=event.response.id)
            )
            self._current_generation = _DiscardedGeneration()
            logger.warning("discarding response that arrived after it was timed out or interrupted")
            return

        self._current_generation = _ResponseGeneration(
            message_ch=utils.aio.Chan(),
            function_ch=utils.aio.Chan(),
            messages={},
            _created_timestamp=time.time(),
            _done_fut=asyncio.Future(),
        )

        generation_ev = llm.GenerationCreatedEvent(
            message_stream=self._current_generation.message_ch,
            function_stream=self._current_generation.function_ch,
            user_initiated=False,
            response_id=event.response.id,
        )

        if client_event_id and (fut := self._response_created_futures.pop(client_event_id, None)):
            if not fut.done():
                generation_ev.user_initiated = True
                fut.set_result(generation_ev)
            else:
                logger.warning("response of generate_reply received after it's timed out.")

        self.emit("generation_created", generation_ev)

    def _handle_response_output_item_added(self, event: ResponseOutputItemAddedEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"
        assert (item_id := event.item.id) is not None, "item.id is None"
        assert (item_type := event.item.type) is not None, "item.type is None"

        if item_type == "message":
            item_generation = _MessageGeneration(
                message_id=item_id,
                text_ch=utils.aio.Chan(),
                audio_ch=utils.aio.Chan(),
                modalities=asyncio.Future(),
            )
            if not self._realtime_model.capabilities.audio_output:
                item_generation.audio_ch.close()
                item_generation.modalities.set_result(["text"])

            self._current_generation.message_ch.send_nowait(
                llm.MessageGeneration(
                    message_id=item_id,
                    text_stream=item_generation.text_ch,
                    audio_stream=item_generation.audio_ch,
                    modalities=item_generation.modalities,
                )
            )
            self._current_generation.messages[item_id] = item_generation

    def _handle_response_content_part_added(self, event: ResponseContentPartAddedEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"
        assert (item_id := event.item_id) is not None, "item_id is None"
        assert (item_type := event.part.type) is not None, "part.type is None"

        if item_type == "text" and self._realtime_model.capabilities.audio_output:
            logger.warning(
                f"Text response received from {self._realtime_model._provider_label} in audio modality."
            )

        with contextlib.suppress(asyncio.InvalidStateError):
            self._current_generation.messages[item_id].modalities.set_result(
                ["text"] if item_type == "text" else ["audio", "text"]
            )

    def _handle_conversion_item_added(self, event: ConversationItemAdded) -> None:
        assert event.item.id is not None, "item.id is None"

        if event.previous_item_id and not self._remote_chat_ctx.get(event.previous_item_id):
            # the server can anchor to an item it just deleted; the item belongs at the tail
            logger.warning(
                f"{self._realtime_model._provider_label} anchored an item to one it is no longer "
                "tracking, appending it instead",
                extra={"item_id": event.item.id, "previous_item_id": event.previous_item_id},
            )
            event.previous_item_id = self._remote_chat_ctx.tail_id

        try:
            lk_item = openai_item_to_livekit_item(event.item)
            self._remote_chat_ctx.insert(event.previous_item_id, lk_item)
            self.emit(
                "remote_item_added",
                llm.RemoteItemAddedEvent(previous_item_id=event.previous_item_id, item=lk_item),
            )
        except ValueError as e:
            logger.warning(
                f"failed to insert item `{event.item.id}`: {str(e)}",
            )

        if fut := self._item_create_future.pop(event.item.id, None):
            if fut.done():
                logger.error(f"item create future for `{event.item.id}` was already settled")
            else:
                fut.set_result(None)

    def _handle_conversion_item_deleted(self, event: ConversationItemDeletedEvent) -> None:
        assert event.item_id is not None, "item_id is None"

        self._input_transcript_accumulators.pop(event.item_id, None)
        self._input_speech_started_at.pop(event.item_id, None)

        try:
            self._remote_chat_ctx.delete(event.item_id)
        except ValueError as e:
            logger.warning(
                f"failed to delete item `{event.item_id}`: {str(e)}",
            )

        if fut := self._item_delete_future.pop(event.item_id, None):
            if fut.done():
                logger.error(f"item delete future for `{event.item_id}` was already settled")
            else:
                fut.set_result(None)

    def _handle_conversion_item_input_audio_transcription_delta(
        self, event: ConversationItemInputAudioTranscriptionDeltaEvent
    ) -> None:
        if not event.delta:
            return

        content_index = event.content_index or 0
        by_index = self._input_transcript_accumulators.setdefault(event.item_id, {})
        accumulated = by_index.get(content_index, "") + event.delta
        by_index[content_index] = accumulated

        self.emit(
            "input_audio_transcription_completed",
            llm.InputTranscriptionCompleted(
                item_id=event.item_id, transcript=accumulated, is_final=False
            ),
        )

    def _clear_transcript_accumulator(self, item_id: str, content_index: int) -> str | None:
        by_index = self._input_transcript_accumulators.get(item_id)
        if by_index is None:
            return None
        partial = by_index.pop(content_index, None)
        if not by_index:
            self._input_transcript_accumulators.pop(item_id, None)
        return partial

    def _handle_conversion_item_input_audio_transcription_completed(
        self, event: ConversationItemInputAudioTranscriptionCompletedEvent
    ) -> None:
        self._clear_transcript_accumulator(event.item_id, event.content_index or 0)

        confidence = calculate_confidence_from_logprobs(event.logprobs)

        if remote_item := self._remote_chat_ctx.get(event.item_id):
            assert isinstance(remote_item.item, llm.ChatMessage)
            remote_item.item.content.append(event.transcript)
            remote_item.item.transcript_confidence = confidence

        self.emit(
            "input_audio_transcription_completed",
            llm.InputTranscriptionCompleted(
                item_id=event.item_id,
                transcript=event.transcript,
                is_final=True,
                confidence=confidence,
                turn_started_at=self._input_speech_started_at.pop(event.item_id, None),
            ),
        )

        self._emit_transcription_metrics(event)

    def _handle_conversion_item_input_audio_transcription_failed(
        self, event: ConversationItemInputAudioTranscriptionFailedEvent
    ) -> None:
        logger.error(
            f"{self._realtime_model._provider_label} failed to transcribe input audio",
            extra={"error": event.error},
        )

        # close any open partial stream so consumers waiting for is_final don't hang
        partial = self._clear_transcript_accumulator(event.item_id, event.content_index or 0)
        turn_started_at = self._input_speech_started_at.pop(event.item_id, None)
        if partial is None:
            return
        self.emit(
            "input_audio_transcription_completed",
            llm.InputTranscriptionCompleted(
                item_id=event.item_id,
                transcript=partial,
                is_final=True,
                turn_started_at=turn_started_at,
            ),
        )

    def _emit_transcription_metrics(
        self, event: ConversationItemInputAudioTranscriptionCompletedEvent
    ) -> None:
        usage = _coerce_transcription_usage(event.usage)
        if usage is None:
            return

        transcription_opts = self._opts.input_audio_transcription
        transcription_model = transcription_opts.model if transcription_opts else None
        metadata = Metadata(
            model_name=transcription_model,
            model_provider=self._realtime_model.provider,
        )

        if isinstance(usage, UsageTranscriptTextUsageTokens):
            details = usage.input_token_details
            input_audio_tokens = (
                details.audio_tokens if details and details.audio_tokens is not None else 0
            )
            stt_metrics = STTMetrics(
                request_id=event.event_id,
                timestamp=time.time(),
                duration=0.0,
                label=self._realtime_model.label,
                audio_duration=0.0,
                streamed=True,
                input_tokens=usage.input_tokens,
                output_tokens=usage.output_tokens,
                total_tokens=usage.total_tokens,
                input_audio_tokens=input_audio_tokens,
                metadata=metadata,
            )
            self.emit("metrics_collected", stt_metrics)
        elif isinstance(usage, UsageTranscriptTextUsageDuration):
            stt_metrics = STTMetrics(
                request_id=event.event_id,
                timestamp=time.time(),
                duration=0.0,
                label=self._realtime_model.label,
                audio_duration=usage.seconds,
                streamed=True,
                metadata=metadata,
            )
            self.emit("metrics_collected", stt_metrics)

    def _handle_response_text_delta(self, event: ResponseTextDeltaEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"
        item_generation = self._current_generation.messages[event.item_id]
        if (
            item_generation.audio_ch.closed
            and self._current_generation._first_token_timestamp is None
        ):
            # only if audio is not available
            self._current_generation._first_token_timestamp = time.time()

        item_generation.text_ch.send_nowait(event.delta)
        item_generation.audio_transcript += event.delta

    def _handle_response_text_done(self, event: ResponseTextDoneEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"

    def _handle_response_audio_transcript_delta(self, event: dict[str, Any]) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"

        item_id = event["item_id"]
        delta = event["delta"]

        if (start_time := event.get("start_time")) is not None:
            delta = io.TimedString(delta, start_time=start_time)

        item_generation = self._current_generation.messages[item_id]
        item_generation.text_ch.send_nowait(delta)
        item_generation.audio_transcript += delta

    def _handle_response_audio_delta(self, event: ResponseAudioDeltaEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"
        item_generation = self._current_generation.messages[event.item_id]
        if self._current_generation._first_token_timestamp is None:
            self._current_generation._first_token_timestamp = time.time()

        if not item_generation.modalities.done():
            item_generation.modalities.set_result(["audio", "text"])

        data = base64.b64decode(event.delta)
        item_generation.audio_ch.send_nowait(
            rtc.AudioFrame(
                data=data,
                sample_rate=SAMPLE_RATE,
                num_channels=NUM_CHANNELS,
                samples_per_channel=len(data) // 2,
            )
        )

    def _handle_response_audio_done(self, _: ResponseAudioDoneEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"

    def _handle_response_output_item_done(self, event: ResponseOutputItemDoneEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            return
        assert self._current_generation is not None, "current_generation is None"
        assert (item_id := event.item.id) is not None, "item.id is None"
        assert (item_type := event.item.type) is not None, "item.type is None"

        if item_type == "function_call" and isinstance(
            event.item, RealtimeConversationItemFunctionCall
        ):
            self._handle_function_call(event.item)

        elif item_type == "message":
            item_generation = self._current_generation.messages[item_id]
            item_generation.text_ch.close()
            item_generation.audio_ch.close()
            if not item_generation.modalities.done():
                # in case message modalities is not set, this shouldn't happen
                item_generation.modalities.set_result(self._opts.modalities)

    def _handle_function_call(self, item: RealtimeConversationItemFunctionCall) -> None:
        assert isinstance(self._current_generation, _ResponseGeneration), (
            "current_generation is None"
        )

        assert item.id is not None, "item.id is None"
        assert item.call_id is not None, "call_id is None"
        assert item.name is not None, "name is None"
        assert item.arguments is not None, "arguments is None"

        self._current_generation.function_ch.send_nowait(
            llm.FunctionCall(
                id=item.id,
                call_id=item.call_id,
                name=item.name,
                arguments=item.arguments,
            )
        )

    def _handle_response_done(self, event: ResponseDoneEvent) -> None:
        if isinstance(self._current_generation, _DiscardedGeneration):
            self._current_generation = None
            return

        if self._current_generation is None:
            return  # OpenAI has a race condition where we could receive response.done without any previous response.created (This happens generally during interruption)  # noqa: E501

        assert self._current_generation is not None, "current_generation is None"

        created_timestamp = self._current_generation._created_timestamp
        first_token_timestamp = self._current_generation._first_token_timestamp

        for generation in self._current_generation.messages.values():
            if not generation.modalities.done():
                generation.modalities.set_result(self._opts.modalities)

        for item_id, item_generation in self._current_generation.messages.items():
            if (remote_item := self._remote_chat_ctx.get(item_id)) and isinstance(
                remote_item.item, llm.ChatMessage
            ):
                remote_item.item.content.append(item_generation.audio_transcript)

        self._current_generation._close()

        with contextlib.suppress(asyncio.InvalidStateError):
            if event.response.status in ("failed", "incomplete"):
                details = event.response.status_details
                msg = f"response {event.response.status}"
                if details and details.error:
                    msg = f"{msg}: [{details.error.type}] {details.error.code}"
                elif details and details.reason:
                    msg = f"{msg}: {details.reason}"
                self._current_generation._done_fut.set_exception(llm.RealtimeError(msg))
            else:
                self._current_generation._done_fut.set_result(None)

        self._current_generation = None

        # calculate metrics
        usage = (
            event.response.usage.model_dump(exclude_defaults=True) if event.response.usage else {}
        )
        ttft = first_token_timestamp - created_timestamp if first_token_timestamp else -1
        duration = time.time() - created_timestamp
        metrics = RealtimeModelMetrics(
            timestamp=created_timestamp,
            request_id=event.response.id or "",
            ttft=ttft,
            duration=duration,
            cancelled=event.response.status == "cancelled",
            label=self._realtime_model.label,
            input_tokens=usage.get("input_tokens", 0),
            output_tokens=usage.get("output_tokens", 0),
            total_tokens=usage.get("total_tokens", 0),
            tokens_per_second=usage.get("output_tokens", 0) / duration if duration > 0 else 0,
            input_token_details=RealtimeModelMetrics.InputTokenDetails(
                audio_tokens=usage.get("input_token_details", {}).get("audio_tokens", 0),
                cached_tokens=usage.get("input_token_details", {}).get("cached_tokens", 0),
                text_tokens=usage.get("input_token_details", {}).get("text_tokens", 0),
                cached_tokens_details=RealtimeModelMetrics.CachedTokenDetails(
                    text_tokens=usage.get("input_token_details", {})
                    .get("cached_tokens_details", {})
                    .get("text_tokens", 0),
                    audio_tokens=usage.get("input_token_details", {})
                    .get("cached_tokens_details", {})
                    .get("audio_tokens", 0),
                    image_tokens=usage.get("input_token_details", {})
                    .get("cached_tokens_details", {})
                    .get("image_tokens", 0),
                ),
                image_tokens=usage.get("input_token_details", {}).get("image_tokens", 0),
            ),
            output_token_details=RealtimeModelMetrics.OutputTokenDetails(
                text_tokens=usage.get("output_token_details", {}).get("text_tokens", 0),
                audio_tokens=usage.get("output_token_details", {}).get("audio_tokens", 0),
                image_tokens=usage.get("output_token_details", {}).get("image_tokens", 0),
            ),
            metadata=Metadata(
                model_name=self._realtime_model.model, model_provider=self._realtime_model.provider
            ),
        )
        self.emit("metrics_collected", metrics)
        self._handle_response_done_but_not_complete(event)

    def _handle_response_done_but_not_complete(self, event: ResponseDoneEvent) -> None:
        """Handle response done but not complete, i.e. cancelled, incomplete or failed.

        For example this method will emit an error if we receive a "failed" status, e.g.
        with type "invalid_request_error" due to code "inference_rate_limit_exceeded".

        In other failures it will emit a debug level log.
        """
        if event.response.status == "completed":
            return

        provider_label = self._realtime_model._provider_label
        if event.response.status == "failed":
            if event.response.status_details and hasattr(event.response.status_details, "error"):
                error_type = getattr(event.response.status_details.error, "type", "unknown")
                error_body = event.response.status_details.error
                message = f"{provider_label} response failed with error type: {error_type}"
            else:
                error_body = None
                message = f"{provider_label} response failed with unknown error"
            # failures are largely undocumented by openai, so we assume optimistically
            # recoverable unless the code is a known-fatal one (quota / auth / billing),
            # which is raised so the recv loop breaks and _main_task stops reconnecting
            recoverable = not self._is_fatal_error(error_body)
            error = APIError(
                message=message,
                body=error_body,
                retryable=recoverable,
            )
            if not recoverable:
                raise error
            self._emit_error(error, recoverable=True)
        elif event.response.status in {"cancelled", "incomplete"}:
            status_details = event.response.status_details
            if isinstance(status_details, str):
                status_type = status_details
                status_reason = None
            else:
                status_type = status_details.type if status_details else None
                status_reason = status_details.reason if status_details else None
            logger.debug(
                "%s response done but not complete with status: %s (type=%s, reason=%s)",
                provider_label,
                event.response.status,
                status_type,
                status_reason,
                extra={
                    "event_id": event.response.id,
                    "event_response_status": event.response.status,
                    "event_response_status_type": status_type,
                    "event_response_status_reason": status_reason,
                },
            )
        else:
            logger.debug("Unknown response status: %s", event.response.status)

    def _is_fatal_error(self, error: object | None) -> bool:
        return _is_fatal_error(error)

    def _handle_error(self, event: RealtimeErrorEvent) -> None:
        if event_id := event.error.event_id:
            # a rejected item event gets no deleted/added reply, so fail its future rather than
            # leave update_chat_ctx to stall inside the speech that awaits it
            if fut := self._chat_ctx_event_futures.pop(event_id, None):
                if not fut.done():
                    # a duplicate id means the item is already there, as the create wanted
                    if event.error.code == "item_create_duplicate_item_id":
                        fut.set_result(None)
                    else:
                        fut.set_exception(llm.RealtimeError(event.error.message))
                # a terminal one still has to end the session, whatever it came in reply to
                if not self._is_fatal_error(event.error):
                    return
            # a rejected response.create gets no response.created; fail its future now
            # instead of orphaning it until the 10s timeout (still emitted/raised below)
            elif fut := self._response_created_futures.pop(event_id, None):
                if not fut.done():
                    fut.set_exception(llm.RealtimeError(event.error.message, code=event.error.code))

        if event.error.message.startswith("Cancellation failed"):
            return

        if event.error.code == "input_audio_buffer_commit_empty" and (
            self._opts.turn_detection is not None
        ):
            # the server VAD commits each segment itself, ours lands on an emptied buffer
            return

        provider_label = self._realtime_model._provider_label
        logger.error(
            f"{provider_label} returned an error: {event.error}",
            extra={"error": event.error},
        )
        recoverable = not self._is_fatal_error(event.error)
        error = APIError(
            message=f"{provider_label} returned an error",
            body=event.error,
            retryable=recoverable,
        )
        if not recoverable:
            # terminal (e.g. insufficient_quota): raise instead of emitting; the recv loop
            # re-raises it so _main_task emits it with recoverable=False and stops
            # reconnecting
            raise error
        self._emit_error(error, recoverable=True)

        # response errors are handled by _handle_response_done via _done_fut.
        # error events here are for non-response errors (e.g. invalid request).

    def _emit_error(self, error: Exception, recoverable: bool) -> None:
        self.emit(
            "error",
            llm.RealtimeModelError(
                timestamp=time.time(),
                label=self._realtime_model._label,
                error=error,
                recoverable=recoverable,
            ),
        )

A session for the OpenAI Realtime API.

This class is used to interact with the OpenAI Realtime API. It is responsible for sending events to the OpenAI Realtime API and receiving events from it.

It exposes two more events: - openai_server_event_received: expose the raw server events from the OpenAI Realtime API - openai_client_event_queued: expose the raw client events sent to the OpenAI Realtime API

Ancestors

  • livekit.agents.llm.realtime.RealtimeSession
  • abc.ABC
  • EventEmitter
  • typing.Generic

Subclasses

Instance variables

prop capabilities : llm.RealtimeCapabilities
Expand source code
@property
def capabilities(self) -> llm.RealtimeCapabilities:
    return self._capabilities

Capabilities of the session.

Defaults to the parent model's capabilities. Adapters that swap the underlying model mid-session override this to report the currently active model's capabilities.

prop chat_ctx : llm.ChatContext
Expand source code
@property
def chat_ctx(self) -> llm.ChatContext:
    return self._remote_chat_ctx.to_chat_ctx()
prop has_active_generation : bool
Expand source code
@property
def has_active_generation(self) -> bool:
    return self._current_generation is not None or len(self._response_created_futures) > 0
prop tools : llm.ToolContext
Expand source code
@property
def tools(self) -> llm.ToolContext:
    return self._tools.copy()

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    self._closing = True
    self._close_current_generation("session closed")
    self._msg_ch.close()
    await self._main_atask
def clear_audio(self) ‑> None
Expand source code
def clear_audio(self) -> None:
    self.send_event(InputAudioBufferClearEvent(type="input_audio_buffer.clear"))
    self._pushed_duration_s = 0
def commit_audio(self) ‑> None
Expand source code
def commit_audio(self) -> None:
    if self._pushed_duration_s > 0.1:  # OpenAI requires at least 100ms of audio
        self.send_event(InputAudioBufferCommitEvent(type="input_audio_buffer.commit"))
        self._pushed_duration_s = 0
def generate_reply(self,
*,
instructions: NotGivenOr[str] = NOT_GIVEN,
tool_choice: NotGivenOr[llm.ToolChoice] = NOT_GIVEN,
tools: NotGivenOr[list[llm.Tool]] = NOT_GIVEN) ‑> _asyncio.Future[livekit.agents.llm.realtime.GenerationCreatedEvent]
Expand source code
def generate_reply(
    self,
    *,
    instructions: NotGivenOr[str] = NOT_GIVEN,
    tool_choice: NotGivenOr[llm.ToolChoice] = NOT_GIVEN,
    tools: NotGivenOr[list[llm.Tool]] = NOT_GIVEN,
) -> asyncio.Future[llm.GenerationCreatedEvent]:
    event_id = utils.shortuuid("response_create_")
    fut = asyncio.Future[llm.GenerationCreatedEvent]()
    self._response_created_futures[event_id] = fut

    if is_given(instructions) and self._instructions:
        # in OpenAI realtime, the session-level instructions are completely replaced
        # by the new instructions for this response
        instructions = f"{self._instructions}\n{instructions}"

    params = RealtimeResponseCreateParams(
        instructions=instructions or None,
        metadata={"client_event_id": event_id},
    )
    if is_given(tool_choice):
        params.tool_choice = to_oai_tool_choice(tool_choice)
    if is_given(tools):
        params.tools = self._convert_tools_to_oai(tools)  # type: ignore

    self.send_event(
        ResponseCreateEvent(type="response.create", event_id=event_id, response=params)
    )

    def _on_timeout() -> None:
        self._response_created_futures.pop(event_id, None)
        if fut and not fut.done():
            # discard the response if the server still creates it after the timeout
            self._discarded_event_ids.add(event_id)
            fut.set_exception(llm.RealtimeError("generate_reply timed out."))

    handle = asyncio.get_event_loop().call_later(10.0, _on_timeout)

    def _on_fut_done(f: asyncio.Future[llm.GenerationCreatedEvent]) -> None:
        handle.cancel()
        self._response_created_futures.pop(event_id, None)
        if f.cancelled():
            # response.create was already sent; cancel the response server-side
            self.send_event(ResponseCancelEvent(type="response.cancel"))
            # the cancel above is a no-op if the response isn't created yet; discard it by id
            # when it arrives
            self._discarded_event_ids.add(event_id)

    fut.add_done_callback(_on_fut_done)
    return fut
def interrupt(self) ‑> None
Expand source code
def interrupt(self) -> None:
    if not self.has_active_generation:
        return
    self.send_event(ResponseCancelEvent(type="response.cancel"))
def push_audio(self, frame: rtc.AudioFrame) ‑> None
Expand source code
def push_audio(self, frame: rtc.AudioFrame) -> None:
    for f in self._resample_audio(frame):
        data = f.data.tobytes()
        for nf in self._bstream.write(data):
            self.send_event(
                InputAudioBufferAppendEvent(
                    type="input_audio_buffer.append",
                    audio=base64.b64encode(nf.data).decode("utf-8"),
                )
            )
            self._pushed_duration_s += nf.duration
def push_video(self, frame: rtc.VideoFrame) ‑> None
Expand source code
def push_video(self, frame: rtc.VideoFrame) -> None:
    message = llm.ChatMessage(
        role="user",
        content=[llm.ImageContent(image=frame)],
    )
    oai_item = livekit_item_to_openai_item(message)
    self.send_event(
        ConversationItemCreateEvent(
            type="conversation.item.create",
            item=oai_item,
            event_id=utils.shortuuid("video_"),
        )
    )
def send_event(self, event: RealtimeClientEvent | dict[str, Any]) ‑> None
Expand source code
def send_event(self, event: RealtimeClientEvent | dict[str, Any]) -> None:
    with contextlib.suppress(utils.aio.channel.ChanClosed):
        self._msg_ch.send_nowait(event)
def truncate(self,
*,
message_id: str,
modalities: "list[Literal['text', 'audio']]",
audio_end_ms: int,
audio_transcript: NotGivenOr[str] = NOT_GIVEN) ‑> None
Expand source code
def truncate(
    self,
    *,
    message_id: str,
    modalities: list[Literal["text", "audio"]],
    audio_end_ms: int,
    audio_transcript: NotGivenOr[str] = NOT_GIVEN,
) -> None:
    if "audio" in modalities:
        if audio_end_ms > 0:
            self.send_event(
                ConversationItemTruncateEvent(
                    type="conversation.item.truncate",
                    content_index=0,
                    item_id=message_id,
                    audio_end_ms=audio_end_ms,
                )
            )
        else:
            self.send_event(
                ConversationItemDeleteEvent(
                    type="conversation.item.delete",
                    item_id=message_id,
                    event_id=utils.shortuuid("chat_ctx_delete_"),
                )
            )
    elif utils.is_given(audio_transcript):
        # sync the forwarded text to the remote chat ctx
        chat_ctx = self.chat_ctx.copy(
            exclude_handoff=True,
            exclude_config_update=True,
        )
        if (idx := chat_ctx.index_by_id(message_id)) is not None:
            new_item = copy.copy(chat_ctx.items[idx])
            assert new_item.type == "message"

            new_item.content = [audio_transcript]
            chat_ctx.items[idx] = new_item
            events = self._create_update_chat_ctx_events(chat_ctx)
            for ev in events:
                self.send_event(ev)
async def update_chat_ctx(self, chat_ctx: llm.ChatContext) ‑> None
Expand source code
async def update_chat_ctx(self, chat_ctx: llm.ChatContext) -> None:
    async with self._update_chat_ctx_lock:
        chat_ctx = chat_ctx.copy(
            exclude_handoff=True,
            exclude_config_update=True,
        )
        # only remove the instructions but keep other system messages
        remove_instructions(chat_ctx)

        events = self._create_update_chat_ctx_events(chat_ctx)
        futs: list[asyncio.Future[None]] = []
        self._chat_ctx_event_futures = {}

        for ev in events:
            futs.append(f := asyncio.Future[None]())
            if isinstance(ev, ConversationItemDeleteEvent):
                self._item_delete_future[ev.item_id] = f
            else:
                assert ev.item.id is not None
                self._item_create_future[ev.item.id] = f

            # an updated item sends a delete and a create under the same id, so only the
            # event id tells a rejection which of the two it answers
            if ev.event_id:
                self._chat_ctx_event_futures[ev.event_id] = f
            self.send_event(ev)

        if not futs:
            return
        try:
            results = await asyncio.wait_for(
                asyncio.gather(*futs, return_exceptions=True), timeout=5.0
            )
        except asyncio.TimeoutError:
            raise llm.RealtimeError("update_chat_ctx timed out.") from None
        finally:
            self._chat_ctx_event_futures = {}
            for ev in events:
                if isinstance(ev, ConversationItemDeleteEvent):
                    self._item_delete_future.pop(ev.item_id, None)
                else:
                    assert ev.item.id is not None
                    self._item_create_future.pop(ev.item.id, None)

        # a rejected item is not worth failing the turn over
        if rejected := [str(r) for r in results if isinstance(r, BaseException)]:
            logger.warning(
                f"{self._realtime_model._provider_label} rejected part of a chat context update",  # noqa: E501
                extra={"errors": rejected},
            )
async def update_instructions(self, instructions: str) ‑> None
Expand source code
async def update_instructions(self, instructions: str) -> None:
    self.send_event(
        self._wrap_session_update(
            event_id=utils.shortuuid("instructions_update_"),
            session=RealtimeSessionCreateRequest.model_construct(
                type="realtime",
                instructions=instructions,
            ),
        )
    )
    self._instructions = instructions
def update_options(self,
*,
tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
voice: NotGivenOr[str] = NOT_GIVEN,
turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | None] = NOT_GIVEN,
max_response_output_tokens: "NotGivenOr[int | Literal['inf'] | None]" = NOT_GIVEN,
input_audio_transcription: NotGivenOr[AudioTranscription | None] = NOT_GIVEN,
input_audio_noise_reduction: NotGivenOr[NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None] = NOT_GIVEN,
speed: NotGivenOr[float] = NOT_GIVEN,
tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN) ‑> None
Expand source code
def update_options(
    self,
    *,
    tool_choice: NotGivenOr[llm.ToolChoice | None] = NOT_GIVEN,
    voice: NotGivenOr[str] = NOT_GIVEN,
    turn_detection: NotGivenOr[RealtimeAudioInputTurnDetection | None] = NOT_GIVEN,
    max_response_output_tokens: NotGivenOr[int | Literal["inf"] | None] = NOT_GIVEN,
    input_audio_transcription: NotGivenOr[AudioTranscription | None] = NOT_GIVEN,
    input_audio_noise_reduction: NotGivenOr[
        NoiseReductionType | NoiseReduction | InputAudioNoiseReduction | None
    ] = NOT_GIVEN,
    speed: NotGivenOr[float] = NOT_GIVEN,
    tracing: NotGivenOr[Tracing | None] = NOT_GIVEN,
    truncation: NotGivenOr[RealtimeTruncation | None] = NOT_GIVEN,
    reasoning: NotGivenOr[RealtimeReasoning | None] = NOT_GIVEN,
) -> None:
    session = RealtimeSessionCreateRequest(type="realtime")
    has_changes = False

    if is_given(tool_choice):
        current_oai = to_oai_tool_choice(self._opts.tool_choice)
        next_oai = to_oai_tool_choice(tool_choice)
        self._opts.tool_choice = tool_choice
        if current_oai != next_oai:
            session.tool_choice = next_oai
            has_changes = True

    if is_given(max_response_output_tokens):
        if self._opts.max_response_output_tokens != max_response_output_tokens:
            session.max_output_tokens = max_response_output_tokens
            has_changes = True
        self._opts.max_response_output_tokens = max_response_output_tokens

    if is_given(tracing):
        if self._opts.tracing != tracing:
            session.tracing = tracing  # type: ignore[assignment]
            has_changes = True
        self._opts.tracing = tracing

    if is_given(truncation):
        if self._opts.truncation != truncation:
            session.truncation = truncation
            has_changes = True
        self._opts.truncation = truncation

    if is_given(reasoning):
        if self._opts.reasoning != reasoning:
            # setting reasoning to None clears it server-side
            session.reasoning = reasoning
            has_changes = True
        self._opts.reasoning = reasoning

    has_audio_config = False
    audio_output = RealtimeAudioConfigOutput()
    audio_input = RealtimeAudioConfigInput()
    audio_config = RealtimeAudioConfig(output=audio_output, input=audio_input)

    if is_given(voice):
        if self._opts.voice != voice:
            audio_output.voice = voice
            has_audio_config = True
        self._opts.voice = voice

    if is_given(turn_detection):
        if self._opts.turn_detection != turn_detection:
            audio_input.turn_detection = turn_detection
            has_audio_config = True
        self._opts.turn_detection = turn_detection
        self._capabilities.turn_detection = _server_turn_taking_enabled(turn_detection)

    if is_given(input_audio_transcription):
        if self._opts.input_audio_transcription != input_audio_transcription:
            audio_input.transcription = input_audio_transcription
            has_audio_config = True
        self._opts.input_audio_transcription = input_audio_transcription
        self._capabilities.user_transcription = input_audio_transcription is not None

    if is_given(input_audio_noise_reduction):
        input_audio_noise_reduction = to_noise_reduction(input_audio_noise_reduction)
        if self._opts.input_audio_noise_reduction != input_audio_noise_reduction:
            audio_input.noise_reduction = input_audio_noise_reduction
            has_audio_config = True
        self._opts.input_audio_noise_reduction = input_audio_noise_reduction

    if is_given(speed):
        if self._opts.speed != speed:
            audio_output.speed = speed
            has_audio_config = True
        self._opts.speed = speed

    if has_audio_config:
        session.audio = audio_config
        has_changes = True

    if has_changes:
        self.send_event(
            self._wrap_session_update(
                event_id=utils.shortuuid("options_update_"), session=session
            )
        )
async def update_tools(self, tools: list[llm.Tool]) ‑> None
Expand source code
async def update_tools(self, tools: list[llm.Tool]) -> None:
    async with self._update_fnc_ctx_lock:
        ev = self._create_tools_update_event(tools)
        self.send_event(ev)

        retained_tool_names: set[str] = set()
        for t in ev["session"]["tools"]:
            if name := t.get("name"):
                retained_tool_names.add(name)
            # TODO(dz): handle MCP tools
        retained_tools = [
            tool
            for tool in tools
            if (
                isinstance(tool, (llm.FunctionTool, llm.RawFunctionTool))
                and tool.info.name in retained_tool_names
            )
            or isinstance(tool, llm.ProviderTool)
        ]
        self._tools = llm.ToolContext(retained_tools)

Inherited members

class ResponsesDelegationOptions (*args, **kwargs)
Expand source code
class ResponsesDelegationOptions(TypedDict, total=False):
    """The backend Responses model delegated work runs on, under ``delegation="responses"``.

    A key left unset is not sent, and the service's own default applies.
    """

    model: str
    """Responses model slug; ``gpt-5.6-luna`` when unset. On Azure, the name of a Responses
    deployment in the same resource, which is required there."""
    instructions: str
    """Instructions for the backend model, distinct from the voice model's."""
    tool_choice: llm.ToolChoice | None
    parallel_tool_calls: bool
    reasoning: Reasoning
    """Responses reasoning settings, for example ``{"effort": "medium"}``."""
    text: ResponseTextConfigParam
    """Responses text settings, for example ``{"verbosity": "low"}``."""
    service_tier: types.ServiceTier
    max_output_tokens: int
    """Upper bound on the tokens one backend response may generate; at least 16."""

The backend Responses model delegated work runs on, under delegation="responses".

A key left unset is not sent, and the service's own default applies.

Ancestors

  • builtins.dict

Class variables

var instructions : str

Instructions for the backend model, distinct from the voice model's.

var max_output_tokens : int

Upper bound on the tokens one backend response may generate; at least 16.

var model : str

Responses model slug; gpt-5.6-luna when unset. On Azure, the name of a Responses deployment in the same resource, which is required there.

var parallel_tool_calls : bool
var reasoning : openai.types.shared_params.reasoning.Reasoning

Responses reasoning settings, for example {"effort": "medium"}.

var service_tier : Literal['auto', 'default', 'flex', 'priority', 'ultrafast']
var text : openai.types.responses.response_text_config_param.ResponseTextConfigParam

Responses text settings, for example {"verbosity": "low"}.

var tool_choice : livekit.agents.llm.tool_context.NamedToolChoice | Literal['auto', 'required', 'none'] | None