Files
VoiceAgent/audio_device_monitor.py

358 lines
16 KiB
Python

"""macOS default audio-device monitoring.
The monitor deliberately depends on a small backend protocol. Core Audio invokes
listeners on an arbitrary thread, while VoiceAgent consumes snapshots on its
asyncio loop; :class:`AudioDeviceMonitor` is the lifecycle and thread boundary
between those two worlds. The backend can be replaced by a deterministic fake
in unit tests.
"""
from __future__ import annotations
import asyncio
import ctypes
import ctypes.util
import inspect
import sys
from dataclasses import dataclass, field
from enum import Enum
from typing import Callable, Protocol
class DeviceChangeReason(str, Enum):
INITIAL = "initial"
DEFAULT_CHANGED = "default_changed"
DEVICE_ADDED = "device_added"
DEVICE_REMOVED = "device_removed"
DEVICE_RECONFIGURED = "device_reconfigured"
PROFILE_CHANGED = "profile_changed"
@dataclass(frozen=True)
class AudioDevice:
uid: str
name: str
can_input: bool
can_output: bool
transport: str = "unknown"
alive: bool = True
@dataclass(frozen=True)
class AudioDeviceSnapshot:
generation: int
default_input_uid: str | None
default_output_uid: str | None
devices: dict[str, AudioDevice] = field(default_factory=dict)
reason: DeviceChangeReason = DeviceChangeReason.INITIAL
Listener = Callable[[], None]
class AudioDeviceBackend(Protocol):
"""The minimal Core Audio surface required by the monitor."""
def enumerate_devices(self) -> list[AudioDevice]: ...
def default_input_uid(self) -> str | None: ...
def default_output_uid(self) -> str | None: ...
def add_listener(self, callback: Listener) -> object: ...
def remove_listener(self, token: object) -> None: ...
class AudioDeviceMonitor:
"""Coalesces Core Audio notifications and dispatches snapshots on one loop."""
def __init__(self, backend: AudioDeviceBackend, *, debounce_seconds: float = 0.05,
loop: asyncio.AbstractEventLoop | None = None):
self._backend = backend
self._debounce_seconds = debounce_seconds
self._loop = loop
self._callback: Callable[[AudioDeviceSnapshot], object] | None = None
self._snapshot = AudioDeviceSnapshot(0, None, None)
self._generation = 0
self._listener_tokens: list[object] = []
self._debounce_handle: asyncio.TimerHandle | None = None
self._running = False
self._refresh_scheduled = False
self._refresh_task: asyncio.Task[None] | None = None
async def start(self, on_change: Callable[[AudioDeviceSnapshot], object]) -> None:
"""Register listeners once and publish an initial snapshot.
``on_change`` may be synchronous or return an awaitable. It is always
called on the loop used by this monitor, never on a Core Audio thread.
"""
if self._running:
return
self._loop = self._loop or asyncio.get_running_loop()
self._callback = on_change
self._running = True
try:
# Register all listeners before the initial read so a concurrent
# device change cannot be missed.
self._listener_tokens = [self._backend.add_listener(self._on_backend_event) for _ in range(1)]
await self._refresh(DeviceChangeReason.INITIAL)
except Exception:
await self.stop()
raise
async def stop(self) -> None:
"""Stop idempotently and ensure no callback can be scheduled afterward."""
self._running = False
if self._debounce_handle:
self._debounce_handle.cancel()
self._debounce_handle = None
current_task = asyncio.current_task()
if (self._refresh_task and not self._refresh_task.done()
and self._refresh_task is not current_task):
self._refresh_task.cancel()
try:
await self._refresh_task
except asyncio.CancelledError:
pass
self._refresh_task = None
tokens, self._listener_tokens = self._listener_tokens, []
for token in tokens:
self._backend.remove_listener(token)
self._callback = None
self._refresh_scheduled = False
def snapshot(self) -> AudioDeviceSnapshot:
return self._snapshot
def _on_backend_event(self) -> None:
"""Core Audio callback entry point; safe to call from any thread."""
if not self._running or not self._loop:
return
self._loop.call_soon_threadsafe(self._schedule_refresh)
def _schedule_refresh(self) -> None:
if not self._running or self._refresh_scheduled:
return
self._refresh_scheduled = True
if self._debounce_handle:
self._debounce_handle.cancel()
self._debounce_handle = self._loop.call_later(self._debounce_seconds, self._start_refresh)
def _start_refresh(self) -> None:
self._debounce_handle = None
self._refresh_scheduled = False
if self._running:
self._refresh_task = self._loop.create_task(self._refresh(self._infer_reason()))
self._refresh_task.add_done_callback(self._consume_refresh_failure)
@staticmethod
def _consume_refresh_failure(task: asyncio.Task[None]) -> None:
"""Retrieve scheduled refresh failures after cleanup has completed."""
if not task.cancelled():
task.exception()
def _infer_reason(self) -> DeviceChangeReason:
# The backend intentionally keeps the callback payload-free. The
# resulting snapshot is authoritative; callers can inspect UID/device
# differences. A generic device reconfiguration is safest here.
return DeviceChangeReason.DEFAULT_CHANGED
async def _refresh(self, reason: DeviceChangeReason) -> None:
if not self._running:
return
devices = {d.uid: d for d in self._backend.enumerate_devices() if d.uid and d.alive}
new_input = self._backend.default_input_uid()
new_output = self._backend.default_output_uid()
# An unavailable default is represented as None rather than a stale UID.
if new_input not in devices or not devices[new_input].can_input:
new_input = None
if new_output not in devices or not devices[new_output].can_output:
new_output = None
old = self._snapshot
if (old.default_input_uid == new_input and old.default_output_uid == new_output
and old.devices == devices and old.generation != 0):
return
if old.generation and reason != DeviceChangeReason.INITIAL:
old_uids, new_uids = set(old.devices), set(devices)
if old.default_input_uid != new_input or old.default_output_uid != new_output:
reason = DeviceChangeReason.DEFAULT_CHANGED
elif new_uids - old_uids:
reason = DeviceChangeReason.DEVICE_ADDED
elif old_uids - new_uids:
reason = DeviceChangeReason.DEVICE_REMOVED
elif any(old.devices[uid] != device for uid, device in devices.items()
if uid in old.devices):
reason = DeviceChangeReason.PROFILE_CHANGED
self._generation += 1
self._snapshot = AudioDeviceSnapshot(self._generation, new_input, new_output, devices, reason)
callback = self._callback
if callback and self._running:
try:
result = callback(self._snapshot)
if inspect.isawaitable(result):
await result
except Exception:
# A failed consumer must not leave a native listener active with
# an unusable callback. ``stop`` handles the current refresh
# task specially so this cleanup is safe from inside _refresh.
await self.stop()
raise
class MacOSCoreAudioBackend:
"""Core Audio backend hook.
PyObjC's CoreAudio listener ABI differs between macOS releases. Keeping the
native adapter behind this class lets packaging provide the matching adapter
without exposing it to the async monitor or its tests.
"""
def __init__(self, adapter):
self._adapter = adapter
def enumerate_devices(self) -> list[AudioDevice]:
return list(self._adapter.enumerate_devices())
def default_input_uid(self) -> str | None:
return self._adapter.default_input_uid()
def default_output_uid(self) -> str | None:
return self._adapter.default_output_uid()
def add_listener(self, callback: Listener) -> object:
return self._adapter.add_device_listener(callback)
def remove_listener(self, token: object) -> None:
self._adapter.remove_device_listener(token)
class NativeMacOSCoreAudioAdapter:
"""Native Core Audio adapter using stable device UIDs, never PortAudio IDs."""
_SYSTEM_OBJECT = 1
_GLOBAL = int.from_bytes(b"glob", "big")
_INPUT_SCOPE = int.from_bytes(b"inpt", "big")
_OUTPUT_SCOPE = int.from_bytes(b"outp", "big")
_DEFAULT_INPUT = int.from_bytes(b"dIn ", "big")
_DEFAULT_OUTPUT = int.from_bytes(b"dOut", "big")
_DEVICES = int.from_bytes(b"dev#", "big")
_UID = int.from_bytes(b"uid ", "big")
_NAME = int.from_bytes(b"lnam", "big")
_ALIVE = int.from_bytes(b"livn", "big")
_TRANSPORT = int.from_bytes(b"tran", "big")
_STREAMS = int.from_bytes(b"stm#", "big")
class _Address(ctypes.Structure):
_fields_ = [("selector", ctypes.c_uint32), ("scope", ctypes.c_uint32), ("element", ctypes.c_uint32)]
def __init__(self):
if sys.platform != "darwin":
raise RuntimeError("Core Audio is only available on macOS")
core_audio = ctypes.util.find_library("CoreAudio")
core_foundation = ctypes.util.find_library("CoreFoundation")
if not core_audio or not core_foundation:
raise RuntimeError("CoreAudio.framework is unavailable")
self._lib = ctypes.CDLL(core_audio)
self._cf = ctypes.CDLL(core_foundation)
self._listener_type = ctypes.CFUNCTYPE(ctypes.c_int32, ctypes.c_uint32, ctypes.c_uint32,
ctypes.POINTER(self._Address), ctypes.c_void_p)
self._lib.AudioObjectGetPropertyData.argtypes = [ctypes.c_uint32, ctypes.POINTER(self._Address), ctypes.c_uint32, ctypes.c_void_p, ctypes.POINTER(ctypes.c_uint32), ctypes.c_void_p]
self._lib.AudioObjectGetPropertyData.restype = ctypes.c_int32
self._lib.AudioObjectGetPropertyDataSize.argtypes = [ctypes.c_uint32, ctypes.POINTER(self._Address), ctypes.c_uint32, ctypes.c_void_p, ctypes.POINTER(ctypes.c_uint32)]
self._lib.AudioObjectGetPropertyDataSize.restype = ctypes.c_int32
self._lib.AudioObjectAddPropertyListener.argtypes = [ctypes.c_uint32, ctypes.POINTER(self._Address), self._listener_type, ctypes.c_void_p]
self._lib.AudioObjectRemovePropertyListener.argtypes = [ctypes.c_uint32, ctypes.POINTER(self._Address), self._listener_type, ctypes.c_void_p]
self._cf.CFStringGetCString.argtypes = [ctypes.c_void_p, ctypes.c_char_p, ctypes.c_long, ctypes.c_uint32]
self._callbacks = {}
def _address(self, selector, scope=None):
return self._Address(selector, self._GLOBAL if scope is None else scope, 0)
def _value(self, object_id, selector, value_type, scope=None):
value, size = value_type(), ctypes.c_uint32(ctypes.sizeof(value_type))
status = self._lib.AudioObjectGetPropertyData(object_id, ctypes.byref(self._address(selector, scope)), 0, None, ctypes.byref(size), ctypes.byref(value))
if status:
raise OSError(f"AudioObjectGetPropertyData({selector}) failed: {status}")
return value.value
def _string(self, object_id, selector):
ref = self._value(object_id, selector, ctypes.c_void_p)
if not ref:
return ""
buffer = ctypes.create_string_buffer(1024)
if not self._cf.CFStringGetCString(ref, buffer, len(buffer), 0x08000100):
return ""
return buffer.value.decode("utf-8", "replace")
def _device_ids(self):
address, size = self._address(self._DEVICES), ctypes.c_uint32()
status = self._lib.AudioObjectGetPropertyDataSize(self._SYSTEM_OBJECT, ctypes.byref(address), 0, None, ctypes.byref(size))
if status:
raise OSError(f"Audio device enumeration failed: {status}")
devices = (ctypes.c_uint32 * (size.value // ctypes.sizeof(ctypes.c_uint32)))()
status = self._lib.AudioObjectGetPropertyData(self._SYSTEM_OBJECT, ctypes.byref(address), 0, None, ctypes.byref(size), devices)
if status:
raise OSError(f"Audio device enumeration failed: {status}")
return list(devices)
def _has_streams(self, device_id, scope):
address, size = self._address(self._STREAMS, scope), ctypes.c_uint32()
status = self._lib.AudioObjectGetPropertyDataSize(device_id, ctypes.byref(address), 0, None, ctypes.byref(size))
return not status and bool(size.value)
def enumerate_devices(self):
devices = []
for device_id in self._device_ids():
try:
uid = self._string(device_id, self._UID)
if not uid:
continue
devices.append(AudioDevice(uid, self._string(device_id, self._NAME) or uid,
self._has_streams(device_id, self._INPUT_SCOPE), self._has_streams(device_id, self._OUTPUT_SCOPE),
self._fourcc(self._value(device_id, self._TRANSPORT, ctypes.c_uint32)), bool(self._value(device_id, self._ALIVE, ctypes.c_uint32))))
except OSError:
continue
return devices
@staticmethod
def _fourcc(value):
return value.to_bytes(4, "big").decode("ascii", "replace").strip() or "unknown"
def default_input_uid(self):
return self._uid_for_id(self._value(self._SYSTEM_OBJECT, self._DEFAULT_INPUT, ctypes.c_uint32))
def default_output_uid(self):
return self._uid_for_id(self._value(self._SYSTEM_OBJECT, self._DEFAULT_OUTPUT, ctypes.c_uint32))
def _uid_for_id(self, device_id):
return self._string(device_id, self._UID) if device_id else None
def add_device_listener(self, callback):
addresses = [(self._SYSTEM_OBJECT, self._address(selector)) for selector in (self._DEFAULT_INPUT, self._DEFAULT_OUTPUT, self._DEVICES)]
addresses += [(device_id, self._address(selector, scope)) for device_id in self._device_ids() for selector, scope in ((self._ALIVE, None), (self._STREAMS, self._INPUT_SCOPE), (self._STREAMS, self._OUTPUT_SCOPE))]
native_callback = self._listener_type(lambda *_: (callback(), 0)[1])
registered = []
try:
for object_id, address in addresses:
status = self._lib.AudioObjectAddPropertyListener(object_id, ctypes.byref(address), native_callback, None)
if status:
raise OSError(f"AudioObjectAddPropertyListener failed: {status}")
registered.append((object_id, address))
except Exception:
for object_id, address in registered:
self._lib.AudioObjectRemovePropertyListener(object_id, ctypes.byref(address), native_callback, None)
raise
token = (native_callback, registered)
self._callbacks[id(token)] = token
return token
def remove_device_listener(self, token):
native_callback, addresses = token
for object_id, address in addresses:
self._lib.AudioObjectRemovePropertyListener(object_id, ctypes.byref(address), native_callback, None)
self._callbacks.pop(id(token), None)
def create_macos_audio_monitor(*, adapter=None, **kwargs) -> AudioDeviceMonitor:
"""Create a native Core Audio monitor (or a supplied test adapter)."""
if adapter is None:
adapter = NativeMacOSCoreAudioAdapter()
return AudioDeviceMonitor(MacOSCoreAudioBackend(adapter), **kwargs)