/
/
1"""
2Play queue command handlers for Plex remote control.
3
4Handles playMedia, createPlayQueue and refreshPlayQueue HTTP commands.
5"""
6
7from __future__ import annotations
8
9import asyncio
10import logging
11import re
12from typing import TYPE_CHECKING, Any
13from urllib.parse import unquote
14
15from aiohttp import web
16from music_assistant_models.enums import QueueOption
17from plexapi.playqueue import PlayQueue
18
19from .parsing import plex_item_fields
20
21if TYPE_CHECKING:
22 from music_assistant.providers.plex import PlexProvider
23
24LOGGER = logging.getLogger(__name__)
25
26# Maximum number of tracks loaded from a Plex play queue into Music Assistant.
27# Plex play queues can be very large (whole libraries); capping keeps queue loading
28# and MA<->Plex sync responsive.
29MAX_QUEUE_ITEMS = 100
30
31
32class QueueCommandsMixin:
33 """Mixin providing HTTP handlers for Plex play queue commands."""
34
35 if TYPE_CHECKING:
36 provider: PlexProvider
37 _ma_player_id: str | None
38 _updating_from_plex: bool
39 play_queue_id: str | None
40 play_queue_version: int
41 play_queue_item_ids: dict[int, int]
42 _last_synced_ma_queue_length: int
43 _last_synced_ma_queue_keys: list[str]
44
45 async def _broadcast_timeline(self) -> None: ...
46
47 async def _seek_to_offset_after_playback(self, player_id: str, offset: int) -> None: ...
48
49 async def _load_remaining_queue_tracks(
50 self,
51 player_id: str,
52 playqueue: PlayQueue,
53 selected_offset: int,
54 shuffle: bool,
55 ) -> None: ...
56
57 async def _replace_entire_queue(self, player_id: str, playqueue: PlayQueue) -> None: ...
58
59 async def _replace_remaining_queue(
60 self, player_id: str, playqueue: PlayQueue, current_index: int
61 ) -> None: ...
62
63 def _remember_synced_queue(
64 self, player_id: str, keys: list[str] | None = None
65 ) -> list[str]: ...
66
67 async def _create_plex_playqueue_from_ma(self) -> None: ...
68
69 async def _play_single_track(self, player_id: str, key: str) -> None:
70 """
71 Resolve a single Plex track key and play it as a fresh (replaced) queue.
72
73 Used as the fallback whenever loading a full Plex play queue fails.
74
75 :param player_id: The Music Assistant player ID.
76 :param key: The Plex track key to play.
77 """
78 track = await self.provider.get_track(key)
79 await self.provider.mass.player_queues.play_media(
80 queue_id=player_id,
81 media=track,
82 option=QueueOption.REPLACE,
83 )
84
85 async def _fetch_full_play_queue(self, queue_id: str) -> PlayQueue | None:
86 """
87 Fetch a PlayQueue, paginating past the server's per-request window cap.
88
89 The Plex server caps each response to ~200 items regardless of the requested
90 window size. We use playQueueTotalCount to detect truncation and keep fetching
91 forward pages until we have every item, up to :data:`MAX_QUEUE_ITEMS`.
92
93 :param queue_id: The Plex PlayQueue ID to fetch.
94 :return: A PlayQueue whose items list contains the tracks (capped), or None.
95 """
96 page_size = 200
97 plex_server = self.provider._plex_server
98
99 def fetch_initial() -> PlayQueue:
100 # own=True transfers ownership of this queue to MA so we can control it fully.
101 return PlayQueue.get(plex_server, playQueueID=queue_id, own=True, window=page_size)
102
103 playqueue = await asyncio.to_thread(fetch_initial)
104
105 # Avoid truthiness checks on the PlayQueue object: plexapi's __len__ returns
106 # playQueueTotalCount, which is None for track-radio (and other on-the-fly)
107 # queues, so `not playqueue` would raise a TypeError.
108 if playqueue is None or not playqueue.items:
109 return playqueue
110
111 all_items = list(playqueue.items)
112 seen_ids = {item.playQueueItemID for item in all_items}
113
114 # playQueueTotalCount is absent for track-radio queues; fall back to the
115 # number of items we already have so pagination simply stops.
116 total_count = playqueue.playQueueTotalCount or len(all_items)
117 target_count = min(total_count, MAX_QUEUE_ITEMS)
118 while len(all_items) < target_count:
119 last_id = all_items[-1].playQueueItemID
120
121 def fetch_next(_last_id: int = last_id) -> PlayQueue:
122 # own=False — we already claimed ownership on the first fetch.
123 # includeBefore=False — only items strictly after center are returned.
124 return PlayQueue.get(
125 plex_server,
126 playQueueID=queue_id,
127 own=False,
128 center=_last_id,
129 window=page_size,
130 includeBefore=False,
131 )
132
133 next_page = await asyncio.to_thread(fetch_next)
134 if next_page is None or not next_page.items:
135 break
136
137 new_items = [i for i in next_page.items if i.playQueueItemID not in seen_ids]
138 if not new_items:
139 break
140
141 all_items.extend(new_items)
142 seen_ids.update(i.playQueueItemID for i in new_items)
143
144 if len(all_items) > MAX_QUEUE_ITEMS:
145 LOGGER.info(
146 "Capping Plex play queue from %d to %d items", len(all_items), MAX_QUEUE_ITEMS
147 )
148 all_items = all_items[:MAX_QUEUE_ITEMS]
149
150 # Patch the cached items property on the PlayQueue object.
151 playqueue.__dict__["items"] = all_items
152 return playqueue
153
154 def _selected_item_index(self, playqueue: PlayQueue) -> int:
155 """Return the selected item's index within the fetched queue window."""
156 selected_id = getattr(playqueue, "playQueueSelectedItemID", None)
157 if selected_id is not None:
158 for index, item in enumerate(playqueue.items):
159 if getattr(item, "playQueueItemID", None) == selected_id:
160 return index
161 selected_offset = getattr(playqueue, "playQueueSelectedItemOffset", 0) or 0
162 return selected_offset if 0 <= selected_offset < len(playqueue.items) else 0
163
164 def _source_key_from_play_queue_uri(self, source_uri: str) -> str | None:
165 """
166 Extract a Plex library key from a playQueueSourceURI string.
167
168 Plex encodes the source as ``library:///directory/ENCODED_PATH``. We decode it
169 to a plain library path that :meth:`_resolve_plex_item` can use.
170
171 :param source_uri: The playQueueSourceURI from a Plex PlayQueue.
172 :return: A Plex library key (e.g. ``/library/playlists/5329``), or None if unparsable.
173 """
174 prefix = "library:///directory/"
175 if not source_uri.startswith(prefix):
176 return None
177
178 path = unquote(source_uri[len(prefix) :]).lstrip("/")
179 if not path:
180 return None
181
182 path = "/" + path.split("?")[0]
183
184 for suffix in ("/children", "/allLeaves", "/items"):
185 if path.endswith(suffix):
186 path = path[: -len(suffix)]
187 break
188
189 return path if path and path != "/" else None
190
191 async def _play_from_shuffled_source(
192 self,
193 player_id: str,
194 source_uri: str,
195 offset: int,
196 ) -> bool:
197 """
198 Load a shuffled PlayQueue's source in original order, then defer shuffle to MA.
199
200 When Plex reports a shuffled PlayQueue the items are already in Plex's shuffled
201 order. Loading them directly would cause MA to shuffle an already-shuffled list.
202 Instead we parse the source URI, load the source collection unshuffled into MA,
203 and schedule a deferred task that applies MA's own shuffle and then recreates the
204 Plex PlayQueue so Plexamp sees the new order.
205
206 :param player_id: The Music Assistant player ID.
207 :param source_uri: The playQueueSourceURI from the Plex PlayQueue.
208 :param offset: Starting position in milliseconds.
209 :return: True if handled, False if the caller should fall back to regular loading.
210 """
211 source_key = self._source_key_from_play_queue_uri(source_uri)
212 if not source_key:
213 return False
214
215 try:
216 LOGGER.info(f"Shuffled queue detected — loading source in original order: {source_key}")
217 source_media = await self._resolve_plex_item(source_key)
218 await self.provider.mass.player_queues.play_media(
219 queue_id=player_id,
220 media=source_media,
221 option=QueueOption.REPLACE,
222 )
223 if offset > 0:
224 await self._seek_to_offset_after_playback(player_id, offset)
225 await self._broadcast_timeline()
226
227 async def _apply_shuffle_deferred() -> None:
228 await self.provider.mass.player_queues.set_shuffle(player_id, True)
229 await asyncio.sleep(0.2)
230 await self._create_plex_playqueue_from_ma()
231 self._remember_synced_queue(player_id)
232
233 self.provider.mass.create_task(_apply_shuffle_deferred())
234 return True
235 except Exception as e:
236 LOGGER.debug(f"Could not resolve source for shuffled queue, falling back: {e}")
237 return False
238
239 async def _resolve_plex_item(self, key: str) -> Any:
240 """
241 Resolve a Plex key to a Music Assistant media item.
242
243 :param key: The Plex key to resolve.
244 """
245 if "/library/metadata/" in key:
246 try:
247 return await self.provider.get_track(key)
248 except Exception as exc:
249 LOGGER.debug(f"Failed to resolve Plex item as track for key '{key}': {exc}")
250
251 try:
252 return await self.provider.get_album(key)
253 except Exception as exc:
254 LOGGER.debug(f"Failed to resolve Plex item as album for key '{key}': {exc}")
255
256 try:
257 return await self.provider.get_artist(key)
258 except Exception:
259 raise ValueError(f"Could not resolve Plex item: {key}") from None
260
261 elif "/playlists/" in key:
262 return await self.provider.get_playlist(key)
263 else:
264 raise ValueError(f"Unknown Plex key format: {key}")
265
266 async def _play_from_plex_queue(
267 self,
268 player_id: str,
269 container_key: str,
270 starting_key: str | None,
271 shuffle: bool,
272 offset: int,
273 ) -> None:
274 """
275 Fetch a Plex PlayQueue and start playback, loading remaining tracks in the background.
276
277 :param player_id: The Music Assistant player ID.
278 :param container_key: The Plex container key (e.g. /playQueues/123).
279 :param starting_key: Fallback track key if queue fetch fails.
280 :param shuffle: Whether shuffle is enabled.
281 :param offset: Starting position in milliseconds.
282 """
283 try:
284 LOGGER.info(f"Fetching play queue: {container_key}")
285
286 queue_id_match = re.search(r"/playQueues/(\d+)", container_key)
287 if not queue_id_match:
288 raise ValueError(f"Invalid container_key format: {container_key}")
289
290 queue_id = queue_id_match.group(1)
291
292 playqueue = await self._fetch_full_play_queue(queue_id)
293
294 if playqueue is not None and playqueue.items:
295 selected_offset = self._selected_item_index(playqueue)
296 LOGGER.info(f"PlayQueue selected item index: {selected_offset}")
297
298 # When Plex reports a shuffled queue, load the original source into MA
299 # unshuffled and let MA apply its own shuffle, then propagate back to Plex.
300 if playqueue.playQueueShuffled and getattr(playqueue, "playQueueSourceURI", None):
301 if await self._play_from_shuffled_source(
302 player_id, playqueue.playQueueSourceURI, offset
303 ):
304 return
305
306 self.play_queue_item_ids = {}
307
308 first_item = playqueue.items[selected_offset]
309 first_track_key, first_play_queue_item_id = plex_item_fields(first_item)
310
311 if not first_track_key:
312 LOGGER.error("No valid first track in play queue")
313 if starting_key:
314 await self._play_single_track(player_id, starting_key)
315 return
316
317 try:
318 first_track = await self.provider.get_track(first_track_key)
319 LOGGER.info(f"Starting playback with first track: {first_track.name}")
320
321 if first_play_queue_item_id:
322 self.play_queue_item_ids[0] = first_play_queue_item_id
323
324 await self.provider.mass.player_queues.play_media(
325 queue_id=player_id,
326 media=first_track,
327 option=QueueOption.REPLACE,
328 )
329
330 if offset > 0:
331 await self._seek_to_offset_after_playback(player_id, offset)
332
333 await self._broadcast_timeline()
334
335 # Use the queue's own shuffle flag as the authoritative state,
336 # falling back to the request parameter if not set.
337 effective_shuffle = playqueue.playQueueShuffled or shuffle
338 self.provider.mass.create_task(
339 self._load_remaining_queue_tracks(
340 player_id, playqueue, selected_offset, effective_shuffle
341 )
342 )
343
344 except Exception:
345 LOGGER.exception("Error starting playback with first track")
346 if starting_key:
347 await self._play_single_track(player_id, starting_key)
348 else:
349 LOGGER.error("Play queue is empty or could not be fetched")
350 if starting_key:
351 await self._play_single_track(player_id, starting_key)
352
353 except Exception:
354 LOGGER.exception("Error playing from queue")
355 if starting_key:
356 await self._play_single_track(player_id, starting_key)
357
358 async def handle_play_media(self, request: web.Request) -> web.Response:
359 """
360 Handle playMedia command from Plex controller.
361
362 Plexamp sends various parameters:
363 - key: The item to play (track, album, playlist, etc.)
364 - containerKey: The container context (play queue)
365 - offset: Starting position in milliseconds
366 - shuffle: Whether to shuffle
367 - repeat: Repeat mode
368 """
369 self._updating_from_plex = True
370 try:
371 key = request.query.get("key")
372 container_key = request.query.get("containerKey")
373 offset = int(request.query.get("offset", 0))
374 shuffle = request.query.get("shuffle", "0") == "1"
375
376 if not key:
377 return web.Response(
378 status=400, text="Missing required 'key' parameter for playMedia command"
379 )
380
381 LOGGER.info(
382 f"Received playMedia command - key: {key}, "
383 f"containerKey: {container_key}, offset: {offset}ms"
384 )
385
386 player_id = self._ma_player_id
387 if not player_id:
388 return web.Response(status=500, text="No player assigned to this server")
389
390 if container_key and "/playQueues/" in container_key:
391 queue_id_match = re.search(r"/playQueues/(\d+)", container_key)
392 if queue_id_match:
393 self.play_queue_id = queue_id_match.group(1)
394 self.play_queue_version = 1
395 LOGGER.info(f"Playing from queue: {container_key} starting at {key}")
396 await self._play_from_plex_queue(player_id, container_key, key, shuffle, offset)
397 else:
398 self.play_queue_id = None
399 self.play_queue_item_ids = {}
400 media = await self._resolve_plex_item(key)
401 await self.provider.mass.player_queues.play_media(
402 queue_id=player_id,
403 media=media,
404 option=QueueOption.REPLACE,
405 )
406 elif container_key:
407 self.play_queue_id = None
408 self.play_queue_item_ids = {}
409 media_to_play = await self._resolve_plex_item(container_key)
410 await self.provider.mass.player_queues.play_media(
411 queue_id=player_id,
412 media=media_to_play,
413 option=QueueOption.REPLACE,
414 )
415 else:
416 self.play_queue_id = None
417 self.play_queue_item_ids = {}
418 media = await self._resolve_plex_item(key)
419 await self.provider.mass.player_queues.play_media(
420 queue_id=player_id,
421 media=media,
422 option=QueueOption.REPLACE,
423 )
424
425 # Always sync shuffle state so that a previously enabled MA shuffle
426 # does not reorder an unshuffled Plex queue.
427 await self.provider.mass.player_queues.set_shuffle(player_id, shuffle)
428
429 if offset > 0:
430 await self._seek_to_offset_after_playback(player_id, offset)
431
432 await self._broadcast_timeline()
433 return web.Response(status=200)
434
435 except Exception:
436 LOGGER.exception("Error handling playMedia")
437 return web.Response(status=500, text="Internal error")
438 finally:
439 self._updating_from_plex = False
440
441 async def handle_create_play_queue(self, request: web.Request) -> web.Response:
442 """
443 Handle createPlayQueue command from Plex controller.
444
445 Creates a new play queue from a URI (album, playlist, artist tracks, etc.)
446 and optionally applies shuffle.
447 """
448 self._updating_from_plex = True
449 try:
450 uri = request.query.get("uri")
451 shuffle = request.query.get("shuffle", "0") == "1"
452 continuous = request.query.get("continuous", "0") == "1"
453
454 if not uri:
455 return web.Response(status=400, text="Missing 'uri' parameter")
456
457 LOGGER.info(f"Received createPlayQueue command - uri: {uri}, shuffle: {shuffle}")
458
459 player_id = self._ma_player_id
460 if not player_id:
461 return web.Response(status=500, text="No player assigned to this server")
462
463 def create_queue() -> PlayQueue:
464 item = self.provider._plex_server.fetchItem(uri)
465 return PlayQueue.create(
466 self.provider._plex_server,
467 item,
468 shuffle=1 if shuffle else 0,
469 continuous=1 if continuous else 0,
470 )
471
472 playqueue = await asyncio.to_thread(create_queue)
473
474 if playqueue is not None and playqueue.items:
475 self.play_queue_id = str(playqueue.playQueueID)
476 self.play_queue_version = 1
477
478 if len(playqueue.items) > MAX_QUEUE_ITEMS:
479 LOGGER.info(
480 "Capping created Plex play queue from %d to %d items",
481 len(playqueue.items),
482 MAX_QUEUE_ITEMS,
483 )
484 playqueue.__dict__["items"] = playqueue.items[:MAX_QUEUE_ITEMS]
485
486 LOGGER.info(
487 f"Created play queue {self.play_queue_id} with {len(playqueue.items)} items"
488 )
489
490 self.play_queue_item_ids = {}
491 first_item = playqueue.items[0]
492 first_track_key, first_play_queue_item_id = plex_item_fields(first_item)
493
494 if not first_track_key:
495 LOGGER.error("No valid first track in created play queue")
496 return web.Response(status=500, text="Failed to load tracks from play queue")
497
498 try:
499 first_track = await self.provider.get_track(first_track_key)
500 LOGGER.info(f"Starting playback with first track: {first_track.name}")
501
502 if first_play_queue_item_id:
503 self.play_queue_item_ids[0] = first_play_queue_item_id
504
505 await self.provider.mass.player_queues.play_media(
506 queue_id=player_id,
507 media=first_track,
508 option=QueueOption.REPLACE,
509 )
510
511 if len(playqueue.items) > 1:
512 self.provider.mass.create_task(
513 self._load_remaining_queue_tracks(player_id, playqueue, 0, shuffle)
514 )
515
516 await self._broadcast_timeline()
517 return web.Response(status=200)
518
519 except Exception:
520 LOGGER.exception("Error starting playback with first track")
521 return web.Response(status=500, text="Failed to start playback")
522 else:
523 LOGGER.error("Failed to create play queue or queue is empty")
524 return web.Response(status=500, text="Failed to create play queue")
525
526 except Exception:
527 LOGGER.exception("Error handling createPlayQueue")
528 return web.Response(status=500, text="Internal error")
529 finally:
530 self._updating_from_plex = False
531
532 async def handle_refresh_play_queue(self, request: web.Request) -> web.Response:
533 """
534 Handle refreshPlayQueue command from Plex controller.
535
536 Called when the play queue is modified (items added, removed, reordered).
537 Syncs the updated queue state to MA while preserving current playback.
538 """
539 self._updating_from_plex = True
540 try:
541 play_queue_id = request.query.get("playQueueID")
542
543 if not play_queue_id:
544 return web.Response(status=400, text="Missing 'playQueueID' parameter")
545
546 LOGGER.info(
547 f"Received refreshPlayQueue command - playQueueID: {play_queue_id}, "
548 f"params: {dict(request.query)}"
549 )
550
551 if self.play_queue_id != play_queue_id:
552 LOGGER.warning(
553 f"Refresh requested for queue {play_queue_id} but active queue is "
554 f"{self.play_queue_id}"
555 )
556 return web.Response(
557 status=409,
558 text=(
559 f"Requested playQueueID {play_queue_id} does not match "
560 f"active queue {self.play_queue_id}"
561 ),
562 )
563
564 self.play_queue_version += 1
565
566 playqueue = await self._fetch_full_play_queue(play_queue_id)
567
568 if playqueue is None or not playqueue.items:
569 LOGGER.error("Failed to refresh play queue - queue is empty or not found")
570 return web.Response(status=404, text="Play queue not found")
571
572 player_id = self._ma_player_id
573 if not player_id:
574 LOGGER.error("No player assigned to this server")
575 return web.Response(status=500, text="No player assigned")
576
577 ma_queue = self.provider.mass.player_queues.get(player_id)
578 if not ma_queue:
579 LOGGER.error(f"MA queue not found for player {player_id}")
580 return web.Response(status=500, text="MA queue not found")
581
582 current_index = ma_queue.current_index
583 ma_queue_items = self.provider.mass.player_queues.items(player_id)
584 ma_queue_count = len(ma_queue_items) if ma_queue_items else 0
585
586 LOGGER.debug(
587 f"Queue refresh: Current index={current_index}, "
588 f"MA has {ma_queue_count} items, Plex has {len(playqueue.items)} items"
589 )
590
591 if current_index is None:
592 LOGGER.debug("No track currently playing, replacing entire queue")
593 await self._replace_entire_queue(player_id, playqueue)
594 else:
595 LOGGER.debug(
596 f"Track at index {current_index} is playing, "
597 f"replacing only items after current track"
598 )
599 await self._replace_remaining_queue(player_id, playqueue, current_index)
600
601 # Sync shuffle state from Plex to MA.
602 await self.provider.mass.player_queues.set_shuffle(
603 player_id, playqueue.playQueueShuffled
604 )
605
606 LOGGER.info(
607 f"Refreshed play queue {play_queue_id} - now has {len(playqueue.items)} items"
608 )
609
610 self._remember_synced_queue(player_id)
611
612 return web.Response(status=200)
613
614 except Exception:
615 LOGGER.exception("Error handling refreshPlayQueue")
616 return web.Response(status=500, text="Internal error")
617 finally:
618 self._updating_from_plex = False
619