/
/
1"""
2Spotify Soloist backend for the Spotify Connect provider.
3
4Soloist is Spotify's official headless Connect client for Linux. The daemon is
5installed/updated through ``SoloistBinaryManager``, plays its audio into a
6private PulseAudio capture sink (read back by MA as a named pipe) and is driven
7over its local WebSocket API (``SoloistClient``).
8
9SECURITY NOTE: the daemon takes the user's personal API key on its command
10line; nothing in this module may ever log the process argv.
11"""
12
13from __future__ import annotations
14
15import asyncio
16import re
17from asyncio import FIRST_COMPLETED
18from contextlib import suppress
19from pathlib import Path
20from typing import TYPE_CHECKING, Any, Final
21
22from aiohttp import ClientError
23from music_assistant_models.enums import ContentType, RepeatMode, StreamType
24from music_assistant_models.errors import AudioError
25from music_assistant_models.media_items import AudioFormat
26
27from music_assistant.helpers.process import AsyncProcess
28from music_assistant.helpers.pulse_capture import (
29 CAPTURE_CHANNELS,
30 CAPTURE_SAMPLE_RATE,
31 PipeSink,
32 get_pulse_capture_server,
33)
34from music_assistant.providers.spotify_connect.base import (
35 AUDIO_QUALITY_HIGH,
36 AUDIO_QUALITY_LOSSLESS,
37 AUDIO_QUALITY_NORMAL,
38 AUDIO_QUALITY_VERY_HIGH,
39 SpotifyConnectBackend,
40)
41from music_assistant.providers.spotify_connect.models import (
42 BackendEvent,
43 BackendEventType,
44 BackendPlaybackOptions,
45 BackendQueueEntry,
46 BackendQueueState,
47 BackendStreamSource,
48 BackendTrackMetadata,
49 QueueEntrySource,
50)
51
52from .runtime import (
53 EXIT_CODE_BUILD_EXPIRED,
54 BuildExpiredError,
55 SoloistAuthState,
56 SoloistBinaryManager,
57 SoloistClient,
58 SoloistContextChanged,
59 SoloistDeviceChanged,
60 SoloistError,
61 SoloistErrorMessage,
62 SoloistOptionsChanged,
63 SoloistPlaybackState,
64 SoloistPositionSync,
65 SoloistQueueChanged,
66 SoloistTrackChanged,
67 SoloistVolumeChanged,
68)
69
70if TYPE_CHECKING:
71 import logging
72
73 from music_assistant.helpers.pulse_capture import PulseCaptureServer
74 from music_assistant.mass import MusicAssistant
75 from music_assistant.providers.spotify_connect.models import (
76 AudioChunkReader,
77 BackendEventCallback,
78 )
79
80 from .runtime import SoloistEntity, SoloistEvent, SoloistPlaybackOptions, SoloistQueueEntry
81
82# How Spotify-side volume changes are handled (see set_volume).
83VOLUME_MODE_PLAYER_ONLY: Final = "player_only"
84VOLUME_MODE_SYNC_SPOTIFY: Final = "sync_spotify"
85
86# Bounded soloist playback cache (docs: 0 = unlimited, otherwise at least 100 MB).
87CACHE_SIZE_MB: Final = 512
88
89# Supervised-restart policy: give up after this many consecutive daemon
90# failures; the counter resets once the events websocket delivers again.
91MAX_RESTART_ATTEMPTS: Final = 5
92RESTART_DELAY_S: Final = 2
93
94# How long the daemon's log reader may keep draining buffered output after the
95# process itself exited.
96DAEMON_LOG_DRAIN_TIMEOUT_S: Final = 5
97
98# Proactive binary refresh interval: soloist builds expire 90 days after their
99# build date, so a daily check swaps in a fresh build long before a long-lived
100# instance would hit the expiry.
101BINARY_REFRESH_INTERVAL_S: Final = 24 * 3600
102
103# How often the supervisor checks whether the pulse daemon restarted (which
104# invalidates the capture sink), so the sink is replaced proactively instead
105# of on the (side-effect-free) stream request.
106GENERATION_WATCH_INTERVAL_S: Final = 5
107
108# The classic desktop-client prefs keys (bare key=value lines) through which the
109# engine's audio behavior is controlled. The per-user prefs override the global
110# prefs per key, so both stores are (re)written before every daemon spawn.
111# NOTE: audio.crossfade.time_v2 is in MILLISECONDS; sub-second values silently
112# disable crossfade (verified empirically), so the key is only written when
113# crossfade is enabled (>= 1000 ms).
114_PREF_CROSSFADE: Final = "audio.crossfade_v2"
115_PREF_CROSSFADE_TIME: Final = "audio.crossfade.time_v2"
116_PREF_NORMALIZE: Final = "audio.normalize_v2"
117# The engine reads the metered variant on a metered connection and the
118# non-metered one otherwise; both are written so the tier holds either way.
119# The "migrated" marker is what makes the engine honor the non-metered key
120# instead of deriving it once from the metered one.
121_PREF_QUALITY: Final = "audio.play_bitrate_enumeration"
122_PREF_QUALITY_NON_METERED: Final = "audio.play_bitrate_non_metered_enumeration"
123_PREF_QUALITY_MIGRATED: Final = "audio.play_bitrate_non_metered_migrated"
124_MANAGED_PREFS: Final = (
125 _PREF_CROSSFADE,
126 _PREF_CROSSFADE_TIME,
127 _PREF_NORMALIZE,
128 _PREF_QUALITY,
129 _PREF_QUALITY_NON_METERED,
130 _PREF_QUALITY_MIGRATED,
131)
132
133# Quality tier -> the engine's bitrate enumeration value. Measured against
134# build 1.3.7.349 on a 4:20 track (bytes fetched for the whole file): 2 and 3
135# deliver ~96 and ~160 kbps, 4 ~320 kbps and 5 lossless FLAC (~810 kbps).
136# 5 is the ceiling â values outside 1-5 are rejected and silently fall back to
137# ~160 kbps, so an unknown tier must never reach the prefs file.
138_QUALITY_VALUES: Final[dict[str, int]] = {
139 AUDIO_QUALITY_NORMAL: 2,
140 AUDIO_QUALITY_HIGH: 3,
141 AUDIO_QUALITY_VERY_HIGH: 4,
142 AUDIO_QUALITY_LOSSLESS: 5,
143}
144
145# playback_state/playback_changed status values mapped to normalized events;
146# undocumented values degrade to OTHER.
147_STATUS_EVENTS: Final[dict[str, BackendEventType]] = {
148 "playing": BackendEventType.PLAYING,
149 "paused": BackendEventType.PAUSED,
150 "buffering": BackendEventType.BUFFERING,
151 "stopped": BackendEventType.STOPPED,
152 "idle": BackendEventType.STOPPED,
153}
154
155
156class SoloistBackend(SpotifyConnectBackend):
157 """
158 Spotify Connect backend wrapping a supervised Spotify Soloist daemon.
159
160 The daemon plays into a private PulseAudio capture sink whose FIFO is
161 consumed by the streams controller as a named pipe; control and state flow
162 over the daemon's local WebSocket API.
163 """
164
165 def __init__( # noqa: PLR0913
166 self,
167 mass: MusicAssistant,
168 *,
169 instance_id: str,
170 publish_name: str,
171 name: str,
172 logger: logging.Logger,
173 event_callback: BackendEventCallback,
174 api_key: str,
175 consent: bool,
176 volume_mode: str = VOLUME_MODE_PLAYER_ONLY,
177 crossfade_ms: int = 0,
178 loudness_normalization: bool = True,
179 audio_quality: str = AUDIO_QUALITY_LOSSLESS,
180 ) -> None:
181 """
182 Initialize the backend (cheap; the daemon is launched in ``start``).
183
184 :param mass: The MusicAssistant instance.
185 :param instance_id: The owning provider's instance id; keys the
186 data/cache dirs and the capture sink name.
187 :param publish_name: Device name advertised to the Spotify app.
188 :param name: Display name of the owning provider instance (log messages).
189 :param logger: Logger to use for diagnostics.
190 :param event_callback: Awaited with a normalized BackendEvent for every
191 state change the daemon reports.
192 :param api_key: The user's personal Spotify Soloist API key (secret,
193 kept out of all logs).
194 :param consent: Whether the user consented to downloading the soloist
195 binary from Spotify's CDN.
196 :param volume_mode: VOLUME_MODE_PLAYER_ONLY to let MA/the player own the
197 volume exclusively, VOLUME_MODE_SYNC_SPOTIFY to mirror the Spotify
198 app's volume onto the MA player.
199 :param crossfade_ms: Crossfade duration between tracks in milliseconds
200 (0 disables crossfade).
201 :param loudness_normalization: Whether Spotify's loudness normalization
202 should be applied to the audio.
203 :param audio_quality: Ceiling for the streaming quality Spotify is asked
204 to deliver (one of the AUDIO_QUALITY_* tiers).
205 """
206 self.mass = mass
207 self.logger = logger
208 self.name = name
209 self._publish_name = publish_name
210 self._event_callback = event_callback
211 self._api_key = api_key
212 self._consent = consent
213 self._volume_mode = volume_mode
214 self._crossfade_ms = crossfade_ms
215 self._loudness_normalization = loudness_normalization
216 self._audio_quality = audio_quality
217 self._data_dir = Path(mass.storage_path) / "spotify_connect" / instance_id / "soloist-data"
218 self._cache_dir = Path(mass.cache_path) / instance_id / "soloist-cache"
219 # PA sink names end up in space-delimited module arguments and env vars
220 self._sink_prefix = re.sub(r"[^A-Za-z0-9_.-]", "_", instance_id)
221 self._binary: Path | None = None
222 # digest of the build the running daemon was spawned from; the shared
223 # install can move ahead of it when a sibling instance updates first
224 self._build_sha: str | None = None
225 self._server: PulseCaptureServer | None = None
226 self._sink: PipeSink | None = None
227 self._sink_generation: int = -1
228 # serializes sink (re)creation against teardown in stop()
229 self._sink_lock = asyncio.Lock()
230 self._client: SoloistClient | None = None
231 self._stop_called: bool = False
232 self._daemon_task: asyncio.Task[None] | None = None
233 self._events_task: asyncio.Task[None] | None = None
234 self._refresh_task: asyncio.Task[None] | None = None
235 self._watcher_task: asyncio.Task[None] | None = None
236 self._proc: AsyncProcess | None = None
237 self._restart_error_count = 0
238 # set when a daemon close is intentional (sink replaced), so the
239 # supervisor respawns immediately instead of counting a failure
240 self._respawn_requested: bool = False
241 # serializes all volume handling (sink compensation + event forwarding)
242 self._volume_lock = asyncio.Lock()
243 # serializes the two-command repeat sequences of concurrent set_repeat calls
244 self._repeat_lock = asyncio.Lock()
245 # last volume reported by the daemon (None until the first event)
246 self._spotify_volume: int | None = None
247 # guards the player_only 100%-pin so overlapping resets are not issued
248 self._pin_in_flight: bool = False
249 # whether an auth_state ever reported a completed login; a fresh daemon
250 # reports logged_in=False while advertising for pairing, which is normal
251 self._was_logged_in: bool = False
252 self._last_context_uri: str | None = None
253 self._last_track_uri: str | None = None
254 # the capture sink delivers fixed s32le/44.1kHz/2ch PCM (the pulse
255 # capture format) â that is what actually arrives on the named pipe.
256 # Soloist decodes internally and never exposes the source codec or
257 # quality, so this capture format doubles as the display format: MA's
258 # input really is 32-bit PCM (24-bit lossless fits losslessly), while
259 # Spotify's upstream quality stays unknowable either way.
260 self._capture_format = AudioFormat(
261 content_type=ContentType.PCM_S32LE,
262 codec_type=ContentType.PCM_S32LE,
263 sample_rate=CAPTURE_SAMPLE_RATE,
264 bit_depth=32,
265 channels=CAPTURE_CHANNELS,
266 )
267 self._audio_format = self._capture_format
268
269 @property
270 def audio_format(self) -> AudioFormat:
271 """Return the source audio format (advertised to clients for display)."""
272 return self._audio_format
273
274 @property
275 def decoded_audio_format(self) -> AudioFormat:
276 """Return the PCM format the capture sink's named pipe delivers."""
277 return self._capture_format
278
279 @property
280 def stream_ends_on_pause(self) -> bool:
281 """The pipe sink delivers silence on pause; the provider stops the player."""
282 return False
283
284 @property
285 def supports_queue_control(self) -> bool:
286 """The soloist session implements the queue verbs and queue/options events."""
287 return True
288
289 async def start(self) -> None:
290 """Start the backend and its supervised soloist daemon."""
291 # setup errors (unsupported platform, missing consent, download failure,
292 # expired build) propagate so the provider load fails with a clear error
293 manager = SoloistBinaryManager(self.mass)
294 self._binary = await manager.ensure_fresh(self._consent)
295 self._build_sha = manager.diagnostics().get("sha256")
296 self._server = await get_pulse_capture_server(self.mass).acquire()
297 try:
298 await self._ensure_fresh_sink()
299
300 def _prepare_dirs() -> None:
301 self._data_dir.mkdir(parents=True, exist_ok=True)
302 # the data dir persists the Spotify device identity and login
303 # session; keep it readable by the MA user only
304 self._data_dir.chmod(0o700)
305 self._cache_dir.mkdir(parents=True, exist_ok=True)
306
307 await asyncio.to_thread(_prepare_dirs)
308 self._client = SoloistClient(self.mass, self._data_dir, self.logger)
309 # Two self-healing supervisors: one keeps the daemon process alive,
310 # the other keeps the events websocket connected (reconnecting
311 # across daemon restarts).
312 self._daemon_task = self.mass.create_task(self._daemon_runner())
313 self._events_task = self.mass.create_task(self._events_runner())
314 # Two housekeeping loops: a daily binary refresh (builds expire 90
315 # days after their build date) and a watcher replacing the capture
316 # sink after a pulse daemon restart.
317 self._refresh_task = self.mass.create_task(self._binary_refresh_loop())
318 self._watcher_task = self.mass.create_task(self._generation_watcher())
319 except BaseException:
320 # a failed startup aborts the provider load before unload() would
321 # ever run â release everything acquired so far ourselves
322 with suppress(Exception):
323 await self.stop()
324 raise
325
326 async def stop(self) -> None:
327 """Stop the daemon, its supervisors and the capture resources (idempotent)."""
328 self._stop_called = True
329 for task in (self._events_task, self._daemon_task, self._refresh_task, self._watcher_task):
330 if task and not task.done():
331 task.cancel()
332 with suppress(asyncio.CancelledError):
333 await task
334 self._events_task = None
335 self._daemon_task = None
336 self._refresh_task = None
337 self._watcher_task = None
338 if (proc := self._proc) is not None:
339 self._proc = None
340 await proc.close()
341 # the lock lets an in-flight _ensure_fresh_sink finish before the
342 # capture resources go away (later callers fail on the stop flag)
343 async with self._sink_lock:
344 if (sink := self._sink) is not None:
345 self._sink = None
346 await sink.unload()
347 if (server := self._server) is not None:
348 self._server = None
349 await server.release()
350
351 async def get_stream_source(self) -> BackendStreamSource:
352 """
353 Return the NAMED_PIPE stream source delivering the capture sink's PCM.
354
355 Side-effect-free (it also runs from queue preload): the sink is only
356 read here â (re)creation is owned by the daemon supervisor and the
357 generation watcher.
358
359 ffmpeg reads the FIFO directly and is the single pacing owner:
360 ``-readrate 1`` paces the read to realtime with a small initial burst
361 as jitter headroom â nothing else may sleep-pace this audio path.
362
363 :raises AudioError: No usable capture sink is currently available.
364 """
365 sink = self._sink
366 server = self._server
367 if sink is None or server is None or self._sink_generation != server.generation:
368 raise AudioError("Spotify Connect capture sink is not available")
369 return BackendStreamSource(
370 stream_type=StreamType.NAMED_PIPE,
371 path=str(sink.fifo_path),
372 extra_input_args=["-readrate", "1", "-readrate_initial_burst", "0.5"],
373 )
374
375 def get_audio_reader(self) -> AudioChunkReader | None:
376 """Return None: the audio is delivered through the NAMED_PIPE stream source."""
377 return None
378
379 async def play(self, uri: str, *, skip_to_uri: str | None = None) -> None:
380 """
381 Start playing a Spotify URI/context, making this device the active one.
382
383 :param uri: Spotify URI (track, album, playlist, ...) â typically a context.
384 :param skip_to_uri: Ignored â soloist's play command has no skip-to-track
385 option, so playback starts at the beginning of the context.
386 """
387 assert self._client is not None
388 # a bare play starts local playback without a Connect transfer, leaving
389 # the Spotify apps unaware; claim active device status first
390 await self._client.activate(await_result=True)
391 await self._client.play(uri)
392 if skip_to_uri:
393 self.logger.debug(
394 "skip_to_uri is not supported by soloist; starting %s from the beginning", uri
395 )
396
397 async def resume(self) -> None:
398 """Resume playback on the active session."""
399 assert self._client is not None
400 # re-claim active device status first (idempotent when already active):
401 # a resume after a deactivate would otherwise start local playback
402 # without a Connect transfer, leaving the Spotify apps unaware
403 await self._client.activate(await_result=True)
404 await self._client.resume()
405
406 async def pause(self) -> None:
407 """Pause playback on the active session."""
408 assert self._client is not None
409 await self._client.pause()
410
411 async def deactivate(self) -> None:
412 """Release this device as the active Spotify Connect device."""
413 assert self._client is not None
414 # pause first so the session's resume position is preserved; tolerate
415 # a rejected or unacknowledged pause and still give up the device
416 with suppress(SoloistError, TimeoutError):
417 await self._client.pause(await_result=True)
418 await self._client.deactivate()
419
420 async def next(self) -> None:
421 """Skip to the next track."""
422 assert self._client is not None
423 await self._client.skip_next()
424
425 async def previous(self) -> None:
426 """Skip to the previous track (or rewind the current one)."""
427 assert self._client is not None
428 await self._client.skip_prev()
429
430 async def seek(self, position_ms: int) -> None:
431 """
432 Seek to an absolute position in the current track.
433
434 :param position_ms: Target position in milliseconds.
435 """
436 assert self._client is not None
437 await self._client.seek(position_ms)
438
439 async def set_volume(self, volume: int) -> None:
440 """
441 Set the Spotify-side playback volume.
442
443 :param volume: Absolute volume as a 0-100 percentage. In player-only
444 mode the value is ignored and the daemon stays pinned at 100%.
445 """
446 assert self._client is not None
447 if self._volume_mode == VOLUME_MODE_SYNC_SPOTIFY:
448 # the daemon echoes the change back as a volume_changed event; the
449 # provider's dedupe keeps that echo from bouncing back to the player
450 async with self._volume_lock:
451 await self._client.set_volume(volume)
452 return
453 # player_only: the MA player owns the volume, so the daemon is pinned at
454 # 100% (unity PCM into the sink); only (re)send the pin when it is known
455 # (or possibly) off
456 async with self._volume_lock:
457 if self._spotify_volume == 100:
458 return
459 await self._client.set_volume(100)
460
461 async def add_to_queue(self, uri: str) -> None:
462 """
463 Add a track to the session's play queue.
464
465 :param uri: Spotify track URI to queue.
466 """
467 assert self._client is not None
468 await self._client.add_to_queue(uri)
469
470 async def set_shuffle(self, enabled: bool) -> None:
471 """
472 Enable or disable shuffle on the active session.
473
474 :param enabled: True to enable shuffle, False to disable it.
475 """
476 assert self._client is not None
477 await self._client.set_shuffle(enabled)
478
479 async def set_repeat(self, repeat: RepeatMode) -> None:
480 """
481 Set the repeat mode on the active session.
482
483 Awaits the engine's acknowledgement of both underlying commands, so
484 this call can block and raise; never call it from the backend event
485 callback (the acknowledgements arrive on the same loop and the wait
486 could only time out).
487
488 :param repeat: OFF for no repeat, ONE for the current track, ALL for
489 the playing context.
490 """
491 assert self._client is not None
492 if repeat == RepeatMode.UNKNOWN:
493 raise ValueError("cannot apply an unknown repeat mode")
494 # soloist models repeat as two independent booleans (track/context) with
495 # both-true undefined: the lock serializes concurrent callers and the
496 # awaited acks order the pair, so one flag is always cleared before the
497 # other is raised
498 async with self._repeat_lock:
499 if repeat == RepeatMode.ONE:
500 await self._client.set_repeat_context(False, await_result=True)
501 await self._client.set_repeat_track(True, await_result=True)
502 return
503 await self._client.set_repeat_track(False, await_result=True)
504 await self._client.set_repeat_context(repeat == RepeatMode.ALL, await_result=True)
505
506 async def request_queue(self, limit: int = 10) -> None:
507 """
508 Ask the session to (re)emit its queue view (arrives as a QUEUE_CHANGED event).
509
510 :param limit: Maximum number of upcoming entries the snapshot should
511 include.
512 """
513 assert self._client is not None
514 await self._client.get_queue(limit)
515
516 async def _ensure_fresh_sink(self) -> PipeSink:
517 """
518 Return the capture sink, replacing it when the pulse daemon restarted.
519
520 Only called from the supervisor paths (daemon spawn), never from the
521 side-effect-free stream request.
522
523 :raises AudioError: The backend is (being) stopped.
524 """
525 async with self._sink_lock:
526 # checked under the lock: a concurrent stop() sets the flag before
527 # it waits for the lock to tear down the capture resources
528 if self._stop_called or self._server is None:
529 raise AudioError("Spotify Connect backend is stopping")
530 if self._sink is not None and self._sink_generation == self._server.generation:
531 return self._sink
532 if self._sink is not None:
533 # sinks do not survive a daemon restart; this only cleans up the FIFO
534 await self._sink.unload()
535 self._sink = await PipeSink.create(self._server, self._sink_prefix)
536 self._sink_generation = self._server.generation
537 # the daemon plays into the sink named in its spawn env (PULSE_SINK),
538 # so a running process must respawn against the recreated sink
539 await self._close_daemon_for_respawn()
540 return self._sink
541
542 async def _recover_sink(self) -> None:
543 """
544 Fail-closed recovery: drop the sink and daemon so the supervisor rebuilds both.
545
546 Used when the current sink can no longer be trusted (its compensation
547 state is unknown after a failed volume call, or the pulse daemon
548 restarted underneath it): audio through such a sink risks a stale
549 reciprocal gain of up to 100x, so both the sink and the daemon are
550 torn down and recreated by the daemon supervisor.
551 """
552 async with self._sink_lock:
553 if self._stop_called:
554 return
555 if (sink := self._sink) is not None:
556 self._sink = None
557 with suppress(Exception):
558 await sink.unload()
559 await self._close_daemon_for_respawn()
560
561 async def _close_daemon_for_respawn(self) -> None:
562 """Close a running daemon so its supervisor respawns it right away (not a failure)."""
563 if (proc := self._proc) is not None:
564 self._respawn_requested = True
565 await proc.close()
566
567 async def _generation_watcher(self) -> None:
568 """Proactively replace the capture sink when the pulse daemon restarted."""
569 while True:
570 await asyncio.sleep(GENERATION_WATCH_INTERVAL_S)
571 server = self._server
572 if server is None or self._sink is None:
573 continue
574 if self._sink_generation != server.generation:
575 self.logger.info("Pulse daemon restart detected; recreating the capture sink")
576 await self._recover_sink()
577
578 def _daemon_args(self) -> list[str]:
579 """
580 Build the soloist daemon's argv.
581
582 SECURITY: the argv carries the user's API key â it must never be logged
583 or end up in any error message.
584 """
585 assert self._binary is not None
586 return [
587 str(self._binary),
588 "--device-name",
589 self._publish_name,
590 "--api-key",
591 self._api_key,
592 "--data-dir",
593 str(self._data_dir),
594 "--cache-dir",
595 str(self._cache_dir),
596 # bounded playback cache (0 would be unlimited)
597 "--cache-size",
598 str(CACHE_SIZE_MB),
599 # start at 100%: MA (or the sink compensation) owns the real volume
600 "--initial-volume",
601 "100",
602 # local WebSocket API on a free loopback port; the daemon publishes
603 # the actual endpoint in its data dir where SoloistClient finds it
604 "--ws",
605 "127.0.0.1:0",
606 ]
607
608 def _write_audio_prefs(self) -> None:
609 """
610 Write the configured audio behavior into the engine's prefs stores (blocking).
611
612 Both the global prefs and every existing per-user prefs file are updated:
613 per-user values override the global ones per key, and a per-user file only
614 appears after an account paired â the global store covers that account's
615 first session until the next daemon (re)spawn refreshes both.
616
617 Best-effort: a write failure must never block playback, so errors are
618 logged per store and the daemon spawns with the engine's previous settings.
619 """
620 settings_dir = self._data_dir / "settings"
621 prefs_files = [settings_dir / "prefs"]
622 quality = _QUALITY_VALUES.get(self._audio_quality, _QUALITY_VALUES[AUDIO_QUALITY_LOSSLESS])
623 managed_lines = [
624 f"{_PREF_CROSSFADE}={'true' if self._crossfade_ms else 'false'}",
625 f"{_PREF_NORMALIZE}={'true' if self._loudness_normalization else 'false'}",
626 f"{_PREF_QUALITY}={quality}",
627 f"{_PREF_QUALITY_NON_METERED}={quality}",
628 f"{_PREF_QUALITY_MIGRATED}=true",
629 ]
630 if self._crossfade_ms:
631 managed_lines.insert(1, f"{_PREF_CROSSFADE_TIME}={self._crossfade_ms}")
632 try:
633 users_dir = settings_dir / "Users"
634 if users_dir.is_dir():
635 prefs_files += [
636 user_dir / "prefs" for user_dir in users_dir.iterdir() if user_dir.is_dir()
637 ]
638 except OSError as err:
639 self.logger.warning("Failed to list the Spotify per-user settings: %s", err)
640 for prefs_file in prefs_files:
641 try:
642 lines = []
643 if prefs_file.is_file():
644 lines = [
645 line
646 for line in prefs_file.read_text(encoding="utf-8").splitlines()
647 if line.split("=", 1)[0] not in _MANAGED_PREFS
648 ]
649 prefs_file.parent.mkdir(parents=True, exist_ok=True)
650 # the stores also carry engine-owned keys, so replace atomically
651 # (a truncated in-place write would lose those too)
652 tmp_file = prefs_file.with_suffix(".tmp")
653 tmp_file.write_text("\n".join([*lines, *managed_lines]) + "\n", encoding="utf-8")
654 tmp_file.replace(prefs_file)
655 except (OSError, UnicodeDecodeError) as err:
656 self.logger.warning(
657 "Failed to write the Spotify audio settings to %s: %s", prefs_file, err
658 )
659
660 async def _daemon_runner(self) -> None:
661 """Run and supervise the soloist daemon, restarting (and refreshing) as needed."""
662 # Loop forever; stop() cancels this task and the explicit stop-check below
663 # handles a graceful exit without a restart.
664 while True:
665 proc: AsyncProcess | None = None
666 returncode: int | None = None
667 try:
668 # the sink may have been replaced (pulse restart) while the
669 # daemon was down; never spawn against a stale sink
670 sink = await self._ensure_fresh_sink()
671 server = self._server
672 assert server is not None # guaranteed by _ensure_fresh_sink
673 # the engine only reads its prefs at startup (and scrubs foreign
674 # keys from the global store when it rewrites it), so refresh the
675 # audio settings on every spawn, while the daemon is down
676 await asyncio.to_thread(self._write_audio_prefs)
677 # the explicit process name keeps AsyncProcess logging free of
678 # the argv (which carries the API key)
679 self._proc = proc = AsyncProcess(
680 self._daemon_args(),
681 # the daemon writes all of its logging to stdout and only
682 # ever puts argument-parsing complaints on stderr, so the
683 # two are merged into one captured stream. Capturing is
684 # what makes the redaction below reachable at all: an
685 # unset stdout is inherited, which would leak the daemon's
686 # output straight to the server console instead.
687 stdout=True,
688 stderr=asyncio.subprocess.STDOUT,
689 name=f"soloist[{self.name}]",
690 env=server.child_env(sink.sink_name),
691 )
692 await proc.start()
693 self.logger.info("Started Spotify Connect background daemon [%s]", self.name)
694 await self._reset_volume_state(sink)
695 await self._await_daemon_exit(proc)
696 except asyncio.CancelledError:
697 raise
698 except Exception as err:
699 self.logger.warning("soloist daemon error [%s]: %s", self.name, err)
700 finally:
701 if proc:
702 await proc.close()
703 returncode = proc.returncode
704 # The daemon â and thus the Spotify session â is gone. Tell the
705 # provider so a dead/restarting daemon isn't treated as active and
706 # controllable; a fresh 'active' event re-establishes it on reconnect.
707 self._proc = None
708 try:
709 await self._event_callback(BackendEvent(BackendEventType.CONNECTION_LOST))
710 except Exception:
711 # never let a callback error replace a propagating
712 # cancellation or kill the daemon supervisor
713 self.logger.exception("Error while handling daemon exit")
714 if self._stop_called:
715 break
716 if self._respawn_requested:
717 # intentional close (the sink was replaced): respawn right away,
718 # this is not a daemon failure
719 self._respawn_requested = False
720 continue
721 self.logger.info("Spotify Connect background daemon stopped for %s", self.name)
722 if returncode == EXIT_CODE_BUILD_EXPIRED and not await self._refresh_expired_binary():
723 return
724 self._restart_error_count += 1
725 if self._restart_error_count >= MAX_RESTART_ATTEMPTS:
726 await self._event_callback(
727 BackendEvent(
728 BackendEventType.FATAL_ERROR,
729 # fatal errors are plain (non-localized) strings for now,
730 # matching the go-librespot backend
731 error="soloist daemon failed to start multiple times.",
732 )
733 )
734 return
735 await asyncio.sleep(RESTART_DELAY_S)
736
737 async def _await_daemon_exit(self, proc: AsyncProcess) -> None:
738 """
739 Wait for the daemon to exit, forwarding its log for as long as it runs.
740
741 :param proc: The running daemon process.
742 """
743 # The log is drained by a side task rather than inline: a close() from
744 # one of the other supervisors (sink replacement, binary refresh) locks
745 # readers out of the process streams for good, so waiting on the reader
746 # would hang here and the daemon would never be respawned.
747 log_task = asyncio.create_task(self._log_daemon_output(proc))
748 wait_task = asyncio.create_task(proc.wait())
749 try:
750 # Watch both: nothing else drains the daemon's stdout, so a reader
751 # that died would leave the daemon blocked on a full pipe and this
752 # wait would never return.
753 await asyncio.wait((wait_task, log_task), return_when=FIRST_COMPLETED)
754 reader_error = (
755 log_task.exception() if log_task.done() and not log_task.cancelled() else None
756 )
757 if reader_error is not None:
758 self.logger.error(
759 "soloist log reader failed [%s]: %s; restarting the daemon",
760 self.name,
761 reader_error,
762 )
763 await proc.close()
764 await wait_task
765 # an exited daemon still has its last (often most telling) lines in
766 # the stream buffer; the shield keeps the reader alive across the
767 # timeout so it can drain them
768 with suppress(TimeoutError):
769 await asyncio.wait_for(asyncio.shield(log_task), DAEMON_LOG_DRAIN_TIMEOUT_S)
770 finally:
771 # a reader locked out by a close() from another supervisor never
772 # ends on its own; joining it consumes its outcome the way the
773 # other process readers in the codebase do
774 for task in (wait_task, log_task):
775 task.cancel()
776 for task in (wait_task, log_task):
777 with suppress(asyncio.CancelledError, Exception):
778 await task
779
780 async def _log_daemon_output(self, proc: AsyncProcess) -> None:
781 """
782 Forward the daemon's log lines to our logger until its output ends.
783
784 :param proc: The running daemon process.
785 """
786 async for line in proc.iter_stdout():
787 # the third-party binary's own output may echo argv (which carries
788 # the api key), so redact it before logging
789 text = line.replace(self._api_key, "<redacted>") if self._api_key else line
790 self.logger.debug("[%s] %s", self.name, text)
791
792 async def _reset_volume_state(self, sink: PipeSink) -> None:
793 """
794 Realign the volume bookkeeping with a freshly spawned daemon.
795
796 :param sink: The capture sink the daemon plays into.
797 """
798 # the daemon always starts at --initial-volume 100; without this reset a
799 # stale reciprocal sink compensation from before a crash would amplify
800 # and clip the captured audio (sync_spotify mode)
801 async with self._volume_lock:
802 self._spotify_volume = 100
803 try:
804 await sink.set_volume(100)
805 except Exception as err:
806 # fail closed: a stale reciprocal gain may still be active on
807 # the sink â recreate sink and daemon rather than clipping
808 self.logger.warning(
809 "Failed to reset capture sink volume (%s); recreating capture sink", err
810 )
811 await self._recover_sink()
812
813 async def _refresh_expired_binary(self) -> bool:
814 """
815 Replace the expired soloist build before the next daemon restart.
816
817 :return: True when the supervisor may restart the daemon, False when a
818 fatal error was reported and the supervisor must stop.
819 """
820 self.logger.warning("soloist build expired; looking for a replacement build")
821 try:
822 # force: the daemon itself reported expiry, so the recently-verified
823 # fast path must not hand back the same binary
824 manager = SoloistBinaryManager(self.mass)
825 self._binary = await manager.ensure_fresh(self._consent, force=True)
826 self._build_sha = manager.diagnostics().get("sha256")
827 except BuildExpiredError:
828 await self._event_callback(
829 BackendEvent(
830 BackendEventType.FATAL_ERROR,
831 error=(
832 "The Spotify Soloist build expired and no replacement could be "
833 "installed. Check the server's internet connection and reload "
834 "this provider."
835 ),
836 )
837 )
838 return False
839 except SoloistError as err:
840 # transient refresh problem: keep restarting, bounded by the failure cap
841 self.logger.warning("Unable to refresh the expired soloist build: %s", err)
842 return True
843
844 async def _binary_refresh_loop(self) -> None:
845 """Periodically refresh the soloist binary ahead of its 90-day build expiry."""
846 while True:
847 await asyncio.sleep(BINARY_REFRESH_INTERVAL_S)
848 try:
849 manager = SoloistBinaryManager(self.mass)
850 # a replaced build keeps the same install path, so compare the
851 # install metadata's digest against the build this daemon runs
852 # (a sibling instance may have updated the shared install)
853 binary = await manager.ensure_fresh(self._consent)
854 new_sha = manager.diagnostics().get("sha256")
855 if binary == self._binary and new_sha == self._build_sha:
856 continue
857 self.logger.info("A fresh soloist build was installed; restarting the daemon")
858 self._binary = binary
859 self._build_sha = new_sha
860 await self._close_daemon_for_respawn()
861 except asyncio.CancelledError:
862 raise
863 except Exception as err:
864 self.logger.warning("Periodic soloist binary refresh failed: %s", err)
865
866 async def _events_runner(self) -> None:
867 """Keep the soloist events websocket connected, reconnecting as needed."""
868 assert self._client is not None
869 while not self._stop_called:
870 try:
871 if not await self._client.wait_until_ready():
872 await asyncio.sleep(RESTART_DELAY_S)
873 continue
874 await self._client.listen_events(self._handle_event)
875 except asyncio.CancelledError:
876 raise
877 except (TimeoutError, OSError, ClientError, SoloistError) as err:
878 # ordinary connection drop; the loop reconnects quietly
879 self.logger.debug("soloist events websocket dropped: %s", err)
880 except Exception:
881 # a defect in event handling must not kill the control plane,
882 # but unlike a connection drop it has to surface loudly
883 self.logger.exception("Unexpected error while handling soloist events")
884 if not self._stop_called:
885 await asyncio.sleep(RESTART_DELAY_S)
886
887 async def _handle_event(self, event: SoloistEvent) -> None:
888 """Adapt a raw soloist event and emit its normalized counterpart."""
889 self.logger.debug("Received %s event [%s]", event.type, self.name)
890 # A delivered event means the websocket â and thus the daemon â is
891 # healthy: reset the restart backoff counter the daemon supervisor uses.
892 # (The endpoint files the events runner polls can be stale leftovers
893 # from a previous run, so the reset cannot happen on wait_until_ready.)
894 self._restart_error_count = 0
895 if isinstance(event.data, SoloistVolumeChanged):
896 await self._handle_volume_changed(event.data.volume)
897 return
898 if (
899 isinstance(event.data, SoloistPlaybackState)
900 and event.data.volume is not None
901 and event.data.volume != self._spotify_volume
902 ):
903 # a playback_state snapshot (e.g. right after a websocket reconnect)
904 # carries the daemon's current volume; resync the pin/compensation
905 # before forwarding the playback event itself
906 await self._handle_volume_changed(event.data.volume)
907 if self._sink is None:
908 # fail-closed recovery tore down the capture path; drop this
909 # stale snapshot â the respawned daemon reports fresh state
910 return
911 if (
912 isinstance(event.data, SoloistPlaybackState)
913 and (item := event.data.item) is not None
914 and item.uri != self._last_track_uri
915 ):
916 # the snapshot describes a track we have no metadata for yet (an
917 # already-playing session at (re)connect); emit its metadata so it
918 # does not stay stale until the next track_changed
919 await self._event_callback(
920 self._make_event(BackendEventType.METADATA, metadata=_entity_metadata(item))
921 )
922 if isinstance(event.data, SoloistPlaybackState) and event.data.options is not None:
923 # a state snapshot/delta carrying the playback options doubles as an
924 # options report; emit a separate OPTIONS_CHANGED so consumers only
925 # ever need to watch one event type for options
926 await self._event_callback(
927 self._make_event(
928 BackendEventType.OPTIONS_CHANGED,
929 options=_backend_options(event.data.options),
930 )
931 )
932 await self._event_callback(self._translate_event(event))
933
934 async def _handle_volume_changed(self, volume: int) -> None:
935 """
936 Apply a Spotify-side volume change according to the configured volume mode.
937
938 :param volume: The reported volume as a 0-100 percentage.
939 """
940 async with self._volume_lock:
941 self._spotify_volume = volume
942 if self._volume_mode != VOLUME_MODE_SYNC_SPOTIFY:
943 # player_only: MA/the player owns the volume. Keep the daemon
944 # pinned at 100% so the captured PCM stays at unity gain, and
945 # never forward VOLUME events (they would fight the MA volume).
946 if volume != 100 and self._client is not None and not self._pin_in_flight:
947 self._pin_in_flight = True
948 try:
949 await self._client.set_volume(100)
950 except Exception as err:
951 self.logger.debug("Failed to reset soloist volume: %s", err)
952 # mark the volume unknown so the next snapshot
953 # reporting the same value still retries the pin
954 self._spotify_volume = None
955 finally:
956 self._pin_in_flight = False
957 return
958 # sync_spotify: the daemon attenuates the PCM it plays into the sink
959 # with Spotify's cubic volume curve; undo that with the reciprocal
960 # raw sink gain (sink_pct = 10000 / spotify_pct) so the FIFO always
961 # carries unity-gain audio, then forward the volume for the MA
962 # player to apply (subject to the provider's dedupe/grace policy).
963 # NOTE: the percentage reciprocal is the exact linear inverse
964 # because pulse's software volume is cubic in the percentage too
965 # (pa_sw_volume_to_linear(p) = (p/100)^3) â validated by capture
966 # measurement, do not "fix" this to an explicit cube.
967 if self._sink is not None:
968 # explicit zero handling: no reciprocal exists, silence the sink
969 sink_pct = 0.0 if volume <= 0 else round(10000 / volume, 2)
970 try:
971 await self._sink.set_volume(sink_pct)
972 except Exception as err:
973 # fail closed: the compensation state is now unknown, so
974 # recreate sink and daemon and do NOT forward the volume â
975 # the player must not adopt a value whose compensation
976 # never applied
977 self.logger.warning(
978 "Failed to set capture sink volume (%s); recreating capture sink", err
979 )
980 await self._recover_sink()
981 return
982 await self._event_callback(
983 self._make_event(BackendEventType.VOLUME, volume=max(0, volume))
984 )
985
986 def _translate_event(self, event: SoloistEvent) -> BackendEvent:
987 """Map a raw soloist event onto the normalized BackendEvent model."""
988 data = event.data
989 if isinstance(data, SoloistAuthState):
990 if not data.logged_in:
991 if self._was_logged_in:
992 # an established login was lost mid-session: real auth loss
993 self._was_logged_in = False
994 return self._make_event(BackendEventType.AUTH_REQUIRED)
995 # a fresh daemon reports logged_in=False while advertising for
996 # pairing; that is the normal pre-pairing state, not an auth loss
997 return self._make_event(BackendEventType.SESSION_INACTIVE)
998 self._was_logged_in = True
999 return self._make_event(
1000 BackendEventType.SESSION_ACTIVE
1001 if data.is_active
1002 else BackendEventType.SESSION_INACTIVE
1003 )
1004 if isinstance(data, SoloistDeviceChanged):
1005 return self._make_event(
1006 BackendEventType.SESSION_ACTIVE
1007 if data.is_active
1008 else BackendEventType.SESSION_INACTIVE
1009 )
1010 if isinstance(data, SoloistPlaybackState):
1011 # covers both the playback_state snapshot and the playback_changed delta
1012 self._cache_uris(track=data.item, context=data.context)
1013 return self._make_event(_STATUS_EVENTS.get(data.status, BackendEventType.OTHER))
1014 if isinstance(data, SoloistTrackChanged) and data.item is not None:
1015 self._cache_uris(track=data.item)
1016 return self._make_event(BackendEventType.METADATA, metadata=_entity_metadata(data.item))
1017 if isinstance(data, SoloistPositionSync):
1018 return self._make_event(
1019 BackendEventType.POSITION, position=data.position.position_ms // 1000
1020 )
1021 if isinstance(data, SoloistErrorMessage):
1022 return self._make_event(BackendEventType.ERROR, error=data.message)
1023 if isinstance(data, SoloistQueueChanged):
1024 return self._make_event(
1025 BackendEventType.QUEUE_CHANGED,
1026 queue=BackendQueueState(
1027 previous=_queue_entries(data.previous),
1028 upcoming=_queue_entries(data.upcoming),
1029 ),
1030 )
1031 if isinstance(data, SoloistOptionsChanged):
1032 return self._make_event(
1033 BackendEventType.OPTIONS_CHANGED, options=_backend_options(data.options)
1034 )
1035 if isinstance(data, SoloistContextChanged):
1036 self._cache_uris(context=data.context)
1037 # context changes, command acks and unknown events
1038 return self._make_event(BackendEventType.OTHER)
1039
1040 def _make_event(self, event_type: BackendEventType, **fields: Any) -> BackendEvent:
1041 """Build a BackendEvent carrying the latest known context/track uris."""
1042 return BackendEvent(
1043 event_type,
1044 context_uri=self._last_context_uri,
1045 track_uri=self._last_track_uri,
1046 **fields,
1047 )
1048
1049 def _cache_uris(
1050 self, *, track: SoloistEntity | None = None, context: SoloistEntity | None = None
1051 ) -> None:
1052 """Remember the latest context/track uris seen on the event stream."""
1053 if track is not None and track.uri:
1054 self._last_track_uri = track.uri
1055 if context is not None and context.uri:
1056 self._last_context_uri = context.uri
1057
1058
1059def _entity_metadata(item: SoloistEntity) -> BackendTrackMetadata:
1060 """
1061 Extract normalized track metadata from an entity's decorations.
1062
1063 Mapped from the observed 1.3.7 wire format: the artist lives under
1064 ``creators[].entity``, the album under ``parent.entity`` and artwork
1065 under ``visual_identity.cover[]`` â all traversed defensively since the
1066 decorations bag is extensible.
1067
1068 :param item: The soloist entity (track) to extract the metadata from.
1069 """
1070 decorations = item.decorations or {}
1071 playback = _as_dict(decorations.get("playback"))
1072 duration_ms = playback.get("duration_ms")
1073 return BackendTrackMetadata(
1074 track_uri=item.uri,
1075 title=_identity_title(decorations),
1076 artist=_creator_name(decorations.get("creators")),
1077 album=_nested_entity_name(decorations.get("parent")),
1078 image_url=_cover_url(decorations),
1079 duration=int(duration_ms) // 1000 if isinstance(duration_ms, int | float) else None,
1080 )
1081
1082
1083def _queue_entries(entries: list[SoloistQueueEntry]) -> list[BackendQueueEntry]:
1084 """
1085 Map soloist queue entries onto normalized queue entries.
1086
1087 Entries without a resolvable uri are skipped.
1088
1089 :param entries: The soloist queue entries (previous or upcoming listing).
1090 """
1091 result: list[BackendQueueEntry] = []
1092 for entry in entries:
1093 if entry.item is None or not entry.item.uri:
1094 continue
1095 result.append(
1096 BackendQueueEntry(
1097 uid=entry.uid,
1098 uri=entry.item.uri,
1099 source=QueueEntrySource(entry.source),
1100 name=_identity_title(entry.item.decorations or {}),
1101 )
1102 )
1103 return result
1104
1105
1106# soloist's repeat vocabulary mapped onto the MA repeat modes
1107_REPEAT_MODES: Final = {
1108 "off": RepeatMode.OFF,
1109 "context": RepeatMode.ALL,
1110 "track": RepeatMode.ONE,
1111}
1112
1113
1114def _backend_options(options: SoloistPlaybackOptions) -> BackendPlaybackOptions:
1115 """Map soloist playback options onto the normalized model."""
1116 return BackendPlaybackOptions(
1117 shuffle=options.shuffle,
1118 repeat=_REPEAT_MODES.get(options.repeat, RepeatMode.UNKNOWN),
1119 )
1120
1121
1122def _as_dict(value: Any) -> dict[str, Any]:
1123 """Return the value if it is a dict, an empty dict otherwise."""
1124 return value if isinstance(value, dict) else {}
1125
1126
1127def _as_str(value: Any) -> str | None:
1128 """Return the value if it is a non-empty string, None otherwise."""
1129 return value if isinstance(value, str) and value else None
1130
1131
1132def _identity_title(decorations: dict[str, Any]) -> str | None:
1133 """Return the display title from a decorations bag's identity."""
1134 identity = _as_dict(decorations.get("identity"))
1135 return _as_str(identity.get("name")) or _as_str(identity.get("title"))
1136
1137
1138def _entity_name(value: Any) -> str | None:
1139 """Return the display name of a nested entity-like decoration value."""
1140 if isinstance(value, str):
1141 return value or None
1142 if isinstance(value, dict):
1143 return _as_str(value.get("name")) or _as_str(_as_dict(value.get("identity")).get("name"))
1144 return None
1145
1146
1147def _nested_entity_name(value: Any) -> str | None:
1148 """Return the identity name of a ``{"entity": {...}}`` decoration (album parent)."""
1149 entity = _as_dict(_as_dict(value).get("entity"))
1150 return _entity_name(_as_dict(entity.get("decorations"))) or _entity_name(entity)
1151
1152
1153def _creator_name(value: Any) -> str | None:
1154 """Return the name of the first creator credit (``creators[].entity``)."""
1155 if isinstance(value, list) and value:
1156 return _nested_entity_name(value[0])
1157 return None
1158
1159
1160def _cover_url(decorations: dict[str, Any]) -> str | None:
1161 """
1162 Return the artwork url from ``visual_identity.cover[]``.
1163
1164 Prefers the large rendition; falls back to the last (largest) entry.
1165 """
1166 covers = _as_dict(decorations.get("visual_identity")).get("cover")
1167 if not isinstance(covers, list) or not covers:
1168 return None
1169 for entry in covers:
1170 if isinstance(entry, dict) and entry.get("size") == "large":
1171 if url := _as_str(entry.get("url")):
1172 return url
1173 for entry in reversed(covers):
1174 if isinstance(entry, dict) and (url := _as_str(entry.get("url"))):
1175 return url
1176 return None
1177