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