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