/
/
/
1"""Embedded HTTP server for the MSX Bridge Provider."""
2
3from __future__ import annotations
4
5import asyncio
6import contextlib
7import functools
8import hashlib
9import io
10import json
11import logging
12import secrets
13import time
14from html import escape as html_escape
15from pathlib import Path
16from types import SimpleNamespace
17from typing import TYPE_CHECKING, Any, NamedTuple, cast
18from urllib.parse import quote, urlsplit, urlunsplit
19
20import aiohttp
21from aiohttp import WSMsgType, web
22from music_assistant_models.enums import ContentType
23from music_assistant_models.errors import InvalidProviderURI
24from music_assistant_models.media_items import AudioFormat, Track
25
26from music_assistant.constants import SENDSPIN_SERVER_PORT
27from music_assistant.controllers.streams.audio_processing import get_media_session_id
28from music_assistant.controllers.streams.constants import (
29 SINGLE_ITEM_READRATE,
30 SINGLE_ITEM_READRATE_INITIAL_BURST,
31)
32from music_assistant.controllers.webserver.helpers.auth_middleware import ImpersonatedUser
33from music_assistant.helpers.ffmpeg import get_ffmpeg_stream
34from music_assistant.helpers.uri import parse_uri
35from music_assistant.helpers.util import join_task
36
37from .constants import (
38 CONF_SHOW_STOP_NOTIFICATION,
39 DEFAULT_SHOW_STOP_NOTIFICATION,
40 MSX_PLAYER_ID_PREFIX,
41 PLAYER_ID_SANITIZE_RE,
42 PRE_BUFFER_BYTES,
43)
44from .mappers import (
45 append_device_param,
46 get_image_url,
47 map_album_to_msx,
48 map_artist_to_msx,
49 map_playlist_to_msx,
50 map_track_to_msx,
51 map_tracks_to_msx_playlist,
52)
53from .models import MsxContent, MsxItem, MsxTemplate
54from .player import MSXPlayer
55
56if TYPE_CHECKING:
57 from collections.abc import Sequence
58
59 from multidict import MultiMapping
60 from music_assistant_models.player import PlayerMedia
61
62 from music_assistant.helpers.dsp import ComplexFilter
63
64 from .provider import MSXBridgeProvider
65
66logger = logging.getLogger(__name__)
67
68STATIC_DIR = Path(__file__).parent / "static"
69
70_KNOWN_EXTENSIONS = (".mp3", ".json", ".flac", ".aac")
71
72PARTY_CACHE_TTL = 10.0
73PARTY_CALL_TIMEOUT = 5.0
74
75# The local proxy modes encode audio themselves, so they carry the core streamserver's
76# pacing ceiling rather than handing a track over as fast as ffmpeg can produce it.
77# See the usage policy note on SINGLE_ITEM_READRATE.
78_READRATE_ARGS = [
79 "-readrate",
80 SINGLE_ITEM_READRATE,
81 "-readrate_initial_burst",
82 SINGLE_ITEM_READRATE_INITIAL_BURST,
83]
84
85
86class PartyInfo(NamedTuple):
87 """Active-party details resolved from the MA Party plugin."""
88
89 join_url: str
90 name: str | None
91 qr_text: str | None
92 qr_version: str
93
94
95def _int_param(query: MultiMapping[str], name: str, default: int, max_val: int = 10000) -> int:
96 """Parse an integer query parameter safely, clamping to [0, max_val]."""
97 try:
98 return max(0, min(int(query.get(name, str(default))), max_val))
99 except ValueError, TypeError:
100 return default
101
102
103async def _is_media_item_uri(uri: str) -> bool:
104 """
105 Check that a caller-supplied uri names a media item rather than a raw stream URL.
106
107 Both spellings of a raw URL â bare, and wrapped as ``builtin://<media_type>/<url>`` â
108 resolve to the builtin provider, which would make the server fetch and play whatever
109 the caller names, so the resolved provider is what decides rather than the uri text.
110 The bridge only ever hands out uris of library or music provider items.
111 """
112 if "://" not in uri:
113 # keeps an item_id-shaped value away from parse_uri's local-file branch
114 return False
115 try:
116 _, provider_instance_id_or_domain, _ = await parse_uri(uri)
117 except InvalidProviderURI:
118 return False
119 return provider_instance_id_or_domain != "builtin"
120
121
122def _is_audio_path(path: str) -> bool:
123 """Check whether the path is one of the audio routes."""
124 return path.startswith(("/stream/", "/msx/audio/"))
125
126
127def _strip_known_extension(value: str) -> str:
128 """Strip only known audio/data extensions from a value."""
129 for ext in _KNOWN_EXTENSIONS:
130 if value.endswith(ext):
131 return value[: -len(ext)]
132 return value
133
134
135@functools.lru_cache(maxsize=4)
136def _render_qr(join_url: str, kind: str) -> bytes:
137 """
138 Render the join URL as a QR image (blocking on a miss; run in a worker thread).
139
140 Results are memoized â the output only changes when the join code rotates.
141 """
142 import segno # noqa: PLC0415 # only needed when the Party plugin is used
143
144 buf = io.BytesIO()
145 segno.make(join_url, error="m").save(buf, kind=kind, scale=8)
146 return buf.getvalue()
147
148
149def _render_qr_cover(join_url: str, cover_bytes: bytes) -> bytes:
150 """Render the QR and composite it onto the cover (blocking; run in a worker thread)."""
151 return _stamp_qr_on_cover(cover_bytes, _render_qr(join_url, "png"))
152
153
154def _stamp_qr_on_cover(cover_bytes: bytes, qr_bytes: bytes) -> bytes:
155 """Composite the QR into the cover's bottom-right corner; returns PNG bytes."""
156 from PIL import Image # noqa: PLC0415 # only needed when the Party plugin is used
157
158 cover = Image.open(io.BytesIO(cover_bytes)).convert("RGB")
159 qr = Image.open(io.BytesIO(qr_bytes)).convert("RGB")
160 # ~28% of the smaller cover side keeps the QR scannable without hiding the art;
161 # NEAREST preserves the hard module edges QR readers need.
162 side = max(48, min(cover.width, cover.height) * 28 // 100)
163 qr = qr.resize((side, side), Image.Resampling.NEAREST)
164 margin = side // 8
165 cover.paste(qr, (cover.width - side - margin, cover.height - side - margin))
166 out = io.BytesIO()
167 cover.save(out, format="PNG")
168 return out.getvalue()
169
170
171def _sort_album_tracks(tracks: list[Any]) -> list[Any]:
172 """
173 Sort album tracks deterministically.
174
175 MA sorts by (disc_number, track_number) but tracks with identical values
176 get non-deterministic ordering between calls. Adding name as a tiebreaker
177 ensures the display page and playlist endpoint always agree on track order.
178 """
179 return sorted(
180 tracks,
181 key=lambda t: (
182 getattr(t, "disc_number", 0) or 0,
183 getattr(t, "track_number", 0) or 0,
184 getattr(t, "name", "") or "",
185 ),
186 )
187
188
189class MSXHTTPServer:
190 """HTTP server that serves MSX bootstrap, library API, and stream proxy."""
191
192 def __init__(self, provider: MSXBridgeProvider, port: int) -> None:
193 """Initialize the HTTP server."""
194 self.provider = provider
195 self.port = port
196 self.app = web.Application(middlewares=[self._cors_middleware])
197 self._runner: web.AppRunner | None = None
198 self._ws_clients: dict[str, set[web.WebSocketResponse]] = {}
199 self._active_stream_tasks: dict[str, set[asyncio.Task[None]]] = {}
200 self._active_stream_transports: dict[str, set[Any]] = {}
201 self._party_cache: tuple[float, PartyInfo | None] | None = None
202 self._qr_cover_cache: dict[tuple[str, str], bytes] = {}
203 self._qr_cover_inflight: dict[tuple[str, str], asyncio.Task[bytes]] = {}
204 self._client_prefixes: dict[str, str] = {}
205 self._setup_routes()
206
207 async def start(self) -> None:
208 """Start the HTTP server."""
209 self._runner = web.AppRunner(self.app)
210 await self._runner.setup()
211 # reuse_address + reuse_port allow fast restart after reload.
212 # 0.0.0.0 is required: MSX TVs on LAN must reach this server by host IP;
213 # binding to 127.0.0.1 would prevent TV connections.
214 site = web.TCPSite(
215 self._runner,
216 "0.0.0.0",
217 self.port,
218 reuse_address=True,
219 reuse_port=True,
220 )
221 await site.start()
222 logger.info("MSX Bridge HTTP server started on port %s", self.port)
223
224 async def stop(self) -> None:
225 """Stop the HTTP server."""
226 # iterate over copies: closing a WS wakes its handler, whose cleanup
227 # discards the WS from these collections mid-iteration
228 for clients in list(self._ws_clients.values()):
229 for ws in list(clients):
230 if not ws.closed:
231 await ws.close()
232 self._ws_clients.clear()
233 for player_id in list(self._active_stream_tasks):
234 self.cancel_streams_for_player(player_id)
235 if self._runner:
236 await self._runner.cleanup()
237 self._runner = None
238 logger.info("MSX Bridge HTTP server stopped")
239
240 def broadcast_play(
241 self,
242 player_id: str,
243 *,
244 title: str | None = None,
245 artist: str | None = None,
246 image_url: str | None = None,
247 duration: int | None = None,
248 next_action: str | None = None,
249 prev_action: str | None = None,
250 ) -> None:
251 """Notify subscribed WebSocket clients to start playback with metadata."""
252 clients = self._ws_clients.get(player_id, set())
253 if not clients:
254 logger.warning(
255 "broadcast_play: no WebSocket clients for player_id=%s (connected: %s)",
256 player_id,
257 list(self._ws_clients.keys()),
258 )
259 return
260 logger.info(
261 "broadcast_play: player_id=%s, sending to %d client(s)",
262 player_id,
263 len(clients),
264 )
265
266 # We always use direct stream for maximum compatibility.
267 play_path = f"/stream/{player_id}?token={self.provider.get_stream_token(player_id)}"
268
269 payload: dict[str, Any] = {
270 "type": "play",
271 "path": play_path,
272 "player_id": player_id,
273 }
274 if title:
275 payload["title"] = title
276 if artist:
277 payload["artist"] = artist
278 if image_url:
279 # During a party the play background carries the join QR (MSX has
280 # no overlays); the endpoint falls back to the original image when
281 # the party is over, so a stale cache entry here is harmless.
282 if self._cached_party() and (client_prefix := self._client_prefixes.get(player_id)):
283 image_url = (
284 f"{client_prefix}/api/party/qr-cover.png?image={quote(image_url, safe='')}"
285 )
286 payload["image_url"] = image_url
287 if duration is not None:
288 payload["duration"] = duration
289 if next_action:
290 payload["next_action"] = next_action
291 if prev_action:
292 payload["prev_action"] = prev_action
293 msg = json.dumps(payload)
294 for ws in list(clients):
295 if not ws.closed:
296 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
297
298 def broadcast_playlist(self, player_id: str, playlist_url: str) -> None:
299 """Notify subscribed WebSocket clients to load an MSX native playlist."""
300 clients = self._ws_clients.get(player_id, set())
301 if not clients:
302 logger.warning(
303 "broadcast_playlist: no WebSocket clients for player_id=%s (connected: %s)",
304 player_id,
305 list(self._ws_clients.keys()),
306 )
307 return
308 logger.info(
309 "broadcast_playlist: player_id=%s, url=%s, sending to %d client(s)",
310 player_id,
311 playlist_url,
312 len(clients),
313 )
314 payload: dict[str, Any] = {
315 "type": "playlist",
316 "url": playlist_url,
317 "player_id": player_id,
318 }
319 msg = json.dumps(payload)
320 for ws in list(clients):
321 if not ws.closed:
322 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
323
324 def broadcast_sendspin(self, player_id: str, url: str) -> None:
325 """Notify WebSocket clients to open the Sendspin kiosk (bridge stream start)."""
326 clients = self._ws_clients.get(player_id, set())
327 if not clients:
328 logger.warning(
329 "broadcast_sendspin: no WebSocket clients for player_id=%s (connected: %s)",
330 player_id,
331 list(self._ws_clients.keys()),
332 )
333 return
334 logger.info(
335 "broadcast_sendspin: player_id=%s, url=%s, sending to %d client(s)",
336 player_id,
337 url,
338 len(clients),
339 )
340 msg = json.dumps({"type": "sendspin", "url": url, "player_id": player_id})
341 for ws in list(clients):
342 if not ws.closed:
343 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
344
345 def broadcast_goto_index(self, player_id: str, index: int) -> None:
346 """Notify subscribed WebSocket clients to jump to a playlist index."""
347 clients = self._ws_clients.get(player_id, set())
348 if not clients:
349 return
350 logger.info(
351 "broadcast_goto_index: player_id=%s, index=%d, sending to %d client(s)",
352 player_id,
353 index,
354 len(clients),
355 )
356 payload: dict[str, Any] = {"type": "goto_index", "index": index}
357 msg = json.dumps(payload)
358 for ws in list(clients):
359 if not ws.closed:
360 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
361
362 def cancel_streams_for_player(self, player_id: str) -> None:
363 """Cancel stream tasks and abort connections for the given player."""
364 tasks = self._active_stream_tasks.pop(player_id, set())
365 transports = self._active_stream_transports.pop(player_id, set())
366 for task in tasks:
367 if not task.done():
368 task.cancel()
369 for transport in transports:
370 with contextlib.suppress(Exception):
371 if transport and hasattr(transport, "abort"):
372 transport.abort()
373 if tasks or transports:
374 logger.debug(
375 "Cancelled %d task(s), aborted %d transport(s) for player %s",
376 len(tasks),
377 len(transports),
378 player_id,
379 )
380
381 def broadcast_pause(self, player_id: str) -> None:
382 """Notify subscribed WebSocket clients to pause playback."""
383 clients = self._ws_clients.get(player_id, set())
384 if not clients:
385 return
386 logger.info(
387 "broadcast_pause: player_id=%s, sending to %d client(s)",
388 player_id,
389 len(clients),
390 )
391 msg = json.dumps({"type": "pause"})
392 for ws in list(clients):
393 if not ws.closed:
394 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
395
396 def broadcast_resume(self, player_id: str) -> None:
397 """Notify subscribed WebSocket clients to resume playback."""
398 clients = self._ws_clients.get(player_id, set())
399 if not clients:
400 return
401 logger.info(
402 "broadcast_resume: player_id=%s, sending to %d client(s)",
403 player_id,
404 len(clients),
405 )
406 msg = json.dumps({"type": "resume"})
407 for ws in list(clients):
408 if not ws.closed:
409 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
410
411 def broadcast_stop(self, player_id: str) -> None:
412 """Notify subscribed WebSocket clients to stop playback."""
413 clients = self._ws_clients.get(player_id, set())
414 if not clients:
415 logger.warning(
416 "broadcast_stop: no WebSocket clients for player_id=%s (connected: %s)",
417 player_id,
418 list(self._ws_clients.keys()),
419 )
420 return
421 logger.info(
422 "broadcast_stop: player_id=%s, sending to %d client(s)",
423 player_id,
424 len(clients),
425 )
426 show_notification = self.provider.config.get_value(
427 CONF_SHOW_STOP_NOTIFICATION, DEFAULT_SHOW_STOP_NOTIFICATION
428 )
429 payload: dict[str, Any] = {
430 "type": "stop",
431 "showNotification": bool(show_notification),
432 }
433 msg = json.dumps(payload)
434 for ws in list(clients):
435 if not ws.closed:
436 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
437
438 def broadcast_seek(self, player_id: str, position_seconds: int) -> None:
439 """Notify subscribed WebSocket clients to seek to a position."""
440 clients = self._ws_clients.get(player_id, set())
441 if not clients:
442 logger.debug("broadcast_seek: no WebSocket clients for player_id=%s", player_id)
443 return
444 msg = json.dumps({"type": "seek", "position": position_seconds})
445 for ws in list(clients):
446 if not ws.closed:
447 self.provider.mass.create_task(self._ws_send(ws, msg, player_id))
448
449 def _setup_routes(self) -> None:
450 """Register all HTTP routes."""
451 self._setup_msx_routes()
452 self._setup_api_routes()
453
454 def _setup_msx_routes(self) -> None:
455 """Register MSX bootstrap, content, and playback routes."""
456 # MSX bootstrap
457 self.app.router.add_get("/", self._handle_root)
458 self.app.router.add_get("/msx/start.json", self._handle_start_json)
459 self.app.router.add_get("/msx/launcher.json", self._handle_launcher_json)
460 self.app.router.add_get("/msx/plugin.html", self._handle_msx_plugin_html)
461 self.app.router.add_get(
462 "/msx/tvx-plugin-module.min.js",
463 self._serve_static("tvx-plugin-module.min.js"),
464 )
465 self.app.router.add_get("/msx/tvx-plugin.min.js", self._serve_static("tvx-plugin.min.js"))
466 self.app.router.add_get("/msx/input.html", self._handle_msx_input_html)
467 self.app.router.add_get("/msx/input.js", self._serve_static("input.js"))
468
469 # MSX content pages (native MSX JSON navigation)
470 self.app.router.add_get("/msx/menu.json", self._handle_msx_menu)
471 self.app.router.add_get("/msx/albums.json", self._handle_msx_albums)
472 self.app.router.add_get("/msx/artists.json", self._handle_msx_artists)
473 self.app.router.add_get("/msx/playlists.json", self._handle_msx_playlists)
474 self.app.router.add_get("/msx/tracks.json", self._handle_msx_tracks)
475 self.app.router.add_get("/msx/recently-played.json", self._handle_msx_recently_played)
476 self.app.router.add_get("/msx/search-page.json", self._handle_msx_search_page)
477 self.app.router.add_get("/msx/search-input.json", self._handle_msx_search_input)
478 self.app.router.add_get("/msx/search.json", self._handle_msx_search)
479 self.app.router.add_get("/msx/party.json", self._handle_msx_party)
480
481 # MSX detail pages
482 self.app.router.add_get("/msx/albums/{item_id}/tracks.json", self._handle_msx_album_tracks)
483 self.app.router.add_get(
484 "/msx/artists/{item_id}/albums.json", self._handle_msx_artist_albums
485 )
486 self.app.router.add_get(
487 "/msx/playlists/{item_id}/tracks.json", self._handle_msx_playlist_tracks
488 )
489
490 # MSX queue playlist (MA queue â MSX native playlist)
491 self.app.router.add_get("/msx/queue-playlist/{player_id}.json", self._handle_queue_playlist)
492
493 # MSX playlist endpoints (native MSX playlist JSON)
494 self.app.router.add_get(
495 "/msx/playlist/album/{item_id}.json", self._handle_msx_album_playlist
496 )
497 self.app.router.add_get(
498 "/msx/playlist/playlist/{item_id}.json", self._handle_msx_playlist_playlist
499 )
500 self.app.router.add_get("/msx/playlist/tracks.json", self._handle_msx_tracks_playlist)
501 self.app.router.add_get(
502 "/msx/playlist/recently-played.json",
503 self._handle_msx_recently_played_playlist,
504 )
505 self.app.router.add_get("/msx/playlist/search.json", self._handle_msx_search_playlist)
506
507 # MSX audio playback
508 self.app.router.add_get("/msx/audio/{player_id}", self._handle_msx_audio)
509 self.app.router.add_get("/msx/audio/{player_id}.mp3", self._handle_msx_audio)
510
511 # Kiosk web player (browser-based, no MSX app needed)
512 self.app.router.add_get("/web", self._handle_web_app)
513 self.app.router.add_static("/web/", STATIC_DIR / "web")
514
515 # Health
516 self.app.router.add_get("/health", self._handle_health)
517
518 # WebSocket for push playback (MA -> MSX)
519 self.app.router.add_get("/ws", self._handle_ws)
520
521 # Stream proxy
522 self.app.router.add_get("/stream/{player_id}", self._handle_stream)
523 self.app.router.add_get("/stream/{player_id}.mp3", self._handle_stream)
524
525 def _setup_api_routes(self) -> None:
526 """Register Library and Playback API routes."""
527 # Library API
528 self.app.router.add_get("/api/albums", self._handle_albums)
529 self.app.router.add_get("/api/albums/{item_id}/tracks", self._handle_album_tracks)
530 self.app.router.add_get("/api/artists", self._handle_artists)
531 self.app.router.add_get("/api/artists/{item_id}/albums", self._handle_artist_albums)
532 self.app.router.add_get("/api/playlists", self._handle_playlists)
533 self.app.router.add_get("/api/playlists/{item_id}/tracks", self._handle_playlist_tracks)
534 self.app.router.add_get("/api/tracks", self._handle_tracks)
535 self.app.router.add_get("/api/search", self._handle_search)
536 self.app.router.add_get("/api/recently-played", self._handle_recently_played)
537 self.app.router.add_get("/api/lyrics/{player_id}", self._handle_lyrics)
538 self.app.router.add_get("/api/queue/{player_id}", self._handle_queue)
539 self.app.router.add_get("/api/party", self._handle_party_status)
540 self.app.router.add_get("/api/party/qr.svg", self._handle_party_qr)
541 self.app.router.add_get("/api/party/qr.png", self._handle_party_qr)
542 self.app.router.add_get("/api/party/qr-cover.png", self._handle_party_qr_cover)
543
544 # Playback control â GET (MSX interaction plugin) + POST (web player,
545 # dashboard). Never wildcard: extra methods only widen the CSRF surface.
546 self.app.router.add_post("/api/play", self._handle_play)
547 for path, handler in (
548 ("/api/pause/{player_id}", self._handle_pause),
549 ("/api/stop/{player_id}", self._handle_stop),
550 ("/api/quick-stop/{player_id}", self._handle_quick_stop),
551 ("/api/next/{player_id}", self._handle_next),
552 ("/api/previous/{player_id}", self._handle_previous),
553 ):
554 self.app.router.add_get(path, handler)
555 self.app.router.add_post(path, handler)
556
557 # --- Server Lifecycle ---
558
559 @web.middleware
560 async def _cors_middleware(self, request: web.Request, handler: Any) -> web.StreamResponse:
561 """
562 Add CORS headers to all responses.
563
564 Wildcard CORS is intentional: this server runs on LAN (default port 8099).
565 The web player (/web) and MSX plugin (/msx/plugin.html) are served from the
566 same origin, so browser playback-control POSTs are always same-origin.
567 MSX TV app only makes GET requests. This matches MA's own webserver pattern.
568
569 The audio routes are the exception and get no header at all: a media element
570 plays a cross-origin source without CORS, and the kiosk visualizer reads the
571 stream same-origin, so withholding it costs nothing and keeps a cross-origin
572 fetch() from reading the audio.
573 """
574 if request.method == "OPTIONS":
575 return web.Response(
576 headers={
577 "Access-Control-Allow-Origin": "*",
578 "Access-Control-Allow-Methods": "GET, POST, OPTIONS",
579 "Access-Control-Allow-Headers": "*",
580 }
581 )
582 response: web.StreamResponse = await handler(request)
583 if not _is_audio_path(request.path):
584 response.headers["Access-Control-Allow-Origin"] = "*"
585 return response
586
587 # --- MSX Bootstrap Routes ---
588
589 async def _handle_root(self, request: web.Request) -> web.Response:
590 """Serve status dashboard."""
591 players = self.provider.players
592 # base is derived from the Host header, so escape it before embedding in HTML
593 prefix = self._get_prefix(request)
594 base = html_escape(prefix)
595 player_rows = []
596 for p in players:
597 row = (
598 f'<li class="player-row"><span>'
599 f"{html_escape(p.display_name)} â {html_escape(p.playback_state.value)}"
600 f"</span>"
601 )
602 row += f'<form method="post" action="{base}/api/quick-stop/{html_escape(p.player_id)}" '
603 row += 'style="display:inline">'
604 row += '<button type="submit" class="btn">Quick stop</button></form></li>'
605 player_rows.append(row)
606 player_info = "".join(player_rows) if player_rows else ""
607
608 # Build URLs
609 safe_host: str = html_escape(request.host) # escape for HTML display
610 _raw_host: str = request.url.host or request.host.split(":")[0] # IPv6-safe, no port
611 hostname = f"[{_raw_host}]" if ":" in _raw_host else _raw_host
612 sendspin_url = f"http://{hostname}:{SENDSPIN_SERVER_PORT}"
613 kiosk_html5_url = f"{base}/web?kiosk=1"
614 # escape the composed URL as a whole: host-derived prefix plus & separators
615 sendspin_query = f"sendspin=1&sendspin_url={quote(sendspin_url, safe='')}"
616 sendspin_web_url = html_escape(f"{prefix}/web?{sendspin_query}")
617 sendspin_kiosk_url = html_escape(f"{prefix}/web?kiosk=1&{sendspin_query}")
618
619 html = f"""<!DOCTYPE html>
620<html>
621<head><title>MSX Bridge</title>
622<style>
623body {{ font-family: system-ui, sans-serif; max-width: 800px; margin: 50px auto; padding: 20px; }}
624.info {{ background: #e3f2fd; padding: 15px; border-radius: 5px; margin: 10px 0; }}
625.info-sendspin {{ background: #e8f5e9; }}
626code {{ background: #f5f5f5; padding: 2px 6px; border-radius: 3px; word-break: break-all; }}
627.player-row {{ display: flex; align-items: center; gap: 12px; margin: 8px 0; list-style: none; }}
628.player-row form {{ margin: 0; }}
629.btn {{ padding: 6px 12px; border-radius: 4px; border: 1px solid #1976d2;
630 background: #1976d2; color: white; cursor: pointer; font-size: 14px; }}
631.btn:hover {{ background: #1565c0; }}
632.link-row {{ margin: 8px 0; }}
633.builder-row {{ margin: 6px 0; }}
634.builder-row label {{ margin-right: 16px; cursor: pointer; }}
635.link-row a {{ color: #1976d2; text-decoration: none; }}
636.link-row a:hover {{ text-decoration: underline; }}
637small {{ color: #666; display: block; margin-top: 4px; }}
638</style>
639</head>
640<body>
641<h1>MSX Music Assistant Bridge</h1>
642
643<div class="info">
644<h3>MSX Setup URL</h3>
645<code>http://{safe_host}/msx/start.json</code>
646</div>
647
648<div class="info">
649<h3>Web Player</h3>
650<div class="link-row">
651<a href="/web">http://{safe_host}/web</a>
652<small>Browser-based player with library navigation (HTTP streaming)</small>
653</div>
654<div class="link-row">
655<a href="{kiosk_html5_url}">Kiosk Mode (HTML5)</a>
656<small>Fullscreen player with WebSocket push - ideal for dedicated displays</small>
657</div>
658</div>
659
660<div class="info info-sendspin">
661<h3>Sendspin Player (Synchronized Audio)</h3>
662<div class="link-row">
663<a href="{sendspin_web_url}">Web Player + Sendspin</a>
664<small>Library navigation with clock-synchronized audio</small>
665</div>
666<div class="link-row">
667<a href="{sendspin_kiosk_url}">Kiosk Mode (Sendspin)</a>
668<small>Fullscreen player with clock-synchronized audio</small>
669</div>
670<div class="link-row" style="margin-top: 12px;">
671<strong>Custom Sendspin URL:</strong><br>
672<code>/web?kiosk=1&sendspin=1&sendspin_url=http://<ma-server>:{SENDSPIN_SERVER_PORT}</code>
673</div>
674</div>
675
676<div class="info">
677<h3>Kiosk URL Builder</h3>
678<div id="kiosk-builder">
679<div class="builder-row">
680<label><input type="radio" name="kiosk-mode" value="html5" checked> HTML5</label>
681<label><input type="radio" name="kiosk-mode" value="sendspin"> Sendspin</label>
682</div>
683<div class="builder-row">
684<label><input type="checkbox" data-kiosk-param="controls" checked> Controls</label>
685<label><input type="checkbox" data-kiosk-param="party" checked> Party QR</label>
686<label><input type="checkbox" data-kiosk-param="viz" checked> Visualizer</label>
687<label><input type="checkbox" data-kiosk-param="lyrics" checked> Lyrics</label>
688</div>
689<div class="link-row">
690<a id="kiosk-builder-link" href="/web?kiosk=1" target="_blank">Open kiosk</a>
691</div>
692<code id="kiosk-builder-url"></code>
693</div>
694<script>
695(function () {{
696 var builder = document.getElementById('kiosk-builder');
697 var link = document.getElementById('kiosk-builder-link');
698 var urlOut = document.getElementById('kiosk-builder-url');
699
700 function rebuild() {{
701 var params = ['kiosk=1'];
702 var mode = builder.querySelector('input[name="kiosk-mode"]:checked').value;
703 if (mode === 'sendspin') {{
704 params.push('sendspin=1');
705 }}
706 var boxes = builder.querySelectorAll('input[data-kiosk-param]');
707 for (var i = 0; i < boxes.length; i++) {{
708 // only non-default choices land in the URL
709 if (!boxes[i].checked) {{
710 params.push(boxes[i].getAttribute('data-kiosk-param') + '=0');
711 }}
712 }}
713 var url = location.origin + '/web?' + params.join('&');
714 link.href = url;
715 urlOut.textContent = url;
716 }}
717
718 builder.addEventListener('change', rebuild);
719 rebuild();
720}})();
721</script>
722</div>
723
724<div class="info">
725<h3>Players</h3>
726<ul>{player_info or "<li>No players registered</li>"}</ul>
727</div>
728</body>
729</html>"""
730 return web.Response(text=html, content_type="text/html")
731
732 async def _handle_start_json(self, request: web.Request) -> web.Response:
733 """Return MSX start configuration pointing to the launcher menu."""
734 prefix = self._get_prefix(request)
735 return web.json_response(
736 {
737 "name": "Music Assistant",
738 "version": "1.0.7",
739 "parameter": f"content:{prefix}/msx/launcher.json",
740 }
741 )
742
743 async def _handle_launcher_json(self, request: web.Request) -> web.Response:
744 """Return MSX launcher page with MSX Player and Web Kiosk options."""
745 prefix = self._get_prefix(request)
746 content = MsxContent(
747 headline="Music Assistant",
748 template=MsxTemplate(
749 type="separate",
750 layout="0,0,2,4",
751 icon="msx-white-soft:music-note",
752 action="content:{context:content}",
753 ),
754 items=[
755 MsxItem(
756 label="MSX Player",
757 icon="msx-white-soft:tv",
758 action=f"menu:request:interaction:init@{prefix}/msx/plugin.html?v=8",
759 ),
760 MsxItem(
761 label="Web Kiosk",
762 icon="msx-white-soft:open-in-browser",
763 action=f"link:{prefix}/web?kiosk=1",
764 ),
765 ],
766 )
767 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
768
769 def _serve_static(self, filename: str) -> Any:
770 """Create a handler that serves a static file from the static directory."""
771 path = STATIC_DIR / filename
772
773 async def handler(_request: web.Request) -> web.FileResponse:
774 return web.FileResponse(path)
775
776 return handler
777
778 async def _handle_msx_plugin_html(self, _request: web.Request) -> web.StreamResponse:
779 """Serve plugin.html with cache-busting headers."""
780 response = cast("web.StreamResponse", web.FileResponse(STATIC_DIR / "plugin.html"))
781 response.headers["Cache-Control"] = "no-cache, no-store, must-revalidate"
782 response.headers["Pragma"] = "no-cache"
783 response.headers["Expires"] = "0"
784 return response
785
786 async def _handle_msx_input_html(self, request: web.Request) -> web.FileResponse:
787 """Serve input.html and ensure player is registered when Search is opened."""
788 await self._ensure_player_for_request(request)
789 return web.FileResponse(STATIC_DIR / "input.html")
790
791 async def _handle_web_app(self, request: web.Request) -> web.Response:
792 """Serve the web player SPA (browser-based, no MSX app needed)."""
793 response = cast("web.Response", web.FileResponse(STATIC_DIR / "web" / "index.html"))
794 response.headers["Cache-Control"] = "no-cache, no-store, must-revalidate"
795 return response
796
797 # --- MSX Content Pages (native MSX JSON) ---
798
799 async def _handle_msx_menu(self, request: web.Request) -> web.Response:
800 """Return the main library menu as an MSX content page."""
801 _, device_param, _ = await self._ensure_player_for_request(request)
802 prefix = self._get_prefix(request)
803 items = [
804 (
805 "Recently played",
806 "msx-white-soft:history",
807 f"{prefix}/msx/recently-played.json",
808 ),
809 ("Albums", "msx-white-soft:album", f"{prefix}/msx/albums.json"),
810 ("Artists", "msx-white-soft:person", f"{prefix}/msx/artists.json"),
811 (
812 "Playlists",
813 "msx-white-soft:playlist-play",
814 f"{prefix}/msx/playlists.json",
815 ),
816 ("Tracks", "msx-white-soft:audiotrack", f"{prefix}/msx/tracks.json"),
817 ("Search", "search", f"{prefix}/msx/search-page.json"),
818 ]
819 if await self._get_active_party() is not None:
820 items.append(("Party", "msx-white-soft:qr-code", f"{prefix}/msx/party.json"))
821 content = MsxContent(
822 headline="Music Assistant",
823 template=MsxTemplate(
824 type="separate",
825 layout="0,0,2,4",
826 icon="msx-white-soft:music-note",
827 action="content:{context:content}",
828 ),
829 items=[
830 MsxItem(
831 label=label,
832 icon=icon,
833 content=append_device_param(url, device_param),
834 )
835 for label, icon, url in items
836 ],
837 )
838 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
839
840 async def _handle_msx_albums(self, request: web.Request) -> web.Response:
841 """Return albums as an MSX content page."""
842 _, device_param, _ = await self._ensure_player_for_request(request)
843 prefix = self._get_prefix(request)
844 limit = _int_param(request.query, "limit", 50)
845 offset = _int_param(request.query, "offset", 0)
846 try:
847 albums = await asyncio.wait_for(
848 self.provider.mass.music.albums.library_items(
849 limit=limit, offset=offset, summary=False
850 ),
851 timeout=10.0,
852 )
853 except Exception:
854 logger.exception("Failed to fetch albums")
855 albums = []
856
857 items = await asyncio.gather(
858 *(map_album_to_msx(a, prefix, self.provider, device_param) for a in albums)
859 )
860 content = MsxContent(
861 headline="Albums",
862 template=MsxTemplate(
863 type="separate",
864 layout="0,0,3,4",
865 color="msx-glass",
866 ),
867 items=items if items else [MsxItem(title="No albums found")],
868 )
869 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
870
871 async def _handle_msx_artists(self, request: web.Request) -> web.Response:
872 """Return artists as an MSX content page."""
873 _, device_param, _ = await self._ensure_player_for_request(request)
874 prefix = self._get_prefix(request)
875 limit = _int_param(request.query, "limit", 50)
876 offset = _int_param(request.query, "offset", 0)
877 try:
878 artists = await asyncio.wait_for(
879 self.provider.mass.music.artists.library_items(
880 limit=limit, offset=offset, summary=False
881 ),
882 timeout=10.0,
883 )
884 except Exception:
885 logger.exception("Failed to fetch artists")
886 artists = []
887
888 items = [map_artist_to_msx(a, prefix, self.provider, device_param) for a in artists]
889 content = MsxContent(
890 headline="Artists",
891 template=MsxTemplate(
892 type="separate",
893 layout="0,0,2,3",
894 color="msx-glass",
895 ),
896 items=items if items else [MsxItem(title="No artists found")],
897 )
898 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
899
900 async def _handle_msx_playlists(self, request: web.Request) -> web.Response:
901 """Return playlists as an MSX content page."""
902 _, device_param, _ = await self._ensure_player_for_request(request)
903 prefix = self._get_prefix(request)
904 limit = _int_param(request.query, "limit", 50)
905 offset = _int_param(request.query, "offset", 0)
906 try:
907 playlists = await asyncio.wait_for(
908 self.provider.mass.music.playlists.library_items(
909 limit=limit, offset=offset, summary=False
910 ),
911 timeout=10.0,
912 )
913 except Exception:
914 logger.exception("Failed to fetch playlists")
915 playlists = []
916
917 items = [map_playlist_to_msx(p, prefix, self.provider, device_param) for p in playlists]
918 content = MsxContent(
919 headline="Playlists",
920 template=MsxTemplate(
921 type="separate",
922 layout="0,0,3,4",
923 color="msx-glass",
924 ),
925 items=items if items else [MsxItem(title="No playlists found")],
926 )
927 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
928
929 async def _handle_msx_tracks(self, request: web.Request) -> web.Response:
930 """Return tracks as an MSX content page."""
931 player_id, device_param, _ = await self._ensure_player_for_request(request)
932 prefix = self._get_prefix(request)
933 limit = _int_param(request.query, "limit", 50)
934 offset = _int_param(request.query, "offset", 0)
935 try:
936 tracks = await asyncio.wait_for(
937 self.provider.mass.music.tracks.library_items(
938 limit=limit, offset=offset, summary=False
939 ),
940 timeout=10.0,
941 )
942 except Exception:
943 logger.exception("Failed to fetch tracks")
944 tracks = []
945
946 playlist_base = f"{prefix}/msx/playlist/tracks.json?limit={limit}&offset={offset}"
947 playlist_base = append_device_param(playlist_base, device_param)
948 items = [
949 map_track_to_msx(
950 t,
951 prefix,
952 player_id,
953 self.provider,
954 device_param,
955 playlist_url=f"{playlist_base}&start={idx}",
956 )
957 for idx, t in enumerate(tracks)
958 ]
959 content = MsxContent(
960 headline="Tracks",
961 template=MsxTemplate(
962 type="default",
963 layout="0,0,6,1",
964 image_width=0.83,
965 color="msx-glass",
966 ),
967 items=items if items else [MsxItem(title="No tracks found")],
968 )
969 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
970
971 async def _handle_msx_recently_played(self, request: web.Request) -> web.Response:
972 """Return recently played tracks as an MSX content page."""
973 player_id, device_param, _ = await self._ensure_player_for_request(request)
974 prefix = self._get_prefix(request)
975 try:
976 tracks = await asyncio.wait_for(
977 self.provider.mass.music.tracks.library_items(
978 limit=50, order_by="last_played", summary=False
979 ),
980 timeout=10.0,
981 )
982 except Exception:
983 logger.exception("Failed to fetch recently played tracks")
984 tracks = []
985 playlist_base = f"{prefix}/msx/playlist/recently-played.json"
986 playlist_base = append_device_param(playlist_base, device_param)
987 items = [
988 map_track_to_msx(
989 t,
990 prefix,
991 player_id,
992 self.provider,
993 device_param,
994 playlist_url=f"{playlist_base}{'&' if '?' in playlist_base else '?'}start={idx}",
995 )
996 for idx, t in enumerate(tracks)
997 ]
998 content = MsxContent(
999 headline="Recently played",
1000 template=MsxTemplate(
1001 type="default",
1002 layout="0,0,6,1",
1003 image_width=0.83,
1004 color="msx-glass",
1005 ),
1006 items=items if items else [MsxItem(title="No recently played tracks")],
1007 )
1008 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1009
1010 async def _handle_msx_search_page(self, request: web.Request) -> web.Response:
1011 """Return a content page whose page-level action launches the Input Plugin keyboard."""
1012 _, device_param, _ = await self._ensure_player_for_request(request)
1013 prefix = self._get_prefix(request)
1014 search_url = append_device_param(
1015 f"{prefix}/msx/search-input.json?q={{INPUT}}", device_param
1016 )
1017 action = (
1018 f"content:request:interaction:"
1019 f"{search_url}"
1020 f"|search:3|en|Search Music||||Search..."
1021 f"@{prefix}/msx/input.html"
1022 )
1023 content = MsxContent(
1024 headline="Search",
1025 action=action,
1026 template=MsxTemplate(
1027 type="separate",
1028 layout="0,0,2,4",
1029 ),
1030 items=[
1031 MsxItem(
1032 title="Search Music",
1033 title_footer="Press OK to open keyboard",
1034 icon="search",
1035 action=action,
1036 )
1037 ],
1038 )
1039 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1040
1041 async def _handle_msx_search_input(self, request: web.Request) -> web.Response:
1042 """Return search results for the MSX Input Plugin (search keyboard)."""
1043 player_id, device_param, _ = await self._ensure_player_for_request(request)
1044 prefix = self._get_prefix(request)
1045 query = request.query.get("q", "")
1046 if not query:
1047 content = MsxContent(
1048 headline="{ico:search} Search",
1049 hint="Type to search...",
1050 template=MsxTemplate(
1051 type="separate",
1052 layout="0,0,2,4",
1053 image_filler="default",
1054 ),
1055 items=[MsxItem(title="Start typing to search")],
1056 )
1057 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1058
1059 limit = _int_param(request.query, "limit", 20)
1060 items = await self._build_search_items(
1061 query,
1062 limit,
1063 player_id,
1064 device_param,
1065 prefix,
1066 )
1067
1068 content = MsxContent(
1069 headline=f'{{ico:search}} "{query}"',
1070 hint=f"Found {len(items)} items",
1071 template=MsxTemplate(
1072 type="separate",
1073 layout="0,0,2,4",
1074 image_filler="default",
1075 ),
1076 items=items if items else [MsxItem(title="No results found")],
1077 )
1078 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1079
1080 async def _handle_msx_search(self, request: web.Request) -> web.Response:
1081 """Return search results as an MSX content page."""
1082 player_id, device_param, _ = await self._ensure_player_for_request(request)
1083 prefix = self._get_prefix(request)
1084 query = request.query.get("q", "")
1085 if not query:
1086 return web.json_response(
1087 MsxContent(
1088 headline="Search",
1089 items=[MsxItem(title="Please enter a search query")],
1090 ).model_dump(by_alias=True, exclude_none=True)
1091 )
1092
1093 limit = _int_param(request.query, "limit", 20)
1094 items = await self._build_search_items(
1095 query,
1096 limit,
1097 player_id,
1098 device_param,
1099 prefix,
1100 )
1101
1102 content = MsxContent(
1103 headline=f"Search: {query}",
1104 template=MsxTemplate(
1105 type="separate",
1106 layout="0,0,2,4",
1107 image_filler="default",
1108 ),
1109 items=items if items else [MsxItem(title="No results found")],
1110 )
1111 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1112
1113 async def _handle_msx_party(self, request: web.Request) -> web.Response:
1114 """Return MSX page with the party QR code, or a hint when no party is active."""
1115 await self._ensure_player_for_request(request)
1116 prefix = self._get_prefix(request)
1117 party = await self._get_active_party()
1118 if party is None:
1119 item = MsxItem(
1120 title="No active party",
1121 label="Enable guest access in the Music Assistant Party plugin",
1122 )
1123 else:
1124 # PNG, not SVG: MSX image slots on older TV engines cannot decode SVG
1125 item = MsxItem(
1126 image=f"{prefix}/api/party/qr.png",
1127 label=party.qr_text or "Scan to join the party",
1128 )
1129 content = MsxContent(
1130 headline=(party.name if party else None) or "Party",
1131 template=MsxTemplate(type="separate", layout="0,0,4,4"),
1132 items=[item],
1133 )
1134 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1135
1136 async def _build_search_items(
1137 self,
1138 query: str,
1139 limit: int,
1140 player_id: str,
1141 device_param: str,
1142 prefix: str,
1143 ) -> list[MsxItem]:
1144 """Build MSX items from search results (shared by search handlers)."""
1145 results = await self.provider.mass.music.search(query, limit=limit)
1146 items: list[MsxItem] = []
1147 for artist in results.artists:
1148 item = map_artist_to_msx(artist, prefix, self.provider, device_param)
1149 item.label = "Artist"
1150 item.icon = "msx-white-soft:person"
1151 items.append(item)
1152 for album in results.albums:
1153 item = await map_album_to_msx(album, prefix, self.provider, device_param)
1154 item.label = f"Album â {getattr(album, 'artist_str', '')}"
1155 item.icon = "msx-white-soft:album"
1156 items.append(item)
1157 playlist_base = f"{prefix}/msx/playlist/search.json?q={quote(query, safe='')}"
1158 playlist_base = append_device_param(playlist_base, device_param)
1159 for idx, track in enumerate(results.tracks):
1160 item = map_track_to_msx(
1161 track,
1162 prefix,
1163 player_id,
1164 self.provider,
1165 device_param,
1166 playlist_url=f"{playlist_base}&start={idx}",
1167 )
1168 item.label = f"Track â {getattr(track, 'artist_str', '')}"
1169 item.icon = "msx-white-soft:audiotrack"
1170 items.append(item)
1171 return items
1172
1173 # --- MSX Detail Pages ---
1174
1175 async def _handle_msx_album_tracks(self, request: web.Request) -> web.Response:
1176 """Return tracks for an album as an MSX content page."""
1177 player_id, device_param, _ = await self._ensure_player_for_request(request)
1178 prefix = self._get_prefix(request)
1179 item_id = request.match_info["item_id"]
1180 provider = request.query.get("provider", "library")
1181 try:
1182 tracks = _sort_album_tracks(
1183 await self.provider.mass.music.albums.tracks(item_id, provider)
1184 )
1185 except Exception:
1186 logger.exception("Failed to fetch tracks for album %s", item_id)
1187 tracks = []
1188 playlist_base = f"{prefix}/msx/playlist/album/{item_id}.json?provider={provider}"
1189 playlist_base = append_device_param(playlist_base, device_param)
1190 items = [
1191 map_track_to_msx(
1192 t,
1193 prefix,
1194 player_id,
1195 self.provider,
1196 device_param,
1197 playlist_url=f"{playlist_base}&start={idx}",
1198 )
1199 for idx, t in enumerate(tracks)
1200 ]
1201 content = MsxContent(
1202 headline="Album Tracks",
1203 template=MsxTemplate(
1204 type="default",
1205 layout="0,0,6,1",
1206 image_width=0.83,
1207 color="msx-glass",
1208 ),
1209 items=items if items else [MsxItem(title="No tracks found")],
1210 )
1211 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1212
1213 async def _handle_msx_artist_albums(self, request: web.Request) -> web.Response:
1214 """Return albums for an artist as an MSX content page."""
1215 _, device_param, _ = await self._ensure_player_for_request(request)
1216 prefix = self._get_prefix(request)
1217 item_id = request.match_info["item_id"]
1218 try:
1219 albums = await self.provider.mass.music.artists.albums(item_id, "library")
1220 except Exception:
1221 logger.exception("Failed to fetch albums for artist %s", item_id)
1222 albums = []
1223
1224 items = await asyncio.gather(
1225 *(map_album_to_msx(a, prefix, self.provider, device_param) for a in albums)
1226 )
1227 content = MsxContent(
1228 headline="Artist Albums",
1229 template=MsxTemplate(
1230 type="default",
1231 layout="0,0,6,2",
1232 image_width=1.5,
1233 color="msx-glass",
1234 ),
1235 items=items if items else [MsxItem(title="No albums found")],
1236 )
1237 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1238
1239 async def _handle_msx_playlist_tracks(self, request: web.Request) -> web.Response:
1240 """Return tracks for a playlist as an MSX content page."""
1241 player_id, device_param, _ = await self._ensure_player_for_request(request)
1242 prefix = self._get_prefix(request)
1243 item_id = request.match_info["item_id"]
1244 try:
1245 tracks = [
1246 t async for t in self.provider.mass.music.playlists.tracks(item_id, "library")
1247 ]
1248 except Exception:
1249 logger.exception("Failed to fetch tracks for playlist %s", item_id)
1250 tracks = []
1251 playlist_base = f"{prefix}/msx/playlist/playlist/{item_id}.json"
1252 playlist_base = append_device_param(playlist_base, device_param)
1253 items = [
1254 map_track_to_msx(
1255 t,
1256 prefix,
1257 player_id,
1258 self.provider,
1259 device_param,
1260 playlist_url=f"{playlist_base}{'&' if '?' in playlist_base else '?'}start={idx}",
1261 )
1262 for idx, t in enumerate(tracks)
1263 ]
1264 content = MsxContent(
1265 headline="Playlist Tracks",
1266 template=MsxTemplate(
1267 type="default",
1268 layout="0,0,6,1",
1269 image_width=0.83,
1270 color="msx-glass",
1271 ),
1272 items=items if items else [MsxItem(title="No tracks found")],
1273 )
1274 return web.json_response(content.model_dump(by_alias=True, exclude_none=True))
1275
1276 # --- MSX Playlist Endpoints ---
1277
1278 async def _handle_msx_album_playlist(self, request: web.Request) -> web.Response:
1279 """Return album tracks as an MSX playlist JSON."""
1280 player_id, device_param, _ = await self._ensure_player_for_request(request)
1281 prefix = self._get_prefix(request)
1282 item_id = request.match_info["item_id"]
1283 provider_name = request.query.get("provider", "library")
1284 start = _int_param(request.query, "start", 0)
1285 try:
1286 tracks = _sort_album_tracks(
1287 await self.provider.mass.music.albums.tracks(item_id, provider_name)
1288 )
1289 except Exception:
1290 logger.exception("Failed to fetch tracks for album playlist %s", item_id)
1291 tracks = []
1292 playlist = map_tracks_to_msx_playlist(
1293 tracks,
1294 start,
1295 prefix,
1296 player_id,
1297 self.provider,
1298 device_param,
1299 qr_cover_base=await self._qr_cover_base(prefix),
1300 )
1301 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1302
1303 async def _handle_msx_playlist_playlist(self, request: web.Request) -> web.Response:
1304 """Return playlist tracks as an MSX playlist JSON."""
1305 player_id, device_param, _ = await self._ensure_player_for_request(request)
1306 prefix = self._get_prefix(request)
1307 item_id = request.match_info["item_id"]
1308 start = _int_param(request.query, "start", 0)
1309 try:
1310 tracks = [
1311 t async for t in self.provider.mass.music.playlists.tracks(item_id, "library")
1312 ]
1313 except Exception:
1314 logger.exception("Failed to fetch tracks for playlist playlist %s", item_id)
1315 tracks = []
1316 playlist = map_tracks_to_msx_playlist(
1317 tracks,
1318 start,
1319 prefix,
1320 player_id,
1321 self.provider,
1322 device_param,
1323 qr_cover_base=await self._qr_cover_base(prefix),
1324 )
1325 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1326
1327 async def _handle_msx_tracks_playlist(self, request: web.Request) -> web.Response:
1328 """Return library tracks as an MSX playlist JSON."""
1329 player_id, device_param, _ = await self._ensure_player_for_request(request)
1330 prefix = self._get_prefix(request)
1331 limit = _int_param(request.query, "limit", 50)
1332 offset = _int_param(request.query, "offset", 0)
1333 start = _int_param(request.query, "start", 0)
1334 tracks = await self.provider.mass.music.tracks.library_items(
1335 limit=limit, offset=offset, summary=False
1336 )
1337 playlist = map_tracks_to_msx_playlist(
1338 list(tracks),
1339 start,
1340 prefix,
1341 player_id,
1342 self.provider,
1343 device_param,
1344 qr_cover_base=await self._qr_cover_base(prefix),
1345 )
1346 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1347
1348 async def _handle_msx_recently_played_playlist(self, request: web.Request) -> web.Response:
1349 """Return recently played tracks as an MSX playlist JSON."""
1350 player_id, device_param, _ = await self._ensure_player_for_request(request)
1351 prefix = self._get_prefix(request)
1352 start = _int_param(request.query, "start", 0)
1353 tracks = await self.provider.mass.music.tracks.library_items(
1354 limit=50, order_by="last_played", summary=False
1355 )
1356 playlist = map_tracks_to_msx_playlist(
1357 list(tracks),
1358 start,
1359 prefix,
1360 player_id,
1361 self.provider,
1362 device_param,
1363 qr_cover_base=await self._qr_cover_base(prefix),
1364 )
1365 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1366
1367 async def _handle_msx_search_playlist(self, request: web.Request) -> web.Response:
1368 """Return search track results as an MSX playlist JSON."""
1369 player_id, device_param, _ = await self._ensure_player_for_request(request)
1370 prefix = self._get_prefix(request)
1371 query = request.query.get("q", "")
1372 start = _int_param(request.query, "start", 0)
1373 if not query:
1374 return web.json_response(
1375 MsxContent(items=[]).model_dump(by_alias=True, exclude_none=True)
1376 )
1377 limit = _int_param(request.query, "limit", 20)
1378 results = await self.provider.mass.music.search(query, limit=limit)
1379 playlist = map_tracks_to_msx_playlist(
1380 list(results.tracks),
1381 start,
1382 prefix,
1383 player_id,
1384 self.provider,
1385 device_param,
1386 qr_cover_base=await self._qr_cover_base(prefix),
1387 )
1388 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1389
1390 # --- MSX Queue Playlist ---
1391
1392 async def _handle_queue_playlist(self, request: web.Request) -> web.Response:
1393 """Return the current MA queue as an MSX native playlist."""
1394 _, device_param, _ = await self._ensure_player_for_request(request)
1395 prefix = self._get_prefix(request)
1396 player_id = request.match_info["player_id"]
1397 queue_id = request.query.get("queue_id", player_id)
1398 start = _int_param(request.query, "start", 0)
1399
1400 try:
1401 queue_items = self.provider.mass.player_queues.items(queue_id)
1402 except Exception:
1403 logger.exception("Failed to fetch queue items for %s", player_id)
1404 queue_items = []
1405
1406 # Convert QueueItems to track-like objects for map_tracks_to_msx_playlist
1407 tracks: list[Any] = []
1408 for qi in queue_items:
1409 mi = getattr(qi, "media_item", None)
1410 tracks.append(
1411 SimpleNamespace(
1412 name=getattr(mi, "name", None) or getattr(qi, "name", "") or "",
1413 uri=getattr(mi, "uri", None) or "",
1414 duration=getattr(mi, "duration", None) or getattr(qi, "duration", 0) or 0,
1415 artist_str=getattr(mi, "artist_str", "") if mi else "",
1416 image=getattr(qi, "image", None),
1417 )
1418 )
1419
1420 playlist = map_tracks_to_msx_playlist(
1421 tracks,
1422 start,
1423 prefix,
1424 player_id,
1425 self.provider,
1426 device_param,
1427 qr_cover_base=await self._qr_cover_base(prefix),
1428 )
1429 return web.json_response(playlist.model_dump(by_alias=True, exclude_none=True))
1430
1431 # --- MSX Audio Playback ---
1432
1433 async def _handle_msx_audio(self, request: web.Request) -> web.StreamResponse:
1434 """Trigger playback via MA queue and stream audio to MSX."""
1435 player_id = _strip_known_extension(request.match_info["player_id"])
1436
1437 uri = request.query.get("uri")
1438 if not uri or not await _is_media_item_uri(uri):
1439 return web.Response(status=400, text="Invalid uri parameter")
1440
1441 from_playlist = request.query.get("from_playlist") == "1"
1442
1443 player = self.provider.mass.players.get_player(player_id)
1444 if not player or not isinstance(player, MSXPlayer):
1445 return web.Response(status=404, text="Player not found")
1446 if rejected := self._reject_invalid_stream_token(request, player_id):
1447 return rejected
1448 self.provider.on_player_activity(player_id)
1449
1450 # When MA is driving the queue (next/prev from MA UI), current_media is
1451 # already set by player.play_media() before the WS goto_index reaches MSX.
1452 # Re-enqueuing would recreate the queue from the track URI, destroying it.
1453 # We verify by checking that current_media's queue item URI matches the
1454 # requested track URI â if not, MSX auto-advanced and we must re-enqueue.
1455 if (
1456 from_playlist
1457 and player._playing_from_queue
1458 and self._current_media_matches_uri(player, uri)
1459 ):
1460 logger.debug("Queue-driven: using current_media for %s", uri)
1461 media = player.current_media
1462 else:
1463 # Suppress WS broadcast when called from MSX playlist to avoid conflicts
1464 if from_playlist:
1465 player._skip_ws_notify = True
1466
1467 # Arm BEFORE enqueuing so wait_for_media() waits for the new track's
1468 # play_media() instead of returning the previous track's media.
1469 player.expect_new_media()
1470 try:
1471 async with ImpersonatedUser(
1472 self.provider.mass, await self.provider.get_owner_username()
1473 ):
1474 await self.provider.mass.player_queues.play_media(player_id, uri)
1475 finally:
1476 if from_playlist:
1477 player._skip_ws_notify = False
1478
1479 # Wait for play_media() to signal media is ready (replaces 10s polling loop)
1480 media = await player.wait_for_media(timeout=10.0)
1481
1482 if not media:
1483 return web.Response(status=504, text="Playback setup timeout")
1484
1485 return await self._serve_audio_stream(
1486 request,
1487 player,
1488 media,
1489 duration=self._resolve_served_duration(media),
1490 )
1491
1492 # --- Audio Streaming Infrastructure ---
1493
1494 def _resolve_served_duration(self, media: PlayerMedia) -> int:
1495 """
1496 Return the length in seconds of the audio served for the given media, or 0 if unknown.
1497
1498 This is what the Content-Length header is derived from, so it describes
1499 the audio we actually serve rather than the media item: starting
1500 playback at a seek position yields a shorter stream.
1501
1502 :param media: The media being served.
1503 """
1504 duration = media.stream_duration or media.duration or 0
1505 if not duration and media.source_id and media.queue_item_id:
1506 queue_item = self.provider.mass.player_queues.get_item(
1507 media.source_id, media.queue_item_id
1508 )
1509 if queue_item:
1510 if queue_item.media_item:
1511 duration = getattr(queue_item.media_item, "duration", None) or duration
1512 if not duration and queue_item.duration:
1513 duration = queue_item.duration
1514 return int(duration)
1515
1516 @staticmethod
1517 def _build_audio_params(
1518 output_format_str: str, duration: int
1519 ) -> tuple[AudioFormat, AudioFormat, dict[str, str]]:
1520 """Build PCM input format, encoded output format, and HTTP headers."""
1521 pcm_format = AudioFormat(
1522 content_type=ContentType.PCM_S16LE,
1523 sample_rate=44100,
1524 bit_depth=16,
1525 channels=2,
1526 )
1527 content_type_map: dict[str, tuple[ContentType, str]] = {
1528 "mp3": (ContentType.MP3, "audio/mpeg"),
1529 "aac": (ContentType.AAC, "audio/aac"),
1530 "flac": (ContentType.FLAC, "audio/flac"),
1531 }
1532 codec, mime_type = content_type_map.get(output_format_str, (ContentType.MP3, "audio/mpeg"))
1533 out_format = AudioFormat(
1534 content_type=codec,
1535 sample_rate=44100,
1536 bit_depth=16,
1537 channels=2,
1538 )
1539 bitrate_map = {"mp3": 40_000, "aac": 32_000}
1540 bytes_per_sec = bitrate_map.get(output_format_str, 0)
1541 headers: dict[str, str] = {
1542 "Content-Type": mime_type,
1543 "Cache-Control": "no-cache",
1544 "Connection": "keep-alive",
1545 "Accept-Ranges": "none",
1546 }
1547 if duration and bytes_per_sec:
1548 capped_duration = min(float(duration), 43200) # cap at 12h
1549 headers["Content-Length"] = str(int(capped_duration * bytes_per_sec))
1550 return pcm_format, out_format, headers
1551
1552 async def _serve_audio_stream(
1553 self,
1554 request: web.Request,
1555 player: MSXPlayer,
1556 media: Any,
1557 duration: int = 0,
1558 ) -> web.StreamResponse:
1559 """
1560 Unified method to stream audio from MA to MSX via ffmpeg.
1561
1562 Supports three modes based on provider configuration:
1563 1. Independent (default): Each player gets its own ffmpeg stream
1564 2. Shared Buffer: Group members share one ffmpeg process via SharedGroupStream
1565 3. MA Redirect: 302 redirect to MA Streamserver (requires MA 2.6+)
1566
1567 Pre-buffers audio data before sending HTTP headers so MSX receives
1568 the response and initial audio burst simultaneously, preventing
1569 stutter/restart from an empty initial buffer.
1570 """
1571 player_id = player.player_id
1572
1573 # --- Mode 1: MA Redirect ---
1574 if self.provider.is_redirect_stream_mode():
1575 redirect_url = await self.provider.get_ma_stream_url(player_id, media)
1576 if redirect_url:
1577 redirect_url = self._rewrite_stream_host(request, redirect_url)
1578 logger.info(
1579 "[StreamMode:redirect] Player %s -> MA Streamserver: %s",
1580 player_id,
1581 redirect_url,
1582 )
1583 raise web.HTTPFound(location=redirect_url)
1584 # Fallback to independent mode if redirect fails
1585 logger.warning(
1586 "[StreamMode:redirect] Failed to get MA URL for %s, "
1587 "falling back to independent mode",
1588 player_id,
1589 )
1590
1591 # Resolve effective output format: per-player config overrides provider default.
1592 # CONF_ENTRY_OUTPUT_CODEC_DEFAULT_MP3 uses key "output_codec"; fall back to
1593 # player.output_format (set from provider-level config during registration).
1594 # Only the proxy paths below need this â in redirect mode the MA streamserver
1595 # applies the same per-player codec config itself.
1596 effective_format = cast(
1597 "str",
1598 player.config.get_value("output_codec", player.output_format),
1599 )
1600
1601 pcm_format, out_format, headers = self._build_audio_params(
1602 effective_format,
1603 duration,
1604 )
1605
1606 # --- Mode 2: Shared Buffer (for groups) ---
1607 group_id = self.provider.get_group_id_for_player(player)
1608 if group_id and self.provider.is_shared_stream_mode():
1609 logger.info(
1610 "[StreamMode:shared] Player %s in group %s, using shared stream",
1611 player_id,
1612 group_id,
1613 )
1614 return await self._serve_shared_stream(
1615 request, player, media, group_id, pcm_format, out_format, headers
1616 )
1617
1618 # --- Mode 3: Independent (default) ---
1619 logger.debug(
1620 "[StreamMode:independent] Serving audio %s: format=%s, duration=%s",
1621 player_id,
1622 effective_format,
1623 duration,
1624 )
1625
1626 audio_source = self.provider.mass.streams.get_stream(
1627 media,
1628 pcm_format,
1629 force_flow_mode=False,
1630 )
1631 output_plan = self.provider.mass.streams.audio.get_player_output_plan(
1632 player_id,
1633 pcm_format,
1634 out_format,
1635 queue_id=getattr(media, "source_id", None),
1636 session_id=get_media_session_id(media),
1637 queue_item_id=getattr(media, "queue_item_id", None),
1638 )
1639
1640 response = web.StreamResponse(status=200, headers=headers)
1641 stream_task: asyncio.Task[None] = asyncio.create_task(
1642 self._stream_with_prebuffer(
1643 request,
1644 response,
1645 player,
1646 headers,
1647 audio_source,
1648 pcm_format,
1649 out_format,
1650 output_plan.filter_params,
1651 )
1652 )
1653 transport = getattr(request, "transport", None)
1654 await self._run_stream_task(player_id, stream_task, transport)
1655
1656 return response
1657
1658 async def _serve_shared_stream(
1659 self,
1660 request: web.Request,
1661 player: MSXPlayer,
1662 media: Any,
1663 group_id: str,
1664 pcm_format: AudioFormat,
1665 out_format: AudioFormat,
1666 headers: dict[str, str],
1667 ) -> web.StreamResponse:
1668 """
1669 Serve audio from a shared group stream.
1670
1671 Multiple players in a group read from the same SharedGroupStream,
1672 which has a single ffmpeg producer.
1673 """
1674 player_id = player.player_id
1675 media_uri = getattr(media, "uri", "") or str(media)
1676
1677 # Check if we need to create a new shared stream (leader creates it)
1678 existing_stream = self.provider._shared_streams.get(group_id)
1679 is_leader = player_id == group_id
1680
1681 if existing_stream and not existing_stream.finished:
1682 # Reuse existing stream
1683 logger.debug(
1684 "[SharedStream] Player %s subscribing to existing stream for group %s",
1685 player_id,
1686 group_id,
1687 )
1688 shared_stream = existing_stream
1689 elif is_leader:
1690 # Leader creates the shared stream
1691 logger.info(
1692 "[SharedStream] Leader %s creating shared stream for group %s",
1693 player_id,
1694 group_id,
1695 )
1696 audio_source = self.provider.mass.streams.get_stream(
1697 media,
1698 pcm_format,
1699 force_flow_mode=False,
1700 )
1701 output_plan = self.provider.mass.streams.audio.get_player_output_plan(
1702 player_id,
1703 pcm_format,
1704 out_format,
1705 queue_id=getattr(media, "source_id", None),
1706 session_id=get_media_session_id(media),
1707 queue_item_id=getattr(media, "queue_item_id", None),
1708 )
1709 # Create ffmpeg chunk generator
1710 audio_chunks = get_ffmpeg_stream(
1711 audio_input=audio_source,
1712 input_format=pcm_format,
1713 output_format=out_format,
1714 filter_params=output_plan.filter_params,
1715 extra_input_args=_READRATE_ARGS,
1716 )
1717 shared_stream = await self.provider.get_or_create_shared_stream(
1718 group_id, media_uri, audio_chunks
1719 )
1720 shared_stream.output_plan = output_plan
1721 else:
1722 # Member but no existing stream - wait briefly for leader
1723 logger.info(
1724 "[SharedStream] Member %s waiting for leader to create stream for group %s",
1725 player_id,
1726 group_id,
1727 )
1728 for _ in range(30): # Wait up to 3 seconds
1729 await asyncio.sleep(0.1)
1730 existing_stream = self.provider._shared_streams.get(group_id)
1731 if existing_stream and not existing_stream.finished:
1732 shared_stream = existing_stream
1733 break
1734 else:
1735 # Timeout - fallback to independent stream
1736 logger.warning(
1737 "[SharedStream] Timeout waiting for leader stream, "
1738 "falling back to independent for %s",
1739 player_id,
1740 )
1741 return await self._serve_independent_stream(
1742 request, player, media, pcm_format, out_format, headers
1743 )
1744
1745 queue_id = getattr(media, "source_id", None)
1746 session_id = get_media_session_id(media)
1747 if (
1748 shared_stream.output_plan is not None
1749 and queue_id is not None
1750 and session_id is not None
1751 ):
1752 self.provider.mass.streams.audio_processing.update_output(
1753 player_id,
1754 shared_stream.output_plan,
1755 queue_id=queue_id,
1756 session_id=session_id,
1757 queue_item_id=getattr(media, "queue_item_id", None),
1758 )
1759
1760 # Subscribe to shared stream
1761 response = web.StreamResponse(status=200, headers=headers)
1762 await response.prepare(request)
1763
1764 total_bytes = 0
1765 try:
1766 async for chunk in shared_stream.subscribe(player_id):
1767 await response.write(chunk)
1768 total_bytes += len(chunk)
1769 except ConnectionResetError, BrokenPipeError, ConnectionAbortedError:
1770 logger.debug(
1771 "[SharedStream] Client %s disconnected after %d bytes",
1772 player_id,
1773 total_bytes,
1774 )
1775 except asyncio.CancelledError:
1776 logger.debug("[SharedStream] Stream cancelled for %s", player_id)
1777 raise
1778
1779 logger.info(
1780 "[SharedStream] Player %s finished, wrote %d bytes",
1781 player_id,
1782 total_bytes,
1783 )
1784 return response
1785
1786 async def _serve_independent_stream(
1787 self,
1788 request: web.Request,
1789 player: MSXPlayer,
1790 media: Any,
1791 pcm_format: AudioFormat,
1792 out_format: AudioFormat,
1793 headers: dict[str, str],
1794 ) -> web.StreamResponse:
1795 """Serve audio via independent ffmpeg stream (fallback)."""
1796 player_id = player.player_id
1797 logger.debug(
1798 "[StreamMode:independent] Fallback stream for %s",
1799 player_id,
1800 )
1801
1802 audio_source = self.provider.mass.streams.get_stream(
1803 media,
1804 pcm_format,
1805 force_flow_mode=False,
1806 )
1807 output_plan = self.provider.mass.streams.audio.get_player_output_plan(
1808 player_id,
1809 pcm_format,
1810 out_format,
1811 queue_id=getattr(media, "source_id", None),
1812 session_id=get_media_session_id(media),
1813 queue_item_id=getattr(media, "queue_item_id", None),
1814 )
1815
1816 response = web.StreamResponse(status=200, headers=headers)
1817 stream_task: asyncio.Task[None] = asyncio.create_task(
1818 self._stream_with_prebuffer(
1819 request,
1820 response,
1821 player,
1822 headers,
1823 audio_source,
1824 pcm_format,
1825 out_format,
1826 output_plan.filter_params,
1827 )
1828 )
1829 transport = getattr(request, "transport", None)
1830 await self._run_stream_task(player_id, stream_task, transport)
1831
1832 return response
1833
1834 async def _stream_with_prebuffer(
1835 self,
1836 request: web.Request,
1837 response: web.StreamResponse,
1838 player: MSXPlayer,
1839 headers: dict[str, str],
1840 audio_source: Any,
1841 pcm_format: AudioFormat,
1842 out_format: AudioFormat,
1843 filter_params: Sequence[str | ComplexFilter],
1844 ) -> None:
1845 """Pre-buffer audio chunks, then send HTTP headers and stream remaining data."""
1846 player_id = player.player_id
1847 chunk_queue: asyncio.Queue[bytes | None] = asyncio.Queue(maxsize=32)
1848
1849 async def producer() -> None:
1850 try:
1851 async for chunk in get_ffmpeg_stream(
1852 audio_input=audio_source,
1853 input_format=pcm_format,
1854 output_format=out_format,
1855 filter_params=filter_params,
1856 extra_input_args=_READRATE_ARGS,
1857 ):
1858 await chunk_queue.put(chunk)
1859 finally:
1860 with contextlib.suppress(asyncio.QueueFull):
1861 chunk_queue.put_nowait(None)
1862
1863 producer_task: asyncio.Task[None] | None = None
1864 total_bytes = 0
1865 try:
1866 producer_task = asyncio.create_task(producer())
1867
1868 # Phase 1: Pre-buffer â collect chunks until we have enough data
1869 pre_buffer: list[bytes] = []
1870 pre_buffer_size = 0
1871 while pre_buffer_size < PRE_BUFFER_BYTES:
1872 chunk = await chunk_queue.get()
1873 if chunk is None:
1874 break
1875 pre_buffer.append(chunk)
1876 pre_buffer_size += len(chunk)
1877
1878 # Re-check: stop may have been called while buffering
1879 if not player.current_media and not pre_buffer:
1880 return
1881
1882 # NOW send HTTP headers + pre-buffer burst
1883 await response.prepare(request)
1884 for buf_chunk in pre_buffer:
1885 await response.write(buf_chunk)
1886 total_bytes += len(buf_chunk)
1887
1888 # If pre-buffer ended with sentinel, we're done
1889 if chunk is None:
1890 return
1891
1892 # Phase 2: Stream remaining chunks normally
1893 while True:
1894 chunk = await chunk_queue.get()
1895 if chunk is None:
1896 break
1897 await response.write(chunk)
1898 total_bytes += len(chunk)
1899 except ConnectionResetError, BrokenPipeError, ConnectionAbortedError:
1900 logger.debug("Client disconnected from stream %s", player_id)
1901 except asyncio.CancelledError:
1902 logger.debug("Stream cancelled for player %s", player_id)
1903 raise
1904 finally:
1905 if producer_task and not producer_task.done():
1906 producer_task.cancel()
1907 with contextlib.suppress(asyncio.CancelledError):
1908 await producer_task
1909 content_length = headers.get("Content-Length")
1910 if content_length:
1911 logger.debug(
1912 "Stream %s: wrote %d bytes, Content-Length=%s, diff=%d",
1913 player_id,
1914 total_bytes,
1915 content_length,
1916 total_bytes - int(content_length),
1917 )
1918 else:
1919 logger.debug("Stream %s finished: wrote %d bytes", player_id, total_bytes)
1920
1921 async def _run_stream_task(
1922 self,
1923 player_id: str,
1924 stream_task: asyncio.Task[None],
1925 transport: Any,
1926 ) -> None:
1927 """Run a stream task with registration and error handling."""
1928 self._register_stream(player_id, stream_task, transport)
1929 try:
1930 await stream_task
1931 except asyncio.CancelledError:
1932 raise
1933 except Exception:
1934 logger.exception("Stream error for player %s", player_id)
1935 finally:
1936 self._unregister_stream(player_id, stream_task, transport)
1937
1938 # --- WebSocket, Broadcast & Health ---
1939
1940 async def _handle_health(self, request: web.Request) -> web.Response:
1941 """Health check endpoint."""
1942 return web.json_response(
1943 {
1944 "status": "ok",
1945 "provider": "msx_bridge",
1946 "players": len(self.provider.players),
1947 }
1948 )
1949
1950 async def _handle_ws(self, request: web.Request) -> web.WebSocketResponse:
1951 """
1952 WebSocket for push playback â clients subscribe by player_id.
1953
1954 Uses the same player_id derivation (device_id or IP) as content and
1955 stream endpoints so broadcast_stop reaches the correct client.
1956 Registers the player in MA on connect so the player appears when MSX starts.
1957 """
1958 ws = web.WebSocketResponse(heartbeat=30)
1959 await ws.prepare(request)
1960
1961 player_id, _, player = await self._ensure_player_for_request(request)
1962 if player_id not in self._ws_clients:
1963 self._ws_clients[player_id] = set()
1964 self._ws_clients[player_id].add(ws)
1965 logger.info(
1966 "WebSocket connected: player_id=%s, clients_for_player=%d, all_players=%s",
1967 player_id,
1968 len(self._ws_clients[player_id]),
1969 list(self._ws_clients.keys()),
1970 )
1971 if player and isinstance(player, MSXPlayer):
1972 player.on_ws_connected()
1973
1974 try:
1975 async for msg in ws:
1976 if msg.type == WSMsgType.TEXT:
1977 self._handle_ws_message(player_id, msg.data)
1978 finally:
1979 self._ws_clients.get(player_id, set()).discard(ws)
1980 if not self._ws_clients.get(player_id):
1981 self._ws_clients.pop(player_id, None)
1982 # Notify the player that its last WS client disconnected
1983 offline_player = self.provider.mass.players.get_player(player_id)
1984 if offline_player and isinstance(offline_player, MSXPlayer):
1985 offline_player.on_ws_disconnected()
1986 logger.debug("WebSocket client disconnected for player %s", player_id)
1987
1988 return ws
1989
1990 def _register_stream(self, player_id: str, task: asyncio.Task[None], transport: Any) -> None:
1991 """Register active stream task and transport for cancel on stop."""
1992 if player_id not in self._active_stream_tasks:
1993 self._active_stream_tasks[player_id] = set()
1994 self._active_stream_transports[player_id] = set()
1995 if task:
1996 self._active_stream_tasks[player_id].add(task)
1997 if transport:
1998 self._active_stream_transports[player_id].add(transport)
1999
2000 def _unregister_stream(self, player_id: str, task: asyncio.Task[None], transport: Any) -> None:
2001 """Unregister stream when done (from finally block)."""
2002 if player_id not in self._active_stream_tasks:
2003 return
2004 if task:
2005 self._active_stream_tasks[player_id].discard(task)
2006 if transport:
2007 self._active_stream_transports[player_id].discard(transport)
2008 if not self._active_stream_tasks[player_id]:
2009 del self._active_stream_tasks[player_id]
2010 del self._active_stream_transports[player_id]
2011
2012 async def _ws_send(
2013 self, ws: web.WebSocketResponse, text: str, player_id: str | None = None
2014 ) -> None:
2015 """Send text to WebSocket; on failure warn and remove the stale client."""
2016 try:
2017 await ws.send_str(text)
2018 except Exception as exc:
2019 logger.warning("WebSocket send failed (player=%s): %s", player_id, exc)
2020 if player_id:
2021 self._ws_clients.get(player_id, set()).discard(ws)
2022
2023 async def _cmd_pause_no_echo(self, player_id: str) -> None:
2024 """Pause player without echoing back to MSX."""
2025 player = self.provider.mass.players.get_player(player_id)
2026 if not (player and isinstance(player, MSXPlayer)):
2027 return
2028 player._skip_ws_notify = True
2029 try:
2030 await self.provider.mass.players.cmd_pause(player_id)
2031 finally:
2032 player._skip_ws_notify = False
2033
2034 async def _cmd_play_no_echo(self, player_id: str) -> None:
2035 """Resume player without echoing back to MSX."""
2036 player = self.provider.mass.players.get_player(player_id)
2037 if not (player and isinstance(player, MSXPlayer)):
2038 return
2039 player._skip_ws_notify = True
2040 try:
2041 await self.provider.mass.players.cmd_play(player_id)
2042 finally:
2043 player._skip_ws_notify = False
2044
2045 def _handle_ws_message(self, player_id: str, data: str) -> None:
2046 """Process an inbound WebSocket message from MSX."""
2047 try:
2048 msg = json.loads(data)
2049 except json.JSONDecodeError, TypeError:
2050 logger.debug("Invalid WS message from %s: %s", player_id, data)
2051 return
2052
2053 msg_type = msg.get("type")
2054 if msg_type == "position":
2055 position = msg.get("position")
2056 if position is not None and isinstance(position, (int, float)):
2057 player = self.provider.mass.players.get_player(player_id)
2058 if player and isinstance(player, MSXPlayer):
2059 player.update_position(float(position))
2060 self.provider.on_player_activity(player_id)
2061 elif msg_type == "pause":
2062 player = self.provider.mass.players.get_player(player_id)
2063 if player and isinstance(player, MSXPlayer):
2064 position = msg.get("position")
2065 if position is not None and isinstance(position, (int, float)):
2066 player.update_position(float(position))
2067 self.provider.mass.create_task(self._cmd_pause_no_echo(player_id))
2068 self.provider.on_player_activity(player_id)
2069 elif msg_type == "resume":
2070 player = self.provider.mass.players.get_player(player_id)
2071 if player and isinstance(player, MSXPlayer):
2072 self.provider.mass.create_task(self._cmd_play_no_echo(player_id))
2073 self.provider.on_player_activity(player_id)
2074 else:
2075 logger.debug("Unknown WS message type from %s: %s", player_id, msg_type)
2076
2077 # --- Stream Proxy ---
2078
2079 async def _handle_stream(self, request: web.Request) -> web.StreamResponse:
2080 """Stream audio from MA to the TV using internal API."""
2081 player_id = _strip_known_extension(request.match_info["player_id"])
2082
2083 player = self.provider.mass.players.get_player(player_id)
2084 if not player or not isinstance(player, MSXPlayer):
2085 return web.Response(status=404, text="Player not found")
2086 if rejected := self._reject_invalid_stream_token(request, player_id):
2087 return rejected
2088 self.provider.on_player_activity(player_id)
2089
2090 media = player.current_media
2091 if not media:
2092 return web.Response(status=404, text="No active stream")
2093
2094 return await self._serve_audio_stream(
2095 request,
2096 player,
2097 media,
2098 duration=self._resolve_served_duration(media),
2099 )
2100
2101 # --- Library API Routes ---
2102
2103 async def _handle_albums(self, request: web.Request) -> web.Response:
2104 """List albums."""
2105 limit = _int_param(request.query, "limit", 50)
2106 offset = _int_param(request.query, "offset", 0)
2107 albums = await self.provider.mass.music.albums.library_items(
2108 limit=limit, offset=offset, summary=False
2109 )
2110 return web.json_response(
2111 {
2112 "items": [
2113 {
2114 "item_id": str(album.item_id),
2115 "name": album.name,
2116 "artist": getattr(album, "artist_str", ""),
2117 "image": get_image_url(album, self.provider),
2118 "uri": album.uri,
2119 }
2120 for album in albums
2121 ],
2122 "total": albums.total if hasattr(albums, "total") else len(albums),
2123 }
2124 )
2125
2126 async def _handle_album_tracks(self, request: web.Request) -> web.Response:
2127 """List tracks for an album."""
2128 item_id = request.match_info["item_id"]
2129 tracks = await self.provider.mass.music.albums.tracks(item_id, "library")
2130 return web.json_response(
2131 {
2132 "items": [self._format_track(track) for track in tracks],
2133 }
2134 )
2135
2136 async def _handle_artists(self, request: web.Request) -> web.Response:
2137 """List artists."""
2138 limit = _int_param(request.query, "limit", 50)
2139 offset = _int_param(request.query, "offset", 0)
2140 artists = await self.provider.mass.music.artists.library_items(
2141 limit=limit, offset=offset, summary=False
2142 )
2143 return web.json_response(
2144 {
2145 "items": [
2146 {
2147 "item_id": str(artist.item_id),
2148 "name": artist.name,
2149 "image": get_image_url(artist, self.provider),
2150 "uri": artist.uri,
2151 }
2152 for artist in artists
2153 ],
2154 "total": artists.total if hasattr(artists, "total") else len(artists),
2155 }
2156 )
2157
2158 async def _handle_artist_albums(self, request: web.Request) -> web.Response:
2159 """List albums for an artist."""
2160 item_id = request.match_info["item_id"]
2161 albums = await self.provider.mass.music.artists.albums(item_id, "library")
2162 return web.json_response(
2163 {
2164 "items": [
2165 {
2166 "item_id": str(album.item_id),
2167 "name": album.name,
2168 "artist": getattr(album, "artist_str", ""),
2169 "image": get_image_url(album, self.provider),
2170 "uri": album.uri,
2171 }
2172 for album in albums
2173 ],
2174 }
2175 )
2176
2177 async def _handle_playlists(self, request: web.Request) -> web.Response:
2178 """List playlists."""
2179 limit = _int_param(request.query, "limit", 50)
2180 offset = _int_param(request.query, "offset", 0)
2181 playlists = await self.provider.mass.music.playlists.library_items(
2182 limit=limit, offset=offset, summary=False
2183 )
2184 return web.json_response(
2185 {
2186 "items": [
2187 {
2188 "item_id": str(playlist.item_id),
2189 "name": playlist.name,
2190 "image": get_image_url(playlist, self.provider),
2191 "uri": playlist.uri,
2192 }
2193 for playlist in playlists
2194 ],
2195 "total": playlists.total if hasattr(playlists, "total") else len(playlists),
2196 }
2197 )
2198
2199 async def _handle_playlist_tracks(self, request: web.Request) -> web.Response:
2200 """List tracks for a playlist."""
2201 item_id = request.match_info["item_id"]
2202 tracks = [t async for t in self.provider.mass.music.playlists.tracks(item_id, "library")]
2203 return web.json_response(
2204 {
2205 "items": [self._format_track(track) for track in tracks],
2206 }
2207 )
2208
2209 async def _handle_tracks(self, request: web.Request) -> web.Response:
2210 """List tracks."""
2211 limit = _int_param(request.query, "limit", 50)
2212 offset = _int_param(request.query, "offset", 0)
2213 tracks = await self.provider.mass.music.tracks.library_items(
2214 limit=limit, offset=offset, summary=False
2215 )
2216 return web.json_response(
2217 {
2218 "items": [self._format_track(track) for track in tracks],
2219 "total": tracks.total if hasattr(tracks, "total") else len(tracks),
2220 }
2221 )
2222
2223 async def _handle_search(self, request: web.Request) -> web.Response:
2224 """Search the music library."""
2225 query = request.query.get("q", "")
2226 if not query:
2227 return web.json_response({"error": "Missing query parameter 'q'"}, status=400)
2228 limit = _int_param(request.query, "limit", 20)
2229 results = await self.provider.mass.music.search(query, limit=limit)
2230 return web.json_response(
2231 {
2232 "artists": [
2233 {
2234 "item_id": str(a.item_id),
2235 "name": a.name,
2236 "image": get_image_url(a, self.provider),
2237 "uri": a.uri,
2238 }
2239 for a in results.artists
2240 ],
2241 "albums": [
2242 {
2243 "item_id": str(a.item_id),
2244 "name": a.name,
2245 "artist": getattr(a, "artist_str", ""),
2246 "image": get_image_url(a, self.provider),
2247 "uri": a.uri,
2248 }
2249 for a in results.albums
2250 ],
2251 "tracks": [self._format_track(t) for t in results.tracks],
2252 "playlists": [
2253 {
2254 "item_id": str(p.item_id),
2255 "name": p.name,
2256 "image": get_image_url(p, self.provider),
2257 "uri": p.uri,
2258 }
2259 for p in results.playlists
2260 ],
2261 }
2262 )
2263
2264 async def _handle_recently_played(self, request: web.Request) -> web.Response:
2265 """Return recently played items."""
2266 limit = _int_param(request.query, "limit", 20)
2267 tracks = await self.provider.mass.music.tracks.library_items(
2268 limit=limit, order_by="last_played", summary=False
2269 )
2270 return web.json_response(
2271 {
2272 "items": [self._format_track(track) for track in tracks],
2273 }
2274 )
2275
2276 async def _handle_lyrics(self, request: web.Request) -> web.Response:
2277 """Return lyrics for the currently playing track on a given player."""
2278 player_id = request.match_info["player_id"]
2279 empty = web.json_response({"lyrics": None, "lrc_lyrics": None})
2280
2281 player = self.provider.mass.players.get_player(player_id)
2282 if not player or not isinstance(player, MSXPlayer):
2283 return empty
2284
2285 media = player.current_media
2286 if not media or not media.source_id or not media.queue_item_id:
2287 return empty
2288
2289 queue_item = self.provider.mass.player_queues.get_item(media.source_id, media.queue_item_id)
2290 if not queue_item or not queue_item.media_item:
2291 return empty
2292
2293 track = queue_item.media_item
2294 if not isinstance(track, Track):
2295 return empty
2296 try:
2297 lyrics, lrc_lyrics = await self.provider.mass.metadata.get_track_lyrics(track)
2298 except Exception:
2299 lyrics, lrc_lyrics = None, None
2300
2301 return web.json_response(
2302 {
2303 "title": getattr(track, "name", ""),
2304 "artist": getattr(track, "artist_str", ""),
2305 "lyrics": lyrics,
2306 "lrc_lyrics": lrc_lyrics,
2307 }
2308 )
2309
2310 async def _handle_queue(self, request: web.Request) -> web.Response:
2311 """Return the current playback queue for a given player."""
2312 player_id = request.match_info["player_id"]
2313
2314 player = self.provider.mass.players.get_player(player_id)
2315 if not player or not isinstance(player, MSXPlayer):
2316 return web.json_response({"items": [], "current_index": -1})
2317
2318 queue_id = player_id
2319 try:
2320 queue_items = self.provider.mass.player_queues.items(queue_id)
2321 except Exception:
2322 logger.debug("Failed to fetch queue items for player %s", player_id, exc_info=True)
2323 queue_items = []
2324
2325 current_uri = None
2326 media = player.current_media
2327 if media and media.source_id and media.queue_item_id:
2328 qi = self.provider.mass.player_queues.get_item(media.source_id, media.queue_item_id)
2329 if qi and qi.media_item:
2330 current_uri = getattr(qi.media_item, "uri", None)
2331
2332 items: list[dict[str, Any]] = []
2333 current_index = -1
2334 for i, qi in enumerate(queue_items):
2335 mi = getattr(qi, "media_item", None)
2336 uri = getattr(mi, "uri", None) or ""
2337 img = None
2338 if hasattr(qi, "image") and qi.image:
2339 img = self.provider.mass.metadata.get_image_url(qi.image)
2340 items.append(
2341 {
2342 "title": getattr(mi, "name", None) or getattr(qi, "name", "") or "",
2343 "artist": getattr(mi, "artist_str", "") if mi else "",
2344 "duration": getattr(mi, "duration", None) or getattr(qi, "duration", 0) or 0,
2345 "image": img,
2346 "uri": uri,
2347 }
2348 )
2349 if current_uri and uri == current_uri and current_index < 0:
2350 current_index = i
2351
2352 return web.json_response({"items": items, "current_index": current_index})
2353
2354 # --- Party Mode ---
2355
2356 def _cached_party(self) -> PartyInfo | None:
2357 """Return the last cached party state without refreshing (sync contexts)."""
2358 return self._party_cache[1] if self._party_cache else None
2359
2360 async def _qr_cover_base(self, prefix: str) -> str | None:
2361 """Return the QR-cover endpoint base when a party is active, else None."""
2362 if await self._get_active_party() is None:
2363 return None
2364 return f"{prefix}/api/party/qr-cover.png"
2365
2366 async def _get_active_party(self) -> PartyInfo | None:
2367 """
2368 Return details of the active party, or None when no party is active.
2369
2370 Never raises: a broken or slow Party plugin degrades to "no party" so
2371 the core UI (menu, kiosk) keeps working. Results are cached briefly.
2372 """
2373 now = time.monotonic()
2374 if self._party_cache is not None and now - self._party_cache[0] < PARTY_CACHE_TTL:
2375 return self._party_cache[1]
2376 info: PartyInfo | None = None
2377 try:
2378 party = cast("Any", self.provider.mass.get_provider("party"))
2379 if party is not None:
2380 join_url = await asyncio.wait_for(party.get_party_url(), PARTY_CALL_TIMEOUT)
2381 if join_url:
2382 config = await asyncio.wait_for(party.get_party_config(), PARTY_CALL_TIMEOUT)
2383 info = PartyInfo(
2384 join_url=join_url,
2385 name=getattr(config, "party_name", None),
2386 qr_text=getattr(config, "qr_text", None),
2387 qr_version=hashlib.sha256(join_url.encode()).hexdigest()[:12],
2388 )
2389 except Exception:
2390 logger.warning("Party plugin status check failed", exc_info=True)
2391 self._party_cache = (now, info)
2392 return info
2393
2394 async def _handle_party_status(self, _request: web.Request) -> web.Response:
2395 """Return party status for the kiosk overlay."""
2396 party = await self._get_active_party()
2397 if party is None:
2398 return web.json_response({"active": False})
2399 # the join URL itself is deliberately not exposed â clients only get the QR
2400 # image URL (relative, so it works behind reverse proxies) plus an opaque
2401 # version so they refetch the image only when the join code rotates
2402 return web.json_response(
2403 {
2404 "active": True,
2405 "name": party.name,
2406 "qr_text": party.qr_text,
2407 "qr_url": "/api/party/qr.svg",
2408 "qr_version": party.qr_version,
2409 }
2410 )
2411
2412 async def _handle_party_qr(self, request: web.Request) -> web.Response:
2413 """Serve the guest join URL as a QR code image (SVG or PNG by route)."""
2414 party = await self._get_active_party()
2415 if party is None:
2416 return web.Response(status=404, text="No active party")
2417 kind = "png" if request.path.endswith(".png") else "svg"
2418 body = await asyncio.to_thread(_render_qr, party.join_url, kind)
2419 return web.Response(
2420 body=body,
2421 content_type="image/png" if kind == "png" else "image/svg+xml",
2422 headers={"Cache-Control": "no-store"},
2423 )
2424
2425 async def _handle_party_qr_cover(self, request: web.Request) -> web.Response:
2426 """
2427 Serve a cover image with the party QR stamped into its corner (PNG).
2428
2429 MSX cannot render overlays, so during a party the playback background
2430 is routed through this endpoint. Degrades to a redirect to the
2431 original image when the party ended, the source is not ours, or the
2432 fetch/composite fails â stale playlist JSON on TVs keeps working.
2433 """
2434 image_url = request.query.get("image", "")
2435 if not image_url:
2436 return web.Response(status=400, text="Missing image parameter")
2437 # Reject non-MA sources outright â redirecting would be an open
2438 # redirect and fetching would be an SSRF proxy.
2439 if not self._is_allowed_cover_source(request, image_url):
2440 return web.Response(status=400, text="Image source not permitted")
2441 party = await self._get_active_party()
2442 if party is None:
2443 raise web.HTTPFound(location=image_url)
2444 cache_key = (image_url, party.qr_version)
2445 if (cached := self._qr_cover_cache.get(cache_key)) is None:
2446 try:
2447 # join: a TV dropping its request must not cancel the shared
2448 # render â late joiners and the cache still get the result
2449 cached = await join_task(self._qr_cover_task(cache_key, image_url, party.join_url))
2450 except Exception as err:
2451 logger.debug("QR cover composite failed for %s: %s", image_url, err)
2452 raise web.HTTPFound(location=image_url) from None
2453 return web.Response(
2454 body=cached,
2455 content_type="image/png",
2456 headers={"Cache-Control": "no-store"},
2457 )
2458
2459 def _qr_cover_task(
2460 self, cache_key: tuple[str, str], image_url: str, join_url: str
2461 ) -> asyncio.Task[bytes]:
2462 """Return the in-flight render task for this cover, starting one if needed."""
2463 if (task := self._qr_cover_inflight.get(cache_key)) is None:
2464 task = asyncio.create_task(self._fetch_and_render_cover(cache_key, image_url, join_url))
2465 self._qr_cover_inflight[cache_key] = task
2466
2467 def _cleanup(finished: asyncio.Task[bytes]) -> None:
2468 self._qr_cover_inflight.pop(cache_key, None)
2469 # consume the exception so a task whose waiters were all
2470 # cancelled never logs "exception was never retrieved"
2471 if not finished.cancelled():
2472 finished.exception()
2473
2474 task.add_done_callback(_cleanup)
2475 return task
2476
2477 async def _fetch_and_render_cover(
2478 self, cache_key: tuple[str, str], image_url: str, join_url: str
2479 ) -> bytes:
2480 """Fetch the cover, composite the QR onto it, and cache the PNG."""
2481 async with self.provider.mass.http_session.get(
2482 image_url,
2483 timeout=aiohttp.ClientTimeout(total=10),
2484 allow_redirects=False,
2485 ) as resp:
2486 if resp.status != 200:
2487 raise ValueError(f"cover fetch returned HTTP {resp.status}")
2488 cover_bytes = await resp.read()
2489 # PIL decode/re-encode blocks; on this loop it would stall audio
2490 # streaming for every player, so hop to a worker thread
2491 rendered = await asyncio.to_thread(_render_qr_cover, join_url, cover_bytes)
2492 # QR rotation changes the cache key; keep the cache tiny and bounded
2493 if len(self._qr_cover_cache) >= 32:
2494 self._qr_cover_cache.clear()
2495 self._qr_cover_cache[cache_key] = rendered
2496 return rendered
2497
2498 @staticmethod
2499 def _rewrite_stream_host(request: web.Request, url: str) -> str:
2500 """
2501 Point a stream URL at the host the client already uses to reach us.
2502
2503 The MA streamserver advertises its own IP, which is unreachable for
2504 the TV when MA runs behind Docker/NAT. The host the TV used for this
2505 request is known-good, so only the URL's host is replaced â scheme,
2506 port, path and query are preserved.
2507 """
2508 client_host = request.url.host
2509 if not client_host:
2510 return url
2511 parts = urlsplit(url)
2512 if ":" in client_host: # IPv6 literals need brackets in a netloc
2513 client_host = f"[{client_host}]"
2514 netloc = f"{client_host}:{parts.port}" if parts.port else client_host
2515 return urlunsplit((parts.scheme, netloc, parts.path, parts.query, parts.fragment))
2516
2517 @staticmethod
2518 def _url_origin(url: str) -> tuple[str, str | None, int | None]:
2519 """Return (scheme, hostname, port); raises ValueError on malformed URLs."""
2520 parts = urlsplit(url)
2521 # .port is lazy and raises on garbage like "host:8095.evil.example"
2522 return (parts.scheme, parts.hostname, parts.port)
2523
2524 def _is_allowed_cover_source(self, request: web.Request, image_url: str) -> bool:
2525 """Only composite covers served by this provider or MA itself (no open proxy)."""
2526 try:
2527 target_origin = self._url_origin(image_url)
2528 except ValueError:
2529 return False
2530 if target_origin[0] not in ("http", "https") or not target_origin[1]:
2531 return False
2532 allowed_bases = [self._get_prefix(request)]
2533 for source in (
2534 getattr(self.provider.mass, "webserver", None),
2535 getattr(self.provider.mass, "streams", None),
2536 ):
2537 base_url = getattr(source, "base_url", None)
2538 if isinstance(base_url, str) and base_url.startswith("http"):
2539 allowed_bases.append(base_url)
2540 # Compare parsed origins, not string prefixes: "http://ma:8095.evil.com"
2541 # must not pass for the allowed base "http://ma:8095".
2542 for base in allowed_bases:
2543 try:
2544 if target_origin == self._url_origin(base):
2545 return True
2546 except ValueError:
2547 continue
2548 return False
2549
2550 # --- Playback Control ---
2551
2552 def _reject_invalid_stream_token(
2553 self, request: web.Request, player_id: str
2554 ) -> web.Response | None:
2555 """
2556 Reject an audio request that does not carry the player's own stream token.
2557
2558 A TV cannot send an auth header, so the token travels in the URL the bridge
2559 itself generated. This stops a request that was never handed out â a web page
2560 firing an <audio> tag at this LAN server. A URL that was handed out stays valid
2561 until the provider reloads, so this is not a defence against a captured URL.
2562 """
2563 expected = self.provider.get_stream_token(player_id)
2564 if not secrets.compare_digest(request.query.get("token", ""), expected):
2565 return web.Response(status=403, text="Invalid or missing stream token")
2566 return None
2567
2568 @staticmethod
2569 def _reject_cross_site(request: web.Request) -> web.Response | None:
2570 """
2571 Reject browser cross-site requests to state-changing endpoints (CSRF guard).
2572
2573 Any web page can fire an unauthenticated GET at this LAN server via an
2574 img/script tag; modern browsers mark such requests with
2575 Sec-Fetch-Site: cross-site. Legitimate callers are same-origin (web
2576 player, MSX interaction plugin, dashboard) or non-browser clients that
2577 omit the header entirely â both pass.
2578 """
2579 if request.headers.get("Sec-Fetch-Site", "").lower() == "cross-site":
2580 return web.json_response({"error": "Cross-site request rejected"}, status=403)
2581 return None
2582
2583 async def _handle_play(self, request: web.Request) -> web.Response:
2584 """Start playback of a track."""
2585 if rejected := self._reject_cross_site(request):
2586 return rejected
2587 try:
2588 body = await request.json()
2589 except Exception:
2590 return web.json_response({"error": "Invalid JSON body"}, status=400)
2591
2592 track_uri = body.get("track_uri")
2593 player_id = body.get("player_id")
2594 # the body is untyped JSON, so the type matters as much as the presence
2595 if not isinstance(track_uri, str) or not isinstance(player_id, str):
2596 return web.json_response({"error": "Invalid track_uri or player_id"}, status=400)
2597 if not track_uri or not player_id:
2598 return web.json_response({"error": "Missing track_uri or player_id"}, status=400)
2599 if not await _is_media_item_uri(track_uri):
2600 return web.json_response({"error": "Invalid track_uri"}, status=400)
2601
2602 if self._get_msx_player(player_id) is None:
2603 return web.json_response({"error": "Unknown MSX player"}, status=404)
2604
2605 async with ImpersonatedUser(self.provider.mass, await self.provider.get_owner_username()):
2606 await self.provider.mass.player_queues.play_media(player_id, track_uri)
2607 return web.json_response({"status": "ok"})
2608
2609 async def _handle_pause(self, request: web.Request) -> web.Response:
2610 """Pause playback."""
2611 if rejected := self._reject_cross_site(request):
2612 return rejected
2613 player_id = _strip_known_extension(request.match_info["player_id"])
2614 if self._get_msx_player(player_id) is None:
2615 return web.json_response({"error": "Unknown MSX player"}, status=404)
2616 self.provider.on_player_activity(player_id)
2617 await self.provider.mass.players.cmd_pause(player_id)
2618 return web.json_response({"status": "ok"})
2619
2620 async def _handle_stop(self, request: web.Request) -> web.Response:
2621 """Stop playback."""
2622 if rejected := self._reject_cross_site(request):
2623 return rejected
2624 player_id = _strip_known_extension(request.match_info["player_id"])
2625 if self._get_msx_player(player_id) is None:
2626 return web.json_response({"error": "Unknown MSX player"}, status=404)
2627 self.provider.on_player_activity(player_id)
2628 await self.provider.mass.players.cmd_stop(player_id)
2629 return web.json_response({"status": "ok"})
2630
2631 async def _handle_quick_stop(self, request: web.Request) -> web.Response:
2632 """Stop playback on MSX immediately (same signal as Disable)."""
2633 if rejected := self._reject_cross_site(request):
2634 return rejected
2635 player_id = _strip_known_extension(request.match_info["player_id"])
2636 if self._get_msx_player(player_id) is None:
2637 return web.json_response({"error": "Unknown MSX player"}, status=404)
2638 self.provider.on_player_activity(player_id)
2639 await self.provider.mass.players.cmd_stop(player_id)
2640 self.provider.notify_play_stopped(player_id)
2641 accept = request.headers.get("Accept", "")
2642 if "text/html" in accept:
2643 return web.Response(status=303, headers={"Location": "/"})
2644 return web.json_response({"status": "ok"})
2645
2646 async def _handle_next(self, request: web.Request) -> web.Response:
2647 """Skip to next track."""
2648 if rejected := self._reject_cross_site(request):
2649 return rejected
2650 player_id = _strip_known_extension(request.match_info["player_id"])
2651 if self._get_msx_player(player_id) is None:
2652 return web.json_response({"error": "Unknown MSX player"}, status=404)
2653 self.provider.on_player_activity(player_id)
2654 await self.provider.mass.players.cmd_next_track(player_id)
2655 return web.json_response({"status": "ok"})
2656
2657 async def _handle_previous(self, request: web.Request) -> web.Response:
2658 """Skip to previous track."""
2659 if rejected := self._reject_cross_site(request):
2660 return rejected
2661 player_id = _strip_known_extension(request.match_info["player_id"])
2662 if self._get_msx_player(player_id) is None:
2663 return web.json_response({"error": "Unknown MSX player"}, status=404)
2664 self.provider.on_player_activity(player_id)
2665 await self.provider.mass.players.cmd_previous_track(player_id)
2666 return web.json_response({"status": "ok"})
2667
2668 # --- Helpers ---
2669
2670 def _get_msx_player(self, player_id: str) -> MSXPlayer | None:
2671 """Return the MSXPlayer for player_id if it belongs to this provider, else None."""
2672 player = self.provider.mass.players.get_player(player_id, raise_unavailable=False)
2673 if isinstance(player, MSXPlayer) and player.provider == self.provider:
2674 return player
2675 return None
2676
2677 def _get_prefix(self, request: web.Request) -> str:
2678 """
2679 Build URL prefix for JSON content, using our known port.
2680
2681 Uses aiohttp's parsed URL host (IPv6-safe, no port) and substitutes
2682 self.port. Note: host is still derived from the Host header; a crafted
2683 header can influence the returned host, but the server binds to 0.0.0.0
2684 so there is no single canonical IP to validate against.
2685 """
2686 host: str = request.url.host or request.host.split(":")[0] # IPv6-safe, no port
2687 host_addr = f"[{host}]" if ":" in host else host # bracket IPv6 literals for URLs
2688 return f"http://{host_addr}:{self.port}"
2689
2690 def _get_player_id_and_device_param(self, request: web.Request) -> tuple[str, str]:
2691 """
2692 Extract player_id and device_id query param from request.
2693
2694 Returns (player_id, device_param) where device_param is e.g. "device_id=xxx"
2695 or "" if using IP fallback.
2696 """
2697 device_id = request.query.get("device_id")
2698 remote_ip = request.remote or "unknown"
2699
2700 if device_id:
2701 device_id = device_id[:64] # clamp before sanitizing (UUIDs are 36 chars)
2702 sanitized = PLAYER_ID_SANITIZE_RE.sub("_", device_id).strip("_") or "device"
2703 player_id = f"{MSX_PLAYER_ID_PREFIX}{sanitized}"
2704 param = f"device_id={quote(device_id, safe='')}"
2705 logger.info(
2706 "[PlayerID] device_id=%s, remote_ip=%s -> player_id=%s",
2707 device_id,
2708 remote_ip,
2709 player_id,
2710 )
2711 else:
2712 ip = remote_ip if remote_ip != "unknown" else "0_0_0_0"
2713 sanitized = PLAYER_ID_SANITIZE_RE.sub("_", ip.replace(".", "_")).strip("_") or "ip"
2714 player_id = f"{MSX_PLAYER_ID_PREFIX}{sanitized}"
2715 param = ""
2716 logger.info(
2717 "[PlayerID] no device_id, remote_ip=%s -> player_id=%s",
2718 remote_ip,
2719 player_id,
2720 )
2721 return player_id, param
2722
2723 async def _ensure_player_for_request(
2724 self, request: web.Request
2725 ) -> tuple[str, str, MSXPlayer | None]:
2726 """
2727 Get or register player for this request.
2728
2729 Returns (player_id, device_param, player).
2730 Player may be None if registration failed.
2731 """
2732 player_id, device_param = self._get_player_id_and_device_param(request)
2733 # Remember how this client reaches us â WS pushes have no request context
2734 self._client_prefixes[player_id] = self._get_prefix(request)
2735 remote_ip = request.remote
2736 # Web player clients pass source=web to distinguish from MSX TV players
2737 prefix_label = "WEB TV" if request.query.get("source") == "web" else "MSX TV"
2738 display_name = self.provider._player_display_name_from_id(
2739 player_id, prefix_label=prefix_label, remote_ip=remote_ip
2740 )
2741 player = await self.provider.get_or_register_player(
2742 player_id, display_name=display_name, ip_address=remote_ip
2743 )
2744 return player_id, device_param, player
2745
2746 def _current_media_matches_uri(self, player: MSXPlayer, track_uri: str) -> bool:
2747 """Check if player's current_media corresponds to the requested track URI."""
2748 media = player.current_media
2749 if not media or not media.source_id or not media.queue_item_id:
2750 return False
2751 queue_item = self.provider.mass.player_queues.get_item(media.source_id, media.queue_item_id)
2752 if queue_item and queue_item.media_item:
2753 return getattr(queue_item.media_item, "uri", None) == track_uri
2754 return False
2755
2756 def _format_track(self, track: Any) -> dict[str, Any]:
2757 """Format a track object for the API response."""
2758 return {
2759 "item_id": str(track.item_id),
2760 "name": track.name,
2761 "artist": getattr(track, "artist_str", ""),
2762 "album": getattr(getattr(track, "album", None), "name", ""),
2763 "duration": getattr(track, "duration", 0),
2764 "image": self.provider.mass.metadata.get_image_url(track.image)
2765 if hasattr(track, "image") and track.image
2766 else None,
2767 "uri": track.uri,
2768 }
2769