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