/
/
1"""Snapcast Player."""
2
3from __future__ import annotations
4
5import asyncio
6from contextlib import suppress
7from typing import TYPE_CHECKING, TypedDict, cast
8
9from music_assistant_models.config_entries import ConfigEntry
10from music_assistant_models.enums import IdentifierType, PlaybackState, PlayerFeature
11from music_assistant_models.player import DeviceInfo, PlayerMedia
12from propcache import under_cached_property as cached_property
13
14from music_assistant.constants import ATTR_ANNOUNCEMENT_IN_PROGRESS, CONF_ENTRY_HTTP_PROFILE_HIDDEN
15from music_assistant.helpers.util import is_valid_mac_address
16from music_assistant.models.player import Player
17from music_assistant.providers.snapcast.constants import (
18 SNAPCLIENT_LIVENESS_POLL_INTERVAL,
19 SNAPCLIENT_STALE_THRESHOLD,
20)
21from music_assistant.providers.snapcast.ma_stream import SnapcastMAStream
22from music_assistant.providers.sync_group.constants import SGP_PREFIX
23
24if TYPE_CHECKING:
25 from music_assistant_models.config_entries import ConfigEntry
26
27 from music_assistant.providers.snapcast.provider import SnapCastProvider
28 from music_assistant.providers.snapcast.snap_cntrl_proto import SnapclientProto, SnapstreamProto
29
30
31class TrackedPlayerState(TypedDict, total=False):
32 """
33 Tracked state for the Snapcast MA player.
34
35 It is used for change detection and state synchronization, and may be
36 partially populated depending on which information is
37 currently available.
38
39 Keys prefixed with ``_attr_`` are exposed as player attributes, while the
40 remaining keys represent internal Snapcast grouping and connection state.
41 """
42
43 # Player attribute fields
44 _attr_name: str
45 _attr_volume_level: float
46 _attr_volume_muted: bool
47 _attr_available: bool
48
49 # snapclient fields
50 connected: bool
51 stream_id: str
52 stream_status: str | None
53 grp_name: str
54 grp_member_ids: list[str]
55 grp_member_avail: list[bool]
56
57
58class SnapCastPlayer(Player):
59 """SnapCastPlayer."""
60
61 def __init__(
62 self,
63 provider: SnapCastProvider,
64 player_id: str,
65 snap_client: SnapclientProto,
66 ) -> None:
67 """Init."""
68 self.snap_client = snap_client
69 super().__init__(provider, player_id)
70
71 # Snapcast stream format is fixed for a provider instance (from advanced settings)
72 stream_format = provider.stream_audio_format
73 self._attr_supported_sample_rates = [(stream_format.sample_rate, stream_format.bit_depth)]
74
75 self._snap_ma_stream: SnapcastMAStream | None = None
76
77 self._update_worker: asyncio.Task[None] | None = None
78 self._poke_evt = asyncio.Event()
79 self._state_update_lock = asyncio.Lock()
80 self._last_tracked_state: TrackedPlayerState | None = None
81 # raw lastSeen value + monotonic time it last advanced (skew-proof liveness)
82 self._last_seen_raw: tuple[int, int] | None = None
83 self._last_seen_changed: float = 0.0
84
85 @property
86 def snap_provider(self) -> SnapCastProvider:
87 """Return the Snapcast provider instance."""
88 return cast("SnapCastProvider", self.provider)
89
90 @property
91 def requires_flow_mode(self) -> bool:
92 """Return if the player requires flow mode."""
93 return True
94
95 @property
96 def synced_to(self) -> str | None:
97 """Return the id of the player this player is synced to (sync leader)."""
98 grp_name = self.snap_group_name
99 if grp_name == self.player_id:
100 # is group leader
101 return None
102
103 grp_player_ids = self._get_player_ids_of_curr_group()
104 if len(grp_player_ids) < 2 or grp_name not in grp_player_ids:
105 return None
106
107 if leader_player := self.mass.players.get_player(grp_name):
108 return grp_name if leader_player.available else None
109
110 return None
111
112 @cached_property
113 def group_members(self) -> list[str]:
114 """Return the group members of the player."""
115 if not self._attr_available:
116 return []
117
118 grp_name = self.snap_group_name
119 if grp_name != self.player_id:
120 # only group leaders can have members
121 return []
122
123 player_ids = self._get_player_ids_of_curr_group()
124 if self.player_id not in player_ids:
125 # should not happen, unless the current
126 # state repr is invalid
127 return []
128
129 player_ids.remove(self.player_id)
130 connected = [
131 player_id
132 for player_id in player_ids
133 if (client := self.snap_provider.get_snap_client(player_id=player_id))
134 and client.connected
135 ]
136 if connected:
137 return [self.player_id, *connected]
138
139 return []
140
141 @property
142 def playback_state(self) -> PlaybackState:
143 """Return the current playback state of the player."""
144 snap_stream = self._get_active_snapstream()
145 if snap_stream is None:
146 return PlaybackState.IDLE
147
148 if snap_stream.identifier == "default" or snap_stream.status == "idle":
149 return PlaybackState.IDLE
150
151 return PlaybackState.PLAYING
152
153 @property
154 def elapsed_time(self) -> float | None:
155 """Return the elapsed time in (fractional) seconds of the current track (if any)."""
156 # using flow-mode, elapsed time will be estimated upstream from 'elapsed_time_last_updated'
157 return 0 if self.active_snap_ma_stream else None
158
159 @property
160 def elapsed_time_last_updated(self) -> float | None:
161 """
162 Return when the elapsed time was last updated.
163
164 return: The (UTC) timestamp when the elapsed time was last updated,
165 or None if it was never updated (or unknown).
166 """
167 # we only update on playback starts
168 if snap_ma_stream := self.active_snap_ma_stream:
169 return snap_ma_stream.playback_started_at
170 return None
171
172 def setup(self) -> None:
173 """Set up player."""
174 self._attr_name = self.snap_client.friendly_name
175 self._attr_available = self.snap_client.connected
176
177 host_dict = self.snap_client._client.get("host", {})
178 os, arch, ip, mac = (host_dict.get(key, "") for key in ["os", "arch", "ip", "mac"])
179 self._attr_device_info = DeviceInfo(
180 model=os,
181 manufacturer=arch,
182 )
183 if ip and (host := self.snap_client._client.get("host")):
184 self._attr_device_info.add_identifier(IdentifierType.IP_ADDRESS, host.get("ip"))
185 # Only add MAC address if it's valid (not 00:00:00:00:00:00)
186 if mac and is_valid_mac_address(mac):
187 self._attr_device_info.add_identifier(IdentifierType.MAC_ADDRESS, mac)
188 self._attr_supported_features = {
189 PlayerFeature.PLAY_MEDIA,
190 PlayerFeature.SET_MEMBERS,
191 PlayerFeature.VOLUME_SET,
192 PlayerFeature.VOLUME_MUTE,
193 PlayerFeature.PLAY_ANNOUNCEMENT,
194 }
195 self._attr_can_group_with = {self.snap_provider.instance_id}
196 # poll to detect abruptly powered-off clients (see _is_alive)
197 self._attr_needs_poll = True
198 self._attr_poll_interval = SNAPCLIENT_LIVENESS_POLL_INTERVAL
199 if not self._update_worker:
200 self._update_worker = self.mass.create_task(self._player_update_worker)
201
202 async def poll(self) -> None:
203 """Poll the snapserver so abruptly powered-off clients are detected."""
204 await self.snap_provider.refresh_server_status()
205 self.poke_player_update()
206
207 async def volume_set(self, volume_level: int) -> None:
208 """Send VOLUME_SET command to given player."""
209 # Use optimistic server state for now
210 # not guaranteed that the client respects it
211 await self.snap_client.set_volume(volume_level)
212
213 async def stop(self) -> None:
214 """Send STOP command to given player."""
215 player_group = await self.snap_provider.ensure_player_owned_group(self.player_id)
216 assert player_group is not None # for type checking
217 await player_group.set_stream("default")
218 if ma_stream := self.active_snap_ma_stream:
219 ma_stream.request_stop_stream()
220 return
221
222 self.poke_player_update()
223
224 async def volume_mute(self, muted: bool) -> None:
225 """Send MUTE command to given player."""
226 # Use optimistic server state for now
227 # not guaranteed that the client respects it
228 # TODO: move this to the snapcast python library
229 vol = self.snap_client._client["config"]["volume"]
230 vol["muted"] = muted
231 res = await self.snap_provider._snapserver.client_volume(self.snap_client.identifier, vol)
232 if res and "muted" in res:
233 self.snap_client._client["config"]["volume"] = res
234 self.snap_client.callback()
235
236 async def set_members(
237 self,
238 player_ids_to_add: list[str] | None = None,
239 player_ids_to_remove: list[str] | None = None,
240 ) -> None:
241 """Handle SET_MEMBERS command on the player."""
242 # get the group owned by this player (identified by the group name)
243 player_group = await self.snap_provider.ensure_player_owned_group(self.player_id)
244
245 if player_group is None:
246 return
247
248 player_group.set_callback(None)
249
250 curr_ma_player_ids = [
251 ma_id
252 for cli_id in player_group.clients
253 if (ma_id := self.snap_provider._get_ma_id(cli_id))
254 ]
255
256 curr_stream_id = player_group.stream
257 sync_group_player: Player | None = None
258 if curr_ma_stream := self.snap_provider.get_snap_ma_stream(curr_stream_id):
259 media = curr_ma_stream.media
260 media_src_id = media.source_id or ""
261 if media_src_id.startswith(SGP_PREFIX):
262 sync_group_player = self.mass.players.get_player(media_src_id)
263 if sync_group_player and self.player_id in (player_ids_to_remove or []):
264 # players in sync_group_player.group_members will be rejoined
265 # remove others first
266 for id_to_remove in player_ids_to_remove or []:
267 if id_to_remove == self.player_id:
268 continue
269 if (
270 id_to_remove in curr_ma_player_ids
271 and id_to_remove not in sync_group_player.group_members
272 ):
273 await self.snap_provider.isolate_player_to_dedicated_group(
274 id_to_remove, target_stream_id="default"
275 )
276
277 # split remaining group into individual groups,
278 # keeps the current stream, set this group to default stream
279 await self.snap_provider.isolate_player_to_dedicated_group(
280 target_player_id=self.player_id,
281 target_stream_id="default",
282 others_stream_id=curr_stream_id,
283 )
284 else:
285 for player_id in player_ids_to_remove or []:
286 if player_id not in curr_ma_player_ids:
287 continue
288 await self.snap_provider.isolate_player_to_dedicated_group(
289 player_id, target_stream_id="default"
290 )
291 curr_ma_player_ids.remove(player_id)
292
293 for ma_id in player_ids_to_add or []:
294 if (
295 snap_id := self.snap_provider._get_snapclient_id(ma_id)
296 ) and ma_id not in curr_ma_player_ids:
297 await player_group.add_client(snap_id)
298
299 # some caller require instant state updates before returning
300 async with self._state_update_lock:
301 if await self._process_snapcast_client_state():
302 self.update_state()
303
304 self.snap_provider._update_group_callbacks(poke=True)
305
306 async def play_media(self, media: PlayerMedia) -> None:
307 """Handle PLAY MEDIA on given player."""
308 if self.synced_to:
309 msg = "A synced player cannot receive play commands directly"
310 raise RuntimeError(msg)
311
312 ma_stream = await self.snap_provider.get_snapcast_media_stream(
313 media, filter_settings_owner=self.player_id
314 )
315
316 if ma_stream is None or ma_stream.stream_id is None:
317 return
318
319 self._snap_ma_stream = ma_stream
320
321 # e.g. DSP settings require a restart
322 await self._snap_ma_stream.start_stream(allow_restart=True)
323
324 # if no announcement is playing we activate the stream now, otherwise it
325 # will be activated by play_announcement when the announcement is over.
326 if not self.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS):
327 player_group = await self.snap_provider.ensure_player_owned_group(self.player_id)
328 assert player_group is not None # for type checking
329 await player_group.set_stream(ma_stream.stream_id)
330
331 self.poke_player_update()
332
333 async def play_announcement(
334 self, announcement: PlayerMedia, volume_level: int | None = None
335 ) -> None:
336 """Handle (provider native) playback of an announcement on given player."""
337 was_synced_to: str | None = self.synced_to
338 orig_volume_level: int | None = self.volume_level
339
340 prev_stream = self.active_snap_ma_stream
341
342 ma_stream = await self.snap_provider.get_snapcast_media_stream(
343 announcement, filter_settings_owner=self.player_id
344 )
345 player_group = await self.snap_provider.ensure_player_owned_group(self.player_id)
346
347 if ma_stream is None or ma_stream.stream_id is None or player_group is None:
348 return
349
350 await player_group.set_stream(ma_stream.stream_id)
351
352 if self.snap_provider._use_builtin_server:
353 await asyncio.sleep(self.snap_provider._snapcast_server_buffer_size / 1000.0)
354
355 if volume_level is not None:
356 await self.volume_set(volume_level)
357
358 await ma_stream.start_stream()
359 await ma_stream.wait_for_stopped()
360
361 if self.volume_level == volume_level and orig_volume_level is not None:
362 await self.volume_set(orig_volume_level)
363
364 if was_synced_to:
365 if (
366 leader_group := await self.snap_provider.ensure_player_owned_group(was_synced_to)
367 ) is None:
368 return
369 await leader_group.add_client(self.snap_client.identifier)
370 else:
371 await player_group.set_stream(
372 prev_stream.stream_id
373 if prev_stream and prev_stream.stream_id is not None
374 else "default"
375 )
376
377 async def get_config_entries(self) -> list[ConfigEntry]:
378 """Player config."""
379 return [
380 # we don't use the http server for streaming
381 CONF_ENTRY_HTTP_PROFILE_HIDDEN,
382 ]
383
384 def _handle_player_update(self, snap_client: SnapclientProto) -> None:
385 """Forward snap_client updates."""
386 self.poke_player_update()
387
388 def poke_player_update(self) -> None:
389 """Signal that a player state update should be processed."""
390 self._poke_evt.set()
391
392 async def _player_update_worker(self) -> None:
393 """Aggregate and process player state update requests."""
394 while True:
395 await self._poke_evt.wait()
396 self._poke_evt.clear()
397 while True:
398 call_update: bool = False
399 try:
400 async with self._state_update_lock:
401 call_update = await self._process_snapcast_client_state()
402 if call_update:
403 self.update_state()
404 except KeyError, AttributeError, TypeError, ValueError:
405 # a failed update must not kill this worker (state would freeze)
406 self.logger.exception(
407 "Error while processing state update for player %s", self.player_id
408 )
409 break
410 if self._poke_evt.is_set():
411 self._poke_evt.clear()
412 continue
413 break
414
415 async def _process_snapcast_client_state(self) -> bool:
416 """
417 Process the latest Snapcast client state and apply changes to this player.
418
419 Returns:
420 True if changes were applied and a state update should be emitted via
421 ``update_state()``; False if no update is necessary (or if required data
422 is temporarily unavailable and the update should be retried later).
423 """
424 snap_group = self.snap_client.group
425 if snap_group is None:
426 # some data syncing error, a client is always a group member
427 # retry again later, don't call update now
428 return False
429
430 stream_id = snap_group.stream
431 snap_stream: SnapstreamProto | None = None
432 with suppress(KeyError):
433 snap_stream = self.snap_provider._snapserver.stream(stream_id)
434
435 members = list(snap_group.clients) # snapshot
436
437 curr_state: TrackedPlayerState = {
438 "_attr_name": self.snap_client.friendly_name,
439 "_attr_volume_level": self.snap_client.volume,
440 "_attr_volume_muted": self.snap_client.muted,
441 "_attr_available": self._is_alive(),
442 "connected": self.snap_client.connected,
443 "stream_id": snap_group.stream,
444 "stream_status": snap_stream.status if snap_stream is not None else None,
445 "grp_name": snap_group.name,
446 "grp_member_ids": members,
447 "grp_member_avail": [
448 pl.available
449 for cl_id in members
450 if (pl_id := self.snap_provider._get_ma_id(cl_id))
451 and (pl := self.mass.players.get_player(pl_id))
452 ],
453 }
454
455 prev_state: TrackedPlayerState = (
456 self._last_tracked_state if self._last_tracked_state is not None else {}
457 )
458 self._last_tracked_state = curr_state
459
460 # change detection for simple attrs
461 changed_attrs = {
462 k: v for k, v in curr_state.items() if k.startswith("_attr_") and prev_state.get(k) != v
463 }
464
465 prev_connected = prev_state.get("connected", False)
466 now_connected = curr_state.get("connected", False)
467 connection_changed = prev_connected != now_connected
468
469 prev_stream_id = prev_state.get("stream_id")
470 curr_stream_id = curr_state["stream_id"]
471 prev_stream_status = prev_state.get("stream_status")
472 curr_stream_status = curr_state.get("stream_status")
473
474 stream_changed = (
475 prev_stream_id != curr_stream_id or prev_stream_status != curr_stream_status
476 )
477
478 grouping_changed = any(
479 prev_state.get(k) != curr_state.get(k)
480 for k in ("grp_name", "grp_member_ids", "grp_member_avail")
481 )
482
483 needs_processing = bool(
484 changed_attrs or grouping_changed or stream_changed or connection_changed
485 )
486 if not needs_processing:
487 return False
488
489 if connection_changed or grouping_changed:
490 self.snap_provider.poke_group_members(snap_group)
491
492 # help cleaning up unused streams
493 if curr_stream_id == "default" or (
494 (my_stream := self._snap_ma_stream)
495 and my_stream.stream_id in {prev_stream_id, curr_stream_id}
496 ):
497 self.snap_provider.update_stream_usage()
498
499 # apply changed attrs
500 for key, value in changed_attrs.items():
501 setattr(self, key, value)
502
503 # finally notify state update once
504 return True
505
506 @property
507 def active_snap_ma_stream(self) -> SnapcastMAStream | None:
508 """Return the MA stream source of the active group."""
509 grp = self.snap_client.group
510 if grp is None or grp.stream is None:
511 return None
512
513 if grp.stream == "default":
514 return None
515
516 return self.snap_provider.get_snap_ma_stream(grp.stream)
517
518 @property
519 def snap_group_name(self) -> str:
520 """Return the name of the active group."""
521 snap_group = self.snap_client.group
522 if snap_group is None:
523 return ""
524 return snap_group.name
525
526 @property
527 def current_media(self) -> PlayerMedia | None:
528 """Return the current media being played by the player."""
529 if snap_ma_stream := self.active_snap_ma_stream:
530 return snap_ma_stream.media
531 return None
532
533 @property
534 def active_source(self) -> str | None:
535 """Return the (id of) the active source of the player."""
536 grp = self.snap_client.group
537 if grp is None or grp.stream is None:
538 return None
539
540 if grp.stream == "default":
541 return None
542
543 if ma_stream := self.snap_provider.get_snap_ma_stream(grp.stream):
544 return ma_stream.source_id
545
546 # external snapcast stream
547 return grp.stream or None
548
549 def _is_alive(self) -> bool:
550 """
551 Return whether the snapclient is connected and recently seen.
552
553 The snapserver advances lastSeen on every Time message a client sends
554 (~1/s while alive), but never marks an abruptly powered-off client
555 disconnected. We treat a connected client whose lastSeen has been frozen
556 longer than the threshold as gone.
557 """
558 now = self.mass.loop.time()
559 last_seen = self.snap_client._client.get("lastSeen")
560 ls = (last_seen["sec"], last_seen["usec"]) if last_seen else None
561 if ls != self._last_seen_raw:
562 self._last_seen_raw = ls
563 self._last_seen_changed = now
564 if not self.snap_client.connected:
565 return False
566 return (now - self._last_seen_changed) < SNAPCLIENT_STALE_THRESHOLD
567
568 def _get_active_snapstream(self) -> SnapstreamProto | None:
569 """Get active stream for given player_id."""
570 if group := self.snap_client.group:
571 with suppress(KeyError):
572 return self.snap_provider._snapserver.stream(group.stream)
573 return None
574
575 def _get_player_ids_of_curr_group(self) -> list[str]:
576 snap_group = self.snap_client.group
577 if snap_group is None:
578 return []
579 return [
580 ma_id
581 for client_id in snap_group.clients
582 if (ma_id := self.snap_provider._get_ma_id(client_id))
583 ]
584
585 def _get_players_of_curr_group(self) -> list[Player]:
586 return [
587 ma_player
588 for ma_id in self._get_player_ids_of_curr_group()
589 if (ma_player := self.mass.players.get_player(ma_id))
590 ]
591