#
# Copyright (c) 2024-2026, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#
"""LiveKit transport implementation for Pipecat.
This module provides comprehensive LiveKit real-time communication integration
including audio streaming, data messaging, participant management, and room
event handling for conversational AI applications.
"""
import asyncio
import json
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
from typing import Any
from loguru import logger
from pydantic import BaseModel
from pipecat.audio.dtmf.types import KeypadEntry
from pipecat.audio.utils import create_stream_resampler
from pipecat.frames.frames import (
AudioRawFrame,
BotConnectedFrame,
CancelFrame,
ClientConnectedFrame,
EndFrame,
Frame,
ImageRawFrame,
InputDTMFFrame,
InputTransportMessageFrame,
InterruptionFrame,
OutputAudioRawFrame,
OutputDTMFFrame,
OutputDTMFUrgentFrame,
OutputImageRawFrame,
OutputTransportMessageFrame,
OutputTransportMessageUrgentFrame,
StartFrame,
UserAudioRawFrame,
UserImageRawFrame,
)
from pipecat.processors.frame_processor import FrameDirection, FrameProcessorSetup
from pipecat.transports.base_input import BaseInputTransport
from pipecat.transports.base_output import BaseOutputTransport
from pipecat.transports.base_transport import BaseTransport, TransportParams
from pipecat.utils.asyncio.task_manager import BaseTaskManager
try:
from livekit import rtc
from livekit.rtc._proto import video_frame_pb2 as proto_video_frame
from tenacity import retry, stop_after_attempt, wait_exponential
except ModuleNotFoundError as e:
logger.error(f"Exception: {e}")
logger.error('In order to use LiveKit, you need to `uv add "pipecat-ai[livekit]"`.')
raise ImportError(f"Missing module: {e}") from e
# DTMF mapping according to RFC 4733
DTMF_CODE_MAP = {
"0": 0,
"1": 1,
"2": 2,
"3": 3,
"4": 4,
"5": 5,
"6": 6,
"7": 7,
"8": 8,
"9": 9,
"*": 10,
"#": 11,
}
# Maps Pipecat's PIL-style color format strings (``OutputImageRawFrame.format``,
# configured via ``TransportParams.video_out_color_format``) to LiveKit's
# ``VideoBufferType`` enum used by ``rtc.VideoFrame`` and its bytes per pixel.
LIVEKIT_VIDEO_BUFFER_TYPES = {
"RGB": (proto_video_frame.VideoBufferType.RGB24, 3),
"RGBA": (proto_video_frame.VideoBufferType.RGBA, 4),
"BGRA": (proto_video_frame.VideoBufferType.BGRA, 4),
"ARGB": (proto_video_frame.VideoBufferType.ARGB, 4),
}
[docs]
@dataclass
class LiveKitOutputTransportMessageFrame(OutputTransportMessageFrame):
"""Frame for transport messages in LiveKit rooms.
Parameters:
participant_id: Optional ID of the participant this message is for/from.
"""
participant_id: str | None = None
[docs]
@dataclass
class LiveKitOutputTransportMessageUrgentFrame(OutputTransportMessageUrgentFrame):
"""Frame for urgent transport messages in LiveKit rooms.
Parameters:
participant_id: Optional ID of the participant this message is for/from.
"""
participant_id: str | None = None
[docs]
class LiveKitParams(TransportParams):
"""Configuration parameters for LiveKit transport.
Video output publishes a single ``"pipecat-video"`` camera track (mirroring how
audio output always publishes one ``"pipecat-audio"`` microphone track) when
``video_out_enabled`` is set. The track is sized using
``video_out_width``/``video_out_height`` and encodes frames according to
``video_out_color_format`` (default ``"RGB"``); ``video_out_framerate``
governs how often ``BaseOutputTransport`` draws frames. Per-destination
video routing (multiple named output tracks, as supported by Daily's
``camera_out_enabled``/``register_video_destination``) is not yet
implemented for LiveKit.
``video_out_codec`` selects the published video codec (``"VP8"``, ``"H264"``,
``"VP9"``, ``"AV1"`` or ``"H265"``); LiveKit picks one when it is unset.
Parameters:
audio_out_queue_size_ms: Buffer size of the outgoing audio source, in milliseconds
(LiveKit's default is 1000).
video_out_max_bitrate: Maximum bitrate of the published video track, in bits
per second, capped at ``video_out_framerate``. LiveKit chooses the encoding
from the track resolution when unset.
"""
audio_out_queue_size_ms: int = 1000
video_out_max_bitrate: int | None = None
[docs]
class LiveKitCallbacks(BaseModel):
"""Callback handlers for LiveKit events.
Parameters:
on_connected: Called when connected to the LiveKit room.
on_disconnected: Called when disconnected from the LiveKit room.
on_participant_connected: Called when a participant joins the room.
on_participant_disconnected: Called when a participant leaves the room.
on_audio_track_subscribed: Called when an audio track is subscribed.
on_audio_track_unsubscribed: Called when an audio track is unsubscribed.
on_data_received: Called when data is received. The sender is None for
packets sent by a server SDK, which LiveKit delivers unattributed.
on_first_participant_joined: Called when the first participant joins.
on_dtmf_event: Called when a SIP DTMF tone is received.
"""
on_connected: Callable[[], Awaitable[None]]
on_disconnected: Callable[[], Awaitable[None]]
on_before_disconnect: Callable[[], Awaitable[None]]
on_participant_connected: Callable[[str], Awaitable[None]]
on_participant_disconnected: Callable[[str], Awaitable[None]]
on_audio_track_subscribed: Callable[[str], Awaitable[None]]
on_audio_track_unsubscribed: Callable[[str], Awaitable[None]]
on_video_track_subscribed: Callable[[str], Awaitable[None]]
on_video_track_unsubscribed: Callable[[str], Awaitable[None]]
on_data_received: Callable[[bytes, str | None], Awaitable[None]]
on_first_participant_joined: Callable[[str], Awaitable[None]]
on_dtmf_event: Callable[[Any], Awaitable[None]]
[docs]
class LiveKitTransportClient:
"""Core client for interacting with LiveKit rooms.
Manages the connection to LiveKit rooms and handles all low-level API interactions
including room management, audio streaming, data messaging, and event handling.
"""
[docs]
def __init__(
self,
url: str,
token: str,
room_name: str,
params: LiveKitParams,
callbacks: LiveKitCallbacks,
transport_name: str,
):
"""Initialize the LiveKit transport client.
Args:
url: LiveKit server URL to connect to.
token: Authentication token for the room.
room_name: Name of the LiveKit room to join.
params: Configuration parameters for the transport.
callbacks: Event callback handlers.
transport_name: Name identifier for the transport.
"""
self._url = url
self._token = token
self._room_name = room_name
self._params = params
self._callbacks = callbacks
self._transport_name = transport_name
self._room: rtc.Room | None = None
self._participant_id: str = ""
self._connected = False
self._disconnect_counter = 0
self._audio_source: rtc.AudioSource | None = None
self._audio_track: rtc.LocalAudioTrack | None = None
self._audio_tracks = {}
self._audio_queue = asyncio.Queue()
# Per-participant ``(AudioStream, Task)`` so unsubscribe can close
# the owned native stream and cancel its producer task instead of
# leaking both on every track republish.
self._audio_streams: dict[str, tuple[rtc.AudioStream, asyncio.Task]] = {}
self._video_source: rtc.VideoSource | None = None
self._video_track: rtc.LocalVideoTrack | None = None
self._video_tracks = {}
self._video_queue = asyncio.Queue()
# Symmetric registry for video streams.
self._video_streams: dict[str, tuple[rtc.VideoStream, asyncio.Task]] = {}
self._other_participant_has_joined = False
self._task_manager: BaseTaskManager | None = None
self._async_lock = asyncio.Lock()
@property
def participant_id(self) -> str:
"""Get the participant ID for this client.
Returns:
The participant ID assigned by LiveKit.
"""
return self._participant_id
@property
def room(self) -> rtc.Room:
"""Get the LiveKit room instance.
Returns:
The LiveKit room object.
Raises:
Exception: If room object is not available.
"""
if not self._room:
raise Exception(f"{self}: missing room object (pipeline not started?)")
return self._room
[docs]
async def setup(self, setup: FrameProcessorSetup):
"""Setup the client with task manager and room initialization.
Args:
setup: The frame processor setup configuration.
"""
if self._task_manager:
return
self._task_manager = setup.task_manager
self._room = rtc.Room(loop=self._task_manager.get_event_loop())
self._out_sample_rate = self._params.audio_out_sample_rate or setup.audio_out_sample_rate
# Set up room event handlers
self.room.on("participant_connected")(self._on_participant_connected_wrapper)
self.room.on("participant_disconnected")(self._on_participant_disconnected_wrapper)
self.room.on("track_subscribed")(self._on_track_subscribed_wrapper)
self.room.on("track_unsubscribed")(self._on_track_unsubscribed_wrapper)
self.room.on("data_received")(self._on_data_received_wrapper)
self.room.on("connected")(self._on_connected_wrapper)
self.room.on("disconnected")(self._on_disconnected_wrapper)
self.room.on("sip_dtmf_received")(self._on_sip_dtmf_received_wrapper)
[docs]
async def cleanup(self):
"""Cleanup client resources."""
await self.disconnect()
[docs]
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=4, max=10))
async def connect(self):
"""Connect to the LiveKit room with retry logic."""
async with self._async_lock:
if self._connected:
# Increment disconnect counter if already connected.
self._disconnect_counter += 1
return
logger.info(f"Connecting to {self._room_name}")
try:
await self.room.connect(
self._url,
self._token,
options=rtc.RoomOptions(auto_subscribe=True),
)
self._participant_id = self.room.local_participant.identity
logger.info(f"Connected to {self._room_name} as {self._participant_id}")
# Set up audio source and track
self._audio_source = rtc.AudioSource(
self._out_sample_rate,
self._params.audio_out_channels,
queue_size_ms=self._params.audio_out_queue_size_ms,
)
self._audio_track = rtc.LocalAudioTrack.create_audio_track(
"pipecat-audio", self._audio_source
)
options = rtc.TrackPublishOptions()
options.source = rtc.TrackSource.SOURCE_MICROPHONE
await self.room.local_participant.publish_track(self._audio_track, options)
# Set up video source and track (only if video output is
# enabled; unlike audio, which is always published).
if self._params.video_out_enabled:
self._video_source = rtc.VideoSource(
self._params.video_out_width, self._params.video_out_height
)
self._video_track = rtc.LocalVideoTrack.create_video_track(
"pipecat-video", self._video_source
)
video_options = self._video_publish_options()
await self.room.local_participant.publish_track(
self._video_track, video_options
)
# Only mark the client connected once its tracks are
# published, so a retry after a failed publish starts over
# instead of returning early without tracks.
self._connected = True
# Increment disconnect counter if we successfully connected.
self._disconnect_counter += 1
await self._callbacks.on_connected()
# Check if there are already participants in the room
participants = self.get_participants()
if participants and not self._other_participant_has_joined:
self._other_participant_has_joined = True
await self._callbacks.on_first_participant_joined(participants[0])
except Exception as e:
logger.error(f"Error connecting to {self._room_name}: {e}")
if not self._connected:
await self._rollback_partial_connect()
raise
def _video_publish_options(self) -> rtc.TrackPublishOptions:
"""Build the publish options for the video track from the params."""
options = rtc.TrackPublishOptions()
options.source = rtc.TrackSource.SOURCE_CAMERA
# LiveKit requires both fields of a video encoding, so the framerate
# is only sent along with a bitrate.
if self._params.video_out_max_bitrate is not None:
options.video_encoding.max_bitrate = self._params.video_out_max_bitrate
options.video_encoding.max_framerate = self._params.video_out_framerate
codec = self._params.video_out_codec
if codec:
try:
options.video_codec = rtc.VideoCodec.Value(codec.upper())
except ValueError:
logger.warning(
f"{self} unsupported video codec for LiveKit output: {codec!r}, "
f"expected one of {list(rtc.VideoCodec.keys())}"
)
return options
async def _rollback_partial_connect(self):
"""Undo a connection attempt that failed before it completed."""
await self._close_output_sources()
try:
await self.room.disconnect()
except Exception as e:
logger.warning(f"{self} error disconnecting after failed connect: {e}")
[docs]
async def disconnect(self):
"""Disconnect from the LiveKit room."""
async with self._async_lock:
# Decrement leave counter when leaving.
self._disconnect_counter -= 1
if not self._connected or self._disconnect_counter > 0:
return
logger.info(f"Disconnecting from {self._room_name}")
await self._callbacks.on_before_disconnect()
# Mark the client disconnected before the room disconnects, so the
# room's own "disconnected" event does not report it a second time.
self._connected = False
await self.room.disconnect()
await self._close_output_sources()
# Close any remaining per-participant streams and cancel their
# producer tasks so they do not outlive the connection.
await self._close_all_streams()
logger.info(f"Disconnected from {self._room_name}")
await self._callbacks.on_disconnected()
async def _close_output_sources(self):
"""Close the published audio and video sources.
``room.disconnect()`` does not release the native source handles, so
each connection would otherwise leave them behind.
"""
audio_source, self._audio_source, self._audio_track = self._audio_source, None, None
video_source, self._video_source, self._video_track = self._video_source, None, None
for source in (audio_source, video_source):
if source is None:
continue
try:
await source.aclose()
except Exception as e:
logger.warning(f"{self} error closing output source: {e}")
[docs]
async def send_data(self, data: bytes, participant_id: str | None = None):
"""Send data to participants in the room.
Args:
data: The data bytes to send.
participant_id: Optional specific participant to send to.
"""
if not self._connected:
return
try:
if participant_id:
await self.room.local_participant.publish_data(
data, reliable=True, destination_identities=[participant_id]
)
else:
await self.room.local_participant.publish_data(data, reliable=True)
except Exception as e:
logger.error(f"Error sending data: {e}")
[docs]
async def send_dtmf(self, digit: str):
r"""Send DTMF tone to the room.
Args:
digit: The DTMF digit to send (0-9, \*, #).
"""
if not self._connected:
return
if digit not in DTMF_CODE_MAP:
logger.warning(f"Invalid DTMF digit: {digit}")
return
code = DTMF_CODE_MAP[digit]
try:
await self.room.local_participant.publish_dtmf(code=code, digit=digit)
except Exception as e:
logger.error(f"Error sending DTMF tone {digit}: {e}")
[docs]
async def publish_audio(self, audio_frame: rtc.AudioFrame) -> bool:
"""Publish an audio frame to the room.
Args:
audio_frame: The LiveKit audio frame to publish.
"""
if not self._connected or not self._audio_source:
return False
try:
await self._audio_source.capture_frame(audio_frame)
return True
except Exception as e:
# When using an audio mixer, the base output transport's
# with_mixer() generator continuously yields frames (mixed with
# background audio) even when no TTS audio is queued. During
# interruptions, the audio task is cancelled and recreated, but
# there is a brief window where the native LiveKit AudioSource
# rejects capture_frame() with an InvalidState error. This is a
# transient condition — the mixer will produce a new frame within
# milliseconds, so we silently drop these frames.
if "InvalidState" not in str(e):
logger.error(f"Error publishing audio: {e}")
return False
[docs]
async def publish_video(self, video_frame: rtc.VideoFrame) -> bool:
"""Publish a video frame to the room.
Args:
video_frame: The LiveKit video frame to publish.
Returns:
True if the video frame was published successfully, False otherwise.
"""
if not self._connected or not self._video_source:
return False
try:
# Unlike ``AudioSource.capture_frame``, ``VideoSource.capture_frame``
# is synchronous in livekit-rtc.
self._video_source.capture_frame(video_frame)
return True
except Exception as e:
logger.error(f"Error publishing video: {e}")
return False
[docs]
def get_participants(self) -> list[str]:
"""Get list of participant IDs in the room.
Returns:
List of participant LiveKit identities.
"""
return [p.identity for p in self.room.remote_participants.values()]
[docs]
async def mute_participant(self, participant_id: str):
"""Stop receiving a specific participant's audio.
LiveKit doesn't let one participant force-mute another's microphone;
this unsubscribes the bot from their audio track instead.
Args:
participant_id: LiveKit identity of the participant to stop
receiving audio from.
"""
participant = self.room.remote_participants.get(participant_id)
if participant:
for publication in participant.track_publications.values():
if publication.kind == rtc.TrackKind.KIND_AUDIO:
publication.set_subscribed(False)
[docs]
async def unmute_participant(self, participant_id: str):
"""Resume receiving a specific participant's audio.
Args:
participant_id: LiveKit identity of the participant to resume
receiving audio from.
"""
participant = self.room.remote_participants.get(participant_id)
if participant:
for publication in participant.track_publications.values():
if publication.kind == rtc.TrackKind.KIND_AUDIO:
publication.set_subscribed(True)
# Wrapper methods for event handlers
def _on_participant_connected_wrapper(self, participant: rtc.RemoteParticipant):
"""Wrapper for participant connected events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_participant_connected(participant),
f"{self}::_async_on_participant_connected",
)
def _on_participant_disconnected_wrapper(self, participant: rtc.RemoteParticipant):
"""Wrapper for participant disconnected events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_participant_disconnected(participant),
f"{self}::_async_on_participant_disconnected",
)
def _on_track_subscribed_wrapper(
self,
track: rtc.Track,
publication: rtc.RemoteTrackPublication,
participant: rtc.RemoteParticipant,
):
"""Wrapper for track subscribed events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_track_subscribed(track, publication, participant),
f"{self}::_async_on_track_subscribed",
)
def _on_track_unsubscribed_wrapper(
self,
track: rtc.Track,
publication: rtc.RemoteTrackPublication,
participant: rtc.RemoteParticipant,
):
"""Wrapper for track unsubscribed events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_track_unsubscribed(track, publication, participant),
f"{self}::_async_on_track_unsubscribed",
)
def _on_data_received_wrapper(self, data: rtc.DataPacket):
"""Wrapper for data received events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_data_received(data),
f"{self}::_async_on_data_received",
)
def _on_connected_wrapper(self):
"""Wrapper for connected events."""
assert self._task_manager is not None
self._task_manager.create_task(self._async_on_connected(), f"{self}::_async_on_connected")
def _on_disconnected_wrapper(self):
"""Wrapper for disconnected events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_disconnected(), f"{self}::_async_on_disconnected"
)
def _on_sip_dtmf_received_wrapper(self, dtmf: rtc.SipDTMF):
"""Wrapper for inbound SIP DTMF events."""
assert self._task_manager is not None
self._task_manager.create_task(
self._async_on_sip_dtmf_received(dtmf),
f"{self}::_async_on_sip_dtmf_received",
)
# Async methods for event handling
async def _async_on_participant_connected(self, participant: rtc.RemoteParticipant):
"""Handle participant connected events."""
logger.info(f"Participant connected: {participant.identity}")
await self._callbacks.on_participant_connected(participant.identity)
if not self._other_participant_has_joined:
self._other_participant_has_joined = True
await self._callbacks.on_first_participant_joined(participant.identity)
async def _async_on_participant_disconnected(self, participant: rtc.RemoteParticipant):
"""Handle participant disconnected events."""
logger.info(f"Participant disconnected: {participant.identity}")
await self._callbacks.on_participant_disconnected(participant.identity)
if len(self.get_participants()) == 0:
self._other_participant_has_joined = False
async def _async_on_track_subscribed(
self,
track: rtc.Track,
publication: rtc.RemoteTrackPublication,
participant: rtc.RemoteParticipant,
):
"""Handle track subscribed events."""
assert self._task_manager is not None
if track.kind == rtc.TrackKind.KIND_AUDIO:
logger.info(
f"Audio track subscribed: {track.sid} from participant {participant.identity}"
)
# If the participant is re-publishing (e.g. mute/unmute cycle),
# close + cancel the previous stream/task before replacing the
# registry entry, so two producers never feed ``_audio_queue``
# for the same participant.
await self._close_audio_stream(participant.identity)
self._audio_tracks[participant.identity] = track
audio_stream = rtc.AudioStream(track)
task = self._task_manager.create_task(
self._process_audio_stream(audio_stream, participant.identity),
f"{self}::_process_audio_stream",
)
self._audio_streams[participant.identity] = (audio_stream, task)
await self._callbacks.on_audio_track_subscribed(participant.identity)
elif track.kind == rtc.TrackKind.KIND_VIDEO:
logger.info(
f"Video track subscribed: {track.sid} from participant {participant.identity}"
)
# Symmetric: clean up any prior video stream/task for the same
# participant before replacing.
await self._close_video_stream(participant.identity)
self._video_tracks[participant.identity] = track
# Only process video stream if video input is enabled to prevent
# unbounded queue growth when there is no consumer for video frames.
if self._params.video_in_enabled:
video_stream = rtc.VideoStream(track)
task = self._task_manager.create_task(
self._process_video_stream(video_stream, participant.identity),
f"{self}::_process_video_stream",
)
self._video_streams[participant.identity] = (video_stream, task)
await self._callbacks.on_video_track_subscribed(participant.identity)
async def _async_on_track_unsubscribed(
self,
track: rtc.Track,
publication: rtc.RemoteTrackPublication,
participant: rtc.RemoteParticipant,
):
"""Handle track unsubscribed events."""
logger.info(f"Track unsubscribed: {publication.sid} from {participant.identity}")
if track.kind == rtc.TrackKind.KIND_AUDIO:
await self._close_audio_stream(participant.identity)
await self._callbacks.on_audio_track_unsubscribed(participant.identity)
elif track.kind == rtc.TrackKind.KIND_VIDEO:
await self._close_video_stream(participant.identity)
await self._callbacks.on_video_track_unsubscribed(participant.identity)
async def _close_audio_stream(self, participant_id: str) -> None:
"""Close a participant's owned audio stream and cancel its producer task.
Idempotent: no-op when there is no registered stream for the
participant.
"""
entry = self._audio_streams.pop(participant_id, None)
if entry is None:
return
stream, task = entry
try:
await asyncio.wait_for(stream.aclose(), timeout=2.0)
except Exception as e:
logger.warning(f"AudioStream.aclose failed for {participant_id}: {e}")
if task is not None and not task.done():
task.cancel()
async def _close_video_stream(self, participant_id: str) -> None:
"""Close a participant's owned video stream and cancel its producer task.
Idempotent: no-op when there is no registered stream for the
participant.
"""
entry = self._video_streams.pop(participant_id, None)
if entry is None:
return
stream, task = entry
try:
await asyncio.wait_for(stream.aclose(), timeout=2.0)
except Exception as e:
logger.warning(f"VideoStream.aclose failed for {participant_id}: {e}")
if task is not None and not task.done():
task.cancel()
async def _close_all_streams(self) -> None:
"""Close every per-participant audio/video stream and cancel its task.
Idempotent: no-op when no streams are registered.
"""
for participant_id in list(self._audio_streams.keys()):
await self._close_audio_stream(participant_id)
for participant_id in list(self._video_streams.keys()):
await self._close_video_stream(participant_id)
async def _async_on_data_received(self, data: rtc.DataPacket):
"""Handle data received events."""
# LiveKit delivers packets sent by a server SDK with no participant.
sender = data.participant.identity if data.participant else None
await self._callbacks.on_data_received(data.data, sender)
async def _async_on_connected(self):
"""Handle connected events."""
await self._callbacks.on_connected()
async def _async_on_disconnected(self, reason=None):
"""Handle disconnected events.
Only a disconnect of a connected client is reported. ``disconnect()``
and a failed ``connect()`` clear the connected flag before disconnecting
the room, so the room's event for those is ignored.
"""
if not self._connected:
return
self._connected = False
logger.info(f"Disconnected from {self._room_name}. Reason: {reason}")
await self._callbacks.on_disconnected()
async def _async_on_sip_dtmf_received(self, dtmf: rtc.SipDTMF):
"""Handle inbound SIP DTMF events from LiveKit telephony."""
participant = getattr(dtmf, "participant", None)
participant_id = getattr(participant, "identity", None) if participant else None
data = {
"tone": dtmf.digit,
"digit": dtmf.digit,
"code": dtmf.code,
"participant_id": participant_id,
}
logger.debug(f"{self} SIP DTMF event: {data}")
await self._callbacks.on_dtmf_event(data)
async def _process_audio_stream(self, audio_stream: rtc.AudioStream, participant_id: str):
"""Process incoming audio stream from a participant."""
logger.info(f"Started processing audio stream for participant {participant_id}")
async for event in audio_stream:
if isinstance(event, rtc.AudioFrameEvent):
await self._audio_queue.put((event, participant_id))
else:
logger.warning(f"Received unexpected event type: {type(event)}")
[docs]
async def get_next_audio_frame(self):
"""Get the next audio frame from the queue."""
while True:
frame, participant_id = await self._audio_queue.get()
yield frame, participant_id
async def _process_video_stream(self, video_stream: rtc.VideoStream, participant_id: str):
"""Process incoming video stream from a participant."""
logger.info(f"Started processing video stream for participant {participant_id}")
async for event in video_stream:
if isinstance(event, rtc.VideoFrameEvent):
await self._video_queue.put((event, participant_id))
else:
logger.warning(f"Received unexpected event type: {type(event)}")
[docs]
async def get_next_video_frame(self):
"""Get the next video frame from the queue."""
while True:
frame, participant_id = await self._video_queue.get()
yield frame, participant_id
def __str__(self):
"""String representation of the LiveKit transport client."""
return f"{self._transport_name}::LiveKitTransportClient"
[docs]
class LiveKitOutputTransport(BaseOutputTransport):
"""Handles outgoing media streams and events to LiveKit rooms.
Manages sending audio and video frames and data messages to LiveKit room
participants, including audio/video format conversion for LiveKit
compatibility. Video output publishes to a single default camera track
when ``LiveKitParams.video_out_enabled`` is set.
"""
[docs]
def __init__(
self,
transport: BaseTransport,
client: LiveKitTransportClient,
params: LiveKitParams,
**kwargs,
):
"""Initialize the LiveKit output transport.
Args:
transport: The parent transport instance.
client: LiveKitTransportClient instance.
params: Configuration parameters.
**kwargs: Additional arguments passed to parent class.
"""
super().__init__(params, **kwargs)
self._transport = transport
self._client = client
# Formats already reported as unsupported, so the error is logged once
# per format instead of once per frame.
self._unsupported_video_formats: set[str | None] = set()
[docs]
async def setup(self, setup: FrameProcessorSetup):
"""Setup the output transport with shared client setup.
Args:
setup: The frame processor setup configuration.
"""
await super().setup(setup)
await self._client.setup(setup)
await self._client.connect()
logger.info("LiveKitOutputTransport connected")
[docs]
async def cleanup(self):
"""Release output transport resources at teardown."""
await super().cleanup()
await self._client.disconnect()
await self._transport.cleanup()
[docs]
async def start(self, frame: StartFrame):
"""Start the output transport.
Args:
frame: The start frame containing initialization parameters.
"""
await super().start(frame)
await self.set_transport_ready(frame)
[docs]
async def stop(self, frame: EndFrame):
"""Stop the output transport and disconnect from LiveKit room.
Args:
frame: The end frame signaling transport shutdown.
"""
await super().stop(frame)
await self._client.disconnect()
logger.info("LiveKitOutputTransport stopped")
[docs]
async def cancel(self, frame: CancelFrame):
"""Cancel the output transport and disconnect from LiveKit room.
Args:
frame: The cancel frame signaling immediate cancellation.
"""
await super().cancel(frame)
await self._client.disconnect()
[docs]
async def process_frame(self, frame: Frame, direction: FrameDirection):
"""Process frames, clearing the LiveKit AudioSource buffer on interruption.
When an InterruptionFrame arrives, any audio already submitted to the
LiveKit AudioSource (but not yet played out) is cleared immediately so
the bot stops speaking without delay.
Args:
frame: The frame to process.
direction: The direction of frame flow in the pipeline.
"""
await super().process_frame(frame, direction)
if isinstance(frame, InterruptionFrame) and self._client._audio_source is not None:
self._client._audio_source.clear_queue()
[docs]
async def send_message(
self, frame: OutputTransportMessageFrame | OutputTransportMessageUrgentFrame
):
"""Send a transport message to participants.
Args:
frame: The transport message frame to send.
"""
message = frame.message
if isinstance(message, dict):
# fix message encoding for dict-like messages, e.g. RTVI messages.
message = json.dumps(message, ensure_ascii=False)
if isinstance(
frame, (LiveKitOutputTransportMessageFrame, LiveKitOutputTransportMessageUrgentFrame)
):
await self._client.send_data(message.encode(), frame.participant_id)
else:
await self._client.send_data(message.encode())
[docs]
async def write_audio_frame(self, frame: OutputAudioRawFrame) -> bool:
"""Write an audio frame to the LiveKit room.
Args:
frame: The audio frame to write.
Returns:
True if the audio frame was written successfully, False otherwise.
"""
livekit_audio = self._convert_pipecat_audio_to_livekit(frame.audio)
return await self._client.publish_audio(livekit_audio)
[docs]
async def write_video_frame(self, frame: OutputImageRawFrame) -> bool:
"""Write a video frame to the LiveKit room's published camera track.
Publishes to the single default video track set up in
``LiveKitTransportClient.connect`` (mirroring how audio always
publishes to one microphone track). Per-destination routing to
multiple named video tracks is not supported yet, so
``frame.transport_destination`` is ignored.
Args:
frame: The video frame to write.
Returns:
True if the video frame was written successfully, False otherwise.
"""
livekit_video = self._convert_pipecat_video_to_livekit(frame)
if livekit_video is None:
return False
return await self._client.publish_video(livekit_video)
def _supports_native_dtmf(self) -> bool:
"""LiveKit supports native DTMF via telephone events.
Returns:
True, as LiveKit supports native DTMF transmission.
"""
return True
async def _write_dtmf_native(self, frame: OutputDTMFFrame | OutputDTMFUrgentFrame):
"""Use LiveKit's native publish_dtmf method for telephone events.
LiveKit's DTMF API sends a single tone per call, so when
``frame.buttons`` contains multiple entries only the first one is
sent.
Args:
frame: The DTMF frame to write.
"""
if not frame.buttons:
return
await self._client.send_dtmf(frame.buttons[0].value)
def _convert_pipecat_audio_to_livekit(self, pipecat_audio: bytes) -> rtc.AudioFrame:
"""Convert Pipecat audio data to LiveKit audio frame."""
bytes_per_sample = 2 # Assuming 16-bit audio
total_samples = len(pipecat_audio) // bytes_per_sample
samples_per_channel = total_samples // self._params.audio_out_channels
return rtc.AudioFrame(
data=pipecat_audio,
sample_rate=self.sample_rate,
num_channels=self._params.audio_out_channels,
samples_per_channel=samples_per_channel,
)
def _convert_pipecat_video_to_livekit(
self, frame: OutputImageRawFrame
) -> rtc.VideoFrame | None:
"""Convert a Pipecat output video frame to a LiveKit video frame.
Returns:
The converted ``rtc.VideoFrame``, or None if ``frame.format`` has
no known LiveKit ``VideoBufferType`` mapping or the image length
does not match its size.
"""
buffer_info = LIVEKIT_VIDEO_BUFFER_TYPES.get(frame.format) if frame.format else None
if buffer_info is None:
if frame.format not in self._unsupported_video_formats:
self._unsupported_video_formats.add(frame.format)
logger.error(
f"{self} unsupported video color format for LiveKit output: {frame.format!r}"
)
return None
buffer_type, bytes_per_pixel = buffer_info
width, height = frame.size
# LiveKit reads ``width * height * bytes_per_pixel`` bytes from the
# buffer without checking its length, so a short buffer crashes the
# process.
expected_length = width * height * bytes_per_pixel
if len(frame.image) != expected_length:
logger.error(
f"{self} video frame of size {width}x{height} and format {frame.format!r} "
f"has {len(frame.image)} bytes, expected {expected_length}"
)
return None
return rtc.VideoFrame(width, height, buffer_type, frame.image)
[docs]
class LiveKitTransport(BaseTransport):
"""Transport implementation for LiveKit real-time communication.
Provides comprehensive LiveKit integration including audio streaming, data
messaging, participant management, and room event handling for conversational
AI applications.
Every ``participant_id`` surfaced by this transport (event args, frame
fields, ``get_participants()``, and the ``get_participant_metadata``/
``mute_participant``/``unmute_participant`` methods) is the participant's
LiveKit *identity* (``rtc.Participant.identity``) — the value set when
minting its access token, and what LiveKit itself keys
``room.remote_participants`` by and expects in ``destination_identities``.
It is not the participant's *SID* (``rtc.Participant.sid``), a
per-connection session id that changes on every reconnect.
Event handlers available:
- on_connected: Called when the bot connects to the room.
- on_disconnected: Called when the bot disconnects from the room.
- on_before_disconnect: [sync] Called just before the bot disconnects.
- on_call_state_updated: Called when the call state changes. Args: (state: str)
- on_first_participant_joined: Called when the first participant joins.
Args: (participant_id: str)
- on_participant_connected: Called when a participant connects.
Args: (participant_id: str)
- on_participant_disconnected: Called when a participant disconnects.
Args: (participant_id: str)
- on_participant_left: Called when a participant leaves.
Args: (participant_id: str, reason: str)
- on_client_connected: Called when a participant connects (alias for
on_participant_connected). Args: (participant: dict)
- on_client_disconnected: Called when a participant disconnects (alias for
on_participant_disconnected). Args: (participant: dict)
- on_audio_track_subscribed: Called when an audio track is subscribed.
Args: (participant_id: str)
- on_audio_track_unsubscribed: Called when an audio track is unsubscribed.
Args: (participant_id: str)
- on_video_track_subscribed: Called when a video track is subscribed.
Args: (participant_id: str)
- on_video_track_unsubscribed: Called when a video track is unsubscribed.
Args: (participant_id: str)
- on_app_message: Called when data is received from a participant. RTVI-compatible version of on_data_received.
Args: (message: Any, sender: str)
- on_data_received: Called when data is received. The participant ID is None
for packets sent by a server SDK, which LiveKit delivers unattributed.
Args: (data: bytes, participant_id: str | None)
- on_dtmf_event: Called when a SIP DTMF tone is received from a participant.
Args: (data: dict) with keys ``tone``/``digit``, ``code``, and
``participant_id``. Also pushes an ``InputDTMFFrame`` so
``DTMFAggregator`` works on LiveKit SIP calls.
Example::
@transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant_id):
await task.queue_frame(TTSSpeakFrame("Hello!"))
@transport.event_handler("on_participant_disconnected")
async def on_participant_disconnected(transport, participant_id):
await task.queue_frame(EndFrame())
"""
[docs]
def __init__(
self,
url: str,
token: str,
room_name: str,
params: LiveKitParams | None = None,
input_name: str | None = None,
output_name: str | None = None,
):
"""Initialize the LiveKit transport.
Args:
url: LiveKit server URL to connect to.
token: Authentication token for the room.
room_name: Name of the LiveKit room to join.
params: Configuration parameters for the transport.
input_name: Optional name for the input transport.
output_name: Optional name for the output transport.
"""
super().__init__(input_name=input_name, output_name=output_name)
callbacks = LiveKitCallbacks(
on_connected=self._on_connected,
on_disconnected=self._on_disconnected,
on_before_disconnect=self._on_before_disconnect,
on_participant_connected=self._on_participant_connected,
on_participant_disconnected=self._on_participant_disconnected,
on_audio_track_subscribed=self._on_audio_track_subscribed,
on_audio_track_unsubscribed=self._on_audio_track_unsubscribed,
on_video_track_subscribed=self._on_video_track_subscribed,
on_video_track_unsubscribed=self._on_video_track_unsubscribed,
on_data_received=self._on_data_received,
on_first_participant_joined=self._on_first_participant_joined,
on_dtmf_event=self._on_dtmf_event,
)
self._params = params or LiveKitParams()
self._client = LiveKitTransportClient(
url, token, room_name, self._params, callbacks, self.name
)
self._input: LiveKitInputTransport | None = None
self._output: LiveKitOutputTransport | None = None
self._register_event_handler("on_connected")
self._register_event_handler("on_disconnected")
self._register_event_handler("on_participant_connected")
self._register_event_handler("on_participant_disconnected")
self._register_event_handler("on_client_connected")
self._register_event_handler("on_client_disconnected")
self._register_event_handler("on_audio_track_subscribed")
self._register_event_handler("on_audio_track_unsubscribed")
self._register_event_handler("on_video_track_subscribed")
self._register_event_handler("on_video_track_unsubscribed")
self._register_event_handler("on_app_message")
self._register_event_handler("on_data_received")
self._register_event_handler("on_first_participant_joined")
self._register_event_handler("on_participant_left")
self._register_event_handler("on_call_state_updated")
self._register_event_handler("on_before_disconnect", sync=True)
self._register_event_handler("on_dtmf_event")
[docs]
def output(self) -> LiveKitOutputTransport:
"""Get the output transport for sending media and events.
Returns:
The LiveKit output transport instance.
"""
if not self._output:
self._output = LiveKitOutputTransport(
self, self._client, self._params, name=self._output_name
)
return self._output
@property
def participant_id(self) -> str:
"""Get the participant ID for this transport.
Returns:
The participant ID assigned by LiveKit.
"""
return self._client.participant_id
[docs]
async def send_audio(self, frame: OutputAudioRawFrame):
"""Send an audio frame to the LiveKit room.
Args:
frame: The audio frame to send.
"""
if self._output:
await self._output.queue_frame(frame, FrameDirection.DOWNSTREAM)
[docs]
def get_participants(self) -> list[str]:
"""Get list of participant IDs in the room.
Returns:
List of participant LiveKit identities.
"""
return self._client.get_participants()
[docs]
async def mute_participant(self, participant_id: str):
"""Stop receiving a specific participant's audio.
Args:
participant_id: LiveKit identity of the participant to stop
receiving audio from.
"""
await self._client.mute_participant(participant_id)
[docs]
async def unmute_participant(self, participant_id: str):
"""Resume receiving a specific participant's audio.
Args:
participant_id: LiveKit identity of the participant to resume
receiving audio from.
"""
await self._client.unmute_participant(participant_id)
async def _on_connected(self):
"""Handle room connected events."""
await self._call_event_handler("on_connected")
if self._input:
await self._input.push_frame(BotConnectedFrame())
async def _on_disconnected(self):
"""Handle room disconnected events."""
await self._call_event_handler("on_disconnected")
async def _on_before_disconnect(self):
"""Handle before disconnection room events."""
await self._call_event_handler("on_before_disconnect")
async def _on_participant_connected(self, participant_id: str):
"""Handle participant connected events."""
await self._call_event_handler("on_participant_connected", participant_id)
# Also call on_client_connected for compatibility with other transports.
# Wrapped as a dict (matching Daily's Mapping[str, Any] shape) so
# drop-in bot templates reading client["id"] work across transports.
await self._call_event_handler("on_client_connected", {"id": participant_id})
if self._input:
await self._input.push_frame(ClientConnectedFrame())
async def _on_participant_disconnected(self, participant_id: str):
"""Handle participant disconnected events."""
await self._call_event_handler("on_participant_disconnected", participant_id)
await self._call_event_handler("on_participant_left", participant_id, "disconnected")
# Also call on_client_disconnected for compatibility with other transports
await self._call_event_handler("on_client_disconnected", {"id": participant_id})
async def _on_audio_track_subscribed(self, participant_id: str):
"""Handle audio track subscribed events."""
await self._call_event_handler("on_audio_track_subscribed", participant_id)
async def _on_audio_track_unsubscribed(self, participant_id: str):
"""Handle audio track unsubscribed events."""
await self._call_event_handler("on_audio_track_unsubscribed", participant_id)
async def _on_video_track_subscribed(self, participant_id: str):
"""Handle video track subscribed events."""
await self._call_event_handler("on_video_track_subscribed", participant_id)
async def _on_video_track_unsubscribed(self, participant_id: str):
"""Handle video track unsubscribed events."""
await self._call_event_handler("on_video_track_unsubscribed", participant_id)
async def _on_data_received(self, data: bytes, participant_id: str | None):
"""Handle data received events."""
try:
message = json.loads(data.decode())
if not isinstance(message, dict):
logger.debug(f"{self} Ignoring non-object JSON data: {message!r}")
message = None
except (UnicodeDecodeError, json.JSONDecodeError) as e:
logger.debug(f"{self} Ignoring non-JSON data from {participant_id}: {e}")
message = None
if message is not None:
if self._input:
await self._input.push_app_message(message, participant_id)
# RTVI compatibility:
await self._call_event_handler("on_app_message", message, participant_id)
# Backwards compatibility with older transports that used on_data_received for app messages
await self._call_event_handler("on_data_received", data, participant_id)
async def _on_dtmf_event(self, data: Any):
"""Handle inbound SIP DTMF events.
Mirrors Daily transport behavior: expose ``on_dtmf_event`` to user code
and push ``InputDTMFFrame`` so ``DTMFAggregator`` can consume digits.
"""
logger.debug(f"{self} DTMF event: {data}")
await self._call_event_handler("on_dtmf_event", data)
tone = data.get("tone") if isinstance(data, dict) else None
if tone is None or self._input is None:
return
try:
button = KeypadEntry(tone)
except ValueError:
logger.warning(f"{self} Ignoring unsupported DTMF tone: {tone!r}")
return
await self._input.push_frame(InputDTMFFrame(button=button))
[docs]
async def send_message(self, message: str, participant_id: str | None = None):
"""Send a message to participants in the room.
Args:
message: The message string to send.
participant_id: Optional specific participant to send to.
"""
if self._output:
frame = LiveKitOutputTransportMessageFrame(
message=message, participant_id=participant_id
)
await self._output.send_message(frame)
[docs]
async def send_message_urgent(self, message: str, participant_id: str | None = None):
"""Send an urgent message to participants in the room.
Args:
message: The urgent message string to send.
participant_id: Optional specific participant to send to.
"""
if self._output:
frame = LiveKitOutputTransportMessageUrgentFrame(
message=message, participant_id=participant_id
)
await self._output.send_message(frame)
[docs]
async def on_room_event(self, event):
"""Handle room events.
Args:
event: The room event to handle.
"""
# Handle room events
pass
[docs]
async def on_participant_event(self, event):
"""Handle participant events.
Args:
event: The participant event to handle.
"""
# Handle participant events
pass
[docs]
async def on_track_event(self, event):
"""Handle track events.
Args:
event: The track event to handle.
"""
# Handle track events
pass
async def _on_call_state_updated(self, state: str):
"""Handle call state update events."""
await self._call_event_handler("on_call_state_updated", state)
async def _on_first_participant_joined(self, participant_id: str):
"""Handle first participant joined events."""
await self._call_event_handler("on_first_participant_joined", participant_id)