/
/
1"""
2Shared helpers for Spotify Soloist, Spotify's official headless Connect client for Linux.
3
4Owned by the Spotify Connect provider and deliberately provider-neutral: the
5Spotify music provider reuses these helpers instead of shipping a second
6implementation. It manages a single shared install of the ``soloist`` binary
7under the server's storage dir and offers a typed client for the daemon's local
8WebSocket API (see https://developer.spotify.com/documentation/soloist).
9
10SECURITY NOTE: the soloist daemon takes the user's personal API key on its
11command line. This module never sees or logs that key, and callers that manage
12the daemon process must equally never log its argv.
13"""
14
15from __future__ import annotations
16
17import asyncio
18import hashlib
19import logging
20import platform
21import re
22import shutil
23import tarfile
24import tempfile
25import time
26from collections import deque
27from collections.abc import Awaitable, Callable
28from contextlib import suppress
29from dataclasses import dataclass, field
30from datetime import UTC, datetime
31from http import HTTPStatus
32from pathlib import Path, PurePosixPath
33from typing import TYPE_CHECKING, Any, Final
34
35import aiofiles
36from aiohttp import ClientError, ClientTimeout, ClientWebSocketResponse, WSMsgType
37from mashumaro import DataClassDictMixin
38from mashumaro.exceptions import InvalidFieldValue, MissingField
39from music_assistant_models.errors import MusicAssistantError
40from yarl import URL
41
42from music_assistant.constants import MASS_LOGGER_NAME
43from music_assistant.helpers.json import json_dumps, json_loads
44from music_assistant.helpers.process import check_output
45
46if TYPE_CHECKING:
47 from music_assistant.mass import MusicAssistant
48
49LOGGER = logging.getLogger(f"{MASS_LOGGER_NAME}.providers.spotify_connect.soloist")
50
51# Official per-architecture release archives (docs: reference/downloads-and-updates).
52CDN_URL_TEMPLATE: Final[str] = "https://soloist-builds.spotifycdn.com/soloist_release_{arch}.tar.gz"
53
54# The daemon publishes its WebSocket endpoint as two small files in its data dir.
55WS_ADDR_FILE: Final[str] = "ws.addr"
56WS_PORT_FILE: Final[str] = "ws.port"
57
58# soloist exits with this code once its build passed the 90-day expiry.
59EXIT_CODE_BUILD_EXPIRED: Final[int] = 10
60
61# platform.machine() values mapped to the CDN artifact architecture;
62# anything else has no official build and is rejected.
63_MACHINE_TO_ARCH: Final[dict[str, str]] = {
64 "aarch64": "arm64",
65 "arm64": "arm64",
66 # the official arm32 build targets ARMv7; older ARM cores are unsupported
67 "armv7l": "arm32",
68 "armv8l": "arm32",
69 "x86_64": "x86_64",
70 "amd64": "x86_64",
71}
72
73# A download may redirect, but only within Spotify's own infrastructure.
74_TRUSTED_DOWNLOAD_DOMAINS: Final[tuple[str, ...]] = ("spotifycdn.com", "spotify.com")
75_MAX_REDIRECTS: Final[int] = 5
76_REDIRECT_STATUSES: Final[frozenset[int]] = frozenset({301, 302, 303, 307, 308})
77
78# Conservative caps so a misbehaving CDN response cannot fill the disk:
79# real release archives are only a few tens of MB.
80_MAX_ARCHIVE_SIZE: Final[int] = 64 * 1024 * 1024
81_MAX_EXTRACTED_SIZE: Final[int] = 128 * 1024 * 1024
82_DOWNLOAD_CHUNK_SIZE: Final[int] = 64 * 1024
83
84# Builds expire 90 days after their build date; start looking for a replacement
85# 14 days ahead of that so there is a comfortable update window.
86_BUILD_EXPIRY_SECONDS: Final[int] = 90 * 24 * 3600
87_REFRESH_AGE_SECONDS: Final[int] = 76 * 24 * 3600
88
89# ELF identity per architecture: (EI_CLASS, e_machine).
90_ELF_IDENT: Final[dict[str, tuple[int, int]]] = {
91 "arm64": (2, 0xB7), # 64-bit, EM_AARCH64
92 "arm32": (1, 0x28), # 32-bit, EM_ARM
93 "x86_64": (2, 0x3E), # 64-bit, EM_X86_64
94}
95
96# Every SoloistBinaryManager instance manages the same shared install paths, so
97# the single-flight install lock is shared process-wide as well.
98_INSTALL_LOCK: Final = asyncio.Lock()
99
100_VERSION_CMD_TIMEOUT: Final[float] = 10.0
101_DOWNLOAD_TIMEOUT: Final[float] = 300.0
102_HEAD_TIMEOUT: Final[float] = 30.0
103
104# --version output is free-form text (the docs give no schema); fish out a
105# version token and an ISO-like timestamp defensively.
106_VERSION_TOKEN_RE: Final[re.Pattern[str]] = re.compile(r"\bv?(\d+\.\d+(?:\.\d+)*)\b")
107_TIMESTAMP_RE: Final[re.Pattern[str]] = re.compile(
108 r"\b(\d{4}-\d{2}-\d{2}(?:[T ]\d{2}:\d{2}(?::\d{2})?(?:\.\d+)?(?:Z|[+-]\d{2}:?\d{2})?)?)\b"
109)
110
111# The docs do not guarantee when the endpoint files appear, so poll for them.
112_ENDPOINT_POLL_INTERVAL: Final[float] = 0.25
113_WS_HEARTBEAT: Final[float] = 30.0
114_COMMAND_RESULT_TIMEOUT: Final[float] = 10.0
115
116
117class SoloistError(MusicAssistantError):
118 """Base error for all soloist helper failures."""
119
120
121class ConsentRequiredError(SoloistError):
122 """A download from Spotify's CDN is needed but the user did not consent (yet)."""
123
124
125class UnsupportedPlatformError(SoloistError):
126 """No official soloist build exists for this platform/architecture."""
127
128
129class DownloadFailedError(SoloistError):
130 """The soloist release archive could not be downloaded."""
131
132
133class InvalidArchiveError(SoloistError):
134 """The downloaded release archive or its binary failed validation."""
135
136
137class BuildExpiredError(SoloistError):
138 """The soloist build passed its 90-day expiry and no replacement is available."""
139
140
141@dataclass
142class SoloistEntity(DataClassDictMixin):
143 """A Spotify entity (track/album/playlist/...) as used in item/context fields."""
144
145 uri: str
146 entity_type: str
147 # decorations is an extensible bag (identity.name, playback.duration_ms, ...)
148 decorations: dict[str, Any] = field(default_factory=dict)
149
150
151@dataclass
152class SoloistPosition(DataClassDictMixin):
153 """Playback position reference point, to be interpolated with speed over time."""
154
155 position_ms: int
156 timestamp_ms: int
157 speed: float = 1.0
158
159
160@dataclass
161class SoloistPlaybackOptions(DataClassDictMixin):
162 """Playback options (shuffle/repeat/speed)."""
163
164 shuffle: bool = False
165 repeat: str = "off"
166 playback_speed: float = 1.0
167 modes: dict[str, Any] = field(default_factory=dict)
168
169
170@dataclass
171class SoloistQueueEntry(DataClassDictMixin):
172 """One entry in the previous/upcoming queue listing."""
173
174 uid: str
175 source: str
176 item: SoloistEntity | None = None
177
178
179@dataclass
180class SoloistAuthState(DataClassDictMixin):
181 """Payload of the ``auth_state`` event."""
182
183 logged_in: bool
184 is_active: bool
185 device_name: str | None = None
186
187
188@dataclass
189class SoloistPlaybackState(DataClassDictMixin):
190 """Payload of the ``playback_state`` snapshot and the ``playback_changed`` delta."""
191
192 status: str
193 item: SoloistEntity | None = None
194 context: SoloistEntity | None = None
195 position: SoloistPosition | None = None
196 volume: int | None = None
197 is_active: bool | None = None
198 options: SoloistPlaybackOptions | None = None
199 available_actions: dict[str, Any] = field(default_factory=dict)
200
201
202@dataclass
203class SoloistTrackChanged(DataClassDictMixin):
204 """Payload of the ``track_changed`` event."""
205
206 item: SoloistEntity | None = None
207
208
209@dataclass
210class SoloistPositionSync(DataClassDictMixin):
211 """Payload of the ``position_sync`` event."""
212
213 position: SoloistPosition
214
215
216@dataclass
217class SoloistVolumeChanged(DataClassDictMixin):
218 """Payload of the ``volume_changed`` event."""
219
220 volume: int
221
222
223@dataclass
224class SoloistDeviceChanged(DataClassDictMixin):
225 """Payload of the ``device_changed`` event."""
226
227 is_active: bool
228
229
230@dataclass
231class SoloistContextChanged(DataClassDictMixin):
232 """Payload of the ``context_changed`` event."""
233
234 context: SoloistEntity | None = None
235
236
237@dataclass
238class SoloistOptionsChanged(DataClassDictMixin):
239 """Payload of the ``options_changed`` event."""
240
241 options: SoloistPlaybackOptions
242
243
244@dataclass
245class SoloistQueueChanged(DataClassDictMixin):
246 """Payload of the ``queue_changed`` event."""
247
248 previous: list[SoloistQueueEntry] = field(default_factory=list)
249 upcoming: list[SoloistQueueEntry] = field(default_factory=list)
250
251
252@dataclass
253class SoloistCommandResult(DataClassDictMixin):
254 """Payload of the ``command_result`` acknowledgement event."""
255
256 command: str
257
258
259@dataclass
260class SoloistErrorMessage(DataClassDictMixin):
261 """Payload of the ``error`` event."""
262
263 message: str
264
265
266SoloistEventData = (
267 SoloistAuthState
268 | SoloistPlaybackState
269 | SoloistTrackChanged
270 | SoloistPositionSync
271 | SoloistVolumeChanged
272 | SoloistDeviceChanged
273 | SoloistContextChanged
274 | SoloistOptionsChanged
275 | SoloistQueueChanged
276 | SoloistCommandResult
277 | SoloistErrorMessage
278)
279
280# Maps documented event types to their payload model;
281# unrecognized types are passed through as generic raw events.
282_EVENT_MODELS: Final[dict[str, type[SoloistEventData]]] = {
283 "auth_state": SoloistAuthState,
284 "playback_state": SoloistPlaybackState,
285 "playback_changed": SoloistPlaybackState,
286 "track_changed": SoloistTrackChanged,
287 "position_sync": SoloistPositionSync,
288 "volume_changed": SoloistVolumeChanged,
289 "device_changed": SoloistDeviceChanged,
290 "context_changed": SoloistContextChanged,
291 "options_changed": SoloistOptionsChanged,
292 "queue_changed": SoloistQueueChanged,
293 "command_result": SoloistCommandResult,
294 "error": SoloistErrorMessage,
295}
296
297
298@dataclass
299class SoloistEvent:
300 """A single event received from the daemon, with its decoded payload when recognized."""
301
302 type: str
303 data: SoloistEventData | None
304 raw: dict[str, Any]
305
306
307# Called with the decoded event for every WebSocket event.
308EventCallback = Callable[[SoloistEvent], Awaitable[None]]
309
310
311class SoloistBinaryManager:
312 """
313 Manages the single shared soloist binary install for all consumers.
314
315 The install lives under ``<storage_path>/soloist``; concurrent callers
316 share one download.
317 """
318
319 def __init__(self, mass: MusicAssistant) -> None:
320 """
321 Initialize the binary manager.
322
323 :param mass: The MusicAssistant instance (for its HTTP session and storage path).
324 """
325 self.mass = mass
326 self._install_dir = Path(mass.storage_path) / "soloist"
327 self._binary_path = self._install_dir / "soloist"
328 self._previous_path = self._install_dir / "soloist.prev"
329 self._metadata_path = self._install_dir / "soloist.meta.json"
330
331 @property
332 def binary_path(self) -> Path:
333 """Path where the soloist binary is (or will be) installed."""
334 return self._binary_path
335
336 async def ensure_binary(self, consent: bool) -> Path:
337 """
338 Return the path to a validated soloist binary, downloading it when needed.
339
340 :param consent: Whether the user consented to downloading the binary from
341 Spotify's CDN. An already-installed valid binary is returned without
342 any network access, regardless of this flag.
343 :raises ConsentRequiredError: A download is needed but consent was not given.
344 :raises UnsupportedPlatformError: No soloist build exists for this platform.
345 :raises DownloadFailedError: The release archive could not be downloaded.
346 :raises InvalidArchiveError: The downloaded archive or binary failed validation.
347 :raises BuildExpiredError: The freshly downloaded build has already expired.
348 """
349 arch = _resolve_architecture()
350 async with _INSTALL_LOCK:
351 if await self._installed_returncode() == 0 and not self._installed_expired():
352 return self._binary_path
353 await self._install_with_consent(consent, arch)
354 return self._binary_path
355
356 async def ensure_fresh(self, consent: bool) -> Path:
357 """
358 Return a validated soloist binary, refreshing it when it is (close to) expiry.
359
360 Builds expire 90 days after their build date: an installed build that
361 already expired is replaced immediately, one nearing expiry only when
362 the CDN offers a different build. When the refresh fails (e.g. offline)
363 a still-valid binary is returned with a warning instead.
364
365 :param consent: Whether the user consented to downloading from Spotify's
366 CDN. Without consent a still-valid installed binary is returned
367 as-is (no proactive refresh).
368 :raises ConsentRequiredError: A download is needed but consent was not given.
369 :raises UnsupportedPlatformError: No soloist build exists for this platform.
370 :raises DownloadFailedError: No usable binary is installed and the download failed.
371 :raises InvalidArchiveError: The downloaded archive or binary failed validation.
372 :raises BuildExpiredError: The installed build expired and no valid
373 replacement could be obtained.
374 """
375 arch = _resolve_architecture()
376 async with _INSTALL_LOCK:
377 returncode = await self._installed_returncode()
378 if returncode not in (0, EXIT_CODE_BUILD_EXPIRED):
379 # missing or broken install: plain (re)install
380 await self._install_with_consent(consent, arch)
381 return self._binary_path
382 expired = returncode == EXIT_CODE_BUILD_EXPIRED or self._installed_expired()
383 if not expired and (not consent or not self._due_for_update()):
384 return self._binary_path
385 if expired and not consent:
386 raise ConsentRequiredError(
387 "Updating the expired soloist binary requires user consent"
388 )
389 if not expired and not await self._update_available(arch):
390 return self._binary_path
391 try:
392 await self._download_and_install(arch)
393 except SoloistError as err:
394 if expired:
395 if isinstance(err, DownloadFailedError):
396 raise BuildExpiredError(
397 "soloist build expired and no replacement could be downloaded"
398 ) from err
399 raise
400 LOGGER.warning("soloist update failed, keeping the current binary: %s", err)
401 return self._binary_path
402
403 def diagnostics(self) -> dict[str, Any]:
404 """
405 Return diagnostic details about the installed binary.
406
407 Contains install/build metadata only; never any key or session material.
408 """
409 metadata = self._read_metadata()
410 if metadata is None:
411 return {"installed": False}
412 return {
413 "installed": True,
414 "sha256": metadata.sha256,
415 "etag": metadata.etag,
416 "version": metadata.version,
417 "version_raw": metadata.version_raw,
418 "installed_at": metadata.installed_at,
419 "build_timestamp": metadata.build_timestamp,
420 # upper-bound estimate when the build timestamp could not be parsed
421 "expires_at": (metadata.build_timestamp or metadata.installed_at)
422 + _BUILD_EXPIRY_SECONDS,
423 }
424
425 async def _install_with_consent(self, consent: bool, arch: str) -> None:
426 """Download and install the binary, requiring user consent first."""
427 if not consent:
428 raise ConsentRequiredError("Downloading the soloist binary requires user consent")
429 await self._download_and_install(arch)
430
431 async def _installed_returncode(self) -> int | None:
432 """Return the ``--version`` exit code of the installed binary, or None if unrunnable."""
433 if not self._binary_path.is_file():
434 return None
435 try:
436 returncode, _ = await check_output(
437 str(self._binary_path), "--version", timeout=_VERSION_CMD_TIMEOUT
438 )
439 except OSError, TimeoutError:
440 return None
441 return returncode
442
443 def _installed_expired(self) -> bool:
444 """
445 Return whether the installed build's known build timestamp has expired.
446
447 Defense in depth next to the exit-code check: whether ``--version``
448 itself reports expiry is not documented, so an expired-by-timestamp
449 build is refused even when the binary still runs. The install time is
450 the fallback anchor — the build is at least as old as its install.
451 """
452 metadata = self._read_metadata()
453 return metadata is not None and _build_expired(
454 metadata.build_timestamp or metadata.installed_at
455 )
456
457 def _due_for_update(self) -> bool:
458 """Return whether the installed build is old enough to look for an update."""
459 metadata = self._read_metadata()
460 if metadata is None:
461 return True
462 anchor = metadata.build_timestamp or metadata.installed_at
463 return (time.time() - anchor) >= _REFRESH_AGE_SECONDS
464
465 async def _update_available(self, arch: str) -> bool:
466 """Return whether the CDN offers a different build than the installed one."""
467 metadata = self._read_metadata()
468 try:
469 remote_etag = await self._fetch_remote_etag(arch)
470 except (ClientError, TimeoutError, OSError, DownloadFailedError) as err:
471 # any failed update check keeps the current (still valid) build
472 LOGGER.warning("Unable to check for a soloist update: %s", err)
473 return False
474 if remote_etag is None:
475 # CDN unreachable/erroring: keep the current build
476 return False
477 if not remote_etag:
478 # reachable but no validator offered: assume an update exists and
479 # let the validated download path decide
480 return True
481 return metadata is None or metadata.etag != remote_etag
482
483 async def _fetch_remote_etag(self, arch: str) -> str | None:
484 """Return the ETag the CDN currently serves for the given architecture, if any."""
485 current_url = CDN_URL_TEMPLATE.format(arch=arch)
486 # follow redirects manually so every hop passes the same host allowlist
487 # the download path enforces
488 for _ in range(_MAX_REDIRECTS + 1):
489 _validate_download_url(current_url)
490 async with self.mass.http_session.head(
491 current_url, allow_redirects=False, timeout=ClientTimeout(total=_HEAD_TIMEOUT)
492 ) as resp:
493 if resp.status in _REDIRECT_STATUSES:
494 if not (location := resp.headers.get("Location")):
495 return None
496 current_url = str(URL(current_url).join(URL(location)))
497 continue
498 if resp.status != HTTPStatus.OK:
499 return None
500 # empty string = reachable but no validator offered
501 return resp.headers.get("ETag", "")
502 return None
503
504 async def _download_and_install(self, arch: str) -> None:
505 """Download, validate and atomically install the binary for the given architecture."""
506 url = CDN_URL_TEMPLATE.format(arch=arch)
507 try:
508 await asyncio.to_thread(self._install_dir.mkdir, parents=True, exist_ok=True)
509 temp_dir = Path(
510 await asyncio.to_thread(tempfile.mkdtemp, prefix=".soloist-", dir=self._install_dir)
511 )
512 except OSError as err:
513 raise DownloadFailedError(f"cannot prepare the soloist install dir: {err}") from err
514 try:
515 archive_path = temp_dir / "soloist.tar.gz"
516 new_binary = temp_dir / "soloist"
517 sha256, etag = await self._download_archive(url, archive_path)
518 await asyncio.to_thread(_extract_binary_from_archive, archive_path, new_binary)
519 await asyncio.to_thread(_validate_elf_header, new_binary, arch)
520 version_raw = await self._swap_in_and_validate(new_binary)
521 metadata = _BinaryMetadata(
522 sha256=sha256,
523 version_raw=version_raw,
524 installed_at=time.time(),
525 etag=etag,
526 version=_parse_version_token(version_raw),
527 build_timestamp=_parse_build_timestamp(version_raw),
528 )
529 # metadata is advisory (diagnostics + refresh hints): the binary is
530 # already validated and installed, so never fail the install over it
531 try:
532 await asyncio.to_thread(self._write_metadata, metadata)
533 except OSError as err:
534 LOGGER.warning("Unable to persist soloist install metadata: %s", err)
535 # stale metadata must not outlive the swap: it would pair the
536 # fresh binary with the previous build's expiry state
537 with suppress(OSError):
538 await asyncio.to_thread(self._metadata_path.unlink, missing_ok=True)
539 LOGGER.info("Installed soloist binary %s (%s)", metadata.version or "unknown", arch)
540 finally:
541 await asyncio.to_thread(shutil.rmtree, temp_dir, ignore_errors=True)
542
543 async def _download_archive(self, url: str, dest: Path) -> tuple[str, str | None]:
544 """
545 Stream the release archive to a local file, following (validated) redirects.
546
547 :return: Tuple of the archive's sha256 hexdigest and its ETag header (if any).
548 """
549 hasher = hashlib.sha256()
550 current_url = url
551 try:
552 for _ in range(_MAX_REDIRECTS + 1):
553 _validate_download_url(current_url)
554 async with self.mass.http_session.get(
555 current_url,
556 allow_redirects=False,
557 timeout=ClientTimeout(total=_DOWNLOAD_TIMEOUT),
558 ) as resp:
559 if resp.status in _REDIRECT_STATUSES:
560 location = resp.headers.get("Location")
561 if not location:
562 raise DownloadFailedError("soloist download redirect has no location")
563 current_url = str(URL(current_url).join(URL(location)))
564 continue
565 if resp.status != HTTPStatus.OK:
566 raise DownloadFailedError(
567 f"soloist download failed with HTTP {resp.status}"
568 )
569 etag = resp.headers.get("ETag")
570 size = 0
571 async with aiofiles.open(dest, "wb") as _file:
572 async for chunk in resp.content.iter_chunked(_DOWNLOAD_CHUNK_SIZE):
573 size += len(chunk)
574 if size > _MAX_ARCHIVE_SIZE:
575 raise DownloadFailedError("soloist archive exceeds the size limit")
576 hasher.update(chunk)
577 await _file.write(chunk)
578 return (hasher.hexdigest(), etag)
579 raise DownloadFailedError("too many redirects while downloading soloist")
580 except (ClientError, TimeoutError, OSError) as err:
581 raise DownloadFailedError(f"soloist download failed: {err}") from err
582
583 async def _swap_in_and_validate(self, new_binary: Path) -> str:
584 """
585 Atomically install the new binary and verify it actually runs.
586
587 Shielded from cancellation so the install always reaches a consistent
588 end state (committed or rolled back) even when the caller goes away.
589
590 :return: The raw ``--version`` output of the installed binary.
591 """
592 inner = asyncio.ensure_future(self._swap_in_and_validate_inner(new_binary))
593 try:
594 return await asyncio.shield(inner)
595 except asyncio.CancelledError:
596 # wait for the shielded install to reach its consistent end state
597 # (committed or rolled back) so the install lock and temp dir stay
598 # held until the shared install is no longer being mutated
599 while not inner.done():
600 with suppress(asyncio.CancelledError):
601 await asyncio.shield(inner)
602 _consume_install_result(inner)
603 raise
604
605 async def _swap_in_and_validate_inner(self, new_binary: Path) -> str:
606 """Perform the swap/validate/rollback sequence (see _swap_in_and_validate)."""
607
608 def _swap_in() -> None:
609 new_binary.chmod(0o755)
610 if self._binary_path.exists():
611 self._binary_path.replace(self._previous_path)
612 new_binary.replace(self._binary_path)
613
614 def _rollback() -> None:
615 if self._previous_path.exists():
616 self._previous_path.replace(self._binary_path)
617 else:
618 self._binary_path.unlink(missing_ok=True)
619
620 def _commit() -> None:
621 self._previous_path.unlink(missing_ok=True)
622
623 try:
624 await asyncio.to_thread(_swap_in)
625 except OSError as err:
626 # keep SoloistError semantics so ensure_fresh's keep-the-old-binary
627 # fallback also covers filesystem failures during the swap
628 with suppress(OSError):
629 await asyncio.to_thread(_rollback)
630 raise InvalidArchiveError(f"failed to install the soloist binary: {err}") from err
631 try:
632 returncode, output = await check_output(
633 str(self._binary_path), "--version", timeout=_VERSION_CMD_TIMEOUT
634 )
635 except (OSError, TimeoutError) as err:
636 with suppress(OSError):
637 await asyncio.to_thread(_rollback)
638 raise InvalidArchiveError(f"soloist binary failed to run: {err}") from err
639 if returncode == EXIT_CODE_BUILD_EXPIRED:
640 with suppress(OSError):
641 await asyncio.to_thread(_rollback)
642 raise BuildExpiredError("the downloaded soloist build has already expired")
643 if returncode != 0:
644 with suppress(OSError):
645 await asyncio.to_thread(_rollback)
646 raise InvalidArchiveError(f"soloist --version exited with code {returncode}")
647 version_raw = output.decode("utf-8", errors="replace").strip()
648 if _build_expired(_parse_build_timestamp(version_raw)):
649 with suppress(OSError):
650 await asyncio.to_thread(_rollback)
651 raise BuildExpiredError("the downloaded soloist build has already expired")
652 await asyncio.to_thread(_commit)
653 return version_raw
654
655 def _read_metadata(self) -> _BinaryMetadata | None:
656 """Return the persisted install metadata, if present and readable."""
657 try:
658 raw = json_loads(self._metadata_path.read_text(encoding="utf-8"))
659 if not isinstance(raw, dict):
660 return None
661 return _BinaryMetadata.from_dict(raw)
662 except OSError, ValueError, TypeError, MissingField, InvalidFieldValue:
663 return None
664
665 def _write_metadata(self, metadata: _BinaryMetadata) -> None:
666 """Persist install metadata next to the binary."""
667 self._metadata_path.write_text(json_dumps(metadata.to_dict()), encoding="utf-8")
668
669
670class SoloistClient:
671 """
672 Client for the local WebSocket API of a running soloist daemon.
673
674 The daemon publishes its (loopback) endpoint as ``ws.addr``/``ws.port``
675 files in its data directory; this client discovers it from there. Commands
676 are sent over the events connection, so :meth:`listen_events` must be
677 running for the command senders to work.
678 """
679
680 def __init__(self, mass: MusicAssistant, data_dir: Path, logger: logging.Logger) -> None:
681 """
682 Initialize the client.
683
684 :param mass: The MusicAssistant instance (for its shared HTTP session).
685 :param data_dir: The daemon's data directory (holds the endpoint files).
686 :param logger: Logger to use for diagnostics.
687 """
688 self.mass = mass
689 self.data_dir = data_dir
690 self.logger = logger
691 self._ws: ClientWebSocketResponse | None = None
692 self._pending_results: dict[str, deque[asyncio.Future[None]]] = {}
693
694 @property
695 def connected(self) -> bool:
696 """Whether the events WebSocket is currently connected."""
697 return self._ws is not None and not self._ws.closed
698
699 async def wait_until_ready(self, timeout: float = 30.0) -> bool:
700 """
701 Poll the daemon's data directory until the WebSocket endpoint is published.
702
703 :param timeout: Maximum seconds to wait for the endpoint files.
704 :return: True once the endpoint is known, False if the timeout elapses.
705 """
706 loop = asyncio.get_running_loop()
707 deadline = loop.time() + timeout
708 while True:
709 if await asyncio.to_thread(self._read_endpoint) is not None:
710 return True
711 if loop.time() >= deadline:
712 return False
713 await asyncio.sleep(_ENDPOINT_POLL_INTERVAL)
714
715 async def listen_events(self, on_event: EventCallback) -> None:
716 """
717 Connect to the daemon's WebSocket API and dispatch events until it closes.
718
719 Returns normally when the connection is closed by either side; raises on
720 connection errors so the caller can implement a reconnect loop.
721 Malformed or unrecognized events are tolerated and never interrupt the
722 stream.
723
724 :param on_event: Coroutine called with a :class:`SoloistEvent` per event.
725 :raises SoloistError: The daemon has not published its endpoint (yet).
726 """
727 if self._ws is not None and not self._ws.closed:
728 # pending acks are correlated per connection; a second concurrent
729 # listener would cross-resolve them
730 raise SoloistError("listen_events is already running")
731 if (endpoint := await asyncio.to_thread(self._read_endpoint)) is None:
732 raise SoloistError("soloist has not published its WebSocket endpoint (yet)")
733 addr, port = endpoint
734 host = f"[{addr}]" if ":" in addr else addr
735 ws: ClientWebSocketResponse | None = None
736 try:
737 async with self.mass.http_session.ws_connect(
738 f"ws://{host}:{port}/", heartbeat=_WS_HEARTBEAT
739 ) as ws:
740 self._ws = ws
741 self.logger.debug("Connected to the soloist websocket at %s:%s", addr, port)
742 async for msg in ws:
743 if msg.type == WSMsgType.TEXT:
744 await self._handle_message(msg.data, on_event)
745 elif msg.type in (WSMsgType.CLOSE, WSMsgType.CLOSING, WSMsgType.CLOSED):
746 break
747 elif msg.type == WSMsgType.ERROR:
748 raise ws.exception() or ClientError("websocket error")
749 finally:
750 # only tear down our own connection state: a reconnecting caller may
751 # already have a newer listen_events connection registered
752 if ws is not None and self._ws is ws:
753 self._ws = None
754 self._fail_pending_results()
755
756 async def play(self, uri: str | None = None, *, await_result: bool = False) -> None:
757 """
758 Start playback of a Spotify URI/context, or resume the current one.
759
760 :param uri: Spotify URI (track/album/playlist/...); omit to resume.
761 :param await_result: Wait for the daemon's command acknowledgement.
762 """
763 fields: dict[str, Any] = {"uri": uri} if uri is not None else {}
764 await self._send_command("play", await_result=await_result, **fields)
765
766 async def resume(self, *, await_result: bool = False) -> None:
767 """Resume playback of the current context."""
768 await self._send_command("play", await_result=await_result)
769
770 async def pause(self, *, await_result: bool = False) -> None:
771 """Pause playback."""
772 await self._send_command("pause", await_result=await_result)
773
774 async def skip_next(self, *, await_result: bool = False) -> None:
775 """Skip to the next track."""
776 await self._send_command("skip_next", await_result=await_result)
777
778 async def skip_prev(self, *, await_result: bool = False) -> None:
779 """Skip to the previous track (or rewind the current one)."""
780 await self._send_command("skip_prev", await_result=await_result)
781
782 async def seek(self, position_ms: int, *, await_result: bool = False) -> None:
783 """
784 Seek to an absolute position in the current track.
785
786 :param position_ms: Target position in milliseconds.
787 """
788 await self._send_command("seek", await_result=await_result, position_ms=max(0, position_ms))
789
790 async def set_volume(self, volume: int, *, await_result: bool = False) -> None:
791 """
792 Set the absolute playback volume.
793
794 :param volume: Volume from 0 to 100.
795 """
796 await self._send_command(
797 "set_volume", await_result=await_result, volume=max(0, min(100, volume))
798 )
799
800 async def activate(self, *, await_result: bool = False) -> None:
801 """Make this daemon the active Spotify Connect device."""
802 await self._send_command("activate", await_result=await_result)
803
804 async def deactivate(self, *, await_result: bool = False) -> None:
805 """Release this daemon as the active Spotify Connect device."""
806 await self._send_command("deactivate", await_result=await_result)
807
808 async def set_shuffle(self, enabled: bool, *, await_result: bool = False) -> None:
809 """Enable or disable shuffle."""
810 await self._send_command("set_shuffle", await_result=await_result, enabled=enabled)
811
812 async def set_repeat_context(self, enabled: bool, *, await_result: bool = False) -> None:
813 """Enable or disable repeating the current context."""
814 await self._send_command("set_repeat_context", await_result=await_result, enabled=enabled)
815
816 async def set_repeat_track(self, enabled: bool, *, await_result: bool = False) -> None:
817 """Enable or disable repeating the current track."""
818 await self._send_command("set_repeat_track", await_result=await_result, enabled=enabled)
819
820 async def add_to_queue(self, uri: str, *, await_result: bool = False) -> None:
821 """
822 Add a track to the play queue.
823
824 :param uri: Spotify track URI.
825 """
826 await self._send_command("add_to_queue", await_result=await_result, uri=uri)
827
828 async def get_state(self) -> None:
829 """Request a full ``playback_state`` snapshot (answered as an event)."""
830 await self._send_command("get_state")
831
832 async def get_auth_state(self) -> None:
833 """Request the current ``auth_state`` (answered as an event)."""
834 await self._send_command("get_auth_state")
835
836 async def get_queue(self, limit: int = 10) -> None:
837 """
838 Request the play queue (answered as a ``queue_changed`` event).
839
840 :param limit: Maximum number of upcoming tracks to return.
841 """
842 await self._send_command("get_queue", limit=limit)
843
844 async def _send_command(
845 self,
846 command: str,
847 *,
848 await_result: bool = False,
849 timeout: float = _COMMAND_RESULT_TIMEOUT,
850 **fields: Any,
851 ) -> None:
852 """
853 Send a command frame, optionally waiting for its ``command_result`` ack.
854
855 Acks carry only the command name (no request id), so concurrent calls of
856 the same command resolve in FIFO order with no per-call correlation, and
857 an ack for a fire-and-forget send of the same command can resolve an
858 awaited call early — avoid mixing both styles for one command.
859 Never wait for an ack from inside an ``on_event`` callback: the ack is
860 delivered by the same event loop, so the wait can only time out.
861
862 :raises SoloistError: The WebSocket is not connected (or closed while waiting).
863 :raises TimeoutError: The acknowledgement did not arrive within the timeout.
864 """
865 ws = self._ws
866 if ws is None or ws.closed:
867 raise SoloistError("soloist websocket is not connected")
868 future: asyncio.Future[None] | None = None
869 if await_result:
870 future = asyncio.get_running_loop().create_future()
871 self._pending_results.setdefault(command, deque()).append(future)
872 try:
873 await ws.send_json({"type": "command", "command": command, **fields})
874 if future is not None:
875 await asyncio.wait_for(future, timeout)
876 finally:
877 if future is not None:
878 self._discard_pending_result(command, future)
879
880 async def _handle_message(self, raw: str, on_event: EventCallback) -> None:
881 """Decode a single WebSocket text frame and dispatch it to the callback."""
882 try:
883 payload = json_loads(raw)
884 except ValueError:
885 self.logger.debug("Ignoring non-JSON websocket message: %s", raw)
886 return
887 if not isinstance(payload, dict) or not isinstance(event_type := payload.get("type"), str):
888 self.logger.debug("Ignoring websocket message without an event type: %s", raw)
889 return
890 data: SoloistEventData | None = None
891 if (model := _EVENT_MODELS.get(event_type)) is not None:
892 try:
893 data = model.from_dict(payload)
894 except Exception as err:
895 # tolerance to malformed events is a hard requirement: whatever
896 # the decode error, log it and keep the stream alive
897 self.logger.debug("Ignoring malformed %s event: %s (%s)", event_type, raw, err)
898 return
899 if event_type == "command_result" and isinstance(data, SoloistCommandResult):
900 self._resolve_pending_result(data.command)
901 await on_event(SoloistEvent(type=event_type, data=data, raw=payload))
902
903 def _read_endpoint(self) -> tuple[str, int] | None:
904 """Read the WebSocket endpoint published in the daemon's data directory."""
905 try:
906 addr = (self.data_dir / WS_ADDR_FILE).read_text(encoding="utf-8").strip()
907 port = int((self.data_dir / WS_PORT_FILE).read_text(encoding="utf-8").strip())
908 except OSError, ValueError:
909 return None
910 if not addr or not 0 < port <= 65535:
911 return None
912 return (addr, port)
913
914 def _resolve_pending_result(self, command: str) -> None:
915 """Resolve the oldest result waiter for the given command, if any."""
916 if (queue := self._pending_results.get(command)) is None:
917 return
918 while queue:
919 future = queue.popleft()
920 if not future.done():
921 future.set_result(None)
922 break
923 if not queue:
924 self._pending_results.pop(command, None)
925
926 def _discard_pending_result(self, command: str, future: asyncio.Future[None]) -> None:
927 """Drop a result waiter that is no longer interested in its result."""
928 if (queue := self._pending_results.get(command)) is None:
929 return
930 with suppress(ValueError):
931 queue.remove(future)
932 if not queue:
933 self._pending_results.pop(command, None)
934
935 def _fail_pending_results(self) -> None:
936 """Fail all outstanding result waiters (connection lost)."""
937 for queue in self._pending_results.values():
938 for future in queue:
939 if not future.done():
940 future.set_exception(SoloistError("websocket connection closed"))
941 self._pending_results.clear()
942
943
944@dataclass
945class _BinaryMetadata(DataClassDictMixin):
946 """Install metadata persisted next to the soloist binary."""
947
948 sha256: str
949 version_raw: str
950 installed_at: float
951 etag: str | None = None
952 version: str | None = None
953 build_timestamp: float | None = None
954
955
956def _resolve_architecture() -> str:
957 """Return the CDN artifact architecture for the current platform."""
958 system = platform.system()
959 if system != "Linux":
960 raise UnsupportedPlatformError(f"soloist is only available for Linux, not {system}")
961 machine = platform.machine().lower()
962 if (arch := _MACHINE_TO_ARCH.get(machine)) is None:
963 raise UnsupportedPlatformError(f"no soloist build is available for {machine}")
964 return arch
965
966
967def _validate_download_url(url: str) -> None:
968 """Reject download URLs that are insecure or outside Spotify's infrastructure."""
969 parsed = URL(url)
970 host = (parsed.host or "").lower()
971 if parsed.scheme != "https" or not host:
972 raise DownloadFailedError(f"refusing insecure soloist download url: {url}")
973 if not any(
974 host == domain or host.endswith(f".{domain}") for domain in _TRUSTED_DOWNLOAD_DOMAINS
975 ):
976 raise DownloadFailedError(f"refusing soloist download from untrusted host: {host}")
977
978
979def _extract_binary_from_archive(archive_path: Path, dest_path: Path) -> None:
980 """
981 Validate the release archive and extract its single soloist binary (blocking).
982
983 :raises InvalidArchiveError: The archive is corrupt, unsafe, or does not
984 contain exactly one regular ``soloist`` file.
985 """
986 try:
987 with tarfile.open(archive_path, mode="r:gz") as tar:
988 binary_member: tarfile.TarInfo | None = None
989 total_size = 0
990 for member in tar:
991 name = PurePosixPath(member.name)
992 if name.is_absolute() or ".." in name.parts:
993 raise InvalidArchiveError(f"unsafe path in soloist archive: {member.name}")
994 if member.isdir():
995 continue
996 if not member.isreg():
997 raise InvalidArchiveError(
998 f"unsupported member type in soloist archive: {member.name}"
999 )
1000 total_size += member.size
1001 if total_size > _MAX_EXTRACTED_SIZE:
1002 raise InvalidArchiveError("soloist archive exceeds the extracted size limit")
1003 if name.name != "soloist":
1004 raise InvalidArchiveError(f"unexpected file in soloist archive: {member.name}")
1005 if binary_member is not None:
1006 raise InvalidArchiveError("soloist archive contains multiple binaries")
1007 binary_member = member
1008 if binary_member is None:
1009 raise InvalidArchiveError("soloist archive does not contain the soloist binary")
1010 source = tar.extractfile(binary_member)
1011 if source is None:
1012 raise InvalidArchiveError("soloist archive member is not extractable")
1013 with source, dest_path.open("wb") as dest:
1014 shutil.copyfileobj(source, dest, _DOWNLOAD_CHUNK_SIZE)
1015 except (tarfile.TarError, EOFError, OSError) as err:
1016 raise InvalidArchiveError(f"invalid soloist archive: {err}") from err
1017
1018
1019def _validate_elf_header(binary_path: Path, arch: str) -> None:
1020 """
1021 Verify the extracted binary is a little-endian ELF for the expected architecture.
1022
1023 :raises InvalidArchiveError: The file is not an ELF binary for the given arch.
1024 """
1025 with binary_path.open("rb") as _file:
1026 header = _file.read(20)
1027 if len(header) < 20 or header[:4] != b"\x7fELF":
1028 raise InvalidArchiveError("soloist binary is not an ELF executable")
1029 # all supported targets are little-endian (EI_DATA == 1)
1030 if header[5] != 1:
1031 raise InvalidArchiveError("soloist binary has an unexpected ELF byte order")
1032 ei_class, e_machine = _ELF_IDENT[arch]
1033 if header[4] != ei_class or int.from_bytes(header[18:20], "little") != e_machine:
1034 raise InvalidArchiveError(f"soloist binary does not match architecture {arch}")
1035
1036
1037def _consume_install_result(fut: asyncio.Future[str]) -> None:
1038 """Log the outcome of an install that finished after its caller was cancelled."""
1039 if (exc := fut.exception()) is not None and not isinstance(exc, asyncio.CancelledError):
1040 LOGGER.warning("soloist install finished with an error after cancellation: %s", exc)
1041
1042
1043def _build_expired(build_timestamp: float | None) -> bool:
1044 """Return whether a build timestamp (when known) has passed the 90-day expiry."""
1045 return build_timestamp is not None and build_timestamp + _BUILD_EXPIRY_SECONDS <= time.time()
1046
1047
1048def _parse_version_token(version_output: str) -> str | None:
1049 """Best-effort extraction of a version number from free-form ``--version`` output."""
1050 match = _VERSION_TOKEN_RE.search(version_output)
1051 return match.group(1) if match else None
1052
1053
1054def _parse_build_timestamp(version_output: str) -> float | None:
1055 """Best-effort extraction of the build timestamp from free-form ``--version`` output."""
1056 for match in _TIMESTAMP_RE.finditer(version_output):
1057 candidate = match.group(1).replace("Z", "+00:00")
1058 try:
1059 parsed = datetime.fromisoformat(candidate)
1060 except ValueError:
1061 continue
1062 if parsed.tzinfo is None:
1063 parsed = parsed.replace(tzinfo=UTC)
1064 return parsed.timestamp()
1065 return None
1066