435 lines
19 KiB
Python
435 lines
19 KiB
Python
"""Pipecat local-audio transport backed by sounddevice instead of PyAudio.
|
|
|
|
Pipecat ships `pipecat.transports.local.audio.LocalAudioTransport`, but it needs
|
|
PyAudio, which has no macOS wheel and must be compiled against a Homebrew
|
|
portaudio. This machine blocks both, so we talk to the same portaudio through
|
|
sounddevice, whose wheel bundles a prebuilt dylib. The frame contract is
|
|
identical to the upstream transport.
|
|
"""
|
|
|
|
import asyncio
|
|
import inspect
|
|
import sys
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
|
|
import sounddevice as sd
|
|
from loguru import logger
|
|
|
|
from pipecat.frames.frames import InputAudioRawFrame, OutputAudioRawFrame, StartFrame
|
|
from pipecat.processors.frame_processor import FrameProcessor
|
|
from pipecat.transports.base_input import BaseInputTransport
|
|
from pipecat.transports.base_output import BaseOutputTransport
|
|
from pipecat.transports.base_transport import BaseTransport, TransportParams
|
|
from audio_device_monitor import AudioDeviceSnapshot, create_macos_audio_monitor
|
|
|
|
|
|
class SoundDeviceTransportParams(TransportParams):
|
|
"""Configuration for the sounddevice transport.
|
|
|
|
Parameters:
|
|
input_device: sounddevice device index or name substring. None uses the default.
|
|
output_device: sounddevice device index or name substring. None uses the default.
|
|
"""
|
|
|
|
input_device: int | str | None = None
|
|
output_device: int | str | None = None
|
|
|
|
|
|
class SoundDeviceInputTransport(BaseInputTransport):
|
|
"""Captures microphone audio and pushes it into the pipeline."""
|
|
|
|
_params: SoundDeviceTransportParams
|
|
_transport: "SoundDeviceTransport"
|
|
|
|
def __init__(self, params: SoundDeviceTransportParams):
|
|
super().__init__(params)
|
|
self._in_stream: sd.RawInputStream | None = None
|
|
self._sample_rate = 0
|
|
self._stream_generation = 0
|
|
self._stream_lock = asyncio.Lock()
|
|
|
|
async def start(self, frame: StartFrame):
|
|
await super().start(frame)
|
|
|
|
if self._in_stream:
|
|
return
|
|
|
|
self._sample_rate = self._params.audio_in_sample_rate or frame.audio_in_sample_rate
|
|
blocksize = int(self._sample_rate / 100) * 2 # 20ms
|
|
|
|
await self._open_stream()
|
|
|
|
await self.set_transport_ready(frame)
|
|
|
|
if hasattr(self, "_transport"):
|
|
await self._transport.start_device_monitor(self)
|
|
|
|
async def reopen(self, *, device=None):
|
|
"""Recreate an unset-device stream so PortAudio resolves the new default."""
|
|
if not self._in_stream:
|
|
return
|
|
async with self._stream_lock:
|
|
old_stream = self._in_stream
|
|
old_device = old_stream.device
|
|
# Fail before disrupting an otherwise healthy conversation whenever
|
|
# PortAudio can already tell us that the new default is unavailable.
|
|
sd.check_input_settings(
|
|
device=self._params.input_device if device is None else device, samplerate=self._sample_rate,
|
|
channels=self._params.audio_in_channels, dtype="int16",
|
|
)
|
|
self._stream_generation += 1 # makes callbacks from the old stream inert
|
|
try:
|
|
old_stream.stop()
|
|
old_stream.close()
|
|
self._in_stream = None
|
|
await self._open_stream(device=device)
|
|
except Exception as exc:
|
|
logger.warning(f"Audio input route change failed; restoring prior stream: {type(exc).__name__}")
|
|
try:
|
|
self._in_stream = None
|
|
await self._open_stream(device=old_device)
|
|
except Exception as restore_exc:
|
|
logger.error(f"Audio input fallback unavailable: {type(restore_exc).__name__}")
|
|
raise
|
|
|
|
async def _open_stream(self, *, device=None):
|
|
blocksize = int(self._sample_rate / 100) * 2
|
|
self._stream_generation += 1
|
|
generation = self._stream_generation
|
|
stream = sd.RawInputStream(
|
|
samplerate=self._sample_rate, blocksize=blocksize,
|
|
device=self._params.input_device if device is None else device,
|
|
channels=self._params.audio_in_channels, dtype="int16",
|
|
callback=lambda *args: self._audio_in_callback(generation, *args),
|
|
)
|
|
try:
|
|
stream.start()
|
|
except Exception:
|
|
stream.close()
|
|
raise
|
|
self._in_stream = stream
|
|
device_name = sd.query_devices(self._in_stream.device, "input")["name"]
|
|
logger.info(f"Microphone: {device_name} @ {self._sample_rate} Hz")
|
|
|
|
async def cleanup(self):
|
|
await super().cleanup()
|
|
async with self._stream_lock:
|
|
self._stream_generation += 1
|
|
if self._in_stream:
|
|
self._in_stream.stop()
|
|
self._in_stream.close()
|
|
self._in_stream = None
|
|
if hasattr(self, "_transport"):
|
|
await self._transport.stop_device_monitor(self)
|
|
|
|
def _audio_in_callback(self, generation, indata, frame_count, time_info, status):
|
|
if generation != self._stream_generation:
|
|
return
|
|
if status:
|
|
logger.trace(f"Audio input status: {status}")
|
|
|
|
frame = InputAudioRawFrame(
|
|
audio=bytes(indata),
|
|
sample_rate=self._sample_rate,
|
|
num_channels=self._params.audio_in_channels,
|
|
)
|
|
|
|
asyncio.run_coroutine_threadsafe(
|
|
self._push_audio_frame_if_current(generation, frame), self.get_event_loop()
|
|
)
|
|
|
|
async def _push_audio_frame_if_current(self, generation, frame):
|
|
"""Serialize frame delivery with replacement so a closed route cannot leak audio."""
|
|
async with self._stream_lock:
|
|
if generation != self._stream_generation:
|
|
return
|
|
await self.push_audio_frame(frame)
|
|
|
|
|
|
class SoundDeviceOutputTransport(BaseOutputTransport):
|
|
"""Plays pipeline audio out through the speakers."""
|
|
|
|
_params: SoundDeviceTransportParams
|
|
|
|
def __init__(self, params: SoundDeviceTransportParams):
|
|
super().__init__(params)
|
|
self._out_stream: sd.RawOutputStream | None = None
|
|
self._sample_rate = 0
|
|
# Writes are serialized by the pipeline, so one worker is enough.
|
|
self._executor = ThreadPoolExecutor(max_workers=1)
|
|
self._stream_lock = asyncio.Lock()
|
|
|
|
async def start(self, frame: StartFrame):
|
|
await super().start(frame)
|
|
|
|
if self._out_stream:
|
|
return
|
|
|
|
self._sample_rate = self._params.audio_out_sample_rate or frame.audio_out_sample_rate
|
|
await self._open_stream()
|
|
|
|
await self.set_transport_ready(frame)
|
|
if hasattr(self, "_transport"):
|
|
await self._transport.start_device_monitor(self)
|
|
|
|
async def reopen(self, *, device=None):
|
|
"""Recreate an unset-device stream so PortAudio resolves the new default."""
|
|
if not self._out_stream:
|
|
return
|
|
async with self._stream_lock:
|
|
old_stream = self._out_stream
|
|
old_device = old_stream.device
|
|
# Keep the current output route intact when the selected default
|
|
# cannot satisfy this stream's negotiated format.
|
|
sd.check_output_settings(
|
|
device=self._params.output_device if device is None else device, samplerate=self._sample_rate,
|
|
channels=self._params.audio_out_channels, dtype="int16",
|
|
)
|
|
try:
|
|
old_stream.stop()
|
|
old_stream.close()
|
|
self._out_stream = None
|
|
await self._open_stream(device=device)
|
|
except Exception as exc:
|
|
logger.warning(f"Audio output route change failed; restoring prior stream: {type(exc).__name__}")
|
|
try:
|
|
self._out_stream = None
|
|
await self._open_stream(device=old_device)
|
|
except Exception as restore_exc:
|
|
logger.error(f"Audio output fallback unavailable: {type(restore_exc).__name__}")
|
|
raise
|
|
|
|
async def _open_stream(self, *, device=None):
|
|
stream = sd.RawOutputStream(
|
|
samplerate=self._sample_rate,
|
|
device=self._params.output_device if device is None else device,
|
|
channels=self._params.audio_out_channels, dtype="int16",
|
|
)
|
|
try:
|
|
stream.start()
|
|
except Exception:
|
|
stream.close()
|
|
raise
|
|
self._out_stream = stream
|
|
device_name = sd.query_devices(self._out_stream.device, "output")["name"]
|
|
logger.info(f"Speaker: {device_name} @ {self._sample_rate} Hz")
|
|
|
|
async def cleanup(self):
|
|
await super().cleanup()
|
|
async with self._stream_lock:
|
|
if self._out_stream:
|
|
self._out_stream.stop()
|
|
self._out_stream.close()
|
|
self._out_stream = None
|
|
if hasattr(self, "_transport"):
|
|
await self._transport.stop_device_monitor(self)
|
|
|
|
async def write_audio_frame(self, frame: OutputAudioRawFrame) -> bool:
|
|
async with self._stream_lock:
|
|
if not self._out_stream:
|
|
return False
|
|
await self.get_event_loop().run_in_executor(
|
|
self._executor, self._out_stream.write, frame.audio
|
|
)
|
|
return True
|
|
|
|
|
|
class SoundDeviceTransport(BaseTransport):
|
|
"""Local microphone + speaker transport."""
|
|
|
|
def __init__(self, params: SoundDeviceTransportParams, *, device_monitor=None,
|
|
device_event_sink=None):
|
|
super().__init__()
|
|
self._params = params
|
|
self._input: SoundDeviceInputTransport | None = None
|
|
self._output: SoundDeviceOutputTransport | None = None
|
|
self._device_monitor = device_monitor
|
|
self._device_event_sink = device_event_sink
|
|
self._last_snapshot: AudioDeviceSnapshot | None = None
|
|
self._restart_lock = asyncio.Lock()
|
|
self._monitor_started = False
|
|
self._monitor_owners: set[object] = set()
|
|
self._monitor_lock = asyncio.Lock()
|
|
|
|
def _monitor_if_needed(self):
|
|
if self._device_monitor is not None:
|
|
return self._device_monitor
|
|
if sys.platform != "darwin":
|
|
return None
|
|
if self._params.input_device is None or self._params.output_device is None:
|
|
try:
|
|
self._device_monitor = create_macos_audio_monitor()
|
|
except Exception as exc:
|
|
logger.warning(f"Audio default monitoring unavailable: {exc}")
|
|
return self._device_monitor
|
|
|
|
async def start_device_monitor(self, owner=None):
|
|
"""Keep the shared monitor running while any transport side is active."""
|
|
owner = self if owner is None else owner
|
|
async with self._monitor_lock:
|
|
self._monitor_owners.add(owner)
|
|
if self._monitor_started:
|
|
return
|
|
monitor = self._monitor_if_needed()
|
|
if monitor:
|
|
await monitor.start(self._on_device_change)
|
|
self._monitor_started = True
|
|
|
|
async def stop_device_monitor(self, owner=None):
|
|
"""Release one transport side; stop only after the final release."""
|
|
owner = self if owner is None else owner
|
|
async with self._monitor_lock:
|
|
self._monitor_owners.discard(owner)
|
|
if self._monitor_owners or not self._device_monitor or not self._monitor_started:
|
|
return
|
|
await self._device_monitor.stop()
|
|
self._monitor_started = False
|
|
|
|
async def _on_device_change(self, snapshot: AudioDeviceSnapshot):
|
|
if self._last_snapshot and snapshot.generation <= self._last_snapshot.generation:
|
|
logger.debug(f"Ignoring stale audio route event generation={snapshot.generation}")
|
|
return
|
|
if self._device_event_sink:
|
|
try:
|
|
result = self._device_event_sink(snapshot)
|
|
if inspect.isawaitable(result):
|
|
await result
|
|
except Exception as exc:
|
|
logger.warning(f"Audio device event sink failed: {type(exc).__name__}")
|
|
old = self._last_snapshot
|
|
self._last_snapshot = snapshot
|
|
if old is None:
|
|
return
|
|
input_changed = old.default_input_uid != snapshot.default_input_uid
|
|
output_changed = old.default_output_uid != snapshot.default_output_uid
|
|
if not (input_changed or output_changed):
|
|
return
|
|
async with self._restart_lock:
|
|
if input_changed and self._params.input_device is None and self._input:
|
|
try:
|
|
await self._reopen_default(self._input, snapshot, snapshot.default_input_uid, "input")
|
|
except Exception as exc:
|
|
logger.warning(f"Audio input route refresh failed: {type(exc).__name__}")
|
|
if output_changed and self._params.output_device is None and self._output:
|
|
try:
|
|
await self._reopen_default(self._output, snapshot, snapshot.default_output_uid, "output")
|
|
except Exception as exc:
|
|
logger.warning(f"Audio output route refresh failed: {type(exc).__name__}")
|
|
|
|
async def _reopen_default(self, stream, snapshot: AudioDeviceSnapshot, uid: str | None, direction: str) -> None:
|
|
# Empty snapshots are supported for legacy/injected monitors. Native
|
|
# snapshots always carry devices and therefore get an explicit index.
|
|
if not snapshot.devices:
|
|
await stream.reopen()
|
|
return
|
|
await stream.reopen(device=self._portaudio_device(snapshot, uid, direction))
|
|
|
|
@staticmethod
|
|
def _portaudio_device(snapshot: AudioDeviceSnapshot, uid: str | None, direction: str) -> int:
|
|
"""Map Core Audio's current default to an explicit PortAudio index.
|
|
|
|
``device=None`` in a long-lived sounddevice process retains PortAudio's
|
|
startup default. Reopening with the current index is what makes a
|
|
default-device event actually move the live stream.
|
|
"""
|
|
if not uid or uid not in snapshot.devices:
|
|
raise RuntimeError(f"No available Core Audio default {direction} device")
|
|
native = snapshot.devices[uid]
|
|
capability = "max_input_channels" if direction == "input" else "max_output_channels"
|
|
matches = [index for index, candidate in enumerate(sd.query_devices())
|
|
if candidate["name"] == native.name and candidate[capability] > 0]
|
|
if len(matches) != 1:
|
|
raise RuntimeError(f"No unique PortAudio {direction} device for {native.name!r}: {matches}")
|
|
return matches[0]
|
|
|
|
@staticmethod
|
|
def available_devices() -> list[dict]:
|
|
return [
|
|
{"id": index, "name": device["name"], "input": bool(device["max_input_channels"]),
|
|
"output": bool(device["max_output_channels"])}
|
|
for index, device in enumerate(sd.query_devices())
|
|
if device["max_input_channels"] or device["max_output_channels"]
|
|
]
|
|
|
|
@classmethod
|
|
def _select_device(cls, request: int | str, direction: str) -> int:
|
|
capability = "input" if direction == "input" else "output"
|
|
devices = cls.available_devices()
|
|
if isinstance(request, int) or (isinstance(request, str) and request.isdecimal()):
|
|
index = int(request)
|
|
if any(device["id"] == index and device[capability] for device in devices):
|
|
return index
|
|
else:
|
|
needle = str(request).casefold().strip()
|
|
matches = [device["id"] for device in devices
|
|
if device[capability] and needle in device["name"].casefold()]
|
|
if len(matches) == 1:
|
|
return matches[0]
|
|
raise ValueError(f"No unique available {direction} device matches {request!r}")
|
|
|
|
async def set_runtime_device(self, direction: str, request: int | str | None) -> dict:
|
|
"""Pin one route live, or pass ``default``/None to follow macOS again."""
|
|
if direction not in {"input", "output"}:
|
|
raise ValueError("direction must be input or output")
|
|
following_default = request is None or str(request).casefold().strip() in {"default", "mac default", "system default"}
|
|
selected = None if following_default else self._select_device(request, direction) # type: ignore[arg-type]
|
|
previous = self._params.input_device if direction == "input" else self._params.output_device
|
|
async with self._restart_lock:
|
|
stream = self._input if direction == "input" else self._output
|
|
try:
|
|
# Temporarily clear this pin so an all-pinned transport can
|
|
# create its native monitor and take a fresh default snapshot.
|
|
if following_default:
|
|
if direction == "input":
|
|
self._params.input_device = None
|
|
else:
|
|
self._params.output_device = None
|
|
await self.start_device_monitor()
|
|
if stream:
|
|
if following_default:
|
|
snapshot = self._last_snapshot
|
|
if not snapshot:
|
|
raise RuntimeError(f"No macOS default-{direction} snapshot is available yet")
|
|
uid = snapshot.default_input_uid if direction == "input" else snapshot.default_output_uid
|
|
await self._reopen_default(stream, snapshot, uid, direction)
|
|
else:
|
|
await stream.reopen(device=selected)
|
|
if direction == "input":
|
|
self._params.input_device = selected
|
|
else:
|
|
self._params.output_device = selected
|
|
except Exception:
|
|
if direction == "input":
|
|
self._params.input_device = previous
|
|
else:
|
|
self._params.output_device = previous
|
|
raise
|
|
name = "macOS default" if following_default else next(device["name"] for device in self.available_devices() if device["id"] == selected)
|
|
logger.info(f"Runtime {direction} device changed to {name}")
|
|
return {"direction": direction, "device": selected, "name": name, "following_default": following_default}
|
|
|
|
def input(self) -> FrameProcessor:
|
|
if not self._input:
|
|
self._input = SoundDeviceInputTransport(self._params)
|
|
self._input._transport = self
|
|
return self._input
|
|
|
|
def output(self) -> FrameProcessor:
|
|
if not self._output:
|
|
self._output = SoundDeviceOutputTransport(self._params)
|
|
self._output._transport = self
|
|
return self._output
|
|
|
|
|
|
def list_devices() -> str:
|
|
"""Render the audio devices sounddevice can see, for `--list-devices`."""
|
|
lines = []
|
|
for index, device in enumerate(sd.query_devices()):
|
|
capability = []
|
|
if device["max_input_channels"]:
|
|
capability.append(f"in:{device['max_input_channels']}")
|
|
if device["max_output_channels"]:
|
|
capability.append(f"out:{device['max_output_channels']}")
|
|
lines.append(f" [{index}] {device['name']} ({', '.join(capability)})")
|
|
return "\n".join(lines)
|