music-assistant-server

33.6 KBPY
provider.py
33.6 KB849 lines • python
1"""MSX Bridge Player Provider implementation."""
2
3from __future__ import annotations
4
5import asyncio
6import contextlib
7import hashlib
8import hmac
9import logging
10import secrets
11import time
12from collections import deque
13from collections.abc import AsyncIterator
14from typing import TYPE_CHECKING, Any, cast
15
16from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
17from music_assistant_models.enums import ConfigEntryType
18
19from music_assistant.models.player_provider import PlayerProvider
20
21from .constants import (
22    CONF_ENABLE_GROUPING,
23    CONF_ENABLE_SENDSPIN_BRIDGE,
24    CONF_GROUP_STREAM_MODE,
25    CONF_HTTP_PORT,
26    CONF_OUTPUT_FORMAT,
27    CONF_PLAYER_IDLE_TIMEOUT,
28    CONF_SHOW_STOP_NOTIFICATION,
29    DEFAULT_ENABLE_GROUPING,
30    DEFAULT_ENABLE_SENDSPIN_BRIDGE,
31    DEFAULT_GROUP_STREAM_MODE,
32    DEFAULT_HTTP_PORT,
33    DEFAULT_OUTPUT_FORMAT,
34    DEFAULT_PLAYER_IDLE_TIMEOUT,
35    DEFAULT_SHOW_STOP_NOTIFICATION,
36    GROUP_STREAM_MODE_INDEPENDENT,
37    GROUP_STREAM_MODE_REDIRECT,
38    GROUP_STREAM_MODE_SHARED,
39    MSX_PLAYER_ID_PREFIX,
40)
41from .http_server import MSXHTTPServer
42from .player import MSXPlayer
43
44if TYPE_CHECKING:
45    from music_assistant.controllers.streams.audio_processing import AudioOutputPlan
46
47    from .sendspin_bridge import MSXSendspinBridgeManager
48
49logger = logging.getLogger(__name__)
50
51
52class SharedGroupStream:
53    """
54    Shared audio stream for a player group.
55
56    One ffmpeg process produces audio, multiple TV clients read from a shared buffer.
57    Late joiners receive buffered data first (catch-up), then live chunks.
58    """
59
60    def __init__(self, group_id: str, media_uri: str) -> None:
61        """Initialize shared stream for a group."""
62        self.group_id = group_id
63        self.media_uri = media_uri
64        self.buffer: deque[bytes] = deque(maxlen=512)  # ~15s @ 40KB/s MP3
65        self.subscribers: dict[str, asyncio.Queue[bytes | None]] = {}
66        self.producer_task: asyncio.Task[None] | None = None
67        self.started = asyncio.Event()
68        self.finished = False
69        self.producer_error: Exception | None = None
70        self.output_plan: AudioOutputPlan | None = None
71        self._lock = asyncio.Lock()
72        self._total_bytes = 0
73        self._start_time: float = 0
74
75        logger.info(
76            "[SharedStream:%s] Created for media_uri=%s",
77            self.group_id,
78            self.media_uri[:80] if self.media_uri else "N/A",
79        )
80
81    async def start(
82        self,
83        audio_chunks: AsyncIterator[bytes],
84    ) -> None:
85        """Start producing audio from the given chunk iterator."""
86        logger.info(
87            "[SharedStream:%s] Starting producer task",
88            self.group_id,
89        )
90        self._start_time = time.monotonic()
91        self.producer_task = asyncio.create_task(self._produce(audio_chunks))
92
93    async def subscribe(self, player_id: str) -> AsyncIterator[bytes]:
94        """
95        Subscribe to stream, get buffered + live chunks.
96
97        Yields:
98            Audio chunks (bytes). First yields catch-up buffer, then live chunks.
99        """
100        # Large queue to handle slow readers (TV with weak WiFi)
101        q: asyncio.Queue[bytes | None] = asyncio.Queue(maxsize=512)
102
103        bytes_sent = 0
104        chunks_sent = 0
105
106        try:
107            # Wait for stream to start before registering so we don't receive
108            # chunks that are already in the catch-up buffer via both paths.
109            try:
110                await asyncio.wait_for(self.started.wait(), timeout=15.0)
111            except TimeoutError:
112                logger.error(
113                    "[SharedStream:%s] Timeout waiting for stream start for %s",
114                    self.group_id,
115                    player_id,
116                )
117                raise
118
119            if self.producer_error is not None:
120                logger.error(
121                    "[SharedStream:%s] Producer failed before subscriber %s could join: %s",
122                    self.group_id,
123                    player_id,
124                    self.producer_error,
125                )
126                raise self.producer_error
127
128            # Phase 1: Snapshot buffer and register for live chunks atomically.
129            # Holding the lock ensures the producer cannot distribute a new chunk
130            # between the snapshot and the registration, eliminating the window
131            # where the first chunk would appear in both the catch-up buffer and
132            # the subscriber's live queue.
133            # Also capture self.finished under the lock: if the producer already
134            # completed before we registered, we must self-signal EOF so that
135            # the live-stream loop below does not block forever.
136            async with self._lock:
137                buffer_snapshot = list(self.buffer)
138                already_finished = self.finished
139                self.subscribers[player_id] = q
140                if already_finished:
141                    q.put_nowait(None)
142                subscriber_count = len(self.subscribers)
143
144            logger.info(
145                "[SharedStream:%s] Subscriber %s joined (total: %d)",
146                self.group_id,
147                player_id,
148                subscriber_count,
149            )
150
151            buffer_bytes = sum(len(c) for c in buffer_snapshot)
152            logger.debug(
153                "[SharedStream:%s] Sending %d catch-up chunks (%d bytes) to %s",
154                self.group_id,
155                len(buffer_snapshot),
156                buffer_bytes,
157                player_id,
158            )
159            for chunk in buffer_snapshot:
160                yield chunk
161                bytes_sent += len(chunk)
162                chunks_sent += 1
163
164            # Phase 2: Live stream
165            while True:
166                next_chunk = await q.get()
167                if next_chunk is None:
168                    logger.debug(
169                        "[SharedStream:%s] EOF received for subscriber %s",
170                        self.group_id,
171                        player_id,
172                    )
173                    break
174                yield next_chunk
175                bytes_sent += len(next_chunk)
176                chunks_sent += 1
177
178        finally:
179            async with self._lock:
180                self.subscribers.pop(player_id, None)
181                remaining = len(self.subscribers)
182
183            logger.info(
184                "[SharedStream:%s] Subscriber %s left after %d chunks, %d bytes (remaining: %d)",
185                self.group_id,
186                player_id,
187                chunks_sent,
188                bytes_sent,
189                remaining,
190            )
191
192    async def stop(self) -> None:
193        """Stop the stream and clean up."""
194        logger.info(
195            "[SharedStream:%s] Stopping (total: %d bytes)",
196            self.group_id,
197            self._total_bytes,
198        )
199        if self.producer_task and not self.producer_task.done():
200            self.producer_task.cancel()
201            with contextlib.suppress(asyncio.CancelledError):
202                await self.producer_task
203        self.finished = True
204
205    @property
206    def subscriber_count(self) -> int:
207        """Return current subscriber count."""
208        return len(self.subscribers)
209
210    async def _produce(self, audio_chunks: AsyncIterator[bytes]) -> None:
211        """Read from ffmpeg and distribute to all subscribers."""
212        try:
213            chunk_count = 0
214            async for chunk in audio_chunks:
215                chunk_count += 1
216                self._total_bytes += len(chunk)
217                async with self._lock:
218                    self.buffer.append(chunk)
219                    for player_id, q in list(self.subscribers.items()):
220                        try:
221                            q.put_nowait(chunk)
222                        except asyncio.QueueFull:
223                            logger.warning(
224                                "[SharedStream:%s] Queue full for subscriber %s, dropping chunk %d",
225                                self.group_id,
226                                player_id,
227                                chunk_count,
228                            )
229
230                if not self.started.is_set():
231                    # Signal that stream has started (first chunk received)
232                    self.started.set()
233                    logger.debug(
234                        "[SharedStream:%s] First chunk received, signaling started",
235                        self.group_id,
236                    )
237
238            logger.info(
239                "[SharedStream:%s] Producer finished: %d chunks, %d bytes, %.1fs",
240                self.group_id,
241                chunk_count,
242                self._total_bytes,
243                time.monotonic() - self._start_time,
244            )
245        except asyncio.CancelledError:
246            logger.debug("[SharedStream:%s] Producer cancelled", self.group_id)
247            raise
248        except Exception as exc:
249            logger.exception("[SharedStream:%s] Producer error", self.group_id)
250            self.producer_error = exc
251        finally:
252            self.finished = True
253            # Ensure subscribers waiting on `started` can proceed even if no chunks were produced
254            if not self.started.is_set():
255                self.started.set()
256                logger.debug(
257                    "[SharedStream:%s] Producer completed without first chunk, signaling started",
258                    self.group_id,
259                )
260            # Signal EOF to all subscribers
261            async with self._lock:
262                for player_id, q in list(self.subscribers.items()):
263                    with contextlib.suppress(asyncio.QueueFull):
264                        q.put_nowait(None)
265                    logger.debug(
266                        "[SharedStream:%s] Sent EOF to subscriber %s",
267                        self.group_id,
268                        player_id,
269                    )
270
271
272class MSXBridgeProvider(PlayerProvider):
273    """Player Provider that bridges Music Assistant to Smart TVs via MSX."""
274
275    http_server: MSXHTTPServer | None = None
276    grouping_enabled: bool = True
277    group_stream_mode: str = DEFAULT_GROUP_STREAM_MODE
278    sendspin_bridge_enabled: bool = False
279    bridge_manager: MSXSendspinBridgeManager | None = None
280    _player_last_activity: dict[str, float]
281    _pending_unregisters: dict[str, asyncio.Event]
282    _stream_token_secret: bytes
283    _timeout_task: asyncio.Task[None] | None = None
284    _owner_username: str | None = None
285    _shared_streams: dict[str, SharedGroupStream]  # group_id -> SharedGroupStream
286    _shared_stream_lock: asyncio.Lock
287    _background_tasks: set[asyncio.Task[None]]  # fire-and-forget tasks (unregister, stream stop)
288
289    def __init__(self, *args: Any, **kwargs: Any) -> None:
290        """Initialize the provider."""
291        super().__init__(*args, **kwargs)
292        self._player_last_activity = {}
293        self._pending_unregisters = {}
294        # one secret per provider instance; the per-player tokens derive from it
295        self._stream_token_secret = secrets.token_bytes(32)
296        self._shared_streams = {}
297        self._shared_stream_lock = asyncio.Lock()
298        self._background_tasks = set()
299
300    async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
301        """Return Config entries to configure this provider."""
302        return (
303            ConfigEntry(
304                key=CONF_HTTP_PORT,
305                type=ConfigEntryType.INTEGER,
306                required=True,
307                default_value=str(DEFAULT_HTTP_PORT),
308            ),
309            ConfigEntry(
310                key=CONF_OUTPUT_FORMAT,
311                type=ConfigEntryType.STRING,
312                required=True,
313                default_value=DEFAULT_OUTPUT_FORMAT,
314            ),
315            ConfigEntry(
316                key=CONF_PLAYER_IDLE_TIMEOUT,
317                type=ConfigEntryType.INTEGER,
318                required=False,
319                default_value=str(DEFAULT_PLAYER_IDLE_TIMEOUT),
320            ),
321            ConfigEntry(
322                key=CONF_SHOW_STOP_NOTIFICATION,
323                type=ConfigEntryType.BOOLEAN,
324                required=False,
325                default_value=DEFAULT_SHOW_STOP_NOTIFICATION,
326            ),
327            ConfigEntry(
328                key=CONF_ENABLE_GROUPING,
329                type=ConfigEntryType.BOOLEAN,
330                required=False,
331                default_value=DEFAULT_ENABLE_GROUPING,
332            ),
333            ConfigEntry(
334                key=CONF_ENABLE_SENDSPIN_BRIDGE,
335                type=ConfigEntryType.BOOLEAN,
336                required=False,
337                default_value=DEFAULT_ENABLE_SENDSPIN_BRIDGE,
338            ),
339            ConfigEntry(
340                key=CONF_GROUP_STREAM_MODE,
341                type=ConfigEntryType.STRING,
342                required=False,
343                default_value=DEFAULT_GROUP_STREAM_MODE,
344                options=[
345                    ConfigValueOption(
346                        GROUP_STREAM_MODE_INDEPENDENT,
347                    ),
348                    ConfigValueOption(
349                        GROUP_STREAM_MODE_SHARED,
350                    ),
351                    ConfigValueOption(
352                        GROUP_STREAM_MODE_REDIRECT,
353                    ),
354                ],
355            ),
356        )
357
358    async def handle_async_init(self) -> None:
359        """Handle async initialization — start embedded HTTP server."""
360        raw_port = cast("int", self.config.get_value(CONF_HTTP_PORT, DEFAULT_HTTP_PORT))
361        port = max(1, min(65535, int(raw_port)))
362        self.grouping_enabled = bool(
363            self.config.get_value(CONF_ENABLE_GROUPING, DEFAULT_ENABLE_GROUPING)
364        )
365        self.group_stream_mode = cast(
366            "str",
367            self.config.get_value(CONF_GROUP_STREAM_MODE, DEFAULT_GROUP_STREAM_MODE),
368        )
369        self.sendspin_bridge_enabled = bool(
370            self.config.get_value(CONF_ENABLE_SENDSPIN_BRIDGE, DEFAULT_ENABLE_SENDSPIN_BRIDGE)
371        )
372        # The Sendspin bridge rides on MA's Sendspin provider; an install that
373        # ships no Sendspin provider can't import the manager. That's a valid
374        # setup — degrade to "no bridge" rather than failing to load.
375        try:
376            self.bridge_manager = self._make_bridge_manager()
377        except ImportError:
378            self.bridge_manager = None
379            if self.sendspin_bridge_enabled:
380                self.logger.warning(
381                    "Sendspin bridge enabled but the Sendspin provider is not available; "
382                    "the bridge is disabled"
383                )
384        self.http_server = MSXHTTPServer(self, port)
385        await self.http_server.start()
386        self.logger.info(
387            "MSX Bridge provider initialized, HTTP server on port %s, group_stream_mode=%s",
388            port,
389            self.group_stream_mode,
390        )
391
392    async def loaded_in_mass(self) -> None:
393        """Start idle timeout task after provider is loaded."""
394        await super().loaded_in_mass()
395        self._timeout_task = self.mass.create_task(self._run_idle_timeout_loop())
396        self.logger.info("MSX Bridge provider loaded — players register on demand")
397
398    async def unload(self, is_removed: bool = False) -> None:
399        """Handle unload — stop timeout task, HTTP server, then unregister players."""
400        if self._timeout_task and not self._timeout_task.done():
401            self._timeout_task.cancel()
402            with contextlib.suppress(asyncio.CancelledError):
403                await self._timeout_task
404            self._timeout_task = None
405
406        # Cancel and await any in-flight unregister tasks before proceeding
407        for task in list(self._background_tasks):
408            if not task.done():
409                task.cancel()
410        if self._background_tasks:
411            await asyncio.gather(*self._background_tasks, return_exceptions=True)
412        self._background_tasks.clear()
413
414        # Cleanup shared streams
415        await self.cleanup_shared_streams()
416
417        if self.bridge_manager:
418            await self.bridge_manager.close()
419            self.bridge_manager = None
420
421        if self.http_server:
422            await self.http_server.stop()
423        for player in list(self.players):
424            try:
425                self.logger.debug("Unloading player %s", player.display_name)
426                await self.mass.players.unregister(player.player_id)
427            except Exception:
428                self.logger.exception("Error unregistering player %s", player.player_id)
429        self._player_last_activity.clear()
430        self.logger.info("MSX Bridge provider unloaded")
431
432    async def get_owner_username(self) -> str | None:
433        """Resolve and cache the first enabled user's username for playlog attribution."""
434        if self._owner_username is None:
435            try:
436                users = await self.mass.webserver.auth.list_users()
437                for user in users:
438                    if user.enabled and user.username:
439                        self._owner_username = user.username
440                        self.logger.debug("Resolved owner username: %s", self._owner_username)
441                        break
442            except Exception as err:
443                self.logger.warning("Could not resolve owner username: %s", err)
444        return self._owner_username
445
446    async def discover_players(self) -> None:
447        """Discover players — MSX players are registered on demand when TVs connect."""
448
449    async def get_or_register_player(
450        self,
451        player_id: str,
452        display_name: str | None = None,
453        ip_address: str | None = None,
454    ) -> MSXPlayer | None:
455        """
456        Get or register an MSX player for the given player_id.
457
458        Returns the player, or None if registration failed.
459        """
460        # Wait for any pending unregister to complete (race condition handling)
461        if pending_event := self._pending_unregisters.get(player_id):
462            self.logger.debug("Waiting for pending unregister of %s before registering", player_id)
463            await pending_event.wait()
464        existing = self.mass.players.get_player(player_id, raise_unavailable=False)
465        if existing and isinstance(existing, MSXPlayer):
466            if ip_address and not existing.device_info.ip_address:
467                existing.device_info.ip_address = ip_address
468            self.on_player_activity(player_id)
469            return existing
470        output_format = cast(
471            "str", self.config.get_value(CONF_OUTPUT_FORMAT, DEFAULT_OUTPUT_FORMAT)
472        )
473        name = display_name or self._player_display_name_from_id(player_id)
474        player = MSXPlayer(
475            provider=self,
476            player_id=player_id,
477            name=name,
478            output_format=output_format,
479            grouping_enabled=self.grouping_enabled,
480            ip_address=ip_address,
481        )
482        await self.mass.players.register(player)
483        self._player_last_activity[player_id] = time.monotonic()
484        self.logger.info("Registered MSX player: %s (%s)", name, player_id)
485        if self.bridge_manager:
486            await self.bridge_manager.evaluate_bridge(player)
487        return player
488
489    def on_player_activity(self, player_id: str) -> None:
490        """Record activity for a player (extends idle timeout)."""
491        # Monotonic: a wall-clock NTP step must not age players past the cutoff
492        self._player_last_activity[player_id] = time.monotonic()
493
494    def on_player_disabled(self, player_id: str) -> None:
495        """
496        Handle player disabled: do not unregister (base would unregister).
497
498        MSX players are registered on demand; unregister on disable would remove them
499        from the list. On enable, discovery is empty so the player would not come back
500        until the TV reconnects. We keep the player registered but disabled so it stays
501        visible in the list when re-enabled.
502
503        Still stop playback on TV by broadcasting stop and cancelling streams.
504        """
505        if self.http_server:
506            self.http_server.broadcast_stop(player_id)
507            self.http_server.cancel_streams_for_player(player_id)
508        # Do NOT call super() — base PlayerProvider unregisters the player here.
509
510    def on_player_enabled(self, player_id: str) -> None:
511        """Handle player enabled: no-op, player already registered."""
512        # Player was never unregistered (see on_player_disabled), so nothing to do.
513
514    async def remove_player(self, player_id: str) -> None:
515        """
516        Remove (delete) a player from this provider.
517
518        Called when user chooses to remove the player from MA.
519        This fully unregisters the player. It will reappear if the TV reconnects.
520        """
521        if self.http_server:
522            self.http_server.broadcast_stop(player_id)
523            self.http_server.cancel_streams_for_player(player_id)
524        await self._handle_player_unregister(player_id)
525        self.logger.info("Player %s removed by user", player_id)
526
527    def notify_play_started(
528        self,
529        player_id: str,
530        *,
531        title: str | None = None,
532        artist: str | None = None,
533        image_url: str | None = None,
534        duration: int | None = None,
535        next_action: str | None = None,
536        prev_action: str | None = None,
537    ) -> None:
538        """Notify WebSocket clients that playback started (for MA -> MSX push)."""
539        if self.http_server:
540            self.http_server.broadcast_play(
541                player_id,
542                title=title,
543                artist=artist,
544                image_url=image_url,
545                duration=duration,
546                next_action=next_action,
547                prev_action=prev_action,
548            )
549
550    def notify_play_playlist(
551        self,
552        player_id: str,
553        start_index: int = 0,
554        queue_id: str | None = None,
555    ) -> None:
556        """Notify WebSocket clients to play an MSX native playlist from the MA queue."""
557        if self.http_server:
558            qid = queue_id or player_id
559            url = f"/msx/queue-playlist/{player_id}.json?start={start_index}&queue_id={qid}"
560            self.http_server.broadcast_playlist(player_id, url)
561
562    def notify_goto_index(self, player_id: str, index: int) -> None:
563        """Notify WebSocket clients to jump to a specific playlist index."""
564        if self.http_server:
565            self.http_server.broadcast_goto_index(player_id, index)
566
567    def notify_play_paused(self, player_id: str) -> None:
568        """Notify WebSocket clients that playback is paused (MA pause -> MSX)."""
569        if self.http_server:
570            self.http_server.broadcast_pause(player_id)
571
572    def notify_play_resumed(self, player_id: str) -> None:
573        """Notify WebSocket clients that playback resumed (MA resume -> MSX)."""
574        if self.http_server:
575            self.http_server.broadcast_resume(player_id)
576
577    def notify_play_stopped(self, player_id: str) -> None:
578        """
579        Notify WebSocket clients that playback stopped (MA stop -> MSX).
580
581        Sends broadcast_stop + cancel_streams twice — same as Disable flow, which
582        stops playback on MSX instantly (vs single signal with ~30s delay).
583        """
584        server = self.http_server
585        if not server:
586            return
587
588        def _send() -> None:
589            server.broadcast_stop(player_id)
590            server.cancel_streams_for_player(player_id)
591
592        _send()
593        _send()
594
595    def notify_seek(self, player_id: str, position_seconds: int) -> None:
596        """Notify WebSocket clients to seek to position (MA seek -> MSX)."""
597        if self.http_server:
598            self.http_server.broadcast_seek(player_id, position_seconds)
599
600    # --- Group Stream Management ---
601
602    def is_shared_stream_mode(self) -> bool:
603        """Check if shared buffer stream mode is enabled."""
604        return self.group_stream_mode == GROUP_STREAM_MODE_SHARED
605
606    def is_redirect_stream_mode(self) -> bool:
607        """
608        Check if MA redirect stream mode is enabled.
609
610        In redirect mode the TV is 302-redirected to the MA Streamserver
611        (``resolve_stream_url``) instead of being served by the local
612        proxy/ffmpeg pipeline. See also ``get_ma_stream_url()``.
613        """
614        return self.group_stream_mode == GROUP_STREAM_MODE_REDIRECT
615
616    def get_stream_token(self, player_id: str) -> str:
617        """
618        Return the token that authorizes the audio routes for the given player.
619
620        Derived rather than stored, so a caller cannot grow provider state by asking for
621        tokens under new player ids. It stays the same for the provider's lifetime: an
622        idle TV is unregistered after the configured timeout, and changing the token there
623        would strand the URLs a long-running kiosk has already cached.
624
625        :param player_id: The player to build an audio URL for.
626        """
627        digest = hmac.new(self._stream_token_secret, player_id.encode(), hashlib.sha256)
628        return digest.hexdigest()[:32]
629
630    def get_group_id_for_player(self, player: MSXPlayer) -> str | None:
631        """
632        Get group ID if player is in a group (as leader or member).
633
634        Returns:
635            group_id if player is grouped, None if solo player.
636        """
637        # If player is synced to another (member), use leader's ID as group
638        if player.synced_to:
639            logger.debug(
640                "[GroupStream] Player %s is member of group %s",
641                player.player_id,
642                player.synced_to,
643            )
644            return player.synced_to
645
646        # If player has group members (is leader), use own ID as group
647        if player.group_members and len(player.group_members) > 1:
648            logger.debug(
649                "[GroupStream] Player %s is leader of group with %d members",
650                player.player_id,
651                len(player.group_members),
652            )
653            return player.player_id
654
655        # Solo player
656        return None
657
658    async def get_or_create_shared_stream(
659        self,
660        group_id: str,
661        media_uri: str,
662        audio_chunks: AsyncIterator[bytes],
663    ) -> SharedGroupStream:
664        """
665        Get existing shared stream or create new one for the group.
666
667        Args:
668            group_id: ID of the group (leader's player_id)
669            media_uri: URI of the media being streamed
670            audio_chunks: Async iterator yielding encoded audio chunks
671
672        Returns:
673            SharedGroupStream instance
674        """
675        # Serialize check-and-create: without the lock, two concurrent callers
676        # replacing an old stream both pass the "existing" check while awaiting
677        # existing.stop(), creating two producers — one orphaned.
678        async with self._shared_stream_lock:
679            existing = self._shared_streams.get(group_id)
680
681            # Reuse existing if same media and not finished
682            if existing and not existing.finished and existing.media_uri == media_uri:
683                logger.info(
684                    "[GroupStream] Reusing existing shared stream for group %s (subscribers: %d)",
685                    group_id,
686                    existing.subscriber_count,
687                )
688                return existing
689
690            # Clean up old stream if exists
691            if existing:
692                logger.info(
693                    "[GroupStream] Replacing old shared stream for group %s "
694                    "(old_uri=%s, new_uri=%s)",
695                    group_id,
696                    existing.media_uri[:50] if existing.media_uri else "N/A",
697                    media_uri[:50] if media_uri else "N/A",
698                )
699                await existing.stop()
700
701            # Create new shared stream
702            logger.info(
703                "[GroupStream] Creating new shared stream for group %s, uri=%s",
704                group_id,
705                media_uri[:80] if media_uri else "N/A",
706            )
707            stream = SharedGroupStream(group_id, media_uri)
708            await stream.start(audio_chunks)
709            self._shared_streams[group_id] = stream
710
711            return stream
712
713    def remove_shared_stream(self, group_id: str) -> None:
714        """Remove and cleanup shared stream for a group."""
715        if stream := self._shared_streams.pop(group_id, None):
716            logger.info("[GroupStream] Removed shared stream for group %s", group_id)
717            task = self.mass.create_task(stream.stop())
718            self._background_tasks.add(task)
719            task.add_done_callback(self._background_tasks.discard)
720
721    async def get_ma_stream_url(self, player_id: str, media: Any) -> str | None:
722        """
723        Resolve the direct MA Streamserver URL for the given media.
724
725        Used by redirect stream mode: the TV fetches audio straight from the
726        MA Streamserver, which applies the player's own codec config and DSP —
727        no local proxy/ffmpeg involved.
728
729        :param player_id: The MSX player requesting the stream.
730        :param media: PlayerMedia to resolve the stream URL for.
731        :return: Direct URL to the MA Streamserver, or None when resolution
732            fails (the caller falls back to the local proxy pipeline).
733        """
734        if not media:
735            logger.debug("[MARedirect] No media provided")
736            return None
737        try:
738            stream_url: str = await self.mass.streams.resolve_stream_url(player_id, media)
739        except Exception as err:
740            logger.warning("[MARedirect] Failed to resolve MA stream URL: %s", err, exc_info=True)
741            return None
742        # MA returns a flow URL (continuous whole-queue stream) when e.g. crossfade
743        # is enabled and the player lacks gapless support. That breaks the MSX
744        # per-track model (progress display, auto-advance re-enqueue), so serve
745        # such tracks through the local per-track proxy instead.
746        if "/flow/" in stream_url:
747            logger.debug(
748                "[MARedirect] Flow-mode URL not usable for MSX per-track playback, "
749                "falling back to proxy: %s",
750                stream_url,
751            )
752            return None
753        logger.debug("[MARedirect] Resolved MA stream URL: %s", stream_url)
754        return stream_url
755
756    async def cleanup_shared_streams(self) -> None:
757        """Cleanup all shared streams (called on unload)."""
758        for group_id, stream in list(self._shared_streams.items()):
759            logger.debug("[GroupStream] Cleaning up stream for group %s", group_id)
760            await stream.stop()
761        self._shared_streams.clear()
762
763    def _player_display_name_from_id(
764        self, player_id: str, prefix_label: str = "MSX TV", remote_ip: str | None = None
765    ) -> str:
766        """Build a unique display name from player_id for the MA UI."""
767        prefix = MSX_PLAYER_ID_PREFIX
768        suffix = player_id.removeprefix(prefix)
769        if not suffix:
770            return prefix_label
771        # IP-based: msx_192_168_10_15 → "MSX TV (192.168.10.15)"
772        if "_" in suffix:
773            parts = suffix.split("_")
774            if all(p.isdigit() for p in parts):
775                return f"{prefix_label} ({'.'.join(parts)})"
776        # UUID-based: msx_msx_bc93ce1d_491d_4d95_9430_2fbeabb5ce1b → "MSX TV (bc93)"
777        # Show only first 4 chars of UUID for readability, plus IP if available
778        if suffix.startswith("msx_") and len(suffix) > 12:
779            uuid_part = suffix[4:8]  # First 4 chars after "msx_"
780            if remote_ip:
781                return f"{prefix_label} ({uuid_part}) [{remote_ip}]"
782            return f"{prefix_label} ({uuid_part})"
783        # Fallback: truncate long suffixes
784        if len(suffix) > 12:
785            if remote_ip:
786                return f"{prefix_label} ({suffix[:8]}...) [{remote_ip}]"
787            return f"{prefix_label} ({suffix[:8]}...)"
788        if remote_ip:
789            return f"{prefix_label} ({suffix}) [{remote_ip}]"
790        return f"{prefix_label} ({suffix})"
791
792    async def _handle_player_unregister(self, player_id: str) -> None:
793        """Unregister a player with race-condition handling."""
794        self.logger.debug("Unregistering MSX player %s", player_id)
795        unregister_event = asyncio.Event()
796        self._pending_unregisters[player_id] = unregister_event
797        try:
798            if self.bridge_manager:
799                await self.bridge_manager.remove_bridge(player_id, permanent=True)
800            await self.mass.players.unregister(player_id)
801        finally:
802            self._pending_unregisters.pop(player_id, None)
803            self._player_last_activity.pop(player_id, None)
804            unregister_event.set()
805
806    async def _run_idle_timeout_loop(self) -> None:
807        """Background task: unregister players idle longer than configured timeout."""
808        timeout_minutes = max(
809            1,
810            min(
811                1440,
812                int(
813                    cast(
814                        "int",
815                        self.config.get_value(
816                            CONF_PLAYER_IDLE_TIMEOUT, DEFAULT_PLAYER_IDLE_TIMEOUT
817                        ),
818                    )
819                ),
820            ),
821        )
822        interval_seconds = 60
823        while not self.mass.closing:
824            try:
825                await asyncio.sleep(interval_seconds)
826            except asyncio.CancelledError:
827                break
828            now = time.monotonic()
829            cutoff = now - (timeout_minutes * 60)
830            for player in list(self.players):
831                if not isinstance(player, MSXPlayer):
832                    continue
833                last = self._player_last_activity.get(player.player_id, 0)
834                if last > 0 and last < cutoff:
835                    self.logger.info(
836                        "Unregistering idle MSX player %s (no activity for %d min)",
837                        player.player_id,
838                        timeout_minutes,
839                    )
840                    task = self.mass.create_task(self._handle_player_unregister(player.player_id))
841                    self._background_tasks.add(task)
842                    task.add_done_callback(self._background_tasks.discard)
843
844    def _make_bridge_manager(self) -> MSXSendspinBridgeManager:
845        """Import and construct the Sendspin bridge manager (raises ImportError if absent)."""
846        from .sendspin_bridge import MSXSendspinBridgeManager  # noqa: PLC0415
847
848        return MSXSendspinBridgeManager(self)
849