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