/
/
1"""MSX Player implementation."""
2
3from __future__ import annotations
4
5import asyncio
6import time
7from typing import TYPE_CHECKING, Any, cast
8
9from music_assistant_models.enums import PlaybackState, PlayerFeature, PlayerType
10from music_assistant_models.errors import PlayerUnavailableError
11from music_assistant_models.player import DeviceInfo
12
13from music_assistant.constants import CONF_ENTRY_OUTPUT_CODEC_DEFAULT_MP3
14from music_assistant.models.player import Player, PlayerMedia
15
16if TYPE_CHECKING:
17 from music_assistant_models.config_entries import ConfigEntry
18
19 from .provider import MSXBridgeProvider
20
21
22class MSXPlayer(Player):
23 """Represents a Smart TV running MSX as a Music Assistant player."""
24
25 current_stream_url: str | None = None
26 output_format: str = "mp3"
27 _skip_ws_notify: bool = False
28 _propagating: bool = False
29 _playing_from_queue: bool = False
30 _queue_source_id: str | None = None
31 _playlist_offset: int = 0
32 _playlist_size: int = 0
33 _media_ready: asyncio.Event
34 _attr_elapsed_time: float | None = None
35 _attr_elapsed_time_last_updated: float | None = None
36 _last_ws_position: float | None = None
37 _ws_ever_connected: bool = False
38
39 def __init__(
40 self,
41 provider: MSXBridgeProvider,
42 player_id: str,
43 name: str = "MSX TV",
44 output_format: str = "mp3",
45 *,
46 grouping_enabled: bool = True,
47 ip_address: str | None = None,
48 ) -> None:
49 """Initialize the MSX Player."""
50 super().__init__(provider, player_id)
51 self._attr_name = name
52 self._attr_type = PlayerType.PLAYER
53 self._attr_supported_features = {
54 PlayerFeature.PLAY_MEDIA,
55 PlayerFeature.PAUSE,
56 PlayerFeature.SEEK,
57 PlayerFeature.VOLUME_SET,
58 }
59 if grouping_enabled:
60 self._attr_supported_features.add(PlayerFeature.SET_MEMBERS)
61 self._attr_can_group_with = {provider.instance_id}
62 self._attr_device_info = DeviceInfo(
63 model="Smart TV (MSX)",
64 manufacturer="MSX Bridge",
65 )
66 if ip_address:
67 self._attr_device_info.ip_address = ip_address
68 self._attr_available = True
69 self._attr_powered = True
70 self._attr_volume_level = 100
71 self.output_format = output_format
72 self._media_ready = asyncio.Event()
73
74 @property
75 def requires_flow_mode(self) -> bool:
76 """MSX plays individual tracks â flow mode breaks progress tracking."""
77 return False
78
79 @property
80 def needs_poll(self) -> bool:
81 """Return if the player needs to be polled for state updates."""
82 return True
83
84 @property
85 def poll_interval(self) -> int:
86 """Return poll interval in seconds."""
87 return 5 if self.playback_state == PlaybackState.PLAYING else 30
88
89 async def get_config_entries(self) -> list[ConfigEntry]:
90 """Return per-player config entries â codec is configurable per TV."""
91 return [CONF_ENTRY_OUTPUT_CODEC_DEFAULT_MP3]
92
93 def on_ws_connected(self) -> None:
94 """Mark player as available when a WebSocket client connects."""
95 self._ws_ever_connected = True
96 if not self._attr_available:
97 self._attr_available = True
98 self.update_state()
99
100 def on_ws_disconnected(self) -> None:
101 """
102 Mark player unavailable when last WebSocket client disconnects while playing.
103
104 If the player was playing when the TV dropped the WS connection,
105 mark it unavailable so MA reflects the actual state.
106 """
107 if self._attr_playback_state == PlaybackState.PLAYING:
108 self._attr_available = False
109 self.update_state()
110
111 async def play_media(self, media: PlayerMedia) -> None:
112 """Handle PLAY MEDIA command â store stream URL for the TV to fetch."""
113 self.logger.info("play_media on %s: uri=%s", self.display_name, media.uri)
114 self.current_stream_url = media.uri
115 self._attr_current_media = media
116 self._media_ready.set()
117 self._attr_playback_state = PlaybackState.PLAYING
118 self._attr_elapsed_time = 0.0
119 self._attr_elapsed_time_last_updated = time.time()
120 self._last_ws_position = None
121 self.update_state()
122
123 if not self._skip_ws_notify:
124 self._notify_msx_playback(media)
125
126 await self._propagate_to_group_members("play_media", media=media)
127
128 async def set_members(
129 self,
130 player_ids_to_add: list[str] | None = None,
131 player_ids_to_remove: list[str] | None = None,
132 ) -> None:
133 """Handle SET_MEMBERS â update group membership."""
134 for pid in player_ids_to_remove or []:
135 if pid in self._attr_group_members:
136 self._attr_group_members.remove(pid)
137 for pid in player_ids_to_add or []:
138 if pid != self.player_id and pid not in self._attr_group_members:
139 other = self.mass.players.get_player(pid)
140 if other and isinstance(other, MSXPlayer):
141 self._attr_group_members.append(pid)
142
143 # Normalize group membership: leader must be first when grouped,
144 # and the list must be empty when no other members exist.
145 members_except_self = [pid for pid in self._attr_group_members if pid != self.player_id]
146 if not members_except_self:
147 self._attr_group_members = []
148 else:
149 self._attr_group_members = [self.player_id, *members_except_self]
150
151 self.update_state()
152
153 async def play(self) -> None:
154 """Handle PLAY (resume) command."""
155 self.logger.info("play (resume) on %s", self.display_name)
156 if self._attr_playback_state == PlaybackState.PAUSED:
157 await self._resume_from_pause()
158 return
159 self._attr_playback_state = PlaybackState.PLAYING
160 self._attr_elapsed_time_last_updated = time.time()
161 self.update_state()
162 await self._propagate_to_group_members("play")
163
164 async def pause(self) -> None:
165 """Handle PAUSE command â pause playback on MSX, keep stream alive for resume."""
166 self.logger.info("pause on %s", self.display_name)
167 # Snapshot the elapsed time before pausing
168 if self._attr_elapsed_time is not None and self._attr_elapsed_time_last_updated is not None:
169 self._attr_elapsed_time += time.time() - self._attr_elapsed_time_last_updated
170 self._attr_playback_state = PlaybackState.PAUSED
171 self._attr_elapsed_time_last_updated = time.time()
172 self.update_state()
173 if not self._skip_ws_notify:
174 cast("MSXBridgeProvider", self.provider).notify_play_paused(self.player_id)
175 await self._propagate_to_group_members("pause")
176
177 async def stop(self) -> None:
178 """Handle STOP command."""
179 self.logger.info("stop on %s", self.display_name)
180 self._attr_playback_state = PlaybackState.IDLE
181 self._attr_current_media = None
182 self._attr_elapsed_time = None
183 self._attr_elapsed_time_last_updated = None
184 self._last_ws_position = None
185 self.current_stream_url = None
186 self._playing_from_queue = False
187 self._queue_source_id = None
188 self._playlist_offset = 0
189 self._playlist_size = 0
190 self.update_state()
191 provider = cast("MSXBridgeProvider", self.provider)
192 provider.notify_play_stopped(self.player_id)
193 await self._propagate_to_group_members("stop")
194
195 async def volume_set(self, volume_level: int) -> None:
196 """Handle VOLUME_SET command."""
197 self._attr_volume_level = volume_level
198 self.update_state()
199
200 async def seek(self, position_seconds: int) -> None:
201 """Handle SEEK command â send seek action to MSX player via WebSocket."""
202 self._attr_elapsed_time = float(position_seconds)
203 self._attr_elapsed_time_last_updated = time.time()
204 self._last_ws_position = None
205 self.update_state()
206 if not self._skip_ws_notify:
207 cast("MSXBridgeProvider", self.provider).notify_seek(self.player_id, position_seconds)
208
209 def update_position(self, position: float) -> None:
210 """
211 Update elapsed time from a WebSocket position report.
212
213 Only accepts updates while PLAYING â late reports arriving after
214 pause() would overwrite the correctly accumulated elapsed_time.
215 """
216 if self._attr_playback_state != PlaybackState.PLAYING:
217 return
218 normalized = max(0.0, float(position))
219 duration = self._served_duration()
220 if duration is not None:
221 normalized = min(normalized, duration)
222 self._attr_elapsed_time = normalized
223 # elapsed_time_last_updated is compared against time.time() by MA core
224 # (corrected_elapsed_time) â must stay wall-clock. The WS staleness
225 # marker is provider-internal â monotonic, immune to NTP steps.
226 self._attr_elapsed_time_last_updated = time.time()
227 self._last_ws_position = time.monotonic()
228 self.update_state()
229
230 async def poll(self) -> None:
231 """
232 Poll player for state updates.
233
234 Raises PlayerUnavailableError if the player was marked unavailable
235 (e.g. WS disconnected while playing â TV likely went offline).
236
237 If a recent WebSocket position report was received (within 10s),
238 skip wall-clock increment â the WS data is more accurate.
239 """
240 if not self._attr_available:
241 raise PlayerUnavailableError(
242 f"MSX TV {self.display_name} is offline (WebSocket disconnected)",
243 translation_key="player_offline",
244 translation_owner=self.translation_owner,
245 translation_args=[self.display_name],
246 )
247 if (
248 self._attr_playback_state == PlaybackState.PLAYING
249 and self._attr_elapsed_time is not None
250 and self._attr_elapsed_time_last_updated is not None
251 ):
252 # Skip wall-clock update if WS reported position recently
253 if self._last_ws_position and (time.monotonic() - self._last_ws_position) < 10:
254 return
255 now = time.time()
256 delta = now - self._attr_elapsed_time_last_updated
257 new_elapsed = max(0.0, float(self._attr_elapsed_time) + float(delta))
258 duration = self._served_duration()
259 if duration is not None:
260 new_elapsed = min(new_elapsed, duration)
261 self._attr_elapsed_time = new_elapsed
262 self._attr_elapsed_time_last_updated = now
263 self.update_state()
264
265 def expect_new_media(self) -> None:
266 """
267 Arm wait_for_media() to wait for the NEXT play_media() call.
268
269 Call this before initiating playback that will (asynchronously) invoke
270 play_media(). Without arming, wait_for_media() would return the stale
271 current_media left over from a previous track.
272 """
273 self._media_ready.clear()
274
275 async def wait_for_media(self, timeout: float = 10.0) -> PlayerMedia | None:
276 """
277 Wait for play_media() to set current_media, with timeout.
278
279 Fast path: current_media already set and not armed via expect_new_media()
280 â return immediately. Slow path: wait for the next play_media() to signal.
281 After stop(), _attr_current_media is None â this method returns None even
282 if the event happens to still be set.
283 """
284 if self._attr_current_media is not None and self._media_ready.is_set():
285 return self._attr_current_media
286 if not self._media_ready.is_set():
287 try:
288 await asyncio.wait_for(self._media_ready.wait(), timeout=timeout)
289 except TimeoutError:
290 return None
291 return self._attr_current_media
292
293 def _notify_msx_playback(self, media: PlayerMedia) -> None:
294 """Send WS notification to MSX about the new playback state."""
295 source_id = media.source_id
296 is_queue_backed = bool(source_id and media.queue_item_id)
297 is_same_queue = self._playing_from_queue and self._queue_source_id == source_id
298 provider = cast("MSXBridgeProvider", self.provider)
299
300 if is_queue_backed and is_same_queue and source_id:
301 self._notify_same_queue(provider, source_id)
302 elif is_queue_backed and source_id:
303 self._notify_new_queue(provider, source_id)
304 else:
305 # Queue-backed playback renders from the MSX native playlist, which carries
306 # its own per-track metadata; only standalone media needs it pushed here.
307 next_action = f"request:interaction:/api/next/{self.player_id}"
308 prev_action = f"request:interaction:/api/previous/{self.player_id}"
309 provider.notify_play_started(
310 self.player_id,
311 title=media.title,
312 artist=media.artist,
313 image_url=media.image_url,
314 duration=media.stream_duration or media.duration,
315 next_action=next_action,
316 prev_action=prev_action,
317 )
318
319 def _notify_same_queue(self, provider: MSXBridgeProvider, source_id: str) -> None:
320 """Handle same-queue playback: goto index or re-send if queue changed."""
321 queue = self.mass.player_queues.get(source_id)
322 ma_index = getattr(queue, "current_index", 0) if queue else 0
323 try:
324 current_size = len(self.mass.player_queues.items(source_id))
325 except Exception:
326 self.logger.debug("Failed to get queue size for %s", source_id, exc_info=True)
327 current_size = self._playlist_size
328 if current_size != self._playlist_size:
329 self._playlist_size = current_size
330 self._playlist_offset = ma_index
331 provider.notify_play_playlist(self.player_id, ma_index, queue_id=source_id)
332 else:
333 if self._playlist_size > 0:
334 msx_index = (ma_index - self._playlist_offset) % self._playlist_size
335 else:
336 msx_index = ma_index
337 provider.notify_goto_index(self.player_id, msx_index)
338
339 def _notify_new_queue(self, provider: MSXBridgeProvider, source_id: str) -> None:
340 """Send full MSX native playlist for a new queue."""
341 queue = self.mass.player_queues.get(source_id)
342 start_index = getattr(queue, "current_index", 0) if queue else 0
343 try:
344 self._playlist_size = len(self.mass.player_queues.items(source_id))
345 except Exception:
346 self.logger.debug("Failed to get queue size for %s", source_id, exc_info=True)
347 self._playlist_size = 0
348 self._playlist_offset = start_index
349 self._queue_source_id = source_id
350 provider.notify_play_playlist(self.player_id, start_index, queue_id=source_id)
351 self._playing_from_queue = True
352
353 def _served_duration(self) -> float | None:
354 """
355 Return the length in seconds of the audio served to the TV, if known.
356
357 The TV reports its position within that audio, which is shorter than the
358 media item itself when playback starts at a seek position.
359 """
360 if (media := self._attr_current_media) is None:
361 return None
362 duration = media.stream_duration or media.duration
363 if not isinstance(duration, (int, float)) or duration <= 0:
364 return None
365 return float(duration)
366
367 def _get_group_member_ids(self) -> list[str]:
368 """
369 Get IDs of group members (excluding self).
370
371 Only returns members when this player is the sync leader.
372 MA's SyncGroupPlayer forwards play_media to the sync leader,
373 whose group_members already contains all SyncGroup members.
374 """
375 if self.synced_to is not None:
376 return []
377 return [x for x in self.group_members if x != self.player_id]
378
379 async def _propagate_to_group_members(self, command: str, **kwargs: Any) -> None:
380 """Propagate command to group members in parallel when we are the leader."""
381 # Skip if grouping is disabled at provider level
382 provider = cast("MSXBridgeProvider", self.provider)
383 if not provider.grouping_enabled:
384 return
385 # Prevent infinite recursion if member.play_media triggers propagation back
386 if self._propagating:
387 return
388 self._propagating = True
389 try:
390 tasks: list[asyncio.Task[None]] = []
391 for member_id in self._get_group_member_ids():
392 member = self.mass.players.get_player(member_id)
393 if not member or not isinstance(member, MSXPlayer) or not member.available:
394 continue
395 tasks.append(asyncio.create_task(self._propagate_single(member, command, **kwargs)))
396 if tasks:
397 await asyncio.gather(*tasks, return_exceptions=True)
398 finally:
399 self._propagating = False
400
401 async def _propagate_single(self, member: MSXPlayer, command: str, **kwargs: Any) -> None:
402 """Propagate a single command to one group member."""
403 try:
404 if command == "play_media":
405 media = kwargs.get("media")
406 if media:
407 # Call member.play_media directly â mass.players.play_media
408 # would redirect synced/grouped players back to the leader
409 await member.play_media(media)
410 elif command == "stop":
411 await member.stop()
412 elif command == "pause":
413 await member.pause()
414 elif command == "play":
415 await member.play()
416 except Exception:
417 self.logger.warning(
418 "Failed to propagate %s to member %s",
419 command,
420 member.player_id,
421 exc_info=True,
422 )
423
424 async def _resume_from_pause(self) -> None:
425 """
426 Resume playback after pause â tell MSX to unpause its native player.
427
428 Note: the HTTP audio stream stays open during pause. For short pauses
429 the chunk buffer (maxsize=32) absorbs the gap. Long pauses (minutes)
430 may cause stream starvation â ffmpeg backs up, and MSX may get silence
431 or a playback error on resume. A reconnect mechanism would be needed
432 for reliable long-pause support.
433 """
434 self._attr_playback_state = PlaybackState.PLAYING
435 self._attr_elapsed_time_last_updated = time.time()
436 self._last_ws_position = None
437 self.update_state()
438 if not self._skip_ws_notify:
439 cast("MSXBridgeProvider", self.provider).notify_play_resumed(self.player_id)
440 await self._propagate_to_group_members("play")
441