/
/
/
1"""
2Helpers for hosting a shared listening experience.
3
4Provides the SharedPlaybackSession abstraction that plugin providers (e.g. the
5party plugin) build on to let a group of guests listen to the same queue.
6Two modes are supported:
7
8- VENUE: an existing real player owns the queue and plays out loud;
9 guests may optionally listen in on their own device when the venue
10 player supports grouping with it.
11- REMOTE: a hidden Sendspin virtual player owns the queue and leads the
12 group; every guest's web player can be attached, so all playback
13 happens on the guests' own devices (silent-disco style).
14
15NOTE: the virtual player backing a REMOTE session lives in memory of the
16Sendspin provider. When that provider reloads, the session is gone and the
17owning plugin is responsible for re-creating it; passing the same session_id
18to :meth:`SharedPlaybackSession.create_remote` yields the same player_id.
19"""
20
21from __future__ import annotations
22
23import asyncio
24import logging
25from enum import StrEnum
26from functools import partial
27from typing import TYPE_CHECKING, cast
28
29from music_assistant_models.enums import PlayerFeature
30from music_assistant_models.errors import (
31 MusicAssistantError,
32 SetupFailedError,
33 UnsupportedFeaturedException,
34)
35
36from music_assistant.helpers.util import join_task
37
38if TYPE_CHECKING:
39 from music_assistant.mass import MusicAssistant
40 from music_assistant.models.player import Player
41 from music_assistant.providers.sendspin.provider import SendspinProvider
42
43SENDSPIN_DOMAIN = "sendspin"
44REMOTE_CREATION_CLEANUP_TIMEOUT = 15.0
45REMOTE_REMOVAL_CLEANUP_DELAYS = (0.0, 1.0, 5.0)
46
47LOGGER = logging.getLogger(__name__)
48
49
50class SharedPlaybackMode(StrEnum):
51 """Mode of a shared playback session."""
52
53 VENUE = "venue"
54 REMOTE = "remote"
55
56
57def is_remote_session_host(mass: MusicAssistant, player_id: str) -> bool:
58 """
59 Return whether the given player is the virtual host of a REMOTE session.
60
61 :param mass: MusicAssistant instance.
62 :param player_id: Player to inspect.
63 """
64 sendspin = cast("SendspinProvider | None", mass.get_provider(SENDSPIN_DOMAIN))
65 return sendspin is not None and sendspin.is_virtual_player(player_id)
66
67
68class SharedPlaybackSession:
69 """
70 A player/queue that hosts a shared listening experience.
71
72 Use the :meth:`create_venue` or :meth:`create_remote` factory to create a
73 session; the owning plugin drives playback on :attr:`queue_id` and calls
74 :meth:`close` when the session ends.
75 """
76
77 def __init__(self, mass: MusicAssistant, mode: SharedPlaybackMode, player_id: str) -> None:
78 """Initialize the session. Use the create_venue/create_remote factories instead."""
79 self.mass = mass
80 self._mode = mode
81 self._player_id = player_id
82 self._guest_listeners: set[str] = set()
83
84 @classmethod
85 async def create_venue(
86 cls, mass: MusicAssistant, venue_player_id: str
87 ) -> SharedPlaybackSession:
88 """
89 Create a session hosted by an existing (real) player.
90
91 :param mass: MusicAssistant instance.
92 :param venue_player_id: The player_id of the player that owns the queue
93 and plays out loud.
94 :raises SetupFailedError: If the venue player is unknown.
95 :return: The created session.
96 """
97 if mass.players.get_player(venue_player_id) is None:
98 raise SetupFailedError(f"Venue player {venue_player_id} is not available")
99 return cls(mass, SharedPlaybackMode.VENUE, venue_player_id)
100
101 @classmethod
102 async def create_remote(
103 cls,
104 mass: MusicAssistant,
105 owner_instance_id: str,
106 display_name: str,
107 session_id: str | None = None,
108 ) -> SharedPlaybackSession:
109 """
110 Create a session hosted by a hidden Sendspin virtual player.
111
112 :param mass: MusicAssistant instance.
113 :param owner_instance_id: Instance id of the plugin provider that owns
114 the session (the virtual player is removed when it unloads).
115 :param display_name: Human readable name for the virtual player.
116 :param session_id: Optional stable id for the virtual player so the
117 owner can re-create the session with the same player_id.
118 :raises SetupFailedError: If the Sendspin provider is not loaded.
119 :return: The created session.
120 """
121 sendspin = cast("SendspinProvider | None", mass.get_provider(SENDSPIN_DOMAIN))
122 if sendspin is None:
123 raise SetupFailedError("The Sendspin provider is required for a remote session")
124 creation = sendspin.create_virtual_player(
125 owner_instance_id=owner_instance_id,
126 display_name=display_name,
127 player_id=session_id,
128 )
129 try:
130 creation_task = mass.create_task(creation, eager_start=False)
131 except Exception:
132 creation.close()
133 raise
134 cleanup_required = asyncio.get_running_loop().create_future()
135 cleanup = cls._cleanup_cancelled_remote_creation(
136 mass,
137 sendspin,
138 creation_task,
139 cleanup_required,
140 )
141 try:
142 mass.create_task(cleanup, eager_start=False)
143 except Exception:
144 cleanup.close()
145 cls._cancel_and_observe_creation(mass, sendspin, creation_task)
146 raise
147 try:
148 # join: cancelling the caller must not abort the creation halfway, or the
149 # cleanup task can no longer remove the player it left behind
150 player_id = await join_task(creation_task)
151 except asyncio.CancelledError:
152 cleanup_required.set_result(True)
153 raise
154 except Exception:
155 cleanup_required.set_result(False)
156 raise
157 cleanup_required.set_result(False)
158 return cls(mass, SharedPlaybackMode.REMOTE, player_id)
159
160 @property
161 def mode(self) -> SharedPlaybackMode:
162 """Return the mode of this session."""
163 return self._mode
164
165 @property
166 def player_id(self) -> str:
167 """Return the player_id of the player that hosts this session."""
168 return self._player_id
169
170 @property
171 def queue_id(self) -> str:
172 """Return the queue_id of the queue that hosts this session."""
173 # a player-owned queue always has the same id as the player
174 return self._player_id
175
176 def can_listen_in(self, web_player_id: str) -> bool:
177 """
178 Return whether the given guest web player can listen in on this session.
179
180 :param web_player_id: The player_id of the guest's web player.
181 """
182 if (host_player := self._get_host_player()) is None:
183 return False
184 if PlayerFeature.SET_MEMBERS not in host_player.state.supported_features:
185 return False
186 # state.can_group_with handles all protocol expansion and translation,
187 # for both a real venue player and a (virtual) Sendspin host player
188 return (
189 web_player_id in host_player.state.can_group_with
190 or web_player_id in host_player.state.group_members
191 )
192
193 async def add_guest_listener(self, web_player_id: str) -> None:
194 """
195 Attach a guest's web player to this session so it plays the same audio.
196
197 :param web_player_id: The player_id of the guest's web player.
198 :raises UnsupportedFeaturedException: If the session host does not
199 support grouping with the given player.
200 """
201 if not self.can_listen_in(web_player_id):
202 raise UnsupportedFeaturedException(
203 f"Player {web_player_id} can not listen in on this session"
204 )
205 await self.mass.players.cmd_set_members(self._player_id, player_ids_to_add=[web_player_id])
206 self._guest_listeners.add(web_player_id)
207
208 async def restore_guest_listeners(self) -> None:
209 """
210 Restore tracked guest listeners missing from the host player's group.
211
212 Missing or temporarily incompatible guest players remain tracked so a
213 later playback transition can restore them after they reconnect.
214 """
215 host_player = self._get_host_player()
216 if host_player is None or not self._guest_listeners:
217 return
218 group_members = set(host_player.state.group_members)
219 for web_player_id in sorted(self._guest_listeners):
220 if web_player_id in group_members:
221 continue
222 guest_player = self.mass.players.get_player(web_player_id)
223 if (
224 guest_player is None
225 or not guest_player.state.available
226 or not self.can_listen_in(web_player_id)
227 ):
228 continue
229 try:
230 await self.mass.players.cmd_set_members(
231 self._player_id,
232 player_ids_to_add=[web_player_id],
233 )
234 except MusicAssistantError as err:
235 LOGGER.warning(
236 "Could not restore guest listener %s to shared playback session %s: %s",
237 web_player_id,
238 self._player_id,
239 err,
240 )
241 continue
242 group_members.add(web_player_id)
243
244 async def remove_guest_listener(self, web_player_id: str) -> None:
245 """
246 Detach a guest's web player from this session.
247
248 :param web_player_id: The player_id of the guest's web player.
249 """
250 self._guest_listeners.discard(web_player_id)
251 if self._get_host_player() is None:
252 return
253 await self.mass.players.cmd_set_members(
254 self._player_id, player_ids_to_remove=[web_player_id]
255 )
256
257 async def close(self) -> None:
258 """
259 Tear down the session.
260
261 In REMOTE mode the virtual player (and its queue) is removed entirely.
262 In VENUE mode only the guest listeners added through this session are
263 detached; the venue player itself is left untouched.
264 """
265 if self._mode == SharedPlaybackMode.REMOTE:
266 sendspin = cast("SendspinProvider | None", self.mass.get_provider(SENDSPIN_DOMAIN))
267 if sendspin is not None and sendspin.is_virtual_player(self._player_id):
268 await sendspin.remove_virtual_player(self._player_id)
269 self._guest_listeners.clear()
270 return
271 if self._guest_listeners and self._get_host_player() is not None:
272 await self.mass.players.cmd_set_members(
273 self._player_id, player_ids_to_remove=list(self._guest_listeners)
274 )
275 self._guest_listeners.clear()
276
277 @classmethod
278 async def _cleanup_cancelled_remote_creation(
279 cls,
280 mass: MusicAssistant,
281 sendspin: SendspinProvider,
282 creation_task: asyncio.Task[str],
283 cleanup_required: asyncio.Future[bool],
284 ) -> None:
285 """
286 Clean up a remote virtual player when its session creation is cancelled.
287
288 :param mass: MusicAssistant instance.
289 :param sendspin: Sendspin provider that owns the virtual player.
290 :param creation_task: In-flight virtual-player creation task.
291 :param cleanup_required: Signal indicating whether cleanup is needed.
292 """
293 if not await asyncio.shield(cleanup_required):
294 return
295 try:
296 done, _ = await asyncio.wait(
297 (creation_task,),
298 timeout=REMOTE_CREATION_CLEANUP_TIMEOUT,
299 )
300 if not done:
301 LOGGER.warning("Timed out waiting for cancelled remote session creation")
302 cls._cancel_and_observe_creation(mass, sendspin, creation_task)
303 return
304 except asyncio.CancelledError:
305 cls._cancel_and_observe_creation(mass, sendspin, creation_task)
306 raise
307
308 try:
309 player_id = creation_task.result()
310 except asyncio.CancelledError:
311 return
312 except Exception as err:
313 LOGGER.debug("Cancelled remote session creation failed: %s", err)
314 return
315
316 await cls._cleanup_cancelled_remote_player(sendspin, player_id)
317
318 @staticmethod
319 async def _cleanup_cancelled_remote_player(
320 sendspin: SendspinProvider,
321 player_id: str,
322 ) -> None:
323 """
324 Remove a virtual player left by cancelled remote session creation.
325
326 :param sendspin: Sendspin provider that owns the virtual player.
327 :param player_id: Virtual player to remove.
328 """
329 last_error: Exception | None = None
330 for delay in REMOTE_REMOVAL_CLEANUP_DELAYS:
331 if delay:
332 await asyncio.sleep(delay)
333 try:
334 if not sendspin.is_virtual_player(player_id):
335 return
336 # awaited to completion on purpose: a timeout is no reliable bound on
337 # the teardown - parts of it swallow the cancellation (see
338 # AsyncProcess.close), and one that does land leaves the player
339 # half torn down for the next attempt to trip over
340 await sendspin.remove_virtual_player(player_id)
341 return
342 except Exception as err:
343 last_error = err
344 LOGGER.warning(
345 "Could not clean up cancelled remote session %s: %s",
346 player_id,
347 last_error,
348 )
349
350 @classmethod
351 def _cancel_and_observe_creation(
352 cls,
353 mass: MusicAssistant,
354 sendspin: SendspinProvider,
355 task: asyncio.Task[str],
356 ) -> None:
357 """Cancel virtual-player creation and observe its eventual result."""
358 task.cancel()
359 cls._observe_late_remote_creation(mass, sendspin, task)
360
361 @classmethod
362 def _observe_late_remote_creation(
363 cls,
364 mass: MusicAssistant,
365 sendspin: SendspinProvider,
366 task: asyncio.Task[str],
367 ) -> None:
368 """Observe creation after bounded cleanup stops waiting for it."""
369 task.add_done_callback(partial(cls._handle_late_remote_creation, mass, sendspin))
370
371 @classmethod
372 def _handle_late_remote_creation(
373 cls,
374 mass: MusicAssistant,
375 sendspin: SendspinProvider,
376 task: asyncio.Task[str],
377 ) -> None:
378 """Schedule cleanup when cancelled creation eventually returns a player."""
379 if task.cancelled():
380 return
381 try:
382 player_id = task.result()
383 except Exception as err:
384 LOGGER.debug("Cancelled remote session creation failed: %s", err)
385 return
386 cleanup = cls._cleanup_cancelled_remote_player(sendspin, player_id)
387 try:
388 mass.create_task(cleanup, eager_start=False)
389 except Exception as err:
390 cleanup.close()
391 LOGGER.warning(
392 "Could not schedule cancelled remote session cleanup for %s: %s",
393 player_id,
394 err,
395 )
396
397 def _get_host_player(self) -> Player | None:
398 """Return the (available) player hosting this session, if any."""
399 player = self.mass.players.get_player(self._player_id)
400 if player is None or not player.state.available:
401 return None
402 return player
403