/
/
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 # a stale pid file from a hard shutdown makes pulse refuse to start
314 # ("daemon already running") when that pid is reused by any process
315 (self._base_dir / "pid").unlink(missing_ok=True)
316 # a hard shutdown skips sink cleanup; sweep leftover FIFOs from
317 # previous runs (a fresh daemon has no modules yet)
318 for stale_fifo in self._base_dir.glob("*.pcm"):
319 stale_fifo.unlink(missing_ok=True)
320 self._config_path.write_text(config_text, encoding="utf-8")
321
322 await asyncio.to_thread(_prepare)
323 # runtime/state dirs are redirected into the private dir via the child
324 # environment only; os.environ is never mutated
325 proc = AsyncProcess(
326 [
327 "pulseaudio",
328 "--daemonize=no",
329 "-n",
330 "--exit-idle-time=-1",
331 f"--file={self._config_path}",
332 ],
333 stderr=True,
334 name="pulse-capture",
335 env={
336 "XDG_RUNTIME_DIR": str(self._base_dir),
337 "PULSE_RUNTIME_PATH": str(self._base_dir),
338 "PULSE_STATE_PATH": str(self._base_dir),
339 },
340 )
341 self._proc = proc
342 try:
343 await proc.start()
344 await self._wait_ready(proc)
345 # verify the daemon actually accepts connections before declaring
346 # ready; the construction cannot be interrupted mid-flight, so on
347 # cancellation close whatever the worker thread still produced to
348 # avoid leaking a live threaded mainloop
349 ctor = asyncio.ensure_future(asyncio.to_thread(PAVolumeController, self.server_address))
350 try:
351 controller = await asyncio.shield(ctor)
352 except asyncio.CancelledError:
353 ctor.add_done_callback(_close_controller_result)
354 raise
355 except BaseException:
356 if self._proc is proc:
357 self._proc = None
358 with suppress(Exception):
359 await proc.close()
360 raise
361 # publish the new controller and generation before disposing of the old
362 # controller, so teardown can always reach the live one even when the
363 # disposal await is cancelled
364 old_controller = self._controller
365 self._controller = controller
366 self._generation += 1
367 if old_controller is not None:
368 with suppress(Exception):
369 await asyncio.to_thread(old_controller.close)
370 LOGGER.debug(
371 "Private PulseAudio capture daemon ready on %s (generation %d)",
372 self.server_address,
373 self._generation,
374 )
375
376 async def _wait_ready(self, proc: AsyncProcess) -> None:
377 """Wait for the daemon's native socket to appear."""
378 try:
379 async with asyncio.timeout(_READY_TIMEOUT):
380 while not await asyncio.to_thread(self._socket_path.exists):
381 if proc.returncode is not None:
382 raise RuntimeError(
383 f"pulseaudio exited during startup (code {proc.returncode})"
384 )
385 await asyncio.sleep(_READY_POLL_INTERVAL)
386 except TimeoutError:
387 raise RuntimeError("Timeout waiting for the pulseaudio daemon socket") from None
388
389 async def _supervise(self) -> None:
390 """Restart the daemon (with bounded backoff) until the supervisor is cancelled."""
391 backoff = _RESTART_BACKOFF_INITIAL
392 while True:
393 if (proc := self._proc) is not None:
394 try:
395 async for line in proc.iter_stderr():
396 LOGGER.debug("pulseaudio: %s", line)
397 except Exception as err:
398 LOGGER.debug("pulseaudio log reader stopped: %s", err)
399 await proc.close()
400 if self._proc is proc:
401 self._proc = None
402 LOGGER.warning(
403 "Private PulseAudio capture daemon exited unexpectedly, restarting in %.1fs",
404 backoff,
405 )
406 await asyncio.sleep(backoff)
407 try:
408 await self._launch_daemon()
409 except Exception as err:
410 backoff = min(backoff * 2, _RESTART_BACKOFF_MAX)
411 LOGGER.error("Failed to restart the PulseAudio capture daemon: %s", err)
412 continue
413 backoff = _RESTART_BACKOFF_INITIAL
414
415 async def _load_module(self, module_name: str, argument: str) -> int | None:
416 """Load a PA module on the private daemon (blocking libpulse call in a thread)."""
417 controller = self._require_controller()
418 return await asyncio.to_thread(controller.load_module, module_name, argument)
419
420 async def _unload_module(self, module_index: int) -> bool:
421 """Unload a PA module from the private daemon."""
422 controller = self._require_controller()
423 return await asyncio.to_thread(controller.unload_module, module_index)
424
425 async def _set_sink_volume_raw(self, sink_name: str, volume_pct: float) -> bool:
426 """Set raw (linear) volume on a sink of the private daemon."""
427 controller = self._require_controller()
428 return await asyncio.to_thread(controller.set_sink_volume_raw, sink_name, volume_pct)
429
430 def _require_controller(self) -> PAVolumeController:
431 """Return the connected controller or raise if the server is not running."""
432 if (controller := self._controller) is None:
433 raise RuntimeError("Pulse capture server is not running")
434 return controller
435
436
437class PipeSink:
438 """
439 One isolated module-pipe-sink capture sink on a :class:`PulseCaptureServer`.
440
441 The PA daemon creates a FIFO at :attr:`fifo_path` that delivers the sink's
442 audio as raw PCM in the fixed capture format (CAPTURE_SAMPLE_FORMAT /
443 CAPTURE_SAMPLE_RATE / CAPTURE_CHANNELS); read it with
444 ``helpers.named_pipe.read_named_pipe`` or by pointing ffmpeg at it.
445
446 Both consumer lifecycles are supported: a single long-lived sink per
447 provider instance, or a fresh sink per stream (create -> suspend/resume ->
448 unload). A sink does not survive a daemon restart: when the server's
449 ``generation`` no longer matches the value snapshotted at creation, drop
450 this instance and create a new one.
451 """
452
453 def __init__(
454 self,
455 server: PulseCaptureServer,
456 sink_name: str,
457 fifo_path: Path,
458 module_index: int,
459 generation: int,
460 ) -> None:
461 """Initialize the sink. Use the async :meth:`create` factory instead."""
462 self._server = server
463 self._sink_name = sink_name
464 self._fifo_path = fifo_path
465 self._module_index: int | None = module_index
466 self._generation = generation
467
468 @classmethod
469 async def create(cls, server: PulseCaptureServer, name_prefix: str) -> PipeSink:
470 """
471 Create a new uniquely-named pipe sink on the given capture server.
472
473 The pipe-sink module creates the FIFO file itself.
474
475 :param server: An acquired PulseCaptureServer.
476 :param name_prefix: Prefix for the generated sink name (e.g. a provider
477 instance id); a short unique suffix is appended.
478 """
479 sink_name = f"{name_prefix}_{uuid.uuid4().hex[:8]}"
480 fifo_path = server._base_dir / f"{sink_name}.pcm"
481 argument = (
482 f"sink_name={sink_name} file={fifo_path} "
483 f"format={CAPTURE_SAMPLE_FORMAT} rate={CAPTURE_SAMPLE_RATE} "
484 f"channels={CAPTURE_CHANNELS}"
485 )
486 # snapshot before the load: a restart during the load would otherwise
487 # pair a module index from the dead daemon with the new generation,
488 # and a later unload could hit an unrelated module on the replacement
489 generation = server.generation
490 module_index = await server._load_module("module-pipe-sink", argument)
491 if module_index is None:
492 raise RuntimeError(f"Failed to load module-pipe-sink for {sink_name}")
493 if server.generation != generation:
494 raise RuntimeError(f"capture daemon restarted while creating sink {sink_name}")
495 return cls(server, sink_name, fifo_path, module_index, generation)
496
497 @property
498 def sink_name(self) -> str:
499 """The PA sink name (pass to PulseCaptureServer.child_env for the client)."""
500 return self._sink_name
501
502 @property
503 def fifo_path(self) -> Path:
504 """Path of the FIFO delivering this sink's PCM audio."""
505 return self._fifo_path
506
507 async def set_volume(self, volume_pct: float) -> None:
508 """
509 Set the sink's raw (linear) volume.
510
511 A sink from a previous daemon generation ignores the call (recreate the
512 sink after a restart).
513
514 :param volume_pct: 100 is unity gain; values above 100 amplify (e.g. 400
515 for reciprocal cubic compensation). Clamped to MAX_RAW_VOLUME_PCT.
516 """
517 if self._generation != self._server.generation:
518 LOGGER.debug("Ignoring volume for stale sink %s", self._sink_name)
519 return
520 if not await self._server._set_sink_volume_raw(self._sink_name, volume_pct):
521 raise RuntimeError(f"Failed to set volume on capture sink {self._sink_name}")
522
523 async def suspend(self) -> None:
524 """Suspend the sink (its FIFO stops producing audio until resumed)."""
525 await self._set_suspended(True)
526
527 async def resume(self) -> None:
528 """Resume a suspended sink."""
529 await self._set_suspended(False)
530
531 async def unload(self) -> None:
532 """
533 Unload the sink's module and remove its FIFO (idempotent).
534
535 A sink whose daemon has restarted since creation is already gone; only
536 the leftover FIFO file is cleaned up in that case.
537 """
538 module_index = self._module_index
539 self._module_index = None
540 if module_index is not None and self._generation == self._server.generation:
541 # best effort: a failure usually means the daemon/module is already
542 # gone, and a restart reclaims all modules anyway
543 if not await self._server._unload_module(module_index):
544 LOGGER.warning("Failed to unload capture sink module %s", self._sink_name)
545 with suppress(OSError):
546 await asyncio.to_thread(self._fifo_path.unlink)
547
548 async def _set_suspended(self, suspended: bool) -> None:
549 """Toggle the sink's suspend state via pactl."""
550 if self._generation != self._server.generation:
551 LOGGER.debug("Ignoring suspend toggle for stale sink %s", self._sink_name)
552 return
553 returncode, output = await check_output(
554 "pactl",
555 "--server",
556 self._server.server_address,
557 "suspend-sink",
558 self._sink_name,
559 "1" if suspended else "0",
560 env={"PULSE_SERVER": self._server.server_address},
561 timeout=5,
562 )
563 if returncode != 0:
564 LOGGER.warning(
565 "pactl suspend-sink %s %d failed: %s",
566 self._sink_name,
567 int(suspended),
568 output.decode("utf-8", errors="replace").strip(),
569 )
570
571
572class PAVolumeController:
573 """
574 Shared libpulse connection for PA sink volume and module control.
575
576 One instance is shared per PA server. All calls are blocking and must be
577 invoked via run_in_executor/to_thread from async code.
578 """
579
580 def __init__(self, server: str | None = None) -> None:
581 """
582 Connect to PulseAudio and start the threaded mainloop.
583
584 :param server: PA server address to connect to (e.g. "unix:<socket>").
585 Uses env/default socket discovery when omitted.
586 """
587 self._lib = _get_full_lib()
588 self._lock = threading.Lock()
589 self._mainloop = self._lib.pa_threaded_mainloop_new()
590 if not self._mainloop:
591 raise OSError("pa_threaded_mainloop_new returned NULL")
592
593 api = self._lib.pa_threaded_mainloop_get_api(self._mainloop)
594 self._context = self._lib.pa_context_new(api, b"music-assistant-volume")
595 if not self._context:
596 self._lib.pa_threaded_mainloop_free(self._mainloop)
597 self._mainloop = None
598 raise OSError("pa_context_new returned NULL")
599
600 self._ready = threading.Event()
601 self._failed = threading.Event()
602
603 def _state_cb_impl(_ctx: int, _userdata: int) -> None:
604 state = self._lib.pa_context_get_state(self._context)
605 if state == PA_CONTEXT_READY:
606 self._ready.set()
607 elif state in (PA_CONTEXT_FAILED, PA_CONTEXT_TERMINATED):
608 self._failed.set()
609
610 self._state_cb = _CONTEXT_NOTIFY_CB(_state_cb_impl) # keep reference alive â GC
611 self._lib.pa_context_set_state_callback(self._context, self._state_cb, None)
612
613 pulse_server = server or get_default_pulse_server()
614 ret = self._lib.pa_context_connect(
615 self._context,
616 pulse_server.encode() if pulse_server else None,
617 PA_CONTEXT_NOAUTOSPAWN,
618 None,
619 )
620 if ret < 0:
621 self.close()
622 raise OSError(f"pa_context_connect failed (ret={ret})")
623
624 self._lib.pa_threaded_mainloop_start(self._mainloop)
625
626 if not self._ready.wait(timeout=5.0):
627 self.close()
628 raise OSError("Timed out connecting to PulseAudio for volume control")
629
630 def set_sink_volume(self, sink_name: str, volume_pct: int, channels: int = 2) -> bool:
631 """
632 Set hardware volume on a named PA sink.
633
634 :param sink_name: PA sink name as returned by ``enumerate_pa_sinks()``.
635 :param volume_pct: Volume level 0-100, mapped through an exponential
636 audio taper curve before being sent to PA.
637 :param channels: Channel count for the PA volume structure. Should
638 match the sink's actual channel count.
639 :returns: True if PA reported success.
640 """
641 amplitude = volume_pct_to_amplitude(volume_pct)
642 # cube root counteracts PA's own cubic volume curve
643 pa_vol = round(PA_VOLUME_NORM * amplitude ** (1.0 / 3.0))
644 return self._apply_sink_volume(sink_name, pa_vol, channels)
645
646 def set_sink_volume_raw(self, sink_name: str, volume_pct: float, channels: int = 2) -> bool:
647 """
648 Set raw (linear) volume on a named PA sink, without the audio taper.
649
650 :param sink_name: PA sink name.
651 :param volume_pct: Linear percentage where 100 maps exactly onto
652 PA_VOLUME_NORM (0 dB). Values above 100 amplify (e.g. 400 for
653 reciprocal cubic compensation); clamped to MAX_RAW_VOLUME_PCT.
654 :param channels: Channel count for the PA volume structure.
655 :returns: True if PA reported success.
656 """
657 volume_pct = min(max(volume_pct, 0.0), MAX_RAW_VOLUME_PCT)
658 pa_vol = round(PA_VOLUME_NORM * volume_pct / 100.0)
659 return self._apply_sink_volume(sink_name, pa_vol, channels)
660
661 def load_module(self, module_name: str, argument: str) -> int | None:
662 """
663 Load a PulseAudio module (e.g. module-remap-sink) via libpulse.
664
665 :param module_name: PA module name, e.g. "module-remap-sink".
666 :param argument: Module argument string, e.g.
667 "sink_name=Foo master=bar channels=2 master_channel_map=...
668 channel_map=front-left,front-right remix=no".
669
670 Blocks (up to ~2s) for PA's response.
671 :returns: The loaded module's index, or None on failure/timeout.
672 """
673 with self._lock:
674 if self._failed.is_set() or not self._mainloop or not self._context:
675 return None
676
677 done = threading.Event()
678 result: dict[str, int] = {}
679
680 def _index_cb_impl(_ctx: int, idx: int, _userdata: int) -> None:
681 result["index"] = idx
682 done.set()
683
684 index_cb = _CONTEXT_INDEX_CB(_index_cb_impl)
685
686 self._lib.pa_threaded_mainloop_lock(self._mainloop)
687 try:
688 op = self._lib.pa_context_load_module(
689 self._context,
690 module_name.encode(),
691 argument.encode(),
692 index_cb,
693 None,
694 )
695 if not op:
696 return None
697 finally:
698 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
699
700 if not done.wait(timeout=2.0):
701 self._cancel_operation(op)
702 return None
703 self._lib.pa_operation_unref(op)
704 idx = result.get("index", PA_INVALID_INDEX)
705 return None if idx == PA_INVALID_INDEX else idx
706
707 def unload_module(self, module_index: int) -> bool:
708 """
709 Unload a previously-loaded PulseAudio module by index.
710
711 Blocks (up to ~2s) for PA's response.
712 :returns: True if PA reported success.
713 """
714 with self._lock:
715 if self._failed.is_set() or not self._mainloop or not self._context:
716 return False
717
718 done = threading.Event()
719 result: dict[str, int] = {}
720
721 def _success_cb_impl(_ctx: int, success: int, _userdata: int) -> None:
722 result["success"] = success
723 done.set()
724
725 success_cb = _CONTEXT_SUCCESS_CB(_success_cb_impl)
726
727 self._lib.pa_threaded_mainloop_lock(self._mainloop)
728 try:
729 op = self._lib.pa_context_unload_module(
730 self._context, module_index, success_cb, None
731 )
732 if not op:
733 return False
734 finally:
735 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
736
737 if not done.wait(timeout=2.0):
738 self._cancel_operation(op)
739 return False
740 self._lib.pa_operation_unref(op)
741 return bool(result.get("success", 0))
742
743 def close(self) -> None:
744 """Disconnect and tear down the mainloop."""
745 with self._lock:
746 if self._context:
747 self._lib.pa_context_disconnect(self._context)
748 self._lib.pa_context_unref(self._context)
749 self._context = None
750 if self._mainloop:
751 self._lib.pa_threaded_mainloop_stop(self._mainloop)
752 self._lib.pa_threaded_mainloop_free(self._mainloop)
753 self._mainloop = None
754
755 def _cancel_operation(self, op: int) -> None:
756 """
757 Detach a timed-out operation's callback before dropping the reference.
758
759 PA may still deliver the response later from the mainloop thread; without
760 the cancel it would invoke the (garbage-collected) ctypes trampoline of a
761 callback that went out of scope â undefined behavior.
762 """
763 self._lib.pa_operation_cancel(op)
764 self._lib.pa_operation_unref(op)
765
766 def _apply_sink_volume(self, sink_name: str, pa_volume: int, channels: int) -> bool:
767 """
768 Send an already-mapped raw PA volume to a named sink.
769
770 :returns: True if PA reported success.
771 """
772 with self._lock:
773 if self._failed.is_set() or not self._mainloop or not self._context:
774 return False
775 cvol = _PACVolume()
776 self._lib.pa_cvolume_set(ctypes.byref(cvol), channels, pa_volume)
777
778 done = threading.Event()
779 result: dict[str, int] = {}
780
781 def _success_cb_impl(_ctx: int, success: int, _userdata: int) -> None:
782 result["success"] = success
783 done.set()
784
785 success_cb = _CONTEXT_SUCCESS_CB(_success_cb_impl)
786
787 self._lib.pa_threaded_mainloop_lock(self._mainloop)
788 try:
789 op = self._lib.pa_context_set_sink_volume_by_name(
790 self._context,
791 sink_name.encode(),
792 ctypes.byref(cvol),
793 success_cb,
794 None,
795 )
796 if not op:
797 return False
798 finally:
799 self._lib.pa_threaded_mainloop_unlock(self._mainloop)
800
801 if not done.wait(timeout=_SET_VOLUME_TIMEOUT):
802 self._cancel_operation(op)
803 return False
804 self._lib.pa_operation_unref(op)
805 return bool(result.get("success", 0))
806
807
808# One shared server per MusicAssistant instance; weak keys so a discarded mass
809# (tests) does not pin its server object forever.
810_servers: weakref.WeakKeyDictionary[MusicAssistant, PulseCaptureServer] = (
811 weakref.WeakKeyDictionary()
812)
813
814
815class _PACVolume(ctypes.Structure):
816 _fields_: ClassVar = [
817 ("channels", ctypes.c_uint8),
818 ("values", ctypes.c_uint32 * PA_CHANNELS_MAX),
819 ]
820
821
822def _log_supervisor_exit(task: asyncio.Task[None]) -> None:
823 """Surface a supervisor that died on an unexpected error (restarts stop with it)."""
824 if not task.cancelled() and task.exception() is not None:
825 LOGGER.error("PulseAudio capture supervisor died unexpectedly: %s", task.exception())
826
827
828def _close_controller_result(fut: asyncio.Future[PAVolumeController]) -> None:
829 """Close a controller whose awaiting task was cancelled mid-construction."""
830 with suppress(Exception):
831 fut.result().close()
832
833
834def _load_full_lib() -> ctypes.CDLL:
835 """
836 Load and configure libpulse for use by PAVolumeController.
837
838 Called once per process; the result is cached by _get_full_lib().
839 """
840 lib = ctypes.CDLL("libpulse.so.0")
841
842 lib.pa_threaded_mainloop_new.restype = ctypes.c_void_p
843 lib.pa_threaded_mainloop_get_api.restype = ctypes.c_void_p
844 lib.pa_threaded_mainloop_get_api.argtypes = [ctypes.c_void_p]
845 lib.pa_threaded_mainloop_start.restype = ctypes.c_int
846 lib.pa_threaded_mainloop_start.argtypes = [ctypes.c_void_p]
847 lib.pa_threaded_mainloop_stop.argtypes = [ctypes.c_void_p]
848 lib.pa_threaded_mainloop_free.argtypes = [ctypes.c_void_p]
849 lib.pa_threaded_mainloop_lock.argtypes = [ctypes.c_void_p]
850 lib.pa_threaded_mainloop_unlock.argtypes = [ctypes.c_void_p]
851
852 lib.pa_context_new.restype = ctypes.c_void_p
853 lib.pa_context_new.argtypes = [ctypes.c_void_p, ctypes.c_char_p]
854 lib.pa_context_set_state_callback.argtypes = [
855 ctypes.c_void_p,
856 _CONTEXT_NOTIFY_CB,
857 ctypes.c_void_p,
858 ]
859 lib.pa_context_connect.restype = ctypes.c_int
860 lib.pa_context_connect.argtypes = [
861 ctypes.c_void_p,
862 ctypes.c_char_p,
863 ctypes.c_int,
864 ctypes.c_void_p,
865 ]
866 lib.pa_context_get_state.restype = ctypes.c_int
867 lib.pa_context_get_state.argtypes = [ctypes.c_void_p]
868 lib.pa_context_disconnect.argtypes = [ctypes.c_void_p]
869 lib.pa_context_unref.argtypes = [ctypes.c_void_p]
870
871 lib.pa_cvolume_set.restype = ctypes.c_void_p
872 lib.pa_cvolume_set.argtypes = [ctypes.c_void_p, ctypes.c_uint, ctypes.c_uint32]
873
874 lib.pa_context_set_sink_volume_by_name.restype = ctypes.c_void_p
875 lib.pa_context_set_sink_volume_by_name.argtypes = [
876 ctypes.c_void_p,
877 ctypes.c_char_p,
878 ctypes.c_void_p,
879 _CONTEXT_SUCCESS_CB,
880 ctypes.c_void_p,
881 ]
882 lib.pa_operation_unref.argtypes = [ctypes.c_void_p]
883 lib.pa_operation_cancel.argtypes = [ctypes.c_void_p]
884
885 lib.pa_context_load_module.restype = ctypes.c_void_p
886 lib.pa_context_load_module.argtypes = [
887 ctypes.c_void_p,
888 ctypes.c_char_p,
889 ctypes.c_char_p,
890 _CONTEXT_INDEX_CB,
891 ctypes.c_void_p,
892 ]
893 lib.pa_context_unload_module.restype = ctypes.c_void_p
894 lib.pa_context_unload_module.argtypes = [
895 ctypes.c_void_p,
896 ctypes.c_uint32,
897 _CONTEXT_SUCCESS_CB,
898 ctypes.c_void_p,
899 ]
900 return lib
901
902
903_full_lib: ctypes.CDLL | None = None
904
905
906def _get_full_lib() -> ctypes.CDLL:
907 global _full_lib # noqa: PLW0603
908 if _full_lib is None:
909 _full_lib = _load_full_lib()
910 return _full_lib
911