/
/
1"""Main Spotify provider implementation."""
2
3from __future__ import annotations
4
5import asyncio
6import os
7import time
8from collections import OrderedDict
9from collections.abc import AsyncGenerator, Sequence
10from dataclasses import dataclass
11from datetime import datetime
12from typing import Any, cast
13
14import aiohttp
15from music_assistant_models.config_entries import ConfigEntry
16from music_assistant_models.enums import (
17 ConfigEntryType,
18 ContentType,
19 ImageType,
20 MediaType,
21 ProviderFeature,
22 StreamType,
23)
24from music_assistant_models.errors import (
25 AudioError,
26 LoginFailed,
27 MediaNotFoundError,
28 ProviderUnavailableError,
29 RateLimited,
30 ResourceTemporarilyUnavailable,
31 UnsupportedFeaturedException,
32)
33from music_assistant_models.media_items import (
34 Album,
35 Artist,
36 Audiobook,
37 AudioFormat,
38 BrowseFolder,
39 ItemMapping,
40 MediaItemImage,
41 MediaItemType,
42 Playlist,
43 Podcast,
44 PodcastEpisode,
45 ProviderMapping,
46 SearchResults,
47 Track,
48 UniqueList,
49)
50from music_assistant_models.media_items.metadata import MediaItemChapter
51from music_assistant_models.streamdetails import StreamDetails
52from orjson import JSONDecodeError
53
54from music_assistant.constants import CONF_ENTRY_UNOFFICIAL_PROVIDER
55from music_assistant.controllers.cache import use_cache
56from music_assistant.helpers.app_vars import app_var
57from music_assistant.helpers.json import SerializableType, json_loads
58from music_assistant.helpers.throttle_retry import ThrottlerManager, throttle_with_retries
59from music_assistant.helpers.util import lock
60from music_assistant.models.music_provider import MusicProvider
61
62from .constants import (
63 CONF_CLIENT_ID,
64 CONF_LIBRESPOT_CREDENTIALS,
65 CONF_REFRESH_TOKEN_DEV,
66 CONF_REFRESH_TOKEN_GLOBAL,
67 CONF_SYNC_AUDIOBOOK_PROGRESS,
68 CONF_SYNC_PODCAST_PROGRESS,
69 CREDENTIALS_FILE,
70 LIKED_SONGS_FAKE_PLAYLIST_ID_PREFIX,
71)
72from .helpers import get_librespot_binary, get_spotify_token
73from .parsers import (
74 parse_album,
75 parse_artist,
76 parse_audiobook,
77 parse_playlist,
78 parse_podcast,
79 parse_podcast_episode,
80 parse_track,
81)
82from .streaming import LibrespotStreamer
83
84_PLAYLIST_PAGINATION_STATE_LIMIT = 32
85
86
87class NotModifiedError(Exception):
88 """Exception raised when a resource has not been modified."""
89
90
91@dataclass(slots=True)
92class _PlaylistPaginationState:
93 """Hold the synchronization and metadata snapshot for one playlist endpoint."""
94
95 lock: asyncio.Lock
96 snapshot: dict[str, Any] | None = None
97
98
99class SpotifyProvider(MusicProvider):
100 """Implementation of a Spotify MusicProvider."""
101
102 # Global session (MA's client ID) - always present
103 _auth_info_global: dict[str, Any] | None = None
104 # Developer session (user's custom client ID) - optional
105 _auth_info_dev: dict[str, Any] | None = None
106 _sp_user: dict[str, Any] | None = None
107 _librespot_bin: str | None = None
108 _audiobooks_supported = False
109 _playlist_pagination_states: OrderedDict[tuple[str, bool], _PlaylistPaginationState]
110 # True if user has configured a custom client ID with valid authentication
111 dev_session_active: bool = False
112 throttler: ThrottlerManager
113
114 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
115 """
116 Return Config entries to setup this provider.
117
118 Authentication is handled by the setup flow (see setup_flow.py); only the genuine
119 options are configurable here.
120 """
121 # audiobook progress sync is only offered where the account's region supports audiobooks
122 audiobooks_supported = bool(getattr(self, "audiobooks_supported", False))
123 return (
124 CONF_ENTRY_UNOFFICIAL_PROVIDER,
125 ConfigEntry(
126 key=CONF_SYNC_PODCAST_PROGRESS,
127 type=ConfigEntryType.BOOLEAN,
128 default_value=True,
129 category="sync_options",
130 ),
131 ConfigEntry(
132 key=CONF_SYNC_AUDIOBOOK_PROGRESS,
133 type=ConfigEntryType.BOOLEAN,
134 default_value=False,
135 category="sync_options",
136 hidden=not audiobooks_supported,
137 ),
138 )
139
140 async def handle_async_init(self) -> None:
141 """Handle async initialization of the provider."""
142 self.cache_dir = os.path.join(self.mass.cache_path, self.instance_id)
143 self._playlist_pagination_states = OrderedDict()
144 # Default throttler for global session (heavy rate limited)
145 self.throttler = ThrottlerManager(rate_limit=1, period=2)
146 self.streamer = LibrespotStreamer(self)
147
148 # check if we have a librespot binary for this arch
149 self._librespot_bin = await get_librespot_binary()
150 # playback authorization is independent of the Web API tokens
151 await self._setup_librespot_auth()
152 # try login which will raise if it fails (logs in global session)
153 await self.login()
154
155 # Check if user has a custom client ID with valid dev token
156 client_id = self.get_setup_value(CONF_CLIENT_ID)
157 dev_token = self.get_setup_value(CONF_REFRESH_TOKEN_DEV)
158
159 if client_id and dev_token and self._sp_user:
160 await self.login_dev()
161 # Verify user matches
162 userinfo = await self._get_data("me", use_global_session=False)
163 if userinfo["id"] != self._sp_user["id"]:
164 raise LoginFailed(
165 "Developer session must use the same Spotify account as the main session."
166 )
167 # loosen the throttler when a custom client id is used
168 self.throttler = ThrottlerManager(rate_limit=45, period=30)
169 self.dev_session_active = True
170 self.logger.info("Developer Spotify session active.")
171
172 self._audiobooks_supported = await self._test_audiobook_support()
173 if not self._audiobooks_supported:
174 self.logger.info(
175 "Audiobook support disabled: Audiobooks are not available in your region. "
176 "See https://support.spotify.com/us/authors/article/audiobooks-availability/ "
177 "for supported countries."
178 )
179
180 @property
181 def audiobooks_supported(self) -> bool:
182 """Check if audiobooks are supported for this user/region."""
183 return self._audiobooks_supported
184
185 @property
186 def audiobook_progress_sync_enabled(self) -> bool:
187 """Check if audiobook progress sync is enabled."""
188 return bool(self.config.get_value(CONF_SYNC_AUDIOBOOK_PROGRESS, False))
189
190 @property
191 def podcast_progress_sync_enabled(self) -> bool:
192 """Check if played status sync is enabled."""
193 value = self.config.get_value(CONF_SYNC_PODCAST_PROGRESS, True)
194 return bool(value) if value is not None else True
195
196 @property
197 def supported_features(self) -> set[ProviderFeature]:
198 """Return the features supported by this Provider."""
199 features = self._supported_features.copy()
200 # Add audiobook features if enabled
201 if self.audiobooks_supported:
202 features.add(ProviderFeature.LIBRARY_AUDIOBOOKS)
203 features.add(ProviderFeature.LIBRARY_AUDIOBOOKS_EDIT)
204 return features
205
206 @property
207 def instance_name_postfix(self) -> str | None:
208 """Return a (default) instance name postfix for this provider instance."""
209 if self._sp_user:
210 return str(self._sp_user["display_name"])
211 return None
212
213 async def get_diagnostics(self) -> dict[str, SerializableType]:
214 """Return diagnostics info for this provider to include in diagnostics reports."""
215 return {
216 "logged_in": self._sp_user is not None,
217 "token_expires_in_sec": (
218 round(self._auth_info_global["expires_at"] - time.time())
219 if self._auth_info_global
220 else None
221 ),
222 "dev_session_active": self.dev_session_active,
223 "librespot_available": self._librespot_bin is not None,
224 "audiobooks_supported": self._audiobooks_supported,
225 }
226
227 ## Library retrieval methods (generators)
228 async def get_library_artists(self) -> AsyncGenerator[Artist]:
229 """Retrieve library artists from spotify."""
230 endpoint = "me/following"
231 while True:
232 spotify_artists = await self._get_data(
233 endpoint,
234 type="artist",
235 limit=50,
236 )
237 for item in spotify_artists["artists"]["items"]:
238 if item and item["id"]:
239 yield parse_artist(item, self)
240 if spotify_artists["artists"]["next"]:
241 endpoint = spotify_artists["artists"]["next"]
242 endpoint = endpoint.replace("https://api.spotify.com/v1/", "")
243 else:
244 break
245
246 async def get_library_albums(self) -> AsyncGenerator[Album]:
247 """Retrieve library albums from the provider."""
248 async for item in self._get_all_items("me/albums"):
249 if item["album"] and item["album"]["id"]:
250 yield parse_album(item["album"], self)
251
252 async def get_library_tracks(self) -> AsyncGenerator[Track]:
253 """Retrieve library tracks from the provider."""
254 async for item in self._get_all_items("me/tracks"):
255 if item and item["track"]["id"]:
256 yield parse_track(item["track"], self)
257
258 async def get_library_podcasts(self) -> AsyncGenerator[Podcast]:
259 """Retrieve library podcasts from spotify."""
260 async for item in self._get_all_items("me/shows"):
261 if item["show"] and item["show"]["id"]:
262 show_obj = item["show"]
263 # Filter out audiobooks - they have a distinctive description format
264 description = show_obj.get("description", "")
265 if description.startswith("Author(s):") and "Narrator(s):" in description:
266 continue
267 yield parse_podcast(show_obj, self)
268
269 async def get_library_audiobooks(self) -> AsyncGenerator[Audiobook]:
270 """Retrieve library audiobooks from spotify."""
271 if not self.audiobooks_supported:
272 return
273 async for item in self._get_all_items("me/audiobooks"):
274 if item and item["id"]:
275 # Parse the basic audiobook
276 audiobook = parse_audiobook(item, self)
277 # Add chapters from Spotify API data
278 await self._add_audiobook_chapters(audiobook)
279 yield audiobook
280
281 async def get_library_playlists(self) -> AsyncGenerator[Playlist]:
282 """
283 Retrieve playlists from the provider.
284
285 Note: We use the global session here because playlists like "Daily Mix"
286 are only returned when using the non-dev (global) token.
287 """
288 yield await self._get_liked_songs_playlist()
289 async for item in self._get_all_items("me/playlists", use_global_session=True):
290 if item and item["id"]:
291 yield parse_playlist(item, self)
292
293 async def browse(self, path: str) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
294 """
295 Browse Spotify items, including curated sections (new releases, genres & moods).
296
297 :param path: The path to browse (e.g. provider_id:// or provider_id://new-releases).
298 """
299 path_parts = path.split("://")[1].split("/") if "://" in path else []
300 subpath = path_parts[0] if path_parts else None
301 sub_subpath = path_parts[1] if len(path_parts) > 1 else None
302 locale = self.mass.metadata.locale
303
304 if subpath == "new-releases":
305 return await self._get_new_releases()
306
307 if subpath == "categories" and sub_subpath:
308 return await self._get_category_playlists(sub_subpath, locale)
309
310 if subpath == "categories":
311 return await self._get_categories(locale)
312
313 # For root path, add curated folders on top of standard library folders.
314 # At the root the path always ends in "://", so curated paths can be appended directly.
315 if not subpath:
316 curated: list[BrowseFolder] = [
317 BrowseFolder(
318 item_id="new-releases",
319 provider=self.instance_id,
320 path=f"{path}new-releases",
321 name="New Releases",
322 translation_key="new_releases",
323 is_playable=True,
324 ),
325 BrowseFolder(
326 item_id="categories",
327 provider=self.instance_id,
328 path=f"{path}categories",
329 name="Genres & Moods",
330 translation_key="genres_and_moods",
331 is_playable=False,
332 ),
333 ]
334 standard = await super().browse(path)
335 return [*curated, *standard]
336
337 return await super().browse(path)
338
339 @use_cache()
340 async def search(
341 self, search_query: str, media_types: list[MediaType] | None = None, limit: int = 5
342 ) -> SearchResults:
343 """
344 Perform search on musicprovider.
345
346 :param search_query: Search query.
347 :param media_types: A list of media_types to include.
348 :param limit: Number of items to return in the search (per type).
349 """
350 searchresult = SearchResults()
351 if media_types is None:
352 return searchresult
353
354 searchtype = self._build_search_types(media_types)
355 if not searchtype:
356 return searchresult
357
358 search_query = search_query.replace("'", "")
359 offset = 0
360 page_limit = min(limit, 10)
361
362 while True:
363 api_result = await self._get_data(
364 "search", q=search_query, type=searchtype, limit=page_limit, offset=offset
365 )
366 items_received = self._process_search_results(api_result, searchresult)
367
368 offset += page_limit
369 if offset >= limit or items_received < page_limit:
370 break
371
372 return searchresult
373
374 @use_cache()
375 async def get_artist(self, prov_artist_id: str) -> Artist:
376 """Get full artist details by id."""
377 artist_obj = await self._get_data(f"artists/{prov_artist_id}")
378 return parse_artist(artist_obj, self)
379
380 @use_cache()
381 async def get_album(self, prov_album_id: str) -> Album:
382 """Get full album details by id."""
383 album_obj = await self._get_data(f"albums/{prov_album_id}")
384 return parse_album(album_obj, self)
385
386 @use_cache()
387 async def get_track(self, prov_track_id: str) -> Track:
388 """Get full track details by id."""
389 track_obj = await self._get_data(f"tracks/{prov_track_id}")
390 return parse_track(track_obj, self)
391
392 @use_cache()
393 async def get_playlist(self, prov_playlist_id: str) -> Playlist:
394 """Get full playlist details by id."""
395 if prov_playlist_id == self._get_liked_songs_playlist_id():
396 return await self._get_liked_songs_playlist()
397
398 # Check cache to see if this playlist requires global token
399 use_global = await self._playlist_requires_global_token(prov_playlist_id)
400 if use_global:
401 playlist_obj = await self._get_data(
402 f"playlists/{prov_playlist_id}", use_global_session=True
403 )
404 return parse_playlist(playlist_obj, self)
405
406 # Try with dev token first (if available), fallback to global on 400 error
407 # Some playlists like Spotify-owned (Daily Mix) or Liked Songs only work with global token
408 try:
409 playlist_obj = await self._get_data(f"playlists/{prov_playlist_id}")
410 return parse_playlist(playlist_obj, self)
411 except MediaNotFoundError:
412 if self.dev_session_active:
413 # Remember that this playlist requires global token
414 await self._set_playlist_requires_global_token(prov_playlist_id)
415 playlist_obj = await self._get_data(
416 f"playlists/{prov_playlist_id}", use_global_session=True
417 )
418 return parse_playlist(playlist_obj, self)
419 raise
420
421 @use_cache()
422 async def get_podcast(self, prov_podcast_id: str) -> Podcast:
423 """Get full podcast details by id."""
424 podcast_obj = await self._get_data(f"shows/{prov_podcast_id}")
425 if not podcast_obj:
426 raise MediaNotFoundError(f"Podcast not found: {prov_podcast_id}")
427 return parse_podcast(podcast_obj, self)
428
429 @use_cache()
430 async def get_audiobook(self, prov_audiobook_id: str) -> Audiobook:
431 """Get full audiobook details by id."""
432 if not self.audiobooks_supported:
433 raise UnsupportedFeaturedException("Audiobooks are not supported with this account")
434
435 audiobook_obj = await self._get_data(f"audiobooks/{prov_audiobook_id}")
436 if not audiobook_obj:
437 raise MediaNotFoundError(f"Audiobook not found: {prov_audiobook_id}")
438
439 # Parse basic audiobook without chapters first
440 audiobook = parse_audiobook(audiobook_obj, self)
441
442 # Add chapters from Spotify API data
443 await self._add_audiobook_chapters(audiobook)
444
445 # Note: Resume position will be handled by MA's internal system
446 # which calls get_resume_position() when needed
447
448 return audiobook
449
450 async def get_podcast_episodes(self, prov_podcast_id: str) -> AsyncGenerator[PodcastEpisode]:
451 """Get all podcast episodes."""
452 podcast = await self.get_podcast(prov_podcast_id)
453
454 # Get (cached) episode data
455 episodes_data = await self._get_podcast_episodes_data(prov_podcast_id)
456
457 # Parse and yield episodes with position
458 for idx, episode_data in enumerate(episodes_data):
459 episode = parse_podcast_episode(episode_data, self, podcast)
460 episode.position = idx + 1
461
462 # Set played status if sync is enabled and resume data exists
463 if self.podcast_progress_sync_enabled and "resume_point" in episode_data:
464 resume_point = episode_data["resume_point"]
465 fully_played = resume_point.get("fully_played", False)
466 position_ms = resume_point.get("resume_position_ms", 0)
467
468 episode.fully_played = fully_played or None
469 episode.resume_position_ms = position_ms if position_ms > 0 else None
470
471 yield episode
472
473 @use_cache(86400) # 24 hours
474 async def get_podcast_episode(self, prov_episode_id: str) -> PodcastEpisode:
475 """Get full podcast episode details by id."""
476 episode_obj = await self._get_data(f"episodes/{prov_episode_id}", market="from_token")
477 if not episode_obj:
478 raise MediaNotFoundError(f"Episode not found: {prov_episode_id}")
479 return parse_podcast_episode(episode_obj, self)
480
481 async def get_resume_position(
482 self, item_id: str, media_type: MediaType
483 ) -> tuple[bool, int, datetime | None]:
484 """Get resume position for episode/audiobook from Spotify."""
485 if media_type == MediaType.PODCAST_EPISODE:
486 if not self.podcast_progress_sync_enabled:
487 raise NotImplementedError("Spotify podcast resume sync disabled in settings")
488
489 try:
490 episode_obj = await self._get_data(f"episodes/{item_id}", market="from_token")
491 except MediaNotFoundError:
492 raise NotImplementedError("Episode not found on Spotify")
493 except (ResourceTemporarilyUnavailable, aiohttp.ClientError) as e:
494 self.logger.debug(f"Error fetching episode {item_id}: {e}")
495 raise NotImplementedError("Unable to fetch episode data from Spotify")
496
497 if (
498 not episode_obj
499 or "resume_point" not in episode_obj
500 or not episode_obj["resume_point"]
501 ):
502 raise NotImplementedError("No resume point data from Spotify")
503
504 resume_point = episode_obj["resume_point"]
505 fully_played = resume_point.get("fully_played", False)
506 position_ms = resume_point.get("resume_position_ms", 0)
507 return fully_played, position_ms, None
508
509 if media_type == MediaType.AUDIOBOOK:
510 if not self.audiobooks_supported:
511 raise NotImplementedError("Audiobook support is disabled")
512 if not self.audiobook_progress_sync_enabled:
513 raise NotImplementedError("Spotify audiobook resume sync disabled in settings")
514
515 try:
516 chapters_data = await self._get_audiobook_chapters_data(item_id)
517 if not chapters_data:
518 raise NotImplementedError("No chapters data available")
519
520 total_position_ms = 0
521 fully_played = True
522
523 for chapter in chapters_data:
524 resume_point = chapter.get("resume_point", {})
525 chapter_fully_played = resume_point.get("fully_played", False)
526 chapter_position_ms = resume_point.get("resume_position_ms", 0)
527
528 if chapter_fully_played:
529 total_position_ms += chapter.get("duration_ms", 0)
530 elif chapter_position_ms > 0:
531 total_position_ms += chapter_position_ms
532 fully_played = False
533 break
534 else:
535 fully_played = False
536 break
537
538 return fully_played, total_position_ms, None
539
540 except (MediaNotFoundError, ResourceTemporarilyUnavailable, aiohttp.ClientError) as e:
541 self.logger.debug(f"Failed to get audiobook resume position for {item_id}: {e}")
542 raise NotImplementedError("Unable to get audiobook resume position from Spotify")
543
544 else:
545 raise NotImplementedError(f"Resume position not supported for {media_type}")
546
547 async def on_played(
548 self,
549 media_type: MediaType,
550 prov_item_id: str,
551 fully_played: bool,
552 position: int,
553 media_item: MediaItemType,
554 is_playing: bool = False,
555 ) -> None:
556 """
557 Call when an episode/audiobook is played in MA.
558
559 MA automatically handles internal position tracking - this method is for
560 provider-specific actions like syncing to external services.
561 """
562 if media_type == MediaType.PODCAST_EPISODE:
563 if not isinstance(media_item, PodcastEpisode):
564 return
565
566 # Log the playback for monitoring/debugging
567 safe_position = position or 0
568 if media_item.duration > 0:
569 completion_percentage = (safe_position / media_item.duration) * 100
570 else:
571 completion_percentage = 0
572
573 self.logger.debug(
574 f"Episode played: {prov_item_id} at {safe_position}s "
575 f"({completion_percentage:.1f}%, fully_played: {fully_played})"
576 )
577
578 # Note: No API exists to sync playback position back to Spotify for episodes
579 # MA handles all internal position tracking automatically
580
581 elif media_type == MediaType.AUDIOBOOK:
582 if not isinstance(media_item, Audiobook):
583 return
584
585 # Log the playback for monitoring/debugging
586 safe_position = position or 0
587 if media_item.duration > 0:
588 completion_percentage = (safe_position / media_item.duration) * 100
589 else:
590 completion_percentage = 0
591
592 self.logger.debug(
593 f"Audiobook played: {prov_item_id} at {safe_position}s "
594 f"({completion_percentage:.1f}%, fully_played: {fully_played})"
595 )
596
597 # Note: No API exists to sync playback position back to Spotify for audiobooks
598 # MA handles all internal position tracking automatically
599
600 # The resume position will be automatically updated by MA's internal tracking
601 # and will be retrieved via get_audiobook() which combines MA + Spotify positions
602
603 @use_cache(86400 * 365, allow_expired_cache=True) # 1 year - album track listings are immutable
604 async def get_album_tracks(self, prov_album_id: str) -> list[Track]:
605 """Get all album tracks for given album id."""
606 return [
607 parse_track(item, self)
608 async for item in self._get_all_items(f"albums/{prov_album_id}/tracks")
609 if item["id"]
610 ]
611
612 @use_cache(3600 * 3, allow_expired_cache=True) # 3 hours
613 async def get_playlist_tracks(self, prov_playlist_id: str, page: int = 0) -> list[Track]:
614 """Get playlist tracks."""
615 is_liked_songs = prov_playlist_id == self._get_liked_songs_playlist_id()
616 uri = "me/tracks" if is_liked_songs else f"playlists/{prov_playlist_id}/items"
617
618 # Liked songs always require global session
619 # For other playlists, call get_playlist first to trigger the fallback logic
620 # and populate the cache for which token to use
621 if is_liked_songs:
622 use_global = True
623 else:
624 # This call is cached and will determine/cache if global token is needed
625 await self.get_playlist(prov_playlist_id)
626 use_global = await self._playlist_requires_global_token(prov_playlist_id)
627
628 page_size = 50
629 offset = page * page_size
630 known_global = use_global
631
632 while True:
633 try:
634 meta = await self._get_playlist_pagination_meta(uri, page, use_global)
635 cache_checksum = meta["etag"]
636 total = meta["total"]
637
638 # Spotify has started returning 5xx for offset >= total on some
639 # playlists (notably algorithmic ones like Daily Mix). The retry
640 # storm that follows surfaces as "No playable items found".
641 if total and offset >= total:
642 spotify_result = {"total": total, "items": []}
643 else:
644 spotify_result = await self._get_data_with_caching(
645 uri,
646 cache_checksum,
647 limit=page_size,
648 offset=offset,
649 use_global_session=use_global,
650 )
651 break
652 except MediaNotFoundError:
653 if use_global or not self.dev_session_active:
654 raise
655 # Development Mode exposes metadata but restricts items for non-owned playlists.
656 use_global = True
657
658 if use_global and not known_global:
659 await self._set_playlist_requires_global_token(prov_playlist_id)
660
661 result: list[Track] = []
662 total = spotify_result.get("total", 0)
663 items = spotify_result.get("items", [])
664 # playlists/{id}/items is transitioning from item["track"] to item["item"]
665 # during Spotify's Feb 2026 rollout, so accept either shape.
666 for index, item in enumerate(items, 1):
667 # Spotify wraps/recycles items for offsets beyond the playlist size,
668 # so we need to break when we've reached the total.
669 if (offset + index) > total:
670 break
671 track_data = item and (item.get("item") or item.get("track"))
672 if not (track_data and track_data.get("id")):
673 continue
674 track = parse_track(track_data, self)
675 track.position = offset + index
676 result.append(track)
677 return result
678
679 @use_cache(86400 * 14, allow_expired_cache=True) # 14 days
680 async def get_artist_albums(self, prov_artist_id: str) -> list[Album]:
681 """Get a list of all albums for the given artist."""
682 try:
683 return [
684 parse_album(item, self)
685 async for item in self._get_all_items(
686 f"artists/{prov_artist_id}/albums?include_groups=album,single,compilation",
687 limit=10,
688 )
689 if (item and item["id"])
690 ]
691 except MediaNotFoundError:
692 self.logger.warning("Unable to fetch albums for artist %s", prov_artist_id)
693 return []
694
695 @use_cache(86400 * 14, allow_expired_cache=True) # 14 days
696 async def get_artist_toptracks(self, prov_artist_id: str) -> list[Track]:
697 """Get a list of 10 most popular tracks for the given artist."""
698 try:
699 artist = await self.get_artist(prov_artist_id)
700 endpoint = f"artists/{prov_artist_id}/top-tracks"
701 items = await self._get_data(endpoint)
702 return [
703 parse_track(item, self, artist=artist)
704 for item in items["tracks"]
705 if (item and item["id"])
706 ]
707 except MediaNotFoundError:
708 self.logger.warning(
709 "Top tracks search for artist %s appears to have been removed by Spotify for this account.",
710 prov_artist_id,
711 )
712 return []
713
714 async def library_add(self, item: MediaItemType) -> bool:
715 """Add item to library."""
716 uri_type_map = {
717 MediaType.ARTIST: "artist",
718 MediaType.ALBUM: "album",
719 MediaType.TRACK: "track",
720 MediaType.PLAYLIST: "playlist",
721 MediaType.PODCAST: "show",
722 MediaType.AUDIOBOOK: "audiobook",
723 }
724 if item.media_type == MediaType.AUDIOBOOK and not self.audiobooks_supported:
725 return False
726 uri_type = uri_type_map.get(item.media_type)
727 if not uri_type:
728 return False
729 uri = f"spotify:{uri_type}:{item.item_id}"
730 await self._put_data("me/library", uris=uri)
731 return True
732
733 async def library_remove(self, prov_item_id: str, media_type: MediaType) -> bool:
734 """Remove item from library."""
735 uri_type_map = {
736 MediaType.ARTIST: "artist",
737 MediaType.ALBUM: "album",
738 MediaType.TRACK: "track",
739 MediaType.PLAYLIST: "playlist",
740 MediaType.PODCAST: "show",
741 MediaType.AUDIOBOOK: "audiobook",
742 }
743 if media_type == MediaType.AUDIOBOOK and not self.audiobooks_supported:
744 return False
745 uri_type = uri_type_map.get(media_type)
746 if not uri_type:
747 return False
748 uri = f"spotify:{uri_type}:{prov_item_id}"
749 await self._delete_data("me/library", uris=uri)
750 return True
751
752 async def add_playlist_tracks(self, prov_playlist_id: str, prov_track_ids: list[str]) -> None:
753 """Add track(s) to playlist."""
754 track_uris = [f"spotify:track:{track_id}" for track_id in prov_track_ids]
755 data = {"uris": track_uris}
756 await self._post_data(f"playlists/{prov_playlist_id}/items", data=data)
757
758 async def remove_playlist_tracks(
759 self, prov_playlist_id: str, positions_to_remove: tuple[int, ...]
760 ) -> None:
761 """Remove track(s) from playlist."""
762 track_uris = []
763 for pos in positions_to_remove:
764 uri = f"playlists/{prov_playlist_id}/items"
765 spotify_result = await self._get_data(uri, limit=1, offset=pos - 1)
766 for item in spotify_result["items"]:
767 track_data = item and (item.get("item") or item.get("track"))
768 if not (track_data and track_data.get("id")):
769 continue
770 track_uris.append({"uri": f"spotify:track:{track_data['id']}"})
771 data = {"items": track_uris}
772 await self._delete_data(f"playlists/{prov_playlist_id}/items", data=data)
773
774 async def create_playlist(self, name: str, media_types: set[MediaType]) -> Playlist:
775 """Create a new playlist on provider with given name."""
776 data = {"name": name, "public": False}
777 new_playlist = await self._post_data("me/playlists", data=data)
778 self._fix_create_playlist_api_bug(new_playlist)
779 return parse_playlist(new_playlist, self)
780
781 @use_cache(86400 * 14, allow_expired_cache=True) # 14 days
782 async def get_similar_tracks(self, prov_track_id: str, limit: int = 25) -> list[Track]:
783 """Retrieve a dynamic list of tracks based on the provided item."""
784 # Recommendations endpoint is only available on global session (not developer API)
785 # https://developer.spotify.com/blog/2024-11-27-changes-to-the-web-api
786 endpoint = "recommendations"
787 items = await self._get_data(
788 endpoint, seed_tracks=prov_track_id, limit=limit, use_global_session=True
789 )
790 return [parse_track(item, self) for item in items["tracks"] if (item and item["id"])]
791
792 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
793 """Return content details for the given track/episode/audiobook when it will be streamed."""
794 if media_type == MediaType.AUDIOBOOK and self.audiobooks_supported:
795 chapters_data = await self._get_audiobook_chapters_data(item_id)
796 if not chapters_data:
797 raise MediaNotFoundError(f"No chapters found for audiobook {item_id}")
798
799 # Calculate total duration and convert to seconds for StreamDetails
800 total_duration_ms = sum(chapter.get("duration_ms", 0) for chapter in chapters_data)
801 duration_seconds = total_duration_ms // 1000
802
803 # Create chapter URIs for streaming
804 chapter_uris = []
805 for chapter in chapters_data:
806 chapter_id = chapter["id"]
807 chapter_uri = f"spotify://episode:{chapter_id}"
808 chapter_uris.append(chapter_uri)
809
810 return StreamDetails(
811 item_id=item_id,
812 provider=self.instance_id,
813 media_type=MediaType.AUDIOBOOK,
814 audio_format=AudioFormat(content_type=ContentType.OGG, bit_rate=320),
815 stream_type=StreamType.CUSTOM,
816 allow_seek=True,
817 can_seek=True,
818 duration=duration_seconds,
819 data={"chapters": chapter_uris, "chapters_data": chapters_data},
820 )
821
822 # For all other media types (tracks, podcast episodes)
823 return StreamDetails(
824 item_id=item_id,
825 provider=self.instance_id,
826 media_type=media_type,
827 audio_format=AudioFormat(content_type=ContentType.OGG, bit_rate=320),
828 stream_type=StreamType.CUSTOM,
829 allow_seek=True,
830 can_seek=True,
831 )
832
833 async def get_audio_stream(
834 self, streamdetails: StreamDetails, seek_position: int = 0
835 ) -> AsyncGenerator[bytes]:
836 """Get audio stream from Spotify via librespot."""
837 if streamdetails.media_type == MediaType.AUDIOBOOK and isinstance(streamdetails.data, dict):
838 chapter_uris = streamdetails.data.get("chapters", [])
839 chapters_data = streamdetails.data.get("chapters_data", [])
840
841 # Calculate which chapter to start from based on seek_position
842 seek_position_ms = seek_position * 1000
843 current_seek_ms = seek_position_ms
844 start_chapter = 0
845
846 if seek_position > 0 and chapters_data:
847 accumulated_duration_ms = 0
848
849 for i, chapter_data in enumerate(chapters_data):
850 chapter_duration_ms = chapter_data.get("duration_ms", 0)
851
852 if accumulated_duration_ms + chapter_duration_ms > seek_position_ms:
853 start_chapter = i
854 current_seek_ms = seek_position_ms - accumulated_duration_ms
855 break
856 accumulated_duration_ms += chapter_duration_ms
857 else:
858 start_chapter = len(chapter_uris) - 1
859 current_seek_ms = 0
860
861 # Convert back to seconds for librespot
862 current_seek_seconds = int(current_seek_ms // 1000)
863
864 # Stream chapters starting from the calculated position
865 consecutive_failures = 0
866 for i in range(start_chapter, len(chapter_uris)):
867 chapter_uri = chapter_uris[i]
868 chapter_seek = current_seek_seconds if i == start_chapter else 0
869
870 try:
871 chunk_count = 0
872 async for chunk in self.streamer.stream_spotify_uri(chapter_uri, chapter_seek):
873 yield chunk
874 chunk_count += 1
875 if chunk_count > 0:
876 consecutive_failures = 0
877 except Exception as e:
878 self.logger.warning("Chapter %s streaming failed", i + 1)
879 consecutive_failures += 1
880 if consecutive_failures >= 3:
881 raise AudioError("Audiobook streaming failed") from e
882 continue
883 else:
884 # Handle normal tracks and podcast episodes
885 async for chunk in self.streamer.get_audio_stream(streamdetails, seek_position):
886 yield chunk
887
888 @lock
889 async def login(self, force_refresh: bool = False) -> dict[str, Any]:
890 """
891 Log-in Spotify global session and return Auth/token info.
892
893 This uses MA's global client ID which has full API access but heavy rate limits.
894 """
895 # return the cached access token while it is still valid (refreshed before expiry)
896 if (
897 not force_refresh
898 and self._auth_info_global
899 and (self._auth_info_global["expires_at"] > (time.time() + 600))
900 ):
901 return self._auth_info_global
902 # read the refresh token from the persisted store rather than the in-memory config copy,
903 # which can lag a rotation and would make us refresh with a stale (revoked) token
904 if not (refresh_token := self._stored_refresh_token(CONF_REFRESH_TOKEN_GLOBAL)):
905 raise LoginFailed("Authentication required")
906
907 try:
908 auth_info = await get_spotify_token(
909 self.mass.http_session,
910 app_var("spotify_client_id"), # Always use MA's global client ID
911 refresh_token,
912 "global",
913 )
914 self.logger.debug("Successfully refreshed global access token")
915 except LoginFailed as err:
916 if "revoked" in str(err) or "invalid_grant" in str(err):
917 # Spotify rotates the refresh token on refresh and revokes the previous one.
918 # If the stored token was rotated while this refresh was in flight, the token
919 # we tried is merely stale, so keep the newer one instead of forcing re-auth.
920 if not self._refresh_token_superseded(CONF_REFRESH_TOKEN_GLOBAL, refresh_token):
921 self._update_setup_data(CONF_REFRESH_TOKEN_GLOBAL, None)
922 if self.available:
923 self.unload_with_error(err)
924 elif self.available:
925 self.mass.create_task(self.mass.unload_provider_with_error(self.instance_id, err))
926 raise
927
928 # make sure that our updated creds get stored in memory + config
929 self._auth_info_global = auth_info
930 # Spotify revokes the previous refresh token only when it rotates one, so on rotation
931 # persist immediately to ensure the new token survives a crash within the debounced-save
932 # window and avoids a forced re-auth; an unchanged token uses the normal debounced save.
933 token_rotated = auth_info["refresh_token"] != refresh_token
934 self._update_setup_data(
935 CONF_REFRESH_TOKEN_GLOBAL,
936 auth_info["refresh_token"],
937 immediate=token_rotated,
938 )
939
940 # get logged-in user info
941 if not self._sp_user:
942 self._sp_user = userinfo = await self._get_data(
943 "me", auth_info=auth_info, use_global_session=True
944 )
945 if country := userinfo.get("country"):
946 self.mass.metadata.set_default_preferred_language(country)
947 self.logger.info("Successfully logged in to Spotify as %s", userinfo["display_name"])
948 return auth_info
949
950 @lock
951 async def login_dev(self, force_refresh: bool = False) -> dict[str, Any]:
952 """
953 Log-in Spotify developer session and return Auth/token info.
954
955 This uses the user's custom client ID which has less rate limits but limited API access.
956 """
957 # return the cached access token while it is still valid (refreshed before expiry)
958 if (
959 not force_refresh
960 and self._auth_info_dev
961 and (self._auth_info_dev["expires_at"] > (time.time() + 600))
962 ):
963 return self._auth_info_dev
964 # read the refresh token from the persisted store rather than the in-memory config copy,
965 # which can lag a rotation and would make us refresh with a stale (revoked) token
966 refresh_token = self._stored_refresh_token(CONF_REFRESH_TOKEN_DEV)
967 client_id = self.get_setup_value(CONF_CLIENT_ID)
968 if not refresh_token or not client_id:
969 raise LoginFailed("Developer authentication not configured")
970
971 try:
972 auth_info = await get_spotify_token(
973 self.mass.http_session,
974 cast("str", client_id),
975 refresh_token,
976 "developer",
977 )
978 self.logger.debug("Successfully refreshed developer access token")
979 except LoginFailed as err:
980 if "revoked" in str(err) or "invalid_grant" in str(err):
981 # Spotify rotates the refresh token on refresh and revokes the previous one.
982 # If the stored token was rotated while this refresh was in flight, the token
983 # we tried is merely stale, so keep the newer one instead of forcing re-auth.
984 if not self._refresh_token_superseded(CONF_REFRESH_TOKEN_DEV, refresh_token):
985 self._update_setup_data(CONF_REFRESH_TOKEN_DEV, None)
986 self._update_setup_data(CONF_CLIENT_ID, None)
987 # Don't unload - we can still use the global session
988 self.dev_session_active = False
989 self.logger.warning(str(err))
990 raise
991
992 # make sure that our updated creds get stored in memory + config
993 self._auth_info_dev = auth_info
994 # Spotify revokes the previous refresh token only when it rotates one, so on rotation
995 # persist immediately to ensure the new token survives a crash within the debounced-save
996 # window and avoids a forced re-auth; an unchanged token uses the normal debounced save.
997 token_rotated = auth_info["refresh_token"] != refresh_token
998 self._update_setup_data(
999 CONF_REFRESH_TOKEN_DEV,
1000 auth_info["refresh_token"],
1001 immediate=token_rotated,
1002 )
1003
1004 self.logger.info("Successfully logged in to Spotify developer session")
1005 return auth_info
1006
1007 def _build_search_types(self, media_types: list[MediaType]) -> str:
1008 """Build comma-separated search types string from media types."""
1009 searchtypes = []
1010 if MediaType.ARTIST in media_types:
1011 searchtypes.append("artist")
1012 if MediaType.ALBUM in media_types:
1013 searchtypes.append("album")
1014 if MediaType.TRACK in media_types:
1015 searchtypes.append("track")
1016 if MediaType.PLAYLIST in media_types:
1017 searchtypes.append("playlist")
1018 if MediaType.PODCAST in media_types:
1019 searchtypes.append("show")
1020 if MediaType.AUDIOBOOK in media_types and self.audiobooks_supported:
1021 searchtypes.append("audiobook")
1022 return ",".join(searchtypes)
1023
1024 def _process_search_results(
1025 self, api_result: dict[str, Any], searchresult: SearchResults
1026 ) -> int:
1027 """
1028 Process API search results and update searchresult object.
1029
1030 Returns the total number of items received.
1031 """
1032 items_received = 0
1033
1034 if "artists" in api_result:
1035 artists = [
1036 parse_artist(item, self)
1037 for item in api_result["artists"]["items"]
1038 if (item and item["id"] and item["name"])
1039 ]
1040 searchresult.artists = [*searchresult.artists, *artists]
1041 items_received += len(api_result["artists"]["items"])
1042
1043 if "albums" in api_result:
1044 albums = [
1045 parse_album(item, self)
1046 for item in api_result["albums"]["items"]
1047 if (item and item["id"])
1048 ]
1049 searchresult.albums = [*searchresult.albums, *albums]
1050 items_received += len(api_result["albums"]["items"])
1051
1052 if "tracks" in api_result:
1053 tracks = [
1054 parse_track(item, self)
1055 for item in api_result["tracks"]["items"]
1056 if (item and item["id"])
1057 ]
1058 searchresult.tracks = [*searchresult.tracks, *tracks]
1059 items_received += len(api_result["tracks"]["items"])
1060
1061 if "playlists" in api_result:
1062 playlists = [
1063 parse_playlist(item, self)
1064 for item in api_result["playlists"]["items"]
1065 if (item and item["id"])
1066 ]
1067 searchresult.playlists = [*searchresult.playlists, *playlists]
1068 items_received += len(api_result["playlists"]["items"])
1069
1070 if "shows" in api_result:
1071 podcasts = []
1072 for item in api_result["shows"]["items"]:
1073 if not (item and item["id"]):
1074 continue
1075 # Filter out audiobooks - they have a distinctive description format
1076 description = item.get("description", "")
1077 if description.startswith("Author(s):") and "Narrator(s):" in description:
1078 continue
1079 podcasts.append(parse_podcast(item, self))
1080 searchresult.podcasts = [*searchresult.podcasts, *podcasts]
1081 items_received += len(api_result["shows"]["items"])
1082
1083 if "audiobooks" in api_result and self.audiobooks_supported:
1084 audiobooks = [
1085 parse_audiobook(item, self)
1086 for item in api_result["audiobooks"]["items"]
1087 if (item and item["id"])
1088 ]
1089 searchresult.audiobooks = [*searchresult.audiobooks, *audiobooks]
1090 items_received += len(api_result["audiobooks"]["items"])
1091
1092 return items_received
1093
1094 async def _setup_librespot_auth(self) -> None:
1095 """
1096 Install the stored playback credential into librespot's cache directory.
1097
1098 :raises LoginFailed: When no playback credential is configured, which requires the
1099 user to re-run the setup flow.
1100 """
1101 if self._librespot_bin is None:
1102 raise LoginFailed("Librespot binary not available")
1103 credentials = self.get_setup_value(CONF_LIBRESPOT_CREDENTIALS)
1104 if not credentials:
1105 # Spotify's login5 refuses credentials minted with any client id other than the one
1106 # librespot presents, so installs predating the dedicated playback credential (and
1107 # anything cached from before) cannot stream and must authorize playback again.
1108 raise LoginFailed(
1109 "Spotify playback authorization required",
1110 translation_key="playback_auth_required",
1111 translation_owner="provider.spotify",
1112 )
1113 await asyncio.to_thread(self._write_librespot_credentials, self.cache_dir, str(credentials))
1114
1115 @staticmethod
1116 def _write_librespot_credentials(cache_dir: str, credentials: str) -> None:
1117 """Write the stored credential to librespot's cache, replacing any stale one."""
1118 os.makedirs(cache_dir, exist_ok=True)
1119 credentials_file = os.path.join(cache_dir, CREDENTIALS_FILE)
1120 with open(credentials_file, "w", encoding="utf-8") as fileobj:
1121 fileobj.write(credentials)
1122
1123 async def _get_auth_info(self, use_global_session: bool = False) -> dict[str, Any]:
1124 """
1125 Get auth info for API requests, preferring dev session if available.
1126
1127 :param use_global_session: Force use of global session (for features not available on dev).
1128 """
1129 if use_global_session or not self.dev_session_active:
1130 return await self.login()
1131
1132 # Try dev session first
1133 try:
1134 return await self.login_dev()
1135 except LoginFailed:
1136 # Fall back to global session
1137 self.logger.debug("Falling back to global session after dev session failure")
1138 return await self.login()
1139
1140 def _get_liked_songs_playlist_id(self) -> str:
1141 return f"{LIKED_SONGS_FAKE_PLAYLIST_ID_PREFIX}-{self.instance_id}"
1142
1143 @use_cache(86400, allow_expired_cache=True) # 24h; serve stale + refresh in background
1144 async def _get_new_releases(self) -> list[Album]:
1145 """Get Spotify's curated 'new releases' albums."""
1146 try:
1147 result = await self._get_data("browse/new-releases", limit=50)
1148 except MediaNotFoundError:
1149 return []
1150 return [
1151 parse_album(item, self)
1152 for item in result.get("albums", {}).get("items", [])
1153 if item and item.get("id")
1154 ]
1155
1156 @use_cache(86400 * 7, allow_expired_cache=True) # 7d; serve stale + refresh in background
1157 async def _get_categories(self, locale: str) -> list[BrowseFolder]:
1158 """Get Spotify's curated browse categories (genres & moods) as browse folders."""
1159 try:
1160 result = await self._get_data("browse/categories", locale=locale, limit=50)
1161 except MediaNotFoundError:
1162 return []
1163 return [
1164 BrowseFolder(
1165 item_id=cat["id"],
1166 provider=self.instance_id,
1167 path=f"{self.instance_id}://categories/{cat['id']}",
1168 name=cat["name"],
1169 is_playable=False,
1170 )
1171 for cat in result.get("categories", {}).get("items", [])
1172 if cat and cat.get("id") and cat.get("name")
1173 ]
1174
1175 @use_cache(86400, allow_expired_cache=True) # 24h; serve stale + refresh in background
1176 async def _get_category_playlists(self, category_id: str, locale: str) -> list[Playlist]:
1177 """Get the playlists for a single Spotify browse category."""
1178 try:
1179 result = await self._get_data(
1180 f"browse/categories/{category_id}/playlists",
1181 locale=locale,
1182 limit=50,
1183 use_global_session=True,
1184 )
1185 except MediaNotFoundError:
1186 return []
1187 return [
1188 parse_playlist(item, self)
1189 for item in result.get("playlists", {}).get("items", [])
1190 if item and item.get("id") and item.get("name")
1191 ]
1192
1193 async def _get_liked_songs_playlist(self) -> Playlist:
1194 if self._sp_user is None:
1195 raise LoginFailed("User info not available - not logged in")
1196
1197 liked_songs = Playlist(
1198 item_id=self._get_liked_songs_playlist_id(),
1199 provider=self.instance_id,
1200 name=f"Liked Songs {self._sp_user['display_name']}",
1201 translation_key="liked_songs",
1202 translation_params=[self._sp_user["display_name"]],
1203 owner=self._sp_user["display_name"],
1204 provider_mappings={
1205 ProviderMapping(
1206 item_id=self._get_liked_songs_playlist_id(),
1207 provider_domain=self.domain,
1208 provider_instance=self.instance_id,
1209 url="https://open.spotify.com/collection/tracks",
1210 is_unique=True, # liked songs is user-specific
1211 )
1212 },
1213 )
1214
1215 liked_songs.is_editable = False # TODO Editing requires special endpoints
1216
1217 # Add image to the playlist metadata
1218 image = MediaItemImage(
1219 type=ImageType.THUMB,
1220 path="https://misc.scdn.co/liked-songs/liked-songs-64.png",
1221 provider=self.instance_id,
1222 remotely_accessible=True,
1223 )
1224 if liked_songs.metadata.images is None:
1225 liked_songs.metadata.images = UniqueList([image])
1226 else:
1227 liked_songs.metadata.add_image(image)
1228
1229 return liked_songs
1230
1231 async def _get_playlist_pagination_meta(
1232 self, endpoint: str, page: int, use_global_session: bool
1233 ) -> dict[str, Any]:
1234 """
1235 Return pagination metadata for a Spotify playlist traversal.
1236
1237 :param endpoint: Spotify API endpoint for the playlist items.
1238 :param page: Requested playlist page.
1239 :param use_global_session: Whether the global Spotify session is required.
1240 """
1241 state_key = (endpoint, use_global_session)
1242 if state := self._playlist_pagination_states.get(state_key):
1243 self._playlist_pagination_states.move_to_end(state_key)
1244 else:
1245 state = _PlaylistPaginationState(lock=asyncio.Lock())
1246 self._playlist_pagination_states[state_key] = state
1247 while len(self._playlist_pagination_states) > _PLAYLIST_PAGINATION_STATE_LIMIT:
1248 self._playlist_pagination_states.popitem(last=False)
1249
1250 observed_snapshot = state.snapshot
1251 async with state.lock:
1252 snapshot = state.snapshot
1253 # A concurrent page may have populated this snapshot while this call waited.
1254 if snapshot and (page > 0 or snapshot is not observed_snapshot):
1255 return snapshot
1256
1257 if page == 0:
1258 state.snapshot = None
1259 meta = await self._get_paginated_meta(
1260 endpoint,
1261 limit=1,
1262 offset=0,
1263 use_global_session=use_global_session,
1264 )
1265 state.snapshot = meta
1266 return meta
1267
1268 async def _playlist_requires_global_token(self, prov_playlist_id: str) -> bool:
1269 """
1270 Check if a playlist requires global token (cached).
1271
1272 :param prov_playlist_id: The Spotify playlist ID.
1273 :returns: True if the playlist requires global token.
1274 """
1275 cache_key = f"playlist_global_token_{prov_playlist_id}"
1276 return bool(await self.mass.cache.get(cache_key, provider=self.instance_id))
1277
1278 async def _set_playlist_requires_global_token(self, prov_playlist_id: str) -> None:
1279 """
1280 Mark a playlist as requiring global token in cache.
1281
1282 :param prov_playlist_id: The Spotify playlist ID.
1283 """
1284 cache_key = f"playlist_global_token_{prov_playlist_id}"
1285 # Cache for 90 days - playlist ownership doesn't change
1286 await self.mass.cache.set(cache_key, True, provider=self.instance_id, expiration=86400 * 90)
1287
1288 async def _add_audiobook_chapters(self, audiobook: Audiobook) -> None:
1289 """Add chapter metadata to an audiobook from Spotify API data."""
1290 try:
1291 chapters_data = await self._get_audiobook_chapters_data(audiobook.item_id)
1292 if chapters_data:
1293 chapters = []
1294 total_duration_seconds = 0.0
1295
1296 for idx, chapter in enumerate(chapters_data):
1297 duration_ms = chapter.get("duration_ms", 0)
1298 duration_seconds = duration_ms / 1000.0
1299
1300 chapter_obj = MediaItemChapter(
1301 position=idx + 1,
1302 name=chapter.get("name", f"Chapter {idx + 1}"),
1303 start=total_duration_seconds,
1304 end=total_duration_seconds + duration_seconds,
1305 )
1306 chapters.append(chapter_obj)
1307 total_duration_seconds += duration_seconds
1308
1309 audiobook.metadata.chapters = chapters
1310 audiobook.duration = int(total_duration_seconds)
1311
1312 except (MediaNotFoundError, ResourceTemporarilyUnavailable, ProviderUnavailableError) as e:
1313 self.logger.warning(f"Failed to get chapters for audiobook {audiobook.item_id}: {e}")
1314
1315 @use_cache(43200) # 12 hours - balances freshness with performance
1316 async def _get_podcast_episodes_data(self, prov_podcast_id: str) -> list[dict[str, Any]]:
1317 """
1318 Get raw episode data from Spotify API (cached).
1319
1320 :param prov_podcast_id: Spotify podcast ID.
1321 """
1322 episodes_data: list[dict[str, Any]] = []
1323
1324 try:
1325 async for item in self._get_all_items(
1326 f"shows/{prov_podcast_id}/episodes", market="from_token"
1327 ):
1328 if item and item.get("id"):
1329 episodes_data.append(item)
1330 except MediaNotFoundError:
1331 self.logger.warning("Podcast %s not found", prov_podcast_id)
1332 return []
1333 except ResourceTemporarilyUnavailable as err:
1334 self.logger.warning(
1335 "Temporary error fetching episodes for %s: %s", prov_podcast_id, err
1336 )
1337 raise
1338
1339 return episodes_data
1340
1341 @use_cache(7200) # 2 hours - shorter cache for resume point data
1342 async def _get_audiobook_chapters_data(self, prov_audiobook_id: str) -> list[dict[str, Any]]:
1343 """
1344 Get raw chapter data from Spotify API (cached).
1345
1346 :param prov_audiobook_id: Spotify audiobook ID.
1347 """
1348 chapters_data: list[dict[str, Any]] = []
1349
1350 try:
1351 async for item in self._get_all_items(
1352 f"audiobooks/{prov_audiobook_id}/chapters", market="from_token"
1353 ):
1354 if item and item.get("id"):
1355 chapters_data.append(item)
1356 except MediaNotFoundError:
1357 self.logger.warning("Audiobook %s not found", prov_audiobook_id)
1358 return []
1359 except ResourceTemporarilyUnavailable as err:
1360 self.logger.warning(
1361 "Temporary error fetching chapters for %s: %s", prov_audiobook_id, err
1362 )
1363 raise
1364
1365 return chapters_data
1366
1367 async def _get_all_items(
1368 self, endpoint: str, key: str = "items", limit: int = 50, **kwargs: Any
1369 ) -> AsyncGenerator[dict[str, Any]]:
1370 """Get all items from a paged list."""
1371 offset = 0
1372 # single request to fetch the etag (used as cache checksum) and total
1373 meta = await self._get_cached_paginated_meta(endpoint, limit=1, offset=0, **kwargs)
1374 cache_checksum = meta["etag"]
1375 total = meta["total"]
1376 while True:
1377 # Avoid requesting beyond the known end. Spotify can return 5xx
1378 # for offset >= total on some endpoints (e.g. algorithmic playlists).
1379 if total and offset >= total:
1380 break
1381 result = await self._get_data_with_caching(
1382 endpoint, cache_checksum=cache_checksum, limit=limit, offset=offset, **kwargs
1383 )
1384 offset += limit
1385 if not result or key not in result or not result[key]:
1386 break
1387 for item in result[key]:
1388 yield item
1389 if len(result[key]) < limit:
1390 break
1391
1392 async def _get_data_with_caching(
1393 self, endpoint: str, cache_checksum: str | None, **kwargs: Any
1394 ) -> dict[str, Any]:
1395 """Get data from api with caching."""
1396 cache_key_parts = [endpoint]
1397 for key in sorted(kwargs.keys()):
1398 cache_key_parts.append(f"{key}{kwargs[key]}")
1399 cache_key = ".".join(map(str, cache_key_parts))
1400 if cached := await self.mass.cache.get(
1401 cache_key, provider=self.instance_id, checksum=cache_checksum, allow_bypass=False
1402 ):
1403 return cast("dict[str, Any]", cached)
1404 result = await self._get_data(endpoint, **kwargs)
1405 await self.mass.cache.set(
1406 cache_key, result, provider=self.instance_id, checksum=cache_checksum
1407 )
1408 return result
1409
1410 @use_cache(120, allow_bypass=False) # short cache: repeated traversals reuse metadata
1411 async def _get_cached_paginated_meta(self, endpoint: str, **kwargs: Any) -> dict[str, Any]:
1412 """Get cached pagination metadata for a paginated API endpoint."""
1413 return await self._get_paginated_meta(endpoint, **kwargs)
1414
1415 async def _get_paginated_meta(self, endpoint: str, **kwargs: Any) -> dict[str, Any]:
1416 """Get etag and total item count for a paginated api endpoint."""
1417 _res = await self._get_data(endpoint, **kwargs)
1418 return {"etag": _res.get("etag"), "total": _res.get("total", 0)}
1419
1420 @throttle_with_retries
1421 async def _get_data(self, endpoint: str, **kwargs: Any) -> dict[str, Any]:
1422 """
1423 Get data from api.
1424
1425 :param endpoint: API endpoint to call.
1426 :param use_global_session: Force use of global session (for features not available on dev).
1427 """
1428 url = f"https://api.spotify.com/v1/{endpoint}"
1429 kwargs["market"] = "from_token"
1430 kwargs["country"] = "from_token"
1431 use_global_session = kwargs.pop("use_global_session", False)
1432 if not (auth_info := kwargs.pop("auth_info", None)):
1433 auth_info = await self._get_auth_info(use_global_session=use_global_session)
1434 headers = {"Authorization": f"Bearer {auth_info['access_token']}"}
1435 locale = self.mass.metadata.locale.replace("_", "-")
1436 language = locale.split("-")[0]
1437 headers["Accept-Language"] = f"{locale}, {language};q=0.9, *;q=0.5"
1438 self.logger.debug("handling get data %s with kwargs %s", url, kwargs)
1439 async with (
1440 self.mass.http_session.get(
1441 url,
1442 headers=headers,
1443 params=kwargs,
1444 timeout=aiohttp.ClientTimeout(total=120),
1445 ) as response,
1446 ):
1447 # handle spotify rate limiter
1448 if response.status == 429:
1449 backoff_time = int(response.headers["Retry-After"])
1450 raise RateLimited("Spotify Rate Limiter", backoff_time=backoff_time)
1451 # handle temporary server error
1452 if response.status in (502, 503):
1453 raise ResourceTemporarilyUnavailable(backoff_time=30)
1454
1455 # handle token expired, raise ResourceTemporarilyUnavailable
1456 # so it will be retried (and the token refreshed)
1457 if response.status == 401:
1458 if use_global_session or not self.dev_session_active:
1459 self._auth_info_global = None
1460 else:
1461 self._auth_info_dev = None
1462 raise ResourceTemporarilyUnavailable("Token expired", backoff_time=1)
1463
1464 if response.status in (400, 403, 404):
1465 try:
1466 error = await response.json(loads=json_loads)
1467 message = error.get("error", {}).get("message") or response.reason
1468 except aiohttp.ContentTypeError, JSONDecodeError:
1469 message = (await response.text()) or response.reason
1470
1471 self.logger.debug(
1472 "Spotify API error: endpoint=%s, status=%s, reason=%s, message=%s",
1473 endpoint,
1474 response.status,
1475 response.reason,
1476 message,
1477 )
1478
1479 raise MediaNotFoundError(f"{endpoint} not found")
1480
1481 response.raise_for_status()
1482 result: dict[str, Any] = await response.json(loads=json_loads)
1483 if etag := response.headers.get("ETag"):
1484 result["etag"] = etag
1485 return result
1486
1487 @throttle_with_retries
1488 async def _delete_data(self, endpoint: str, data: Any = None, **kwargs: Any) -> None:
1489 """Delete data from api."""
1490 url = f"https://api.spotify.com/v1/{endpoint}"
1491 use_global_session = kwargs.pop("use_global_session", False)
1492 if not (auth_info := kwargs.pop("auth_info", None)):
1493 auth_info = await self._get_auth_info(use_global_session=use_global_session)
1494 headers = {"Authorization": f"Bearer {auth_info['access_token']}"}
1495 async with self.mass.http_session.delete(
1496 url, headers=headers, params=kwargs, json=data, ssl=True
1497 ) as response:
1498 # handle spotify rate limiter
1499 if response.status == 429:
1500 backoff_time = int(response.headers["Retry-After"])
1501 raise RateLimited("Spotify Rate Limiter", backoff_time=backoff_time)
1502 # handle token expired, raise ResourceTemporarilyUnavailable
1503 # so it will be retried (and the token refreshed)
1504 if response.status == 401:
1505 if use_global_session or not self.dev_session_active:
1506 self._auth_info_global = None
1507 else:
1508 self._auth_info_dev = None
1509 raise ResourceTemporarilyUnavailable("Token expired", backoff_time=1)
1510 # handle temporary server error
1511 if response.status in (502, 503):
1512 raise ResourceTemporarilyUnavailable(backoff_time=30)
1513 response.raise_for_status()
1514
1515 @throttle_with_retries
1516 async def _put_data(self, endpoint: str, data: Any = None, **kwargs: Any) -> None:
1517 """Put data on api."""
1518 url = f"https://api.spotify.com/v1/{endpoint}"
1519 use_global_session = kwargs.pop("use_global_session", False)
1520 if not (auth_info := kwargs.pop("auth_info", None)):
1521 auth_info = await self._get_auth_info(use_global_session=use_global_session)
1522 headers = {"Authorization": f"Bearer {auth_info['access_token']}"}
1523 async with self.mass.http_session.put(
1524 url, headers=headers, params=kwargs, json=data, ssl=True
1525 ) as response:
1526 # handle spotify rate limiter
1527 if response.status == 429:
1528 backoff_time = int(response.headers["Retry-After"])
1529 raise RateLimited("Spotify Rate Limiter", backoff_time=backoff_time)
1530 # handle token expired, raise ResourceTemporarilyUnavailable
1531 # so it will be retried (and the token refreshed)
1532 if response.status == 401:
1533 if use_global_session or not self.dev_session_active:
1534 self._auth_info_global = None
1535 else:
1536 self._auth_info_dev = None
1537 raise ResourceTemporarilyUnavailable("Token expired", backoff_time=1)
1538
1539 # handle temporary server error
1540 if response.status in (502, 503):
1541 raise ResourceTemporarilyUnavailable(backoff_time=30)
1542 response.raise_for_status()
1543
1544 @throttle_with_retries
1545 async def _post_data(
1546 self, endpoint: str, data: Any = None, want_result: bool = True, **kwargs: Any
1547 ) -> dict[str, Any]:
1548 """Post data on api."""
1549 url = f"https://api.spotify.com/v1/{endpoint}"
1550 use_global_session = kwargs.pop("use_global_session", False)
1551 if not (auth_info := kwargs.pop("auth_info", None)):
1552 auth_info = await self._get_auth_info(use_global_session=use_global_session)
1553 headers = {"Authorization": f"Bearer {auth_info['access_token']}"}
1554 async with self.mass.http_session.post(
1555 url, headers=headers, params=kwargs, json=data, ssl=True
1556 ) as response:
1557 # handle spotify rate limiter
1558 if response.status == 429:
1559 backoff_time = int(response.headers["Retry-After"])
1560 raise RateLimited("Spotify Rate Limiter", backoff_time=backoff_time)
1561 # handle token expired, raise ResourceTemporarilyUnavailable
1562 # so it will be retried (and the token refreshed)
1563 if response.status == 401:
1564 if use_global_session or not self.dev_session_active:
1565 self._auth_info_global = None
1566 else:
1567 self._auth_info_dev = None
1568 raise ResourceTemporarilyUnavailable("Token expired", backoff_time=1)
1569 # handle temporary server error
1570 if response.status in (502, 503):
1571 raise ResourceTemporarilyUnavailable(backoff_time=30)
1572 response.raise_for_status()
1573 if not want_result:
1574 return {}
1575 result: dict[str, Any] = await response.json(loads=json_loads)
1576 return result
1577
1578 def _fix_create_playlist_api_bug(self, playlist_obj: dict[str, Any]) -> None:
1579 """Fix spotify API bug where incorrect owner id is returned from Create Playlist."""
1580 if self._sp_user is None:
1581 raise LoginFailed("User info not available - not logged in")
1582
1583 if playlist_obj["owner"]["id"] != self._sp_user["id"]:
1584 playlist_obj["owner"]["id"] = self._sp_user["id"]
1585 playlist_obj["owner"]["display_name"] = self._sp_user["display_name"]
1586 else:
1587 self.logger.warning(
1588 "FIXME: Spotify have fixed their Create Playlist API, this fix can be removed."
1589 )
1590
1591 async def _test_audiobook_support(self) -> bool:
1592 """Test if audiobooks are supported in user's region."""
1593 try:
1594 await self._get_data("me/audiobooks", limit=1)
1595 return True
1596 except aiohttp.ClientResponseError as e:
1597 if e.status == 403:
1598 return False # Not available
1599 raise # Re-raise other HTTP errors
1600 except MediaNotFoundError, ProviderUnavailableError:
1601 return False
1602
1603 def _stored_refresh_token(self, key: str) -> str | None:
1604 """
1605 Return the currently persisted refresh token, or None if not set.
1606
1607 Reads through the live setup_data (kept in sync with a just-rotated token) so a
1608 refresh never uses a stale, revoked token from a lagging in-memory config copy.
1609
1610 :param key: Setup data key of the refresh token to read.
1611 """
1612 token = self.get_setup_value(key)
1613 return cast("str", token) if token else None
1614
1615 def _refresh_token_superseded(self, key: str, used_token: str) -> bool:
1616 """
1617 Return whether the stored refresh token differs from the one just used.
1618
1619 :param key: Config key of the refresh token to check.
1620 :param used_token: The refresh token value that was just used to refresh.
1621 """
1622 stored_token = self._stored_refresh_token(key)
1623 if not stored_token:
1624 return False
1625 return stored_token != used_token
1626