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