/
/
/
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 # the deferred task below applies MA's shuffle; loading the source already
223 # shuffled would leave the queue's own order out of step with the Plex one
224 shuffle=False,
225 )
226 if offset > 0:
227 await self._seek_to_offset_after_playback(player_id, offset)
228 await self._broadcast_timeline()
229
230 async def _apply_shuffle_deferred() -> None:
231 await self.provider.mass.player_queues.set_shuffle(player_id, True)
232 await asyncio.sleep(0.2)
233 await self._create_plex_playqueue_from_ma()
234 self._remember_synced_queue(player_id)
235
236 self.provider.mass.create_task(_apply_shuffle_deferred())
237 return True
238 except Exception as e:
239 LOGGER.debug(f"Could not resolve source for shuffled queue, falling back: {e}")
240 return False
241
242 async def _resolve_plex_item(self, key: str) -> Any:
243 """
244 Resolve a Plex key to a Music Assistant media item.
245
246 :param key: The Plex key to resolve.
247 """
248 if "/library/metadata/" in key:
249 try:
250 return await self.provider.get_track(key)
251 except Exception as exc:
252 LOGGER.debug(f"Failed to resolve Plex item as track for key '{key}': {exc}")
253
254 try:
255 return await self.provider.get_album(key)
256 except Exception as exc:
257 LOGGER.debug(f"Failed to resolve Plex item as album for key '{key}': {exc}")
258
259 try:
260 return await self.provider.get_artist(key)
261 except Exception:
262 raise ValueError(f"Could not resolve Plex item: {key}") from None
263
264 elif "/playlists/" in key:
265 return await self.provider.get_playlist(key)
266 else:
267 raise ValueError(f"Unknown Plex key format: {key}")
268
269 async def _play_from_plex_queue(
270 self,
271 player_id: str,
272 container_key: str,
273 starting_key: str | None,
274 shuffle: bool,
275 offset: int,
276 ) -> None:
277 """
278 Fetch a Plex PlayQueue and start playback, loading remaining tracks in the background.
279
280 :param player_id: The Music Assistant player ID.
281 :param container_key: The Plex container key (e.g. /playQueues/123).
282 :param starting_key: Fallback track key if queue fetch fails.
283 :param shuffle: Whether shuffle is enabled.
284 :param offset: Starting position in milliseconds.
285 """
286 try:
287 LOGGER.info(f"Fetching play queue: {container_key}")
288
289 queue_id_match = re.search(r"/playQueues/(\d+)", container_key)
290 if not queue_id_match:
291 raise ValueError(f"Invalid container_key format: {container_key}")
292
293 queue_id = queue_id_match.group(1)
294
295 playqueue = await self._fetch_full_play_queue(queue_id)
296
297 if playqueue is not None and playqueue.items:
298 selected_offset = self._selected_item_index(playqueue)
299 LOGGER.info(f"PlayQueue selected item index: {selected_offset}")
300
301 # When Plex reports a shuffled queue, load the original source into MA
302 # unshuffled and let MA apply its own shuffle, then propagate back to Plex.
303 if playqueue.playQueueShuffled and getattr(playqueue, "playQueueSourceURI", None):
304 if await self._play_from_shuffled_source(
305 player_id, playqueue.playQueueSourceURI, offset
306 ):
307 return
308
309 self.play_queue_item_ids = {}
310
311 first_item = playqueue.items[selected_offset]
312 first_track_key, first_play_queue_item_id = plex_item_fields(first_item)
313
314 if not first_track_key:
315 LOGGER.error("No valid first track in play queue")
316 if starting_key:
317 await self._play_single_track(player_id, starting_key)
318 return
319
320 try:
321 first_track = await self.provider.get_track(first_track_key)
322 LOGGER.info(f"Starting playback with first track: {first_track.name}")
323
324 if first_play_queue_item_id:
325 self.play_queue_item_ids[0] = first_play_queue_item_id
326
327 await self.provider.mass.player_queues.play_media(
328 queue_id=player_id,
329 media=first_track,
330 option=QueueOption.REPLACE,
331 # Plex owns the order and a shuffled play queue already arrives
332 # shuffled, so the remaining tracks must be appended to an unshuffled
333 # queue; _load_remaining_queue_tracks sets the flag once they landed
334 shuffle=False,
335 )
336
337 if offset > 0:
338 await self._seek_to_offset_after_playback(player_id, offset)
339
340 await self._broadcast_timeline()
341
342 # Use the queue's own shuffle flag as the authoritative state,
343 # falling back to the request parameter if not set.
344 effective_shuffle = playqueue.playQueueShuffled or shuffle
345 self.provider.mass.create_task(
346 self._load_remaining_queue_tracks(
347 player_id, playqueue, selected_offset, effective_shuffle
348 )
349 )
350
351 except Exception:
352 LOGGER.exception("Error starting playback with first track")
353 if starting_key:
354 await self._play_single_track(player_id, starting_key)
355 else:
356 LOGGER.error("Play queue is empty or could not be fetched")
357 if starting_key:
358 await self._play_single_track(player_id, starting_key)
359
360 except Exception:
361 LOGGER.exception("Error playing from queue")
362 if starting_key:
363 await self._play_single_track(player_id, starting_key)
364
365 async def handle_play_media(self, request: web.Request) -> web.Response:
366 """
367 Handle playMedia command from Plex controller.
368
369 Plexamp sends various parameters:
370 - key: The item to play (track, album, playlist, etc.)
371 - containerKey: The container context (play queue)
372 - offset: Starting position in milliseconds
373 - shuffle: Whether to shuffle
374 - repeat: Repeat mode
375 """
376 self._updating_from_plex = True
377 try:
378 key = request.query.get("key")
379 container_key = request.query.get("containerKey")
380 offset = int(request.query.get("offset", 0))
381 shuffle = request.query.get("shuffle", "0") == "1"
382
383 if not key:
384 return web.Response(
385 status=400, text="Missing required 'key' parameter for playMedia command"
386 )
387
388 LOGGER.info(
389 f"Received playMedia command - key: {key}, "
390 f"containerKey: {container_key}, offset: {offset}ms"
391 )
392
393 player_id = self._ma_player_id
394 if not player_id:
395 return web.Response(status=500, text="No player assigned to this server")
396
397 if container_key and "/playQueues/" in container_key:
398 queue_id_match = re.search(r"/playQueues/(\d+)", container_key)
399 if queue_id_match:
400 self.play_queue_id = queue_id_match.group(1)
401 self.play_queue_version = 1
402 LOGGER.info(f"Playing from queue: {container_key} starting at {key}")
403 await self._play_from_plex_queue(player_id, container_key, key, shuffle, offset)
404 else:
405 self.play_queue_id = None
406 self.play_queue_item_ids = {}
407 media = await self._resolve_plex_item(key)
408 await self.provider.mass.player_queues.play_media(
409 queue_id=player_id,
410 media=media,
411 option=QueueOption.REPLACE,
412 shuffle=shuffle,
413 )
414 elif container_key:
415 self.play_queue_id = None
416 self.play_queue_item_ids = {}
417 media_to_play = await self._resolve_plex_item(container_key)
418 await self.provider.mass.player_queues.play_media(
419 queue_id=player_id,
420 media=media_to_play,
421 option=QueueOption.REPLACE,
422 shuffle=shuffle,
423 )
424 else:
425 self.play_queue_id = None
426 self.play_queue_item_ids = {}
427 media = await self._resolve_plex_item(key)
428 await self.provider.mass.player_queues.play_media(
429 queue_id=player_id,
430 media=media,
431 option=QueueOption.REPLACE,
432 shuffle=shuffle,
433 )
434
435 # Always sync shuffle state so that a previously enabled MA shuffle
436 # does not reorder an unshuffled Plex queue.
437 await self.provider.mass.player_queues.set_shuffle(player_id, shuffle)
438
439 if offset > 0:
440 await self._seek_to_offset_after_playback(player_id, offset)
441
442 await self._broadcast_timeline()
443 return web.Response(status=200)
444
445 except Exception:
446 LOGGER.exception("Error handling playMedia")
447 return web.Response(status=500, text="Internal error")
448 finally:
449 self._updating_from_plex = False
450
451 async def handle_create_play_queue(self, request: web.Request) -> web.Response:
452 """
453 Handle createPlayQueue command from Plex controller.
454
455 Creates a new play queue from a URI (album, playlist, artist tracks, etc.)
456 and optionally applies shuffle.
457 """
458 self._updating_from_plex = True
459 try:
460 uri = request.query.get("uri")
461 shuffle = request.query.get("shuffle", "0") == "1"
462 continuous = request.query.get("continuous", "0") == "1"
463
464 if not uri:
465 return web.Response(status=400, text="Missing 'uri' parameter")
466
467 LOGGER.info(f"Received createPlayQueue command - uri: {uri}, shuffle: {shuffle}")
468
469 player_id = self._ma_player_id
470 if not player_id:
471 return web.Response(status=500, text="No player assigned to this server")
472
473 def create_queue() -> PlayQueue:
474 item = self.provider._plex_server.fetchItem(uri)
475 return PlayQueue.create(
476 self.provider._plex_server,
477 item,
478 shuffle=1 if shuffle else 0,
479 continuous=1 if continuous else 0,
480 )
481
482 playqueue = await asyncio.to_thread(create_queue)
483
484 if playqueue is not None and playqueue.items:
485 self.play_queue_id = str(playqueue.playQueueID)
486 self.play_queue_version = 1
487
488 if len(playqueue.items) > MAX_QUEUE_ITEMS:
489 LOGGER.info(
490 "Capping created Plex play queue from %d to %d items",
491 len(playqueue.items),
492 MAX_QUEUE_ITEMS,
493 )
494 playqueue.__dict__["items"] = playqueue.items[:MAX_QUEUE_ITEMS]
495
496 LOGGER.info(
497 f"Created play queue {self.play_queue_id} with {len(playqueue.items)} items"
498 )
499
500 self.play_queue_item_ids = {}
501 first_item = playqueue.items[0]
502 first_track_key, first_play_queue_item_id = plex_item_fields(first_item)
503
504 if not first_track_key:
505 LOGGER.error("No valid first track in created play queue")
506 return web.Response(status=500, text="Failed to load tracks from play queue")
507
508 try:
509 first_track = await self.provider.get_track(first_track_key)
510 LOGGER.info(f"Starting playback with first track: {first_track.name}")
511
512 if first_play_queue_item_id:
513 self.play_queue_item_ids[0] = first_play_queue_item_id
514
515 await self.provider.mass.player_queues.play_media(
516 queue_id=player_id,
517 media=first_track,
518 option=QueueOption.REPLACE,
519 # as above: the tracks that follow are appended in Plex's order
520 shuffle=False,
521 )
522
523 if len(playqueue.items) > 1:
524 self.provider.mass.create_task(
525 self._load_remaining_queue_tracks(player_id, playqueue, 0, shuffle)
526 )
527
528 await self._broadcast_timeline()
529 return web.Response(status=200)
530
531 except Exception:
532 LOGGER.exception("Error starting playback with first track")
533 return web.Response(status=500, text="Failed to start playback")
534 else:
535 LOGGER.error("Failed to create play queue or queue is empty")
536 return web.Response(status=500, text="Failed to create play queue")
537
538 except Exception:
539 LOGGER.exception("Error handling createPlayQueue")
540 return web.Response(status=500, text="Internal error")
541 finally:
542 self._updating_from_plex = False
543
544 async def handle_refresh_play_queue(self, request: web.Request) -> web.Response:
545 """
546 Handle refreshPlayQueue command from Plex controller.
547
548 Called when the play queue is modified (items added, removed, reordered).
549 Syncs the updated queue state to MA while preserving current playback.
550 """
551 self._updating_from_plex = True
552 try:
553 play_queue_id = request.query.get("playQueueID")
554
555 if not play_queue_id:
556 return web.Response(status=400, text="Missing 'playQueueID' parameter")
557
558 LOGGER.info(
559 f"Received refreshPlayQueue command - playQueueID: {play_queue_id}, "
560 f"params: {dict(request.query)}"
561 )
562
563 if self.play_queue_id != play_queue_id:
564 LOGGER.warning(
565 f"Refresh requested for queue {play_queue_id} but active queue is "
566 f"{self.play_queue_id}"
567 )
568 return web.Response(
569 status=409,
570 text=(
571 f"Requested playQueueID {play_queue_id} does not match "
572 f"active queue {self.play_queue_id}"
573 ),
574 )
575
576 self.play_queue_version += 1
577
578 playqueue = await self._fetch_full_play_queue(play_queue_id)
579
580 if playqueue is None or not playqueue.items:
581 LOGGER.error("Failed to refresh play queue - queue is empty or not found")
582 return web.Response(status=404, text="Play queue not found")
583
584 player_id = self._ma_player_id
585 if not player_id:
586 LOGGER.error("No player assigned to this server")
587 return web.Response(status=500, text="No player assigned")
588
589 ma_queue = self.provider.mass.player_queues.get(player_id)
590 if not ma_queue:
591 LOGGER.error(f"MA queue not found for player {player_id}")
592 return web.Response(status=500, text="MA queue not found")
593
594 current_index = ma_queue.current_index
595 ma_queue_items = self.provider.mass.player_queues.items(player_id)
596 ma_queue_count = len(ma_queue_items) if ma_queue_items else 0
597
598 LOGGER.debug(
599 f"Queue refresh: Current index={current_index}, "
600 f"MA has {ma_queue_count} items, Plex has {len(playqueue.items)} items"
601 )
602
603 if current_index is None:
604 LOGGER.debug("No track currently playing, replacing entire queue")
605 await self._replace_entire_queue(player_id, playqueue)
606 else:
607 LOGGER.debug(
608 f"Track at index {current_index} is playing, "
609 f"replacing only items after current track"
610 )
611 await self._replace_remaining_queue(player_id, playqueue, current_index)
612
613 # Sync shuffle state from Plex to MA.
614 await self.provider.mass.player_queues.set_shuffle(
615 player_id, playqueue.playQueueShuffled
616 )
617
618 LOGGER.info(
619 f"Refreshed play queue {play_queue_id} - now has {len(playqueue.items)} items"
620 )
621
622 self._remember_synced_queue(player_id)
623
624 return web.Response(status=200)
625
626 except Exception:
627 LOGGER.exception("Error handling refreshPlayQueue")
628 return web.Response(status=500, text="Internal error")
629 finally:
630 self._updating_from_plex = False
631