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