/
/
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 ProviderFeature,
29 QueueOption,
30 RepeatMode,
31 SourceControl,
32)
33from music_assistant_models.errors import (
34 AudioError,
35 InsufficientPermissions,
36 InvalidCommand,
37 InvalidDataError,
38 MediaNotFoundError,
39 PlayerUnavailableError,
40 QueueEmpty,
41)
42from music_assistant_models.media_items import (
43 Audiobook,
44 AudioSource,
45 ItemMapping,
46 MediaItemType,
47 PlayableMediaItemType,
48 Playlist,
49 PodcastEpisode,
50 SoundEffect,
51 Track,
52)
53from music_assistant_models.player_queue import PlayerQueue
54
55from music_assistant.constants import (
56 ATTR_ANNOUNCEMENT_IN_PROGRESS,
57 MASS_LOGO_ONLINE,
58 PLAYLIST_MEDIA_TYPES,
59)
60from music_assistant.controllers.player_queues.autoplay import Autoplay
61from music_assistant.controllers.player_queues.config import (
62 core_config_entries,
63 queue_config_entries,
64)
65from music_assistant.controllers.player_queues.constants import (
66 CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
67 CACHE_CATEGORY_PLAYER_QUEUE_STATE,
68 PLAYBACK_START_TIMEOUT,
69 QUEUE_CACHE_SAVE_DELAY,
70)
71from music_assistant.controllers.player_queues.helpers import (
72 get_current_playback_speed,
73 handle_play_action,
74 is_dynamic_source,
75)
76from music_assistant.controllers.player_queues.managed_pool import ManagedPool
77from music_assistant.controllers.player_queues.media_resolver import MediaResolver
78from music_assistant.controllers.player_queues.playback_tracker import PlaybackTrackerMixin
79from music_assistant.controllers.player_queues.queue_loader import QueueLoaderMixin
80from music_assistant.controllers.player_queues.smart_shuffle import SmartShuffle
81from music_assistant.controllers.player_queues.state import PlayerQueueData
82from music_assistant.controllers.player_queues.stream_feeder import StreamFeederMixin
83from music_assistant.controllers.webserver.helpers.auth_middleware import get_current_user
84from music_assistant.helpers.api import api_command
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 (delegated := self._get_delegated_source(queue_id)) is not None:
802 audio_source, _caps, provider = delegated
803 if not audio_source.can_next_previous:
804 raise InvalidCommand("Cannot skip: the external session does not support skipping")
805 await provider.on_source_control(audio_source.item_id, SourceControl.NEXT)
806 return
807 self._set_transitioning(queue_id, True)
808 idx = self._queue_data[queue_id].queue.current_index
809 if idx is None:
810 self.logger.warning("Queue %s has no current index", queue.display_name)
811 self._set_transitioning(queue_id, False)
812 return
813 next_index = self._get_next_index(queue_id, idx, True)
814 if next_index is None:
815 self._set_transitioning(queue_id, False)
816 return
817
818 # immediately update current item so UI shows the new track right away
819 queue.current_index = next_index
820 queue.current_item = self.get_item(queue_id, next_index)
821 queue.elapsed_time = 0
822 queue.elapsed_time_last_updated = time.time()
823 self.signal_update(queue_id)
824 if queue_player := self.mass.players.get_player(queue_id, True):
825 queue_player.update_state()
826
827 # debounce rapid next button presses using call_later
828 self.mass.call_later(
829 1,
830 self.play_index,
831 queue_id,
832 next_index,
833 task_id=f"queue_play_index_{queue_id}",
834 )
835
836 @api_command("player_queues/previous", required_scope=Scope.QUEUES_CONTROL)
837 @handle_play_action
838 async def previous(self, queue_id: str) -> None:
839 """
840 Handle PREVIOUS TRACK command for given queue.
841
842 :param queue_id: queue_id of the queue to handle the command.
843 """
844 self._check_player_permission(queue_id)
845 if (queue := self.get(queue_id)) is None or not queue.active:
846 raise InvalidCommand(f"Queue {queue_id} is not active")
847 if (delegated := self._get_delegated_source(queue_id)) is not None:
848 audio_source, _caps, provider = delegated
849 if not audio_source.can_next_previous:
850 raise InvalidCommand("Cannot skip: the external session does not support skipping")
851 await provider.on_source_control(audio_source.item_id, SourceControl.PREVIOUS)
852 return
853 self._set_transitioning(queue_id, True)
854 current_index = self._queue_data[queue_id].queue.current_index
855 if current_index is None:
856 self._set_transitioning(queue_id, False)
857 return
858 prev_index = int(current_index)
859 # restart current track if elapsed > 5s, otherwise go to previous
860 if self._queue_data[queue_id].queue.elapsed_time < 5:
861 prev_index = max(current_index - 1, 0)
862
863 # immediately update current item so UI shows the new track right away
864 queue.current_index = prev_index
865 queue.current_item = self.get_item(queue_id, prev_index)
866 queue.elapsed_time = 0
867 queue.elapsed_time_last_updated = time.time()
868 self.signal_update(queue_id)
869 if queue_player := self.mass.players.get_player(queue_id, True):
870 queue_player.update_state()
871
872 # debounce rapid previous button presses using call_later
873 self.mass.call_later(
874 1,
875 self.play_index,
876 queue_id,
877 prev_index,
878 task_id=f"queue_play_index_{queue_id}",
879 )
880
881 @api_command("player_queues/skip", required_scope=Scope.QUEUES_CONTROL)
882 async def skip(self, queue_id: str, seconds: int = 10) -> None:
883 """
884 Handle SKIP command for given queue.
885
886 - queue_id: queue_id of the queue to handle the command.
887 - seconds: number of seconds to skip in track. Use negative value to skip back.
888 """
889 if (queue := self.get(queue_id)) is None or not queue.active:
890 raise InvalidCommand(f"Queue {queue_id} is not active")
891 await self.seek(queue_id, int(self._queue_data[queue_id].queue.elapsed_time + seconds))
892
893 @api_command("player_queues/seek", required_scope=Scope.QUEUES_CONTROL)
894 async def seek(self, queue_id: str, position: int = 10) -> None:
895 """
896 Handle SEEK command for given queue.
897
898 - queue_id: queue_id of the queue to handle the command.
899 - position: position in seconds to seek to in the current playing item.
900 """
901 if (queue := self.get(queue_id)) is None or not queue.active:
902 raise InvalidCommand(f"Queue {queue_id} is not active")
903 if (delegated := self._get_delegated_source(queue_id)) is not None:
904 audio_source, _caps, provider = delegated
905 if not audio_source.can_seek:
906 raise InvalidCommand("Cannot seek: the external session does not support seeking")
907 position = max(0, int(position))
908 current_item = queue.current_item
909 await provider.on_source_control(audio_source.item_id, SourceControl.SEEK, position)
910 # publish the seek target so the progress bar does not snap back to the
911 # last position report (mirrors the non-delegated path below) â unless a
912 # concurrent play/stop replaced the current item during the forward
913 if queue.current_item is current_item:
914 queue.elapsed_time = position
915 queue.elapsed_time_last_updated = time.time()
916 self.signal_update(queue_id)
917 return
918 queue_player = self.mass.players.get_player(queue_id, True)
919 if queue_player is None:
920 raise PlayerUnavailableError(f"Player {queue_id} is not available")
921 if not queue.current_item:
922 raise InvalidCommand(f"Queue {queue_player.state.name} has no item(s) loaded.")
923 if not queue.current_item.duration:
924 raise InvalidCommand("Can not seek items without duration.")
925 position = max(0, int(position))
926 if position > queue.current_item.duration:
927 raise InvalidCommand("Can not seek outside of duration range.")
928 if queue.current_index is None:
929 raise InvalidCommand(f"Queue {queue_player.state.name} has no current index.")
930 # Publish the seek target before rebuilding the stream to prevent progress snapback.
931 queue.elapsed_time = position
932 queue.elapsed_time_last_updated = time.time()
933 self.signal_update(queue_id)
934 await self.play_index(queue_id, queue.current_index, seek_position=position)
935
936 @api_command("player_queues/resume", required_scope=Scope.QUEUES_CONTROL)
937 @handle_play_action
938 async def resume(self, queue_id: str, fade_in: bool | None = None) -> None:
939 """
940 Handle RESUME command for given queue.
941
942 - queue_id: queue_id of the queue to handle the command.
943 """
944 self._check_player_permission(queue_id)
945 queue = self._queue_data[queue_id].queue
946 queue_items = self._queue_data[queue_id].items
947 resume_item = queue.current_item
948 if queue.state == PlaybackState.PLAYING:
949 # resume requested while already playing,
950 # use current position as resume position
951 resume_pos = queue.corrected_elapsed_time
952 fade_in = False
953 else:
954 resume_pos = queue.resume_pos or queue.elapsed_time
955
956 if queue.ended and len(queue_items) > 0:
957 # the queue played to its end and is parked on its last item,
958 # so pressing play starts it over from the beginning
959 resume_item = queue_items[0]
960 resume_pos = 0
961 elif not resume_item and queue.current_index is not None and len(queue_items) > 0:
962 resume_item = self.get_item(queue_id, queue.current_index)
963 resume_pos = 0
964 elif not resume_item and queue.current_index is None and len(queue_items) > 0:
965 # items available in queue but no previous track, start at 0
966 resume_item = self.get_item(queue_id, 0)
967 resume_pos = 0
968
969 if resume_item is not None:
970 queue_player = self.mass.players.get_player(queue_id)
971 if queue_player is None:
972 raise PlayerUnavailableError(f"Player {queue_id} is not available")
973 if (
974 fade_in is None
975 and queue_player.state.playback_state == PlaybackState.IDLE
976 and (time.time() - queue.elapsed_time_last_updated) > 60
977 ):
978 # enable fade in effect if the player is idle for a while
979 fade_in = resume_pos > 0
980 if resume_item.media_type == MediaType.RADIO:
981 # we're not able to skip in online radio so this is pointless
982 resume_pos = 0
983 await self.play_index(
984 queue_id, resume_item.queue_item_id, int(resume_pos), fade_in or False
985 )
986 else:
987 msg = f"Resume queue requested but queue {queue.display_name} is empty"
988 raise QueueEmpty(msg)
989
990 @api_command("player_queues/play_index", required_scope=Scope.QUEUES_CONTROL)
991 @handle_play_action
992 async def play_index( # noqa: PLR0915
993 self,
994 queue_id: str,
995 index: int | str,
996 seek_position: int = 0,
997 fade_in: bool = False,
998 ) -> None:
999 """Play item at index (or item_id) X in queue."""
1000 self._check_player_permission(queue_id)
1001 # cancel any pending play_index calls for this queue to prevent conflicts
1002 self.mass.cancel_timer(f"queue_play_index_{queue_id}")
1003 # we set a flag to notify the update logic that we're transitioning to a new track
1004 self._set_transitioning(queue_id, True)
1005 try:
1006 queue_data = self._queue_data[queue_id]
1007 queue = queue_data.queue
1008 queue.resume_pos = 0
1009 # A queue picked up from its end plays its items over from the start, so a resume point
1010 # left on an audiobook/episode must not pull it back to where it was left off. The flag
1011 # itself is only cleared once an item actually loaded below, so a start that never got
1012 # off the ground leaves the queue finished instead of stranding it without a position.
1013 restarting_ended_queue = queue.ended
1014 if isinstance(index, str):
1015 temp_index = self.index_by_id(queue_id, index)
1016 if temp_index is None:
1017 raise InvalidDataError(f"Item {index} not found in queue")
1018 index = temp_index
1019 # At this point index is guaranteed to be int
1020 queue.index_in_buffer = index
1021 queue_data.flow_mode_stream_log = []
1022 queue_data.flow_buffer_completed = None
1023 queue_data.flow_queue_exhausted = None
1024 target_player = self.mass.players.get_player(queue_id)
1025 if target_player is None:
1026 raise PlayerUnavailableError(f"Player {queue_id} is not available")
1027 queue_data.next_item_id_enqueued = None
1028 # always update session id when we start a new playback session
1029 queue_data.session_id = shortuuid.random(length=8)
1030 self.mass.streams.audio_processing.start_session(
1031 queue_id,
1032 queue_data.session_id,
1033 )
1034 # handle resume point of audiobook(chapter) or podcast(episode)
1035 if (
1036 not seek_position
1037 and not restarting_ended_queue
1038 and (queue_item := self.get_item(queue_id, index))
1039 and (resume_position_ms := getattr(queue_item.media_item, "resume_position_ms", 0))
1040 ):
1041 # the client may have fetched the item before its duration was known
1042 await self._restore_probed_duration(queue_item)
1043 if queue_item.duration or getattr(queue_item.media_item, "duration", 0):
1044 seek_position = max(0, int((resume_position_ms - 500) / 1000))
1045 else:
1046 # seeking needs a duration, which is determined while streaming
1047 self.logger.debug(
1048 "Can not resume %s at %ss: its duration is not known (yet)",
1049 queue_item.name,
1050 int(resume_position_ms / 1000),
1051 )
1052
1053 # restore the persisted playback speed for a freshly queued audiobook/episode
1054 # (an in-session item already carries its speed in extra_attributes)
1055 if (
1056 (queue_item := self.get_item(queue_id, index))
1057 and queue_item.media_item is not None
1058 and queue_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1059 and "playback_speed" not in queue_item.extra_attributes
1060 ):
1061 stored_speed = await self.mass.music.get_playback_speed(
1062 cast("Audiobook | PodcastEpisode", queue_item.media_item),
1063 userid=queue_data.userid,
1064 )
1065 if stored_speed != 1.0:
1066 queue_item.extra_attributes["playback_speed"] = stored_speed
1067
1068 # try to load the item, retry with next item if it fails
1069 for attempt in range(5):
1070 try:
1071 queue_item = self.get_item(queue_id, index)
1072 if not queue_item:
1073 continue # guard
1074 await self._load_item(
1075 queue_item,
1076 self._get_next_index(queue_id, index),
1077 is_start=True,
1078 seek_position=seek_position if attempt == 0 else 0,
1079 fade_in=fade_in if attempt == 0 else False,
1080 )
1081 # if we reach this point, loading the item succeeded, break the loop
1082 queue.current_index = index
1083 queue.current_item = queue_item
1084 # playback is under way, so the queue is no longer sitting at its end
1085 queue.ended = False
1086 # reset the elapsed clock together with the item switch (like
1087 # next/previous do), so queue updates signaled before the player
1088 # reports position don't carry the previous item's elapsed_time
1089 queue.elapsed_time = seek_position if attempt == 0 else 0
1090 queue.elapsed_time_last_updated = time.time()
1091 break
1092 except (MediaNotFoundError, AudioError) as err:
1093 item_name = queue_item.name if queue_item else "unknown"
1094 if isinstance(err, ProviderStreamLimitError):
1095 # the requested item is playable, its provider is just at capacity:
1096 # report that instead of silently advancing to another item
1097 self.logger.error("%s", err)
1098 await self.stop(queue_id)
1099 raise
1100 # Only MediaNotFoundError (item unreachable) is persistent;
1101 # keep AudioError items available so a retry can resurface
1102 # the same actionable error.
1103 if queue_item and isinstance(err, MediaNotFoundError):
1104 queue_item.available = False
1105 next_index = self._get_next_index(queue_id, index, allow_repeat=False)
1106 if next_index is None:
1107 # Surface an AudioError's own (actionable) message;
1108 # MediaNotFoundError gets the generic wording.
1109 if isinstance(err, AudioError) and str(err):
1110 msg = str(err)
1111 else:
1112 msg = f"Playback failed for {item_name} - no more tracks available"
1113 self.logger.error(msg)
1114 await self.stop(queue_id)
1115 raise MediaNotFoundError(msg) from err
1116 self.logger.warning(
1117 "Skipping unplayable item %s",
1118 item_name,
1119 )
1120 index = next_index
1121 else:
1122 # all attempts to find a playable item failed
1123 await self.stop(queue_id)
1124 raise MediaNotFoundError("No playable item found to start playback")
1125
1126 # Reset flow_mode - the streams controller will set it if flow mode is used.
1127 queue.flow_mode = False
1128 player_media = await self.player_media_from_queue_item(queue_item)
1129 # Hold the play action until the player confirms playback so the UI keeps
1130 # showing the command as in progress instead of falling back to a play button
1131 # for the time the player still needs to connect and start. The queue update
1132 # for the new item goes out first, so the item shows while it is starting.
1133 async with self.mass.players.wait_for_player_update(
1134 queue_id,
1135 attribute_name="playback_state",
1136 attribute_value=PlaybackState.PLAYING,
1137 timeout=PLAYBACK_START_TIMEOUT,
1138 ):
1139 await self.mass.players.play_media(queue_id, player_media)
1140 queue.current_index = index
1141 queue.current_item = queue_item
1142 self.signal_update(queue_id)
1143 finally:
1144 self._set_transitioning(queue_id, False)
1145
1146 @api_command("player_queues/transfer", required_scope=Scope.QUEUES_CONTROL)
1147 async def transfer_queue(
1148 self,
1149 source_queue_id: str,
1150 target_queue_id: str,
1151 auto_play: bool | None = None,
1152 ) -> None:
1153 """Transfer queue to another queue."""
1154 if not (source_queue := self.get(source_queue_id)):
1155 raise PlayerUnavailableError(f"Queue {source_queue_id} is not available")
1156 if not (target_queue := self.get(target_queue_id)):
1157 raise PlayerUnavailableError(f"Queue {target_queue_id} is not available")
1158 if auto_play is None:
1159 auto_play = source_queue.state == PlaybackState.PLAYING
1160
1161 target_player = self.mass.players.get_player(target_queue_id)
1162 if target_player is None:
1163 raise PlayerUnavailableError(f"Player {target_queue_id} is not available")
1164 if target_player.state.active_group or target_player.state.synced_to:
1165 # edge case: the user wants to move playback from the group as a whole, to a single
1166 # player in the group or it is grouped and the command targeted at the single player.
1167 # We need to dissolve the group/sync first, and wait for the state to actually
1168 # propagate before we hand the queue over to the target player.
1169 group_id = target_player.state.active_group or target_player.state.synced_to
1170 assert group_id is not None # checked in if condition above
1171 # For an ad-hoc sync group (target is a sync member of a regular leader),
1172 # ungroup the target itself so only it is freed - ungrouping the leader would
1173 # transfer leadership to a remaining member and recurse back into this method.
1174 # For a virtual group player (active_group), release the group so its static
1175 # members are handled correctly.
1176 ungroup_target = (
1177 target_queue_id
1178 if target_player.state.synced_to and not target_player.state.active_group
1179 else group_id
1180 )
1181 async with self.mass.players.wait_for_player_update(
1182 target_queue_id,
1183 attribute_name=(
1184 "active_group" if target_player.state.active_group else "synced_to"
1185 ),
1186 attribute_value=None,
1187 timeout=5,
1188 ):
1189 await self.mass.players.cmd_ungroup(ungroup_target)
1190
1191 # capture source state before stopping (stop resets these)
1192 source_items = self._queue_data[source_queue_id].items
1193 if source_queue.state == PlaybackState.PLAYING:
1194 # use the live playback clock while actively playing
1195 source_resume_pos = int(source_queue.corrected_elapsed_time)
1196 else:
1197 # when not playing the live clock is stale, so use the stored resume position
1198 source_resume_pos = int(source_queue.resume_pos or source_queue.elapsed_time or 0)
1199 source_current_index = source_queue.current_index
1200 source_current_item = source_queue.current_item
1201
1202 # stop the source player synchronously to prevent the async stop from
1203 # clear() racing with the target's sync group formation/protocol switching
1204 if source_queue.state != PlaybackState.IDLE:
1205 await self.stop(source_queue_id)
1206
1207 target_queue.repeat_mode = source_queue.repeat_mode
1208 target_queue.shuffle_enabled = source_queue.shuffle_enabled
1209 target_queue.crossfade_enabled = source_queue.crossfade_enabled
1210 # refresh the derived smart-fades indicator for the target's own config/availability
1211 target_queue.smart_fades_active = self.mass.streams.is_smart_fades_active(target_queue)
1212 target_queue.autoplay_enabled = source_queue.autoplay_enabled
1213 self._queue_data[target_queue_id].source_items = list(
1214 self._queue_data[source_queue_id].source_items
1215 )
1216 target_queue.sources = list(source_queue.sources)
1217 target_queue.is_dynamic = source_queue.is_dynamic
1218 target_queue.smart_shuffle_active = self.is_smart_shuffle_active(target_queue)
1219 self._queue_data[target_queue_id].enqueued_media_items = list(
1220 self._queue_data[source_queue_id].enqueued_media_items
1221 )
1222 target_queue.resume_pos = source_resume_pos
1223 # the target's own contents are about to be overwritten by the transferred ones
1224 target_outgoing_sources = self._audio_sources_in(target_queue_id)
1225 target_queue.current_index = source_current_index
1226 if source_current_item:
1227 target_queue.current_item = source_current_item
1228 target_queue.current_item.queue_id = target_queue_id
1229 # every live source the queue holds moves with it, not just the one that was playing
1230 transferred_sources = self._audio_sources_in(source_queue_id)
1231 self._clear(source_queue_id, skip_stop=True)
1232
1233 await self.load(target_queue_id, source_items, keep_remaining=False, keep_played=False)
1234 for item in source_items:
1235 item.queue_id = target_queue_id
1236 self.update_items(target_queue_id, source_items)
1237 # a live source the target was holding has just been displaced by the transferred queue
1238 self._notify_audio_source_replaced(target_queue_id, target_outgoing_sources)
1239 await self._notify_audio_source_transferred(
1240 transferred_sources, source_queue_id, target_queue_id
1241 )
1242 if auto_play:
1243 await self.resume(target_queue_id)
1244
1245 # Interaction with player
1246
1247 async def on_player_register(self, player: Player) -> None:
1248 """Register PlayerQueue for given player/queue id."""
1249 queue_id = player.player_id
1250 queue_data: PlayerQueueData | None = None
1251 # try to restore previous state
1252 try:
1253 if prev_state := await self.mass.cache.get(
1254 key=queue_id,
1255 provider=self.domain,
1256 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1257 ):
1258 prev_items = await self.mass.cache.get(
1259 key=queue_id,
1260 provider=self.domain,
1261 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1262 default=[],
1263 )
1264 queue_data = PlayerQueueData.from_cache(prev_state, prev_items)
1265 except Exception as err:
1266 self.logger.warning(
1267 "Failed to restore the queue(items) for %s - %s",
1268 player.state.name,
1269 str(err),
1270 )
1271 # Reset to clean state on failure
1272 queue_data = None
1273 if queue_data is None:
1274 queue_data = PlayerQueueData(
1275 queue=PlayerQueue(
1276 queue_id=queue_id,
1277 active=False,
1278 display_name=player.state.name,
1279 available=player.state.available,
1280 # Autoplay starts out on for a brand new queue; the player's own Autoplay
1281 # switch owns it from here on (and is restored above for a queue we know)
1282 autoplay_enabled=True,
1283 items=0,
1284 )
1285 )
1286
1287 self._queue_data[queue_id] = queue_data
1288 # always call update to calculate state etc
1289 self.on_player_update(player, {})
1290 self.mass.signal_event(EventType.QUEUE_ADDED, object_id=queue_id, data=queue_data.queue)
1291
1292 def on_player_update(
1293 self,
1294 player: Player,
1295 changed_values: dict[str, tuple[Any, Any]],
1296 ) -> None:
1297 """
1298 Call when a PlayerQueue needs to be updated (e.g. when player updates).
1299
1300 NOTE: This is called every second if the player is playing.
1301 """
1302 if player.type == PlayerType.PROTOCOL:
1303 # protocol players do not have a queue on their own
1304 return
1305 queue_id = player.player_id
1306 if (queue := self.get(queue_id)) is None:
1307 # race condition
1308 return
1309 if player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS):
1310 # do nothing while the announcement is in progress
1311 return
1312 # determine if this queue is currently active for this player
1313 queue.active = player.state.active_source in (queue.queue_id, None)
1314 if not queue.active and self._queue_data[queue_id].prev_state is None:
1315 queue.state = PlaybackState.IDLE
1316 # return early if the queue is not active and we have no previous state
1317 return
1318 if self._queue_data[queue_id].transitioning:
1319 # we're currently transitioning to a new track,
1320 # ignore updates from the player during this time
1321 return
1322 # queue is active and preflight checks passed, update the queue details
1323 self._update_queue_from_player(player)
1324
1325 def on_player_elapsed_time_corrected(self, player: Player) -> None:
1326 """Correct the queue's timing base if the player's real elapsed_time diverged."""
1327 if player.type == PlayerType.PROTOCOL:
1328 return
1329 queue_id = player.player_id
1330 if (queue := self.get(queue_id)) is None:
1331 return
1332 if not queue.active:
1333 return
1334 player_elapsed = player.state.corrected_elapsed_time
1335 if player_elapsed is None:
1336 return
1337 now = time.time()
1338 # queue.elapsed_time is stored in media-time so it can be displayed and
1339 # used as a resume position directly. The player reports stream-time
1340 # (post-atempo), so we scale by the current item's playback_speed.
1341 speed = get_current_playback_speed(queue)
1342 if queue.flow_mode:
1343 # _get_flow_queue_stream_index returns media-time in the current item
1344 # using each playlog entry's recorded speed.
1345 _, elapsed_time = self._get_flow_queue_stream_index(queue, player)
1346 else:
1347 elapsed_time = player_elapsed * speed
1348 if queue.current_item and queue.current_item.streamdetails:
1349 if seek_pos := queue.current_item.streamdetails.seek_position:
1350 elapsed_time += seek_pos
1351 queue.elapsed_time = elapsed_time
1352 queue.elapsed_time_last_updated = now
1353 queue.playback_speed = speed
1354 self.mass.signal_event(
1355 EventType.QUEUE_TIME_UPDATED,
1356 object_id=queue_id,
1357 data=queue.elapsed_time,
1358 )
1359
1360 def on_player_remove(self, player_id: str, permanent: bool) -> None:
1361 """Call when a player is removed from the registry."""
1362 self.mass.streams.audio_processing.clear(player_id)
1363 # cancel any pending play_index calls for this queue to prevent conflicts
1364 self.mass.cancel_timer(f"queue_play_index_{player_id}")
1365 # cancel a pending debounced cache write AND an already-started one, so neither can
1366 # recreate a deleted entry after the player is gone (the timer becomes a task once it fires)
1367 self.mass.cancel_timer(f"save_queue_cache_{player_id}")
1368 self.mass.cancel_task(f"save_queue_cache_{player_id}")
1369 self._set_transitioning(player_id, False)
1370 if permanent:
1371 self.purge_saved_queue(player_id)
1372 self._queue_data.pop(player_id, None)
1373 self._managed_pool.forget(player_id)
1374
1375 def purge_saved_queue(self, queue_id: str) -> None:
1376 """Delete the persisted state and items of the given queue."""
1377 for category in (CACHE_CATEGORY_PLAYER_QUEUE_STATE, CACHE_CATEGORY_PLAYER_QUEUE_ITEMS):
1378 # a removal runs both the player teardown and the config cleanup, so keep the
1379 # delete to one task per category instead of one per caller
1380 self.mass.create_task(
1381 self.mass.cache.delete(
1382 key=queue_id,
1383 provider=self.domain,
1384 category=category,
1385 ),
1386 task_id=f"purge_saved_queue_{queue_id}_{category}",
1387 )
1388
1389 async def load_next_queue_item(
1390 self,
1391 queue_id: str,
1392 current_item_id: str,
1393 ) -> QueueItem:
1394 """
1395 Call when a player wants the next queue item to play.
1396
1397 Raises QueueEmpty if there are no more tracks left.
1398 """
1399 queue = self.get(queue_id)
1400 if not queue:
1401 msg = f"PlayerQueue {queue_id} is not available"
1402 raise PlayerUnavailableError(msg)
1403 cur_index = self.index_by_id(queue_id, current_item_id)
1404 if cur_index is None:
1405 # this is just a guard for bad data
1406 raise QueueEmpty("Invalid item id for queue given.")
1407 next_item: QueueItem | None = None
1408 idx = 0
1409 while True:
1410 next_index = self._get_next_index(queue_id, cur_index + idx)
1411 if next_index is None:
1412 raise QueueEmpty("No more tracks left in the queue.")
1413 queue_item = self.get_item(queue_id, next_index)
1414 if queue_item is None:
1415 raise QueueEmpty("No more tracks left in the queue.")
1416 if idx >= 10:
1417 # we only allow 10 retries to prevent infinite loops
1418 raise QueueEmpty("No more (playable) tracks left in the queue.")
1419 try:
1420 await self._load_item(queue_item, next_index)
1421 # we're all set, this is our next item
1422 next_item = queue_item
1423 break
1424 except ProviderStreamLimitError:
1425 # transient source capacity, do not burn a playable item over it
1426 raise
1427 except MediaNotFoundError, AudioError:
1428 # No stream details found, skip this QueueItem
1429 self.logger.warning(
1430 "Skipping unplayable item %s (%s)", queue_item.name, queue_item.uri
1431 )
1432 queue_item.available = False
1433 idx += 1
1434 if idx != 0:
1435 # we skipped some items, signal a queue items update
1436 self.update_items(queue_id, self._queue_data[queue_id].items)
1437 if next_item is None:
1438 raise QueueEmpty("No more (playable) tracks left in the queue.")
1439
1440 # carry playback_speed forward across consecutive audiobook/podcast items
1441 current_item = self.get_item(queue_id, current_item_id)
1442 if (
1443 current_item
1444 and current_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1445 and next_item.media_type in (MediaType.AUDIOBOOK, MediaType.PODCAST_EPISODE)
1446 ):
1447 next_item.extra_attributes["playback_speed"] = current_item.extra_attributes.get(
1448 "playback_speed", 1.0
1449 )
1450
1451 return next_item
1452
1453 def track_loaded_in_buffer(self, queue_id: str, item_id: str) -> None:
1454 """Call when a player has (started) loading a track in the buffer."""
1455 queue = self.get(queue_id)
1456 if not queue:
1457 msg = f"PlayerQueue {queue_id} is not available"
1458 raise PlayerUnavailableError(msg)
1459 # store the index of the item that is currently (being) loaded in the buffer
1460 # which helps us a bit to determine how far the player has buffered ahead
1461 current_index = self.index_by_id(queue_id, item_id)
1462 queue.index_in_buffer = current_index
1463 self.logger.debug("PlayerQueue %s loaded item %s in buffer", queue.display_name, item_id)
1464 self.signal_update(queue_id)
1465 # preload next streamdetails
1466 self._preload_next_item(queue_id, item_id)
1467 # clean up stale audio buffers for old queue items to prevent memory leaks
1468 if current_index is not None:
1469 self.mass.create_task(self._cleanup_stale_queue_buffers(queue_id, current_index))
1470
1471 def queue_buffer_completed(self, queue_id: str, queue_exhausted: bool) -> None:
1472 """
1473 Call when the flow stream has finished generating all audio data for a queue.
1474
1475 At this point all audio data for the queue has been passed to the encoding pipeline.
1476 The player will go idle once it finishes playing the remaining buffered audio.
1477
1478 We start a background task that waits for the player to go idle and checks if new
1479 items have been added to the queue in the meantime, resuming playback if so.
1480
1481 :param queue_id: The queue ID.
1482 :param queue_exhausted: Whether the flow ended because the queue ran out of items,
1483 as opposed to ending early to restart on a format change or a live item.
1484 """
1485 queue = self.get(queue_id)
1486 if not queue:
1487 return
1488 self.logger.debug("Queue flow buffer completed for %s", queue.display_name)
1489
1490 # capture session_id so we can bail out if playback restarts
1491 queue_data = self._queue_data[queue_id]
1492 original_session_id = queue_data.session_id
1493 # record so player providers can detect flow EOF without an idle report
1494 if original_session_id is not None:
1495 queue_data.flow_buffer_completed = original_session_id
1496 if queue_exhausted:
1497 queue_data.flow_queue_exhausted = original_session_id
1498
1499 async def _resume_on_idle() -> None:
1500 # wait for the player to finish playing the buffered audio and go idle
1501 idle_detected = False
1502 for _ in range(60):
1503 await asyncio.sleep(1)
1504 if not queue.active or queue_data.session_id != original_session_id:
1505 return
1506 if queue.state == PlaybackState.IDLE:
1507 idle_detected = True
1508 break
1509 if not idle_detected:
1510 return
1511 # player went idle, give it a brief moment to settle
1512 await asyncio.sleep(1)
1513 if queue.state != PlaybackState.IDLE or queue_data.session_id != original_session_id:
1514 return
1515 # check if new items were added to the queue after the flow stream ended
1516 if queue.current_index is not None and (
1517 next_item := self.get_next_item(queue_id, queue.current_index)
1518 ):
1519 next_index = self.index_by_id(queue_id, next_item.queue_item_id)
1520 if next_index is not None:
1521 self.logger.info(
1522 "Resuming playback after flow stream completed for %s",
1523 queue.display_name,
1524 )
1525 await self.play_index(queue_id, next_index)
1526
1527 task_id = f"queue_buffer_completed_{queue_id}"
1528 self.mass.create_task(_resume_on_idle(), task_id=task_id)
1529
1530 def flow_stream_finished(self, queue_id: str) -> bool:
1531 """
1532 Return whether the flow stream for the current playback session is fully generated.
1533
1534 Lets player providers detect flow EOF when the device does not report idle
1535 (e.g. a Cast group that underruns the LIVE flow stream and keeps reporting playing).
1536
1537 :param queue_id: The queue ID.
1538 """
1539 queue_data = self.queue_data_or_none(queue_id)
1540 if queue_data is None or queue_data.session_id is None:
1541 return False
1542 return queue_data.flow_buffer_completed == queue_data.session_id
1543
1544 def flow_queue_exhausted(self, queue_id: str, session_id: str) -> bool:
1545 """
1546 Return whether the given flow stream session played the queue to its end.
1547
1548 False while a session is still streaming, and for a flow stream that ended early
1549 to be restarted (a format change or a live item), where the player is expected to
1550 pick up the next stream right away.
1551
1552 :param queue_id: The queue ID.
1553 :param session_id: The stream session to check.
1554 """
1555 queue_data = self.queue_data_or_none(queue_id)
1556 if queue_data is None or queue_data.session_id != session_id:
1557 return False
1558 return queue_data.flow_queue_exhausted == session_id
1559
1560 # Main queue manipulation methods
1561
1562 async def load(
1563 self,
1564 queue_id: str,
1565 queue_items: list[QueueItem],
1566 insert_at_index: int = 0,
1567 keep_remaining: bool = True,
1568 keep_played: bool = True,
1569 shuffle: bool = False,
1570 pin_first: bool = False,
1571 ) -> None:
1572 """
1573 Load new items at index.
1574
1575 - queue_id: id of the queue to process this request.
1576 - queue_items: a list of QueueItems
1577 - insert_at_index: insert the item(s) at this index
1578 - keep_remaining: keep the remaining items after the insert
1579 - shuffle: (re)shuffle the items after insert index
1580 - pin_first: keep the first item at the insert index instead of letting the shuffle
1581 move it; only meaningful together with shuffle
1582 """
1583 prev_items = self._queue_data[queue_id].items[:insert_at_index] if keep_played else []
1584 next_items = queue_items
1585
1586 # if keep_remaining, append the old 'next' items
1587 if keep_remaining:
1588 next_items += self._queue_data[queue_id].items[insert_at_index:]
1589
1590 # we set the original insert order as attribute so we can un-shuffle
1591 for index, item in enumerate(next_items):
1592 item.sort_index += insert_at_index + index
1593 # (re)shuffle the final batch if needed: smart shuffle when enabled, else pure random
1594 if shuffle:
1595 queue = self._queue_data[queue_id].queue
1596 # a user-picked item must stay the one that plays, so hold it out of the shuffle
1597 pinned = next_items[:1] if pin_first else []
1598 shuffled = next_items[1:] if pin_first else next_items
1599 if self._smart_shuffle.is_enabled(queue_id):
1600 shuffled = await self._smart_shuffle.arrange(queue, shuffled)
1601 else:
1602 shuffled = random.sample(shuffled, len(shuffled))
1603 next_items = pinned + shuffled
1604 self.update_items(queue_id, prev_items + next_items)
1605
1606 def update_items(self, queue_id: str, queue_items: list[QueueItem]) -> None:
1607 """Update the existing queue items, mostly caused by reordering."""
1608 self._queue_data[queue_id].items = queue_items
1609 queue = self._queue_data[queue_id].queue
1610 queue.items = len(self._queue_data[queue_id].items)
1611 self.signal_update(queue_id, True)
1612 if (
1613 queue.state == PlaybackState.PLAYING
1614 and queue.index_in_buffer is not None
1615 and queue.index_in_buffer == queue.current_index
1616 ):
1617 # if the queue is playing,
1618 # ensure to (re)queue the next track because it might have changed
1619 # note that we only do this if the player has loaded the current track
1620 # if not, we wait until it has loaded to prevent conflicts
1621 if next_item := self.get_next_item(queue_id, queue.index_in_buffer):
1622 self._enqueue_next_item(queue_id, next_item)
1623
1624 # Helper methods
1625
1626 def get_item(self, queue_id: str, item_id_or_index: int | str | None) -> QueueItem | None:
1627 """Get queue item by index or item_id."""
1628 if item_id_or_index is None:
1629 return None
1630 if (queue_data := self._queue_data.get(queue_id)) is None:
1631 return None
1632 queue_items = queue_data.items
1633 if isinstance(item_id_or_index, int) and len(queue_items) > item_id_or_index:
1634 return queue_items[item_id_or_index]
1635 if isinstance(item_id_or_index, str):
1636 return next((x for x in queue_items if x.queue_item_id == item_id_or_index), None)
1637 return None
1638
1639 def signal_update(self, queue_id: str, items_changed: bool = False) -> None:
1640 """Signal state changed of given queue."""
1641 if (queue_data := self._queue_data.get(queue_id)) is None:
1642 return
1643 queue = queue_data.queue
1644 # tell clients who owns the queue's ordering: the AudioSource uri while
1645 # delegated (see _get_delegated_source), None when MA owns it
1646 delegated = self._get_delegated_source(queue_id)
1647 queue.queue_owner = delegated[0].uri if delegated is not None else None
1648 # a mirrored shuffle write (streams controller) changes what smart shuffle
1649 # resolves to, so refresh the derived flag with every signaled update
1650 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1651 if items_changed:
1652 queue_data.items_cache_dirty = True
1653 self.mass.signal_event(EventType.QUEUE_ITEMS_UPDATED, object_id=queue_id, data=queue)
1654 self.mass.streams.audio_processing.prune(queue_id)
1655 # always send the base event
1656 self.mass.signal_event(EventType.QUEUE_UPDATED, object_id=queue_id, data=queue)
1657 # also signal update to the player itself so it can update its current_media
1658 self.mass.players.trigger_player_update(queue_id)
1659 # persist the (settings-bearing) queue state, debounced so a burst of updates or the
1660 # per-track updates during playback collapse into a single cache write
1661 self.mass.call_later(
1662 QUEUE_CACHE_SAVE_DELAY,
1663 self._save_queue_to_cache,
1664 queue_id,
1665 task_id=f"save_queue_cache_{queue_id}",
1666 )
1667
1668 def index_by_id(self, queue_id: str, queue_item_id: str) -> int | None:
1669 """Get index by queue_item_id."""
1670 if (queue_data := self._queue_data.get(queue_id)) is None:
1671 return None
1672 for index, item in enumerate(queue_data.items):
1673 if item.queue_item_id == queue_item_id:
1674 return index
1675 return None
1676
1677 async def get_tracks_for_playback(self, media_item: MediaItemType) -> list[Track]:
1678 """
1679 Return the playable tracks a media item resolves to, honoring the user's selection prefs.
1680
1681 :param media_item: The media item to resolve to playable tracks.
1682 """
1683 return await self._media_resolver.get_tracks_for_playback(media_item)
1684
1685 async def get_playlist_tracks(
1686 self, playlist: Playlist, start_item: str | None = None, sort_by: str | None = None
1687 ) -> list[PlaylistPlayableItem]:
1688 """
1689 Return the playable tracks for a playlist, honoring the user's selection prefs.
1690
1691 :param playlist: The playlist to resolve.
1692 :param start_item: Optional item URI to start the playlist from.
1693 :param sort_by: Optional sort key for the returned tracks.
1694 """
1695 return await self._media_resolver.get_playlist_tracks(playlist, start_item, sort_by)
1696
1697 async def get_dynamic_source_tracks(self, item: MediaItemType) -> list[Track]:
1698 """
1699 Return a fresh batch of tracks for a dynamic source (a dynamic playlist or radio station).
1700
1701 :param item: The dynamic playlist or radio station to fetch the next batch for.
1702 """
1703 return await self._media_resolver.get_dynamic_source_tracks(item)
1704
1705 def recency_windows(self) -> RecencyWindows:
1706 """Return the configured recency windows (a global setting; used for recency-aware gating)."""
1707 return self._smart_shuffle.windows()
1708
1709 async def player_media_from_queue_item(self, queue_item: QueueItem) -> PlayerMedia:
1710 """
1711 Parse PlayerMedia from QueueItem.
1712
1713 :param queue_item: The queue item to create media from.
1714 """
1715 queue_data = self._queue_data[queue_item.queue_id]
1716 stream_duration: int | None = None
1717 if queue_item.streamdetails:
1718 # prefer netto duration
1719 duration = queue_item.streamdetails.duration or queue_item.duration
1720 if duration and queue_item.streamdetails.seek_position:
1721 # the audio handed to the player starts at the seek position, so it is
1722 # shorter than the media item itself. seeking to (or past) the end
1723 # leaves no stream to describe, so the full length is kept instead.
1724 remaining = int(duration - queue_item.streamdetails.seek_position)
1725 stream_duration = remaining if remaining > 0 else None
1726 else:
1727 duration = queue_item.duration
1728 if queue_data.session_id is None:
1729 raise InvalidDataError("Queue session_id is None")
1730 media = PlayerMedia(
1731 uri=queue_item.uri,
1732 media_type=queue_item.media_type,
1733 title=queue_item.name,
1734 image_url=MASS_LOGO_ONLINE,
1735 duration=duration,
1736 stream_duration=stream_duration,
1737 source_id=queue_item.queue_id,
1738 queue_item_id=queue_item.queue_item_id,
1739 queue_session_id=queue_data.session_id,
1740 custom_data={
1741 "original_uri": queue_item.uri,
1742 },
1743 )
1744 if queue_item.media_item:
1745 media.title = queue_item.media_item.name
1746 media.artist = getattr(queue_item.media_item, "artist_str", "")
1747 media.album = (
1748 album.name if (album := getattr(queue_item.media_item, "album", None)) else ""
1749 )
1750 if queue_item.image:
1751 # the image format needs to be 512x512 jpeg for maximum compatibility with players
1752 # we prefer the imageproxy on the streamserver here because this request is sent
1753 # to the player itself which may not be able to reach the regular webserver
1754 media.image_url = self.mass.metadata.get_image_url(
1755 queue_item.image, size=512, image_format="jpeg", prefer_stream_server=True
1756 )
1757 return media
1758
1759 def get_next_item(self, queue_id: str, cur_index: int | str) -> QueueItem | None:
1760 """Return next QueueItem for given queue."""
1761 index: int
1762 if isinstance(cur_index, str):
1763 resolved_index = self.index_by_id(queue_id, cur_index)
1764 if resolved_index is None:
1765 return None # guard
1766 index = resolved_index
1767 else:
1768 index = cur_index
1769 # At this point index is guaranteed to be int
1770 for skip in range(5):
1771 if (next_index := self._get_next_index(queue_id, index + skip)) is None:
1772 break
1773 next_item = self.get_item(queue_id, next_index)
1774 if next_item is None:
1775 continue
1776 if not next_item.available:
1777 # ensure that we skip unavailable items (set by load_next track logic)
1778 continue
1779 return next_item
1780 return None
1781
1782 def store_sources(self, queue: PlayerQueue, items: list[MediaItemType]) -> None:
1783 """
1784 Hold the queue's full dynamic-source items server-side and project them onto `sources`.
1785
1786 :param queue: The queue whose sources are being set.
1787 :param items: The full source media items; an empty list clears the queue's sources.
1788 """
1789 self._queue_data[queue.queue_id].source_items = items
1790 # keep every occurrence server-side (a source added more than once weights it up in the
1791 # managed pool), but expose only the distinct container sources on the wire for clients to
1792 # show. Individual items (tracks, live radio streams, podcast episodes, ...) are omitted; see
1793 # `_WIRE_SOURCE_MEDIA_TYPES`. Autoplay/pool refill reads the full `source_items` above, not
1794 # this projected list, so it is unaffected.
1795 seen: set[str] = set()
1796 sources: list[ItemMapping] = []
1797 for item in items:
1798 if item.media_type not in _WIRE_SOURCE_MEDIA_TYPES and not is_dynamic_source(item):
1799 continue
1800 mapping = ItemMapping.from_item(item)
1801 if mapping.uri and mapping.uri in seen:
1802 continue
1803 if mapping.uri:
1804 seen.add(mapping.uri)
1805 sources.append(mapping)
1806 queue.sources = sources
1807 # release any materialized finite-source state whose source is no longer present
1808 self._managed_pool.retain(
1809 queue.queue_id, {item.uri for item in items if item.uri is not None}
1810 )
1811
1812 async def _save_queue_to_cache(self, queue_id: str) -> None:
1813 """Persist the queue's state (and its items when changed) to the cache."""
1814 if (queue_data := self._queue_data.get(queue_id)) is None:
1815 return
1816 try:
1817 # persistent so a cache clear/reset does not wipe the user's queues; the default
1818 # expiration still applies but is refreshed on every write. Skip the state write when its
1819 # persist-worthy content is unchanged (i.e. only playback progress advanced).
1820 state = queue_data.to_cache()
1821 significant = queue_data.cache_significant(state)
1822 if significant != queue_data.last_saved_state:
1823 await self.mass.cache.set(
1824 key=queue_id,
1825 data=state,
1826 provider=self.domain,
1827 category=CACHE_CATEGORY_PLAYER_QUEUE_STATE,
1828 persistent=True,
1829 )
1830 queue_data.last_saved_state = significant
1831 if queue_data.items_cache_dirty:
1832 # only cache items with a valid media_item
1833 await self.mass.cache.set(
1834 key=queue_id,
1835 data=queue_data.items_to_cache(),
1836 provider=self.domain,
1837 category=CACHE_CATEGORY_PLAYER_QUEUE_ITEMS,
1838 persistent=True,
1839 )
1840 queue_data.items_cache_dirty = False
1841 except Exception as err:
1842 self.logger.warning("Failed to persist the queue for %s - %s", queue_id, err)
1843
1844 def _check_player_permission(self, queue_id: str) -> None:
1845 """
1846 Check if the current user has permission to control this player/queue.
1847
1848 :param queue_id: The queue/player ID to check access for.
1849 :raises InsufficientPermissions: If the user lacks access.
1850 """
1851 current_user = get_current_user()
1852 if (
1853 current_user
1854 and current_user.player_filter
1855 and queue_id not in current_user.player_filter
1856 ):
1857 msg = f"{current_user.username} does not have access to player {queue_id}"
1858 raise InsufficientPermissions(msg)
1859
1860 @handle_play_action
1861 async def _handle_play(self, queue_id: str) -> None:
1862 """Handle play without acquiring the queue lock."""
1863 queue_player = self.mass.players.get_player(queue_id, True)
1864 if queue_player is None:
1865 raise PlayerUnavailableError(f"Player {queue_id} is not available")
1866 if (queue := self.get(queue_id)) and queue.active and queue.state == PlaybackState.PAUSED:
1867 # forward the actual play/unpause command to the player,
1868 # holding the action until the player confirms it resumed playback
1869 async with self.mass.players.wait_for_player_update(
1870 queue_id,
1871 attribute_name="playback_state",
1872 attribute_value=PlaybackState.PLAYING,
1873 timeout=PLAYBACK_START_TIMEOUT,
1874 ):
1875 await queue_player.play()
1876 return
1877 # player is not paused, perform resume instead
1878 await self.resume(queue_id)
1879
1880 def _set_transitioning(self, queue_id: str, value: bool) -> None:
1881 """Mark (or clear) whether a queue is mid-transition (no-op if it is not registered)."""
1882 if (queue_data := self._queue_data.get(queue_id)) is not None:
1883 queue_data.transitioning = value
1884
1885 def _clear(self, queue_id: str, skip_stop: bool = False) -> None:
1886 """Drop the queue's items and playback position, leaving user settings untouched."""
1887 queue = self._queue_data[queue_id].queue
1888 self.mass.streams.audio_processing.clear(queue_id)
1889 self.store_sources(queue, [])
1890 if queue.is_dynamic:
1891 # Dynamic sources impose shuffle, so clearing the source clears that shuffle too.
1892 queue.shuffle_enabled = False
1893 queue.is_dynamic = False
1894 # dropping the dynamic source changes what smart shuffle resolves to, so the derived
1895 # flag has to follow or clients keep showing a smart mix on a plain queue
1896 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
1897 queue.ended = False
1898 if queue.state != PlaybackState.IDLE and not skip_stop:
1899 self.mass.create_task(self.stop(queue_id))
1900 queue.current_index = None
1901 queue.current_item = None
1902 queue.elapsed_time = 0
1903 queue.elapsed_time_last_updated = time.time()
1904 queue.index_in_buffer = None
1905 self.mass.create_task(self._cleanup_queue_audio_data(queue_id))
1906 self.update_items(queue_id, [])
1907
1908 def _audio_source_plugin(
1909 self, queue_item: QueueItem | None
1910 ) -> tuple[PluginProvider, str] | None:
1911 """Return the plugin owning this item's live AudioSource, with the source id it exposes."""
1912 if queue_item is None or (media_item := queue_item.media_item) is None:
1913 return None
1914 if media_item.media_type != MediaType.AUDIO_SOURCE:
1915 return None
1916 if not isinstance(prov := self.mass.get_provider(media_item.provider), PluginProvider):
1917 return None
1918 return prov, media_item.item_id
1919
1920 def _audio_sources_in(self, queue_id: str) -> dict[str, tuple[PluginProvider, str]]:
1921 """
1922 Map every live AudioSource the queue holds to its owning plugin, keyed by media uri.
1923
1924 Covers the whole queue rather than just what is playing: an option that starts media
1925 alongside a live source leaves that source behind as an ordinary item, and it still has
1926 to be released when it eventually goes.
1927
1928 :param queue_id: The queue to look through.
1929 :return: The owning plugin and the source id it exposes, per AudioSource media uri.
1930 """
1931 owners: dict[str, tuple[PluginProvider, str]] = {}
1932 if (queue_data := self._queue_data.get(queue_id)) is None:
1933 return owners
1934 for item in queue_data.items:
1935 media_item = item.media_item
1936 if media_item is None or media_item.media_type != MediaType.AUDIO_SOURCE:
1937 continue
1938 if (uri := media_item.uri) is None or uri in owners:
1939 continue
1940 if (owner := self._audio_source_plugin(item)) is not None:
1941 owners[uri] = owner
1942 return owners
1943
1944 def _notify_audio_source_removed(self, queue_id: str) -> None:
1945 """Tell the owning plugins that their AudioSources are dropped along with this queue."""
1946 for prov, source_id in self._audio_sources_in(queue_id).values():
1947 self.mass.create_task(prov.on_source_removed(source_id, queue_id))
1948
1949 def _notify_audio_source_replaced(
1950 self, queue_id: str, outgoing: dict[str, tuple[PluginProvider, str]]
1951 ) -> None:
1952 """
1953 Tell the owning plugins which of their AudioSources the queue's new contents pushed out.
1954
1955 A source that is still among the queue's items is left in place, so the same source
1956 coming back and contents that keep it alongside them both release nothing.
1957
1958 :param queue_id: The queue whose contents were taken over.
1959 :param outgoing: The queue's live AudioSources from before the takeover.
1960 """
1961 if not outgoing or (queue_data := self._queue_data.get(queue_id)) is None:
1962 return
1963 remaining = {
1964 item.media_item.uri for item in queue_data.items if item.media_item is not None
1965 }
1966 for uri, (prov, source_id) in outgoing.items():
1967 if uri not in remaining:
1968 self.mass.create_task(prov.on_source_removed(source_id, queue_id))
1969
1970 async def _notify_audio_source_transferred(
1971 self,
1972 transferred: dict[str, tuple[PluginProvider, str]],
1973 from_queue_id: str,
1974 to_queue_id: str,
1975 ) -> None:
1976 """
1977 Tell the owning plugins that their AudioSources moved to another queue.
1978
1979 Awaited rather than dispatched so the plugins have settled the handover before the target
1980 queue is resumed. A plugin that raises must not break the transfer.
1981
1982 :param transferred: The live AudioSources the source queue handed over.
1983 :param from_queue_id: The queue that gave the sources up.
1984 :param to_queue_id: The queue that took them over.
1985 """
1986 if not transferred or (queue_data := self._queue_data.get(to_queue_id)) is None:
1987 return
1988 arrived = {item.media_item.uri for item in queue_data.items if item.media_item is not None}
1989 for uri, (prov, source_id) in transferred.items():
1990 if uri not in arrived:
1991 continue
1992 try:
1993 await prov.on_source_transferred(source_id, from_queue_id, to_queue_id)
1994 except Exception:
1995 self.logger.warning(
1996 "on_source_transferred raised for provider %s source %s",
1997 prov.instance_id,
1998 source_id,
1999 exc_info=True,
2000 )
2001
2002 def _reset_shuffle(self, queue_id: str) -> None:
2003 """Switch shuffle off."""
2004 queue = self._queue_data[queue_id].queue
2005 if not queue.shuffle_enabled:
2006 return
2007 queue.shuffle_enabled = False
2008 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
2009 self.signal_update(queue_id)
2010
2011 async def _apply_shuffle(
2012 self, queue_id: str, option: QueueOption, shuffle: bool | None
2013 ) -> None:
2014 """
2015 Settle the queue's shuffle state for a play command before its items are resolved.
2016
2017 :param queue_id: The queue the media is played on.
2018 :param option: The enqueue option this command resolved to.
2019 :param shuffle: The state to put the queue's shuffle in; None to leave it as it is.
2020 """
2021 queue = self._queue_data[queue_id].queue
2022 if queue.is_dynamic and option in (
2023 QueueOption.PLAY,
2024 QueueOption.REPLACE,
2025 QueueOption.REPLACE_NEXT,
2026 ):
2027 # These are the options that replace the queue's sources, so the smart mix may be on
2028 # its way out - and its shuffle is never the user's own (a dynamic queue's toggle is
2029 # locked), so it must not outlive the source that imposed it. Recorded directly
2030 # because set_shuffle refuses a queue that is still a smart mix, and the items are
2031 # resolved against this flag. The state is provisional until the sources are known:
2032 # `_enter_dynamic_mode` forces shuffle back on if the queue stays dynamic.
2033 if option == QueueOption.REPLACE_NEXT:
2034 # staging leaves the shuffle the user chose alone, so it never carries a request
2035 # of its own to honour here
2036 queue.shuffle_enabled = False
2037 else:
2038 queue.shuffle_enabled = bool(shuffle)
2039 return
2040 if shuffle is None or option not in (QueueOption.PLAY, QueueOption.REPLACE):
2041 # nothing to settle: the media brings no order of its own to protect, or the option
2042 # only stages items for later and leaves the queue's shuffle state alone
2043 return
2044 if queue.shuffle_enabled == shuffle:
2045 return
2046 if self._get_delegated_source(queue_id) is not None:
2047 # the queue is still delegated to the external session this play is replacing,
2048 # so set_shuffle would forward the toggle to that session; apply the state and
2049 # the tail re-order locally instead â MA-owned items behind the session may
2050 # sit in shuffled order and must not survive an ordered play unshuffled-in-name-only
2051 await self._apply_local_shuffle(queue_id, shuffle)
2052 return
2053 # routed through set_shuffle so switching shuffle off also restores the order of
2054 # the items that stay in the queue: a play keeps them, and a tail left in shuffled
2055 # order behind a queue that now reads unshuffled would contradict its own flag
2056 await self.set_shuffle(queue_id, shuffle)
2057
2058 async def _apply_mirrored_shuffle(
2059 self, queue_id: str, source_id: str, provider_instance: str, shuffle_enabled: bool
2060 ) -> None:
2061 """
2062 Apply a session-mirrored shuffle change, serialized against playback commands.
2063
2064 Scheduled as a task by the streams controller; the delegation and source
2065 identity are re-validated under the player lock so a session that ended (or
2066 changed) between the options event and this task running cannot re-order an
2067 unrelated queue.
2068
2069 :param queue_id: The queue the session mirror applies to.
2070 :param source_id: The AudioSource.item_id that reported the change.
2071 :param provider_instance: The provider instance id that reported the change.
2072 :param shuffle_enabled: The session's shuffle state to mirror.
2073 """
2074 async with self.mass.players.get_player_lock(queue_id):
2075 delegated = self._get_delegated_source(queue_id)
2076 if delegated is None:
2077 return
2078 audio_source = delegated[0]
2079 if audio_source.item_id != source_id or audio_source.provider != provider_instance:
2080 return
2081 if self._queue_data[queue_id].queue.shuffle_enabled == shuffle_enabled:
2082 return
2083 await self._apply_local_shuffle(queue_id, shuffle_enabled)
2084
2085 async def _apply_local_shuffle(self, queue_id: str, shuffle_enabled: bool) -> None:
2086 """
2087 Record the queue's shuffle state and re-order the un-played tail accordingly.
2088
2089 :param queue_id: The queue to apply the shuffle state to.
2090 :param shuffle_enabled: The shuffle state to record and apply to the tail.
2091 """
2092 queue = self._queue_data[queue_id].queue
2093 queue.shuffle_enabled = shuffle_enabled
2094 queue.smart_shuffle_active = self.is_smart_shuffle_active(queue)
2095 queue_items = self._queue_data[queue_id].items
2096 cur_index = (
2097 queue.index_in_buffer if queue.index_in_buffer is not None else queue.current_index
2098 )
2099 if cur_index is not None:
2100 next_index = cur_index + 1
2101 next_items = queue_items[next_index:]
2102 else:
2103 next_items = []
2104 next_index = 0
2105 if not shuffle_enabled:
2106 # shuffle disabled, try to restore original sort order of the remaining items
2107 next_items.sort(key=lambda x: x.sort_index, reverse=False)
2108 await self.load(
2109 queue_id=queue_id,
2110 queue_items=next_items,
2111 insert_at_index=next_index,
2112 keep_remaining=False,
2113 shuffle=shuffle_enabled,
2114 )
2115
2116 def _get_delegated_source(
2117 self, queue_id: str
2118 ) -> tuple[AudioSource, SourceQueueCapabilities, PluginProvider] | None:
2119 """
2120 Return the AudioSource owning the queue's commands, its capabilities and owning plugin.
2121
2122 While the queue's current item is an AudioSource declaring ``queue_capabilities``,
2123 the external session owns the queue: commands are forwarded to the owning plugin
2124 and the mirrored options event updates the queue state afterwards. This matches
2125 the player-layer proxy: no playback-state gate, because a plugin may map a paused
2126 session onto a stopped player (Spotify Connect does) â a genuinely dead session
2127 is the owning plugin's call, surfaced as its localized not-active error. Returns
2128 None â queue commands then apply to the MA queue as usual â when the current item
2129 is not such an AudioSource (a transport-only source keeps queue commands with MA)
2130 or when the owning plugin provider is no longer available.
2131
2132 :param queue_id: The queue to inspect.
2133 """
2134 if (queue_data := self._queue_data.get(queue_id)) is None:
2135 return None
2136 current_item = queue_data.queue.current_item
2137 if current_item is None or current_item.media_item is None:
2138 return None
2139 media_item = current_item.media_item
2140 if not isinstance(media_item, AudioSource):
2141 return None
2142 caps = media_item.queue_capabilities
2143 if caps is None:
2144 return None
2145 provider = self.mass.get_provider(media_item.provider)
2146 if not isinstance(provider, PluginProvider):
2147 return None
2148 # Belt-and-suspenders: a queue item carrying media_type=AUDIO_SOURCE
2149 # can only have come from a provider that declared the feature, but
2150 # a feature flag flipped off at runtime (provider reload, config
2151 # change) would leave on_source_control raising NotImplementedError.
2152 if ProviderFeature.AUDIO_SOURCE not in provider.supported_features:
2153 return None
2154 return media_item, caps, provider
2155