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