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