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