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