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