/
/
1"""
2MusicAssistant Player Queues Controller.
3
4Handles all logic to PLAY Media Items, provided by Music Providers to supported players.
5
6It is loosely coupled to the MusicAssistant Music Controller and Player Controller.
7A Music Assistant Player always has a PlayerQueue associated with it
8which holds the queue items and state.
9
10The PlayerQueue is in that case the active source of the player,
11but it can also be something else, hence the loose coupling.
12"""
13
14from __future__ import annotations
15
16import asyncio
17import random
18import time
19from typing import TYPE_CHECKING, Any, Final, cast
20
21import shortuuid
22from music_assistant_models.auth import Scope
23from music_assistant_models.enums import (
24 EventType,
25 MediaType,
26 PlaybackState,
27 PlayerType,
28 QueueOption,
29 RepeatMode,
30 SourceControl,
31)
32from music_assistant_models.errors import (
33 AudioError,
34 InsufficientPermissions,
35 InvalidCommand,
36 InvalidDataError,
37 MediaNotFoundError,
38 PlayerUnavailableError,
39 QueueEmpty,
40)
41from music_assistant_models.media_items import (
42 Audiobook,
43 AudioSource,
44 ItemMapping,
45 MediaItemType,
46 PlayableMediaItemType,
47 Playlist,
48 PodcastEpisode,
49 SoundEffect,
50 Track,
51)
52from music_assistant_models.player_queue import PlayerQueue
53
54from music_assistant.constants import (
55 ATTR_ANNOUNCEMENT_IN_PROGRESS,
56 MASS_LOGO_ONLINE,
57 PLAYLIST_MEDIA_TYPES,
58)
59from music_assistant.controllers.player_queues.autoplay import Autoplay
60from music_assistant.controllers.player_queues.config import (
61 core_config_entries,
62 queue_config_entries,
63)
64from music_assistant.controllers.player_queues.constants import (
65 CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
66 CACHE_CATEGORY_PLAYER_QUEUE_STATE,
67 PLAYBACK_START_TIMEOUT,
68 QUEUE_CACHE_SAVE_DELAY,
69)
70from music_assistant.controllers.player_queues.helpers import (
71 get_current_playback_speed,
72 handle_play_action,
73 is_dynamic_source,
74)
75from music_assistant.controllers.player_queues.managed_pool import ManagedPool
76from music_assistant.controllers.player_queues.media_resolver import MediaResolver
77from music_assistant.controllers.player_queues.playback_tracker import PlaybackTrackerMixin
78from music_assistant.controllers.player_queues.queue_loader import QueueLoaderMixin
79from music_assistant.controllers.player_queues.smart_shuffle import SmartShuffle
80from music_assistant.controllers.player_queues.state import PlayerQueueData
81from music_assistant.controllers.player_queues.stream_feeder import StreamFeederMixin
82from music_assistant.controllers.webserver.helpers.auth_middleware import get_current_user
83from music_assistant.helpers.api import api_command
84from music_assistant.helpers.player import get_queue_audio_source
85from music_assistant.models.music_provider import ProviderStreamLimitError
86from music_assistant.models.player import Player, PlayerMedia
87from music_assistant.models.plugin import PluginProvider
88
89if TYPE_CHECKING:
90 from collections.abc import Iterator
91
92 from music_assistant_models import BackgroundTask
93 from music_assistant_models.config_entries import (
94 ConfigEntry,
95 ConfigValueOption,
96 CoreConfig,
97 )
98 from music_assistant_models.media_items import SourceQueueCapabilities
99 from music_assistant_models.queue_item import QueueItem
100
101 from music_assistant import MusicAssistant
102 from music_assistant.constants import PlaylistPlayableItem
103 from music_assistant.controllers.music.recency import RecencyWindows
104 from music_assistant.helpers.json import SerializableType
105 from music_assistant.models.player import Player
106
107
108# the container media types worth surfacing as a queue "source" for clients to display. Individual
109# items (single tracks, radio streams, podcast episodes, live audio sources, ...) carry no grouping
110# and only clutter the "playing from" representation, so they are omitted from the wire `sources`.
111_WIRE_SOURCE_MEDIA_TYPES: Final = frozenset(
112 {
113 MediaType.ARTIST,
114 MediaType.ALBUM,
115 MediaType.PLAYLIST,
116 MediaType.PODCAST,
117 MediaType.AUDIOBOOK,
118 }
119)
120
121
122class PlayerQueuesController(QueueLoaderMixin, PlaybackTrackerMixin, StreamFeederMixin):
123 """
124 Controller holding all logic to enqueue music for players.
125
126 The loading, playback-tracking and stream-feeding logic lives in mixins (over the shared base);
127 this class owns the public API surface, the per-queue records and the stateful helper services.
128 """
129
130 def __init__(self, mass: MusicAssistant) -> None:
131 """Initialize core controller."""
132 super().__init__(mass)
133 # server-side per-queue records, keyed by queue_id; each bundles the wire PlayerQueue with
134 # its items, dynamic-source items and runtime-only state (see PlayerQueueData)
135 self._queue_data: dict[str, PlayerQueueData] = {}
136 # stateful helper services (own per-queue state + lifecycle), constructed with self
137 self._autoplay = Autoplay(self)
138 self._smart_shuffle = SmartShuffle(self)
139 self._managed_pool = ManagedPool(self)
140 self._media_resolver = MediaResolver(self)
141 self.manifest.name = "Player Queues controller"
142 self.manifest.description = (
143 "Music Assistant's core controller which manages the queues for all players."
144 )
145 self.manifest.icon = "playlist-music"
146
147 async def close(self) -> None:
148 """Cleanup on exit."""
149 # stop all playback
150 for queue in self.all():
151 if queue.state in (PlaybackState.PLAYING, PlaybackState.PAUSED):
152 await self.stop(queue.queue_id)
153 # flush any pending (debounced) state writes so the latest queue survives shutdown/update
154 for queue in self.all():
155 self.mass.cancel_timer(f"save_queue_cache_{queue.queue_id}")
156 await self._save_queue_to_cache(queue.queue_id)
157
158 async def get_diagnostics(self) -> dict[str, SerializableType]:
159 """Return diagnostics info for this controller to include in diagnostics reports."""
160 queues = [queue_data.queue for queue_data in self._queue_data.values()]
161 by_state: dict[str, int] = {}
162 for queue in queues:
163 by_state[queue.state.value] = by_state.get(queue.state.value, 0) + 1
164 return {
165 "total": len(queues),
166 "active": sum(queue.active for queue in queues),
167 "by_state": by_state,
168 "flow_mode_active": sum(queue.flow_mode for queue in queues),
169 "dynamic_mode_active": sum(queue.is_dynamic for queue in queues),
170 "total_items": sum(queue.items for queue in queues),
171 }
172
173 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
174 """Return the core-module (global) config entries: the queue-controller defaults."""
175 # kept cheap (no library lookup): the config controller populates the global autoplay
176 # playlist dropdown for the UI, so this stays fast on the config value/parse path
177 return core_config_entries(self.mass)
178
179 async def update_config(self, config: CoreConfig, changed_keys: set[str]) -> None:
180 """Apply a global queue-settings change: refresh derived per-queue state and notify clients."""
181 await super().update_config(config, changed_keys)
182 if not any(key.startswith("values/") for key in changed_keys):
183 return
184 # queues that follow a changed global value may flip their derived indicators, so refresh
185 # and signal them (mirrors what save_player_queue_config does for a single queue)
186 for queue in self.all():
187 queue.smart_fades_active = self.mass.streams.is_smart_fades_active(queue)
188 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
189 self.signal_update(queue.queue_id)
190
191 def get_queue_config_entries(
192 self, playlist_options: list[ConfigValueOption] | None = None
193 ) -> list[ConfigEntry]:
194 """
195 Return the per-queue config entries.
196
197 The autoplay_mode select disables the 'similar' option when no provider can supply
198 similar tracks. The crossfade_mode select's options and default depend on whether smart
199 fades are available: the smart option is disabled (shown but not selectable) and the
200 default falls back to standard crossfade when smart fades can't be used on this server.
201
202 :param playlist_options: Library playlists to offer for the 'playlist' autoplay mode.
203 Only populated when serving the entries to the UI; the parse path can omit it.
204 """
205 return queue_config_entries(self.mass, playlist_options)
206
207 def __iter__(self) -> Iterator[PlayerQueue]:
208 """Iterate over (available) players."""
209 return iter(queue_data.queue for queue_data in self._queue_data.values())
210
211 @api_command("player_queues/all", required_scope=Scope.QUEUES_READ)
212 def all(self) -> tuple[PlayerQueue, ...]:
213 """Return all registered PlayerQueues."""
214 return tuple(queue_data.queue for queue_data in self._queue_data.values())
215
216 @api_command("player_queues/get", required_scope=Scope.QUEUES_READ)
217 def get(self, queue_id: str) -> PlayerQueue | None:
218 """Return PlayerQueue by queue_id or None if not found."""
219 queue_data = self._queue_data.get(queue_id)
220 return queue_data.queue if queue_data else None
221
222 def queue_data(self, queue_id: str) -> PlayerQueueData:
223 """
224 Return the server-side record for a queue (raises if the queue is unknown).
225
226 Internal accessor for the stateful helper services so they reach per-queue state through
227 the controller rather than its private store.
228 """
229 return self._queue_data[queue_id]
230
231 def queue_data_or_none(self, queue_id: str) -> PlayerQueueData | None:
232 """Return the server-side record for a queue, or None if it is not registered."""
233 return self._queue_data.get(queue_id)
234
235 @api_command("player_queues/items", required_scope=Scope.QUEUES_READ)
236 def items(self, queue_id: str, limit: int = 500, offset: int = 0) -> list[QueueItem]:
237 """Return all QueueItems for given PlayerQueue."""
238 if (queue_data := self._queue_data.get(queue_id)) is None:
239 return []
240 return queue_data.items[offset : offset + limit]
241
242 @api_command("player_queues/get_active_queue", required_scope=Scope.QUEUES_READ)
243 def get_active_queue(self, player_id: str) -> PlayerQueue | None:
244 """Return the current active/synced queue for a player."""
245 if player := self.mass.players.get_player(player_id):
246 return self.mass.players.get_active_queue(player)
247 return None
248
249 # Queue commands
250
251 @api_command("player_queues/shuffle", required_scope=Scope.QUEUES_CONTROL)
252 async def set_shuffle(self, queue_id: str, shuffle_enabled: bool) -> None:
253 """
254 Configure shuffle setting on the the queue.
255
256 On a session-owned queue the command is forwarded to the external session
257 and may raise InvalidCommand when the session does not support shuffle.
258 """
259 queue = self._queue_data[queue_id].queue
260 if (delegated := self._get_delegated_source(queue_id)) is not None:
261 audio_source, caps, provider = delegated
262 # always forward: the mirrored state lags the session (options echo), so an
263 # equality check here would drop a quick second toggle as a false no-change
264 if not caps.can_shuffle:
265 raise InvalidCommand(
266 "Cannot change shuffle: the external session does not support it"
267 )
268 await provider.on_source_control(
269 audio_source.item_id, SourceControl.SHUFFLE, shuffle_enabled
270 )
271 return
272 if queue.is_dynamic:
273 # a dynamic queue is an always-on, recency-orchestrated smart mix; manual shuffle
274 # (and plain linear order) have no meaning here so the toggle is locked
275 raise InvalidCommand("Cannot change shuffle while the queue is in dynamic mode")
276 if queue.shuffle_enabled == shuffle_enabled:
277 return # no change
278 await self._apply_local_shuffle(queue_id, shuffle_enabled)
279
280 def is_smart_shuffle_active(self, queue: PlayerQueue) -> bool:
281 """
282 Return whether smart shuffle is currently in effect for the queue.
283
284 A dynamic queue is always an orchestrated smart mix (the managed pool), so it always counts
285 as active; otherwise smart shuffle is active when shuffle is on and the per-queue
286 smart-shuffle setting is enabled.
287
288 :param queue: The queue to evaluate.
289 """
290 if queue.is_dynamic:
291 return True
292 return queue.shuffle_enabled and self._smart_shuffle.is_enabled(queue.queue_id)
293
294 @api_command("player_queues/autoplay", required_scope=Scope.QUEUES_CONTROL)
295 def set_autoplay(self, queue_id: str, autoplay_enabled: bool) -> None:
296 """Configure Autoplay setting on the queue."""
297 queue_data = self._queue_data[queue_id]
298 queue = queue_data.queue
299 queue.autoplay_enabled = autoplay_enabled
300 # if we're already at/near the end of the queue, kick off a refill right away
301 # (an active dynamic source manages its own refills, so leave it be)
302 if (
303 queue.autoplay_enabled
304 and not queue.is_dynamic
305 and queue.current_index is not None
306 and (queue.items - queue.current_index) < 5
307 ):
308 task_id = f"fill_autoplay_tracks_{queue_id}"
309 self.mass.call_later(5, self._fill_autoplay_tracks, queue_id, task_id=task_id)
310 self.signal_update(queue_id=queue_id)
311
312 @api_command(
313 "player_queues/dont_stop_the_music", required_scope=Scope.QUEUES_CONTROL, alias=True
314 )
315 def set_dont_stop_the_music(self, queue_id: str, dont_stop_the_music_enabled: bool) -> None:
316 """Backwards-compatible alias for the autoplay command, used by older clients."""
317 self.set_autoplay(queue_id, dont_stop_the_music_enabled)
318
319 @api_command("player_queues/repeat", required_scope=Scope.QUEUES_CONTROL)
320 async def set_repeat(self, queue_id: str, repeat_mode: RepeatMode) -> None:
321 """
322 Configure repeat setting on the the queue.
323
324 On a session-owned queue the command is forwarded to the external session
325 and may raise InvalidCommand when the session does not support repeat.
326 """
327 queue = self._queue_data[queue_id].queue
328 if (delegated := self._get_delegated_source(queue_id)) is not None:
329 audio_source, caps, provider = delegated
330 # always forward: the mirrored state lags the session (options echo), so an
331 # equality check here would drop a quick second toggle as a false no-change
332 if repeat_mode == RepeatMode.UNKNOWN:
333 raise InvalidCommand("Cannot set an unknown repeat mode")
334 if not caps.can_repeat:
335 raise InvalidCommand(
336 "Cannot change repeat: the external session does not support it"
337 )
338 await provider.on_source_control(
339 audio_source.item_id, SourceControl.REPEAT, repeat_mode
340 )
341 return
342 if queue.is_dynamic:
343 # a dynamic queue is an always-on flowing mix of its sources; repeat has no meaning here
344 raise InvalidCommand("Cannot change repeat while the queue is in dynamic mode")
345 if queue.repeat_mode == repeat_mode:
346 return # no change
347 queue.repeat_mode = repeat_mode
348 self.signal_update(queue_id)
349 if (
350 queue.state == PlaybackState.PLAYING
351 and queue.index_in_buffer is not None
352 and queue.index_in_buffer == queue.current_index
353 ):
354 # if the queue is playing,
355 # ensure to (re)queue the next track because it might have changed
356 # note that we only do this if the player has loaded the current track
357 # if not, we wait until it has loaded to prevent conflicts
358 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
359 self._enqueue_next_item(queue_id, next_item)
360
361 @api_command("player_queues/crossfade", required_scope=Scope.QUEUES_CONTROL)
362 def set_crossfade(self, queue_id: str, crossfade_enabled: bool) -> None:
363 """Enable or disable crossfade on the queue."""
364 queue = self._queue_data[queue_id].queue
365 if queue.crossfade_enabled == crossfade_enabled:
366 return # no change
367 queue.crossfade_enabled = crossfade_enabled
368 # refresh the derived smart-fades indicator so the update we signal reflects the new state
369 queue.smart_fades_active = self.mass.streams.is_smart_fades_active(queue)
370 self.signal_update(queue_id)
371 if (
372 queue.state == PlaybackState.PLAYING
373 and queue.index_in_buffer is not None
374 and queue.index_in_buffer == queue.current_index
375 ):
376 # re-enqueue the next track so the new crossfade behaviour applies to the
377 # upcoming transition (only when the player has already loaded the current track)
378 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
379 self._enqueue_next_item(queue_id, next_item)
380
381 @api_command("player_queues/overlay", required_scope=Scope.QUEUES_CONTROL)
382 async def set_overlay(
383 self,
384 queue_id: str,
385 enabled: bool | None = None,
386 source: str | None = None,
387 volume: int | None = None,
388 ) -> None:
389 """
390 Configure the audio overlay for the given queue.
391
392 The audio overlay mixes a looping sound effect (e.g. rain or white noise)
393 into the queue's audio stream. Changes take effect immediately: if the
394 queue is playing, playback is restarted from the current position.
395
396 :param queue_id: queue_id of the queue to configure.
397 :param enabled: Enable or disable the audio overlay. Omit to leave unchanged.
398 :param source: URI of the sound effect item to mix in. Omit to leave unchanged.
399 :param volume: Overlay loudness relative to the music in percent
400 (0-200, 100 = equally loud). Omit to leave unchanged.
401 """
402 queue = self._queue_data[queue_id].queue
403 changed = audible_change = False
404 if source is not None:
405 item = await self.mass.music.get_item_by_uri(source)
406 if item.media_type != MediaType.SOUND_EFFECT:
407 raise InvalidDataError("Audio overlay source must be a sound effect item")
408 mapping = ItemMapping.from_item(cast("SoundEffect", item))
409 if queue.overlay_source != mapping:
410 queue.overlay_source = mapping
411 changed = True
412 audible_change = queue.overlay_enabled
413 if volume is not None:
414 if not (0 <= volume <= 200):
415 raise InvalidDataError(f"Overlay volume must be between 0 and 200, got {volume}")
416 if queue.overlay_volume != volume:
417 queue.overlay_volume = volume
418 changed = True
419 audible_change |= queue.overlay_enabled
420 if enabled is not None and queue.overlay_enabled != enabled:
421 if enabled and queue.overlay_source is None:
422 raise InvalidCommand("Can not enable audio overlay: no overlay source selected")
423 queue.overlay_enabled = enabled
424 changed = audible_change = True
425 if not changed:
426 return
427 self.signal_update(queue_id)
428 if audible_change and queue.state == PlaybackState.PLAYING:
429 # restart playback from the current position so the change is heard
430 # immediately instead of after the player's audio buffer drains
431 await self.resume(queue_id)
432
433 # Two timebases are used in this controller when variable playback speed is in
434 # effect (atempo applied server-side):
435 # "stream-time" â seconds of audio the player has played (post-atempo).
436 # "media-time" â seconds of the original content the listener has heard.
437 # What the user expects to see on the progress bar and what
438 # we use for resume positions.
439 # Conversion: media-time = stream-time x playback_speed.
440 @api_command("player_queues/set_playback_speed", required_scope=Scope.QUEUES_CONTROL)
441 async def set_playback_speed(
442 self, queue_id: str, speed: float, queue_item_id: str | None = None
443 ) -> None:
444 """
445 Set the playback speed for the given queue item.
446
447 Variable playback speed is supported only for audiobooks and podcast episodes.
448
449 If queue_item_id is not provided,
450 the speed will be set for the current item in the queue.
451
452 :param queue_id: queue_id of the queue to configure.
453 :param speed: playback speed multiplier (0.5 to 3.0). 1.0 = normal speed.
454 """
455 if not (0.5 <= speed <= 3.0):
456 raise InvalidDataError(f"Playback speed must be between 0.5 and 3.0, got {speed}")
457 queue = self._queue_data[queue_id].queue
458 if not queue.current_item:
459 raise QueueEmpty("Cannot set playback speed: queue is empty")
460 queue_item_id = queue_item_id or queue.current_item.queue_item_id
461 queue_item = self.get_item(queue_id, queue_item_id)
462 if not queue_item:
463 raise InvalidDataError(f"Queue item {queue_item_id} not found in queue")
464 if queue_item.media_type not in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE):
465 raise InvalidCommand(
466 "Variable playback speed is only supported for audiobooks and podcast episodes"
467 )
468 if not queue_item.duration:
469 raise InvalidCommand("Cannot set playback speed for items with unknown duration")
470 current_speed = float(queue_item.extra_attributes.get("playback_speed") or 1.0)
471 if abs(current_speed - speed) < 0.001:
472 return # no change
473 # use extra_attributes of the queue item to store the playback speed
474 queue_item.extra_attributes["playback_speed"] = speed
475 # mirror onto the queue so corrected_elapsed_time advances in media-time
476 # immediately, before the next on_player_elapsed_time_corrected snapshot.
477 if queue.current_item and queue.current_item.queue_item_id == queue_item_id:
478 # close off the wallclock seconds that already ticked by at the old speed
479 # before switching, so corrected_elapsed_time doesn't multiply them by the new speed
480 if queue.state == PlaybackState.PLAYING:
481 queue.elapsed_time = queue.corrected_elapsed_time
482 queue.elapsed_time_last_updated = time.time()
483 queue.playback_speed = speed
484 self.signal_update(queue_id)
485 if queue.state == PlaybackState.PLAYING:
486 await self.resume(queue_id)
487
488 @api_command(
489 "player_queues/play_media", required_scope=Scope.QUEUES_CONTROL, allow_impersonation=True
490 )
491 async def play_media(
492 self,
493 queue_id: str,
494 media: MediaItemType | ItemMapping | str | list[MediaItemType | ItemMapping | str],
495 option: QueueOption | None = None,
496 radio_mode: bool = False,
497 start_item: PlayableMediaItemType | str | None = None,
498 sort_by: str | None = None,
499 start_from_beginning: bool = False,
500 shuffle: bool | None = None,
501 ) -> None:
502 """
503 Play media item(s) on the given queue.
504
505 :param queue_id: The queue_id of the queue to play media on.
506 :param media: Media that should be played (MediaItem(s) and/or uri's).
507 :param option: Which enqueue mode to use.
508 :param radio_mode: Deprecated â translated to a radio_playlist:// dynamic playlist;
509 prefer enqueuing that URI directly.
510 :param start_item: Optional item to start the playlist or album from.
511 :param sort_by: Optional sort key to order tracks before applying start_item.
512 :param start_from_beginning: Start a podcast episode at position 0, ignoring any
513 saved resume position. The stored progress itself is left untouched.
514 :param shuffle: Play the media shuffled (or explicitly in order). Only applies to the
515 options that start playing right away (play/replace), and never to a dynamic source
516 (an always-on smart mix). Omit to follow the queue's own shuffle setting, which media
517 with an order of its own (album, podcast, episode, audiobook, audio source) switches
518 off; the first item of a batch decides for the whole batch.
519 """
520 self._check_player_permission(queue_id)
521 if not self.get(queue_id):
522 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
523 # the live sources the queue holds may not survive what is being started; remembered
524 # here so their plugins can be told once we know which of them did
525 outgoing_sources = self._audio_sources_in(queue_id)
526 try:
527 # Lock is acquired by the @handle_play_action decorator on the internal handler
528 await self._handle_play_media(
529 queue_id,
530 media,
531 option,
532 radio_mode,
533 start_item,
534 sort_by,
535 start_from_beginning,
536 shuffle,
537 )
538 finally:
539 # also when the media failed to load: the queue may already have given up its items
540 # for it, and the source is just as gone either way
541 self._notify_audio_source_replaced(queue_id, outgoing_sources)
542
543 @api_command("player_queues/move_item", required_scope=Scope.QUEUES_CONTROL)
544 def move_item(self, queue_id: str, queue_item_id: str, pos_shift: int = 1) -> None:
545 """
546 Move queue item x up/down the queue.
547
548 - queue_id: id of the queue to process this request.
549 - queue_item_id: the item_id of the queueitem that needs to be moved.
550 - pos_shift: move item x positions down if positive value
551 - pos_shift: move item x positions up if negative value
552 - pos_shift: move item to top of queue as next item if 0.
553 """
554 queue = self._queue_data[queue_id].queue
555 item_index = self.index_by_id(queue_id, queue_item_id)
556 if item_index is None:
557 raise InvalidDataError(f"Item {queue_item_id} not found in queue")
558 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
559 msg = f"{item_index} is already played/buffered"
560 raise IndexError(msg)
561
562 queue_items = self._queue_data[queue_id].items
563 queue_items = queue_items.copy()
564
565 if pos_shift == 0 and queue.state == PlaybackState.PLAYING:
566 new_index = (queue.current_index or 0) + 1
567 elif pos_shift == 0:
568 new_index = queue.current_index or 0
569 else:
570 new_index = item_index + pos_shift
571 if (new_index < (queue.current_index or 0)) or (new_index > len(queue_items)):
572 return
573 # move the item in the list
574 queue_items.insert(new_index, queue_items.pop(item_index))
575 self.update_items(queue_id, queue_items)
576
577 @api_command("player_queues/move_item_end", required_scope=Scope.QUEUES_CONTROL)
578 def move_item_end(self, queue_id: str, queue_item_id: str) -> None:
579 """
580 Move queue item to the end the queue.
581
582 - queue_id: id of the queue to process this request.
583 - queue_item_id: the item_id of the queueitem that needs to be moved.
584 """
585 queue = self._queue_data[queue_id].queue
586 item_index = self.index_by_id(queue_id, queue_item_id)
587 if item_index is None:
588 raise InvalidDataError(f"Item {queue_item_id} not found in queue")
589 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
590 msg = f"{item_index} is already played/buffered"
591 raise IndexError(msg)
592
593 queue_items = self._queue_data[queue_id].items
594 if item_index == (len(queue_items) - 1):
595 return
596 queue_items = queue_items.copy()
597
598 new_index = len(self._queue_data[queue_id].items) - 1
599
600 # move the item in the list
601 queue_items.insert(new_index, queue_items.pop(item_index))
602 self.update_items(queue_id, queue_items)
603
604 @api_command("player_queues/delete_item", required_scope=Scope.QUEUES_CONTROL)
605 def delete_item(self, queue_id: str, item_id_or_index: int | str) -> None:
606 """Delete item (by id or index) from the queue."""
607 if isinstance(item_id_or_index, str):
608 item_index = self.index_by_id(queue_id, item_id_or_index)
609 if item_index is None:
610 raise InvalidDataError(f"Item {item_id_or_index} not found in queue")
611 else:
612 item_index = item_id_or_index
613 queue = self._queue_data[queue_id].queue
614 if queue.index_in_buffer is not None and item_index <= queue.index_in_buffer:
615 # ignore request if track already loaded in the buffer
616 # the frontend should guard so this is just in case
617 self.logger.warning("delete requested for item already loaded in buffer")
618 return
619 queue_items = self._queue_data[queue_id].items.copy()
620 queue_items.pop(item_index)
621 self.update_items(queue_id, queue_items)
622
623 @api_command("player_queues/clear", required_scope=Scope.QUEUES_CONTROL)
624 def clear(self, queue_id: str, skip_stop: bool = False) -> None:
625 """Clear all items in the queue, switching shuffle off with them."""
626 # Only the explicit clear notifies from here, never `_clear` itself: a replace or a
627 # transfer clears as a first step and may re-select the very same source right after,
628 # where releasing it would cost the user their playback position. Those two decide for
629 # themselves, in `play_media` and `transfer_queue`.
630 self._notify_audio_source_removed(queue_id)
631 self._clear(queue_id, skip_stop)
632 # clearing is an explicit "start over" gesture by the user, so a shuffle that belonged to
633 # the discarded content must not carry over into whatever is played next
634 self._reset_shuffle(queue_id)
635
636 def mark_ended(self, queue_id: str) -> None:
637 """
638 Mark a queue as played to its end, keeping its items so it can be replayed.
639
640 The playback position is parked on the last item rather than cleared: a null index is
641 indistinguishable from a queue that was loaded but never started, and an index past the
642 end is silently misread by everything that does arithmetic on it. `ended` is what tells
643 clients the queue finished, and pressing play starts it over from the first item.
644
645 :param queue_id: The queue_id of the queue that reached its end.
646 """
647 queue_data = self._queue_data[queue_id]
648 queue = queue_data.queue
649 if not queue_data.items:
650 # nothing to replay, so there is nothing to advertise as finished either
651 self._clear(queue_id)
652 return
653 self.mass.streams.audio_processing.clear(queue_id)
654 queue.ended = True
655 queue.current_index = len(queue_data.items) - 1
656 queue.current_item = queue_data.items[-1]
657 queue.next_item = None
658 queue.elapsed_time = 0
659 queue.elapsed_time_last_updated = time.time()
660 queue.index_in_buffer = None
661 queue.resume_pos = 0
662 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
663 self.signal_update(queue_id)
664
665 @api_command("player_queues/save_as_playlist", required_scope=Scope.LIBRARY_WRITE)
666 async def save_as_playlist(self, queue_id: str, name: str) -> BackgroundTask:
667 """
668 Save the current queue items as a new playlist.
669
670 :param queue_id: The queue_id of the queue to save.
671 :param name: The name for the new playlist.
672 """
673 if not self.get(queue_id):
674 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
675 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
676 if not queue_items:
677 raise QueueEmpty("Cannot save an empty queue as a playlist.")
678 # collect URIs from queue items that are playlist-compatible
679 uris: list[str] = []
680 for item in queue_items:
681 if item.uri and item.media_type in PLAYLIST_MEDIA_TYPES:
682 uris.append(item.uri)
683 if not uris:
684 raise InvalidDataError("No valid items in queue to save as playlist.")
685 playlist = await self.mass.music.playlists.create_playlist(name)
686 return await self.mass.music.playlists.add_playlist_tracks(playlist.item_id, uris)
687
688 @api_command("player_queues/stop", required_scope=Scope.QUEUES_CONTROL)
689 @handle_play_action
690 async def stop(self, queue_id: str) -> None:
691 """
692 Handle STOP command for given queue.
693
694 - queue_id: queue_id of the playerqueue to handle the command.
695 """
696 self._check_player_permission(queue_id)
697 # cancel any pending play_index calls for this queue to prevent conflicts
698 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
699 # cancel in-flight preload/enqueue-next so it can't enqueue after stop
700 self.mass.cancel_task(f"preload_next_item_{queue_id}")
701 self.mass.cancel_timer(f"enqueue_next_item_{queue_id}")
702 self.mass.cancel_task(f"enqueue_next_item_{queue_id}")
703 self._set_transitioning(queue_id, False)
704 queue_data = self._queue_data[queue_id]
705 session_id = queue_data.session_id
706 queue_player = self.mass.players.get_player(queue_id, True)
707 if queue_player is None:
708 raise PlayerUnavailableError(f"Player {queue_id} is not available")
709 if (queue := self.get(queue_id)) and queue.active:
710 if queue.state == PlaybackState.PLAYING:
711 queue.resume_pos = int(queue.corrected_elapsed_time)
712 # Use internal handler to avoid circular redirect:
713 # public cmd_stop redirects to queue.stop when a queue is active,
714 # which would loop back here indefinitely.
715 await self.mass.players._handle_cmd_stop(queue_id)
716 if queue_data.session_id == session_id:
717 queue_data.session_id = None
718 self.mass.streams.audio_processing.clear(queue_id, session_id)
719 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
720
721 @api_command("player_queues/play", required_scope=Scope.QUEUES_CONTROL)
722 async def play(self, queue_id: str) -> None:
723 """
724 Handle PLAY command for given queue.
725
726 :param queue_id: queue_id of the playerqueue to handle the command.
727 """
728 self._check_player_permission(queue_id)
729 if not self.get(queue_id):
730 raise PlayerUnavailableError(f"Queue {queue_id} is not available")
731 await self._handle_play(queue_id)
732
733 @api_command("player_queues/pause", required_scope=Scope.QUEUES_CONTROL)
734 async def pause(self, queue_id: str) -> None:
735 """
736 Handle PAUSE command for given queue.
737
738 - queue_id: queue_id of the playerqueue to handle the command.
739 """
740 self._check_player_permission(queue_id)
741 # cancel any pending play_index calls for this queue to prevent conflicts
742 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
743 self._set_transitioning(queue_id, False)
744 if not (queue := self.get(queue_id)):
745 return
746 queue_active = queue.active
747 if queue.active and queue.state == PlaybackState.PLAYING:
748 queue.resume_pos = int(queue.corrected_elapsed_time)
749 # Use internal handler to avoid circular redirect
750 # (cmd_pause redirects to queue.pause, which calls cmd_pause again)
751 await self.mass.players._handle_cmd_pause(queue_id)
752
753 async def _watch_pause(player: Player) -> None:
754 count = 0
755 # wait for pause
756 while count < 5 and player.state.playback_state == PlaybackState.PLAYING:
757 count += 1
758 await asyncio.sleep(1)
759 # wait for unpause
760 if player.state.playback_state != PlaybackState.PAUSED:
761 return
762 count = 0
763 while count < 30 and player.state.playback_state == PlaybackState.PAUSED:
764 count += 1
765 await asyncio.sleep(1)
766 # if player is still paused when the limit is reached, send stop
767 if player.state.playback_state == PlaybackState.PAUSED:
768 await self.stop(queue_id)
769
770 # we auto stop a player from paused when its paused for 30 seconds
771 if (
772 queue_active
773 and (queue_player := self.mass.players.get_player(queue_id))
774 and not queue_player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS)
775 ):
776 self.mass.create_task(_watch_pause(queue_player))
777
778 @api_command("player_queues/play_pause", required_scope=Scope.QUEUES_CONTROL)
779 async def play_pause(self, queue_id: str) -> None:
780 """
781 Toggle play/pause on given playerqueue.
782
783 - queue_id: queue_id of the queue to handle the command.
784 """
785 if (queue := self.get(queue_id)) and queue.state == PlaybackState.PLAYING:
786 await self.pause(queue_id)
787 return
788 await self.play(queue_id)
789
790 @api_command("player_queues/next", required_scope=Scope.QUEUES_CONTROL)
791 @handle_play_action
792 async def next(self, queue_id: str) -> None:
793 """
794 Handle NEXT TRACK command for given queue.
795
796 :param queue_id: queue_id of the queue to handle the command.
797 """
798 self._check_player_permission(queue_id)
799 if (queue := self.get(queue_id)) is None or not queue.active:
800 raise InvalidCommand(f"Queue {queue_id} is not active")
801 if (active := self._get_current_audio_source(queue_id)) is not None:
802 audio_source, provider = active
803 # gate on the per-action flag alone so a transport-only source (no
804 # queue_capabilities) also skips within its own session
805 if audio_source.can_next_previous:
806 await provider.on_source_control(audio_source.item_id, SourceControl.NEXT)
807 return
808 if audio_source.queue_capabilities is not None:
809 raise InvalidCommand("Cannot skip: the external session does not support skipping")
810 # a live source without skip support: fall through to the MA index walk
811 self._set_transitioning(queue_id, True)
812 idx = self._queue_data[queue_id].queue.current_index
813 if idx is None:
814 self.logger.warning("Queue %s has no current index", queue.display_name)
815 self._set_transitioning(queue_id, False)
816 return
817 next_index = self._get_next_index(queue_id, idx, True)
818 if next_index is None:
819 self._set_transitioning(queue_id, False)
820 return
821
822 # immediately update current item so UI shows the new track right away
823 queue.current_index = next_index
824 queue.current_item = self.get_item(queue_id, next_index)
825 queue.elapsed_time = 0
826 queue.elapsed_time_last_updated = time.time()
827 self.signal_update(queue_id)
828 if queue_player := self.mass.players.get_player(queue_id, True):
829 queue_player.update_state()
830
831 # debounce rapid next button presses using call_later
832 self.mass.call_later(
833 1,
834 self.play_index,
835 queue_id,
836 next_index,
837 task_id=f"queue_play_index_{queue_id}",
838 )
839
840 @api_command("player_queues/previous", required_scope=Scope.QUEUES_CONTROL)
841 @handle_play_action
842 async def previous(self, queue_id: str) -> None:
843 """
844 Handle PREVIOUS TRACK command for given queue.
845
846 :param queue_id: queue_id of the queue to handle the command.
847 """
848 self._check_player_permission(queue_id)
849 if (queue := self.get(queue_id)) is None or not queue.active:
850 raise InvalidCommand(f"Queue {queue_id} is not active")
851 if (active := self._get_current_audio_source(queue_id)) is not None:
852 audio_source, provider = active
853 # gate on the per-action flag alone so a transport-only source (no
854 # queue_capabilities) also skips within its own session
855 if audio_source.can_next_previous:
856 await provider.on_source_control(audio_source.item_id, SourceControl.PREVIOUS)
857 return
858 if audio_source.queue_capabilities is not None:
859 raise InvalidCommand("Cannot skip: the external session does not support skipping")
860 # a live source without skip support: fall through to the MA index walk
861 self._set_transitioning(queue_id, True)
862 current_index = self._queue_data[queue_id].queue.current_index
863 if current_index is None:
864 self._set_transitioning(queue_id, False)
865 return
866 prev_index = int(current_index)
867 # restart current track if elapsed > 5s, otherwise go to previous
868 if self._queue_data[queue_id].queue.elapsed_time < 5:
869 prev_index = max(current_index - 1, 0)
870
871 # immediately update current item so UI shows the new track right away
872 queue.current_index = prev_index
873 queue.current_item = self.get_item(queue_id, prev_index)
874 queue.elapsed_time = 0
875 queue.elapsed_time_last_updated = time.time()
876 self.signal_update(queue_id)
877 if queue_player := self.mass.players.get_player(queue_id, True):
878 queue_player.update_state()
879
880 # debounce rapid previous button presses using call_later
881 self.mass.call_later(
882 1,
883 self.play_index,
884 queue_id,
885 prev_index,
886 task_id=f"queue_play_index_{queue_id}",
887 )
888
889 @api_command("player_queues/skip", required_scope=Scope.QUEUES_CONTROL)
890 async def skip(self, queue_id: str, seconds: int = 10) -> None:
891 """
892 Handle SKIP command for given queue.
893
894 - queue_id: queue_id of the queue to handle the command.
895 - seconds: number of seconds to skip in track. Use negative value to skip back.
896 """
897 if (queue := self.get(queue_id)) is None or not queue.active:
898 raise InvalidCommand(f"Queue {queue_id} is not active")
899 await self.seek(queue_id, int(self._queue_data[queue_id].queue.elapsed_time + seconds))
900
901 @api_command("player_queues/seek", required_scope=Scope.QUEUES_CONTROL)
902 async def seek(self, queue_id: str, position: int = 10) -> None:
903 """
904 Handle SEEK command for given queue.
905
906 - queue_id: queue_id of the queue to handle the command.
907 - position: position in seconds to seek to in the current playing item.
908 """
909 if (queue := self.get(queue_id)) is None or not queue.active:
910 raise InvalidCommand(f"Queue {queue_id} is not active")
911 if (active := self._get_current_audio_source(queue_id)) is not None:
912 audio_source, provider = active
913 # gate on the per-action flag alone so a transport-only source (no
914 # queue_capabilities) also seeks within its own session
915 if audio_source.can_seek:
916 position = max(0, int(position))
917 current_item = queue.current_item
918 await provider.on_source_control(audio_source.item_id, SourceControl.SEEK, position)
919 # publish the seek target so the progress bar does not snap back to the
920 # last position report (mirrors the non-delegated path below) â unless a
921 # concurrent play/stop replaced the current item during the forward
922 if queue.current_item is current_item:
923 queue.elapsed_time = position
924 queue.elapsed_time_last_updated = time.time()
925 self.signal_update(queue_id)
926 return
927 if audio_source.queue_capabilities is not None:
928 raise InvalidCommand("Cannot seek: the external session does not support seeking")
929 # a non-seekable live source falls through: the duration guard below rejects it
930 queue_player = self.mass.players.get_player(queue_id, True)
931 if queue_player is None:
932 raise PlayerUnavailableError(f"Player {queue_id} is not available")
933 if not queue.current_item:
934 raise InvalidCommand(f"Queue {queue_player.state.name} has no item(s) loaded.")
935 if not queue.current_item.duration:
936 raise InvalidCommand("Can not seek items without duration.")
937 position = max(0, int(position))
938 if position > queue.current_item.duration:
939 raise InvalidCommand("Can not seek outside of duration range.")
940 if queue.current_index is None:
941 raise InvalidCommand(f"Queue {queue_player.state.name} has no current index.")
942 # Publish the seek target before rebuilding the stream to prevent progress snapback.
943 queue.elapsed_time = position
944 queue.elapsed_time_last_updated = time.time()
945 self.signal_update(queue_id)
946 await self.play_index(queue_id, queue.current_index, seek_position=position)
947
948 @api_command("player_queues/resume", required_scope=Scope.QUEUES_CONTROL)
949 @handle_play_action
950 async def resume(self, queue_id: str, fade_in: bool | None = None) -> None:
951 """
952 Handle RESUME command for given queue.
953
954 - queue_id: queue_id of the queue to handle the command.
955 """
956 self._check_player_permission(queue_id)
957 queue = self._queue_data[queue_id].queue
958 queue_items = self._queue_data[queue_id].items
959 resume_item = queue.current_item
960 if queue.state == PlaybackState.PLAYING:
961 # resume requested while already playing,
962 # use current position as resume position
963 resume_pos = queue.corrected_elapsed_time
964 fade_in = False
965 else:
966 resume_pos = queue.resume_pos or queue.elapsed_time
967
968 if queue.ended and len(queue_items) > 0:
969 # the queue played to its end and is parked on its last item,
970 # so pressing play starts it over from the beginning
971 resume_item = queue_items[0]
972 resume_pos = 0
973 elif not resume_item and queue.current_index is not None and len(queue_items) > 0:
974 resume_item = self.get_item(queue_id, queue.current_index)
975 resume_pos = 0
976 elif not resume_item and queue.current_index is None and len(queue_items) > 0:
977 # items available in queue but no previous track, start at 0
978 resume_item = self.get_item(queue_id, 0)
979 resume_pos = 0
980
981 if resume_item is not None:
982 queue_player = self.mass.players.get_player(queue_id)
983 if queue_player is None:
984 raise PlayerUnavailableError(f"Player {queue_id} is not available")
985 if (
986 fade_in is None
987 and queue_player.state.playback_state == PlaybackState.IDLE
988 and (time.time() - queue.elapsed_time_last_updated) > 60
989 ):
990 # enable fade in effect if the player is idle for a while
991 fade_in = resume_pos > 0
992 if resume_item.media_type == MediaType.RADIO:
993 # we're not able to skip in online radio so this is pointless
994 resume_pos = 0
995 await self.play_index(
996 queue_id, resume_item.queue_item_id, int(resume_pos), fade_in or False
997 )
998 else:
999 msg = f"Resume queue requested but queue {queue.display_name} is empty"
1000 raise QueueEmpty(msg)
1001
1002 @api_command("player_queues/play_index", required_scope=Scope.QUEUES_CONTROL)
1003 @handle_play_action
1004 async def play_index( # noqa: PLR0915
1005 self,
1006 queue_id: str,
1007 index: int | str,
1008 seek_position: int = 0,
1009 fade_in: bool = False,
1010 ) -> None:
1011 """Play item at index (or item_id) X in queue."""
1012 self._check_player_permission(queue_id)
1013 # cancel any pending play_index calls for this queue to prevent conflicts
1014 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
1015 # we set a flag to notify the update logic that we're transitioning to a new track
1016 self._set_transitioning(queue_id, True)
1017 try:
1018 queue_data = self._queue_data[queue_id]
1019 queue = queue_data.queue
1020 queue.resume_pos = 0
1021 # A queue picked up from its end plays its items over from the start, so a resume point
1022 # left on an audiobook/episode must not pull it back to where it was left off. The flag
1023 # itself is only cleared once an item actually loaded below, so a start that never got
1024 # off the ground leaves the queue finished instead of stranding it without a position.
1025 restarting_ended_queue = queue.ended
1026 if isinstance(index, str):
1027 temp_index = self.index_by_id(queue_id, index)
1028 if temp_index is None:
1029 raise InvalidDataError(f"Item {index} not found in queue")
1030 index = temp_index
1031 # At this point index is guaranteed to be int
1032 queue.index_in_buffer = index
1033 queue_data.flow_mode_stream_log = []
1034 queue_data.flow_buffer_completed = None
1035 queue_data.flow_queue_exhausted = None
1036 target_player = self.mass.players.get_player(queue_id)
1037 if target_player is None:
1038 raise PlayerUnavailableError(f"Player {queue_id} is not available")
1039 queue_data.next_item_id_enqueued = None
1040 # always update session id when we start a new playback session
1041 queue_data.session_id = shortuuid.random(length=8)
1042 self.mass.streams.audio_processing.start_session(
1043 queue_id,
1044 queue_data.session_id,
1045 )
1046 # handle resume point of audiobook(chapter) or podcast(episode)
1047 if (
1048 not seek_position
1049 and not restarting_ended_queue
1050 and (queue_item := self.get_item(queue_id, index))
1051 and (resume_position_ms := getattr(queue_item.media_item, "resume_position_ms", 0))
1052 ):
1053 # the client may have fetched the item before its duration was known
1054 await self._restore_probed_duration(queue_item)
1055 if queue_item.duration or getattr(queue_item.media_item, "duration", 0):
1056 seek_position = max(0, int((resume_position_ms - 500) / 1000))
1057 else:
1058 # seeking needs a duration, which is determined while streaming
1059 self.logger.debug(
1060 "Can not resume %s at %ss: its duration is not known (yet)",
1061 queue_item.name,
1062 int(resume_position_ms / 1000),
1063 )
1064
1065 # restore the persisted playback speed for a freshly queued audiobook/episode
1066 # (an in-session item already carries its speed in extra_attributes)
1067 if (
1068 (queue_item := self.get_item(queue_id, index))
1069 and queue_item.media_item is not None
1070 and queue_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1071 and "playback_speed" not in queue_item.extra_attributes
1072 ):
1073 stored_speed = await self.mass.music.get_playback_speed(
1074 cast("Audiobook | PodcastEpisode", queue_item.media_item),
1075 userid=queue_data.userid,
1076 )
1077 if stored_speed != 1.0:
1078 queue_item.extra_attributes["playback_speed"] = stored_speed
1079
1080 # try to load the item, retry with next item if it fails
1081 for attempt in range(5):
1082 try:
1083 queue_item = self.get_item(queue_id, index)
1084 if not queue_item:
1085 continue # guard
1086 await self._load_item(
1087 queue_item,
1088 self._get_next_index(queue_id, index),
1089 is_start=True,
1090 seek_position=seek_position if attempt == 0 else 0,
1091 fade_in=fade_in if attempt == 0 else False,
1092 )
1093 # if we reach this point, loading the item succeeded, break the loop
1094 queue.current_index = index
1095 queue.current_item = queue_item
1096 # playback is under way, so the queue is no longer sitting at its end
1097 queue.ended = False
1098 # reset the elapsed clock together with the item switch (like
1099 # next/previous do), so queue updates signaled before the player
1100 # reports position don't carry the previous item's elapsed_time
1101 queue.elapsed_time = seek_position if attempt == 0 else 0
1102 queue.elapsed_time_last_updated = time.time()
1103 break
1104 except (MediaNotFoundError, AudioError) as err:
1105 item_name = queue_item.name if queue_item else "unknown"
1106 if isinstance(err, ProviderStreamLimitError):
1107 # the requested item is playable, its provider is just at capacity:
1108 # report that instead of silently advancing to another item
1109 self.logger.error("%s", err)
1110 await self.stop(queue_id)
1111 raise
1112 # Only MediaNotFoundError (item unreachable) is persistent;
1113 # keep AudioError items available so a retry can resurface
1114 # the same actionable error.
1115 if queue_item and isinstance(err, MediaNotFoundError):
1116 queue_item.available = False
1117 next_index = self._get_next_index(queue_id, index, allow_repeat=False)
1118 if next_index is None:
1119 # Surface an AudioError's own (actionable) message;
1120 # MediaNotFoundError gets the generic wording.
1121 if isinstance(err, AudioError) and str(err):
1122 msg = str(err)
1123 else:
1124 msg = f"Playback failed for {item_name} - no more tracks available"
1125 self.logger.error(msg)
1126 await self.stop(queue_id)
1127 raise MediaNotFoundError(msg) from err
1128 self.logger.warning(
1129 "Skipping unplayable item %s",
1130 item_name,
1131 )
1132 index = next_index
1133 else:
1134 # all attempts to find a playable item failed
1135 await self.stop(queue_id)
1136 raise MediaNotFoundError("No playable item found to start playback")
1137
1138 # Reset flow_mode - the streams controller will set it if flow mode is used.
1139 queue.flow_mode = False
1140 player_media = await self.player_media_from_queue_item(queue_item)
1141 # Hold the play action until the player confirms playback so the UI keeps
1142 # showing the command as in progress instead of falling back to a play button
1143 # for the time the player still needs to connect and start. The queue update
1144 # for the new item goes out first, so the item shows while it is starting.
1145 async with self.mass.players.wait_for_player_update(
1146 queue_id,
1147 attribute_name="playback_state",
1148 attribute_value=PlaybackState.PLAYING,
1149 timeout=PLAYBACK_START_TIMEOUT,
1150 ):
1151 await self.mass.players.play_media(queue_id, player_media)
1152 queue.current_index = index
1153 queue.current_item = queue_item
1154 self.signal_update(queue_id)
1155 finally:
1156 self._set_transitioning(queue_id, False)
1157
1158 @api_command("player_queues/transfer", required_scope=Scope.QUEUES_CONTROL)
1159 async def transfer_queue(
1160 self,
1161 source_queue_id: str,
1162 target_queue_id: str,
1163 auto_play: bool | None = None,
1164 ) -> None:
1165 """Transfer queue to another queue."""
1166 if not (source_queue := self.get(source_queue_id)):
1167 raise PlayerUnavailableError(f"Queue {source_queue_id} is not available")
1168 if not (target_queue := self.get(target_queue_id)):
1169 raise PlayerUnavailableError(f"Queue {target_queue_id} is not available")
1170 if auto_play is None:
1171 auto_play = source_queue.state == PlaybackState.PLAYING
1172
1173 target_player = self.mass.players.get_player(target_queue_id)
1174 if target_player is None:
1175 raise PlayerUnavailableError(f"Player {target_queue_id} is not available")
1176 if target_player.state.active_group or target_player.state.synced_to:
1177 # edge case: the user wants to move playback from the group as a whole, to a single
1178 # player in the group or it is grouped and the command targeted at the single player.
1179 # We need to dissolve the group/sync first, and wait for the state to actually
1180 # propagate before we hand the queue over to the target player.
1181 group_id = target_player.state.active_group or target_player.state.synced_to
1182 assert group_id is not None # checked in if condition above
1183 # For an ad-hoc sync group (target is a sync member of a regular leader),
1184 # ungroup the target itself so only it is freed - ungrouping the leader would
1185 # transfer leadership to a remaining member and recurse back into this method.
1186 # For a virtual group player (active_group), release the group so its static
1187 # members are handled correctly.
1188 ungroup_target = (
1189 target_queue_id
1190 if target_player.state.synced_to and not target_player.state.active_group
1191 else group_id
1192 )
1193 async with self.mass.players.wait_for_player_update(
1194 target_queue_id,
1195 attribute_name=(
1196 "active_group" if target_player.state.active_group else "synced_to"
1197 ),
1198 attribute_value=None,
1199 timeout=5,
1200 ):
1201 await self.mass.players.cmd_ungroup(ungroup_target)
1202
1203 # capture source state before stopping (stop resets these)
1204 source_items = self._queue_data[source_queue_id].items
1205 if source_queue.state == PlaybackState.PLAYING:
1206 # use the live playback clock while actively playing
1207 source_resume_pos = int(source_queue.corrected_elapsed_time)
1208 else:
1209 # when not playing the live clock is stale, so use the stored resume position
1210 source_resume_pos = int(source_queue.resume_pos or source_queue.elapsed_time or 0)
1211 source_current_index = source_queue.current_index
1212 source_current_item = source_queue.current_item
1213
1214 # stop the source player synchronously to prevent the async stop from
1215 # clear() racing with the target's sync group formation/protocol switching
1216 if source_queue.state != PlaybackState.IDLE:
1217 await self.stop(source_queue_id)
1218
1219 target_queue.repeat_mode = source_queue.repeat_mode
1220 target_queue.shuffle_enabled = source_queue.shuffle_enabled
1221 target_queue.crossfade_enabled = source_queue.crossfade_enabled
1222 # refresh the derived smart-fades indicator for the target's own config/availability
1223 target_queue.smart_fades_active = self.mass.streams.is_smart_fades_active(target_queue)
1224 target_queue.autoplay_enabled = source_queue.autoplay_enabled
1225 self._queue_data[target_queue_id].source_items = list(
1226 self._queue_data[source_queue_id].source_items
1227 )
1228 target_queue.sources = list(source_queue.sources)
1229 target_queue.is_dynamic = source_queue.is_dynamic
1230 target_queue.smart_shuffle_active = self.is_smart_shuffle_active(target_queue)
1231 self._queue_data[target_queue_id].enqueued_media_items = list(
1232 self._queue_data[source_queue_id].enqueued_media_items
1233 )
1234 target_queue.resume_pos = source_resume_pos
1235 # the target's own contents are about to be overwritten by the transferred ones
1236 target_outgoing_sources = self._audio_sources_in(target_queue_id)
1237 target_queue.current_index = source_current_index
1238 if source_current_item:
1239 target_queue.current_item = source_current_item
1240 target_queue.current_item.queue_id = target_queue_id
1241 # every live source the queue holds moves with it, not just the one that was playing
1242 transferred_sources = self._audio_sources_in(source_queue_id)
1243 self._clear(source_queue_id, skip_stop=True)
1244
1245 await self.load(target_queue_id, source_items, keep_remaining=False, keep_played=False)
1246 for item in source_items:
1247 item.queue_id = target_queue_id
1248 self.update_items(target_queue_id, source_items)
1249 # a live source the target was holding has just been displaced by the transferred queue
1250 self._notify_audio_source_replaced(target_queue_id, target_outgoing_sources)
1251 await self._notify_audio_source_transferred(
1252 transferred_sources, source_queue_id, target_queue_id
1253 )
1254 if auto_play:
1255 await self.resume(target_queue_id)
1256
1257 # Interaction with player
1258
1259 async def on_player_register(self, player: Player) -> None:
1260 """Register PlayerQueue for given player/queue id."""
1261 queue_id = player.player_id
1262 queue_data: PlayerQueueData | None = None
1263 # try to restore previous state
1264 try:
1265 if prev_state := await self.mass.cache.get(
1266 key=queue_id,
1267 provider=self.domain,
1268 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1269 ):
1270 prev_items = await self.mass.cache.get(
1271 key=queue_id,
1272 provider=self.domain,
1273 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1274 default=[],
1275 )
1276 queue_data = PlayerQueueData.from_cache(prev_state, prev_items)
1277 except Exception as err:
1278 self.logger.warning(
1279 "Failed to restore the queue(items) for %s - %s",
1280 player.state.name,
1281 str(err),
1282 )
1283 # Reset to clean state on failure
1284 queue_data = None
1285 if queue_data is None:
1286 queue_data = PlayerQueueData(
1287 queue=PlayerQueue(
1288 queue_id=queue_id,
1289 active=False,
1290 display_name=player.state.name,
1291 available=player.state.available,
1292 # Autoplay starts out on for a brand new queue; the player's own Autoplay
1293 # switch owns it from here on (and is restored above for a queue we know)
1294 autoplay_enabled=True,
1295 items=0,
1296 )
1297 )
1298
1299 self._queue_data[queue_id] = queue_data
1300 # always call update to calculate state etc
1301 self.on_player_update(player, {})
1302 self.mass.signal_event(EventType.QUEUE_ADDED, object_id=queue_id, data=queue_data.queue)
1303
1304 def on_player_update(
1305 self,
1306 player: Player,
1307 changed_values: dict[str, tuple[Any, Any]],
1308 ) -> None:
1309 """
1310 Call when a PlayerQueue needs to be updated (e.g. when player updates).
1311
1312 NOTE: This is called every second if the player is playing.
1313 """
1314 if player.type == PlayerType.PROTOCOL:
1315 # protocol players do not have a queue on their own
1316 return
1317 queue_id = player.player_id
1318 if (queue := self.get(queue_id)) is None:
1319 # race condition
1320 return
1321 if player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS):
1322 # do nothing while the announcement is in progress
1323 return
1324 # determine if this queue is currently active for this player
1325 queue.active = player.state.active_source in (queue.queue_id, None)
1326 if not queue.active and self._queue_data[queue_id].prev_state is None:
1327 queue.state = PlaybackState.IDLE
1328 # return early if the queue is not active and we have no previous state
1329 return
1330 if self._queue_data[queue_id].transitioning:
1331 # we're currently transitioning to a new track,
1332 # ignore updates from the player during this time
1333 return
1334 # queue is active and preflight checks passed, update the queue details
1335 self._update_queue_from_player(player)
1336
1337 def on_player_elapsed_time_corrected(self, player: Player) -> None:
1338 """Correct the queue's timing base if the player's real elapsed_time diverged."""
1339 if player.type == PlayerType.PROTOCOL:
1340 return
1341 queue_id = player.player_id
1342 if (queue := self.get(queue_id)) is None:
1343 return
1344 if not queue.active:
1345 return
1346 player_elapsed = player.state.corrected_elapsed_time
1347 if player_elapsed is None:
1348 return
1349 now = time.time()
1350 # queue.elapsed_time is stored in media-time so it can be displayed and
1351 # used as a resume position directly. The player reports stream-time
1352 # (post-atempo), so we scale by the current item's playback_speed.
1353 speed = get_current_playback_speed(queue)
1354 if queue.flow_mode:
1355 # _get_flow_queue_stream_index returns media-time in the current item
1356 # using each playlog entry's recorded speed.
1357 _, elapsed_time = self._get_flow_queue_stream_index(queue, player)
1358 else:
1359 elapsed_time = player_elapsed * speed
1360 if queue.current_item and queue.current_item.streamdetails:
1361 if seek_pos := queue.current_item.streamdetails.seek_position:
1362 elapsed_time += seek_pos
1363 queue.elapsed_time = elapsed_time
1364 queue.elapsed_time_last_updated = now
1365 queue.playback_speed = speed
1366 self.mass.signal_event(
1367 EventType.QUEUE_TIME_UPDATED,
1368 object_id=queue_id,
1369 data=queue.elapsed_time,
1370 )
1371
1372 def on_player_remove(self, player_id: str, permanent: bool) -> None:
1373 """Call when a player is removed from the registry."""
1374 self.mass.streams.audio_processing.clear(player_id)
1375 # cancel any pending play_index calls for this queue to prevent conflicts
1376 self.mass.cancel_timer(f"queue_play_index_{player_id}")
1377 # cancel a pending debounced cache write AND an already-started one, so neither can
1378 # recreate a deleted entry after the player is gone (the timer becomes a task once it fires)
1379 self.mass.cancel_timer(f"save_queue_cache_{player_id}")
1380 self.mass.cancel_task(f"save_queue_cache_{player_id}")
1381 self._set_transitioning(player_id, False)
1382 if permanent:
1383 self.purge_saved_queue(player_id)
1384 self._queue_data.pop(player_id, None)
1385 self._managed_pool.forget(player_id)
1386
1387 def purge_saved_queue(self, queue_id: str) -> None:
1388 """Delete the persisted state and items of the given queue."""
1389 for category in (CACHE_CATEGORY_PLAYER_QUEUE_STATE, CACHE_CATEGORY_PLAYER_QUEUE_ITEMS):
1390 # a removal runs both the player teardown and the config cleanup, so keep the
1391 # delete to one task per category instead of one per caller
1392 self.mass.create_task(
1393 self.mass.cache.delete(
1394 key=queue_id,
1395 provider=self.domain,
1396 category=category,
1397 ),
1398 task_id=f"purge_saved_queue_{queue_id}_{category}",
1399 )
1400
1401 async def load_next_queue_item(
1402 self,
1403 queue_id: str,
1404 current_item_id: str,
1405 ) -> QueueItem:
1406 """
1407 Call when a player wants the next queue item to play.
1408
1409 Raises QueueEmpty if there are no more tracks left.
1410 """
1411 queue = self.get(queue_id)
1412 if not queue:
1413 msg = f"PlayerQueue {queue_id} is not available"
1414 raise PlayerUnavailableError(msg)
1415 cur_index = self.index_by_id(queue_id, current_item_id)
1416 if cur_index is None:
1417 # this is just a guard for bad data
1418 raise QueueEmpty("Invalid item id for queue given.")
1419 next_item: QueueItem | None = None
1420 idx = 0
1421 while True:
1422 next_index = self._get_next_index(queue_id, cur_index + idx)
1423 if next_index is None:
1424 raise QueueEmpty("No more tracks left in the queue.")
1425 queue_item = self.get_item(queue_id, next_index)
1426 if queue_item is None:
1427 raise QueueEmpty("No more tracks left in the queue.")
1428 if idx >= 10:
1429 # we only allow 10 retries to prevent infinite loops
1430 raise QueueEmpty("No more (playable) tracks left in the queue.")
1431 try:
1432 await self._load_item(queue_item, next_index)
1433 # we're all set, this is our next item
1434 next_item = queue_item
1435 break
1436 except ProviderStreamLimitError:
1437 # transient source capacity, do not burn a playable item over it
1438 raise
1439 except MediaNotFoundError, AudioError:
1440 # No stream details found, skip this QueueItem
1441 self.logger.warning(
1442 "Skipping unplayable item %s (%s)", queue_item.name, queue_item.uri
1443 )
1444 queue_item.available = False
1445 idx += 1
1446 if idx != 0:
1447 # we skipped some items, signal a queue items update
1448 self.update_items(queue_id, self._queue_data[queue_id].items)
1449 if next_item is None:
1450 raise QueueEmpty("No more (playable) tracks left in the queue.")
1451
1452 # carry playback_speed forward across consecutive audiobook/podcast items
1453 current_item = self.get_item(queue_id, current_item_id)
1454 if (
1455 current_item
1456 and current_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1457 and next_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1458 ):
1459 next_item.extra_attributes["playback_speed"] = current_item.extra_attributes.get(
1460 "playback_speed", 1.0
1461 )
1462
1463 return next_item
1464
1465 def track_loaded_in_buffer(self, queue_id: str, item_id: str) -> None:
1466 """Call when a player has (started) loading a track in the buffer."""
1467 queue = self.get(queue_id)
1468 if not queue:
1469 msg = f"PlayerQueue {queue_id} is not available"
1470 raise PlayerUnavailableError(msg)
1471 # store the index of the item that is currently (being) loaded in the buffer
1472 # which helps us a bit to determine how far the player has buffered ahead
1473 current_index = self.index_by_id(queue_id, item_id)
1474 queue.index_in_buffer = current_index
1475 self.logger.debug("PlayerQueue %s loaded item %s in buffer", queue.display_name, item_id)
1476 self.signal_update(queue_id)
1477 # preload next streamdetails
1478 self._preload_next_item(queue_id, item_id)
1479 # clean up stale audio buffers for old queue items to prevent memory leaks
1480 if current_index is not None:
1481 self.mass.create_task(self._cleanup_stale_queue_buffers(queue_id, current_index))
1482
1483 def queue_buffer_completed(self, queue_id: str, queue_exhausted: bool) -> None:
1484 """
1485 Call when the flow stream has finished generating all audio data for a queue.
1486
1487 At this point all audio data for the queue has been passed to the encoding pipeline.
1488 The player will go idle once it finishes playing the remaining buffered audio.
1489
1490 We start a background task that waits for the player to go idle and checks if new
1491 items have been added to the queue in the meantime, resuming playback if so.
1492
1493 :param queue_id: The queue ID.
1494 :param queue_exhausted: Whether the flow ended because the queue ran out of items,
1495 as opposed to ending early to restart on a format change or a live item.
1496 """
1497 queue = self.get(queue_id)
1498 if not queue:
1499 return
1500 self.logger.debug("Queue flow buffer completed for %s", queue.display_name)
1501
1502 # capture session_id so we can bail out if playback restarts
1503 queue_data = self._queue_data[queue_id]
1504 original_session_id = queue_data.session_id
1505 # record so player providers can detect flow EOF without an idle report
1506 if original_session_id is not None:
1507 queue_data.flow_buffer_completed = original_session_id
1508 if queue_exhausted:
1509 queue_data.flow_queue_exhausted = original_session_id
1510
1511 async def _resume_on_idle() -> None:
1512 # wait for the player to finish playing the buffered audio and go idle
1513 idle_detected = False
1514 for _ in range(60):
1515 await asyncio.sleep(1)
1516 if not queue.active or queue_data.session_id != original_session_id:
1517 return
1518 if queue.state == PlaybackState.IDLE:
1519 idle_detected = True
1520 break
1521 if not idle_detected:
1522 return
1523 # player went idle, give it a brief moment to settle
1524 await asyncio.sleep(1)
1525 if queue.state != PlaybackState.IDLE or queue_data.session_id != original_session_id:
1526 return
1527 # check if new items were added to the queue after the flow stream ended
1528 if queue.current_index is not None and (
1529 next_item := self.get_next_item(queue_id, queue.current_index)
1530 ):
1531 next_index = self.index_by_id(queue_id, next_item.queue_item_id)
1532 if next_index is not None:
1533 self.logger.info(
1534 "Resuming playback after flow stream completed for %s",
1535 queue.display_name,
1536 )
1537 await self.play_index(queue_id, next_index)
1538
1539 task_id = f"queue_buffer_completed_{queue_id}"
1540 self.mass.create_task(_resume_on_idle(), task_id=task_id)
1541
1542 def flow_stream_finished(self, queue_id: str) -> bool:
1543 """
1544 Return whether the flow stream for the current playback session is fully generated.
1545
1546 Lets player providers detect flow EOF when the device does not report idle
1547 (e.g. a Cast group that underruns the LIVE flow stream and keeps reporting playing).
1548
1549 :param queue_id: The queue ID.
1550 """
1551 queue_data = self.queue_data_or_none(queue_id)
1552 if queue_data is None or queue_data.session_id is None:
1553 return False
1554 return queue_data.flow_buffer_completed == queue_data.session_id
1555
1556 def flow_queue_exhausted(self, queue_id: str, session_id: str) -> bool:
1557 """
1558 Return whether the given flow stream session played the queue to its end.
1559
1560 False while a session is still streaming, and for a flow stream that ended early
1561 to be restarted (a format change or a live item), where the player is expected to
1562 pick up the next stream right away.
1563
1564 :param queue_id: The queue ID.
1565 :param session_id: The stream session to check.
1566 """
1567 queue_data = self.queue_data_or_none(queue_id)
1568 if queue_data is None or queue_data.session_id != session_id:
1569 return False
1570 return queue_data.flow_queue_exhausted == session_id
1571
1572 # Main queue manipulation methods
1573
1574 async def load(
1575 self,
1576 queue_id: str,
1577 queue_items: list[QueueItem],
1578 insert_at_index: int = 0,
1579 keep_remaining: bool = True,
1580 keep_played: bool = True,
1581 shuffle: bool = False,
1582 pin_first: bool = False,
1583 ) -> None:
1584 """
1585 Load new items at index.
1586
1587 - queue_id: id of the queue to process this request.
1588 - queue_items: a list of QueueItems
1589 - insert_at_index: insert the item(s) at this index
1590 - keep_remaining: keep the remaining items after the insert
1591 - shuffle: (re)shuffle the items after insert index
1592 - pin_first: keep the first item at the insert index instead of letting the shuffle
1593 move it; only meaningful together with shuffle
1594 """
1595 prev_items = self._queue_data[queue_id].items[:insert_at_index] if keep_played else []
1596 next_items = queue_items
1597
1598 # if keep_remaining, append the old 'next' items
1599 if keep_remaining:
1600 next_items += self._queue_data[queue_id].items[insert_at_index:]
1601
1602 # we set the original insert order as attribute so we can un-shuffle
1603 for index, item in enumerate(next_items):
1604 item.sort_index += insert_at_index + index
1605 # (re)shuffle the final batch if needed: smart shuffle when enabled, else pure random
1606 if shuffle:
1607 queue = self._queue_data[queue_id].queue
1608 # a user-picked item must stay the one that plays, so hold it out of the shuffle
1609 pinned = next_items[:1] if pin_first else []
1610 shuffled = next_items[1:] if pin_first else next_items
1611 if self._smart_shuffle.is_enabled(queue_id):
1612 shuffled = await self._smart_shuffle.arrange(queue, shuffled)
1613 else:
1614 shuffled = random.sample(shuffled, len(shuffled))
1615 next_items = pinned + shuffled
1616 self.update_items(queue_id, prev_items + next_items)
1617
1618 def update_items(self, queue_id: str, queue_items: list[QueueItem]) -> None:
1619 """Update the existing queue items, mostly caused by reordering."""
1620 self._queue_data[queue_id].items = queue_items
1621 queue = self._queue_data[queue_id].queue
1622 queue.items = len(self._queue_data[queue_id].items)
1623 self.signal_update(queue_id, True)
1624 if (
1625 queue.state == PlaybackState.PLAYING
1626 and queue.index_in_buffer is not None
1627 and queue.index_in_buffer == queue.current_index
1628 ):
1629 # if the queue is playing,
1630 # ensure to (re)queue the next track because it might have changed
1631 # note that we only do this if the player has loaded the current track
1632 # if not, we wait until it has loaded to prevent conflicts
1633 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
1634 self._enqueue_next_item(queue_id, next_item)
1635
1636 # Helper methods
1637
1638 def get_item(self, queue_id: str, item_id_or_index: int | str | None) -> QueueItem | None:
1639 """Get queue item by index or item_id."""
1640 if item_id_or_index is None:
1641 return None
1642 if (queue_data := self._queue_data.get(queue_id)) is None:
1643 return None
1644 queue_items = queue_data.items
1645 if isinstance(item_id_or_index, int) and len(queue_items) > item_id_or_index:
1646 return queue_items[item_id_or_index]
1647 if isinstance(item_id_or_index, str):
1648 return next((x for x in queue_items if x.queue_item_id == item_id_or_index), None)
1649 return None
1650
1651 def signal_update(self, queue_id: str, items_changed: bool = False) -> None:
1652 """Signal state changed of given queue."""
1653 if (queue_data := self._queue_data.get(queue_id)) is None:
1654 return
1655 queue = queue_data.queue
1656 # tell clients who owns the queue's ordering: the AudioSource uri while
1657 # delegated (see _get_delegated_source), None when MA owns it
1658 delegated = self._get_delegated_source(queue_id)
1659 queue.queue_owner = delegated[0].uri if delegated is not None else None
1660 # a mirrored shuffle write (streams controller) changes what smart shuffle
1661 # resolves to, so refresh the derived flag with every signaled update
1662 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1663 if items_changed:
1664 queue_data.items_cache_dirty = True
1665 self.mass.signal_event(EventType.QUEUE_ITEMS_UPDATED, object_id=queue_id, data=queue)
1666 self.mass.streams.audio_processing.prune(queue_id)
1667 # always send the base event
1668 self.mass.signal_event(EventType.QUEUE_UPDATED, object_id=queue_id, data=queue)
1669 # also signal update to the player itself so it can update its current_media
1670 self.mass.players.trigger_player_update(queue_id)
1671 # persist the (settings-bearing) queue state, debounced so a burst of updates or the
1672 # per-track updates during playback collapse into a single cache write
1673 self.mass.call_later(
1674 QUEUE_CACHE_SAVE_DELAY,
1675 self._save_queue_to_cache,
1676 queue_id,
1677 task_id=f"save_queue_cache_{queue_id}",
1678 )
1679
1680 def index_by_id(self, queue_id: str, queue_item_id: str) -> int | None:
1681 """Get index by queue_item_id."""
1682 if (queue_data := self._queue_data.get(queue_id)) is None:
1683 return None
1684 for index, item in enumerate(queue_data.items):
1685 if item.queue_item_id == queue_item_id:
1686 return index
1687 return None
1688
1689 async def get_tracks_for_playback(self, media_item: MediaItemType) -> list[Track]:
1690 """
1691 Return the playable tracks a media item resolves to, honoring the user's selection prefs.
1692
1693 :param media_item: The media item to resolve to playable tracks.
1694 """
1695 return await self._media_resolver.get_tracks_for_playback(media_item)
1696
1697 async def get_playlist_tracks(
1698 self, playlist: Playlist, start_item: str | None = None, sort_by: str | None = None
1699 ) -> list[PlaylistPlayableItem]:
1700 """
1701 Return the playable tracks for a playlist, honoring the user's selection prefs.
1702
1703 :param playlist: The playlist to resolve.
1704 :param start_item: Optional item URI to start the playlist from.
1705 :param sort_by: Optional sort key for the returned tracks.
1706 """
1707 return await self._media_resolver.get_playlist_tracks(playlist, start_item, sort_by)
1708
1709 async def get_dynamic_source_tracks(self, item: MediaItemType) -> list[Track]:
1710 """
1711 Return a fresh batch of tracks for a dynamic source (a dynamic playlist or radio station).
1712
1713 :param item: The dynamic playlist or radio station to fetch the next batch for.
1714 """
1715 return await self._media_resolver.get_dynamic_source_tracks(item)
1716
1717 def recency_windows(self) -> RecencyWindows:
1718 """Return the configured recency windows (a global setting; used for recency-aware gating)."""
1719 return self._smart_shuffle.windows()
1720
1721 async def player_media_from_queue_item(self, queue_item: QueueItem) -> PlayerMedia:
1722 """
1723 Parse PlayerMedia from QueueItem.
1724
1725 :param queue_item: The queue item to create media from.
1726 """
1727 queue_data = self._queue_data[queue_item.queue_id]
1728 stream_duration: int | None = None
1729 if queue_item.streamdetails:
1730 # prefer netto duration
1731 duration = queue_item.streamdetails.duration or queue_item.duration
1732 if duration and queue_item.streamdetails.seek_position:
1733 # the audio handed to the player starts at the seek position, so it is
1734 # shorter than the media item itself. seeking to (or past) the end
1735 # leaves no stream to describe, so the full length is kept instead.
1736 remaining = int(duration - queue_item.streamdetails.seek_position)
1737 stream_duration = remaining if remaining > 0 else None
1738 else:
1739 duration = queue_item.duration
1740 if queue_data.session_id is None:
1741 raise InvalidDataError("Queue session_id is None")
1742 media = PlayerMedia(
1743 uri=queue_item.uri,
1744 media_type=queue_item.media_type,
1745 title=queue_item.name,
1746 image_url=MASS_LOGO_ONLINE,
1747 duration=duration,
1748 stream_duration=stream_duration,
1749 source_id=queue_item.queue_id,
1750 queue_item_id=queue_item.queue_item_id,
1751 queue_session_id=queue_data.session_id,
1752 custom_data={
1753 "original_uri": queue_item.uri,
1754 },
1755 )
1756 if queue_item.media_item:
1757 media.title = queue_item.media_item.name
1758 media.artist = getattr(queue_item.media_item, "artist_str", "")
1759 media.album = (
1760 album.name if (album := getattr(queue_item.media_item, "album", None)) else ""
1761 )
1762 if queue_item.image:
1763 # the image format needs to be 512x512 jpeg for maximum compatibility with players
1764 # we prefer the imageproxy on the streamserver here because this request is sent
1765 # to the player itself which may not be able to reach the regular webserver
1766 media.image_url = self.mass.metadata.get_image_url(
1767 queue_item.image, size=512, image_format="jpeg", prefer_stream_server=True
1768 )
1769 return media
1770
1771 def get_next_item(self, queue_id: str, cur_index: int | str) -> QueueItem | None:
1772 """Return next QueueItem for given queue."""
1773 index: int
1774 if isinstance(cur_index, str):
1775 resolved_index = self.index_by_id(queue_id, cur_index)
1776 if resolved_index is None:
1777 return None # guard
1778 index = resolved_index
1779 else:
1780 index = cur_index
1781 # At this point index is guaranteed to be int
1782 for skip in range(5):
1783 if (next_index := self._get_next_index(queue_id, index + skip)) is None:
1784 break
1785 next_item = self.get_item(queue_id, next_index)
1786 if next_item is None:
1787 continue
1788 if not next_item.available:
1789 # ensure that we skip unavailable items (set by load_next track logic)
1790 continue
1791 return next_item
1792 return None
1793
1794 def store_sources(self, queue: PlayerQueue, items: list[MediaItemType]) -> None:
1795 """
1796 Hold the queue's full dynamic-source items server-side and project them onto `sources`.
1797
1798 :param queue: The queue whose sources are being set.
1799 :param items: The full source media items; an empty list clears the queue's sources.
1800 """
1801 self._queue_data[queue.queue_id].source_items = items
1802 # keep every occurrence server-side (a source added more than once weights it up in the
1803 # managed pool), but expose only the distinct container sources on the wire for clients to
1804 # show. Individual items (tracks, live radio streams, podcast episodes, ...) are omitted; see
1805 # `_WIRE_SOURCE_MEDIA_TYPES`. Autoplay/pool refill reads the full `source_items` above, not
1806 # this projected list, so it is unaffected.
1807 seen: set[str] = set()
1808 sources: list[ItemMapping] = []
1809 for item in items:
1810 if item.media_type not in _WIRE_SOURCE_MEDIA_TYPES and not is_dynamic_source(item):
1811 continue
1812 mapping = ItemMapping.from_item(item)
1813 if mapping.uri and mapping.uri in seen:
1814 continue
1815 if mapping.uri:
1816 seen.add(mapping.uri)
1817 sources.append(mapping)
1818 queue.sources = sources
1819 # release any materialized finite-source state whose source is no longer present
1820 self._managed_pool.retain(
1821 queue.queue_id, {item.uri for item in items if item.uri is not None}
1822 )
1823
1824 async def _save_queue_to_cache(self, queue_id: str) -> None:
1825 """Persist the queue's state (and its items when changed) to the cache."""
1826 if (queue_data := self._queue_data.get(queue_id)) is None:
1827 return
1828 try:
1829 # persistent so a cache clear/reset does not wipe the user's queues; the default
1830 # expiration still applies but is refreshed on every write. Skip the state write when its
1831 # persist-worthy content is unchanged (i.e. only playback progress advanced).
1832 state = queue_data.to_cache()
1833 significant = queue_data.cache_significant(state)
1834 if significant != queue_data.last_saved_state:
1835 await self.mass.cache.set(
1836 key=queue_id,
1837 data=state,
1838 provider=self.domain,
1839 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1840 persistent=True,
1841 )
1842 queue_data.last_saved_state = significant
1843 if queue_data.items_cache_dirty:
1844 # only cache items with a valid media_item
1845 await self.mass.cache.set(
1846 key=queue_id,
1847 data=queue_data.items_to_cache(),
1848 provider=self.domain,
1849 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1850 persistent=True,
1851 )
1852 queue_data.items_cache_dirty = False
1853 except Exception as err:
1854 self.logger.warning("Failed to persist the queue for %s - %s", queue_id, err)
1855
1856 def _check_player_permission(self, queue_id: str) -> None:
1857 """
1858 Check if the current user has permission to control this player/queue.
1859
1860 :param queue_id: The queue/player ID to check access for.
1861 :raises InsufficientPermissions: If the user lacks access.
1862 """
1863 current_user = get_current_user()
1864 if (
1865 current_user
1866 and current_user.player_filter
1867 and queue_id not in current_user.player_filter
1868 ):
1869 msg = f"{current_user.username} does not have access to player {queue_id}"
1870 raise InsufficientPermissions(msg)
1871
1872 @handle_play_action
1873 async def _handle_play(self, queue_id: str) -> None:
1874 """Handle play without acquiring the queue lock."""
1875 queue_player = self.mass.players.get_player(queue_id, True)
1876 if queue_player is None:
1877 raise PlayerUnavailableError(f"Player {queue_id} is not available")
1878 if (queue := self.get(queue_id)) and queue.active and queue.state == PlaybackState.PAUSED:
1879 # forward the actual play/unpause command to the player,
1880 # holding the action until the player confirms it resumed playback
1881 async with self.mass.players.wait_for_player_update(
1882 queue_id,
1883 attribute_name="playback_state",
1884 attribute_value=PlaybackState.PLAYING,
1885 timeout=PLAYBACK_START_TIMEOUT,
1886 ):
1887 await queue_player.play()
1888 return
1889 # player is not paused, perform resume instead
1890 await self.resume(queue_id)
1891
1892 def _set_transitioning(self, queue_id: str, value: bool) -> None:
1893 """Mark (or clear) whether a queue is mid-transition (no-op if it is not registered)."""
1894 if (queue_data := self._queue_data.get(queue_id)) is not None:
1895 queue_data.transitioning = value
1896
1897 def _clear(self, queue_id: str, skip_stop: bool = False) -> None:
1898 """Drop the queue's items and playback position, leaving user settings untouched."""
1899 queue = self._queue_data[queue_id].queue
1900 self.mass.streams.audio_processing.clear(queue_id)
1901 self.store_sources(queue, [])
1902 if queue.is_dynamic:
1903 # Dynamic sources impose shuffle, so clearing the source clears that shuffle too.
1904 queue.shuffle_enabled = False
1905 queue.is_dynamic = False
1906 # dropping the dynamic source changes what smart shuffle resolves to, so the derived
1907 # flag has to follow or clients keep showing a smart mix on a plain queue
1908 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1909 queue.ended = False
1910 if queue.state != PlaybackState.IDLE and not skip_stop:
1911 self.mass.create_task(self.stop(queue_id))
1912 queue.current_index = None
1913 queue.current_item = None
1914 queue.elapsed_time = 0
1915 queue.elapsed_time_last_updated = time.time()
1916 queue.index_in_buffer = None
1917 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
1918 self.update_items(queue_id, [])
1919
1920 def _audio_source_plugin(
1921 self, queue_item: QueueItem | None
1922 ) -> tuple[PluginProvider, str] | None:
1923 """Return the plugin owning this item's live AudioSource, with the source id it exposes."""
1924 if queue_item is None or (media_item := queue_item.media_item) is None:
1925 return None
1926 if media_item.media_type != MediaType.AUDIO_SOURCE:
1927 return None
1928 if not isinstance(prov := self.mass.get_provider(media_item.provider), PluginProvider):
1929 return None
1930 return prov, media_item.item_id
1931
1932 def _audio_sources_in(self, queue_id: str) -> dict[str, tuple[PluginProvider, str]]:
1933 """
1934 Map every live AudioSource the queue holds to its owning plugin, keyed by media uri.
1935
1936 Covers the whole queue rather than just what is playing: an option that starts media
1937 alongside a live source leaves that source behind as an ordinary item, and it still has
1938 to be released when it eventually goes.
1939
1940 :param queue_id: The queue to look through.
1941 :return: The owning plugin and the source id it exposes, per AudioSource media uri.
1942 """
1943 owners: dict[str, tuple[PluginProvider, str]] = {}
1944 if (queue_data := self._queue_data.get(queue_id)) is None:
1945 return owners
1946 for item in queue_data.items:
1947 media_item = item.media_item
1948 if media_item is None or media_item.media_type != MediaType.AUDIO_SOURCE:
1949 continue
1950 if (uri := media_item.uri) is None or uri in owners:
1951 continue
1952 if (owner := self._audio_source_plugin(item)) is not None:
1953 owners[uri] = owner
1954 return owners
1955
1956 def _notify_audio_source_removed(self, queue_id: str) -> None:
1957 """Tell the owning plugins that their AudioSources are dropped along with this queue."""
1958 for prov, source_id in self._audio_sources_in(queue_id).values():
1959 self.mass.create_task(prov.on_source_removed(source_id, queue_id))
1960
1961 def _notify_audio_source_replaced(
1962 self, queue_id: str, outgoing: dict[str, tuple[PluginProvider, str]]
1963 ) -> None:
1964 """
1965 Tell the owning plugins which of their AudioSources the queue's new contents pushed out.
1966
1967 A source that is still among the queue's items is left in place, so the same source
1968 coming back and contents that keep it alongside them both release nothing.
1969
1970 :param queue_id: The queue whose contents were taken over.
1971 :param outgoing: The queue's live AudioSources from before the takeover.
1972 """
1973 if not outgoing or (queue_data := self._queue_data.get(queue_id)) is None:
1974 return
1975 remaining = {
1976 item.media_item.uri for item in queue_data.items if item.media_item is not None
1977 }
1978 for uri, (prov, source_id) in outgoing.items():
1979 if uri not in remaining:
1980 self.mass.create_task(prov.on_source_removed(source_id, queue_id))
1981
1982 async def _notify_audio_source_transferred(
1983 self,
1984 transferred: dict[str, tuple[PluginProvider, str]],
1985 from_queue_id: str,
1986 to_queue_id: str,
1987 ) -> None:
1988 """
1989 Tell the owning plugins that their AudioSources moved to another queue.
1990
1991 Awaited rather than dispatched so the plugins have settled the handover before the target
1992 queue is resumed. A plugin that raises must not break the transfer.
1993
1994 :param transferred: The live AudioSources the source queue handed over.
1995 :param from_queue_id: The queue that gave the sources up.
1996 :param to_queue_id: The queue that took them over.
1997 """
1998 if not transferred or (queue_data := self._queue_data.get(to_queue_id)) is None:
1999 return
2000 arrived = {item.media_item.uri for item in queue_data.items if item.media_item is not None}
2001 for uri, (prov, source_id) in transferred.items():
2002 if uri not in arrived:
2003 continue
2004 try:
2005 await prov.on_source_transferred(source_id, from_queue_id, to_queue_id)
2006 except Exception:
2007 self.logger.warning(
2008 "on_source_transferred raised for provider %s source %s",
2009 prov.instance_id,
2010 source_id,
2011 exc_info=True,
2012 )
2013
2014 def _reset_shuffle(self, queue_id: str) -> None:
2015 """Switch shuffle off."""
2016 queue = self._queue_data[queue_id].queue
2017 if not queue.shuffle_enabled:
2018 return
2019 queue.shuffle_enabled = False
2020 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
2021 self.signal_update(queue_id)
2022
2023 async def _apply_shuffle(
2024 self, queue_id: str, option: QueueOption, shuffle: bool | None
2025 ) -> None:
2026 """
2027 Settle the queue's shuffle state for a play command before its items are resolved.
2028
2029 :param queue_id: The queue the media is played on.
2030 :param option: The enqueue option this command resolved to.
2031 :param shuffle: The state to put the queue's shuffle in; None to leave it as it is.
2032 """
2033 queue = self._queue_data[queue_id].queue
2034 if queue.is_dynamic and option in (
2035 QueueOption.PLAY,
2036 QueueOption.REPLACE,
2037 QueueOption.REPLACE_NEXT,
2038 ):
2039 # These are the options that replace the queue's sources, so the smart mix may be on
2040 # its way out - and its shuffle is never the user's own (a dynamic queue's toggle is
2041 # locked), so it must not outlive the source that imposed it. Recorded directly
2042 # because set_shuffle refuses a queue that is still a smart mix, and the items are
2043 # resolved against this flag. The state is provisional until the sources are known:
2044 # `_enter_dynamic_mode` forces shuffle back on if the queue stays dynamic.
2045 if option == QueueOption.REPLACE_NEXT:
2046 # staging leaves the shuffle the user chose alone, so it never carries a request
2047 # of its own to honour here
2048 queue.shuffle_enabled = False
2049 else:
2050 queue.shuffle_enabled = bool(shuffle)
2051 return
2052 if shuffle is None or option not in (QueueOption.PLAY, QueueOption.REPLACE):
2053 # nothing to settle: the media brings no order of its own to protect, or the option
2054 # only stages items for later and leaves the queue's shuffle state alone
2055 return
2056 if queue.shuffle_enabled == shuffle:
2057 return
2058 if self._get_delegated_source(queue_id) is not None:
2059 # the queue is still delegated to the external session this play is replacing,
2060 # so set_shuffle would forward the toggle to that session; apply the state and
2061 # the tail re-order locally instead â MA-owned items behind the session may
2062 # sit in shuffled order and must not survive an ordered play unshuffled-in-name-only
2063 await self._apply_local_shuffle(queue_id, shuffle)
2064 return
2065 # routed through set_shuffle so switching shuffle off also restores the order of
2066 # the items that stay in the queue: a play keeps them, and a tail left in shuffled
2067 # order behind a queue that now reads unshuffled would contradict its own flag
2068 await self.set_shuffle(queue_id, shuffle)
2069
2070 async def _apply_mirrored_shuffle(
2071 self, queue_id: str, source_id: str, provider_instance: str, shuffle_enabled: bool
2072 ) -> None:
2073 """
2074 Apply a session-mirrored shuffle change, serialized against playback commands.
2075
2076 Scheduled as a task by the streams controller; the delegation and source
2077 identity are re-validated under the player lock so a session that ended (or
2078 changed) between the options event and this task running cannot re-order an
2079 unrelated queue.
2080
2081 :param queue_id: The queue the session mirror applies to.
2082 :param source_id: The AudioSource.item_id that reported the change.
2083 :param provider_instance: The provider instance id that reported the change.
2084 :param shuffle_enabled: The session's shuffle state to mirror.
2085 """
2086 async with self.mass.players.get_player_lock(queue_id):
2087 delegated = self._get_delegated_source(queue_id)
2088 if delegated is None:
2089 return
2090 audio_source = delegated[0]
2091 if audio_source.item_id != source_id or audio_source.provider != provider_instance:
2092 return
2093 if self._queue_data[queue_id].queue.shuffle_enabled == shuffle_enabled:
2094 return
2095 await self._apply_local_shuffle(queue_id, shuffle_enabled)
2096
2097 async def _apply_local_shuffle(self, queue_id: str, shuffle_enabled: bool) -> None:
2098 """
2099 Record the queue's shuffle state and re-order the un-played tail accordingly.
2100
2101 :param queue_id: The queue to apply the shuffle state to.
2102 :param shuffle_enabled: The shuffle state to record and apply to the tail.
2103 """
2104 queue = self._queue_data[queue_id].queue
2105 queue.shuffle_enabled = shuffle_enabled
2106 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
2107 queue_items = self._queue_data[queue_id].items
2108 cur_index = (
2109 queue.index_in_buffer if queue.index_in_buffer is not None else queue.current_index
2110 )
2111 if cur_index is not None:
2112 next_index = cur_index + 1
2113 next_items = queue_items[next_index:]
2114 else:
2115 next_items = []
2116 next_index = 0
2117 if not shuffle_enabled:
2118 # shuffle disabled, try to restore original sort order of the remaining items
2119 next_items.sort(key=lambda x: x.sort_index, reverse=False)
2120 await self.load(
2121 queue_id=queue_id,
2122 queue_items=next_items,
2123 insert_at_index=next_index,
2124 keep_remaining=False,
2125 shuffle=shuffle_enabled,
2126 )
2127
2128 def _get_delegated_source(
2129 self, queue_id: str
2130 ) -> tuple[AudioSource, SourceQueueCapabilities, PluginProvider] | None:
2131 """
2132 Return the AudioSource owning the queue's commands, its capabilities and owning plugin.
2133
2134 While the queue's current item is an AudioSource declaring ``queue_capabilities``,
2135 the external session owns the queue: shuffle/repeat are forwarded to the owning
2136 plugin and the mirrored options event updates the queue state afterwards. There is
2137 no playback-state gate, because a plugin may map a paused session onto a stopped
2138 player (Spotify Connect does) â a genuinely dead session is the owning plugin's
2139 call, surfaced as its localized not-active error. Returns None â queue commands
2140 then apply to the MA queue as usual â when the current item is not such an
2141 AudioSource (a transport-only source keeps ownership of the queue with MA) or
2142 when the owning plugin provider is no longer available. Transport commands
2143 (next/previous/seek) do not use this gate: see ``_get_current_audio_source``.
2144
2145 :param queue_id: The queue to inspect.
2146 """
2147 if (resolved := self._get_current_audio_source(queue_id)) is None:
2148 return None
2149 audio_source, provider = resolved
2150 if (caps := audio_source.queue_capabilities) is None:
2151 return None
2152 return audio_source, caps, provider
2153
2154 def _get_current_audio_source(self, queue_id: str) -> tuple[AudioSource, PluginProvider] | None:
2155 """
2156 Return the AudioSource current on the queue and its owning PluginProvider.
2157
2158 Unlike ``_get_delegated_source`` this does not require ``queue_capabilities``:
2159 transport commands (next/previous/seek) delegate on the per-action capability
2160 flags alone, so a transport-only source skips/seeks within its own session
2161 via the queue API just like it does via the player-command API.
2162
2163 :param queue_id: The queue to inspect.
2164 """
2165 if (queue_data := self._queue_data.get(queue_id)) is None:
2166 return None
2167 return get_queue_audio_source(self.mass, queue_data.queue)
2168