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