/
/
1"""Sendspin Bridge for Local Audio Out - streams audio to local soundcards."""
2
3from __future__ import annotations
4
5import asyncio
6import sys
7import time
8import uuid
9from contextlib import suppress
10from typing import TYPE_CHECKING, Any, TypeVar, cast
11
12import numpy as np
13from aiosendspin.models.core import ClientHelloPayload
14from aiosendspin.models.core import DeviceInfo as SendspinDeviceInfo
15from aiosendspin.models.player import ClientHelloPlayerSupport, SupportedAudioFormat
16from aiosendspin.models.types import AudioCodec, PlayerCommand
17from music_assistant_models.enums import IdentifierType
18from music_assistant_models.player import DeviceInfo
19
20from music_assistant.helpers.pulse_capture import PAVolumeController, volume_pct_to_amplitude
21from music_assistant.helpers.util import join_task
22from music_assistant.models.player import Player
23from music_assistant.providers.sendspin.bridge_manager import SendspinBridgeManagerBase
24from music_assistant.providers.sendspin.bridge_role import (
25 BRIDGE_BIT_DEPTH,
26 BRIDGE_CHANNELS,
27 BRIDGE_ROLE_ID,
28 BRIDGE_SAMPLE_RATE,
29 BridgePlayerRole,
30)
31from music_assistant.providers.sendspin.helpers import bridge_client_id_from_uuid
32
33from .constants import (
34 AUDIO_BACKEND_ALSA,
35 AUDIO_BACKEND_AUTO,
36 AUDIO_BACKEND_PULSEAUDIO,
37 CACHE_CATEGORY_PREV_STATE,
38 CONF_AUDIO_BACKEND,
39 DEFAULT_BUFFER_FRAMES,
40 DEFAULT_PLAYER_VOLUME,
41 DEVICE_UUID_NAMESPACE,
42 VOLUME_CONTROL_HARDWARE,
43 VOLUME_CONTROL_SOFTWARE,
44)
45
46if sys.platform == "linux":
47 from .pa_simple import (
48 PASimpleStream,
49 enumerate_alsa_devices,
50 enumerate_pa_sinks,
51 suspend_resume_sink,
52 )
53 from .remap_topology import (
54 build_remap_sink_argument,
55 compute_remap_topology,
56 normalize_card_name,
57 )
58
59if TYPE_CHECKING:
60 from collections.abc import Callable
61
62 import sounddevice as sd
63 from aiosendspin.server import (
64 ExternalStreamStartRequest,
65 SendspinClient,
66 SendspinServer,
67 )
68 from aiosendspin.server.roles import AudioChunk
69
70 from music_assistant.providers.sendspin.provider import SendspinProvider
71
72 from .provider import LocalAudioProvider
73
74
75# Chunks arriving more than this many microseconds late are dropped rather than
76# played, since writing stale audio would push subsequent chunks further out of
77# sync with the server timeline.
78_LATE_DROP_THRESHOLD_US = 500_000 # 500 ms
79
80# A sound device that fails again this soon after being reopened is not going to
81# settle by itself (a sink that keeps disappearing, a card whose driver stalls),
82# so the bridge gives up instead of churning the device for the rest of the
83# stream. Measured from the moment the reopen is attempted, which includes the
84# open it is waiting on, so it sits well clear of a full _DEVICE_OPEN_TIMEOUT_SECONDS
85# cycle to leave a slow-reopening device room to prove itself.
86_DEVICE_REOPEN_GUARD_SECONDS: float = 30.0
87
88# How long a writer gets to stop on its own before it is cancelled. It can be
89# parked on a device that stopped responding â mid-write, or waiting out an
90# open â where the graceful stop sentinel never gets read.
91_WRITER_STOP_TIMEOUT_SECONDS: float = 2.0
92
93# An open that has not come back by now is not going to: the device is wedged
94# rather than simply gone. Generous enough that a sink resuming from suspend or
95# a card waking up still makes it.
96_DEVICE_OPEN_TIMEOUT_SECONDS: float = 10.0
97
98_StreamT = TypeVar("_StreamT")
99
100
101def _now_us() -> int:
102 """
103 Return current monotonic time in microseconds.
104
105 Uses CLOCK_MONOTONIC_RAW on Linux to match the aiosendspin server clock,
106 which avoids NTP slewing poisoning playback scheduling.
107 """
108 if sys.platform == "linux" and hasattr(time, "CLOCK_MONOTONIC_RAW"):
109 return time.clock_gettime_ns(time.CLOCK_MONOTONIC_RAW) // 1_000
110 return time.monotonic_ns() // 1_000
111
112
113def get_device_uuid(device_name: str, hostapi_index: int) -> str:
114 """Generate a stable UUID for a local audio device."""
115 return str(uuid.uuid5(DEVICE_UUID_NAMESPACE, f"{device_name}:{hostapi_index}"))
116
117
118class SendspinLocalAudioBridge:
119 """Manages the Sendspin to local soundcard bridge for a single device."""
120
121 def __init__(
122 self,
123 provider: LocalAudioProvider,
124 device_info: dict[str, Any],
125 sendspin_server: SendspinServer,
126 backend: str = "auto",
127 volume_controller: PAVolumeController | None = None,
128 has_remap_children: bool = False,
129 ) -> None:
130 """
131 Initialize the bridge.
132
133 :param provider: The Local Audio provider instance.
134 :param device_info: Device info dict â on Linux PA: from enumerate_pa_sinks();
135 on Linux ALSA or Darwin: from sounddevice.query_devices() with 'index' added.
136 :param sendspin_server: The Sendspin server to register with.
137 :param backend: Resolved audio backend: "pulse", "alsa", or "sounddevice".
138 :param volume_controller: Shared PAVolumeController for hardware sink volume
139 (pulse backend only). If None, falls back to software volume scaling.
140 :param has_remap_children: True if this is a master sink with at least
141 one module-remap-sink.c child in the current enumeration. Ignored
142 for remap sinks themselves.
143 """
144 self.provider = provider
145 self.mass = provider.mass
146 self.sendspin_server = sendspin_server
147 self.device_info = device_info
148 self.device_name: str = device_info["name"]
149 self.display_name: str = device_info.get("description", self.device_name)
150 self.pa_sink_name: str | None = device_info.get("pa_sink_name")
151 self.sample_rate: int = device_info.get("sample_rate", BRIDGE_SAMPLE_RATE)
152 self.bit_depth: int = device_info.get("bit_depth", BRIDGE_BIT_DEPTH)
153 # The PA sink's actual channel count (e.g. 8 for a 7.1 master, 2 for a
154 # remap sink). Used for pa_cvolume_set in PAVolumeController calls â
155 # a mismatch against the sink's real channel_map can leave the
156 # displayed reference_volume updated while soft_volume (real gain)
157 # doesn't change. This is independent of BRIDGE_CHANNELS, which is
158 # the audio *stream's* channel count (always 2).
159 self.pa_channels: int = device_info.get("max_output_channels", BRIDGE_CHANNELS)
160 self.device_index: int | None = device_info.get("index")
161 self.backend: str = backend
162 self.logger = provider.logger.getChild(f"bridge.{self.display_name}")
163
164 # The shared PAVolumeController, if connected for the pulse backend.
165 # Used by _reset_sink_volume() (pin-to-100% fallback, via libpulse
166 # instead of a pactl subprocess) for ALL pulse bridges.
167 self._shared_volume_controller: PAVolumeController | None = (
168 volume_controller if backend == "pulse" else None
169 )
170
171 # Hardware per-player volume control is used for:
172 # - remap sinks (module-remap-sink.c): each has its own independent
173 # PA volume that doesn't affect its master or siblings (validated
174 # empirically), or
175 # - base/master sinks with NO remap-sink children in the current
176 # enumeration (e.g. remap addon disabled, or a standalone card):
177 # nothing shares this sink's master output, so hardware volume is
178 # safe here too.
179 #
180 # A master sink that DOES have remap children must NOT use hardware
181 # volume for itself: its PA volume multiplies into every remap sink
182 # feeding through it, so driving it from *this* player's volume
183 # would silently attenuate every other zone sharing the card. Such
184 # masters stay pinned at 100% (via _reset_sink_volume) and use
185 # software (numpy) volume scaling for their own player.
186 is_remap = device_info.get("is_remap", False)
187 self._volume_controller: PAVolumeController | None = (
188 self._shared_volume_controller if (is_remap or not has_remap_children) else None
189 )
190 self._volume_control_mode = (
191 VOLUME_CONTROL_HARDWARE
192 if self._volume_controller is not None
193 else VOLUME_CONTROL_SOFTWARE
194 )
195
196 # Volume/mute state â owned by the bridge; persisted to MA cache
197 self._volume_level: int = DEFAULT_PLAYER_VOLUME
198 self._volume_muted: bool = False
199 # Stable UUID for this device â set in start()
200 self._device_uuid: str = ""
201
202 self._sendspin_client: SendspinClient | None = None
203 self._bridge_client_id: str | None = None
204 self._bridge_role: BridgePlayerRole | None = None
205 self._is_streaming = False
206 self._logged_chunk_fmt: bool = False
207 # Queue holds (timestamp_us, pcm_bytes) tuples for scheduled playback.
208 # None sentinel signals the writer to stop.
209 self._write_queue: asyncio.Queue[tuple[int, bytes] | None] = asyncio.Queue()
210 # The writer that currently owns the output device. Registration is the
211 # ownership token: a writer only touches shared bridge state while it is
212 # still the one registered here, so a stream that replaced it is left alone.
213 self._writer_task: asyncio.Task[None] | None = None
214 # Monotonic instant of the last output device reopen within the running
215 # writer (None when none happened yet), used to tell a one-off device
216 # failure from one that keeps repeating.
217 self._last_device_reopen: float | None = None
218 # Pre-warmed PA stream â opened during start() so PA stream-open
219 # latency is paid once at provider init, not at first play. Eliminates
220 # the ~200ms cold-start sync offset when all bridges in a group open
221 # their PA streams simultaneously for the first time.
222 self._pa_stream: PASimpleStream | None = None
223 self._lock = asyncio.Lock()
224
225 @property
226 def is_registered(self) -> bool:
227 """Return whether the bridge is registered with Sendspin."""
228 return self._sendspin_client is not None
229
230 async def start(self) -> None:
231 """Register the local audio device as an external Sendspin client."""
232 hostapi_index: int = self.device_info.get("hostapi", 0)
233 self._device_uuid = get_device_uuid(self.device_name, hostapi_index)
234 self._bridge_client_id = bridge_client_id_from_uuid(self._device_uuid)
235
236 if sendspin_prov := self._get_sendspin_provider():
237 sendspin_prov.register_bridge_identifiers(
238 self._bridge_client_id,
239 {IdentifierType.UUID: self._device_uuid},
240 )
241 # The local audio device player (registered under the bare device
242 # uuid) is the player this bridge runs on top of.
243 sendspin_prov.register_bridge_underlying_player(
244 self._bridge_client_id, self._device_uuid
245 )
246
247 # Restore cached volume/mute state from a previous session.
248 # Never restore muted=True on startup: a stale cached mute silently
249 # kills audio every boot with no visible indication in the MA UI.
250 # The user can re-mute intentionally; we should never start muted.
251 if last_state := await self.mass.cache.get(
252 key=self._device_uuid,
253 provider=self.provider.instance_id,
254 category=CACHE_CATEGORY_PREV_STATE,
255 ):
256 self._volume_muted = False
257 self._volume_level = last_state[1]
258
259 # On Linux advertise the sink's native format so MA transcodes correctly.
260 _depths = sorted({self.bit_depth, BRIDGE_BIT_DEPTH}, reverse=True)
261 supported_formats = [
262 SupportedAudioFormat(
263 codec=AudioCodec.PCM,
264 channels=BRIDGE_CHANNELS,
265 sample_rate=self.sample_rate,
266 bit_depth=d,
267 )
268 for d in _depths
269 ]
270
271 hello = ClientHelloPayload(
272 client_id=self._bridge_client_id,
273 name=self.display_name,
274 version=1,
275 supported_roles=[BRIDGE_ROLE_ID, "player@v1"],
276 device_info=SendspinDeviceInfo(
277 product_name=self.display_name,
278 manufacturer="Local Audio",
279 ),
280 player_support=ClientHelloPlayerSupport(
281 supported_formats=supported_formats,
282 buffer_capacity=1_000,
283 supported_commands=[PlayerCommand.VOLUME, PlayerCommand.MUTE],
284 ),
285 )
286
287 self.logger.debug(
288 "Registering Sendspin bridge for %s (client_id=%s)",
289 self.device_name,
290 self._bridge_client_id,
291 )
292
293 self._sendspin_client = self.sendspin_server.register_external_player(
294 hello, on_stream_start=self._on_stream_start
295 )
296
297 for role in self._sendspin_client.roles_by_family("player"):
298 self.logger.debug("Found player role: %s type=%s", role.role_id, type(role).__name__)
299 if isinstance(role, BridgePlayerRole):
300 self._bridge_role = role
301 break
302
303 if self._bridge_role is None:
304 self.logger.error("No BridgePlayerRole found for %s", self.device_name)
305 return
306
307 self._bridge_role.set_callbacks(
308 on_audio_chunk=self._on_audio_chunk,
309 on_volume_change=self._on_volume_change,
310 on_mute_change=self._on_mute_change,
311 on_stream_start=self._on_bridge_stream_start,
312 on_stream_end=self._on_bridge_stream_end,
313 initial_volume=self._volume_level, # restore cached volume for correct slider position
314 )
315 self._bridge_role.setup_audio_requirements(
316 sample_rate=self.sample_rate,
317 bit_depth=self.bit_depth,
318 channels=BRIDGE_CHANNELS,
319 )
320
321 self.logger.info(
322 "Sendspin bridge registered for %s (client_id=%s, volume_control=%s)",
323 self.device_name,
324 self._bridge_client_id,
325 self._volume_control_mode,
326 )
327
328 if self._volume_controller is not None:
329 # Remap sink: apply the restored volume/mute to PA sink hardware
330 # directly so the sink starts at the user's last level.
331 await self._apply_hardware_volume()
332 elif self.backend == "pulse":
333 # Master/base sink (or remap sink if PAVolumeController failed to
334 # connect): pin PA sink volume to 100% so it doesn't attenuate
335 # remap sinks feeding through it; volume is controlled in software
336 # for this player. Also undoes any deferred_volume bleed from
337 # stream volume into sink hardware volume on pa_simple_new open.
338 await self._reset_sink_volume()
339
340 # Pre-warm the PA stream so the latency of pa_simple_new() is paid
341 # once at provider init rather than at first play. Without this, each
342 # bridge in a sync group opens its PA stream on the first play, and the
343 # spread in that per-bridge open latency (~50-200ms) becomes a fixed
344 # sync offset that persists for the session. Pre-warming means all
345 # streams are already open and idle when the first play starts, so
346 # play_at_us scheduling lands all bridges within a much tighter window.
347 if self.backend == "pulse" and self.pa_sink_name:
348 await self._prewarm_pa_stream()
349
350 async def stop(self) -> None:
351 """Stop and unregister the Sendspin bridge."""
352 async with self._lock:
353 await self._stop_streaming()
354 if self._sendspin_client and self._bridge_client_id:
355 await self.sendspin_server.remove_client(self._bridge_client_id)
356 self._sendspin_client = None
357 self._bridge_role = None
358 # Close the pre-warmed stream if it was never consumed by a writer.
359 if self._pa_stream is not None:
360 stream = self._pa_stream
361 self._pa_stream = None
362 with suppress(OSError):
363 await self.mass.loop.run_in_executor(None, stream.close)
364 self.logger.debug("Sendspin bridge stopped for %s", self.device_name)
365
366 def _get_sendspin_provider(self) -> SendspinProvider | None:
367 """Get the Sendspin provider if available."""
368 return cast("SendspinProvider | None", self.mass.get_provider("sendspin"))
369
370 def _on_stream_start(self, request: ExternalStreamStartRequest) -> None:
371 """Handle stream start request from Sendspin server."""
372 self.logger.debug(
373 "Sendspin stream start request for %s (reason=%s)",
374 self.device_name,
375 request.connection_reason,
376 )
377 self._is_streaming = True
378
379 def _on_bridge_stream_start(self) -> None:
380 """Start the audio writer task for a new stream."""
381 if self._writer_task is not None and not self._writer_task.done():
382 self._writer_task.cancel()
383 self._is_streaming = True
384 while not self._write_queue.empty():
385 self._write_queue.get_nowait()
386 # Started on the next loop iteration so the writer is registered before
387 # it runs: everything it touches on the bridge is keyed on being the
388 # registered writer, which an eager start would leave it short of.
389 self._writer_task = self.mass.create_task(self._audio_writer(), eager_start=False)
390 self.logger.info("Bridge writer started for %s", self.device_name)
391
392 def _on_bridge_stream_end(self) -> None:
393 """Stop streaming when the stream ends."""
394 # Clearing the streaming flag is left to the teardown, which is the only
395 # thing that knows whether this end still owns the running writer.
396 self.mass.create_task(self._stop_streaming_locked(self._writer_task))
397 if self._volume_controller is None and self.backend == "pulse":
398 # Master/base sink (software volume control): re-pin PA sink
399 # volume to 100% after playback ends.
400 self.mass.create_task(self._reset_sink_volume())
401
402 async def _prewarm_pa_stream(self) -> None:
403 """
404 Open the PA stream early so stream-open latency doesn't affect sync.
405
406 Opens a PASimpleStream and writes a single silence chunk to trigger
407 the idle->RUNNING ALSA transition (paying the mmap init cost once at
408 provider startup rather than at first play). The stream is then held
409 open and idle â PA may suspend it via suspend-on-idle, but the ALSA
410 driver remains initialised so the first real write reconnects with
411 negligible latency. No keepalive writes are used to avoid interfering
412 with real audio when the writer takes ownership.
413 """
414 if not self.pa_sink_name:
415 return
416 pa_sink_name = self.pa_sink_name
417 bytes_per_sample = self.bit_depth // 8
418 # One chunk of silence to trigger idle->RUNNING and initialise the
419 # ALSA DMA buffer â just enough to pay the mmap init cost.
420 silence = bytes(BRIDGE_CHANNELS * bytes_per_sample * 64) # 64 frames
421 try:
422 stream = await self.mass.loop.run_in_executor(
423 None,
424 lambda: PASimpleStream(
425 sink_name=pa_sink_name,
426 app_name=f"music-assistant-{pa_sink_name}",
427 rate=self.sample_rate,
428 channels=BRIDGE_CHANNELS,
429 bit_depth=self.bit_depth,
430 ),
431 )
432 except OSError as err:
433 self.logger.debug("PA stream pre-warm failed for %s: %s", pa_sink_name, err)
434 return
435 # Write one silence chunk to trigger the ALSA DMA init cycle.
436 try:
437 await self.mass.loop.run_in_executor(None, stream.write, silence)
438 except OSError:
439 with suppress(OSError):
440 await self.mass.loop.run_in_executor(None, stream.close)
441 return
442 self._pa_stream = stream
443 self.logger.debug("PA stream pre-warmed for %s", pa_sink_name)
444
445 async def _reset_sink_volume(self) -> None:
446 """
447 Pin PA sink hardware volume to 100% via the shared PAVolumeController.
448
449 Called at bridge startup and after each playback session for:
450 - Master/base ALSA-card sinks (is_remap=False), where volume control
451 is software (per-player) and the underlying PA sink must stay at
452 unity so it doesn't attenuate every remap sink feeding through it.
453 - The fallback path on remap sinks if PAVolumeController itself
454 failed to connect (self._volume_controller is None in that case
455 too, so this also covers undoing any deferred_volume bleed from
456 stream volume into sink hardware volume).
457 No-op on non-PA backends or if PAVolumeController is unavailable.
458 """
459 if self.backend != "pulse" or not self.pa_sink_name:
460 return
461 if self._shared_volume_controller is None:
462 return
463 pa_sink_name = self.pa_sink_name
464 controller = self._shared_volume_controller
465 try:
466 ok = await self.mass.loop.run_in_executor(
467 None, controller.set_sink_volume, pa_sink_name, 100, self.pa_channels
468 )
469 if not ok:
470 self.logger.warning("Failed to pin PA sink volume to 100%% for %s", pa_sink_name)
471 else:
472 self.logger.debug("Pinned PA sink hardware volume to 100%% for %s", pa_sink_name)
473 except OSError as err:
474 self.logger.warning("Hardware volume control error for %s: %s", pa_sink_name, err)
475
476 def _on_volume_change(self, volume: int) -> None:
477 """Sync volume from Sendspin side to bridge state and persist."""
478 self._volume_level = volume
479 self.mass.create_task(self._save_state())
480 if self._volume_controller is not None:
481 self.mass.create_task(self._apply_hardware_volume())
482
483 def _on_mute_change(self, muted: bool) -> None:
484 """Sync mute from Sendspin side to bridge state and persist."""
485 self._volume_muted = muted
486 self.mass.create_task(self._save_state())
487 if self._volume_controller is not None:
488 self.mass.create_task(self._apply_hardware_volume())
489
490 async def _apply_hardware_volume(self) -> None:
491 """
492 Apply the current volume/mute state to PA sink hardware volume.
493
494 No-op if no PAVolumeController is available (software fallback path).
495 """
496 if self._volume_controller is None or not self.pa_sink_name:
497 return
498 target_pct = 0 if self._volume_muted else self._volume_level
499 self.logger.debug(
500 "apply_hardware_volume %s: target=%d%% (level=%d muted=%s)",
501 self.pa_sink_name,
502 target_pct,
503 self._volume_level,
504 self._volume_muted,
505 )
506 pa_sink_name = self.pa_sink_name
507 controller = self._volume_controller
508 try:
509 ok = await self.mass.loop.run_in_executor(
510 None, controller.set_sink_volume, pa_sink_name, target_pct, self.pa_channels
511 )
512 if not ok:
513 self.logger.warning(
514 "Failed to set hardware volume to %d%% for %s", target_pct, pa_sink_name
515 )
516 except OSError as err:
517 self.logger.warning("Hardware volume control error for %s: %s", pa_sink_name, err)
518
519 def _schedule_hardware_volume_reapply(self) -> None:
520 """Re-apply the cached hardware volume shortly after the sink starts running."""
521 if self._volume_controller is None:
522 return
523
524 # module-device-restore may restore a persisted per-sink-name volume
525 # when the sink transitions idle->RUNNING â which happens on the first
526 # actual write, *after* _apply_hardware_volume() already ran during
527 # start() while the sink was idle, silently reverting it. Re-apply
528 # shortly after the first write so our cached volume wins, once that
529 # one-time revert has had a chance to fire.
530 async def _reapply_hardware_volume_once_active() -> None:
531 await asyncio.sleep(0.3)
532 await self._apply_hardware_volume()
533
534 self.mass.create_task(_reapply_hardware_volume_once_active())
535
536 async def _save_state(self) -> None:
537 """Persist current volume/mute state to MA cache."""
538 if not self._device_uuid:
539 return
540 await self.mass.cache.set(
541 key=self._device_uuid,
542 data=[self._volume_muted, self._volume_level],
543 provider=self.provider.instance_id,
544 category=CACHE_CATEGORY_PREV_STATE,
545 )
546
547 def _on_audio_chunk(self, chunk: AudioChunk) -> None:
548 """Enqueue an incoming audio chunk with its playback timestamp."""
549 if not self._is_streaming:
550 return
551 if not self._logged_chunk_fmt:
552 self.logger.debug(
553 "First chunk: len=%d sample_rate=%d bit_depth=%d timestamp_us=%d",
554 len(chunk.data),
555 self.sample_rate,
556 self.bit_depth,
557 chunk.timestamp_us,
558 )
559 self._logged_chunk_fmt = True
560 self._write_queue.put_nowait((chunk.timestamp_us, chunk.data))
561
562 def _static_delay_us(self) -> int:
563 """Return the configured static playback delay in microseconds."""
564 if self._bridge_role is not None:
565 return self._bridge_role.get_static_delay_us()
566 return 0
567
568 def _apply_format_conversion(self, pcm_data: bytes) -> bytes:
569 """
570 Apply PA format conversion only â no volume scaling.
571
572 Used on the hardware-volume-control path, where mute and volume are
573 handled by PAVolumeController directly on the PA sink. Only the
574 24-bit -> packed s24le repack is still needed here, since PA expects
575 3 bytes/sample but MA delivers 24-bit audio left-justified in 32-bit
576 containers.
577 """
578 if self.bit_depth != 24:
579 return pcm_data
580 samples = np.frombuffer(pcm_data, dtype=np.int32)
581 return samples.view(np.uint8).reshape(-1, 4)[:, 1:].tobytes()
582
583 def _apply_software_volume(self, pcm_data: bytes) -> bytes:
584 """
585 Apply software volume scaling and format conversion.
586
587 Fallback path used when no PAVolumeController is available (or for
588 the ALSA/sounddevice backend, which has no PA hardware volume path).
589 """
590 if self._volume_muted:
591 if self.bit_depth == 24:
592 # PA expects packed s24le: 3 bytes/sample, not 4
593 return b"\x00" * (len(pcm_data) * 3 // 4)
594 return b"\x00" * len(pcm_data)
595 volume = self._volume_level
596 amplitude = volume_pct_to_amplitude(volume)
597 # 1.0 == no-op; skip the scaling pass entirely
598 scale: float | None = None if amplitude >= 1.0 else amplitude
599
600 if self.bit_depth == 32:
601 if scale is None:
602 return pcm_data
603 samples = np.frombuffer(pcm_data, dtype=np.int32).copy()
604 scaled = np.clip(samples.astype(np.float64) * scale, -2147483648, 2147483647)
605 return scaled.astype(np.int32).tobytes()
606
607 if self.bit_depth == 24:
608 # MA delivers 24-bit audio left-justified in 32-bit containers.
609 # Always repack to packed s24le (bytes 1-3 of each int32).
610 samples = np.frombuffer(pcm_data, dtype=np.int32).copy()
611 if scale is not None:
612 samples = np.clip(
613 samples.astype(np.float64) * scale, -2147483648, 2147483647
614 ).astype(np.int32)
615 return samples.view(np.uint8).reshape(-1, 4)[:, 1:].tobytes()
616
617 # 16-bit
618 if scale is None:
619 return pcm_data
620 samples_16 = np.frombuffer(pcm_data, dtype=np.int16).copy()
621 scaled = np.clip(samples_16.astype(np.float64) * scale, -32768, 32767)
622 return scaled.astype(np.int16).tobytes()
623
624 def _prepare_pa_chunk(self, pcm_data: bytes) -> bytes:
625 """
626 Return chunk data ready to be written to the PA sink.
627
628 Volume is applied in software unless a PAVolumeController drives the
629 sink's hardware volume, in which case only the format conversion runs.
630 """
631 if self._volume_controller is not None:
632 return self._apply_format_conversion(pcm_data)
633 return self._apply_software_volume(pcm_data)
634
635 async def _wait_for_chunk_time(self, timestamp_us: int) -> bool:
636 """
637 Sleep until the scheduled playback time for a chunk.
638
639 :param timestamp_us: Server-domain playback timestamp in microseconds.
640 :returns: True if the chunk should be played, False if it arrived too
641 late and should be dropped.
642
643 The bridge runs in-process with the Sendspin server and both use
644 CLOCK_MONOTONIC_RAW (Linux) or monotonic_ns elsewhere, so
645 chunk.timestamp_us is already in the local clock domain â no NTP-style
646 offset conversion is required. The static_delay_us configured by the
647 user is subtracted so that playback on this device is advanced relative
648 to other players, compensating for a slower DAC pipeline.
649 """
650 play_at_us = timestamp_us - self._static_delay_us()
651 now = _now_us()
652 wait_us = play_at_us - now
653 if wait_us < -_LATE_DROP_THRESHOLD_US:
654 self.logger.warning(
655 "Dropping late chunk for %s: %d ms behind schedule",
656 self.device_name,
657 -wait_us // 1_000,
658 )
659 return False
660 if wait_us > 0:
661 await asyncio.sleep(wait_us / 1_000_000)
662 return True
663
664 async def _audio_writer(self) -> None:
665 """Write queued audio to the output device."""
666 # Every writer gets its own reopen budget: a device failure on the
667 # previous stream says nothing about its health on this one.
668 self._last_device_reopen = None
669 if self.backend == "pulse":
670 await self._audio_writer_pulse()
671 else:
672 await self._audio_writer_sounddevice()
673
674 async def _get_pa_stream(self) -> PASimpleStream:
675 """
676 Return a ready PASimpleStream for this bridge's PA sink.
677
678 Uses the pre-warmed stream if available, otherwise opens a new one.
679 """
680 if self._pa_stream is not None:
681 stream = self._pa_stream
682 self._pa_stream = None # take ownership
683 self.logger.debug("Using pre-warmed PA stream for %s", self.pa_sink_name)
684 return stream
685 self.logger.debug(
686 "Opening PA stream: sink=%s rate=%d channels=%d bit_depth=%d",
687 self.pa_sink_name,
688 self.sample_rate,
689 BRIDGE_CHANNELS,
690 self.bit_depth,
691 )
692 assert self.pa_sink_name is not None
693 pa_sink_name = self.pa_sink_name
694 return await self._await_device_open(
695 self.mass.loop.run_in_executor(
696 None,
697 lambda: PASimpleStream(
698 sink_name=pa_sink_name,
699 app_name=f"music-assistant-{pa_sink_name}",
700 rate=self.sample_rate,
701 channels=BRIDGE_CHANNELS,
702 bit_depth=self.bit_depth,
703 ),
704 ),
705 lambda stream: stream.close(),
706 )
707
708 async def _audio_writer_pulse(self) -> None:
709 """Write queued audio to a PA sink via PASimpleStream (Linux)."""
710 stream: PASimpleStream | None = None
711 write_future: asyncio.Future[None] | None = None
712 try:
713 stream = await self._get_pa_stream()
714 self.logger.debug("PA stream ready for %s", self.pa_sink_name)
715 assert stream is not None
716
717 first_chunk_written = False
718
719 while True:
720 item = await self._write_queue.get()
721 if item is None or not self._is_streaming:
722 break
723 timestamp_us, data = item
724 if not await self._wait_for_chunk_time(timestamp_us):
725 continue # late chunk dropped
726 data = self._prepare_pa_chunk(data)
727 try:
728 write_future = self.mass.loop.run_in_executor(None, stream.write, data)
729 await write_future
730 write_future = None
731 except OSError as err:
732 self.logger.error("PA stream error for %s: %s", self.pa_sink_name, err)
733 write_future = None
734 if (stream := await self._reopen_pa_stream(stream)) is None:
735 # closed by the reopen, so the teardown has nothing to release
736 self._abandon_streaming()
737 return
738 # The reopened sink runs through idle->RUNNING again, so the
739 # cached volume has to win over module-device-restore once more.
740 first_chunk_written = False
741 continue
742
743 if not first_chunk_written:
744 first_chunk_written = True
745 self._schedule_hardware_volume_reapply()
746
747 except asyncio.CancelledError:
748 pass
749 except OSError as err:
750 self.logger.error("Failed to open PA stream for %s: %s", self.pa_sink_name, err)
751 self._abandon_streaming()
752 except Exception:
753 # The sink is this player's only audio path, so a writer dying on
754 # anything at all has to take the player out of playback with it.
755 self.logger.exception("PA writer for %s failed", self.pa_sink_name)
756 self._abandon_streaming()
757 finally:
758 await self._finish_pa_writer(stream, write_future)
759
760 async def _finish_pa_writer(
761 self, stream: PASimpleStream | None, write_future: asyncio.Future[None] | None
762 ) -> None:
763 """
764 Release the PA stream a writer owned and hand back its registration.
765
766 :param stream: The stream the writer got, if it got one.
767 :param write_future: A write still in flight, if any.
768 """
769 if self._writer_task is asyncio.current_task():
770 self._is_streaming = False
771 if write_future is not None:
772 # How the write ended does not matter, only that it is no longer
773 # running when the stream underneath it is closed below.
774 with suppress(Exception, asyncio.CancelledError):
775 await asyncio.shield(write_future)
776 if stream is not None:
777 with suppress(OSError):
778 await asyncio.shield(self.mass.loop.run_in_executor(None, stream.close))
779 if self._writer_task is asyncio.current_task():
780 self._writer_task = None
781
782 async def _audio_writer_sounddevice(self) -> None:
783 """Write queued audio to a sounddevice output stream (ALSA on Linux or Darwin)."""
784 import sounddevice as _sd # noqa: PLC0415
785
786 stream: sd.RawOutputStream | None = None
787 try:
788 stream = await self._get_sounddevice_stream()
789 self.logger.debug("sounddevice stream opened for %s", self.device_name)
790 assert stream is not None
791
792 while True:
793 item = await self._write_queue.get()
794 if item is None or not self._is_streaming:
795 break
796 timestamp_us, data = item
797 if not await self._wait_for_chunk_time(timestamp_us):
798 continue # late chunk dropped
799 data = self._apply_software_volume(data)
800 try:
801 await self.mass.loop.run_in_executor(None, stream.write, data)
802 except _sd.PortAudioError as err:
803 self.logger.error("PortAudio error for %s: %s", self.device_name, err)
804 replacement = await self._reopen_sounddevice_stream(stream)
805 if replacement is None:
806 stream = None # closed by the reopen, nothing left to release
807 self._abandon_streaming()
808 return
809 stream = replacement
810 except (_sd.PortAudioError, TimeoutError) as err:
811 self.logger.error("Failed to open sounddevice stream for %s: %s", self.device_name, err)
812 self._abandon_streaming()
813 except Exception:
814 # The device is this player's only audio path, so a writer dying on
815 # anything at all has to take the player out of playback with it.
816 self.logger.exception("Sound device writer for %s failed", self.device_name)
817 self._abandon_streaming()
818 finally:
819 await self._finish_sounddevice_writer(stream)
820
821 async def _finish_sounddevice_writer(self, stream: sd.RawOutputStream | None) -> None:
822 """
823 Release the sound device a writer owned and hand back its registration.
824
825 :param stream: The device the writer got, if it got one.
826 """
827 if self._writer_task is asyncio.current_task():
828 self._is_streaming = False
829 if stream is not None:
830 # Closing waits for the device to play out what it still holds,
831 # which is not something to do on the event loop.
832 with suppress(OSError):
833 await asyncio.shield(
834 self.mass.loop.run_in_executor(None, self._close_sounddevice_stream, stream)
835 )
836 if self._writer_task is asyncio.current_task():
837 self._writer_task = None
838
839 async def _get_sounddevice_stream(self) -> sd.RawOutputStream:
840 """Open a sounddevice output stream for this device, off the event loop."""
841 return await self._await_device_open(
842 self.mass.loop.run_in_executor(None, self._create_sounddevice_stream),
843 lambda stream: self._close_sounddevice_stream(stream, drain=False),
844 )
845
846 async def _await_device_open(
847 self, open_future: asyncio.Future[_StreamT], close_stream: Callable[[_StreamT], None]
848 ) -> _StreamT:
849 """
850 Wait for a device to finish opening, for as long as that is worth doing.
851
852 An open that never comes back would otherwise hold the writer forever
853 while chunks pile up behind a player still reporting playback.
854
855 :param open_future: The open running off the event loop.
856 :param close_stream: How to release the device, used when it arrives
857 after nobody is waiting for it any more.
858 :return: The opened stream.
859 :raises TimeoutError: The device did not open in time.
860 """
861 try:
862 return await join_task(open_future, _DEVICE_OPEN_TIMEOUT_SECONDS)
863 except asyncio.CancelledError:
864 self._release_orphaned_stream(open_future, close_stream)
865 raise
866 except TimeoutError:
867 self._release_orphaned_stream(open_future, close_stream)
868 raise TimeoutError(
869 f"the device did not open within {_DEVICE_OPEN_TIMEOUT_SECONDS:.0f}s"
870 ) from None
871
872 def _release_orphaned_stream(
873 self, open_future: asyncio.Future[_StreamT], close_stream: Callable[[_StreamT], None]
874 ) -> None:
875 """
876 Arrange for a device to be released once an open nobody is waiting for finishes.
877
878 The open runs to completion in its own thread whatever the caller does,
879 and a device nobody hands back keeps every later open of it failing.
880
881 :param open_future: The open that outlived the caller waiting on it.
882 :param close_stream: How to release the device it produces.
883 """
884
885 def _close_device(stream: _StreamT) -> None:
886 # Nobody is left to hand a failed release to, and the executor
887 # future it would surface on is never retrieved.
888 with suppress(Exception):
889 close_stream(stream)
890
891 def _close_when_open(finished: asyncio.Future[_StreamT]) -> None:
892 if finished.cancelled() or finished.exception() is not None:
893 return
894 self.mass.loop.run_in_executor(None, _close_device, finished.result())
895
896 open_future.add_done_callback(_close_when_open)
897
898 def _create_sounddevice_stream(self) -> sd.RawOutputStream:
899 """Open and start a sounddevice output stream, blocking until the device is ready."""
900 import sounddevice as _sd # noqa: PLC0415
901
902 stream = _sd.RawOutputStream(
903 device=self.device_index,
904 samplerate=self.sample_rate,
905 channels=BRIDGE_CHANNELS,
906 dtype="int16",
907 blocksize=DEFAULT_BUFFER_FRAMES,
908 )
909 try:
910 stream.start()
911 except Exception:
912 # PortAudio already handed out the device, so it has to be given
913 # back before the failure propagates.
914 stream.close()
915 raise
916 return stream
917
918 async def _reopen_pa_stream(self, dead_stream: PASimpleStream) -> PASimpleStream | None:
919 """
920 Replace a PA stream that failed mid-playback.
921
922 :param dead_stream: The stream that raised on write. Always closed â
923 libpulse has no reconnect, so a failed connection stays failed.
924 :return: A fresh stream, or None when the device may not be reopened
925 or could not be opened.
926 """
927 with suppress(OSError):
928 await self.mass.loop.run_in_executor(None, dead_stream.close)
929 if not self._start_device_reopen():
930 return None
931 try:
932 stream = await self._get_pa_stream()
933 except OSError as err:
934 # covers the bounded open giving up, which reports a TimeoutError
935 self.logger.error("Could not reopen PA sink %s: %s", self.pa_sink_name, err)
936 return None
937 if self._volume_controller is None:
938 # Opening a stream can bleed its volume into the sink's hardware
939 # volume, which for this path has to stay pinned at unity.
940 await self._reset_sink_volume()
941 self.logger.warning("Reopened PA sink %s after a write failure", self.pa_sink_name)
942 return stream
943
944 async def _reopen_sounddevice_stream(
945 self, dead_stream: sd.RawOutputStream
946 ) -> sd.RawOutputStream | None:
947 """
948 Replace a sounddevice output stream that failed mid-playback.
949
950 :param dead_stream: The stream that raised on write; always closed.
951 :return: A fresh started stream, or None when the device may not be
952 reopened or could not be opened.
953 """
954 import sounddevice as _sd # noqa: PLC0415
955
956 # The device just failed, so closing it can block; it goes off the event
957 # loop rather than stalling every other player behind it.
958 await self.mass.loop.run_in_executor(
959 None, lambda: self._close_sounddevice_stream(dead_stream, drain=False)
960 )
961 if not self._start_device_reopen():
962 return None
963 try:
964 stream = await self._get_sounddevice_stream()
965 except (_sd.PortAudioError, TimeoutError) as err:
966 self.logger.error("Could not reopen sound device %s: %s", self.device_name, err)
967 return None
968 self.logger.warning("Reopened sound device %s after a write failure", self.device_name)
969 return stream
970
971 def _close_sounddevice_stream(self, stream: sd.RawOutputStream, *, drain: bool = True) -> None:
972 """
973 Close a sounddevice output stream.
974
975 :param stream: The stream to close.
976 :param drain: Let buffered audio play out first. Pass False for a device
977 that is no longer responding, where that wait cannot finish.
978 """
979 if drain:
980 with suppress(OSError):
981 stream.stop()
982 with suppress(OSError):
983 stream.close()
984
985 def _start_device_reopen(self) -> bool:
986 """
987 Register an attempt to reopen the output device after a failure.
988
989 :return: True when the attempt may go ahead, False when the device
990 already failed again within the guard window of the previous
991 reopen and the bridge should give up on it.
992 """
993 now = time.monotonic()
994 last_reopen = self._last_device_reopen
995 if last_reopen is not None and now - last_reopen < _DEVICE_REOPEN_GUARD_SECONDS:
996 self.logger.error(
997 "Output device %s failed again within %.0fs of being reopened, "
998 "giving up and taking it out of the Sendspin session",
999 self.device_name,
1000 _DEVICE_REOPEN_GUARD_SECONDS,
1001 )
1002 return False
1003 self._last_device_reopen = now
1004 return True
1005
1006 def _abandon_streaming(self) -> None:
1007 """
1008 Give up on this stream: stop accepting chunks and leave the Sendspin session.
1009
1010 Playback resumes normally on the next stream â the device is only given
1011 up on for the remainder of this one. Called by a writer a newer stream
1012 has already replaced, it does nothing: that stream is not this one's to
1013 give up on.
1014 """
1015 if self._writer_task is not asyncio.current_task():
1016 return
1017 # The device is this player's only audio path and nothing restarts the
1018 # writer within a stream, so leaving is what surfaces the silence
1019 # instead of letting the group hold the player on PLAYING.
1020 self._is_streaming = False
1021 # Leaving ends the Sendspin stream, which unwinds straight back into
1022 # this bridge's own teardown â that must not run inside the writer task
1023 # calling this, which is on its way out.
1024 self.mass.create_task(self._leave_sendspin_session(), eager_start=False)
1025
1026 async def _leave_sendspin_session(self) -> None:
1027 """
1028 Take the bridge out of Sendspin playback, keeping it registered.
1029
1030 A shared group plays on without this device, a solo one stops. The
1031 client stays registered, so the player survives to be grouped again.
1032 """
1033 if not (client := self._sendspin_client):
1034 return
1035 try:
1036 await client.quiesce_to_solo_stopped()
1037 except Exception as err:
1038 self.logger.warning(
1039 "Could not take %s out of its Sendspin session: %s", self.display_name, err
1040 )
1041
1042 async def _stop_streaming_locked(self, writer_task: asyncio.Task[None] | None) -> None:
1043 """
1044 Serialize streaming teardown.
1045
1046 :param writer_task: The writer that was running when teardown was
1047 scheduled. A newer stream having replaced it means this teardown is
1048 stale and is skipped, so the running writer is left alone.
1049 """
1050 async with self._lock:
1051 if self._writer_task is not writer_task:
1052 return
1053 await self._stop_streaming()
1054
1055 async def _stop_streaming(self) -> None:
1056 """Stop streaming (internal, called with lock held)."""
1057 self._is_streaming = False
1058 # Held as a local across the awaits below: a new stream can install its
1059 # own writer meanwhile, and this teardown owns only the one it started on.
1060 if (writer_task := self._writer_task) and not writer_task.done():
1061 if self.backend == "pulse":
1062 # PA writer handles CancelledError cleanly
1063 writer_task.cancel()
1064 with suppress(asyncio.CancelledError, Exception):
1065 await writer_task
1066 else:
1067 # sounddevice (ALSA or Darwin): signal gracefully via None sentinel
1068 while not self._write_queue.empty():
1069 self._write_queue.get_nowait()
1070 self._write_queue.put_nowait(None)
1071 try:
1072 await asyncio.wait_for(writer_task, timeout=_WRITER_STOP_TIMEOUT_SECONDS)
1073 except TimeoutError:
1074 writer_task.cancel()
1075 with suppress(asyncio.CancelledError):
1076 await writer_task
1077 except asyncio.CancelledError:
1078 pass
1079 if self._writer_task is writer_task:
1080 self._writer_task = None
1081 while not self._write_queue.empty():
1082 self._write_queue.get_nowait()
1083
1084
1085class LocalAudioBridgeManager(SendspinBridgeManagerBase[SendspinLocalAudioBridge]):
1086 """Manages Sendspin bridges for all local audio output devices."""
1087
1088 def __init__(self, provider: LocalAudioProvider) -> None:
1089 """Initialize the bridge manager."""
1090 super().__init__(provider)
1091 self._discover_lock = asyncio.Lock()
1092 self._volume_controller: PAVolumeController | None = None
1093 # player_id (device uuid) -> device info dict from the last enumeration
1094 self._devices: dict[str, dict[str, Any]] = {}
1095 self._backend: str = AUDIO_BACKEND_AUTO
1096 # master sink names with at least one remap-sink child (see _remap_master_sinks)
1097 self._remap_masters: set[str] = set()
1098 # sink_name -> PA module index, for sinks local_audio created via
1099 # _ensure_remap_topology(). Unloaded on stop_all().
1100 self._loaded_remap_modules: dict[str, int] = {}
1101
1102 async def discover_and_register(self) -> None:
1103 """Enumerate output devices, register their players and set up the bridges."""
1104 if not self.sendspin_server:
1105 self.logger.debug("Sendspin provider not available, skipping device enumeration")
1106 return
1107
1108 async with self._discover_lock:
1109 # Read audio backend config (Linux only; ignored on Darwin)
1110 configured_backend: str = str(
1111 self.provider.config.get_value(CONF_AUDIO_BACKEND) or AUDIO_BACKEND_AUTO
1112 )
1113 try:
1114 resolved_backend, devices = await self.mass.loop.run_in_executor(
1115 None, self._enumerate_output_devices, configured_backend
1116 )
1117 except Exception as err:
1118 self.logger.warning("Failed to enumerate audio devices: %s", err)
1119 return
1120
1121 self.logger.info(
1122 "Found %d local audio output device(s) via %s backend",
1123 len(devices),
1124 resolved_backend,
1125 )
1126 if not devices:
1127 self.logger.info("No local audio output devices found")
1128 return
1129
1130 await self._ensure_volume_controller(resolved_backend)
1131 if resolved_backend == "pulse":
1132 devices = await self._refresh_after_remap_topology(devices)
1133
1134 self._backend = resolved_backend
1135 self._remap_masters = self._remap_master_sinks(devices)
1136 passthrough_masters = self._remap_masters_with_passthrough(devices)
1137
1138 device_map: dict[str, dict[str, Any]] = {}
1139 for device in devices:
1140 device_name: str = device["name"]
1141 # A master sink fully covered by its own remap-sink topology
1142 # (per-zone sinks plus a full-passthrough
1143 # "_multichannel_stereo" sink) isn't registered as its own
1144 # player â the passthrough sink already provides "play to
1145 # all outputs" with independent hardware volume control.
1146 if not device.get("is_remap", False) and device_name in passthrough_masters:
1147 self.logger.debug(
1148 "Skipping %s â covered by its remap-sink topology",
1149 device.get("description", device_name),
1150 )
1151 continue
1152 device_map[get_device_uuid(device_name, device.get("hostapi", 0))] = device
1153 self._devices = device_map
1154
1155 for device_uuid, device in self._devices.items():
1156 player = self.mass.players.get_player(device_uuid)
1157 if player is None or player.provider is not self.provider:
1158 # not registered yet, or a stale player object left behind by a
1159 # previous provider instance (players survive a provider reload):
1160 # (re)register so the player is bound to this provider instance
1161 player = await self._register_player(
1162 device_uuid, device.get("description", device["name"])
1163 )
1164 if player is None:
1165 # registration skipped - the player is disabled
1166 continue
1167 await self.evaluate_bridge(player)
1168
1169 async def evaluate_bridge(self, player: Player) -> None:
1170 """
1171 Reconcile the Sendspin bridge for a player and update the player's availability.
1172
1173 A local audio player has no playback path other than its bridge, so a
1174 bridge that should exist but failed to start marks the player
1175 unavailable (and a successful rebuild restores it).
1176
1177 :param player: The player to evaluate.
1178 """
1179 await super().evaluate_bridge(player)
1180 if not (self._should_have_bridge(player) and self._lifecycle_allows_bridge(player)):
1181 return
1182 available = self._has_bridge(player.player_id)
1183 if player.available != available:
1184 player._attr_available = available
1185 player.update_state()
1186
1187 async def stop_all(self) -> None:
1188 """Stop all Sendspin bridges and unload any remap sinks we created."""
1189 await super().stop_all()
1190 if self._loaded_remap_modules:
1191 # Give PA a moment to release active streams on the remap sinks
1192 # before unloading their modules â unload-module fails if a stream
1193 # is still open on the sink at the time of the call.
1194 await asyncio.sleep(0.5)
1195 if self._volume_controller is not None:
1196 controller = self._volume_controller
1197 for sink_name, module_index in self._loaded_remap_modules.items():
1198 with suppress(OSError):
1199 ok = await self.mass.loop.run_in_executor(
1200 None, controller.unload_module, module_index
1201 )
1202 if not ok:
1203 self.logger.warning("Failed to unload remap sink %s", sink_name)
1204 self._loaded_remap_modules.clear()
1205 with suppress(OSError):
1206 await self.mass.loop.run_in_executor(None, self._volume_controller.close)
1207 self._volume_controller = None
1208
1209 async def _ensure_volume_controller(self, resolved_backend: str) -> None:
1210 """
1211 Lazily connect the shared PAVolumeController for hardware sink volume.
1212
1213 No-op if not the pulse backend or already connected. Logs and leaves
1214 it as None on failure â bridges fall back to software volume scaling.
1215 """
1216 if resolved_backend != "pulse" or self._volume_controller is not None:
1217 return
1218 try:
1219 self._volume_controller = await self.mass.loop.run_in_executor(None, PAVolumeController)
1220 self.logger.debug("Hardware volume control connected (PAVolumeController)")
1221 except OSError as err:
1222 self.logger.warning(
1223 "Hardware volume control unavailable, falling back to software volume scaling: %s",
1224 err,
1225 )
1226
1227 @staticmethod
1228 def _remap_master_sinks(devices: list[dict[str, Any]]) -> set[str]:
1229 """
1230 Return master sink names with at least one remap-sink child.
1231
1232 Such masters must keep software volume control â their PA volume
1233 multiplies into every remap sink feeding through them. Masters with
1234 no remap children (e.g. remap addon disabled) can safely use
1235 hardware volume control like any standalone sink.
1236 """
1237 return {d["master_device"] for d in devices if d.get("is_remap") and d.get("master_device")}
1238
1239 @staticmethod
1240 def _remap_masters_with_passthrough(devices: list[dict[str, Any]]) -> set[str]:
1241 """
1242 Return master sink names with a full-passthrough remap-sink child.
1243
1244 A passthrough child has the same channel count as its master (a 1:1
1245 remap, e.g. the "_multichannel_stereo" sink from
1246 compute_remap_topology). Such masters are fully covered by their remap-sink topology (per-zone
1247 sinks plus the passthrough for "play to all outputs") and should not
1248 also be registered as their own player.
1249 """
1250 master_channels = {
1251 d["name"]: d.get("max_output_channels", 0) for d in devices if not d.get("is_remap")
1252 }
1253 result: set[str] = set()
1254 for d in devices:
1255 master = d.get("master_device")
1256 if (
1257 d.get("is_remap")
1258 and master
1259 and master_channels.get(master) == d.get("max_output_channels")
1260 ):
1261 result.add(master)
1262 return result
1263
1264 async def _ensure_remap_topology(self, devices: list[dict[str, Any]]) -> bool:
1265 """
1266 Create any missing per-zone/passthrough remap sinks for multi-channel cards.
1267
1268 For each ALSA-card master sink with more than 2 channels, computes
1269 the expected remap-sink topology (see remap_topology module) and
1270 loads module-remap-sink.c for any sink not already present by name.
1271 Idempotent: sinks that already exist (created by a previous run, or
1272 matching this naming convention) are left untouched.
1273
1274 :param devices: Current enumerate_pa_sinks() result.
1275 :returns: True if any new modules were loaded (caller should
1276 re-enumerate to pick up the new sinks).
1277 """
1278 if self._volume_controller is None:
1279 return False
1280 controller = self._volume_controller
1281 existing_names = {d["name"] for d in devices}
1282 loaded_any = False
1283
1284 for device in devices:
1285 # Master sinks are anything that isn't itself a remap-sink
1286 # child (is_remap=False) â this is more robust than matching
1287 # the "driver" field, since PipeWire's PulseAudio-compatibility
1288 # layer reports a generic "PipeWire" driver string for every
1289 # sink it manages (ALSA-backed or not), unlike real PulseAudio
1290 # which reports the literal module name (e.g.
1291 # "module-alsa-card.c"). alsa_card_name and channel_map being
1292 # present (checked below) is what actually confirms this is a
1293 # real ALSA-backed multi-channel sink worth computing a remap
1294 # topology for, rather than e.g. a virtual/monitor/null sink.
1295 if device.get("is_remap"):
1296 continue
1297 channels: int = device.get("max_output_channels", 0)
1298 if channels <= 2:
1299 continue
1300 alsa_card_name = device.get("alsa_card_name")
1301 channel_map: list[str] = device.get("channel_map", [])
1302 if not alsa_card_name or not channel_map:
1303 continue
1304
1305 card_name = normalize_card_name(alsa_card_name)
1306 master_sink_name: str = device["name"]
1307 for spec in compute_remap_topology(card_name, channel_map, channels):
1308 if spec.sink_name in existing_names:
1309 continue
1310 argument = build_remap_sink_argument(spec, master_sink_name)
1311 module_index = await self.mass.loop.run_in_executor(
1312 None, controller.load_module, "module-remap-sink", argument
1313 )
1314 if module_index is None:
1315 self.logger.warning("Failed to create remap sink %s", spec.sink_name)
1316 continue
1317 self.logger.info(
1318 "Created remap sink %s (master=%s, module=%d)",
1319 spec.sink_name,
1320 master_sink_name,
1321 module_index,
1322 )
1323 self._loaded_remap_modules[spec.sink_name] = module_index
1324 existing_names.add(spec.sink_name)
1325 loaded_any = True
1326
1327 # Suspend/resume the master sink to reset the underlying ALSA
1328 # driver state. Specifically works around the snd_ctxfi mmap bug
1329 # (kernel 6.12.x, commit 391e69143d0a) where the X-Fi card's DMA
1330 # stalls after init, causing pa_simple_write to timeout silently.
1331 # Running this on every ALSA-card master at topology creation time
1332 # is safe â it only takes ~0.5s and has no effect on cards that
1333 # don't have the bug.
1334 await self.mass.loop.run_in_executor(None, suspend_resume_sink, master_sink_name)
1335
1336 # Pin the master sink to 100% so it never attenuates remap sinks
1337 # feeding through it. The master has no bridge of its own (it's
1338 # skipped from player registration), so _reset_sink_volume() is
1339 # never called for it from a bridge's start() path â we have to
1340 # do it here explicitly. Without this, module-device-restore or
1341 # any other PA client can leave the master at an arbitrary volume,
1342 # stacking an extra hidden attenuation on top of every remap sink.
1343 with suppress(OSError):
1344 await self.mass.loop.run_in_executor(
1345 None,
1346 controller.set_sink_volume,
1347 master_sink_name,
1348 100,
1349 channels,
1350 )
1351
1352 if loaded_any:
1353 # module-device-restore may restore a persisted per-sink-name
1354 # volume asynchronously shortly after sink creation, which can
1355 # race with and overwrite _apply_hardware_volume() during
1356 # bridge.start(). Give it a moment to settle first so our
1357 # cached-volume application (run after re-enumeration) wins.
1358 await asyncio.sleep(0.3)
1359
1360 return loaded_any
1361
1362 def _bridge_client_id(self, player: Player) -> str | None:
1363 """Return the Sendspin client_id used to bridge a local audio player."""
1364 return bridge_client_id_from_uuid(player.player_id)
1365
1366 def _create_bridge(self, player: Player) -> SendspinLocalAudioBridge:
1367 """Create a bridge instance for a local audio player."""
1368 sendspin_server = self.sendspin_server
1369 assert sendspin_server is not None # guaranteed by _lifecycle_allows_bridge
1370 device = self._devices[player.player_id]
1371 return SendspinLocalAudioBridge(
1372 cast("LocalAudioProvider", self.provider),
1373 device,
1374 sendspin_server,
1375 backend=self._backend,
1376 volume_controller=self._volume_controller,
1377 has_remap_children=device["name"] in self._remap_masters,
1378 )
1379
1380 def _should_have_bridge(self, player: Player) -> bool:
1381 """Return whether a local audio player should have a Sendspin bridge."""
1382 return player.player_id in self._devices
1383
1384 async def _register_player(self, device_uuid: str, display_name: str) -> Player | None:
1385 """
1386 Register a local audio output device as a regular, visible player.
1387
1388 The player has no native playback features of its own â the Sendspin
1389 bridge provides its (single) output protocol â and its enabled flag
1390 is the on/off toggle for the device.
1391
1392 :param device_uuid: The stable device UUID, used as player_id.
1393 :param display_name: The human friendly device name.
1394 :return: The registered player, or None when the player is disabled
1395 (registration of disabled players is skipped).
1396 """
1397 player = Player(self.provider, device_uuid)
1398 player._attr_name = display_name
1399 player._attr_available = True
1400 player._attr_supported_features = set() # no direct playback capability
1401 player._attr_device_info = DeviceInfo(model=display_name, manufacturer="Local Audio")
1402 player._attr_device_info.add_identifier(IdentifierType.UUID, device_uuid)
1403 await self.mass.players.register_or_update(player)
1404 return self.mass.players.get_player(device_uuid)
1405
1406 async def _refresh_after_remap_topology(
1407 self, devices: list[dict[str, Any]]
1408 ) -> list[dict[str, Any]]:
1409 """
1410 Create missing remap sinks and re-enumerate if any were added.
1411
1412 :returns: The original devices list, or a freshly re-enumerated list
1413 if _ensure_remap_topology() loaded any new modules.
1414 """
1415 if not await self._ensure_remap_topology(devices):
1416 return devices
1417 try:
1418 new_devices = await self.mass.loop.run_in_executor(None, enumerate_pa_sinks)
1419 except (FileNotFoundError, RuntimeError) as err:
1420 self.logger.warning("Failed to re-enumerate after creating remap sinks: %s", err)
1421 return devices
1422 self.logger.info(
1423 "Found %d local audio output device(s) after creating remap sinks", len(new_devices)
1424 )
1425 return new_devices
1426
1427 @staticmethod
1428 def _enumerate_output_devices(backend: str) -> tuple[str, list[dict[str, Any]]]:
1429 """
1430 Enumerate available audio output devices for the given backend.
1431
1432 :param backend: One of AUDIO_BACKEND_AUTO / AUDIO_BACKEND_PULSEAUDIO /
1433 AUDIO_BACKEND_ALSA (from constants). On non-Linux the backend is
1434 always resolved to "sounddevice".
1435 :returns: Tuple of (resolved_backend, device_list).
1436 """
1437 if sys.platform != "linux":
1438 import sounddevice as _sd # noqa: PLC0415
1439
1440 devices: list[dict[str, Any]] = []
1441 for idx, dev in enumerate(_sd.query_devices()):
1442 if dev.get("max_output_channels", 0) < 2:
1443 continue
1444 try:
1445 test_stream = _sd.RawOutputStream(
1446 device=idx,
1447 samplerate=BRIDGE_SAMPLE_RATE,
1448 channels=BRIDGE_CHANNELS,
1449 dtype="int16",
1450 )
1451 test_stream.close()
1452 except _sd.PortAudioError:
1453 continue
1454 dev_with_index = dict(dev)
1455 dev_with_index["index"] = idx
1456 devices.append(dev_with_index)
1457 return "sounddevice", devices
1458
1459 # Linux â dispatch by configured backend
1460 if backend == AUDIO_BACKEND_ALSA:
1461 return "alsa", enumerate_alsa_devices()
1462
1463 if backend == AUDIO_BACKEND_PULSEAUDIO:
1464 return "pulse", enumerate_pa_sinks()
1465
1466 # AUTO: try PA, fall back to ALSA
1467 try:
1468 sinks = enumerate_pa_sinks()
1469 if sinks:
1470 return "pulse", sinks
1471 except (FileNotFoundError, RuntimeError): # fmt: skip
1472 pass
1473 return "alsa", enumerate_alsa_devices()
1474