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