/
/
/
1"""Squeezelite Player Provider implementation."""
2
3from __future__ import annotations
4
5import logging
6from typing import TYPE_CHECKING, cast
7
8from aiohttp import web
9from aioslimproto.models import EventType as SlimEventType
10from aioslimproto.models import SlimEvent
11from aioslimproto.server import SlimServer
12from music_assistant_models.config_entries import ConfigEntry
13from music_assistant_models.enums import ConfigEntryType, MediaType
14from music_assistant_models.errors import SetupFailedError
15
16from music_assistant.constants import CONF_PORT, CONF_SYNC_ADJUST, VERBOSE_LOG_LEVEL
17from music_assistant.helpers.audio import get_mime_type
18from music_assistant.helpers.util import is_port_in_use
19from music_assistant.models.player_provider import PlayerProvider
20
21from .constants import (
22 CONF_CLI_JSON_PORT,
23 CONF_CLI_TELNET_PORT,
24 CONF_DISCOVERY,
25 DEFAULT_SLIMPROTO_PORT,
26)
27from .player import SqueezelitePlayer
28
29if TYPE_CHECKING:
30 from aioslimproto.client import SlimClient
31
32
33class SqueezelitePlayerProvider(PlayerProvider):
34 """Player provider for players using slimproto (like Squeezelite)."""
35
36 reload_on_streams_network_change = True
37 slimproto: SlimServer | None = None
38
39 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
40 """Return Config entries to setup this provider."""
41 return (
42 ConfigEntry(
43 key=CONF_CLI_TELNET_PORT,
44 type=ConfigEntryType.INTEGER,
45 default_value=9090,
46 advanced=True,
47 ),
48 ConfigEntry(
49 key=CONF_CLI_JSON_PORT,
50 type=ConfigEntryType.INTEGER,
51 default_value=9000,
52 advanced=True,
53 ),
54 ConfigEntry(
55 key=CONF_DISCOVERY,
56 type=ConfigEntryType.BOOLEAN,
57 default_value=True,
58 advanced=True,
59 ),
60 ConfigEntry(
61 key=CONF_PORT,
62 type=ConfigEntryType.INTEGER,
63 default_value=DEFAULT_SLIMPROTO_PORT,
64 ),
65 )
66
67 async def handle_async_init(self) -> None:
68 """Handle async initialization of the provider."""
69 # set-up aioslimproto logging
70 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
71 logging.getLogger("aioslimproto").setLevel(logging.DEBUG)
72 else:
73 logging.getLogger("aioslimproto").setLevel(self.logger.level + 10)
74
75 # Get all port configurations
76 control_port = cast("int", self.config.get_value(CONF_PORT))
77 telnet_port = cast("int | None", self.config.get_value(CONF_CLI_TELNET_PORT))
78 json_port = cast("int | None", self.config.get_value(CONF_CLI_JSON_PORT))
79
80 # Validate ALL required ports before starting ANY services
81 await self._validate_all_ports(control_port, telnet_port, json_port)
82
83 # create the server here (also validates config and sets up the CLI) but defer
84 # start() to loaded_in_mass, so we subscribe to events before accepting clients
85 self.slimproto = SlimServer(
86 cli_port=telnet_port or None,
87 cli_port_json=json_port or None,
88 ip_address=self.mass.streams.publish_ip,
89 name="Music Assistant",
90 control_port=control_port,
91 )
92
93 async def loaded_in_mass(self) -> None:
94 """Call after the provider has been loaded."""
95 await super().loaded_in_mass()
96 assert self.slimproto is not None # for type checker
97 # subscribe before starting the socket server: aioslimproto does not buffer
98 # events, so a client connecting before we subscribe would be missed entirely
99 self.slimproto.subscribe(self._handle_slimproto_event)
100 try:
101 await self.slimproto.start()
102 except Exception as err:
103 # ports were validated during setup, so a failure here is unlikely
104 self.unload_with_error(err)
105 return
106 self.mass.streams.register_dynamic_route(
107 "/slimproto/multi", self._serve_multi_client_stream
108 )
109 # it seems that WiiM devices do not use the json rpc port that is broadcasted
110 # in the discovery info but instead they just assume that the jsonrpc endpoint
111 # lives on the same server as stream URL. So we need to provide a jsonrpc.js
112 # endpoint that just redirects to the jsonrpc handler within the slimproto package.
113 self.mass.streams.register_dynamic_route(
114 "/jsonrpc.js", self.slimproto.cli._handle_jsonrpc_client
115 )
116
117 async def unload(self, is_removed: bool = False) -> None:
118 """Handle unload/close of the provider."""
119 # Ensure complete cleanup
120 await self._cleanup_server()
121 self.mass.streams.unregister_dynamic_route("/slimproto/multi")
122 self.mass.streams.unregister_dynamic_route("/jsonrpc.js")
123
124 def get_corrected_elapsed_milliseconds(self, slimplayer: SlimClient) -> int:
125 """Return corrected elapsed milliseconds for a slimplayer."""
126 sync_delay = self.mass.config.get_raw_player_config_value(
127 slimplayer.player_id, CONF_SYNC_ADJUST, 0
128 )
129 return int(slimplayer.elapsed_milliseconds - sync_delay)
130
131 async def _validate_all_ports(
132 self, control_port: int, telnet_port: int | None, json_port: int | None
133 ) -> None:
134 """Validate that all required ports are available before starting any services."""
135 ports_to_check = [(control_port, "SlimProto control")]
136
137 if telnet_port and telnet_port > 0:
138 ports_to_check.append((telnet_port, "Telnet CLI"))
139
140 if json_port and json_port > 0:
141 ports_to_check.append((json_port, "JSON-RPC CLI"))
142
143 # Collect all port conflicts before raising any errors
144 occupied_ports = []
145 for port, port_description in ports_to_check:
146 if await is_port_in_use(port):
147 occupied_ports.append(f"{port_description} port {port}")
148
149 # If any ports are occupied, raise a comprehensive error message
150 if occupied_ports:
151 if len(occupied_ports) == 1:
152 msg = f"{occupied_ports[0]} is not available"
153 else:
154 msg = f"Multiple ports are not available: {', '.join(occupied_ports)}"
155 raise SetupFailedError(msg)
156
157 async def _cleanup_server(self) -> None:
158 """Ensure complete cleanup of the SlimProto server on initialization failure."""
159 if self.slimproto:
160 try:
161 await self.slimproto.stop()
162 except Exception as err:
163 self.logger.warning("Error stopping SlimProto server during cleanup: %s", err)
164 finally:
165 self.slimproto = None
166
167 def _handle_slimproto_event(
168 self,
169 event: SlimEvent,
170 ) -> None:
171 """Handle events from SlimProto players."""
172 # Exit early if system is closing or slimproto server is not initialized
173 if self.mass.closing or not self.slimproto:
174 return
175
176 # Handle new player connect (or reconnect of existing player)
177 if event.type == SlimEventType.PLAYER_CONNECTED:
178 slimclient = self.slimproto.get_player(event.player_id)
179 if not slimclient:
180 return # should not happen, but guard anyways
181 player = SqueezelitePlayer(self, event.player_id, slimclient)
182 self.mass.create_task(player.setup())
183 return
184
185 if not (mass_player := self.mass.players.get_player(event.player_id)):
186 return # guard for unknown player
187 player = cast("SqueezelitePlayer", mass_player)
188
189 # Handle player disconnect
190 if event.type == SlimEventType.PLAYER_DISCONNECTED:
191 self.mass.create_task(self.mass.players.unregister(player.player_id))
192 return
193
194 # forward all other events to the player itself
195 player.handle_slim_event(event)
196
197 async def _serve_multi_client_stream(self, request: web.Request) -> web.StreamResponse:
198 """Serve the multi-client flow stream audio to a player."""
199 player_id = request.query.get("player_id")
200 fmt = request.query.get("fmt")
201 child_player_id = request.query.get("child_player_id")
202
203 if not player_id:
204 raise web.HTTPNotFound(reason="Missing player_id parameter")
205 if not fmt:
206 raise web.HTTPNotFound(reason="Missing fmt parameter")
207 if not child_player_id:
208 raise web.HTTPNotFound(reason="Missing child_player_id parameter")
209
210 if not (sync_parent := self.mass.players.get_player(player_id)):
211 raise web.HTTPNotFound(reason=f"Unknown player: {player_id}")
212 sync_parent = cast("SqueezelitePlayer", sync_parent)
213
214 if not (child_player := self.mass.players.get_player(child_player_id)):
215 raise web.HTTPNotFound(reason=f"Unknown player: {child_player_id}")
216
217 if not (stream := sync_parent.multi_client_stream) or stream.done:
218 raise web.HTTPNotFound(reason=f"There is no active stream for {player_id}!")
219
220 resp = web.StreamResponse(
221 status=200,
222 reason="OK",
223 headers={
224 "Content-Type": get_mime_type(fmt),
225 },
226 )
227 await resp.prepare(request)
228
229 # return early if this is not a GET request
230 if request.method != "GET":
231 return resp
232
233 # all checks passed, start streaming!
234 self.logger.debug(
235 "Start serving multi-client flow audio stream to %s",
236 child_player.display_name,
237 )
238
239 output_format = await self.mass.streams.audio.get_output_format(
240 output_format_str=fmt,
241 player=child_player,
242 content_sample_rate=stream.audio_format.sample_rate, # Flow PCM sample rate
243 content_bit_depth=stream.audio_format.bit_depth, # Flow PCM bit depth (32)
244 media_type=MediaType.FLOW_STREAM,
245 )
246 output_plan = self.mass.streams.audio.get_player_output_plan(
247 child_player_id,
248 stream.audio_format,
249 output_format,
250 queue_id=stream.queue_id,
251 session_id=stream.session_id,
252 )
253
254 async for chunk in stream.get_stream(
255 output_format=output_format,
256 filter_params=output_plan.filter_params,
257 ):
258 try:
259 await resp.write(chunk)
260 except BrokenPipeError, ConnectionResetError, ConnectionError:
261 # race condition
262 break
263 return resp
264