/
/
1"""
2Queue loading for the Player Queues controller.
3
4Applies the enqueue option (play/replace/next/add) to a batch of resolved items, loads a single
5media item into the queue, resumes from the play-log when the queue is empty, computes the next
6index, and refills the queue (dynamic managed-pool fill and autoplay fill). Owns no per-queue state;
7it is mixed into the controller and reads/mutates the controller's `PlayerQueueData` records.
8"""
9# ruff: noqa: PLR0915
10
11from __future__ import annotations
12
13import random
14from contextlib import suppress
15from typing import TYPE_CHECKING, cast
16
17from music_assistant_models.enums import (
18 MediaType,
19 PlaybackState,
20 QueueOption,
21 RepeatMode,
22)
23from music_assistant_models.errors import (
24 InvalidDataError,
25 MediaNotFoundError,
26 MusicAssistantError,
27 PlayerUnavailableError,
28)
29from music_assistant_models.media_items import (
30 Album,
31 Audiobook,
32 BrowseFolder,
33 ItemMapping,
34 MediaItemType,
35 PlayableMediaItemType,
36 PodcastEpisode,
37 Track,
38 UniqueList,
39 media_from_dict,
40)
41
42from music_assistant.constants import ATTR_ANNOUNCEMENT_IN_PROGRESS
43from music_assistant.controllers.player_queues.autoplay import (
44 AUTOPLAY_EXCLUDED_MEDIA_TYPES,
45 AUTOPLAY_SERIES_MEDIA_TYPES,
46 AutoplayMode,
47)
48from music_assistant.controllers.player_queues.base import _PlayerQueuesBase
49from music_assistant.controllers.player_queues.constants import (
50 CONF_DEFAULT_ENQUEUE_OPTION_LIVE_SOURCES,
51 MANAGED_POOL_MAX,
52 ORDERED_MEDIA_TYPES,
53 PROBED_DURATION_MEDIA_TYPES,
54)
55from music_assistant.controllers.player_queues.helpers import (
56 build_queue_item,
57 handle_play_action,
58 has_dynamic_source,
59 is_dynamic_source,
60)
61from music_assistant.controllers.player_queues.managed_pool import gate_tracks
62from music_assistant.controllers.webserver.helpers.auth_middleware import (
63 get_current_user,
64 set_current_user,
65)
66from music_assistant.helpers.audio import get_probed_duration, store_probed_duration
67from music_assistant.helpers.throttle_retry import BYPASS_THROTTLER
68from music_assistant.models.music_provider import MusicProvider
69
70if TYPE_CHECKING:
71 from music_assistant_models.media_items.metadata import MediaItemImage
72 from music_assistant_models.queue_item import QueueItem
73
74 from music_assistant.providers.radio_playlist import RadioPlaylistProvider
75
76
77class QueueLoaderMixin(_PlayerQueuesBase):
78 """Load items into a queue: apply the enqueue option, resolve single items, refill the pool."""
79
80 async def _enqueue_with_option(
81 self,
82 queue_id: str,
83 queue_items: list[QueueItem],
84 option: QueueOption | None,
85 pin_first: bool = False,
86 ) -> None:
87 """
88 Load queue items into the queue according to the given enqueue option.
89
90 :param queue_id: The queue to load the items into.
91 :param queue_items: The items to load.
92 :param option: The enqueue option to apply.
93 :param pin_first: The first item was explicitly picked by the user (a start_item), so it
94 must keep its position when the batch is shuffled instead of being moved at random.
95 """
96 queue = self._queue_data[queue_id].queue
97 # A queue that played to its end is finished, so anything enqueued onto it starts a fresh
98 # queue rather than stacking onto the items that already played. Only an explicit ADD keeps
99 # them: there the added items continue the queue from where it ended, and the index is moved
100 # onto the first of them below so pressing play starts there instead of replaying the last
101 # item. ADD never starts playback by itself.
102 continues_ended_queue = queue.ended and option == QueueOption.ADD
103 items_before_add = len(self._queue_data[queue_id].items)
104 if queue.ended and not continues_ended_queue and option != QueueOption.REPLACE:
105 # mechanical clear: the shuffle state for this batch was already settled by the caller.
106 # Replace is exempt: it swaps the whole queue below without ever emptying it.
107 self._clear(queue_id, skip_stop=True)
108 if queue.state in (PlaybackState.PLAYING, PlaybackState.PAUSED):
109 cur_index = (
110 queue.index_in_buffer
111 if queue.index_in_buffer is not None
112 else (queue.current_index if queue.current_index is not None else 0)
113 )
114 else:
115 cur_index = queue.current_index or 0
116 insert_at_index = cur_index + 1
117 shuffle = queue.shuffle_enabled and len(queue_items) > 1
118 # a user-picked start item must be the one that actually starts playing, so keep it in
119 # front of the shuffled rest instead of letting the shuffle move it to a random slot
120 pin_first = pin_first and shuffle
121
122 # handle replace: swap the queue's contents for the new items in one step
123 if option == QueueOption.REPLACE:
124 # Release the audio the outgoing items hold while they are still on the queue: the
125 # track being started needs their source slot, and once they are swapped out nothing
126 # reaches them any more.
127 await self._cleanup_queue_audio_data(queue_id)
128 # the player is still on the old index, so drop it: the swap would otherwise hand it a
129 # "next" item taken from the new list at that position. play_index sets the real one.
130 queue.index_in_buffer = None
131 # playback starts over below, and play_index reads this to decide whether to honour a
132 # stored resume position
133 queue.ended = False
134 if pin_first:
135 await self._load_pinned_first(
136 queue_id,
137 queue_items,
138 insert_at_index=0,
139 keep_remaining=False,
140 keep_played=False,
141 )
142 else:
143 await self.load(
144 queue_id,
145 queue_items=queue_items,
146 keep_remaining=False,
147 keep_played=False,
148 shuffle=shuffle,
149 )
150 await self.play_index(queue_id, 0)
151 return
152 # handle next: add item(s) in the index next to the playing/loaded/buffered index
153 if option == QueueOption.NEXT:
154 if shuffle:
155 # honour "play next" under shuffle: the first new item goes right after the
156 # buffered index so it plays next, the rest of the batch is shuffled into the tail
157 # behind it. insert_at_index is the first un-buffered slot, so the track the player
158 # already prepared for crossfade is left untouched.
159 await self._load_pinned_first(queue_id, queue_items, insert_at_index)
160 else:
161 await self.load(
162 queue_id,
163 queue_items=queue_items,
164 insert_at_index=insert_at_index,
165 shuffle=shuffle,
166 )
167 self._ensure_current_index(queue_id)
168 return
169 if option == QueueOption.REPLACE_NEXT:
170 if pin_first:
171 await self._load_pinned_first(
172 queue_id, queue_items, insert_at_index, keep_remaining=False
173 )
174 else:
175 await self.load(
176 queue_id,
177 queue_items=queue_items,
178 insert_at_index=insert_at_index,
179 keep_remaining=False,
180 shuffle=shuffle,
181 )
182 self._ensure_current_index(queue_id)
183 return
184 # handle play: replace current loaded/playing index with new item(s)
185 if option == QueueOption.PLAY:
186 # an idle/empty queue has no current item to insert after, so insert at and
187 # start from the very first index instead of skipping past it
188 play_at_index = 0 if queue.current_index is None else insert_at_index
189 if pin_first:
190 await self._load_pinned_first(queue_id, queue_items, play_at_index)
191 else:
192 await self.load(
193 queue_id,
194 queue_items=queue_items,
195 insert_at_index=play_at_index,
196 shuffle=shuffle,
197 )
198 next_index = min(play_at_index, len(self._queue_data[queue_id].items) - 1)
199 await self.play_index(queue_id, next_index)
200 return
201 # handle add: add/append item(s) to the remaining queue items
202 if option == QueueOption.ADD:
203 # When shuffling, mix the new items into the not-yet-played tail. While playing,
204 # keep the item right after the buffered one in place: it has already been enqueued
205 # to the player (and prepared for crossfade), so reshuffling it would swap the
206 # upcoming track underneath the player and cause an abrupt, non-crossfaded switch.
207 if not queue.shuffle_enabled:
208 add_at_index = len(self._queue_data[queue_id].items) + 1
209 elif queue.state in (PlaybackState.PLAYING, PlaybackState.PAUSED):
210 add_at_index = insert_at_index + 1
211 else:
212 add_at_index = insert_at_index
213 await self.load(
214 queue_id=queue_id,
215 queue_items=queue_items,
216 insert_at_index=add_at_index,
217 shuffle=queue.shuffle_enabled,
218 )
219 if continues_ended_queue:
220 self._continue_ended_queue(queue_id, items_before_add)
221 return
222 self._ensure_current_index(queue_id)
223
224 async def _load_pinned_first(
225 self,
226 queue_id: str,
227 queue_items: list[QueueItem],
228 insert_at_index: int,
229 keep_remaining: bool = True,
230 keep_played: bool = True,
231 ) -> None:
232 """
233 Insert the first item at the given index and shuffle the rest of the batch behind it.
234
235 :param queue_id: The queue to load the items into.
236 :param queue_items: The items to load; the first one keeps the given index.
237 :param insert_at_index: The index to place the first item at.
238 :param keep_remaining: Keep the queue's existing items from the insert index onwards.
239 :param keep_played: Keep the queue's existing items before the insert index.
240 """
241 # a single load, so the queue is never published holding just the pinned item
242 await self.load(
243 queue_id,
244 queue_items=queue_items,
245 insert_at_index=insert_at_index,
246 keep_remaining=keep_remaining,
247 keep_played=keep_played,
248 shuffle=True,
249 pin_first=True,
250 )
251
252 def _ensure_current_index(self, queue_id: str) -> None:
253 """
254 Point the current index at the first item when the queue does not have one yet.
255
256 NEXT/ADD/REPLACE_NEXT stage items without starting playback; on an empty queue there is no
257 current index, so set it to the first item to give the queue a current item. A queue that
258 already has content keeps its current index untouched (its items are inserted after it).
259
260 :param queue_id: The queue to update.
261 """
262 queue = self._queue_data[queue_id].queue
263 if queue.current_index is not None:
264 return
265 queue.current_index = 0
266 queue.current_item = self.get_item(queue_id, 0)
267 self.signal_update(queue_id)
268
269 def _continue_ended_queue(self, queue_id: str, first_added_index: int) -> None:
270 """
271 Point a finished queue at the first item just added to it, without starting playback.
272
273 The items that already played are kept, so the queue is no longer finished but its position
274 still sits on its old last item. Moving it onto the added items is what makes a play press
275 start there rather than replay the item the queue ended on.
276
277 :param queue_id: The queue that was added to.
278 :param first_added_index: Index of the first of the added items.
279 """
280 queue = self._queue_data[queue_id].queue
281 queue.ended = False
282 if (current_item := self.get_item(queue_id, first_added_index)) is None:
283 return
284 queue.current_index = first_added_index
285 queue.current_item = current_item
286 # ending the queue cleared the next item; refresh it so a batch of added items reports
287 # what follows instead of looking like there is nothing after the first one
288 queue.next_item = self.get_next_item(queue_id, first_added_index)
289 self.signal_update(queue_id)
290
291 async def _load_item(
292 self,
293 queue_item: QueueItem,
294 next_index: int | None,
295 is_start: bool = False,
296 seek_position: int = 0,
297 fade_in: bool = False,
298 ) -> None:
299 """Try to load the stream details for the given queue item."""
300 queue_id = queue_item.queue_id
301 queue = self._queue_data[queue_id].queue
302
303 # we use a contextvar to bypass the throttler for this asyncio task/context
304 # this makes sure that playback has priority over other requests that may be
305 # happening in the background
306 BYPASS_THROTTLER.set(True)
307
308 self.logger.debug(
309 "(pre)loading (next) item for queue %s...",
310 queue.display_name,
311 )
312
313 if not queue_item.available:
314 raise MediaNotFoundError(f"Item {queue_item.uri} is not available")
315
316 # work out if we are playing an album and if we should prefer album
317 # loudness
318 next_track_from_same_album = (
319 next_index is not None
320 and (next_item := self.get_item(queue_id, next_index))
321 and (
322 queue_item.media_item
323 and hasattr(queue_item.media_item, "album")
324 and queue_item.media_item.album
325 and next_item.media_item
326 and hasattr(next_item.media_item, "album")
327 and next_item.media_item.album
328 and queue_item.media_item.album.item_id == next_item.media_item.album.item_id
329 )
330 )
331 current_index = self.index_by_id(queue_id, queue_item.queue_item_id)
332 if current_index is None:
333 previous_track_from_same_album = False
334 else:
335 previous_index = max(current_index - 1, 0)
336 previous_track_from_same_album = (
337 previous_index > 0
338 and (previous_item := self.get_item(queue_id, previous_index)) is not None
339 and previous_item.media_item is not None
340 and hasattr(previous_item.media_item, "album")
341 and previous_item.media_item.album is not None
342 and queue_item.media_item is not None
343 and hasattr(queue_item.media_item, "album")
344 and queue_item.media_item.album is not None
345 and queue_item.media_item.album.item_id == previous_item.media_item.album.item_id
346 )
347 playing_album_tracks = next_track_from_same_album or previous_track_from_same_album
348 if queue_item.media_item and isinstance(queue_item.media_item, Track):
349 album = queue_item.media_item.album
350 # prefer the full library media item so we have all metadata and provider(quality) info
351 # always request the full library item as there might be other qualities available
352 if library_item := await self.mass.music.get_library_item_by_prov_id(
353 queue_item.media_item.media_type,
354 queue_item.media_item.item_id,
355 queue_item.media_item.provider,
356 ):
357 queue_item.media_item = cast("Track", library_item)
358 elif not queue_item.media_item.image or queue_item.media_item.provider.startswith(
359 "ytmusic"
360 ):
361 # Youtube Music has poor thumbs by default, so we always fetch the full item
362 # this also catches the case where they have an unavailable item in a listing
363 fetched_item = await self.mass.music.get_item_by_uri(queue_item.uri)
364 queue_item.media_item = cast("Track", fetched_item)
365
366 # ensure we got the full (original) album set
367 if album and (
368 library_album := await self.mass.music.get_library_item_by_prov_id(
369 album.media_type,
370 album.item_id,
371 album.provider,
372 )
373 ):
374 queue_item.media_item.album = cast("Album", library_album)
375 elif album:
376 # Restore original album if we have no better alternative from the library
377 queue_item.media_item.album = album
378 # prefer album image over track image
379 if queue_item.media_item.album and queue_item.media_item.album.image:
380 org_images: list[MediaItemImage] = queue_item.media_item.metadata.images or []
381 queue_item.media_item.metadata.images = UniqueList(
382 [
383 queue_item.media_item.album.image,
384 *org_images,
385 ]
386 )
387 if is_start:
388 # a track skip should hand its source slot to the item the user is starting
389 await self._abort_superseded_source_buffers(queue_item)
390
391 # Fetch streamdetails (reuses existing if buffer is still valid for the seek).
392 queue_item.streamdetails = await self.mass.streams.audio.get_stream_details(
393 queue_item=queue_item,
394 seek_position=seek_position,
395 fade_in=fade_in,
396 prefer_album_loudness=bool(playing_album_tracks),
397 )
398 # update queue_item.duration from streamdetails if we got a better value
399 self._apply_probed_duration(queue_item)
400
401 # pre-initialize the AudioBuffer so audio is ready
402 # when the player requests it. For the current/first track this ensures
403 # immediate playback start. For preloaded next tracks we skip this and
404 # initialize the buffer ~30s before the current track ends instead.
405 # AudioSource items are realtime/live and bypass the AudioBuffer.
406 if is_start and queue_item.streamdetails.media_type != MediaType.AUDIO_SOURCE:
407 await self.mass.streams.audio.get_audio_buffer(
408 queue_item,
409 seek_position_ms=int(seek_position * 1000),
410 reason="prepare",
411 )
412 # the first chunk is in, so the source has been probed and a duration the
413 # provider did not report is known before playback starts
414 self._apply_probed_duration(queue_item)
415
416 def _apply_probed_duration(self, queue_item: QueueItem) -> None:
417 """
418 Apply a duration determined while streaming to the queue item and its media item.
419
420 :param queue_item: The queue item whose streamdetails to take the duration from.
421 """
422 streamdetails = queue_item.streamdetails
423 if streamdetails is None or not streamdetails.duration:
424 return
425 duration = int(streamdetails.duration)
426 if not self._set_missing_duration(queue_item, duration):
427 return
428 if uri := getattr(queue_item.media_item, "uri", None):
429 # store it so listings and later playbacks have it up front
430 self.mass.create_task(store_probed_duration(self.mass, uri, duration))
431
432 async def _restore_probed_duration(self, queue_item: QueueItem) -> None:
433 """
434 Apply the duration determined during an earlier playback to an item that lacks one.
435
436 :param queue_item: The queue item to fill the duration of.
437 """
438 if queue_item.media_type not in PROBED_DURATION_MEDIA_TYPES:
439 return
440 if not (uri := getattr(queue_item.media_item, "uri", None)):
441 return
442 if queue_item.duration and getattr(queue_item.media_item, "duration", None):
443 return
444 if duration := await get_probed_duration(self.mass, uri):
445 self._set_missing_duration(queue_item, duration)
446
447 def _set_missing_duration(self, queue_item: QueueItem, duration: int) -> bool:
448 """
449 Fill in the duration of a queue item and its media item, leaving known ones alone.
450
451 :param queue_item: The queue item to fill the duration of.
452 :param duration: The duration in seconds.
453 :return: True if the item (or its media item) did not have a duration yet.
454 """
455 if queue_item.media_type not in PROBED_DURATION_MEDIA_TYPES:
456 return False
457 media_item = queue_item.media_item
458 # an ItemMapping or any other reference without a duration is left untouched
459 media_item_duration = getattr(media_item, "duration", None)
460 if queue_item.duration and media_item_duration != 0:
461 return False
462 if not queue_item.duration:
463 queue_item.duration = duration
464 if media_item_duration == 0:
465 media_item.duration = duration # type: ignore[union-attr]
466 self.signal_update(queue_item.queue_id, items_changed=True)
467 return True
468
469 def _get_next_index(
470 self,
471 queue_id: str,
472 cur_index: int | None,
473 is_skip: bool = False,
474 allow_repeat: bool = True,
475 ) -> int | None:
476 """
477 Return the next index for the queue, accounting for repeat settings.
478
479 Will return None if there are no (more) items in the queue.
480 """
481 queue = self._queue_data[queue_id].queue
482 queue_items = self._queue_data[queue_id].items
483 if not queue_items or cur_index is None:
484 # queue is empty
485 return None
486 # handle repeat single track
487 if queue.repeat_mode == RepeatMode.ONE and not is_skip:
488 return cur_index if allow_repeat else None
489 # handle cur_index is last index of the queue
490 if cur_index >= (len(queue_items) - 1):
491 if allow_repeat and queue.repeat_mode == RepeatMode.ALL:
492 # if repeat all is enabled, we simply start again from the beginning
493 return 0
494 return None
495 # all other: just the next index
496 return cur_index + 1
497
498 async def _fill_dynamic_tracks(self, queue_id: str) -> None:
499 """Fill a Queue with (additional) tracks from its dynamic sources."""
500 self.logger.debug(
501 "Filling dynamic tracks for queue %s",
502 queue_id,
503 )
504 queue_data = self._queue_data[queue_id]
505 queue = queue_data.queue
506 # restore the queue owner's user context so provider filters are respected during this
507 # background refill (dynamic-playlist generation honours the current user)
508 playback_user = (
509 await self.mass.webserver.auth.get_user(queue_data.userid)
510 if queue_data.userid
511 else None
512 )
513 set_current_user(playback_user)
514 # Top up from the queue's dynamic sources (dynamic playlists and any mixed-in finite items),
515 # weighted per source and recency-gated. fill() already sizes the batch to the pool target;
516 # the tail cap below is a defensive ceiling so the unplayed tail never grows past
517 # MANAGED_POOL_MAX.
518 pool_tracks = await self._managed_pool.fill(queue_id, is_initial=False)
519 # keep the unplayed tail within the bounded pool size (no current_index => nothing played yet)
520 played = 0 if queue.current_index is None else queue.current_index + 1
521 unplayed = max(len(self._queue_data[queue_id].items) - played, 0)
522 headroom = max(MANAGED_POOL_MAX - unplayed, 0)
523 queue_items = [build_queue_item(queue_id, x) for x in pool_tracks[:headroom] if x.available]
524 if not queue_items:
525 return
526 await self.load(
527 queue_id,
528 queue_items,
529 insert_at_index=len(self._queue_data[queue_id].items) + 1,
530 )
531
532 async def _fill_autoplay_tracks(self, queue_id: str) -> None:
533 """
534 Append more items to a queue that is running low, based on what is ending.
535
536 Autoplay is a single "keep going" switch; what it appends is decided by the media type
537 of the queue's last item, since that is the item the appended items follow.
538 """
539 queue = self.get(queue_id)
540 if queue is None or not queue.autoplay_enabled:
541 return
542 queue_data = self._queue_data[queue_id]
543 if not queue_data.items:
544 return
545 last_item = queue_data.items[-1]
546 if last_item.media_type in AUTOPLAY_EXCLUDED_MEDIA_TYPES:
547 return
548 # Restore the queue owner's user context so provider filters, library access and
549 # resume positions are respected during this background refill, mirroring
550 # _fill_dynamic_tracks.
551 playback_user = (
552 await self.mass.webserver.auth.get_user(queue_data.userid)
553 if queue_data.userid
554 else None
555 )
556 set_current_user(playback_user)
557 if last_item.media_type in AUTOPLAY_SERIES_MEDIA_TYPES:
558 await self._fill_autoplay_next_in_series(queue_id, last_item)
559 return
560 await self._fill_autoplay_music_tracks(queue_id)
561
562 async def _fill_autoplay_next_in_series(self, queue_id: str, last_item: QueueItem) -> None:
563 """
564 Append the episode/book that follows the queue's last item, if there is one.
565
566 Nothing is appended for the last episode of a podcast or a book without a next one in
567 its collection, so the queue simply ends there.
568
569 :param queue_id: The queue to append to.
570 :param last_item: The queue's last item, an audiobook or podcast episode.
571 """
572 queue_data = self._queue_data[queue_id]
573 media_item = last_item.media_item
574 next_item: PodcastEpisode | Audiobook | None
575 try:
576 if isinstance(media_item, PodcastEpisode):
577 next_item = await self._media_resolver.get_next_podcast_episode(
578 media_item, userid=queue_data.userid
579 )
580 elif isinstance(media_item, Audiobook):
581 next_item = await self._media_resolver.get_next_audiobook(
582 media_item, userid=queue_data.userid
583 )
584 else:
585 return
586 except MusicAssistantError as err:
587 self.logger.warning(
588 "Autoplay failed to fetch the item following %s: %s", last_item.name, err
589 )
590 return
591 if next_item is None or not next_item.available:
592 self.logger.debug("Autoplay found nothing to play after %s", last_item.name)
593 return
594 if any(
595 item.media_item and item.media_item.uri == next_item.uri for item in queue_data.items
596 ):
597 # already queued (e.g. the user added it themselves), so there is nothing to do
598 return
599 await self.load(
600 queue_id,
601 [build_queue_item(queue_id, next_item)],
602 insert_at_index=len(queue_data.items) + 1,
603 )
604
605 async def _fill_autoplay_music_tracks(self, queue_id: str) -> None:
606 """Fill a Queue with additional tracks based on the configured Autoplay mode."""
607 queue = self.get(queue_id)
608 if queue is None:
609 return
610 queue_data = self._queue_data[queue_id]
611 if not queue_data.enqueued_media_items:
612 # the music refill needs what the user enqueued as its seed
613 return
614 mode = self._autoplay.resolve_mode(queue_id)
615 self.logger.debug(
616 "Filling autoplay tracks (mode: %s) for queue %s", mode.value, queue.display_name
617 )
618 existing_tracks = {
619 item.media_item
620 for item in self._queue_data[queue_id].items
621 if isinstance(item.media_item, Track)
622 }
623 try:
624 if mode == AutoplayMode.PLAYLIST:
625 tracks = await self._autoplay.get_playlist_tracks(queue, existing_tracks)
626 elif mode == AutoplayMode.LIBRARY:
627 tracks = await self._autoplay.get_library_tracks(queue, existing_tracks)
628 elif mode == AutoplayMode.SIMILAR:
629 tracks = await self._get_similar_tracks(
630 queue_id, seed_items=queue_data.enqueued_media_items
631 )
632 else:
633 # AUTO: try similar tracks first, fall back to the library mix. The similar
634 # fetch raises when no provider can supply base/similar tracks, so suppress
635 # that here to make sure the library fallback still runs.
636 tracks = []
637 with suppress(MusicAssistantError):
638 tracks = await self._get_similar_tracks(
639 queue_id, seed_items=queue_data.enqueued_media_items
640 )
641 if not tracks:
642 tracks = await self._autoplay.get_library_tracks(queue, existing_tracks)
643 except MusicAssistantError as err:
644 self.logger.warning(
645 "Autoplay failed to fetch tracks for queue %s: %s", queue.display_name, err
646 )
647 return
648 # route the autoplay batch through the recency engine so a recently-heard track isn't
649 # immediately re-added (ungated fallback keeps autoplay going if everything is recent)
650 windows = self._smart_shuffle.windows()
651 snapshot = await self.mass.music.recency.snapshot(windows, userid=queue_data.userid)
652 tracks = gate_tracks(
653 [track for track in tracks if isinstance(track, Track)], snapshot, windows
654 )
655 queue_items = [build_queue_item(queue_id, x) for x in tracks if x.available]
656 if not queue_items:
657 self.logger.info("Autoplay found no new tracks to add for queue %s", queue.display_name)
658 return
659 await self.load(
660 queue_id,
661 queue_items,
662 insert_at_index=len(self._queue_data[queue_id].items) + 1,
663 )
664
665 @handle_play_action
666 async def _handle_play_media(
667 self,
668 queue_id: str,
669 media: MediaItemType | ItemMapping | str | list[MediaItemType | ItemMapping | str],
670 option: QueueOption | None = None,
671 radio_mode: bool = False,
672 start_item: PlayableMediaItemType | str | None = None,
673 sort_by: str | None = None,
674 start_from_beginning: bool = False,
675 shuffle: bool | None = None,
676 ) -> None:
677 """Handle play media without acquiring the queue lock."""
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 self._set_transitioning(queue_id, False)
681 # we use a contextvar to bypass the throttler for this asyncio task/context
682 # this makes sure that playback has priority over other requests that may be
683 # happening in the background
684 BYPASS_THROTTLER.set(True)
685 if not (queue := self.get(queue_id)):
686 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
687 queue_data = self._queue_data[queue_id]
688 # always fetch the underlying player so we can raise early if its not available
689 queue_player = self.mass.players.get_player(queue_id, True)
690 assert queue_player is not None # for type checking
691 if queue_player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS):
692 self.logger.warning("Ignore queue command: An announcement is in progress")
693 return
694
695 # save the user requesting the playback (clear it for anonymous playback)
696 playback_user = get_current_user()
697 queue_data.userid = playback_user.user_id if playback_user else None
698 if playback_user:
699 self.logger.debug(
700 "User %s requested playback.", playback_user.display_name or playback_user.username
701 )
702
703 # a single item or list of items may be provided
704 media_list = media if isinstance(media, list) else [media]
705
706 if radio_mode:
707 # radio_mode is deprecated: a "radio" is now a dynamic radio playlist. Translate each
708 # seed into the radio_playlist provider's URI and enqueue those (resolved to dynamic
709 # playlists that self-manage their refills).
710 self.logger.warning(
711 "radio_mode is deprecated; enqueue a radio_playlist:// dynamic playlist instead"
712 )
713 media_list = [
714 seed_uri
715 if (seed_uri := item if isinstance(item, str) else str(item.uri)).startswith(
716 "radio_playlist://"
717 )
718 else f"radio_playlist://playlist/{seed_uri}"
719 for item in media_list
720 ]
721 radio_mode = False
722
723 # Clear the 'enqueued media item' list when a new queue is requested
724 if option not in (QueueOption.ADD, QueueOption.NEXT):
725 queue_data.enqueued_media_items.clear()
726 # An ADD/NEXT onto a queue that is already a managed pool (has a dynamic source): a finite
727 # item is kept only as a source (the bounded pool materializes it) instead of being expanded
728 # into the queue. Any other enqueue (PLAY/REPLACE, or onto a linear queue) expands finite
729 # items normally. Keys off is_dynamic since a finite-only queue records sources too.
730 already_dynamic = queue.is_dynamic and option in (QueueOption.ADD, QueueOption.NEXT)
731
732 media_items: list[MediaItemType] = []
733 source_items: list[MediaItemType] = []
734 shuffle_settled = False
735 # resolve all media items
736 for item in media_list:
737 try:
738 # parse provided uri into a MA MediaItem or Basic QueueItem from URL
739 media_item: MediaItemType | ItemMapping | BrowseFolder
740 if isinstance(item, str):
741 media_item = await self.mass.music.get_item_by_uri(item)
742 elif isinstance(item, dict): # type: ignore[unreachable]
743 # TODO: Investigate why the API parser sometimes passes raw dicts instead of
744 # converting them to MediaItem objects. The parse_value function in api.py
745 # should handle dict-to-object conversion, but dicts are slipping through
746 # in some cases. This is defensive handling for that parser bug.
747 media_item = media_from_dict(item) # type: ignore[unreachable]
748 self.logger.debug("Converted to: %s", type(media_item))
749 else:
750 # item is MediaItemType | ItemMapping at this point
751 media_item = item
752
753 if isinstance(media_item, ItemMapping):
754 # Resolve any ItemMapping to its full media item, exactly as the str-uri
755 # form above already does. Everything below needs the real object: the
756 # enqueued/source bookkeeping only accepts full items (so a mapping would
757 # otherwise never count as a user-initiated play), and the dynamic check
758 # needs details such as a playlist's 'is_dynamic'.
759 if media_item.uri is None:
760 raise InvalidDataError("ItemMapping has no URI")
761 media_item = await self.mass.music.get_item_by_uri(media_item.uri)
762
763 # Save requested media item to play on the queue so we can use it as a seed
764 # for Autoplay's music refill (the podcast/audiobook continuations resolve
765 # their successor from the queue's last item instead).
766 # Use FIFO list to keep track of the last 10 played items
767 # Skip ItemMapping and BrowseFolder - only queue full MediaItemType objects
768 if not isinstance(media_item, BrowseFolder) and (
769 is_dynamic_source(media_item)
770 or media_item.media_type
771 in (MediaType.TRACK, MediaType.ALBUM, MediaType.PLAYLIST, MediaType.ARTIST)
772 ):
773 queue_data.enqueued_media_items.append(media_item)
774 if len(queue_data.enqueued_media_items) > 10:
775 queue_data.enqueued_media_items.pop(0)
776 if is_dynamic_source(media_item):
777 # a dynamic playlist/station is always a self-managing dynamic source
778 source_items.append(media_item)
779
780 # handle default enqueue option if needed
781 if option is None:
782 # Radio + AudioSource share a single "live_sources" enqueue default —
783 # both are live infinite streams where REPLACE is almost always the
784 # right semantic. Other media types use their per-type config key.
785 if media_item.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
786 config_key = CONF_DEFAULT_ENQUEUE_OPTION_LIVE_SOURCES
787 else:
788 config_key = f"default_enqueue_option_{media_item.media_type.value}"
789 config_value = self.get_config_value(config_key, return_type=str)
790 option = QueueOption(config_value)
791
792 # The shuffle state has to be settled before the items are resolved below: a
793 # shuffled queue keeps the items preceding a start_item (chosen track pinned
794 # first) instead of dropping them. The first item that resolves decides for the
795 # whole batch, because it is the only media type known this early.
796 if not shuffle_settled:
797 shuffle_settled = True
798 await self._apply_shuffle(
799 queue_id,
800 option,
801 # an explicit request always wins; only an unset one defers to the
802 # media's own order
803 False
804 if shuffle is None and media_item.media_type in ORDERED_MEDIA_TYPES
805 else shuffle,
806 )
807
808 # collect media_items to play
809 if is_dynamic_source(media_item):
810 # a dynamic playlist/station supplies its own tracks on demand; just mark it
811 # played. The queue goes dynamic below and the bounded pool seeds its batch from
812 # all sources, so there is no need to fetch a batch here.
813 self.mass.create_task(
814 self.mass.music.mark_item_played(
815 media_item,
816 userid=queue_data.userid,
817 queue_id=queue_id,
818 user_initiated=True,
819 )
820 )
821 elif already_dynamic:
822 # feed the already-active pool: keep the finite item as a (materialized) source
823 if not isinstance(media_item, BrowseFolder):
824 source_items.append(media_item)
825 else:
826 # not (yet) a managed pool: record the finite parent as a source (kept for a
827 # later dynamic transition and for similar/autoplay seeds) and expand it into
828 # the linear queue
829 if not isinstance(media_item, BrowseFolder) and media_item.media_type in (
830 MediaType.TRACK,
831 MediaType.ALBUM,
832 MediaType.PLAYLIST,
833 MediaType.ARTIST,
834 ):
835 source_items.append(media_item)
836 # Convert start_item to string URI if needed
837 start_item_uri: str | None = None
838 if isinstance(start_item, str):
839 start_item_uri = start_item
840 elif start_item is not None:
841 start_item_uri = start_item.uri
842 media_items += await self._media_resolver._resolve_media_items(
843 media_item,
844 start_item_uri,
845 userid=queue_data.userid,
846 queue_id=queue_id,
847 sort_by=sort_by,
848 start_from_beginning=start_from_beginning,
849 # under shuffle "start here and play forward" has no meaning, so keep the
850 # whole playlist/album (chosen track first) instead of dropping everything
851 # before it - the chosen track is pinned in front of the shuffled rest
852 keep_preceding_items=queue.shuffle_enabled,
853 )
854
855 except MusicAssistantError as err:
856 # invalid MA uri or item not found error
857 self.logger.warning("Skipping %s: %s", item, str(err))
858
859 if not shuffle_settled and option is not None:
860 # nothing resolved, so no media type ever decided - but the sources are replaced
861 # below all the same, and a dynamic queue's imposed shuffle must not survive that
862 await self._apply_shuffle(queue_id, option, shuffle)
863
864 # overwrite or append the queue's source items
865 replace_sources = option not in (QueueOption.ADD, QueueOption.NEXT)
866 if replace_sources:
867 self.store_sources(queue, source_items)
868 else:
869 self.store_sources(queue, self._queue_data[queue_id].source_items + source_items)
870 source_items = self._queue_data[queue_id].source_items
871 queue.is_dynamic = has_dynamic_source(source_items)
872 # a queue that just gained or lost its dynamic source resolves smart shuffle differently
873 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
874
875 if queue.is_dynamic:
876 # the queue has (or just gained) a dynamic source: (re)build the upcoming tail into a
877 # single bounded, recency-orchestrated mix over ALL sources — existing finite content as
878 # materialized TRACKS seed(s), dynamic playlists as DYNAMIC seed(s). Every add rebuilds
879 # from the buffer position, so the queue stays a fixed-size mix instead of growing by
880 # each added source's own batch.
881 await self._enter_dynamic_mode(queue_id, option)
882 return
883
884 # only add valid/available items
885 queue_items: list[QueueItem] = [
886 build_queue_item(queue_id, cast("PlayableMediaItemType", x))
887 for x in media_items
888 if x and x.available
889 ]
890
891 if not queue_items:
892 raise MediaNotFoundError("No playable items found", translation_key="no_playable_items")
893
894 await self._enqueue_with_option(
895 queue_id, queue_items, option, pin_first=start_item is not None
896 )
897
898 async def _enter_dynamic_mode(self, queue_id: str, option: QueueOption | None) -> None:
899 """
900 (Re)build a queue's upcoming tail into a single bounded managed pool over all its sources.
901
902 Runs whenever an enqueue leaves the queue dynamic — both the first transition and every
903 later add. Keeps the current + already-buffered track(s), drops the rest of the upcoming
904 tail, and replaces it with a bounded, recency-orchestrated mix of all the queue's sources
905 (finite sources materialized as TRACKS seeds, dynamic playlists as DYNAMIC seeds), so the
906 queue stays a fixed-size mix instead of growing by each added source's own batch. Shuffle is
907 enabled implicitly: a dynamic queue is always a smart mix.
908
909 :param queue_id: The queue to (re)build the dynamic pool for.
910 :param option: The enqueue option that triggered the (re)build. PLAY/REPLACE start playback
911 on the rebuilt pool; ADD/NEXT/REPLACE_NEXT stage it without starting playback (behind the
912 current/buffered track, or from the front of an idle/empty queue).
913 """
914 queue_data = self._queue_data[queue_id]
915 queue = queue_data.queue
916 # a dynamic queue is an always-on smart mix; reflect that in the (now locked) shuffle state
917 queue.shuffle_enabled = True
918 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
919 # rebuild from the buffered position so the already-prepared next track is kept and the
920 # crossfade isn't disturbed; fall back to the current index (or the front when idle/empty)
921 base_index = (
922 queue.index_in_buffer if queue.index_in_buffer is not None else queue.current_index
923 )
924 insert_at = 0 if base_index is None else base_index + 1
925 if option == QueueOption.REPLACE:
926 # A replace is a fresh queue, so the pool takes the place of the old items rather than
927 # being appended behind the one that is playing (as PLAY, which shares start_playing,
928 # deliberately does). Zeroed before the truncation below so the pool is sized against
929 # an empty queue and none of the discarded tracks are held back from it.
930 insert_at = 0
931 # as on the linear path: release the outgoing audio while its items are still on the
932 # queue, and drop the stale position
933 await self._cleanup_queue_audio_data(queue_id)
934 queue.index_in_buffer = None
935 queue.ended = False
936 # PLAY/REPLACE start playback on the rebuilt pool; ADD/NEXT/REPLACE_NEXT only stage it and
937 # never start playback (an idle/empty queue stays idle on an add, just like the linear path)
938 start_playing = option in (QueueOption.PLAY, QueueOption.REPLACE)
939 # The tail is dropped before the pool is fetched, so the pool is sized and deduped against
940 # the kept head only (the tail we are discarding must not exclude its own tracks from it).
941 # That leaves the queue holding less than it plays - for a replace, nothing at all - across
942 # the fetch, so hold player reconciliation off until the new items are in: it would
943 # otherwise publish that half-built state, which is exactly the empty queue this avoids.
944 self._set_transitioning(queue_id, True)
945 try:
946 queue_data.items = queue_data.items[:insert_at]
947 queue.items = len(queue_data.items)
948 pool_tracks = await self._managed_pool.fill(queue_id, is_initial=False)
949 queue_items = [
950 build_queue_item(queue_id, track) for track in pool_tracks if track.available
951 ]
952 if not queue_items:
953 raise MediaNotFoundError(
954 "No playable items found", translation_key="no_playable_items"
955 )
956 # the managed pool already interleaved the sources in a recency-aware order; load as-is
957 await self.load(
958 queue_id,
959 queue_items,
960 insert_at_index=insert_at,
961 keep_remaining=False,
962 keep_played=option != QueueOption.REPLACE,
963 )
964 if start_playing:
965 await self.play_index(queue_id, insert_at)
966 else:
967 # give an idle/empty queue a current item without starting playback
968 self._ensure_current_index(queue_id)
969 finally:
970 self._set_transitioning(queue_id, False)
971
972 async def _get_similar_tracks(
973 self,
974 queue_id: str,
975 is_initial: bool = False,
976 seed_items: list[MediaItemType] | None = None,
977 ) -> list[Track]:
978 """
979 Fetch tracks similar to the given seeds (autoplay's similar/continuation mode).
980
981 :param queue_id: The queue to fetch tracks for.
982 :param is_initial: True to interleave the base/seed tracks into the result, False to
983 return only similar tracks.
984 :param seed_items: Explicit seed items to base the tracks on. Defaults to the queue's
985 sources; autoplay passes the enqueued media items instead.
986 """
987 queue_data = self._queue_data[queue_id]
988 queue = queue_data.queue
989 queue_track_items: list[Track] = [
990 q.media_item
991 for q in self._queue_data[queue_id].items
992 if q.media_item and isinstance(q.media_item, Track)
993 ]
994 source_items = (
995 seed_items if seed_items is not None else self._queue_data[queue_id].source_items
996 )
997 if not source_items:
998 # this may happen during race conditions as this method is called delayed
999 return []
1000 self.logger.info(
1001 "Fetching similar tracks for queue %s based on: %s",
1002 queue.display_name,
1003 ", ".join([x.name for x in source_items]),
1004 )
1005
1006 # Get user's preferred provider instances for steering provider selection
1007 preferred_provider_instances: list[str] | None = None
1008 if (
1009 queue_data.userid
1010 and (playback_user := await self.mass.webserver.auth.get_user(queue_data.userid))
1011 and playback_user.provider_filter
1012 ):
1013 preferred_provider_instances = playback_user.provider_filter
1014
1015 # Some providers have very deterministic similar-track algorithms for a single track
1016 # seed. When continuing from a single track on a refill, seed from the play history
1017 # instead so the result keeps varying.
1018 if (
1019 len(source_items) == 1
1020 and source_items[0].media_type == MediaType.TRACK
1021 and not is_initial
1022 and queue_track_items
1023 ):
1024 # Helper samples 5 internally; bound the input.
1025 seeds: list[MediaItemType] = random.sample(
1026 queue_track_items, min(len(queue_track_items), 10)
1027 )
1028 else:
1029 seeds = list(source_items)
1030
1031 radio_prov = self.mass.get_provider("radio_playlist")
1032 if radio_prov is None:
1033 return []
1034 dynamic_tracks = await cast("RadioPlaylistProvider", radio_prov).get_dynamic_tracks(
1035 seeds,
1036 include_base_tracks=is_initial,
1037 target_size=25,
1038 preferred_provider_instances=preferred_provider_instances,
1039 )
1040 # Drop anything already queued/played
1041 queued_set = set(queue_track_items)
1042 return [track for track in dynamic_tracks if track not in queued_set]
1043
1044 async def _abort_superseded_source_buffers(self, queue_item: QueueItem) -> None:
1045 """
1046 Abort the still-filling source buffers of other items in the same queue.
1047
1048 :param queue_item: The queue item that is about to start playing.
1049 """
1050 queue_data = self._queue_data.get(queue_item.queue_id)
1051 items = tuple(queue_data.items) if queue_data else ()
1052 successor: QueueItem | None = None
1053 for index, item in enumerate(items):
1054 if item.queue_item_id == queue_item.queue_item_id and index + 1 < len(items):
1055 successor = items[index + 1]
1056 break
1057 # the started item keeps its own buffer, and its direct successor keeps the prewarm
1058 # for the upcoming crossfade unless the aborts below leave the provider without a slot
1059 spared_item_ids = {queue_item.queue_item_id}
1060 if successor is not None:
1061 spared_item_ids.add(successor.queue_item_id)
1062 for item in items:
1063 if item.queue_item_id in spared_item_ids:
1064 continue
1065 await self._abort_source_buffer(item, queue_item)
1066 if successor is not None:
1067 await self._abort_source_buffer(successor, queue_item, only_when_saturated=True)
1068
1069 async def _abort_source_buffer(
1070 self,
1071 item: QueueItem,
1072 started_item: QueueItem,
1073 only_when_saturated: bool = False,
1074 ) -> None:
1075 """
1076 Cancel one item's still-filling source so its provider stream slot is handed over.
1077
1078 :param item: The queue item whose source buffer should be aborted.
1079 :param started_item: The queue item that is about to start playing.
1080 :param only_when_saturated: Only abort while the provider has no free slot left.
1081 """
1082 if item.streamdetails is None:
1083 return
1084 audio_buffer = item.streamdetails.buffer
1085 if audio_buffer is None or not audio_buffer.is_buffering:
1086 return
1087 provider = self.mass.get_provider(item.streamdetails.provider, return_unavailable=True)
1088 if not isinstance(provider, MusicProvider) or provider.max_concurrent_streams is None:
1089 return
1090 if only_when_saturated:
1091 if provider.has_available_stream_slot:
1092 # an abort above already freed a slot, so this prewarm can stay
1093 return
1094 self.logger.debug(
1095 "Aborting the prewarm of %s: %s has no free stream slot left for %s",
1096 item.name,
1097 provider.name,
1098 started_item.name,
1099 )
1100 else:
1101 self.logger.debug(
1102 "Aborting the source of %s to free a %s stream slot for %s",
1103 item.name,
1104 provider.name,
1105 started_item.name,
1106 )
1107 # the cancelled buffer stays attached: it marks the source as aborted for
1108 # the flow stream's accounting and fails is_valid() for any later reuse
1109 await audio_buffer.clear()
1110