/
/
1"""
2go-librespot backend for the Spotify Connect provider.
3
4go-librespot is driven entirely over its local HTTP+WebSocket API: the WebSocket
5``/events`` stream feeds player/session/metadata/volume state into Music Assistant,
6and transport + volume commands are issued via the REST endpoints. Playback control
7works without a configured Spotify *music* provider or the Spotify Web API.
8"""
9
10from __future__ import annotations
11
12import asyncio
13import json
14import os
15from contextlib import suppress
16from functools import partial
17from pathlib import Path
18from typing import TYPE_CHECKING, Any
19
20from music_assistant_models.enums import ContentType
21from music_assistant_models.media_items import AudioFormat
22
23from music_assistant.helpers.process import AsyncProcess
24from music_assistant.helpers.util import (
25 interface_name_for_ip,
26 is_port_in_use,
27 select_free_port,
28)
29from music_assistant.providers.spotify_connect.clients.go_librespot import GoLibrespotClient
30from music_assistant.providers.spotify_connect.helpers import (
31 generate_device_id,
32 get_go_librespot_binary,
33)
34from music_assistant.providers.spotify_connect.models import (
35 BackendEvent,
36 BackendEventType,
37 BackendTrackMetadata,
38)
39
40from .base import SpotifyConnectBackend
41
42if TYPE_CHECKING:
43 import logging
44
45 from music_assistant.mass import MusicAssistant
46 from music_assistant.providers.spotify_connect.models import (
47 AudioChunkReader,
48 BackendEventCallback,
49 )
50
51# go-librespot volume scale; we pin volume_steps to this so the daemon's 0..max
52# volume maps 1:1 to a 0-100 percentage.
53VOLUME_STEPS = 100
54
55# Read size for pulling PCM off the daemon's stdout.
56STREAM_READ_CHUNK = 16384
57
58# Port range the go-librespot API server binds to (loopback only, one per instance).
59API_PORT_RANGE_START = 38800
60API_PORT_RANGE_END = 38900
61
62
63class GoLibrespotBackend(SpotifyConnectBackend):
64 """Spotify Connect backend wrapping a supervised go-librespot daemon."""
65
66 def __init__(
67 self,
68 mass: MusicAssistant,
69 *,
70 instance_id: str,
71 publish_name: str,
72 name: str,
73 logger: logging.Logger,
74 event_callback: BackendEventCallback,
75 ) -> None:
76 """
77 Initialize the backend (cheap; the daemon is launched in ``start``).
78
79 :param mass: The MusicAssistant instance.
80 :param instance_id: The owning provider's instance id; keys the
81 credential/cache dir and the stable Spotify device id.
82 :param publish_name: Device name advertised to the Spotify app.
83 :param name: Display name of the owning provider instance (log messages).
84 :param logger: Logger to use for diagnostics.
85 :param event_callback: Awaited with a normalized BackendEvent for every
86 state change the daemon reports.
87 """
88 self.mass = mass
89 self.logger = logger
90 self.name = name
91 self._instance_id = instance_id
92 self._publish_name = publish_name
93 self._event_callback = event_callback
94 self.cache_dir = os.path.join(self.mass.cache_path, instance_id)
95 self._binary: str | None = None
96 self._api_port: int = 0
97 self._client: GoLibrespotClient | None = None
98 self._stop_called: bool = False
99 self._daemon_task: asyncio.Task[None] | None = None
100 self._events_task: asyncio.Task[None] | None = None
101 self._proc: AsyncProcess | None = None
102 self._restart_error_count = 0
103 # _audio_format is the original Spotify source codec (Ogg Vorbis 320 kbps),
104 # advertised to clients for display. _decoded_audio_format is the raw PCM
105 # go-librespot actually writes to its stdout after decoding â what the
106 # audio reader yields and what the streams controller hands ffmpeg as
107 # the input format. We always emit the source's own format here; MA is
108 # responsible for converting it to whatever each player needs.
109 self._audio_format = AudioFormat(
110 content_type=ContentType.OGG,
111 codec_type=ContentType.VORBIS,
112 sample_rate=44100,
113 bit_depth=16,
114 channels=2,
115 bit_rate=320,
116 )
117 self._decoded_audio_format = AudioFormat(
118 content_type=ContentType.PCM_S16LE,
119 codec_type=ContentType.PCM_S16LE,
120 sample_rate=44100,
121 bit_depth=16,
122 channels=2,
123 )
124
125 @property
126 def audio_format(self) -> AudioFormat:
127 """Return the source audio format (advertised to clients for display)."""
128 return self._audio_format
129
130 @property
131 def decoded_audio_format(self) -> AudioFormat:
132 """Return the decoded PCM format the audio reader actually delivers."""
133 return self._decoded_audio_format
134
135 async def start(self) -> None:
136 """Start the backend and its supervised go-librespot daemon."""
137 self._binary = get_go_librespot_binary()
138 self._api_port = await select_free_port(
139 API_PORT_RANGE_START, API_PORT_RANGE_END, host="127.0.0.1"
140 )
141 self._client = GoLibrespotClient(
142 self.mass, f"http://127.0.0.1:{self._api_port}", self.logger
143 )
144 # Two self-healing supervisors: one keeps the daemon process alive, the
145 # other keeps the events websocket connected (reconnecting across daemon
146 # restarts). The events runner resets the daemon's restart backoff once
147 # the websocket is healthy again.
148 self._daemon_task = self.mass.create_task(self._daemon_runner())
149 self._events_task = self.mass.create_task(self._events_runner())
150
151 async def stop(self) -> None:
152 """Stop the daemon and all supervisor tasks."""
153 self._stop_called = True
154 for task in (self._events_task, self._daemon_task):
155 if task and not task.done():
156 task.cancel()
157 with suppress(asyncio.CancelledError):
158 await task
159
160 def get_audio_reader(self) -> AudioChunkReader | None:
161 """
162 Return a PCM chunk reader bound to the currently running daemon.
163
164 The reader stays bound to this daemon process: once it exits, the
165 reader returns b"" (clean EOF) even if a restarted daemon is already
166 up. None is returned when no daemon is running.
167 """
168 if (proc := self._proc) is None:
169 return None
170 return partial(proc.read, STREAM_READ_CHUNK)
171
172 async def play(self, uri: str, *, skip_to_uri: str | None = None) -> None:
173 """
174 Start playing a Spotify URI/context, making this device the active one.
175
176 :param uri: Spotify URI (track, album, playlist, ...) â typically a context.
177 :param skip_to_uri: Optional track URI within the context to start at.
178 """
179 assert self._client is not None
180 await self._client.play(uri, skip_to_uri=skip_to_uri)
181
182 async def resume(self) -> None:
183 """Resume playback on the active session."""
184 assert self._client is not None
185 await self._client.resume()
186
187 async def pause(self) -> None:
188 """Pause playback on the active session."""
189 assert self._client is not None
190 await self._client.pause()
191
192 async def next(self) -> None:
193 """Skip to the next track."""
194 assert self._client is not None
195 await self._client.next()
196
197 async def previous(self) -> None:
198 """Skip to the previous track (or rewind the current one)."""
199 assert self._client is not None
200 await self._client.prev()
201
202 async def seek(self, position_ms: int) -> None:
203 """
204 Seek to an absolute position in the current track.
205
206 :param position_ms: Target position in milliseconds.
207 """
208 assert self._client is not None
209 await self._client.seek(position_ms)
210
211 async def set_volume(self, volume: int) -> None:
212 """
213 Set the daemon's playback volume.
214
215 :param volume: Absolute volume as a 0-100 percentage (translated to
216 go-librespot's own volume scale).
217 """
218 assert self._client is not None
219 await self._client.set_volume(round(volume / 100 * VOLUME_STEPS))
220
221 def _write_config(self, source_ip: str | None) -> None:
222 """
223 Write the go-librespot ``config.yml`` for this instance.
224
225 go-librespot reads a YAML config; JSON is valid YAML, so we emit JSON to
226 sidestep an extra dependency and any string-quoting pitfalls (the device
227 name is user-provided). The config dir doubles as the credential/device
228 cache so the Spotify Connect device stays paired across restarts.
229
230 :param source_ip: Local address of the player-facing interface, or None to
231 advertise the Spotify Connect device on all interfaces.
232 """
233 Path(self.cache_dir).mkdir(parents=True, exist_ok=True)
234 config: dict[str, Any] = {
235 "device_name": self._publish_name,
236 "device_type": "speaker",
237 "device_id": generate_device_id(self._instance_id),
238 "bitrate": 320,
239 "audio_backend": "pipe",
240 # write decoded PCM to the daemon's stdout, which we capture and
241 # forward (the process pipe is always attached, so the daemon's
242 # non-blocking pipe open never fails for lack of a reader). s16le is
243 # the Spotify source representation; MA converts it per player.
244 "audio_output_pipe": "/dev/stdout",
245 "audio_output_pipe_format": "s16le",
246 # external_volume: don't let go-librespot attenuate the PCM â MA / the
247 # target player owns the actual volume. We still receive 'volume'
248 # events and push volume back so the Spotify app slider stays in sync.
249 # No initial_volume: go-librespot ignores it with external_volume set;
250 # the provider pushes the player's volume instead.
251 "external_volume": True,
252 "volume_steps": VOLUME_STEPS,
253 "zeroconf_enabled": True,
254 "credentials": {"type": "zeroconf", "zeroconf": {"persist_credentials": True}},
255 "server": {"enabled": True, "address": "127.0.0.1", "port": self._api_port},
256 }
257 # Advertise the Spotify Connect device only on the interface the streams
258 # server binds to, so it lands on the right network on multi-homed hosts.
259 # go-librespot selects advertise interfaces by name, so map the IP to one.
260 if source_ip:
261 if iface_name := interface_name_for_ip(source_ip):
262 config["zeroconf_interfaces_to_advertise"] = [iface_name]
263 else:
264 self.logger.debug(
265 "No interface found for stream bind IP %s; advertising on all interfaces",
266 source_ip,
267 )
268 config_file = os.path.join(self.cache_dir, "config.yml")
269 with open(config_file, "w", encoding="utf-8") as fileobj:
270 json.dump(config, fileobj, indent=2)
271
272 async def _daemon_runner(self) -> None:
273 """Run and supervise the go-librespot daemon, restarting it if it exits."""
274 assert self._binary
275 assert self._client
276 # Loop forever; stop() cancels this task and the explicit stop-check below
277 # handles a graceful exit without a restart.
278 while True:
279 # If the API port was taken while the daemon was down, move to a
280 # fresh port instead of crash-looping on a bind error.
281 if await is_port_in_use(self._api_port, host="127.0.0.1"):
282 self._api_port = await select_free_port(
283 API_PORT_RANGE_START, API_PORT_RANGE_END, host="127.0.0.1"
284 )
285 self._client.base_url = f"http://127.0.0.1:{self._api_port}"
286 self.logger.warning(
287 "API port in use by another process; switching to port %s", self._api_port
288 )
289 self._write_config(await self.mass.streams.get_source_ip())
290 proc: AsyncProcess | None = None
291 try:
292 # stdout carries the decoded PCM (audio_output_pipe=/dev/stdout) and
293 # is consumed through get_audio_reader; stderr carries the daemon's
294 # logs, read here. Because the process pipe owns stdout from spawn
295 # there is always a reader, so go-librespot's non-blocking pipe open
296 # never fails for lack of a consumer.
297 self._proc = proc = AsyncProcess(
298 [self._binary, "--config_dir", self.cache_dir],
299 stdout=True,
300 stderr=True,
301 name=f"go-librespot[{self.name}]",
302 )
303 await proc.start()
304 self.logger.info("Started Spotify Connect background daemon [%s]", self.name)
305 async for line in proc.iter_stderr():
306 self.logger.debug("[%s] %s", self.name, line)
307 except asyncio.CancelledError:
308 raise
309 except Exception as err:
310 self.logger.warning("go-librespot daemon error [%s]: %s", self.name, err)
311 finally:
312 if proc:
313 await proc.close()
314 # The daemon â and thus the Spotify session â is gone. Tell the
315 # provider so a dead/restarting daemon isn't treated as active and
316 # controllable; a fresh 'active' event re-establishes it on reconnect.
317 self._proc = None
318 try:
319 await self._event_callback(BackendEvent(BackendEventType.CONNECTION_LOST))
320 except Exception:
321 # never let a callback error replace a propagating
322 # cancellation or kill the daemon supervisor
323 self.logger.exception("Error while handling daemon exit")
324 if self._stop_called:
325 break
326 self.logger.info("Spotify Connect background daemon stopped for %s", self.name)
327 self._restart_error_count += 1
328 if self._restart_error_count >= 5:
329 await self._event_callback(
330 BackendEvent(
331 BackendEventType.FATAL_ERROR,
332 error="go-librespot daemon failed to start multiple times.",
333 )
334 )
335 return
336 await asyncio.sleep(2)
337
338 async def _events_runner(self) -> None:
339 """Keep the go-librespot events websocket connected, reconnecting as needed."""
340 assert self._client is not None
341 while not self._stop_called:
342 try:
343 if not await self._client.wait_until_ready():
344 await asyncio.sleep(2)
345 continue
346 # A live websocket means the daemon is healthy: reset the restart
347 # backoff counter the daemon supervisor uses.
348 self._restart_error_count = 0
349 await self._client.listen_events(self._handle_event)
350 except asyncio.CancelledError:
351 raise
352 except Exception as err:
353 self.logger.debug("go-librespot events websocket dropped: %s", err)
354 if not self._stop_called:
355 await asyncio.sleep(2)
356
357 async def _handle_event(self, event_type: str, data: dict[str, Any]) -> None:
358 """Translate a single go-librespot websocket event and emit it normalized."""
359 self.logger.debug("Received event [%s]: %s %s", self.name, event_type, data)
360 await self._event_callback(self._translate_event(event_type, data))
361
362 def _translate_event(self, event_type: str, data: dict[str, Any]) -> BackendEvent:
363 """Map a raw go-librespot event onto the normalized BackendEvent model."""
364 # Every event carries the latest context/track uris when present, so the
365 # provider can take playback back when the user moves the active device
366 # away in the Spotify app.
367 context_uri: str | None = data.get("context_uri") or None
368 track_uri: str | None = data.get("uri") or None
369 if event_type == "active":
370 return BackendEvent(
371 BackendEventType.SESSION_ACTIVE, context_uri=context_uri, track_uri=track_uri
372 )
373 if event_type == "inactive":
374 return BackendEvent(
375 BackendEventType.SESSION_INACTIVE, context_uri=context_uri, track_uri=track_uri
376 )
377 if event_type == "playing":
378 return BackendEvent(
379 BackendEventType.PLAYING, context_uri=context_uri, track_uri=track_uri
380 )
381 if event_type == "paused":
382 return BackendEvent(
383 BackendEventType.PAUSED, context_uri=context_uri, track_uri=track_uri
384 )
385 if event_type == "stopped":
386 return BackendEvent(
387 BackendEventType.STOPPED, context_uri=context_uri, track_uri=track_uri
388 )
389 if event_type == "metadata":
390 artists = data.get("artist_names") or []
391 duration_ms = data.get("duration")
392 return BackendEvent(
393 BackendEventType.METADATA,
394 context_uri=context_uri,
395 track_uri=track_uri,
396 metadata=BackendTrackMetadata(
397 track_uri=data.get("uri"),
398 title=data.get("name") or None,
399 artist=artists[0] if artists else None,
400 album=data.get("album_name"),
401 image_url=data.get("album_cover_url"),
402 duration=duration_ms // 1000 if duration_ms else None,
403 position=int(data.get("position", 0)) // 1000,
404 ),
405 )
406 if event_type == "seek":
407 return BackendEvent(
408 BackendEventType.POSITION,
409 context_uri=context_uri,
410 track_uri=track_uri,
411 position=int(data.get("position", 0)) // 1000,
412 )
413 if event_type == "volume" and data.get("value") is not None:
414 # translate the daemon's 0..max scale to a 0-100 percentage
415 max_value = data.get("max") or VOLUME_STEPS
416 return BackendEvent(
417 BackendEventType.VOLUME,
418 context_uri=context_uri,
419 track_uri=track_uri,
420 volume=int(int(data["value"]) / int(max_value) * 100),
421 )
422 return BackendEvent(BackendEventType.OTHER, context_uri=context_uri, track_uri=track_uri)
423