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