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