/
/
1"""
2MusicAssistant Player Queues Controller.
3
4Handles all logic to PLAY Media Items, provided by Music Providers to supported players.
5
6It is loosely coupled to the MusicAssistant Music Controller and Player Controller.
7A Music Assistant Player always has a PlayerQueue associated with it
8which holds the queue items and state.
9
10The PlayerQueue is in that case the active source of the player,
11but it can also be something else, hence the loose coupling.
12"""
13
14from __future__ import annotations
15
16import asyncio
17import random
18import time
19from typing import TYPE_CHECKING, Any, Final, cast
20
21import shortuuid
22from music_assistant_models.auth import Scope
23from music_assistant_models.enums import (
24 EventType,
25 MediaType,
26 PlaybackState,
27 PlayerType,
28 QueueOption,
29 RepeatMode,
30)
31from music_assistant_models.errors import (
32 AudioError,
33 InsufficientPermissions,
34 InvalidCommand,
35 InvalidDataError,
36 MediaNotFoundError,
37 PlayerUnavailableError,
38 QueueEmpty,
39)
40from music_assistant_models.media_items import (
41 Audiobook,
42 ItemMapping,
43 MediaItemType,
44 PlayableMediaItemType,
45 Playlist,
46 PodcastEpisode,
47 SoundEffect,
48 Track,
49)
50from music_assistant_models.player_queue import PlayerQueue
51
52from music_assistant.constants import (
53 ATTR_ANNOUNCEMENT_IN_PROGRESS,
54 MASS_LOGO_ONLINE,
55 PLAYLIST_MEDIA_TYPES,
56)
57from music_assistant.controllers.player_queues.autoplay import Autoplay
58from music_assistant.controllers.player_queues.config import (
59 core_config_entries,
60 queue_config_entries,
61)
62from music_assistant.controllers.player_queues.constants import (
63 CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
64 CACHE_CATEGORY_PLAYER_QUEUE_STATE,
65 PLAYBACK_START_TIMEOUT,
66 QUEUE_CACHE_SAVE_DELAY,
67)
68from music_assistant.controllers.player_queues.helpers import (
69 get_current_playback_speed,
70 handle_play_action,
71 is_dynamic_source,
72)
73from music_assistant.controllers.player_queues.managed_pool import ManagedPool
74from music_assistant.controllers.player_queues.media_resolver import MediaResolver
75from music_assistant.controllers.player_queues.playback_tracker import PlaybackTrackerMixin
76from music_assistant.controllers.player_queues.queue_loader import QueueLoaderMixin
77from music_assistant.controllers.player_queues.smart_shuffle import SmartShuffle
78from music_assistant.controllers.player_queues.state import PlayerQueueData
79from music_assistant.controllers.player_queues.stream_feeder import StreamFeederMixin
80from music_assistant.controllers.webserver.helpers.auth_middleware import get_current_user
81from music_assistant.helpers.api import api_command
82from music_assistant.models.music_provider import ProviderStreamLimitError
83from music_assistant.models.player import Player, PlayerMedia
84
85if TYPE_CHECKING:
86 from collections.abc import Iterator
87
88 from music_assistant_models import BackgroundTask
89 from music_assistant_models.config_entries import (
90 ConfigEntry,
91 ConfigValueOption,
92 CoreConfig,
93 )
94 from music_assistant_models.queue_item import QueueItem
95
96 from music_assistant import MusicAssistant
97 from music_assistant.constants import PlaylistPlayableItem
98 from music_assistant.controllers.music.recency import RecencyWindows
99 from music_assistant.helpers.json import SerializableType
100 from music_assistant.models.player import Player
101
102
103# the container media types worth surfacing as a queue "source" for clients to display. Individual
104# items (single tracks, radio streams, podcast episodes, live audio sources, ...) carry no grouping
105# and only clutter the "playing from" representation, so they are omitted from the wire `sources`.
106_WIRE_SOURCE_MEDIA_TYPES: Final = frozenset(
107 {
108 MediaType.ARTIST,
109 MediaType.ALBUM,
110 MediaType.PLAYLIST,
111 MediaType.PODCAST,
112 MediaType.AUDIOBOOK,
113 }
114)
115
116
117class PlayerQueuesController(QueueLoaderMixin, PlaybackTrackerMixin, StreamFeederMixin):
118 """
119 Controller holding all logic to enqueue music for players.
120
121 The loading, playback-tracking and stream-feeding logic lives in mixins (over the shared base);
122 this class owns the public API surface, the per-queue records and the stateful helper services.
123 """
124
125 def __init__(self, mass: MusicAssistant) -> None:
126 """Initialize core controller."""
127 super().__init__(mass)
128 # server-side per-queue records, keyed by queue_id; each bundles the wire PlayerQueue with
129 # its items, dynamic-source items and runtime-only state (see PlayerQueueData)
130 self._queue_data: dict[str, PlayerQueueData] = {}
131 # stateful helper services (own per-queue state + lifecycle), constructed with self
132 self._autoplay = Autoplay(self)
133 self._smart_shuffle = SmartShuffle(self)
134 self._managed_pool = ManagedPool(self)
135 self._media_resolver = MediaResolver(self)
136 self.manifest.name = "Player Queues controller"
137 self.manifest.description = (
138 "Music Assistant's core controller which manages the queues for all players."
139 )
140 self.manifest.icon = "playlist-music"
141
142 async def close(self) -> None:
143 """Cleanup on exit."""
144 # stop all playback
145 for queue in self.all():
146 if queue.state in (PlaybackState.PLAYING, PlaybackState.PAUSED):
147 await self.stop(queue.queue_id)
148 # flush any pending (debounced) state writes so the latest queue survives shutdown/update
149 for queue in self.all():
150 self.mass.cancel_timer(f"save_queue_cache_{queue.queue_id}")
151 await self._save_queue_to_cache(queue.queue_id)
152
153 async def get_diagnostics(self) -> dict[str, SerializableType]:
154 """Return diagnostics info for this controller to include in diagnostics reports."""
155 queues = [queue_data.queue for queue_data in self._queue_data.values()]
156 by_state: dict[str, int] = {}
157 for queue in queues:
158 by_state[queue.state.value] = by_state.get(queue.state.value, 0) + 1
159 return {
160 "total": len(queues),
161 "active": sum(queue.active for queue in queues),
162 "by_state": by_state,
163 "flow_mode_active": sum(queue.flow_mode for queue in queues),
164 "dynamic_mode_active": sum(queue.is_dynamic for queue in queues),
165 "total_items": sum(queue.items for queue in queues),
166 }
167
168 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
169 """Return the core-module (global) config entries: the queue-controller defaults."""
170 # kept cheap (no library lookup): the config controller populates the global autoplay
171 # playlist dropdown for the UI, so this stays fast on the config value/parse path
172 return core_config_entries(self.mass)
173
174 async def update_config(self, config: CoreConfig, changed_keys: set[str]) -> None:
175 """Apply a global queue-settings change: refresh derived per-queue state and notify clients."""
176 await super().update_config(config, changed_keys)
177 if not any(key.startswith("values/") for key in changed_keys):
178 return
179 # queues that follow a changed global value may flip their derived indicators, so refresh
180 # and signal them (mirrors what save_player_queue_config does for a single queue)
181 for queue in self.all():
182 queue.smart_fades_active = self.mass.streams.is_smart_fades_active(queue)
183 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
184 self.signal_update(queue.queue_id)
185
186 def get_queue_config_entries(
187 self, playlist_options: list[ConfigValueOption] | None = None
188 ) -> list[ConfigEntry]:
189 """
190 Return the per-queue config entries.
191
192 The autoplay_mode select disables the 'similar' option when no provider can supply
193 similar tracks. The crossfade_mode select's options and default depend on whether smart
194 fades are available: the smart option is disabled (shown but not selectable) and the
195 default falls back to standard crossfade when smart fades can't be used on this server.
196
197 :param playlist_options: Library playlists to offer for the 'playlist' autoplay mode.
198 Only populated when serving the entries to the UI; the parse path can omit it.
199 """
200 return queue_config_entries(self.mass, playlist_options)
201
202 def __iter__(self) -> Iterator[PlayerQueue]:
203 """Iterate over (available) players."""
204 return iter(queue_data.queue for queue_data in self._queue_data.values())
205
206 @api_command("player_queues/all", required_scope=Scope.QUEUES_READ)
207 def all(self) -> tuple[PlayerQueue, ...]:
208 """Return all registered PlayerQueues."""
209 return tuple(queue_data.queue for queue_data in self._queue_data.values())
210
211 @api_command("player_queues/get", required_scope=Scope.QUEUES_READ)
212 def get(self, queue_id: str) -> PlayerQueue | None:
213 """Return PlayerQueue by queue_id or None if not found."""
214 queue_data = self._queue_data.get(queue_id)
215 return queue_data.queue if queue_data else None
216
217 def queue_data(self, queue_id: str) -> PlayerQueueData:
218 """
219 Return the server-side record for a queue (raises if the queue is unknown).
220
221 Internal accessor for the stateful helper services so they reach per-queue state through
222 the controller rather than its private store.
223 """
224 return self._queue_data[queue_id]
225
226 def queue_data_or_none(self, queue_id: str) -> PlayerQueueData | None:
227 """Return the server-side record for a queue, or None if it is not registered."""
228 return self._queue_data.get(queue_id)
229
230 @api_command("player_queues/items", required_scope=Scope.QUEUES_READ)
231 def items(self, queue_id: str, limit: int = 500, offset: int = 0) -> list[QueueItem]:
232 """Return all QueueItems for given PlayerQueue."""
233 if (queue_data := self._queue_data.get(queue_id)) is None:
234 return []
235 return queue_data.items[offset : offset + limit]
236
237 @api_command("player_queues/get_active_queue", required_scope=Scope.QUEUES_READ)
238 def get_active_queue(self, player_id: str) -> PlayerQueue | None:
239 """Return the current active/synced queue for a player."""
240 if player := self.mass.players.get_player(player_id):
241 return self.mass.players.get_active_queue(player)
242 return None
243
244 # Queue commands
245
246 @api_command("player_queues/shuffle", required_scope=Scope.QUEUES_CONTROL)
247 async def set_shuffle(self, queue_id: str, shuffle_enabled: bool) -> None:
248 """Configure shuffle setting on the the queue."""
249 queue = self._queue_data[queue_id].queue
250 if queue.is_dynamic:
251 # a dynamic queue is an always-on, recency-orchestrated smart mix; manual shuffle
252 # (and plain linear order) have no meaning here so the toggle is locked
253 raise InvalidCommand("Cannot change shuffle while the queue is in dynamic mode")
254 if queue.shuffle_enabled == shuffle_enabled:
255 return # no change
256 queue.shuffle_enabled = shuffle_enabled
257 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
258 queue_items = self._queue_data[queue_id].items
259 cur_index = (
260 queue.index_in_buffer if queue.index_in_buffer is not None else queue.current_index
261 )
262 if cur_index is not None:
263 next_index = cur_index + 1
264 next_items = queue_items[next_index:]
265 else:
266 next_items = []
267 next_index = 0
268 if not shuffle_enabled:
269 # shuffle disabled, try to restore original sort order of the remaining items
270 next_items.sort(key=lambda x: x.sort_index, reverse=False)
271 await self.load(
272 queue_id=queue_id,
273 queue_items=next_items,
274 insert_at_index=next_index,
275 keep_remaining=False,
276 shuffle=shuffle_enabled,
277 )
278
279 def is_smart_shuffle_active(self, queue: PlayerQueue) -> bool:
280 """
281 Return whether smart shuffle is currently in effect for the queue.
282
283 A dynamic queue is always an orchestrated smart mix (the managed pool), so it always counts
284 as active; otherwise smart shuffle is active when shuffle is on and the per-queue
285 smart-shuffle setting is enabled.
286
287 :param queue: The queue to evaluate.
288 """
289 if queue.is_dynamic:
290 return True
291 return queue.shuffle_enabled and self._smart_shuffle.is_enabled(queue.queue_id)
292
293 @api_command("player_queues/autoplay", required_scope=Scope.QUEUES_CONTROL)
294 def set_autoplay(self, queue_id: str, autoplay_enabled: bool) -> None:
295 """Configure Autoplay setting on the queue."""
296 queue_data = self._queue_data[queue_id]
297 queue = queue_data.queue
298 queue.autoplay_enabled = autoplay_enabled
299 # if we're already at/near the end of the queue, kick off a refill right away
300 # (an active dynamic source manages its own refills, so leave it be)
301 if (
302 queue.autoplay_enabled
303 and not queue.is_dynamic
304 and queue.current_index is not None
305 and (queue.items - queue.current_index) < 5
306 ):
307 task_id = f"fill_autoplay_tracks_{queue_id}"
308 self.mass.call_later(5, self._fill_autoplay_tracks, queue_id, task_id=task_id)
309 self.signal_update(queue_id=queue_id)
310
311 @api_command(
312 "player_queues/dont_stop_the_music", required_scope=Scope.QUEUES_CONTROL, alias=True
313 )
314 def set_dont_stop_the_music(self, queue_id: str, dont_stop_the_music_enabled: bool) -> None:
315 """Backwards-compatible alias for the autoplay command, used by older clients."""
316 self.set_autoplay(queue_id, dont_stop_the_music_enabled)
317
318 @api_command("player_queues/repeat", required_scope=Scope.QUEUES_CONTROL)
319 def set_repeat(self, queue_id: str, repeat_mode: RepeatMode) -> None:
320 """Configure repeat setting on the the queue."""
321 queue = self._queue_data[queue_id].queue
322 if queue.is_dynamic:
323 # a dynamic queue is an always-on flowing mix of its sources; repeat has no meaning here
324 raise InvalidCommand("Cannot change repeat while the queue is in dynamic mode")
325 if queue.repeat_mode == repeat_mode:
326 return # no change
327 queue.repeat_mode = repeat_mode
328 self.signal_update(queue_id)
329 if (
330 queue.state == PlaybackState.PLAYING
331 and queue.index_in_buffer is not None
332 and queue.index_in_buffer == queue.current_index
333 ):
334 # if the queue is playing,
335 # ensure to (re)queue the next track because it might have changed
336 # note that we only do this if the player has loaded the current track
337 # if not, we wait until it has loaded to prevent conflicts
338 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
339 self._enqueue_next_item(queue_id, next_item)
340
341 @api_command("player_queues/crossfade", required_scope=Scope.QUEUES_CONTROL)
342 def set_crossfade(self, queue_id: str, crossfade_enabled: bool) -> None:
343 """Enable or disable crossfade on the queue."""
344 queue = self._queue_data[queue_id].queue
345 if queue.crossfade_enabled == crossfade_enabled:
346 return # no change
347 queue.crossfade_enabled = crossfade_enabled
348 # refresh the derived smart-fades indicator so the update we signal reflects the new state
349 queue.smart_fades_active = self.mass.streams.is_smart_fades_active(queue)
350 self.signal_update(queue_id)
351 if (
352 queue.state == PlaybackState.PLAYING
353 and queue.index_in_buffer is not None
354 and queue.index_in_buffer == queue.current_index
355 ):
356 # re-enqueue the next track so the new crossfade behaviour applies to the
357 # upcoming transition (only when the player has already loaded the current track)
358 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
359 self._enqueue_next_item(queue_id, next_item)
360
361 @api_command("player_queues/overlay", required_scope=Scope.QUEUES_CONTROL)
362 async def set_overlay(
363 self,
364 queue_id: str,
365 enabled: bool | None = None,
366 source: str | None = None,
367 volume: int | None = None,
368 ) -> None:
369 """
370 Configure the audio overlay for the given queue.
371
372 The audio overlay mixes a looping sound effect (e.g. rain or white noise)
373 into the queue's audio stream. Changes take effect immediately: if the
374 queue is playing, playback is restarted from the current position.
375
376 :param queue_id: queue_id of the queue to configure.
377 :param enabled: Enable or disable the audio overlay. Omit to leave unchanged.
378 :param source: URI of the sound effect item to mix in. Omit to leave unchanged.
379 :param volume: Overlay loudness relative to the music in percent
380 (0-200, 100 = equally loud). Omit to leave unchanged.
381 """
382 queue = self._queue_data[queue_id].queue
383 changed = audible_change = False
384 if source is not None:
385 item = await self.mass.music.get_item_by_uri(source)
386 if item.media_type != MediaType.SOUND_EFFECT:
387 raise InvalidDataError("Audio overlay source must be a sound effect item")
388 mapping = ItemMapping.from_item(cast("SoundEffect", item))
389 if queue.overlay_source != mapping:
390 queue.overlay_source = mapping
391 changed = True
392 audible_change = queue.overlay_enabled
393 if volume is not None:
394 if not (0 <= volume <= 200):
395 raise InvalidDataError(f"Overlay volume must be between 0 and 200, got {volume}")
396 if queue.overlay_volume != volume:
397 queue.overlay_volume = volume
398 changed = True
399 audible_change |= queue.overlay_enabled
400 if enabled is not None and queue.overlay_enabled != enabled:
401 if enabled and queue.overlay_source is None:
402 raise InvalidCommand("Can not enable audio overlay: no overlay source selected")
403 queue.overlay_enabled = enabled
404 changed = audible_change = True
405 if not changed:
406 return
407 self.signal_update(queue_id)
408 if audible_change and queue.state == PlaybackState.PLAYING:
409 # restart playback from the current position so the change is heard
410 # immediately instead of after the player's audio buffer drains
411 await self.resume(queue_id)
412
413 # Two timebases are used in this controller when variable playback speed is in
414 # effect (atempo applied server-side):
415 # "stream-time" — seconds of audio the player has played (post-atempo).
416 # "media-time" — seconds of the original content the listener has heard.
417 # What the user expects to see on the progress bar and what
418 # we use for resume positions.
419 # Conversion: media-time = stream-time x playback_speed.
420 @api_command("player_queues/set_playback_speed", required_scope=Scope.QUEUES_CONTROL)
421 async def set_playback_speed(
422 self, queue_id: str, speed: float, queue_item_id: str | None = None
423 ) -> None:
424 """
425 Set the playback speed for the given queue item.
426
427 Variable playback speed is supported only for audiobooks and podcast episodes.
428
429 If queue_item_id is not provided,
430 the speed will be set for the current item in the queue.
431
432 :param queue_id: queue_id of the queue to configure.
433 :param speed: playback speed multiplier (0.5 to 3.0). 1.0 = normal speed.
434 """
435 if not (0.5 <= speed <= 3.0):
436 raise InvalidDataError(f"Playback speed must be between 0.5 and 3.0, got {speed}")
437 queue = self._queue_data[queue_id].queue
438 if not queue.current_item:
439 raise QueueEmpty("Cannot set playback speed: queue is empty")
440 queue_item_id = queue_item_id or queue.current_item.queue_item_id
441 queue_item = self.get_item(queue_id, queue_item_id)
442 if not queue_item:
443 raise InvalidDataError(f"Queue item {queue_item_id} not found in queue")
444 if queue_item.media_type not in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE):
445 raise InvalidCommand(
446 "Variable playback speed is only supported for audiobooks and podcast episodes"
447 )
448 if not queue_item.duration:
449 raise InvalidCommand("Cannot set playback speed for items with unknown duration")
450 current_speed = float(queue_item.extra_attributes.get("playback_speed") or 1.0)
451 if abs(current_speed - speed) < 0.001:
452 return # no change
453 # use extra_attributes of the queue item to store the playback speed
454 queue_item.extra_attributes["playback_speed"] = speed
455 # mirror onto the queue so corrected_elapsed_time advances in media-time
456 # immediately, before the next on_player_elapsed_time_corrected snapshot.
457 if queue.current_item and queue.current_item.queue_item_id == queue_item_id:
458 # close off the wallclock seconds that already ticked by at the old speed
459 # before switching, so corrected_elapsed_time doesn't multiply them by the new speed
460 if queue.state == PlaybackState.PLAYING:
461 queue.elapsed_time = queue.corrected_elapsed_time
462 queue.elapsed_time_last_updated = time.time()
463 queue.playback_speed = speed
464 self.signal_update(queue_id)
465 if queue.state == PlaybackState.PLAYING:
466 await self.resume(queue_id)
467
468 @api_command(
469 "player_queues/play_media", required_scope=Scope.QUEUES_CONTROL, allow_impersonation=True
470 )
471 async def play_media(
472 self,
473 queue_id: str,
474 media: MediaItemType | ItemMapping | str | list[MediaItemType | ItemMapping | str],
475 option: QueueOption | None = None,
476 radio_mode: bool = False,
477 start_item: PlayableMediaItemType | str | None = None,
478 sort_by: str | None = None,
479 start_from_beginning: bool = False,
480 shuffle: bool | None = None,
481 ) -> None:
482 """
483 Play media item(s) on the given queue.
484
485 :param queue_id: The queue_id of the queue to play media on.
486 :param media: Media that should be played (MediaItem(s) and/or uri's).
487 :param option: Which enqueue mode to use.
488 :param radio_mode: Deprecated — translated to a radio_playlist:// dynamic playlist;
489 prefer enqueuing that URI directly.
490 :param start_item: Optional item to start the playlist or album from.
491 :param sort_by: Optional sort key to order tracks before applying start_item.
492 :param start_from_beginning: Start a podcast episode at position 0, ignoring any
493 saved resume position. The stored progress itself is left untouched.
494 :param shuffle: Play the media shuffled (or explicitly in order). Only applies to the
495 options that start playing right away (play/replace), and never to a dynamic source
496 (an always-on smart mix). Omit to follow the queue's own shuffle setting, which media
497 with an order of its own (album, podcast, episode, audiobook, audio source) switches
498 off; the first item of a batch decides for the whole batch.
499 """
500 self._check_player_permission(queue_id)
501 if not self.get(queue_id):
502 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
503 # Lock is acquired by the @handle_play_action decorator on the internal handler
504 await self._handle_play_media(
505 queue_id, media, option, radio_mode, start_item, sort_by, start_from_beginning, shuffle
506 )
507
508 @api_command("player_queues/move_item", required_scope=Scope.QUEUES_CONTROL)
509 def move_item(self, queue_id: str, queue_item_id: str, pos_shift: int = 1) -> None:
510 """
511 Move queue item x up/down the queue.
512
513 - queue_id: id of the queue to process this request.
514 - queue_item_id: the item_id of the queueitem that needs to be moved.
515 - pos_shift: move item x positions down if positive value
516 - pos_shift: move item x positions up if negative value
517 - pos_shift: move item to top of queue as next item if 0.
518 """
519 queue = self._queue_data[queue_id].queue
520 item_index = self.index_by_id(queue_id, queue_item_id)
521 if item_index is None:
522 raise InvalidDataError(f"Item {queue_item_id} not found in queue")
523 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
524 msg = f"{item_index} is already played/buffered"
525 raise IndexError(msg)
526
527 queue_items = self._queue_data[queue_id].items
528 queue_items = queue_items.copy()
529
530 if pos_shift == 0 and queue.state == PlaybackState.PLAYING:
531 new_index = (queue.current_index or 0) + 1
532 elif pos_shift == 0:
533 new_index = queue.current_index or 0
534 else:
535 new_index = item_index + pos_shift
536 if (new_index < (queue.current_index or 0)) or (new_index > len(queue_items)):
537 return
538 # move the item in the list
539 queue_items.insert(new_index, queue_items.pop(item_index))
540 self.update_items(queue_id, queue_items)
541
542 @api_command("player_queues/move_item_end", required_scope=Scope.QUEUES_CONTROL)
543 def move_item_end(self, queue_id: str, queue_item_id: str) -> None:
544 """
545 Move queue item to the end the queue.
546
547 - queue_id: id of the queue to process this request.
548 - queue_item_id: the item_id of the queueitem that needs to be moved.
549 """
550 queue = self._queue_data[queue_id].queue
551 item_index = self.index_by_id(queue_id, queue_item_id)
552 if item_index is None:
553 raise InvalidDataError(f"Item {queue_item_id} not found in queue")
554 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
555 msg = f"{item_index} is already played/buffered"
556 raise IndexError(msg)
557
558 queue_items = self._queue_data[queue_id].items
559 if item_index == (len(queue_items) - 1):
560 return
561 queue_items = queue_items.copy()
562
563 new_index = len(self._queue_data[queue_id].items) - 1
564
565 # move the item in the list
566 queue_items.insert(new_index, queue_items.pop(item_index))
567 self.update_items(queue_id, queue_items)
568
569 @api_command("player_queues/delete_item", required_scope=Scope.QUEUES_CONTROL)
570 def delete_item(self, queue_id: str, item_id_or_index: int | str) -> None:
571 """Delete item (by id or index) from the queue."""
572 if isinstance(item_id_or_index, str):
573 item_index = self.index_by_id(queue_id, item_id_or_index)
574 if item_index is None:
575 raise InvalidDataError(f"Item {item_id_or_index} not found in queue")
576 else:
577 item_index = item_id_or_index
578 queue = self._queue_data[queue_id].queue
579 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
580 # ignore request if track already loaded in the buffer
581 # the frontend should guard so this is just in case
582 self.logger.warning("delete requested for item already loaded in buffer")
583 return
584 queue_items = self._queue_data[queue_id].items.copy()
585 queue_items.pop(item_index)
586 self.update_items(queue_id, queue_items)
587
588 @api_command("player_queues/clear", required_scope=Scope.QUEUES_CONTROL)
589 def clear(self, queue_id: str, skip_stop: bool = False) -> None:
590 """Clear all items in the queue, switching shuffle off with them."""
591 self._clear(queue_id, skip_stop)
592 # clearing is an explicit "start over" gesture by the user, so a shuffle that belonged to
593 # the discarded content must not carry over into whatever is played next
594 self._reset_shuffle(queue_id)
595
596 def mark_ended(self, queue_id: str) -> None:
597 """
598 Mark a queue as played to its end, keeping its items so it can be replayed.
599
600 The playback position is parked on the last item rather than cleared: a null index is
601 indistinguishable from a queue that was loaded but never started, and an index past the
602 end is silently misread by everything that does arithmetic on it. `ended` is what tells
603 clients the queue finished, and pressing play starts it over from the first item.
604
605 :param queue_id: The queue_id of the queue that reached its end.
606 """
607 queue_data = self._queue_data[queue_id]
608 queue = queue_data.queue
609 if not queue_data.items:
610 # nothing to replay, so there is nothing to advertise as finished either
611 self._clear(queue_id)
612 return
613 self.mass.streams.audio_processing.clear(queue_id)
614 queue.ended = True
615 queue.current_index = len(queue_data.items) - 1
616 queue.current_item = queue_data.items[-1]
617 queue.next_item = None
618 queue.elapsed_time = 0
619 queue.elapsed_time_last_updated = time.time()
620 queue.index_in_buffer = None
621 queue.resume_pos = 0
622 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
623 self.signal_update(queue_id)
624
625 @api_command("player_queues/save_as_playlist", required_scope=Scope.LIBRARY_WRITE)
626 async def save_as_playlist(self, queue_id: str, name: str) -> BackgroundTask:
627 """
628 Save the current queue items as a new playlist.
629
630 :param queue_id: The queue_id of the queue to save.
631 :param name: The name for the new playlist.
632 """
633 if not self.get(queue_id):
634 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
635 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
636 if not queue_items:
637 raise QueueEmpty("Cannot save an empty queue as a playlist.")
638 # collect URIs from queue items that are playlist-compatible
639 uris: list[str] = []
640 for item in queue_items:
641 if item.uri and item.media_type in PLAYLIST_MEDIA_TYPES:
642 uris.append(item.uri)
643 if not uris:
644 raise InvalidDataError("No valid items in queue to save as playlist.")
645 playlist = await self.mass.music.playlists.create_playlist(name)
646 return await self.mass.music.playlists.add_playlist_tracks(playlist.item_id, uris)
647
648 @api_command("player_queues/stop", required_scope=Scope.QUEUES_CONTROL)
649 @handle_play_action
650 async def stop(self, queue_id: str) -> None:
651 """
652 Handle STOP command for given queue.
653
654 - queue_id: queue_id of the playerqueue to handle the command.
655 """
656 self._check_player_permission(queue_id)
657 # cancel any pending play_index calls for this queue to prevent conflicts
658 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
659 # cancel in-flight preload/enqueue-next so it can't enqueue after stop
660 self.mass.cancel_task(f"preload_next_item_{queue_id}")
661 self.mass.cancel_timer(f"enqueue_next_item_{queue_id}")
662 self.mass.cancel_task(f"enqueue_next_item_{queue_id}")
663 self._set_transitioning(queue_id, False)
664 queue_data = self._queue_data[queue_id]
665 session_id = queue_data.session_id
666 queue_player = self.mass.players.get_player(queue_id, True)
667 if queue_player is None:
668 raise PlayerUnavailableError(f"Player {queue_id} is not available")
669 if (queue := self.get(queue_id)) and queue.active:
670 if queue.state == PlaybackState.PLAYING:
671 queue.resume_pos = int(queue.corrected_elapsed_time)
672 # Use internal handler to avoid circular redirect:
673 # public cmd_stop redirects to queue.stop when a queue is active,
674 # which would loop back here indefinitely.
675 await self.mass.players._handle_cmd_stop(queue_id)
676 if queue_data.session_id == session_id:
677 queue_data.session_id = None
678 self.mass.streams.audio_processing.clear(queue_id, session_id)
679 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
680
681 @api_command("player_queues/play", required_scope=Scope.QUEUES_CONTROL)
682 async def play(self, queue_id: str) -> None:
683 """
684 Handle PLAY command for given queue.
685
686 :param queue_id: queue_id of the playerqueue to handle the command.
687 """
688 self._check_player_permission(queue_id)
689 if not self.get(queue_id):
690 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
691 await self._handle_play(queue_id)
692
693 @api_command("player_queues/pause", required_scope=Scope.QUEUES_CONTROL)
694 async def pause(self, queue_id: str) -> None:
695 """
696 Handle PAUSE command for given queue.
697
698 - queue_id: queue_id of the playerqueue to handle the command.
699 """
700 self._check_player_permission(queue_id)
701 # cancel any pending play_index calls for this queue to prevent conflicts
702 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
703 self._set_transitioning(queue_id, False)
704 if not (queue := self.get(queue_id)):
705 return
706 queue_active = queue.active
707 if queue.active and queue.state == PlaybackState.PLAYING:
708 queue.resume_pos = int(queue.corrected_elapsed_time)
709 # Use internal handler to avoid circular redirect
710 # (cmd_pause redirects to queue.pause, which calls cmd_pause again)
711 await self.mass.players._handle_cmd_pause(queue_id)
712
713 async def _watch_pause(player: Player) -> None:
714 count = 0
715 # wait for pause
716 while count < 5 and player.state.playback_state == PlaybackState.PLAYING:
717 count += 1
718 await asyncio.sleep(1)
719 # wait for unpause
720 if player.state.playback_state != PlaybackState.PAUSED:
721 return
722 count = 0
723 while count < 30 and player.state.playback_state == PlaybackState.PAUSED:
724 count += 1
725 await asyncio.sleep(1)
726 # if player is still paused when the limit is reached, send stop
727 if player.state.playback_state == PlaybackState.PAUSED:
728 await self.stop(queue_id)
729
730 # we auto stop a player from paused when its paused for 30 seconds
731 if (
732 queue_active
733 and (queue_player := self.mass.players.get_player(queue_id))
734 and not queue_player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS)
735 ):
736 self.mass.create_task(_watch_pause(queue_player))
737
738 @api_command("player_queues/play_pause", required_scope=Scope.QUEUES_CONTROL)
739 async def play_pause(self, queue_id: str) -> None:
740 """
741 Toggle play/pause on given playerqueue.
742
743 - queue_id: queue_id of the queue to handle the command.
744 """
745 if (queue := self.get(queue_id)) and queue.state == PlaybackState.PLAYING:
746 await self.pause(queue_id)
747 return
748 await self.play(queue_id)
749
750 @api_command("player_queues/next", required_scope=Scope.QUEUES_CONTROL)
751 @handle_play_action
752 async def next(self, queue_id: str) -> None:
753 """
754 Handle NEXT TRACK command for given queue.
755
756 :param queue_id: queue_id of the queue to handle the command.
757 """
758 self._check_player_permission(queue_id)
759 if (queue := self.get(queue_id)) is None or not queue.active:
760 raise InvalidCommand(f"Queue {queue_id} is not active")
761 self._set_transitioning(queue_id, True)
762 idx = self._queue_data[queue_id].queue.current_index
763 if idx is None:
764 self.logger.warning("Queue %s has no current index", queue.display_name)
765 self._set_transitioning(queue_id, False)
766 return
767 next_index = self._get_next_index(queue_id, idx, True)
768 if next_index is None:
769 self._set_transitioning(queue_id, False)
770 return
771
772 # immediately update current item so UI shows the new track right away
773 queue.current_index = next_index
774 queue.current_item = self.get_item(queue_id, next_index)
775 queue.elapsed_time = 0
776 queue.elapsed_time_last_updated = time.time()
777 self.signal_update(queue_id)
778 if queue_player := self.mass.players.get_player(queue_id, True):
779 queue_player.update_state()
780
781 # debounce rapid next button presses using call_later
782 self.mass.call_later(
783 1,
784 self.play_index,
785 queue_id,
786 next_index,
787 task_id=f"queue_play_index_{queue_id}",
788 )
789
790 @api_command("player_queues/previous", required_scope=Scope.QUEUES_CONTROL)
791 @handle_play_action
792 async def previous(self, queue_id: str) -> None:
793 """
794 Handle PREVIOUS TRACK command for given queue.
795
796 :param queue_id: queue_id of the queue to handle the command.
797 """
798 self._check_player_permission(queue_id)
799 if (queue := self.get(queue_id)) is None or not queue.active:
800 raise InvalidCommand(f"Queue {queue_id} is not active")
801 self._set_transitioning(queue_id, True)
802 current_index = self._queue_data[queue_id].queue.current_index
803 if current_index is None:
804 self._set_transitioning(queue_id, False)
805 return
806 prev_index = int(current_index)
807 # restart current track if elapsed > 5s, otherwise go to previous
808 if self._queue_data[queue_id].queue.elapsed_time < 5:
809 prev_index = max(current_index - 1, 0)
810
811 # immediately update current item so UI shows the new track right away
812 queue.current_index = prev_index
813 queue.current_item = self.get_item(queue_id, prev_index)
814 queue.elapsed_time = 0
815 queue.elapsed_time_last_updated = time.time()
816 self.signal_update(queue_id)
817 if queue_player := self.mass.players.get_player(queue_id, True):
818 queue_player.update_state()
819
820 # debounce rapid previous button presses using call_later
821 self.mass.call_later(
822 1,
823 self.play_index,
824 queue_id,
825 prev_index,
826 task_id=f"queue_play_index_{queue_id}",
827 )
828
829 @api_command("player_queues/skip", required_scope=Scope.QUEUES_CONTROL)
830 async def skip(self, queue_id: str, seconds: int = 10) -> None:
831 """
832 Handle SKIP command for given queue.
833
834 - queue_id: queue_id of the queue to handle the command.
835 - seconds: number of seconds to skip in track. Use negative value to skip back.
836 """
837 if (queue := self.get(queue_id)) is None or not queue.active:
838 raise InvalidCommand(f"Queue {queue_id} is not active")
839 await self.seek(queue_id, int(self._queue_data[queue_id].queue.elapsed_time + seconds))
840
841 @api_command("player_queues/seek", required_scope=Scope.QUEUES_CONTROL)
842 async def seek(self, queue_id: str, position: int = 10) -> None:
843 """
844 Handle SEEK command for given queue.
845
846 - queue_id: queue_id of the queue to handle the command.
847 - position: position in seconds to seek to in the current playing item.
848 """
849 if (queue := self.get(queue_id)) is None or not queue.active:
850 raise InvalidCommand(f"Queue {queue_id} is not active")
851 queue_player = self.mass.players.get_player(queue_id, True)
852 if queue_player is None:
853 raise PlayerUnavailableError(f"Player {queue_id} is not available")
854 if not queue.current_item:
855 raise InvalidCommand(f"Queue {queue_player.state.name} has no item(s) loaded.")
856 if not queue.current_item.duration:
857 raise InvalidCommand("Can not seek items without duration.")
858 position = max(0, int(position))
859 if position > queue.current_item.duration:
860 raise InvalidCommand("Can not seek outside of duration range.")
861 if queue.current_index is None:
862 raise InvalidCommand(f"Queue {queue_player.state.name} has no current index.")
863 # Publish the seek target before rebuilding the stream to prevent progress snapback.
864 queue.elapsed_time = position
865 queue.elapsed_time_last_updated = time.time()
866 self.signal_update(queue_id)
867 await self.play_index(queue_id, queue.current_index, seek_position=position)
868
869 @api_command("player_queues/resume", required_scope=Scope.QUEUES_CONTROL)
870 @handle_play_action
871 async def resume(self, queue_id: str, fade_in: bool | None = None) -> None:
872 """
873 Handle RESUME command for given queue.
874
875 - queue_id: queue_id of the queue to handle the command.
876 """
877 self._check_player_permission(queue_id)
878 queue = self._queue_data[queue_id].queue
879 queue_items = self._queue_data[queue_id].items
880 resume_item = queue.current_item
881 if queue.state == PlaybackState.PLAYING:
882 # resume requested while already playing,
883 # use current position as resume position
884 resume_pos = queue.corrected_elapsed_time
885 fade_in = False
886 else:
887 resume_pos = queue.resume_pos or queue.elapsed_time
888
889 if queue.ended and len(queue_items) > 0:
890 # the queue played to its end and is parked on its last item,
891 # so pressing play starts it over from the beginning
892 resume_item = queue_items[0]
893 resume_pos = 0
894 elif not resume_item and queue.current_index is not None and len(queue_items) > 0:
895 resume_item = self.get_item(queue_id, queue.current_index)
896 resume_pos = 0
897 elif not resume_item and queue.current_index is None and len(queue_items) > 0:
898 # items available in queue but no previous track, start at 0
899 resume_item = self.get_item(queue_id, 0)
900 resume_pos = 0
901
902 if resume_item is not None:
903 queue_player = self.mass.players.get_player(queue_id)
904 if queue_player is None:
905 raise PlayerUnavailableError(f"Player {queue_id} is not available")
906 if (
907 fade_in is None
908 and queue_player.state.playback_state == PlaybackState.IDLE
909 and (time.time() - queue.elapsed_time_last_updated) > 60
910 ):
911 # enable fade in effect if the player is idle for a while
912 fade_in = resume_pos > 0
913 if resume_item.media_type == MediaType.RADIO:
914 # we're not able to skip in online radio so this is pointless
915 resume_pos = 0
916 await self.play_index(
917 queue_id, resume_item.queue_item_id, int(resume_pos), fade_in or False
918 )
919 else:
920 msg = f"Resume queue requested but queue {queue.display_name} is empty"
921 raise QueueEmpty(msg)
922
923 @api_command("player_queues/play_index", required_scope=Scope.QUEUES_CONTROL)
924 @handle_play_action
925 async def play_index( # noqa: PLR0915
926 self,
927 queue_id: str,
928 index: int | str,
929 seek_position: int = 0,
930 fade_in: bool = False,
931 ) -> None:
932 """Play item at index (or item_id) X in queue."""
933 self._check_player_permission(queue_id)
934 # cancel any pending play_index calls for this queue to prevent conflicts
935 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
936 # we set a flag to notify the update logic that we're transitioning to a new track
937 self._set_transitioning(queue_id, True)
938 try:
939 queue_data = self._queue_data[queue_id]
940 queue = queue_data.queue
941 queue.resume_pos = 0
942 # A queue picked up from its end plays its items over from the start, so a resume point
943 # left on an audiobook/episode must not pull it back to where it was left off. The flag
944 # itself is only cleared once an item actually loaded below, so a start that never got
945 # off the ground leaves the queue finished instead of stranding it without a position.
946 restarting_ended_queue = queue.ended
947 if isinstance(index, str):
948 temp_index = self.index_by_id(queue_id, index)
949 if temp_index is None:
950 raise InvalidDataError(f"Item {index} not found in queue")
951 index = temp_index
952 # At this point index is guaranteed to be int
953 queue.index_in_buffer = index
954 queue_data.flow_mode_stream_log = []
955 queue_data.flow_buffer_completed = None
956 queue_data.flow_queue_exhausted = None
957 target_player = self.mass.players.get_player(queue_id)
958 if target_player is None:
959 raise PlayerUnavailableError(f"Player {queue_id} is not available")
960 queue_data.next_item_id_enqueued = None
961 # always update session id when we start a new playback session
962 queue_data.session_id = shortuuid.random(length=8)
963 self.mass.streams.audio_processing.start_session(
964 queue_id,
965 queue_data.session_id,
966 )
967 # handle resume point of audiobook(chapter) or podcast(episode)
968 if (
969 not seek_position
970 and not restarting_ended_queue
971 and (queue_item := self.get_item(queue_id, index))
972 and (resume_position_ms := getattr(queue_item.media_item, "resume_position_ms", 0))
973 ):
974 # the client may have fetched the item before its duration was known
975 await self._restore_probed_duration(queue_item)
976 if queue_item.duration or getattr(queue_item.media_item, "duration", 0):
977 seek_position = max(0, int((resume_position_ms - 500) / 1000))
978 else:
979 # seeking needs a duration, which is determined while streaming
980 self.logger.debug(
981 "Can not resume %s at %ss: its duration is not known (yet)",
982 queue_item.name,
983 int(resume_position_ms / 1000),
984 )
985
986 # restore the persisted playback speed for a freshly queued audiobook/episode
987 # (an in-session item already carries its speed in extra_attributes)
988 if (
989 (queue_item := self.get_item(queue_id, index))
990 and queue_item.media_item is not None
991 and queue_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
992 and "playback_speed" not in queue_item.extra_attributes
993 ):
994 stored_speed = await self.mass.music.get_playback_speed(
995 cast("Audiobook | PodcastEpisode", queue_item.media_item),
996 userid=queue_data.userid,
997 )
998 if stored_speed != 1.0:
999 queue_item.extra_attributes["playback_speed"] = stored_speed
1000
1001 # try to load the item, retry with next item if it fails
1002 for attempt in range(5):
1003 try:
1004 queue_item = self.get_item(queue_id, index)
1005 if not queue_item:
1006 continue # guard
1007 await self._load_item(
1008 queue_item,
1009 self._get_next_index(queue_id, index),
1010 is_start=True,
1011 seek_position=seek_position if attempt == 0 else 0,
1012 fade_in=fade_in if attempt == 0 else False,
1013 )
1014 # if we reach this point, loading the item succeeded, break the loop
1015 queue.current_index = index
1016 queue.current_item = queue_item
1017 # playback is under way, so the queue is no longer sitting at its end
1018 queue.ended = False
1019 # reset the elapsed clock together with the item switch (like
1020 # next/previous do), so queue updates signaled before the player
1021 # reports position don't carry the previous item's elapsed_time
1022 queue.elapsed_time = seek_position if attempt == 0 else 0
1023 queue.elapsed_time_last_updated = time.time()
1024 break
1025 except (MediaNotFoundError, AudioError) as err:
1026 item_name = queue_item.name if queue_item else "unknown"
1027 if isinstance(err, ProviderStreamLimitError):
1028 # the requested item is playable, its provider is just at capacity:
1029 # report that instead of silently advancing to another item
1030 self.logger.error("%s", err)
1031 await self.stop(queue_id)
1032 raise
1033 # Only MediaNotFoundError (item unreachable) is persistent;
1034 # keep AudioError items available so a retry can resurface
1035 # the same actionable error.
1036 if queue_item and isinstance(err, MediaNotFoundError):
1037 queue_item.available = False
1038 next_index = self._get_next_index(queue_id, index, allow_repeat=False)
1039 if next_index is None:
1040 # Surface an AudioError's own (actionable) message;
1041 # MediaNotFoundError gets the generic wording.
1042 if isinstance(err, AudioError) and str(err):
1043 msg = str(err)
1044 else:
1045 msg = f"Playback failed for {item_name} - no more tracks available"
1046 self.logger.error(msg)
1047 await self.stop(queue_id)
1048 raise MediaNotFoundError(msg) from err
1049 self.logger.warning(
1050 "Skipping unplayable item %s",
1051 item_name,
1052 )
1053 index = next_index
1054 else:
1055 # all attempts to find a playable item failed
1056 await self.stop(queue_id)
1057 raise MediaNotFoundError("No playable item found to start playback")
1058
1059 # Reset flow_mode - the streams controller will set it if flow mode is used.
1060 queue.flow_mode = False
1061 player_media = await self.player_media_from_queue_item(queue_item)
1062 # Hold the play action until the player confirms playback so the UI keeps
1063 # showing the command as in progress instead of falling back to a play button
1064 # for the time the player still needs to connect and start. The queue update
1065 # for the new item goes out first, so the item shows while it is starting.
1066 async with self.mass.players.wait_for_player_update(
1067 queue_id,
1068 attribute_name="playback_state",
1069 attribute_value=PlaybackState.PLAYING,
1070 timeout=PLAYBACK_START_TIMEOUT,
1071 ):
1072 await self.mass.players.play_media(queue_id, player_media)
1073 queue.current_index = index
1074 queue.current_item = queue_item
1075 self.signal_update(queue_id)
1076 finally:
1077 self._set_transitioning(queue_id, False)
1078
1079 @api_command("player_queues/transfer", required_scope=Scope.QUEUES_CONTROL)
1080 async def transfer_queue(
1081 self,
1082 source_queue_id: str,
1083 target_queue_id: str,
1084 auto_play: bool | None = None,
1085 ) -> None:
1086 """Transfer queue to another queue."""
1087 if not (source_queue := self.get(source_queue_id)):
1088 raise PlayerUnavailableError(f"Queue {source_queue_id} is not available")
1089 if not (target_queue := self.get(target_queue_id)):
1090 raise PlayerUnavailableError(f"Queue {target_queue_id} is not available")
1091 if auto_play is None:
1092 auto_play = source_queue.state == PlaybackState.PLAYING
1093
1094 target_player = self.mass.players.get_player(target_queue_id)
1095 if target_player is None:
1096 raise PlayerUnavailableError(f"Player {target_queue_id} is not available")
1097 if target_player.state.active_group or target_player.state.synced_to:
1098 # edge case: the user wants to move playback from the group as a whole, to a single
1099 # player in the group or it is grouped and the command targeted at the single player.
1100 # We need to dissolve the group/sync first, and wait for the state to actually
1101 # propagate before we hand the queue over to the target player.
1102 group_id = target_player.state.active_group or target_player.state.synced_to
1103 assert group_id is not None # checked in if condition above
1104 # For an ad-hoc sync group (target is a sync member of a regular leader),
1105 # ungroup the target itself so only it is freed - ungrouping the leader would
1106 # transfer leadership to a remaining member and recurse back into this method.
1107 # For a virtual group player (active_group), release the group so its static
1108 # members are handled correctly.
1109 ungroup_target = (
1110 target_queue_id
1111 if target_player.state.synced_to and not target_player.state.active_group
1112 else group_id
1113 )
1114 async with self.mass.players.wait_for_player_update(
1115 target_queue_id,
1116 attribute_name=(
1117 "active_group" if target_player.state.active_group else "synced_to"
1118 ),
1119 attribute_value=None,
1120 timeout=5,
1121 ):
1122 await self.mass.players.cmd_ungroup(ungroup_target)
1123
1124 # capture source state before stopping (stop resets these)
1125 source_items = self._queue_data[source_queue_id].items
1126 if source_queue.state == PlaybackState.PLAYING:
1127 # use the live playback clock while actively playing
1128 source_resume_pos = int(source_queue.corrected_elapsed_time)
1129 else:
1130 # when not playing the live clock is stale, so use the stored resume position
1131 source_resume_pos = int(source_queue.resume_pos or source_queue.elapsed_time or 0)
1132 source_current_index = source_queue.current_index
1133 source_current_item = source_queue.current_item
1134
1135 # stop the source player synchronously to prevent the async stop from
1136 # clear() racing with the target's sync group formation/protocol switching
1137 if source_queue.state != PlaybackState.IDLE:
1138 await self.stop(source_queue_id)
1139
1140 target_queue.repeat_mode = source_queue.repeat_mode
1141 target_queue.shuffle_enabled = source_queue.shuffle_enabled
1142 target_queue.crossfade_enabled = source_queue.crossfade_enabled
1143 # refresh the derived smart-fades indicator for the target's own config/availability
1144 target_queue.smart_fades_active = self.mass.streams.is_smart_fades_active(target_queue)
1145 target_queue.autoplay_enabled = source_queue.autoplay_enabled
1146 self._queue_data[target_queue_id].source_items = list(
1147 self._queue_data[source_queue_id].source_items
1148 )
1149 target_queue.sources = list(source_queue.sources)
1150 target_queue.is_dynamic = source_queue.is_dynamic
1151 target_queue.smart_shuffle_active = self.is_smart_shuffle_active(target_queue)
1152 self._queue_data[target_queue_id].enqueued_media_items = list(
1153 self._queue_data[source_queue_id].enqueued_media_items
1154 )
1155 target_queue.resume_pos = source_resume_pos
1156 target_queue.current_index = source_current_index
1157 if source_current_item:
1158 target_queue.current_item = source_current_item
1159 target_queue.current_item.queue_id = target_queue_id
1160 self._clear(source_queue_id, skip_stop=True)
1161
1162 await self.load(target_queue_id, source_items, keep_remaining=False, keep_played=False)
1163 for item in source_items:
1164 item.queue_id = target_queue_id
1165 self.update_items(target_queue_id, source_items)
1166 if auto_play:
1167 await self.resume(target_queue_id)
1168
1169 # Interaction with player
1170
1171 async def on_player_register(self, player: Player) -> None:
1172 """Register PlayerQueue for given player/queue id."""
1173 queue_id = player.player_id
1174 queue_data: PlayerQueueData | None = None
1175 # try to restore previous state
1176 try:
1177 if prev_state := await self.mass.cache.get(
1178 key=queue_id,
1179 provider=self.domain,
1180 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1181 ):
1182 prev_items = await self.mass.cache.get(
1183 key=queue_id,
1184 provider=self.domain,
1185 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1186 default=[],
1187 )
1188 queue_data = PlayerQueueData.from_cache(prev_state, prev_items)
1189 except Exception as err:
1190 self.logger.warning(
1191 "Failed to restore the queue(items) for %s - %s",
1192 player.state.name,
1193 str(err),
1194 )
1195 # Reset to clean state on failure
1196 queue_data = None
1197 if queue_data is None:
1198 queue_data = PlayerQueueData(
1199 queue=PlayerQueue(
1200 queue_id=queue_id,
1201 active=False,
1202 display_name=player.state.name,
1203 available=player.state.available,
1204 # Autoplay starts out on for a brand new queue; the player's own Autoplay
1205 # switch owns it from here on (and is restored above for a queue we know)
1206 autoplay_enabled=True,
1207 items=0,
1208 )
1209 )
1210
1211 self._queue_data[queue_id] = queue_data
1212 # always call update to calculate state etc
1213 self.on_player_update(player, {})
1214 self.mass.signal_event(EventType.QUEUE_ADDED, object_id=queue_id, data=queue_data.queue)
1215
1216 def on_player_update(
1217 self,
1218 player: Player,
1219 changed_values: dict[str, tuple[Any, Any]],
1220 ) -> None:
1221 """
1222 Call when a PlayerQueue needs to be updated (e.g. when player updates).
1223
1224 NOTE: This is called every second if the player is playing.
1225 """
1226 if player.type == PlayerType.PROTOCOL:
1227 # protocol players do not have a queue on their own
1228 return
1229 queue_id = player.player_id
1230 if (queue := self.get(queue_id)) is None:
1231 # race condition
1232 return
1233 if player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS):
1234 # do nothing while the announcement is in progress
1235 return
1236 # determine if this queue is currently active for this player
1237 queue.active = player.state.active_source in (queue.queue_id, None)
1238 if not queue.active and self._queue_data[queue_id].prev_state is None:
1239 queue.state = PlaybackState.IDLE
1240 # return early if the queue is not active and we have no previous state
1241 return
1242 if self._queue_data[queue_id].transitioning:
1243 # we're currently transitioning to a new track,
1244 # ignore updates from the player during this time
1245 return
1246 # queue is active and preflight checks passed, update the queue details
1247 self._update_queue_from_player(player)
1248
1249 def on_player_elapsed_time_corrected(self, player: Player) -> None:
1250 """Correct the queue's timing base if the player's real elapsed_time diverged."""
1251 if player.type == PlayerType.PROTOCOL:
1252 return
1253 queue_id = player.player_id
1254 if (queue := self.get(queue_id)) is None:
1255 return
1256 if not queue.active:
1257 return
1258 player_elapsed = player.state.corrected_elapsed_time
1259 if player_elapsed is None:
1260 return
1261 now = time.time()
1262 # queue.elapsed_time is stored in media-time so it can be displayed and
1263 # used as a resume position directly. The player reports stream-time
1264 # (post-atempo), so we scale by the current item's playback_speed.
1265 speed = get_current_playback_speed(queue)
1266 if queue.flow_mode:
1267 # _get_flow_queue_stream_index returns media-time in the current item
1268 # using each playlog entry's recorded speed.
1269 _, elapsed_time = self._get_flow_queue_stream_index(queue, player)
1270 else:
1271 elapsed_time = player_elapsed * speed
1272 if queue.current_item and queue.current_item.streamdetails:
1273 if seek_pos := queue.current_item.streamdetails.seek_position:
1274 elapsed_time += seek_pos
1275 queue.elapsed_time = elapsed_time
1276 queue.elapsed_time_last_updated = now
1277 queue.playback_speed = speed
1278 self.mass.signal_event(
1279 EventType.QUEUE_TIME_UPDATED,
1280 object_id=queue_id,
1281 data=queue.elapsed_time,
1282 )
1283
1284 def on_player_remove(self, player_id: str, permanent: bool) -> None:
1285 """Call when a player is removed from the registry."""
1286 self.mass.streams.audio_processing.clear(player_id)
1287 # cancel any pending play_index calls for this queue to prevent conflicts
1288 self.mass.cancel_timer(f"queue_play_index_{player_id}")
1289 # cancel a pending debounced cache write AND an already-started one, so neither can
1290 # recreate a deleted entry after the player is gone (the timer becomes a task once it fires)
1291 self.mass.cancel_timer(f"save_queue_cache_{player_id}")
1292 self.mass.cancel_task(f"save_queue_cache_{player_id}")
1293 self._set_transitioning(player_id, False)
1294 if permanent:
1295 self.purge_saved_queue(player_id)
1296 self._queue_data.pop(player_id, None)
1297 self._managed_pool.forget(player_id)
1298
1299 def purge_saved_queue(self, queue_id: str) -> None:
1300 """Delete the persisted state and items of the given queue."""
1301 for category in (CACHE_CATEGORY_PLAYER_QUEUE_STATE, CACHE_CATEGORY_PLAYER_QUEUE_ITEMS):
1302 # a removal runs both the player teardown and the config cleanup, so keep the
1303 # delete to one task per category instead of one per caller
1304 self.mass.create_task(
1305 self.mass.cache.delete(
1306 key=queue_id,
1307 provider=self.domain,
1308 category=category,
1309 ),
1310 task_id=f"purge_saved_queue_{queue_id}_{category}",
1311 )
1312
1313 async def load_next_queue_item(
1314 self,
1315 queue_id: str,
1316 current_item_id: str,
1317 ) -> QueueItem:
1318 """
1319 Call when a player wants the next queue item to play.
1320
1321 Raises QueueEmpty if there are no more tracks left.
1322 """
1323 queue = self.get(queue_id)
1324 if not queue:
1325 msg = f"PlayerQueue {queue_id} is not available"
1326 raise PlayerUnavailableError(msg)
1327 cur_index = self.index_by_id(queue_id, current_item_id)
1328 if cur_index is None:
1329 # this is just a guard for bad data
1330 raise QueueEmpty("Invalid item id for queue given.")
1331 next_item: QueueItem | None = None
1332 idx = 0
1333 while True:
1334 next_index = self._get_next_index(queue_id, cur_index + idx)
1335 if next_index is None:
1336 raise QueueEmpty("No more tracks left in the queue.")
1337 queue_item = self.get_item(queue_id, next_index)
1338 if queue_item is None:
1339 raise QueueEmpty("No more tracks left in the queue.")
1340 if idx >= 10:
1341 # we only allow 10 retries to prevent infinite loops
1342 raise QueueEmpty("No more (playable) tracks left in the queue.")
1343 try:
1344 await self._load_item(queue_item, next_index)
1345 # we're all set, this is our next item
1346 next_item = queue_item
1347 break
1348 except ProviderStreamLimitError:
1349 # transient source capacity, do not burn a playable item over it
1350 raise
1351 except MediaNotFoundError, AudioError:
1352 # No stream details found, skip this QueueItem
1353 self.logger.warning(
1354 "Skipping unplayable item %s (%s)", queue_item.name, queue_item.uri
1355 )
1356 queue_item.available = False
1357 idx += 1
1358 if idx != 0:
1359 # we skipped some items, signal a queue items update
1360 self.update_items(queue_id, self._queue_data[queue_id].items)
1361 if next_item is None:
1362 raise QueueEmpty("No more (playable) tracks left in the queue.")
1363
1364 # carry playback_speed forward across consecutive audiobook/podcast items
1365 current_item = self.get_item(queue_id, current_item_id)
1366 if (
1367 current_item
1368 and current_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1369 and next_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1370 ):
1371 next_item.extra_attributes["playback_speed"] = current_item.extra_attributes.get(
1372 "playback_speed", 1.0
1373 )
1374
1375 return next_item
1376
1377 def track_loaded_in_buffer(self, queue_id: str, item_id: str) -> None:
1378 """Call when a player has (started) loading a track in the buffer."""
1379 queue = self.get(queue_id)
1380 if not queue:
1381 msg = f"PlayerQueue {queue_id} is not available"
1382 raise PlayerUnavailableError(msg)
1383 # store the index of the item that is currently (being) loaded in the buffer
1384 # which helps us a bit to determine how far the player has buffered ahead
1385 current_index = self.index_by_id(queue_id, item_id)
1386 queue.index_in_buffer = current_index
1387 self.logger.debug("PlayerQueue %s loaded item %s in buffer", queue.display_name, item_id)
1388 self.signal_update(queue_id)
1389 # preload next streamdetails
1390 self._preload_next_item(queue_id, item_id)
1391 # clean up stale audio buffers for old queue items to prevent memory leaks
1392 if current_index is not None:
1393 self.mass.create_task(self._cleanup_stale_queue_buffers(queue_id, current_index))
1394
1395 def queue_buffer_completed(self, queue_id: str, queue_exhausted: bool) -> None:
1396 """
1397 Call when the flow stream has finished generating all audio data for a queue.
1398
1399 At this point all audio data for the queue has been passed to the encoding pipeline.
1400 The player will go idle once it finishes playing the remaining buffered audio.
1401
1402 We start a background task that waits for the player to go idle and checks if new
1403 items have been added to the queue in the meantime, resuming playback if so.
1404
1405 :param queue_id: The queue ID.
1406 :param queue_exhausted: Whether the flow ended because the queue ran out of items,
1407 as opposed to ending early to restart on a format change or a live item.
1408 """
1409 queue = self.get(queue_id)
1410 if not queue:
1411 return
1412 self.logger.debug("Queue flow buffer completed for %s", queue.display_name)
1413
1414 # capture session_id so we can bail out if playback restarts
1415 queue_data = self._queue_data[queue_id]
1416 original_session_id = queue_data.session_id
1417 # record so player providers can detect flow EOF without an idle report
1418 if original_session_id is not None:
1419 queue_data.flow_buffer_completed = original_session_id
1420 if queue_exhausted:
1421 queue_data.flow_queue_exhausted = original_session_id
1422
1423 async def _resume_on_idle() -> None:
1424 # wait for the player to finish playing the buffered audio and go idle
1425 idle_detected = False
1426 for _ in range(60):
1427 await asyncio.sleep(1)
1428 if not queue.active or queue_data.session_id != original_session_id:
1429 return
1430 if queue.state == PlaybackState.IDLE:
1431 idle_detected = True
1432 break
1433 if not idle_detected:
1434 return
1435 # player went idle, give it a brief moment to settle
1436 await asyncio.sleep(1)
1437 if queue.state != PlaybackState.IDLE or queue_data.session_id != original_session_id:
1438 return
1439 # check if new items were added to the queue after the flow stream ended
1440 if queue.current_index is not None and (
1441 next_item := self.get_next_item(queue_id, queue.current_index)
1442 ):
1443 next_index = self.index_by_id(queue_id, next_item.queue_item_id)
1444 if next_index is not None:
1445 self.logger.info(
1446 "Resuming playback after flow stream completed for %s",
1447 queue.display_name,
1448 )
1449 await self.play_index(queue_id, next_index)
1450
1451 task_id = f"queue_buffer_completed_{queue_id}"
1452 self.mass.create_task(_resume_on_idle(), task_id=task_id)
1453
1454 def flow_stream_finished(self, queue_id: str) -> bool:
1455 """
1456 Return whether the flow stream for the current playback session is fully generated.
1457
1458 Lets player providers detect flow EOF when the device does not report idle
1459 (e.g. a Cast group that underruns the LIVE flow stream and keeps reporting playing).
1460
1461 :param queue_id: The queue ID.
1462 """
1463 queue_data = self.queue_data_or_none(queue_id)
1464 if queue_data is None or queue_data.session_id is None:
1465 return False
1466 return queue_data.flow_buffer_completed == queue_data.session_id
1467
1468 def flow_queue_exhausted(self, queue_id: str, session_id: str) -> bool:
1469 """
1470 Return whether the given flow stream session played the queue to its end.
1471
1472 False while a session is still streaming, and for a flow stream that ended early
1473 to be restarted (a format change or a live item), where the player is expected to
1474 pick up the next stream right away.
1475
1476 :param queue_id: The queue ID.
1477 :param session_id: The stream session to check.
1478 """
1479 queue_data = self.queue_data_or_none(queue_id)
1480 if queue_data is None or queue_data.session_id != session_id:
1481 return False
1482 return queue_data.flow_queue_exhausted == session_id
1483
1484 # Main queue manipulation methods
1485
1486 async def load(
1487 self,
1488 queue_id: str,
1489 queue_items: list[QueueItem],
1490 insert_at_index: int = 0,
1491 keep_remaining: bool = True,
1492 keep_played: bool = True,
1493 shuffle: bool = False,
1494 pin_first: bool = False,
1495 ) -> None:
1496 """
1497 Load new items at index.
1498
1499 - queue_id: id of the queue to process this request.
1500 - queue_items: a list of QueueItems
1501 - insert_at_index: insert the item(s) at this index
1502 - keep_remaining: keep the remaining items after the insert
1503 - shuffle: (re)shuffle the items after insert index
1504 - pin_first: keep the first item at the insert index instead of letting the shuffle
1505 move it; only meaningful together with shuffle
1506 """
1507 prev_items = self._queue_data[queue_id].items[:insert_at_index] if keep_played else []
1508 next_items = queue_items
1509
1510 # if keep_remaining, append the old 'next' items
1511 if keep_remaining:
1512 next_items += self._queue_data[queue_id].items[insert_at_index:]
1513
1514 # we set the original insert order as attribute so we can un-shuffle
1515 for index, item in enumerate(next_items):
1516 item.sort_index += insert_at_index + index
1517 # (re)shuffle the final batch if needed: smart shuffle when enabled, else pure random
1518 if shuffle:
1519 queue = self._queue_data[queue_id].queue
1520 # a user-picked item must stay the one that plays, so hold it out of the shuffle
1521 pinned = next_items[:1] if pin_first else []
1522 shuffled = next_items[1:] if pin_first else next_items
1523 if self._smart_shuffle.is_enabled(queue_id):
1524 shuffled = await self._smart_shuffle.arrange(queue, shuffled)
1525 else:
1526 shuffled = random.sample(shuffled, len(shuffled))
1527 next_items = pinned + shuffled
1528 self.update_items(queue_id, prev_items + next_items)
1529
1530 def update_items(self, queue_id: str, queue_items: list[QueueItem]) -> None:
1531 """Update the existing queue items, mostly caused by reordering."""
1532 self._queue_data[queue_id].items = queue_items
1533 queue = self._queue_data[queue_id].queue
1534 queue.items = len(self._queue_data[queue_id].items)
1535 self.signal_update(queue_id, True)
1536 if (
1537 queue.state == PlaybackState.PLAYING
1538 and queue.index_in_buffer is not None
1539 and queue.index_in_buffer == queue.current_index
1540 ):
1541 # if the queue is playing,
1542 # ensure to (re)queue the next track because it might have changed
1543 # note that we only do this if the player has loaded the current track
1544 # if not, we wait until it has loaded to prevent conflicts
1545 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
1546 self._enqueue_next_item(queue_id, next_item)
1547
1548 # Helper methods
1549
1550 def get_item(self, queue_id: str, item_id_or_index: int | str | None) -> QueueItem | None:
1551 """Get queue item by index or item_id."""
1552 if item_id_or_index is None:
1553 return None
1554 if (queue_data := self._queue_data.get(queue_id)) is None:
1555 return None
1556 queue_items = queue_data.items
1557 if isinstance(item_id_or_index, int) and len(queue_items) > item_id_or_index:
1558 return queue_items[item_id_or_index]
1559 if isinstance(item_id_or_index, str):
1560 return next((x for x in queue_items if x.queue_item_id == item_id_or_index), None)
1561 return None
1562
1563 def signal_update(self, queue_id: str, items_changed: bool = False) -> None:
1564 """Signal state changed of given queue."""
1565 if (queue_data := self._queue_data.get(queue_id)) is None:
1566 return
1567 queue = queue_data.queue
1568 if items_changed:
1569 queue_data.items_cache_dirty = True
1570 self.mass.signal_event(EventType.QUEUE_ITEMS_UPDATED, object_id=queue_id, data=queue)
1571 self.mass.streams.audio_processing.prune(queue_id)
1572 # always send the base event
1573 self.mass.signal_event(EventType.QUEUE_UPDATED, object_id=queue_id, data=queue)
1574 # also signal update to the player itself so it can update its current_media
1575 self.mass.players.trigger_player_update(queue_id)
1576 # persist the (settings-bearing) queue state, debounced so a burst of updates or the
1577 # per-track updates during playback collapse into a single cache write
1578 self.mass.call_later(
1579 QUEUE_CACHE_SAVE_DELAY,
1580 self._save_queue_to_cache,
1581 queue_id,
1582 task_id=f"save_queue_cache_{queue_id}",
1583 )
1584
1585 def index_by_id(self, queue_id: str, queue_item_id: str) -> int | None:
1586 """Get index by queue_item_id."""
1587 if (queue_data := self._queue_data.get(queue_id)) is None:
1588 return None
1589 for index, item in enumerate(queue_data.items):
1590 if item.queue_item_id == queue_item_id:
1591 return index
1592 return None
1593
1594 async def get_tracks_for_playback(self, media_item: MediaItemType) -> list[Track]:
1595 """
1596 Return the playable tracks a media item resolves to, honoring the user's selection prefs.
1597
1598 :param media_item: The media item to resolve to playable tracks.
1599 """
1600 return await self._media_resolver.get_tracks_for_playback(media_item)
1601
1602 async def get_playlist_tracks(
1603 self, playlist: Playlist, start_item: str | None = None, sort_by: str | None = None
1604 ) -> list[PlaylistPlayableItem]:
1605 """
1606 Return the playable tracks for a playlist, honoring the user's selection prefs.
1607
1608 :param playlist: The playlist to resolve.
1609 :param start_item: Optional item URI to start the playlist from.
1610 :param sort_by: Optional sort key for the returned tracks.
1611 """
1612 return await self._media_resolver.get_playlist_tracks(playlist, start_item, sort_by)
1613
1614 async def get_dynamic_source_tracks(self, item: MediaItemType) -> list[Track]:
1615 """
1616 Return a fresh batch of tracks for a dynamic source (a dynamic playlist or radio station).
1617
1618 :param item: The dynamic playlist or radio station to fetch the next batch for.
1619 """
1620 return await self._media_resolver.get_dynamic_source_tracks(item)
1621
1622 def recency_windows(self) -> RecencyWindows:
1623 """Return the configured recency windows (a global setting; used for recency-aware gating)."""
1624 return self._smart_shuffle.windows()
1625
1626 async def player_media_from_queue_item(self, queue_item: QueueItem) -> PlayerMedia:
1627 """
1628 Parse PlayerMedia from QueueItem.
1629
1630 :param queue_item: The queue item to create media from.
1631 """
1632 queue_data = self._queue_data[queue_item.queue_id]
1633 stream_duration: int | None = None
1634 if queue_item.streamdetails:
1635 # prefer netto duration
1636 duration = queue_item.streamdetails.duration or queue_item.duration
1637 if duration and queue_item.streamdetails.seek_position:
1638 # the audio handed to the player starts at the seek position, so it is
1639 # shorter than the media item itself. seeking to (or past) the end
1640 # leaves no stream to describe, so the full length is kept instead.
1641 remaining = int(duration - queue_item.streamdetails.seek_position)
1642 stream_duration = remaining if remaining > 0 else None
1643 else:
1644 duration = queue_item.duration
1645 if queue_data.session_id is None:
1646 raise InvalidDataError("Queue session_id is None")
1647 media = PlayerMedia(
1648 uri=queue_item.uri,
1649 media_type=queue_item.media_type,
1650 title=queue_item.name,
1651 image_url=MASS_LOGO_ONLINE,
1652 duration=duration,
1653 stream_duration=stream_duration,
1654 source_id=queue_item.queue_id,
1655 queue_item_id=queue_item.queue_item_id,
1656 queue_session_id=queue_data.session_id,
1657 custom_data={
1658 "original_uri": queue_item.uri,
1659 },
1660 )
1661 if queue_item.media_item:
1662 media.title = queue_item.media_item.name
1663 media.artist = getattr(queue_item.media_item, "artist_str", "")
1664 media.album = (
1665 album.name if (album := getattr(queue_item.media_item, "album", None)) else ""
1666 )
1667 if queue_item.image:
1668 # the image format needs to be 512x512 jpeg for maximum compatibility with players
1669 # we prefer the imageproxy on the streamserver here because this request is sent
1670 # to the player itself which may not be able to reach the regular webserver
1671 media.image_url = self.mass.metadata.get_image_url(
1672 queue_item.image, size=512, image_format="jpeg", prefer_stream_server=True
1673 )
1674 return media
1675
1676 def get_next_item(self, queue_id: str, cur_index: int | str) -> QueueItem | None:
1677 """Return next QueueItem for given queue."""
1678 index: int
1679 if isinstance(cur_index, str):
1680 resolved_index = self.index_by_id(queue_id, cur_index)
1681 if resolved_index is None:
1682 return None # guard
1683 index = resolved_index
1684 else:
1685 index = cur_index
1686 # At this point index is guaranteed to be int
1687 for skip in range(5):
1688 if (next_index := self._get_next_index(queue_id, index + skip)) is None:
1689 break
1690 next_item = self.get_item(queue_id, next_index)
1691 if next_item is None:
1692 continue
1693 if not next_item.available:
1694 # ensure that we skip unavailable items (set by load_next track logic)
1695 continue
1696 return next_item
1697 return None
1698
1699 def store_sources(self, queue: PlayerQueue, items: list[MediaItemType]) -> None:
1700 """
1701 Hold the queue's full dynamic-source items server-side and project them onto `sources`.
1702
1703 :param queue: The queue whose sources are being set.
1704 :param items: The full source media items; an empty list clears the queue's sources.
1705 """
1706 self._queue_data[queue.queue_id].source_items = items
1707 # keep every occurrence server-side (a source added more than once weights it up in the
1708 # managed pool), but expose only the distinct container sources on the wire for clients to
1709 # show. Individual items (tracks, live radio streams, podcast episodes, ...) are omitted; see
1710 # `_WIRE_SOURCE_MEDIA_TYPES`. Autoplay/pool refill reads the full `source_items` above, not
1711 # this projected list, so it is unaffected.
1712 seen: set[str] = set()
1713 sources: list[ItemMapping] = []
1714 for item in items:
1715 if item.media_type not in _WIRE_SOURCE_MEDIA_TYPES and not is_dynamic_source(item):
1716 continue
1717 mapping = ItemMapping.from_item(item)
1718 if mapping.uri and mapping.uri in seen:
1719 continue
1720 if mapping.uri:
1721 seen.add(mapping.uri)
1722 sources.append(mapping)
1723 queue.sources = sources
1724 # release any materialized finite-source state whose source is no longer present
1725 self._managed_pool.retain(
1726 queue.queue_id, {item.uri for item in items if item.uri is not None}
1727 )
1728
1729 async def _save_queue_to_cache(self, queue_id: str) -> None:
1730 """Persist the queue's state (and its items when changed) to the cache."""
1731 if (queue_data := self._queue_data.get(queue_id)) is None:
1732 return
1733 try:
1734 # persistent so a cache clear/reset does not wipe the user's queues; the default
1735 # expiration still applies but is refreshed on every write. Skip the state write when its
1736 # persist-worthy content is unchanged (i.e. only playback progress advanced).
1737 state = queue_data.to_cache()
1738 significant = queue_data.cache_significant(state)
1739 if significant != queue_data.last_saved_state:
1740 await self.mass.cache.set(
1741 key=queue_id,
1742 data=state,
1743 provider=self.domain,
1744 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1745 persistent=True,
1746 )
1747 queue_data.last_saved_state = significant
1748 if queue_data.items_cache_dirty:
1749 # only cache items with a valid media_item
1750 await self.mass.cache.set(
1751 key=queue_id,
1752 data=queue_data.items_to_cache(),
1753 provider=self.domain,
1754 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1755 persistent=True,
1756 )
1757 queue_data.items_cache_dirty = False
1758 except Exception as err:
1759 self.logger.warning("Failed to persist the queue for %s - %s", queue_id, err)
1760
1761 def _check_player_permission(self, queue_id: str) -> None:
1762 """
1763 Check if the current user has permission to control this player/queue.
1764
1765 :param queue_id: The queue/player ID to check access for.
1766 :raises InsufficientPermissions: If the user lacks access.
1767 """
1768 current_user = get_current_user()
1769 if (
1770 current_user
1771 and current_user.player_filter
1772 and queue_id not in current_user.player_filter
1773 ):
1774 msg = f"{current_user.username} does not have access to player {queue_id}"
1775 raise InsufficientPermissions(msg)
1776
1777 @handle_play_action
1778 async def _handle_play(self, queue_id: str) -> None:
1779 """Handle play without acquiring the queue lock."""
1780 queue_player = self.mass.players.get_player(queue_id, True)
1781 if queue_player is None:
1782 raise PlayerUnavailableError(f"Player {queue_id} is not available")
1783 if (queue := self.get(queue_id)) and queue.active and queue.state == PlaybackState.PAUSED:
1784 # forward the actual play/unpause command to the player,
1785 # holding the action until the player confirms it resumed playback
1786 async with self.mass.players.wait_for_player_update(
1787 queue_id,
1788 attribute_name="playback_state",
1789 attribute_value=PlaybackState.PLAYING,
1790 timeout=PLAYBACK_START_TIMEOUT,
1791 ):
1792 await queue_player.play()
1793 return
1794 # player is not paused, perform resume instead
1795 await self.resume(queue_id)
1796
1797 def _set_transitioning(self, queue_id: str, value: bool) -> None:
1798 """Mark (or clear) whether a queue is mid-transition (no-op if it is not registered)."""
1799 if (queue_data := self._queue_data.get(queue_id)) is not None:
1800 queue_data.transitioning = value
1801
1802 def _clear(self, queue_id: str, skip_stop: bool = False) -> None:
1803 """Drop the queue's items and playback position, leaving user settings untouched."""
1804 queue = self._queue_data[queue_id].queue
1805 self.mass.streams.audio_processing.clear(queue_id)
1806 self.store_sources(queue, [])
1807 if queue.is_dynamic:
1808 # Dynamic sources impose shuffle, so clearing the source clears that shuffle too.
1809 queue.shuffle_enabled = False
1810 queue.is_dynamic = False
1811 # dropping the dynamic source changes what smart shuffle resolves to, so the derived
1812 # flag has to follow or clients keep showing a smart mix on a plain queue
1813 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1814 queue.ended = False
1815 if queue.state != PlaybackState.IDLE and not skip_stop:
1816 self.mass.create_task(self.stop(queue_id))
1817 queue.current_index = None
1818 queue.current_item = None
1819 queue.elapsed_time = 0
1820 queue.elapsed_time_last_updated = time.time()
1821 queue.index_in_buffer = None
1822 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
1823 self.update_items(queue_id, [])
1824
1825 def _reset_shuffle(self, queue_id: str) -> None:
1826 """Switch shuffle off."""
1827 queue = self._queue_data[queue_id].queue
1828 if not queue.shuffle_enabled:
1829 return
1830 queue.shuffle_enabled = False
1831 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1832 self.signal_update(queue_id)
1833
1834 async def _apply_shuffle(
1835 self, queue_id: str, option: QueueOption, shuffle: bool | None
1836 ) -> None:
1837 """
1838 Settle the queue's shuffle state for a play command before its items are resolved.
1839
1840 :param queue_id: The queue the media is played on.
1841 :param option: The enqueue option this command resolved to.
1842 :param shuffle: The state to put the queue's shuffle in; None to leave it as it is.
1843 """
1844 queue = self._queue_data[queue_id].queue
1845 if queue.is_dynamic and option in (
1846 QueueOption.PLAY,
1847 QueueOption.REPLACE,
1848 QueueOption.REPLACE_NEXT,
1849 ):
1850 # These are the options that replace the queue's sources, so the smart mix may be on
1851 # its way out - and its shuffle is never the user's own (a dynamic queue's toggle is
1852 # locked), so it must not outlive the source that imposed it. Recorded directly
1853 # because set_shuffle refuses a queue that is still a smart mix, and the items are
1854 # resolved against this flag. The state is provisional until the sources are known:
1855 # `_enter_dynamic_mode` forces shuffle back on if the queue stays dynamic.
1856 if option == QueueOption.REPLACE_NEXT:
1857 # staging leaves the shuffle the user chose alone, so it never carries a request
1858 # of its own to honour here
1859 queue.shuffle_enabled = False
1860 else:
1861 queue.shuffle_enabled = bool(shuffle)
1862 return
1863 if shuffle is None or option not in (QueueOption.PLAY, QueueOption.REPLACE):
1864 # nothing to settle: the media brings no order of its own to protect, or the option
1865 # only stages items for later and leaves the queue's shuffle state alone
1866 return
1867 if queue.shuffle_enabled == shuffle:
1868 return
1869 # routed through set_shuffle so switching shuffle off also restores the order of
1870 # the items that stay in the queue: a play keeps them, and a tail left in shuffled
1871 # order behind a queue that now reads unshuffled would contradict its own flag
1872 await self.set_shuffle(queue_id, shuffle)
1873