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