/
/
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 # An ADD/NEXT onto a queue that is already a managed pool (has a dynamic source): a finite
735 # item is kept only as a source (the bounded pool materializes it) instead of being expanded
736 # into the queue. Any other enqueue (PLAY/REPLACE, or onto a linear queue) expands finite
737 # items normally. Keys off is_dynamic since a finite-only queue records sources too.
738 already_dynamic = queue.is_dynamic and option in (QueueOption.ADD, QueueOption.NEXT)
739
740 media_items: list[MediaItemType] = []
741 source_items: list[MediaItemType] = []
742 shuffle_settled = False
743 # resolve all media items
744 for item in media_list:
745 try:
746 # parse provided uri into a MA MediaItem or Basic QueueItem from URL
747 media_item: MediaItemType | ItemMapping | BrowseFolder
748 if isinstance(item, str):
749 media_item = await self.mass.music.get_item_by_uri(item)
750 elif isinstance(item, dict): # type: ignore[unreachable]
751 # TODO: Investigate why the API parser sometimes passes raw dicts instead of
752 # converting them to MediaItem objects. The parse_value function in api.py
753 # should handle dict-to-object conversion, but dicts are slipping through
754 # in some cases. This is defensive handling for that parser bug.
755 media_item = media_from_dict(item) # type: ignore[unreachable]
756 self.logger.debug("Converted to: %s", type(media_item))
757 else:
758 # item is MediaItemType | ItemMapping at this point
759 media_item = item
760
761 if isinstance(media_item, ItemMapping):
762 # Resolve any ItemMapping to its full media item, exactly as the str-uri
763 # form above already does. Everything below needs the real object: the
764 # enqueued/source bookkeeping only accepts full items (so a mapping would
765 # otherwise never count as a user-initiated play), and the dynamic check
766 # needs details such as a playlist's 'is_dynamic'.
767 if media_item.uri is None:
768 raise InvalidDataError("ItemMapping has no URI")
769 media_item = await self.mass.music.get_item_by_uri(media_item.uri)
770
771 # Save requested media item to play on the queue so we can use it as a seed
772 # for Autoplay's music refill (the podcast/audiobook continuations resolve
773 # their successor from the queue's last item instead).
774 # Use FIFO list to keep track of the last 10 played items
775 # Skip ItemMapping and BrowseFolder - only queue full MediaItemType objects
776 if not isinstance(media_item, BrowseFolder) and (
777 is_dynamic_source(media_item)
778 or media_item.media_type
779 in (MediaType.TRACK, MediaType.ALBUM, MediaType.PLAYLIST, MediaType.ARTIST)
780 ):
781 queue_data.enqueued_media_items.append(media_item)
782 if len(queue_data.enqueued_media_items) > 10:
783 queue_data.enqueued_media_items.pop(0)
784 if is_dynamic_source(media_item):
785 # a dynamic playlist/station is always a self-managing dynamic source
786 source_items.append(media_item)
787
788 # handle default enqueue option if needed
789 if option is None:
790 # Radio + AudioSource share a single "live_sources" enqueue default —
791 # both are live infinite streams where REPLACE is almost always the
792 # right semantic. Other media types use their per-type config key.
793 if media_item.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
794 config_key = CONF_DEFAULT_ENQUEUE_OPTION_LIVE_SOURCES
795 else:
796 config_key = f"default_enqueue_option_{media_item.media_type.value}"
797 config_value = self.get_config_value(config_key, return_type=str)
798 option = QueueOption(config_value)
799
800 # The shuffle state has to be settled before the items are resolved below: a
801 # shuffled queue keeps the items preceding a start_item (chosen track pinned
802 # first) instead of dropping them. The first item that resolves decides for the
803 # whole batch, because it is the only media type known this early.
804 if not shuffle_settled:
805 shuffle_settled = True
806 await self._apply_shuffle(
807 queue_id,
808 option,
809 # an explicit request always wins; only an unset one defers to the
810 # media's own order
811 False
812 if shuffle is None and media_item.media_type in ORDERED_MEDIA_TYPES
813 else shuffle,
814 )
815
816 # collect media_items to play
817 if is_dynamic_source(media_item):
818 # a dynamic playlist/station supplies its own tracks on demand; just mark it
819 # played. The queue goes dynamic below and the bounded pool seeds its batch from
820 # all sources, so there is no need to fetch a batch here.
821 self.mass.create_task(
822 self.mass.music.mark_item_played(
823 media_item,
824 userid=queue_data.userid,
825 queue_id=queue_id,
826 user_initiated=True,
827 )
828 )
829 elif already_dynamic:
830 # feed the already-active pool: keep the finite item as a (materialized) source
831 if not isinstance(media_item, BrowseFolder):
832 source_items.append(media_item)
833 else:
834 # not (yet) a managed pool: record the finite parent as a source (kept for a
835 # later dynamic transition and for similar/autoplay seeds) and expand it into
836 # the linear queue
837 if not isinstance(media_item, BrowseFolder) and media_item.media_type in (
838 MediaType.TRACK,
839 MediaType.ALBUM,
840 MediaType.PLAYLIST,
841 MediaType.ARTIST,
842 ):
843 source_items.append(media_item)
844 # Convert start_item to string URI if needed
845 start_item_uri: str | None = None
846 if isinstance(start_item, str):
847 start_item_uri = start_item
848 elif start_item is not None:
849 start_item_uri = start_item.uri
850 media_items += await self._media_resolver._resolve_media_items(
851 media_item,
852 start_item_uri,
853 userid=queue_data.userid,
854 queue_id=queue_id,
855 sort_by=sort_by,
856 start_from_beginning=start_from_beginning,
857 # under shuffle "start here and play forward" has no meaning, so keep the
858 # whole playlist/album (chosen track first) instead of dropping everything
859 # before it - the chosen track is pinned in front of the shuffled rest
860 keep_preceding_items=queue.shuffle_enabled,
861 )
862
863 except MusicAssistantError as err:
864 # invalid MA uri or item not found error
865 self.logger.warning("Skipping %s: %s", item, str(err))
866
867 if not shuffle_settled and option is not None:
868 # nothing resolved, so no media type ever decided - but the sources are replaced
869 # below all the same, and a dynamic queue's imposed shuffle must not survive that
870 await self._apply_shuffle(queue_id, option, shuffle)
871
872 # overwrite or append the queue's source items
873 replace_sources = option not in (QueueOption.ADD, QueueOption.NEXT)
874 if replace_sources:
875 self.store_sources(queue, source_items)
876 else:
877 self.store_sources(queue, self._queue_data[queue_id].source_items + source_items)
878 source_items = self._queue_data[queue_id].source_items
879 queue.is_dynamic = has_dynamic_source(source_items)
880 # a queue that just gained or lost its dynamic source resolves smart shuffle differently
881 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
882
883 if queue.is_dynamic:
884 # the queue has (or just gained) a dynamic source: (re)build the upcoming tail into a
885 # single bounded, recency-orchestrated mix over ALL sources — existing finite content as
886 # materialized TRACKS seed(s), dynamic playlists as DYNAMIC seed(s). Every add rebuilds
887 # from the buffer position, so the queue stays a fixed-size mix instead of growing by
888 # each added source's own batch.
889 await self._enter_dynamic_mode(queue_id, option)
890 return
891
892 # only add valid/available items
893 queue_items: list[QueueItem] = [
894 build_queue_item(queue_id, cast("PlayableMediaItemType", x))
895 for x in media_items
896 if x and x.available
897 ]
898
899 if not queue_items:
900 raise MediaNotFoundError("No playable items found", translation_key="no_playable_items")
901
902 await self._enqueue_with_option(
903 queue_id, queue_items, option, pin_first=start_item is not None
904 )
905
906 async def _enter_dynamic_mode(self, queue_id: str, option: QueueOption | None) -> None:
907 """
908 (Re)build a queue's upcoming tail into a single bounded managed pool over all its sources.
909
910 Runs whenever an enqueue leaves the queue dynamic — both the first transition and every
911 later add. Keeps the current + already-buffered track(s), drops the rest of the upcoming
912 tail, and replaces it with a bounded, recency-orchestrated mix of all the queue's sources
913 (finite sources materialized as TRACKS seeds, dynamic playlists as DYNAMIC seeds), so the
914 queue stays a fixed-size mix instead of growing by each added source's own batch. Shuffle is
915 enabled implicitly: a dynamic queue is always a smart mix.
916
917 :param queue_id: The queue to (re)build the dynamic pool for.
918 :param option: The enqueue option that triggered the (re)build. PLAY/REPLACE start playback
919 on the rebuilt pool; ADD/NEXT/REPLACE_NEXT stage it without starting playback (behind the
920 current/buffered track, or from the front of an idle/empty queue).
921 """
922 queue_data = self._queue_data[queue_id]
923 queue = queue_data.queue
924 # a dynamic queue is an always-on smart mix; reflect that in the (now locked) shuffle state
925 queue.shuffle_enabled = True
926 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
927 # rebuild from the buffered position so the already-prepared next track is kept and the
928 # crossfade isn't disturbed; fall back to the current index (or the front when idle/empty)
929 base_index = (
930 queue.index_in_buffer if queue.index_in_buffer is not None else queue.current_index
931 )
932 insert_at = 0 if base_index is None else base_index + 1
933 if option == QueueOption.REPLACE:
934 # A replace is a fresh queue, so the pool takes the place of the old items rather than
935 # being appended behind the one that is playing (as PLAY, which shares start_playing,
936 # deliberately does). Zeroed before the truncation below so the pool is sized against
937 # an empty queue and none of the discarded tracks are held back from it.
938 insert_at = 0
939 # as on the linear path: release the outgoing audio while its items are still on the
940 # queue, and drop the stale position
941 await self._cleanup_queue_audio_data(queue_id)
942 queue.index_in_buffer = None
943 queue.ended = False
944 # PLAY/REPLACE start playback on the rebuilt pool; ADD/NEXT/REPLACE_NEXT only stage it and
945 # never start playback (an idle/empty queue stays idle on an add, just like the linear path)
946 start_playing = option in (QueueOption.PLAY, QueueOption.REPLACE)
947 # The tail is dropped before the pool is fetched, so the pool is sized and deduped against
948 # the kept head only (the tail we are discarding must not exclude its own tracks from it).
949 # That leaves the queue holding less than it plays - for a replace, nothing at all - across
950 # the fetch, so hold player reconciliation off until the new items are in: it would
951 # otherwise publish that half-built state, which is exactly the empty queue this avoids.
952 self._set_transitioning(queue_id, True)
953 try:
954 queue_data.items = queue_data.items[:insert_at]
955 queue.items = len(queue_data.items)
956 pool_tracks = await self._managed_pool.fill(queue_id, is_initial=False)
957 queue_items = [
958 build_queue_item(queue_id, track) for track in pool_tracks if track.available
959 ]
960 if not queue_items:
961 raise MediaNotFoundError(
962 "No playable items found", translation_key="no_playable_items"
963 )
964 # the managed pool already interleaved the sources in a recency-aware order; load as-is
965 await self.load(
966 queue_id,
967 queue_items,
968 insert_at_index=insert_at,
969 keep_remaining=False,
970 keep_played=option != QueueOption.REPLACE,
971 )
972 if start_playing:
973 await self.play_index(queue_id, insert_at)
974 else:
975 # give an idle/empty queue a current item without starting playback
976 self._ensure_current_index(queue_id)
977 finally:
978 self._set_transitioning(queue_id, False)
979
980 async def _get_similar_tracks(
981 self,
982 queue_id: str,
983 is_initial: bool = False,
984 seed_items: list[MediaItemType] | None = None,
985 ) -> list[Track]:
986 """
987 Fetch tracks similar to the given seeds (autoplay's similar/continuation mode).
988
989 :param queue_id: The queue to fetch tracks for.
990 :param is_initial: True to interleave the base/seed tracks into the result, False to
991 return only similar tracks.
992 :param seed_items: Explicit seed items to base the tracks on. Defaults to the queue's
993 sources; autoplay passes the enqueued media items instead.
994 """
995 queue_data = self._queue_data[queue_id]
996 queue = queue_data.queue
997 queue_track_items: list[Track] = [
998 q.media_item
999 for q in self._queue_data[queue_id].items
1000 if q.media_item and isinstance(q.media_item, Track)
1001 ]
1002 source_items = (
1003 seed_items if seed_items is not None else self._queue_data[queue_id].source_items
1004 )
1005 if not source_items:
1006 # this may happen during race conditions as this method is called delayed
1007 return []
1008 self.logger.info(
1009 "Fetching similar tracks for queue %s based on: %s",
1010 queue.display_name,
1011 ", ".join([x.name for x in source_items]),
1012 )
1013
1014 # Get user's preferred provider instances for steering provider selection
1015 preferred_provider_instances: list[str] | None = None
1016 if (
1017 queue_data.userid
1018 and (playback_user := await self.mass.webserver.auth.get_user(queue_data.userid))
1019 and playback_user.provider_filter
1020 ):
1021 preferred_provider_instances = playback_user.provider_filter
1022
1023 # Some providers have very deterministic similar-track algorithms for a single track
1024 # seed. When continuing from a single track on a refill, seed from the play history
1025 # instead so the result keeps varying.
1026 if (
1027 len(source_items) == 1
1028 and source_items[0].media_type == MediaType.TRACK
1029 and not is_initial
1030 and queue_track_items
1031 ):
1032 # Helper samples 5 internally; bound the input.
1033 seeds: list[MediaItemType] = random.sample(
1034 queue_track_items, min(len(queue_track_items), 10)
1035 )
1036 else:
1037 seeds = list(source_items)
1038
1039 radio_prov = self.mass.get_provider("radio_playlist")
1040 if radio_prov is None:
1041 return []
1042 dynamic_tracks = await cast("RadioPlaylistProvider", radio_prov).get_dynamic_tracks(
1043 seeds,
1044 include_base_tracks=is_initial,
1045 target_size=25,
1046 preferred_provider_instances=preferred_provider_instances,
1047 )
1048 # Drop anything already queued/played
1049 queued_set = set(queue_track_items)
1050 return [track for track in dynamic_tracks if track not in queued_set]
1051
1052 async def _abort_superseded_source_buffers(self, queue_item: QueueItem) -> None:
1053 """
1054 Abort the still-filling source buffers of other items in the same queue.
1055
1056 :param queue_item: The queue item that is about to start playing.
1057 """
1058 queue_data = self._queue_data.get(queue_item.queue_id)
1059 items = tuple(queue_data.items) if queue_data else ()
1060 successor: QueueItem | None = None
1061 for index, item in enumerate(items):
1062 if item.queue_item_id == queue_item.queue_item_id and index + 1 < len(items):
1063 successor = items[index + 1]
1064 break
1065 # the started item keeps its own buffer, and its direct successor keeps the prewarm
1066 # for the upcoming crossfade unless the aborts below leave the provider without a slot
1067 spared_item_ids = {queue_item.queue_item_id}
1068 if successor is not None:
1069 spared_item_ids.add(successor.queue_item_id)
1070 for item in items:
1071 if item.queue_item_id in spared_item_ids:
1072 continue
1073 await self._abort_source_buffer(item, queue_item)
1074 if successor is not None:
1075 await self._abort_source_buffer(successor, queue_item, only_when_saturated=True)
1076
1077 async def _abort_source_buffer(
1078 self,
1079 item: QueueItem,
1080 started_item: QueueItem,
1081 only_when_saturated: bool = False,
1082 ) -> None:
1083 """
1084 Cancel one item's still-filling source so its provider stream slot is handed over.
1085
1086 :param item: The queue item whose source buffer should be aborted.
1087 :param started_item: The queue item that is about to start playing.
1088 :param only_when_saturated: Only abort while the provider has no free slot left.
1089 """
1090 if item.streamdetails is None:
1091 return
1092 audio_buffer = item.streamdetails.buffer
1093 if audio_buffer is None or not audio_buffer.is_buffering:
1094 return
1095 provider = self.mass.get_provider(item.streamdetails.provider, return_unavailable=True)
1096 if not isinstance(provider, MusicProvider) or provider.max_concurrent_streams is None:
1097 return
1098 if only_when_saturated:
1099 if provider.has_available_stream_slot:
1100 # an abort above already freed a slot, so this prewarm can stay
1101 return
1102 self.logger.debug(
1103 "Aborting the prewarm of %s: %s has no free stream slot left for %s",
1104 item.name,
1105 provider.name,
1106 started_item.name,
1107 )
1108 else:
1109 self.logger.debug(
1110 "Aborting the source of %s to free a %s stream slot for %s",
1111 item.name,
1112 provider.name,
1113 started_item.name,
1114 )
1115 # the cancelled buffer stays attached: it marks the source as aborted for
1116 # the flow stream's accounting and fails is_valid() for any later reuse
1117 await audio_buffer.clear()
1118