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