/
/
/
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 # pin the music stream, else the idle timer stops it before the announcement ends
342 if prev_stream:
343 prev_stream.pin(self.player_id)
344
345 try:
346 ma_stream = await self.snap_provider.get_snapcast_media_stream(
347 announcement, filter_settings_owner=self.player_id
348 )
349 player_group = await self.snap_provider.ensure_player_owned_group(self.player_id)
350
351 if ma_stream is None or ma_stream.stream_id is None or player_group is None:
352 return
353
354 await player_group.set_stream(ma_stream.stream_id)
355
356 if self.snap_provider._use_builtin_server:
357 await asyncio.sleep(self.snap_provider._snapcast_server_buffer_size / 1000.0)
358
359 if volume_level is not None:
360 await self.volume_set(volume_level)
361
362 await ma_stream.start_stream()
363 await ma_stream.wait_for_stopped()
364
365 if self.volume_level == volume_level and orig_volume_level is not None:
366 await self.volume_set(orig_volume_level)
367
368 if was_synced_to:
369 if (
370 leader_group := await self.snap_provider.ensure_player_owned_group(
371 was_synced_to
372 )
373 ) is None:
374 return
375 await leader_group.add_client(self.snap_client.identifier)
376 else:
377 await player_group.set_stream(
378 prev_stream.stream_id
379 if prev_stream and prev_stream.stream_id is not None
380 else "default"
381 )
382 finally:
383 if prev_stream:
384 prev_stream.unpin(self.player_id)
385 self.snap_provider.update_stream_usage()
386
387 async def get_config_entries(self) -> list[ConfigEntry]:
388 """Player config."""
389 return [
390 # we don't use the http server for streaming
391 CONF_ENTRY_HTTP_PROFILE_HIDDEN,
392 ]
393
394 def _handle_player_update(self, snap_client: SnapclientProto) -> None:
395 """Forward snap_client updates."""
396 self.poke_player_update()
397
398 def poke_player_update(self) -> None:
399 """Signal that a player state update should be processed."""
400 self._poke_evt.set()
401
402 async def _player_update_worker(self) -> None:
403 """Aggregate and process player state update requests."""
404 while True:
405 await self._poke_evt.wait()
406 self._poke_evt.clear()
407 while True:
408 call_update: bool = False
409 try:
410 async with self._state_update_lock:
411 call_update = await self._process_snapcast_client_state()
412 if call_update:
413 self.update_state()
414 except KeyError, AttributeError, TypeError, ValueError:
415 # a failed update must not kill this worker (state would freeze)
416 self.logger.exception(
417 "Error while processing state update for player %s", self.player_id
418 )
419 break
420 if self._poke_evt.is_set():
421 self._poke_evt.clear()
422 continue
423 break
424
425 async def _process_snapcast_client_state(self) -> bool:
426 """
427 Process the latest Snapcast client state and apply changes to this player.
428
429 Returns:
430 True if changes were applied and a state update should be emitted via
431 ``update_state()``; False if no update is necessary (or if required data
432 is temporarily unavailable and the update should be retried later).
433 """
434 snap_group = self.snap_client.group
435 if snap_group is None:
436 # some data syncing error, a client is always a group member
437 # retry again later, don't call update now
438 return False
439
440 stream_id = snap_group.stream
441 snap_stream: SnapstreamProto | None = None
442 with suppress(KeyError):
443 snap_stream = self.snap_provider._snapserver.stream(stream_id)
444
445 members = list(snap_group.clients) # snapshot
446
447 curr_state: TrackedPlayerState = {
448 "_attr_name": self.snap_client.friendly_name,
449 "_attr_volume_level": self.snap_client.volume,
450 "_attr_volume_muted": self.snap_client.muted,
451 "_attr_available": self._is_alive(),
452 "connected": self.snap_client.connected,
453 "stream_id": snap_group.stream,
454 "stream_status": snap_stream.status if snap_stream is not None else None,
455 "grp_name": snap_group.name,
456 "grp_member_ids": members,
457 "grp_member_avail": [
458 pl.available
459 for cl_id in members
460 if (pl_id := self.snap_provider._get_ma_id(cl_id))
461 and (pl := self.mass.players.get_player(pl_id))
462 ],
463 }
464
465 prev_state: TrackedPlayerState = (
466 self._last_tracked_state if self._last_tracked_state is not None else {}
467 )
468 self._last_tracked_state = curr_state
469
470 # change detection for simple attrs
471 changed_attrs = {
472 k: v for k, v in curr_state.items() if k.startswith("_attr_") and prev_state.get(k) != v
473 }
474
475 prev_connected = prev_state.get("connected", False)
476 now_connected = curr_state.get("connected", False)
477 connection_changed = prev_connected != now_connected
478
479 prev_stream_id = prev_state.get("stream_id")
480 curr_stream_id = curr_state["stream_id"]
481 prev_stream_status = prev_state.get("stream_status")
482 curr_stream_status = curr_state.get("stream_status")
483
484 stream_changed = (
485 prev_stream_id != curr_stream_id or prev_stream_status != curr_stream_status
486 )
487
488 grouping_changed = any(
489 prev_state.get(k) != curr_state.get(k)
490 for k in ("grp_name", "grp_member_ids", "grp_member_avail")
491 )
492
493 needs_processing = bool(
494 changed_attrs or grouping_changed or stream_changed or connection_changed
495 )
496 if not needs_processing:
497 return False
498
499 if connection_changed or grouping_changed:
500 self.snap_provider.poke_group_members(snap_group)
501
502 # help cleaning up unused streams
503 if curr_stream_id == "default" or (
504 (my_stream := self._snap_ma_stream)
505 and my_stream.stream_id in {prev_stream_id, curr_stream_id}
506 ):
507 self.snap_provider.update_stream_usage()
508
509 # apply changed attrs
510 for key, value in changed_attrs.items():
511 setattr(self, key, value)
512
513 # finally notify state update once
514 return True
515
516 @property
517 def active_snap_ma_stream(self) -> SnapcastMAStream | None:
518 """Return the MA stream source of the active group."""
519 grp = self.snap_client.group
520 if grp is None or grp.stream is None:
521 return None
522
523 if grp.stream == "default":
524 return None
525
526 return self.snap_provider.get_snap_ma_stream(grp.stream)
527
528 @property
529 def snap_group_name(self) -> str:
530 """Return the name of the active group."""
531 snap_group = self.snap_client.group
532 if snap_group is None:
533 return ""
534 return snap_group.name
535
536 @property
537 def current_media(self) -> PlayerMedia | None:
538 """Return the current media being played by the player."""
539 if snap_ma_stream := self.active_snap_ma_stream:
540 return snap_ma_stream.media
541 return None
542
543 @property
544 def active_source(self) -> str | None:
545 """Return the (id of) the active source of the player."""
546 grp = self.snap_client.group
547 if grp is None or grp.stream is None:
548 return None
549
550 if grp.stream == "default":
551 return None
552
553 if ma_stream := self.snap_provider.get_snap_ma_stream(grp.stream):
554 return ma_stream.source_id
555
556 # external snapcast stream
557 return grp.stream or None
558
559 def _is_alive(self) -> bool:
560 """
561 Return whether the snapclient is connected and recently seen.
562
563 The snapserver advances lastSeen on every Time message a client sends
564 (~1/s while alive), but never marks an abruptly powered-off client
565 disconnected. We treat a connected client whose lastSeen has been frozen
566 longer than the threshold as gone.
567 """
568 now = self.mass.loop.time()
569 last_seen = self.snap_client._client.get("lastSeen")
570 ls = (last_seen["sec"], last_seen["usec"]) if last_seen else None
571 if ls != self._last_seen_raw:
572 self._last_seen_raw = ls
573 self._last_seen_changed = now
574 if not self.snap_client.connected:
575 return False
576 return (now - self._last_seen_changed) < SNAPCLIENT_STALE_THRESHOLD
577
578 def _get_active_snapstream(self) -> SnapstreamProto | None:
579 """Get active stream for given player_id."""
580 if group := self.snap_client.group:
581 with suppress(KeyError):
582 return self.snap_provider._snapserver.stream(group.stream)
583 return None
584
585 def _get_player_ids_of_curr_group(self) -> list[str]:
586 snap_group = self.snap_client.group
587 if snap_group is None:
588 return []
589 return [
590 ma_id
591 for client_id in snap_group.clients
592 if (ma_id := self.snap_provider._get_ma_id(client_id))
593 ]
594
595 def _get_players_of_curr_group(self) -> list[Player]:
596 return [
597 ma_player
598 for ma_id in self._get_player_ids_of_curr_group()
599 if (ma_player := self.mass.players.get_player(ma_id))
600 ]
601