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