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