/
/
/
1"""
2Sendspin WebSocket proxy handler for Music Assistant.
3
4This module provides an authenticated WebSocket proxy to the internal Sendspin server,
5allowing web clients to connect through the main webserver instead of requiring direct
6access to the Sendspin port.
7"""
8
9from __future__ import annotations
10
11import asyncio
12import contextlib
13import json
14import logging
15from typing import TYPE_CHECKING
16
17import aiohttp
18from aiohttp import ClientConnectorError, WSMsgType, web
19
20from music_assistant.constants import MASS_LOGGER_NAME
21from music_assistant.controllers.webserver.helpers.auth_middleware import (
22 get_authenticated_user,
23 is_request_from_ingress,
24)
25
26if TYPE_CHECKING:
27 from music_assistant_models.auth import User
28
29 from music_assistant.controllers.webserver import WebserverController
30
31LOGGER = logging.getLogger(f"{MASS_LOGGER_NAME}.sendspin_proxy")
32
33
34class SendspinProxyHandler:
35 """Handler for proxying WebSocket connections to the internal Sendspin server."""
36
37 def __init__(self, webserver: WebserverController) -> None:
38 """
39 Initialize the Sendspin proxy handler.
40
41 :param webserver: The webserver controller instance.
42 """
43 self.webserver = webserver
44 self.mass = webserver.mass
45 self.logger = LOGGER
46
47 async def handle_sendspin_proxy(self, request: web.Request) -> web.WebSocketResponse:
48 """
49 Handle incoming WebSocket connection and proxy to internal Sendspin server.
50
51 Authentication is required as the first message. The client must send:
52 {"type": "auth", "token": "<access_token>"}
53
54 After successful authentication, all messages are proxied bidirectionally.
55
56 :param request: The incoming HTTP request to upgrade to WebSocket.
57 :return: The WebSocket response.
58 """
59 wsock = web.WebSocketResponse(heartbeat=25)
60 await wsock.prepare(request)
61
62 self.logger.debug("Sendspin proxy connection from %s", request.remote)
63
64 # Check for ingress authentication (HA handles auth via headers)
65 if is_request_from_ingress(request):
66 user = await get_authenticated_user(request)
67 if not user:
68 self.logger.warning(
69 "Ingress auth failed for sendspin proxy from %s", request.remote
70 )
71 await wsock.close(code=4001, message=b"Ingress authentication failed")
72 return wsock
73 self.logger.debug("Sendspin proxy authenticated via ingress: %s", user.username)
74 else:
75 # Regular auth via first message
76 try:
77 user = await self._authenticate(wsock)
78 if not user:
79 return wsock
80 except TimeoutError:
81 self.logger.warning("Auth timeout for sendspin proxy from %s", request.remote)
82 await wsock.close(code=4001, message=b"Authentication timeout")
83 return wsock
84 except Exception:
85 self.logger.exception("Auth error for sendspin proxy")
86 await wsock.close(code=4001, message=b"Authentication error")
87 return wsock
88
89 # The internal Sendspin server may not be ready yet during startup
90 # (it starts in the provider load phase, after the webserver).
91 # Retry a few times with backoff to handle this race condition.
92 try:
93 internal_ws = None
94 for attempt in range(5):
95 try:
96 internal_ws = await self.mass.http_session.ws_connect(
97 self.webserver.internal_sendspin_url
98 )
99 break
100 except ClientConnectorError:
101 if attempt < 4:
102 await asyncio.sleep(0.5 * (attempt + 1))
103 continue
104 self.logger.exception("Failed to connect to internal Sendspin server")
105 await wsock.close(code=1011, message=b"Internal server error")
106 return wsock
107 if internal_ws is None:
108 raise RuntimeError("Retry loop exited without connecting or returning")
109 except Exception:
110 self.logger.exception("Failed to connect to internal Sendspin server")
111 await wsock.close(code=1011, message=b"Internal server error")
112 return wsock
113 self.logger.debug("Sendspin proxy authenticated and connected for %s", request.remote)
114
115 try:
116 await self._proxy_messages(wsock, internal_ws)
117 finally:
118 if not internal_ws.closed:
119 await internal_ws.close()
120 if not wsock.closed:
121 await wsock.close()
122
123 return wsock
124
125 async def _authenticate(self, wsock: web.WebSocketResponse) -> User | None:
126 """
127 Wait for and validate authentication message.
128
129 :param wsock: The client WebSocket connection.
130 :return: The authenticated user, or None if authentication failed.
131 """
132 async with asyncio.timeout(10):
133 msg = await wsock.receive()
134
135 if msg.type != WSMsgType.TEXT:
136 await wsock.close(code=4001, message=b"Expected text message for auth")
137 return None
138
139 try:
140 auth_data = json.loads(msg.data)
141 except json.JSONDecodeError:
142 await wsock.close(code=4001, message=b"Invalid JSON in auth message")
143 return None
144
145 if auth_data.get("type") != "auth":
146 await wsock.close(code=4001, message=b"First message must be auth")
147 return None
148
149 token = auth_data.get("token")
150 if not token:
151 await wsock.close(code=4001, message=b"Token required in auth message")
152 return None
153
154 user = await self.webserver.auth.authenticate_with_token(token)
155 if not user:
156 await wsock.close(code=4001, message=b"Invalid or expired token")
157 return None
158
159 # Set the sendspin player_id on this session's websocket client(s)
160 # This allows the player controller to auto-whitelist this (web)player
161 # without modifying the user's player_filter list
162 client_id = auth_data.get("client_id")
163 if client_id:
164 self.webserver.set_sendspin_player_for_token(token, client_id)
165 self.logger.debug(
166 "Registered sendspin player %s for a session of user %s",
167 client_id,
168 user.username,
169 )
170
171 self.logger.debug("Sendspin proxy authenticated user: %s", user.username)
172 await wsock.send_str('{"type": "auth_ok"}')
173 return user
174
175 async def _proxy_messages(
176 self,
177 client_ws: web.WebSocketResponse,
178 internal_ws: aiohttp.ClientWebSocketResponse,
179 ) -> None:
180 """
181 Proxy messages bidirectionally between client and internal Sendspin server.
182
183 :param client_ws: The client WebSocket connection.
184 :param internal_ws: The internal Sendspin server WebSocket connection.
185 """
186 client_to_internal = asyncio.create_task(
187 self._forward_client_to_internal(client_ws, internal_ws)
188 )
189 internal_to_client = asyncio.create_task(
190 self._forward_internal_to_client(client_ws, internal_ws)
191 )
192
193 done, pending = await asyncio.wait(
194 [client_to_internal, internal_to_client],
195 return_when=asyncio.FIRST_COMPLETED,
196 )
197
198 for task in pending:
199 task.cancel()
200 peer_results = await asyncio.gather(*pending, return_exceptions=True)
201
202 # collect everything first so cleanup failures cannot mask the primary error
203 unexpected: list[BaseException] = []
204 for task in done:
205 with contextlib.suppress(asyncio.CancelledError):
206 if exc := task.exception():
207 self._collect_proxy_exception(exc, unexpected)
208 for result in peer_results:
209 if isinstance(result, BaseException):
210 self._collect_proxy_exception(result, unexpected)
211 if not unexpected:
212 return
213 for extra in unexpected[1:]:
214 self.logger.warning(
215 "Additional Sendspin proxy error while forwarding: %s",
216 extra,
217 )
218 raise unexpected[0]
219
220 def _collect_proxy_exception(
221 self,
222 exc: BaseException,
223 unexpected: list[BaseException],
224 ) -> None:
225 """Log expected transport disconnects; collect anything else."""
226 if isinstance(exc, asyncio.CancelledError):
227 return
228 if isinstance(
229 exc,
230 (ConnectionError, aiohttp.ClientError, asyncio.IncompleteReadError, EOFError),
231 ):
232 self.logger.debug("Sendspin proxy connection closed while forwarding: %s", exc)
233 return
234 unexpected.append(exc)
235
236 async def _forward_client_to_internal(
237 self,
238 client_ws: web.WebSocketResponse,
239 internal_ws: aiohttp.ClientWebSocketResponse,
240 ) -> None:
241 """
242 Forward messages from client to internal Sendspin server.
243
244 :param client_ws: The client WebSocket connection.
245 :param internal_ws: The internal Sendspin server WebSocket connection.
246 """
247 async for msg in client_ws:
248 if msg.type == WSMsgType.TEXT:
249 await internal_ws.send_str(msg.data)
250 elif msg.type == WSMsgType.BINARY:
251 await internal_ws.send_bytes(msg.data)
252 elif msg.type in (WSMsgType.CLOSE, WSMsgType.CLOSED, WSMsgType.ERROR):
253 if msg.type == WSMsgType.ERROR:
254 self.logger.debug("Sendspin proxy client transport error: %s", msg.data)
255 break
256
257 async def _forward_internal_to_client(
258 self,
259 client_ws: web.WebSocketResponse,
260 internal_ws: aiohttp.ClientWebSocketResponse,
261 ) -> None:
262 """
263 Forward messages from internal Sendspin server to client.
264
265 :param client_ws: The client WebSocket connection.
266 :param internal_ws: The internal Sendspin server WebSocket connection.
267 """
268 async for msg in internal_ws:
269 if msg.type == WSMsgType.TEXT:
270 await client_ws.send_str(msg.data)
271 elif msg.type == WSMsgType.BINARY:
272 await client_ws.send_bytes(msg.data)
273 elif msg.type in (WSMsgType.CLOSE, WSMsgType.CLOSED, WSMsgType.ERROR):
274 if msg.type == WSMsgType.ERROR:
275 self.logger.debug("Sendspin proxy internal transport error: %s", msg.data)
276 break
277