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