/
/
/
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