/
/
1"""
2Private PulseAudio capture server for grabbing PCM audio from client applications.
3
4Runs a minimal, MA-owned classic PulseAudio daemon (no autodetected devices, no
5network) whose module-pipe-sink FIFOs let MA capture the audio output of external
6client processes that can only play via PulseAudio/PipeWire (e.g. the official
7Spotify client). Consumers read a sink's FIFO with
8``helpers.named_pipe.read_named_pipe`` or by pointing ffmpeg at it directly.
9
10Also hosts :class:`PAVolumeController`, the shared ctypes libpulse controller for
11sink volume and module control, usable against any PA server (the private capture
12daemon or the system/host one).
13"""
14
15from __future__ import annotations
16
17import asyncio
18import ctypes
19import logging
20import math
21import os
22import shutil
23import threading
24import uuid
25import weakref
26from contextlib import suppress
27from pathlib import Path
28from typing import TYPE_CHECKING, ClassVar, Final
29
30from music_assistant.constants import MASS_LOGGER_NAME
31from music_assistant.helpers.process import AsyncProcess, check_output, get_subprocess_env
32
33if TYPE_CHECKING:
34 from music_assistant.mass import MusicAssistant
35
36LOGGER = logging.getLogger(f"{MASS_LOGGER_NAME}.helpers.pulse_capture")
37
38# Fixed capture format for all pipe sinks: 32-bit little-endian PCM, CD sample
39# rate, stereo. s32le keeps full precision of any 16/24-bit source even when the
40# sink volume applies (reciprocal) gain before MA reads the samples back.
41CAPTURE_SAMPLE_FORMAT: Final[str] = "s32le"
42CAPTURE_SAMPLE_RATE: Final[int] = 44100
43CAPTURE_CHANNELS: Final[int] = 2
44
45PA_VOLUME_NORM: Final = 65536 # 100% (0 dB) on PA's raw volume scale
46PA_CHANNELS_MAX: Final = 32
47
48# Ceiling for the raw (linear) volume setter. Reciprocal cubic compensation
49# needs up to 10000% (a client stream at 1% compensates with a 100x sink gain);
50# the resulting raw value (100 * PA_VOLUME_NORM) stays far below PA_VOLUME_MAX.
51MAX_RAW_VOLUME_PCT: Final = 10000.0
52
53# Daemon startup: how long to wait for the native socket to appear, and the
54# poll interval while waiting.
55_READY_TIMEOUT: Final = 10.0
56_READY_POLL_INTERVAL: Final = 0.1
57
58# Supervised-restart backoff: starts small so a one-off crash recovers fast,
59# doubles up to the max so a permanently failing daemon doesn't spin.
60_RESTART_BACKOFF_INITIAL: Final = 1.0
61_RESTART_BACKOFF_MAX: Final = 30.0
62
63# set_sink_volume() is called frequently (on every volume/mute change, and
64# at bridge start for every player). A short timeout limits how long a
65# stuck/unresponsive PA call can occupy an executor thread â under normal
66# conditions PA responds in single-digit milliseconds, so 0.5s is generous
67# while bounding the worst case. load_module()/unload_module() (rare,
68# one-time during topology setup/teardown) keep a longer 2.0s timeout since
69# we'd rather wait than have sink creation/cleanup spuriously fail.
70_SET_VOLUME_TIMEOUT: Final = 0.5
71
72PA_CONTEXT_READY: Final = 4
73PA_CONTEXT_FAILED: Final = 5
74PA_CONTEXT_TERMINATED: Final = 6
75PA_CONTEXT_NOAUTOSPAWN: Final = 1
76
77_CONTEXT_NOTIFY_CB = ctypes.CFUNCTYPE(None, ctypes.c_void_p, ctypes.c_void_p)
78_CONTEXT_SUCCESS_CB = ctypes.CFUNCTYPE(None, ctypes.c_void_p, ctypes.c_int, ctypes.c_void_p)
79_CONTEXT_INDEX_CB = ctypes.CFUNCTYPE(None, ctypes.c_void_p, ctypes.c_uint32, ctypes.c_void_p)
80
81PA_INVALID_INDEX: Final = 0xFFFFFFFF
82
83# --- Audio taper curve (dr-lex exponential, with linear roll-off) -----------
84#
85# y = a * e^(b*x) gives constant dB-per-slider-step ("audio taper" /
86# logarithmic potentiometer behavior), unlike a plain linear-amplitude
87# mapping (y = x) where the bottom of the slider is wildly more sensitive
88# than the top. See https://www.dr-lex.be/info-stuff/volumecontrols.html
89#
90# a = 10**(-range_dB/20) sets the amplitude floor; b = ln(1/a) ensures
91# y(1.0) = 1.0 (0dB) at full volume. Below _TAPER_ROLLOFF_X, a linear ramp
92# to (0, 0) is used so volume_pct=0 is true silence rather than asymptoting
93# toward the floor.
94#
95# Reference values for common dB ranges (pick one _TAPER_A and comment out
96# the rest; _TAPER_B recalculates automatically):
97#
98# Range _TAPER_A _TAPER_B MA 70% = Notes
99# 40 dB 0.01 ~4.605 -12 dB receiver / outdoor speakers (current)
100# 50 dB 0.003162 ~5.757 -15 dB medium-range setups
101# 60 dB 0.001 ~6.908 -18 dB consumer headphones / desktop speakers
102# 70 dB 0.000316 ~8.059 -21 dB high-dynamic-range hi-fi systems
103#
104# Used both for PA hardware volume (PAVolumeController.set_sink_volume, after
105# a cube-root step to counteract PA's own cubic volume curve) and for the
106# software PCM-scaling fallback path in local_audio, so the same slider
107# position sounds the same regardless of which volume-control mode is active.
108_TAPER_A: Final = 0.01 # 10**(-40/20) â 40dB range, suits receiver/outdoor setups
109# _TAPER_A: Final = 0.003162 # 10**(-50/20) â 50dB range
110# _TAPER_A: Final = 0.001 # 10**(-60/20) â 60dB range, suits headphones/desktop
111# _TAPER_A: Final = 0.000316 # 10**(-70/20) â 70dB range, hi-fi high dynamic range
112_TAPER_B: Final = math.log(1.0 / _TAPER_A) # recalculates automatically from _TAPER_A
113_TAPER_ROLLOFF_X: Final = 0.10 # below 10% slider, linear ramp to true silence
114
115
116def volume_pct_to_amplitude(volume_pct: int) -> float:
117 """
118 Map a 0-100 volume percentage to a linear amplitude scale factor.
119
120 Uses the dr-lex exponential audio taper (y = a*e^(b*x)) for
121 volume_pct >= 10, giving constant dB change per slider step. Below 10%,
122 a linear ramp to (0, 0) ensures volume_pct=0 produces true silence.
123 """
124 x = max(0, min(volume_pct, 100)) / 100.0
125 if x <= 0:
126 return 0.0
127 if x < _TAPER_ROLLOFF_X:
128 y1 = _TAPER_A * math.exp(_TAPER_B * _TAPER_ROLLOFF_X)
129 return y1 * (x / _TAPER_ROLLOFF_X)
130 return _TAPER_A * math.exp(_TAPER_B * x)
131
132
133def get_default_pulse_server() -> str:
134 """
135 Detect the system's default PulseAudio server address.
136
137 Checked fresh on each call â the socket may not exist at import time but
138 appear later once the audio host/addon has fully started.
139
140 :returns: Server address (env value or "unix:<socket path>"), or an empty
141 string when nothing was found (libpulse then uses its own defaults).
142 """
143 if server := os.environ.get("PULSE_SERVER"):
144 return server
145 for path in (
146 "/run/audio/pulse.sock",
147 "/run/pulse/native",
148 "/var/run/pulse/native",
149 ):
150 if Path(path).exists():
151 return f"unix:{path}"
152 return ""
153
154
155def get_pulse_capture_server(mass: MusicAssistant) -> PulseCaptureServer:
156 """
157 Return the shared PulseCaptureServer instance for this MusicAssistant.
158
159 All consumers must obtain the server through this function so they share a
160 single private daemon, and must pair acquire()/release() around their usage.
161 """
162 if (server := _servers.get(mass)) is None:
163 server = PulseCaptureServer(mass)
164 _servers[mass] = server
165 return server
166
167
168class PulseCaptureServer:
169 """
170 Private, MA-owned classic PulseAudio daemon for audio capture.
171
172 Lazily started on the first :meth:`acquire` and stopped again on the last
173 :meth:`release` (reference-counted) â use :func:`get_pulse_capture_server`
174 to obtain the shared per-process instance. The daemon loads only the native
175 protocol on a unix socket inside a private runtime dir under MA's cache dir;
176 audio clients are pointed at it via :meth:`child_env` and their audio is
177 captured through :class:`PipeSink` FIFOs.
178
179 If the daemon dies it is restarted automatically and :attr:`generation` is
180 bumped. All sinks created against the previous daemon are gone at that
181 point, so consumers must snapshot ``generation`` when creating a sink and
182 recreate their :class:`PipeSink` when it changes.
183 """
184
185 def __init__(self, mass: MusicAssistant) -> None:
186 """
187 Initialize the capture server (the daemon is not started yet).
188
189 :param mass: MusicAssistant instance, used for the private runtime dir
190 under its cache path.
191 """
192 # only the cache path is kept: holding mass itself would pin the
193 # WeakKeyDictionary registry entry (value -> key) forever
194 self._base_dir = Path(mass.cache_path) / "pulse_capture"
195 self._socket_path = self._base_dir / "native"
196 self._config_path = self._base_dir / "pulse_capture.pa"
197 self._lock = asyncio.Lock()
198 self._refcount = 0
199 self._generation = 0
200 self._proc: AsyncProcess | None = None
201 self._controller: PAVolumeController | None = None
202 self._supervisor_task: asyncio.Task[None] | None = None
203
204 @property
205 def generation(self) -> int:
206 """Daemon generation, bumped on every (re)start â snapshot when creating sinks."""
207 return self._generation
208
209 @property
210 def server_address(self) -> str:
211 """PA server address ("unix:<socket path>") of the private daemon."""
212 return f"unix:{self._socket_path}"
213
214 async def acquire(self) -> PulseCaptureServer:
215 """
216 Register a consumer, starting the private daemon on first use.
217
218 Every successful acquire() must be paired with exactly one release().
219
220 :returns: This server instance, for convenient chaining.
221 """
222 async with self._lock:
223 if self._refcount == 0:
224 await self._start()
225 self._refcount += 1
226 return self
227
228 async def release(self) -> None:
229 """
230 Unregister a consumer, stopping the daemon when none remain.
231
232 Safe to call more often than acquire() (extra calls are ignored).
233 """
234 async with self._lock:
235 if self._refcount == 0:
236 return
237 if self._refcount > 1:
238 self._refcount -= 1
239 return
240 # last consumer: run teardown to completion even when this call is
241 # cancelled, and only commit the zero refcount once it finished â
242 # no half-stopped daemon can be left behind or double-started
243 teardown = asyncio.ensure_future(self._stop())
244 try:
245 await asyncio.shield(teardown)
246 except asyncio.CancelledError:
247 while not teardown.done():
248 with suppress(asyncio.CancelledError):
249 await asyncio.shield(teardown)
250 self._refcount = 0
251 raise
252 self._refcount = 0
253
254 def child_env(self, sink_name: str) -> dict[str, str]:
255 """
256 Environment for an audio-client subprocess that must play into a sink.
257
258 :param sink_name: Sink the client's audio must go to (PipeSink.sink_name).
259 :returns: Full subprocess environment with PULSE_SERVER/PULSE_SINK set.
260 """
261 return get_subprocess_env(
262 {
263 "PULSE_SERVER": self.server_address,
264 "PULSE_SINK": sink_name,
265 }
266 )
267
268 async def _start(self) -> None:
269 """Start the daemon and its supervisor. Called with the lock held."""
270 await self._launch_daemon()
271 self._supervisor_task = asyncio.create_task(self._supervise())
272 self._supervisor_task.add_done_callback(_log_supervisor_exit)
273
274 async def _stop(self) -> None:
275 """Stop supervisor and daemon, clean the private dir. Called with the lock held."""
276 if (task := self._supervisor_task) is not None:
277 self._supervisor_task = None
278 task.cancel()
279 with suppress(asyncio.CancelledError):
280 await task
281 if (controller := self._controller) is not None:
282 self._controller = None
283 with suppress(Exception):
284 await asyncio.to_thread(controller.close)
285 if (proc := self._proc) is not None:
286 self._proc = None
287 await proc.close()
288
289 # the daemon's death drops all loaded modules with it, so nothing needs
290 # unloading; only the private runtime dir (socket + leftover FIFOs) remains
291 def _cleanup() -> None:
292 shutil.rmtree(self._base_dir, ignore_errors=True)
293
294 await asyncio.to_thread(_cleanup)
295
296 async def _launch_daemon(self) -> None:
297 """Start the daemon process and wait until it accepts connections."""
298 config_text = (
299 f"load-module module-native-protocol-unix socket={self._socket_path} auth-anonymous=1\n"
300 )
301
302 # PA's .pa config and module argument parsers are space-delimited with no
303 # escaping; refuse early with a clear error instead of mis-parsing later
304 if " " in str(self._base_dir):
305 raise RuntimeError(f"cache path may not contain spaces: {self._base_dir}")
306
307 def _prepare() -> None:
308 self._base_dir.mkdir(parents=True, exist_ok=True)
309 # the daemon accepts anonymous connections on its socket, so the
310 # private dir must never be accessible to other local users
311 self._base_dir.chmod(0o700)
312 self._socket_path.unlink(missing_ok=True)
313 self._config_path.write_text(config_text, encoding="utf-8")
314
315 await asyncio.to_thread(_prepare)
316 # runtime/state dirs are redirected into the private dir via the child
317 # environment only; os.environ is never mutated
318 proc = AsyncProcess(
319 [
320 "pulseaudio",
321 "--daemonize=no",
322 "-n",
323 "--exit-idle-time=-1",
324 f"--file={self._config_path}",
325 ],
326 stderr=True,
327 name="pulse-capture",
328 env={
329 "XDG_RUNTIME_DIR": str(self._base_dir),
330 "PULSE_RUNTIME_PATH": str(self._base_dir),
331 "PULSE_STATE_PATH": str(self._base_dir),
332 },
333 )
334 self._proc = proc
335 try:
336 await proc.start()
337 await self._wait_ready(proc)
338 # verify the daemon actually accepts connections before declaring
339 # ready; the construction cannot be interrupted mid-flight, so on
340 # cancellation close whatever the worker thread still produced to
341 # avoid leaking a live threaded mainloop
342 ctor = asyncio.ensure_future(asyncio.to_thread(PAVolumeController, self.server_address))
343 try:
344 controller = await asyncio.shield(ctor)
345 except asyncio.CancelledError:
346 ctor.add_done_callback(_close_controller_result)
347 raise
348 except BaseException:
349 if self._proc is proc:
350 self._proc = None
351 with suppress(Exception):
352 await proc.close()
353 raise
354 # publish the new controller and generation before disposing of the old
355 # controller, so teardown can always reach the live one even when the
356 # disposal await is cancelled
357 old_controller = self._controller
358 self._controller = controller
359 self._generation += 1
360 if old_controller is not None:
361 with suppress(Exception):
362 await asyncio.to_thread(old_controller.close)
363 LOGGER.debug(
364 "Private PulseAudio capture daemon ready on %s (generation %d)",
365 self.server_address,
366 self._generation,
367 )
368
369 async def _wait_ready(self, proc: AsyncProcess) -> None:
370 """Wait for the daemon's native socket to appear."""
371 try:
372 async with asyncio.timeout(_READY_TIMEOUT):
373 while not await asyncio.to_thread(self._socket_path.exists):
374 if proc.returncode is not None:
375 raise RuntimeError(
376 f"pulseaudio exited during startup (code {proc.returncode})"
377 )
378 await asyncio.sleep(_READY_POLL_INTERVAL)
379 except TimeoutError:
380 raise RuntimeError("Timeout waiting for the pulseaudio daemon socket") from None
381
382 async def _supervise(self) -> None:
383 """Restart the daemon (with bounded backoff) until the supervisor is cancelled."""
384 backoff = _RESTART_BACKOFF_INITIAL
385 while True:
386 if (proc := self._proc) is not None:
387 try:
388 async for line in proc.iter_stderr():
389 LOGGER.debug("pulseaudio: %s", line)
390 except Exception as err:
391 LOGGER.debug("pulseaudio log reader stopped: %s", err)
392 await proc.close()
393 if self._proc is proc:
394 self._proc = None
395 LOGGER.warning(
396 "Private PulseAudio capture daemon exited unexpectedly, restarting in %.1fs",
397 backoff,
398 )
399 await asyncio.sleep(backoff)
400 try:
401 await self._launch_daemon()
402 except Exception as err:
403 backoff = min(backoff * 2, _RESTART_BACKOFF_MAX)
404 LOGGER.error("Failed to restart the PulseAudio capture daemon: %s", err)
405 continue
406 backoff = _RESTART_BACKOFF_INITIAL
407
408 async def _load_module(self, module_name: str, argument: str) -> int | None:
409 """Load a PA module on the private daemon (blocking libpulse call in a thread)."""
410 controller = self._require_controller()
411 return await asyncio.to_thread(controller.load_module, module_name, argument)
412
413 async def _unload_module(self, module_index: int) -> bool:
414 """Unload a PA module from the private daemon."""
415 controller = self._require_controller()
416 return await asyncio.to_thread(controller.unload_module, module_index)
417
418 async def _set_sink_volume_raw(self, sink_name: str, volume_pct: float) -> bool:
419 """Set raw (linear) volume on a sink of the private daemon."""
420 controller = self._require_controller()
421 return await asyncio.to_thread(controller.set_sink_volume_raw, sink_name, volume_pct)
422
423 def _require_controller(self) -> PAVolumeController:
424 """Return the connected controller or raise if the server is not running."""
425 if (controller := self._controller) is None:
426 raise RuntimeError("Pulse capture server is not running")
427 return controller
428
429
430class PipeSink:
431 """
432 One isolated module-pipe-sink capture sink on a :class:`PulseCaptureServer`.
433
434 The PA daemon creates a FIFO at :attr:`fifo_path` that delivers the sink's
435 audio as raw PCM in the fixed capture format (CAPTURE_SAMPLE_FORMAT /
436 CAPTURE_SAMPLE_RATE / CAPTURE_CHANNELS); read it with
437 ``helpers.named_pipe.read_named_pipe`` or by pointing ffmpeg at it.
438
439 Both consumer lifecycles are supported: a single long-lived sink per
440 provider instance, or a fresh sink per stream (create -> suspend/resume ->
441 unload). A sink does not survive a daemon restart: when the server's
442 ``generation`` no longer matches the value snapshotted at creation, drop
443 this instance and create a new one.
444 """
445
446 def __init__(
447 self,
448 server: PulseCaptureServer,
449 sink_name: str,
450 fifo_path: Path,
451 module_index: int,
452 generation: int,
453 ) -> None:
454 """Initialize the sink. Use the async :meth:`create` factory instead."""
455 self._server = server
456 self._sink_name = sink_name
457 self._fifo_path = fifo_path
458 self._module_index: int | None = module_index
459 self._generation = generation
460
461 @classmethod
462 async def create(cls, server: PulseCaptureServer, name_prefix: str) -> PipeSink:
463 """
464 Create a new uniquely-named pipe sink on the given capture server.
465
466 The pipe-sink module creates the FIFO file itself.
467
468 :param server: An acquired PulseCaptureServer.
469 :param name_prefix: Prefix for the generated sink name (e.g. a provider
470 instance id); a short unique suffix is appended.
471 """
472 sink_name = f"{name_prefix}_{uuid.uuid4().hex[:8]}"
473 fifo_path = server._base_dir / f"{sink_name}.pcm"
474 argument = (
475 f"sink_name={sink_name} file={fifo_path} "
476 f"format={CAPTURE_SAMPLE_FORMAT} rate={CAPTURE_SAMPLE_RATE} "
477 f"channels={CAPTURE_CHANNELS}"
478 )
479 # snapshot before the load: a restart during the load would otherwise
480 # pair a module index from the dead daemon with the new generation,
481 # and a later unload could hit an unrelated module on the replacement
482 generation = server.generation
483 module_index = await server._load_module("module-pipe-sink", argument)
484 if module_index is None:
485 raise RuntimeError(f"Failed to load module-pipe-sink for {sink_name}")
486 if server.generation != generation:
487 raise RuntimeError(f"capture daemon restarted while creating sink {sink_name}")
488 return cls(server, sink_name, fifo_path, module_index, generation)
489
490 @property
491 def sink_name(self) -> str:
492 """The PA sink name (pass to PulseCaptureServer.child_env for the client)."""
493 return self._sink_name
494
495 @property
496 def fifo_path(self) -> Path:
497 """Path of the FIFO delivering this sink's PCM audio."""
498 return self._fifo_path
499
500 async def set_volume(self, volume_pct: float) -> None:
501 """
502 Set the sink's raw (linear) volume.
503
504 A sink from a previous daemon generation ignores the call (recreate the
505 sink after a restart).
506
507 :param volume_pct: 100 is unity gain; values above 100 amplify (e.g. 400
508 for reciprocal cubic compensation). Clamped to MAX_RAW_VOLUME_PCT.
509 """
510 if self._generation != self._server.generation:
511 LOGGER.debug("Ignoring volume for stale sink %s", self._sink_name)
512 return
513 if not await self._server._set_sink_volume_raw(self._sink_name, volume_pct):
514 raise RuntimeError(f"Failed to set volume on capture sink {self._sink_name}")
515
516 async def suspend(self) -> None:
517 """Suspend the sink (its FIFO stops producing audio until resumed)."""
518 await self._set_suspended(True)
519
520 async def resume(self) -> None:
521 """Resume a suspended sink."""
522 await self._set_suspended(False)
523
524 async def unload(self) -> None:
525 """
526 Unload the sink's module and remove its FIFO (idempotent).
527
528 A sink whose daemon has restarted since creation is already gone; only
529 the leftover FIFO file is cleaned up in that case.
530 """
531 module_index = self._module_index
532 self._module_index = None
533 if module_index is not None and self._generation == self._server.generation:
534 # best effort: a failure usually means the daemon/module is already
535 # gone, and a restart reclaims all modules anyway
536 if not await self._server._unload_module(module_index):
537 LOGGER.warning("Failed to unload capture sink module %s", self._sink_name)
538 with suppress(OSError):
539 await asyncio.to_thread(self._fifo_path.unlink)
540
541 async def _set_suspended(self, suspended: bool) -> None:
542 """Toggle the sink's suspend state via pactl."""
543 if self._generation != self._server.generation:
544 LOGGER.debug("Ignoring suspend toggle for stale sink %s", self._sink_name)
545 return
546 returncode, output = await check_output(
547 "pactl",
548 "--server",
549 self._server.server_address,
550 "suspend-sink",
551 self._sink_name,
552 "1" if suspended else "0",
553 env={"PULSE_SERVER": self._server.server_address},
554 timeout=5,
555 )
556 if returncode != 0:
557 LOGGER.warning(
558 "pactl suspend-sink %s %d failed: %s",
559 self._sink_name,
560 int(suspended),
561 output.decode("utf-8", errors="replace").strip(),
562 )
563
564
565class PAVolumeController:
566 """
567 Shared libpulse connection for PA sink volume and module control.
568
569 One instance is shared per PA server. All calls are blocking and must be
570 invoked via run_in_executor/to_thread from async code.
571 """
572
573 def __init__(self, server: str | None = None) -> None:
574 """
575 Connect to PulseAudio and start the threaded mainloop.
576
577 :param server: PA server address to connect to (e.g. "unix:<socket>").
578 Uses env/default socket discovery when omitted.
579 """
580 self._lib = _get_full_lib()
581 self._lock = threading.Lock()
582 self._mainloop = self._lib.pa_threaded_mainloop_new()
583 if not self._mainloop:
584 raise OSError("pa_threaded_mainloop_new returned NULL")
585
586 api = self._lib.pa_threaded_mainloop_get_api(self._mainloop)
587 self._context = self._lib.pa_context_new(api, b"music-assistant-volume")
588 if not self._context:
589 self._lib.pa_threaded_mainloop_free(self._mainloop)
590 self._mainloop = None
591 raise OSError("pa_context_new returned NULL")
592
593 self._ready = threading.Event()
594 self._failed = threading.Event()
595
596 def _state_cb_impl(_ctx: int, _userdata: int) -> None:
597 state = self._lib.pa_context_get_state(self._context)
598 if state == PA_CONTEXT_READY:
599 self._ready.set()
600 elif state in (PA_CONTEXT_FAILED, PA_CONTEXT_TERMINATED):
601 self._failed.set()
602
603 self._state_cb = _CONTEXT_NOTIFY_CB(_state_cb_impl) # keep reference alive â GC
604 self._lib.pa_context_set_state_callback(self._context, self._state_cb, None)
605
606 pulse_server = server or get_default_pulse_server()
607 ret = self._lib.pa_context_connect(
608 self._context,
609 pulse_server.encode() if pulse_server else None,
610 PA_CONTEXT_NOAUTOSPAWN,
611 None,
612 )
613 if ret < 0:
614 self.close()
615 raise OSError(f"pa_context_connect failed (ret={ret})")
616
617 self._lib.pa_threaded_mainloop_start(self._mainloop)
618
619 if not self._ready.wait(timeout=5.0):
620 self.close()
621 raise OSError("Timed out connecting to PulseAudio for volume control")
622
623 def set_sink_volume(self, sink_name: str, volume_pct: int, channels: int = 2) -> bool:
624 """
625 Set hardware volume on a named PA sink.
626
627 :param sink_name: PA sink name as returned by ``enumerate_pa_sinks()``.
628 :param volume_pct: Volume level 0-100, mapped through an exponential
629 audio taper curve before being sent to PA.
630 :param channels: Channel count for the PA volume structure. Should
631 match the sink's actual channel count.
632 :returns: True if PA reported success.
633 """
634 amplitude = volume_pct_to_amplitude(volume_pct)
635 # cube root counteracts PA's own cubic volume curve
636 pa_vol = round(PA_VOLUME_NORM * amplitude ** (1.0 / 3.0))
637 return self._apply_sink_volume(sink_name, pa_vol, channels)
638
639 def set_sink_volume_raw(self, sink_name: str, volume_pct: float, channels: int = 2) -> bool:
640 """
641 Set raw (linear) volume on a named PA sink, without the audio taper.
642
643 :param sink_name: PA sink name.
644 :param volume_pct: Linear percentage where 100 maps exactly onto
645 PA_VOLUME_NORM (0 dB). Values above 100 amplify (e.g. 400 for
646 reciprocal cubic compensation); clamped to MAX_RAW_VOLUME_PCT.
647 :param channels: Channel count for the PA volume structure.
648 :returns: True if PA reported success.
649 """
650 volume_pct = min(max(volume_pct, 0.0), MAX_RAW_VOLUME_PCT)
651 pa_vol = round(PA_VOLUME_NORM * volume_pct / 100.0)
652 return self._apply_sink_volume(sink_name, pa_vol, channels)
653
654 def load_module(self, module_name: str, argument: str) -> int | None:
655 """
656 Load a PulseAudio module (e.g. module-remap-sink) via libpulse.
657
658 :param module_name: PA module name, e.g. "module-remap-sink".
659 :param argument: Module argument string, e.g.
660 "sink_name=Foo master=bar channels=2 master_channel_map=...
661 channel_map=front-left,front-right remix=no".
662
663 Blocks (up to ~2s) for PA's response.
664 :returns: The loaded module's index, or None on failure/timeout.
665 """
666 with self._lock:
667 if self._failed.is_set() or not self._mainloop or not self._context:
668 return None
669
670 done = threading.Event()
671 result: dict[str, int] = {}
672
673 def _index_cb_impl(_ctx: int, idx: int, _userdata: int) -> None:
674 result["index"] = idx
675 done.set()
676
677 index_cb = _CONTEXT_INDEX_CB(_index_cb_impl)
678
679 self._lib.pa_threaded_mainloop_lock(self._mainloop)
680 try:
681 op = self._lib.pa_context_load_module(
682 self._context,
683 module_name.encode(),
684 argument.encode(),
685 index_cb,
686 None,
687 )
688 if not op:
689 return None
690 finally:
691 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
692
693 if not done.wait(timeout=2.0):
694 self._cancel_operation(op)
695 return None
696 self._lib.pa_operation_unref(op)
697 idx = result.get("index", PA_INVALID_INDEX)
698 return None if idx == PA_INVALID_INDEX else idx
699
700 def unload_module(self, module_index: int) -> bool:
701 """
702 Unload a previously-loaded PulseAudio module by index.
703
704 Blocks (up to ~2s) for PA's response.
705 :returns: True if PA reported success.
706 """
707 with self._lock:
708 if self._failed.is_set() or not self._mainloop or not self._context:
709 return False
710
711 done = threading.Event()
712 result: dict[str, int] = {}
713
714 def _success_cb_impl(_ctx: int, success: int, _userdata: int) -> None:
715 result["success"] = success
716 done.set()
717
718 success_cb = _CONTEXT_SUCCESS_CB(_success_cb_impl)
719
720 self._lib.pa_threaded_mainloop_lock(self._mainloop)
721 try:
722 op = self._lib.pa_context_unload_module(
723 self._context, module_index, success_cb, None
724 )
725 if not op:
726 return False
727 finally:
728 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
729
730 if not done.wait(timeout=2.0):
731 self._cancel_operation(op)
732 return False
733 self._lib.pa_operation_unref(op)
734 return bool(result.get("success", 0))
735
736 def close(self) -> None:
737 """Disconnect and tear down the mainloop."""
738 with self._lock:
739 if self._context:
740 self._lib.pa_context_disconnect(self._context)
741 self._lib.pa_context_unref(self._context)
742 self._context = None
743 if self._mainloop:
744 self._lib.pa_threaded_mainloop_stop(self._mainloop)
745 self._lib.pa_threaded_mainloop_free(self._mainloop)
746 self._mainloop = None
747
748 def _cancel_operation(self, op: int) -> None:
749 """
750 Detach a timed-out operation's callback before dropping the reference.
751
752 PA may still deliver the response later from the mainloop thread; without
753 the cancel it would invoke the (garbage-collected) ctypes trampoline of a
754 callback that went out of scope â undefined behavior.
755 """
756 self._lib.pa_operation_cancel(op)
757 self._lib.pa_operation_unref(op)
758
759 def _apply_sink_volume(self, sink_name: str, pa_volume: int, channels: int) -> bool:
760 """
761 Send an already-mapped raw PA volume to a named sink.
762
763 :returns: True if PA reported success.
764 """
765 with self._lock:
766 if self._failed.is_set() or not self._mainloop or not self._context:
767 return False
768 cvol = _PACVolume()
769 self._lib.pa_cvolume_set(ctypes.byref(cvol), channels, pa_volume)
770
771 done = threading.Event()
772 result: dict[str, int] = {}
773
774 def _success_cb_impl(_ctx: int, success: int, _userdata: int) -> None:
775 result["success"] = success
776 done.set()
777
778 success_cb = _CONTEXT_SUCCESS_CB(_success_cb_impl)
779
780 self._lib.pa_threaded_mainloop_lock(self._mainloop)
781 try:
782 op = self._lib.pa_context_set_sink_volume_by_name(
783 self._context,
784 sink_name.encode(),
785 ctypes.byref(cvol),
786 success_cb,
787 None,
788 )
789 if not op:
790 return False
791 finally:
792 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
793
794 if not done.wait(timeout=_SET_VOLUME_TIMEOUT):
795 self._cancel_operation(op)
796 return False
797 self._lib.pa_operation_unref(op)
798 return bool(result.get("success", 0))
799
800
801# One shared server per MusicAssistant instance; weak keys so a discarded mass
802# (tests) does not pin its server object forever.
803_servers: weakref.WeakKeyDictionary[MusicAssistant, PulseCaptureServer] = (
804 weakref.WeakKeyDictionary()
805)
806
807
808class _PACVolume(ctypes.Structure):
809 _fields_: ClassVar = [
810 ("channels", ctypes.c_uint8),
811 ("values", ctypes.c_uint32 * PA_CHANNELS_MAX),
812 ]
813
814
815def _log_supervisor_exit(task: asyncio.Task[None]) -> None:
816 """Surface a supervisor that died on an unexpected error (restarts stop with it)."""
817 if not task.cancelled() and task.exception() is not None:
818 LOGGER.error("PulseAudio capture supervisor died unexpectedly: %s", task.exception())
819
820
821def _close_controller_result(fut: asyncio.Future[PAVolumeController]) -> None:
822 """Close a controller whose awaiting task was cancelled mid-construction."""
823 with suppress(Exception):
824 fut.result().close()
825
826
827def _load_full_lib() -> ctypes.CDLL:
828 """
829 Load and configure libpulse for use by PAVolumeController.
830
831 Called once per process; the result is cached by _get_full_lib().
832 """
833 lib = ctypes.CDLL("libpulse.so.0")
834
835 lib.pa_threaded_mainloop_new.restype = ctypes.c_void_p
836 lib.pa_threaded_mainloop_get_api.restype = ctypes.c_void_p
837 lib.pa_threaded_mainloop_get_api.argtypes = [ctypes.c_void_p]
838 lib.pa_threaded_mainloop_start.restype = ctypes.c_int
839 lib.pa_threaded_mainloop_start.argtypes = [ctypes.c_void_p]
840 lib.pa_threaded_mainloop_stop.argtypes = [ctypes.c_void_p]
841 lib.pa_threaded_mainloop_free.argtypes = [ctypes.c_void_p]
842 lib.pa_threaded_mainloop_lock.argtypes = [ctypes.c_void_p]
843 lib.pa_threaded_mainloop_unlock.argtypes = [ctypes.c_void_p]
844
845 lib.pa_context_new.restype = ctypes.c_void_p
846 lib.pa_context_new.argtypes = [ctypes.c_void_p, ctypes.c_char_p]
847 lib.pa_context_set_state_callback.argtypes = [
848 ctypes.c_void_p,
849 _CONTEXT_NOTIFY_CB,
850 ctypes.c_void_p,
851 ]
852 lib.pa_context_connect.restype = ctypes.c_int
853 lib.pa_context_connect.argtypes = [
854 ctypes.c_void_p,
855 ctypes.c_char_p,
856 ctypes.c_int,
857 ctypes.c_void_p,
858 ]
859 lib.pa_context_get_state.restype = ctypes.c_int
860 lib.pa_context_get_state.argtypes = [ctypes.c_void_p]
861 lib.pa_context_disconnect.argtypes = [ctypes.c_void_p]
862 lib.pa_context_unref.argtypes = [ctypes.c_void_p]
863
864 lib.pa_cvolume_set.restype = ctypes.c_void_p
865 lib.pa_cvolume_set.argtypes = [ctypes.c_void_p, ctypes.c_uint, ctypes.c_uint32]
866
867 lib.pa_context_set_sink_volume_by_name.restype = ctypes.c_void_p
868 lib.pa_context_set_sink_volume_by_name.argtypes = [
869 ctypes.c_void_p,
870 ctypes.c_char_p,
871 ctypes.c_void_p,
872 _CONTEXT_SUCCESS_CB,
873 ctypes.c_void_p,
874 ]
875 lib.pa_operation_unref.argtypes = [ctypes.c_void_p]
876 lib.pa_operation_cancel.argtypes = [ctypes.c_void_p]
877
878 lib.pa_context_load_module.restype = ctypes.c_void_p
879 lib.pa_context_load_module.argtypes = [
880 ctypes.c_void_p,
881 ctypes.c_char_p,
882 ctypes.c_char_p,
883 _CONTEXT_INDEX_CB,
884 ctypes.c_void_p,
885 ]
886 lib.pa_context_unload_module.restype = ctypes.c_void_p
887 lib.pa_context_unload_module.argtypes = [
888 ctypes.c_void_p,
889 ctypes.c_uint32,
890 _CONTEXT_SUCCESS_CB,
891 ctypes.c_void_p,
892 ]
893 return lib
894
895
896_full_lib: ctypes.CDLL | None = None
897
898
899def _get_full_lib() -> ctypes.CDLL:
900 global _full_lib # noqa: PLW0603
901 if _full_lib is None:
902 _full_lib = _load_full_lib()
903 return _full_lib
904