/
/
1"""Yandex Music provider implementation."""
2
3from __future__ import annotations
4
5import asyncio
6import hashlib
7import json
8import logging
9import random
10import uuid
11from collections.abc import AsyncGenerator, Sequence
12from io import BytesIO
13from typing import TYPE_CHECKING, Any
14
15from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
16from music_assistant_models.enums import (
17 ConfigEntryType,
18 ImageType,
19 MediaType,
20 ProviderFeature,
21 StreamType,
22)
23from music_assistant_models.errors import (
24 InvalidDataError,
25 LoginFailed,
26 MediaNotFoundError,
27 ProviderUnavailableError,
28 ResourceTemporarilyUnavailable,
29)
30from music_assistant_models.media_items import (
31 Album,
32 Artist,
33 Audiobook,
34 BrowseFolder,
35 ItemMapping,
36 MediaItemChapter,
37 MediaItemImage,
38 MediaItemType,
39 Playlist,
40 Podcast,
41 PodcastEpisode,
42 ProviderMapping,
43 RecommendationFolder,
44 SearchResults,
45 Track,
46 UniqueList,
47)
48from music_assistant_models.streamdetails import StreamDetails
49from PIL import Image as PilImage
50from ya_passport_auth import SecretStr
51
52from music_assistant.constants import CONF_ENTRY_UNOFFICIAL_PROVIDER
53from music_assistant.controllers.cache import use_cache
54from music_assistant.helpers.datetime import utc
55from music_assistant.models.music_provider import MusicProvider
56
57from .api_client import YandexMusicClient
58from .auth import refresh_credentials_via_passport, refresh_music_token
59from .constants import (
60 BROWSE_INITIAL_TRACKS,
61 COLLECTION_FOLDER_ID,
62 CONF_ACTION_DELETE_WAVE_PRESET,
63 CONF_ACTION_SAVE_WAVE_PRESET,
64 CONF_BASE_URL,
65 CONF_LIKED_TRACKS_MAX_TRACKS,
66 CONF_MANUAL_TOKEN,
67 CONF_MY_WAVE_MAX_TRACKS,
68 CONF_QUALITY,
69 CONF_REFRESH_TOKEN,
70 CONF_RESTRICTIVE_RATE_LIMITS,
71 CONF_TOKEN,
72 CONF_WAVE_PRESET_DRAFT_DIVERSITY,
73 CONF_WAVE_PRESET_DRAFT_LANGUAGE,
74 CONF_WAVE_PRESET_DRAFT_MOOD,
75 CONF_WAVE_PRESET_DRAFT_NAME,
76 CONF_WAVE_PRESET_TO_DELETE,
77 CONF_WAVE_PRESETS_DATA,
78 CONF_X_TOKEN,
79 DEFAULT_BASE_URL,
80 DISCOVERY_INITIAL_TRACKS,
81 FOR_YOU_FOLDER_ID,
82 IMAGE_SIZE_MEDIUM,
83 LIKED_BATCH_JITTER_MIN_S,
84 LIKED_BATCH_JITTER_SPAN_S,
85 LIKED_TRACKS_PLAYLIST_ID,
86 LISTENING_HISTORY_FOLDER_ID,
87 MY_WAVE_BATCH_SIZE,
88 MY_WAVE_MODES_FOLDER_ID,
89 MY_WAVE_PLAYLIST_ID,
90 MY_WAVE_PRESETS_FOLDER_ID,
91 MY_WAVES_FOLDER_ID,
92 MY_WAVES_SET_FOLDER_ID,
93 PINNED_ITEMS_FOLDER_ID,
94 PLAYLIST_ID_SPLITTER,
95 QUALITY_BALANCED,
96 QUALITY_EFFICIENT,
97 QUALITY_HIGH,
98 QUALITY_SUPERB,
99 RADIO_FOLDER_ID,
100 RADIO_TRACK_ID_SEP,
101 ROTOR_STATION_MY_WAVE,
102 TAG_CATEGORY_ACTIVITY,
103 TAG_CATEGORY_ERA,
104 TAG_CATEGORY_GENRES,
105 TAG_CATEGORY_MOOD,
106 TAG_CATEGORY_ORDER,
107 TAG_MIXES,
108 TAG_SEASONAL_MAP,
109 TAG_SLUG_CATEGORY,
110 TRACK_BATCH_SIZE,
111 WAVE_CATEGORY_DISPLAY_ORDER,
112 WAVE_MODE_ORDER,
113 WAVE_MODE_PRESETS,
114 WAVE_MODE_SEP,
115 WAVE_PRESET_DIVERSITY_VALUES,
116 WAVE_PRESET_LANGUAGE_VALUES,
117 WAVE_PRESET_MOOD_VALUES,
118 WAVES_FOLDER_ID,
119 WAVES_LANDING_FOLDER_ID,
120)
121from .parsers import (
122 _get_image_url as get_image_url,
123)
124from .parsers import (
125 classify_album,
126 get_canonical_provider_name,
127 parse_album,
128 parse_artist,
129 parse_audiobook,
130 parse_playlist,
131 parse_podcast,
132 parse_podcast_episode,
133 parse_track,
134)
135from .presets import parse_stored_presets
136from .streaming import YandexMusicStreamingManager
137
138if TYPE_CHECKING:
139 from music_assistant_models.config_entries import ConfigActionResult
140 from yandex_music import Album as YandexAlbum
141 from yandex_music import Track as YandexTrack
142
143
144# MediaType sub-paths that MA's default MusicProvider.browse() understands.
145# Used by the Collection dispatcher to delegate nested paths back to core.
146_COLLECTION_SUB_FOLDERS: frozenset[str] = frozenset(
147 {"tracks", "artists", "albums", "playlists", "audiobooks", "podcasts"}
148)
149
150# Collection sub-folder rows: (ProviderFeature, browse sub_id, strings.json label key,
151# is_playable). The sub_id ("tracks") and label key ("my_favorites") differ on purpose so the
152# Collection labels stay distinct from the core "media.folder.*" library labels.
153_COLLECTION_SUBFOLDERS: tuple[tuple[ProviderFeature, str, str, bool], ...] = (
154 (ProviderFeature.LIBRARY_TRACKS, "tracks", "my_favorites", True),
155 (ProviderFeature.LIBRARY_ARTISTS, "artists", "my_artists", True),
156 (ProviderFeature.LIBRARY_ALBUMS, "albums", "my_albums", True),
157 (ProviderFeature.LIBRARY_PLAYLISTS, "playlists", "my_playlists", True),
158 (ProviderFeature.LIBRARY_PODCASTS, "podcasts", "my_podcasts", False),
159 (ProviderFeature.LIBRARY_AUDIOBOOKS, "audiobooks", "my_audiobooks", False),
160)
161
162
163def _media_label_key(slug: str) -> str:
164 """Normalize a tag/category slug into its strings.json authoring key (spaces â underscores)."""
165 return slug.replace(" ", "_")
166
167
168def _split_wave_mode(station_id: str) -> tuple[str, dict[str, str]]:
169 """
170 Split a wave-mode station key into its base station ID and preset settings.
171
172 Keys like ``user:onyourwave#discover`` encode a specific preset on top of
173 the base rotor station. The part before ``#`` is the station ID that goes
174 to Yandex; the part after is a key into WAVE_MODE_PRESETS.
175
176 :param station_id: Station key, with or without a ``#preset`` suffix.
177 :return: Tuple of (base_station_id, settings_dict). The suffix, if
178 present, is always stripped â only the base station goes to
179 Yandex. ``settings_dict`` is the preset's settings when the suffix
180 matches a known WAVE_MODE_PRESETS key, or an empty dict otherwise
181 (unknown suffix â base station fired with no extra seeds).
182 """
183 if WAVE_MODE_SEP not in station_id:
184 return (station_id, {})
185 base, preset = station_id.split(WAVE_MODE_SEP, 1)
186 return (base, dict(WAVE_MODE_PRESETS.get(preset, {})))
187
188
189def _parse_radio_item_id(item_id: str) -> tuple[str, str | None]:
190 """
191 Extract track_id and optional station_id from provider item_id.
192
193 My Wave tracks use item_id format 'track_id@station_id'. Other tracks use
194 plain track_id.
195
196 :param item_id: Provider item_id (may contain RADIO_TRACK_ID_SEP).
197 :return: (track_id, station_id or None).
198 """
199 if RADIO_TRACK_ID_SEP in item_id:
200 parts = item_id.split(RADIO_TRACK_ID_SEP, 1)
201 return (parts[0], parts[1] if len(parts) > 1 else None)
202 return (item_id, None)
203
204
205def _extract_chapter_map_from_album(album: YandexAlbum) -> tuple[list[str], list[int]]:
206 """
207 Flatten an audiobook album's volumes into (chapter_track_ids, chapter_durations_ms).
208
209 Shared by ``_get_audiobook_stream_details`` and ``_resolve_audiobook_chapter_map``
210 so the two code paths can't drift (e.g. when we later filter bad tracks).
211 """
212 chapter_ids: list[str] = []
213 chapter_durations_ms: list[int] = []
214 for disc in album.volumes or []:
215 for track_obj in disc:
216 chapter_ids.append(str(track_obj.id))
217 chapter_durations_ms.append(int(track_obj.duration_ms or 0))
218 return chapter_ids, chapter_durations_ms
219
220
221def _merge_wave_preset(
222 name: str | None,
223 diversity: str | None,
224 mood: str | None,
225 language: str | None,
226 presets_data: str | None,
227) -> str:
228 """
229 Merge the given draft fields into the stored preset list and return the new JSON.
230
231 Overwrites an existing preset with the same name instead of creating a duplicate.
232 Raises ``InvalidDataError`` when the name is blank.
233
234 :param name: Draft preset name; blank/whitespace-only raises.
235 :param diversity: Draft diversity seed ("" / None â omitted).
236 :param mood: Draft mood/energy seed ("" / None â omitted).
237 :param language: Draft language seed ("" / None â omitted).
238 :param presets_data: The current stored presets JSON.
239 """
240 clean_name = name.strip() if isinstance(name, str) else ""
241 if not clean_name:
242 raise InvalidDataError("Please fill the preset name before saving.")
243 presets = parse_stored_presets(presets_data)
244 presets = [p for p in presets if p["name"] != clean_name]
245 new_preset: dict[str, str] = {
246 "name": clean_name,
247 **{
248 api_key: val
249 for val, api_key in (
250 (diversity, "diversity"),
251 (mood, "moodEnergy"),
252 (language, "language"),
253 )
254 if isinstance(val, str) and val
255 },
256 }
257 presets.append(new_preset)
258 return json.dumps(presets, ensure_ascii=False)
259
260
261def _remove_wave_preset(target: str | None, presets_data: str | None) -> str:
262 """
263 Remove the named preset from the stored list and return the new JSON.
264
265 Raises ``InvalidDataError`` when no name is selected. Idempotent â an absent
266 name simply rewrites an unchanged list.
267
268 :param target: Name of the preset to remove; blank/whitespace-only raises.
269 :param presets_data: The current stored presets JSON.
270 """
271 clean_target = target.strip() if isinstance(target, str) else ""
272 if not clean_target:
273 raise InvalidDataError("Please select a preset to delete.")
274 presets = parse_stored_presets(presets_data)
275 presets = [p for p in presets if p["name"] != clean_target]
276 return json.dumps(presets, ensure_ascii=False)
277
278
279def _wave_preset_config_entries(presets_data: str | None) -> list[ConfigEntry]:
280 """
281 Return the wave-preset builder UI (all advanced settings).
282
283 Layout:
284 - Section label showing how many presets are saved.
285 - Four "draft" fields (name + three dropdowns) the user fills in.
286 - "Save preset" action â copies draft into the JSON store.
287 - "Delete preset" dropdown + action (hidden when no presets exist).
288 - Hidden STRING carrying the JSON store itself.
289
290 Number of presets is unbounded; the user never edits JSON directly.
291
292 :param presets_data: The stored wave-presets JSON (``CONF_WAVE_PRESETS_DATA``).
293 """
294 empty_title = "â Default â"
295 diversity_options = [
296 ConfigValueOption(v, title=empty_title if not v else v.title())
297 for v in WAVE_PRESET_DIVERSITY_VALUES
298 ]
299 mood_options = [
300 ConfigValueOption(v, title=empty_title if not v else v.title())
301 for v in WAVE_PRESET_MOOD_VALUES
302 ]
303 language_options = [
304 ConfigValueOption(v, title=empty_title if not v else v.replace("-", " ").title())
305 for v in WAVE_PRESET_LANGUAGE_VALUES
306 ]
307
308 presets = parse_stored_presets(presets_data)
309 has_presets = bool(presets)
310 delete_options = [ConfigValueOption(p["name"], title=p["name"]) for p in presets]
311 if not delete_options:
312 # Empty options can break some frontends; supply a no-op placeholder.
313 delete_options = [ConfigValueOption("")]
314
315 return [
316 ConfigEntry(
317 key="wave_preset_section_label",
318 type=ConfigEntryType.LABEL,
319 translation_key="wave_preset_section_saved" if has_presets else None,
320 translation_params=[str(len(presets))] if has_presets else None,
321 advanced=True,
322 ),
323 ConfigEntry(
324 key=CONF_WAVE_PRESET_DRAFT_NAME,
325 type=ConfigEntryType.STRING,
326 default_value=None,
327 required=False,
328 advanced=True,
329 ),
330 ConfigEntry(
331 key=CONF_WAVE_PRESET_DRAFT_DIVERSITY,
332 type=ConfigEntryType.STRING,
333 options=diversity_options,
334 default_value="",
335 required=False,
336 advanced=True,
337 ),
338 ConfigEntry(
339 key=CONF_WAVE_PRESET_DRAFT_MOOD,
340 type=ConfigEntryType.STRING,
341 options=mood_options,
342 default_value="",
343 required=False,
344 advanced=True,
345 ),
346 ConfigEntry(
347 key=CONF_WAVE_PRESET_DRAFT_LANGUAGE,
348 type=ConfigEntryType.STRING,
349 options=language_options,
350 default_value="",
351 required=False,
352 advanced=True,
353 ),
354 ConfigEntry(
355 key=CONF_ACTION_SAVE_WAVE_PRESET,
356 type=ConfigEntryType.ACTION,
357 action=CONF_ACTION_SAVE_WAVE_PRESET,
358 advanced=True,
359 ),
360 ConfigEntry(
361 key=CONF_WAVE_PRESET_TO_DELETE,
362 type=ConfigEntryType.STRING,
363 options=delete_options,
364 default_value="",
365 required=False,
366 advanced=True,
367 hidden=not has_presets,
368 ),
369 ConfigEntry(
370 key=CONF_ACTION_DELETE_WAVE_PRESET,
371 type=ConfigEntryType.ACTION,
372 action=CONF_ACTION_DELETE_WAVE_PRESET,
373 advanced=True,
374 hidden=not has_presets,
375 ),
376 ConfigEntry(
377 key=CONF_WAVE_PRESETS_DATA,
378 type=ConfigEntryType.STRING,
379 default_value="",
380 required=False,
381 advanced=True,
382 hidden=True,
383 ),
384 ]
385
386
387class _WaveState:
388 """
389 Per-station mutable state for rotor wave playback.
390
391 Holds both the new session-based rotor identifiers (`session_id`) and the
392 legacy stations-based ones (`batch_id`). Call sites prefer `session_id`
393 when present; `batch_id` is still carried because feedback events anchor
394 to a specific batch within the session.
395 """
396
397 def __init__(self) -> None:
398 self.session_id: str | None = None
399 self.batch_id: str | None = None
400 self.last_track_id: str | None = None
401 self.playlist_next_cursor: str | None = None
402 self.seen_track_ids: set[str] = set()
403 self.radio_started_sent: bool = False
404 self.prefetched: list[Any] = []
405 self.settings: dict[str, str] = {}
406 self.lock: asyncio.Lock = asyncio.Lock()
407
408
409class YandexMusicProvider(MusicProvider):
410 """Implementation of a Yandex Music MusicProvider."""
411
412 _client: YandexMusicClient | None = None
413 _streaming: YandexMusicStreamingManager | None = None
414 _wave_states: dict[str, _WaveState] # Per-station state (incl. My Wave)
415 _wave_bg_colors: dict[str, str] # image_url -> hex bg color for transparent covers
416 # Short-lived cache to dedupe the three library syncs (albums/podcasts/audiobooks)
417 # that all derive from the same liked-albums endpoint.
418 _liked_albums_cache: tuple[float, list[YandexAlbum]] | None = None
419 _liked_albums_lock: asyncio.Lock
420 # Per-audiobook cache of (chapter_track_ids, chapter_durations_ms) used to
421 # report playback progress per chapter via play_audio.
422 _audiobook_chapter_cache: dict[str, tuple[list[str], list[int]]]
423 # Stable play_id per audiobook session, cleared in on_streamed.
424 _audiobook_play_ids: dict[str, str]
425
426 @property
427 def client(self) -> YandexMusicClient:
428 """Return the Yandex Music client."""
429 if self._client is None:
430 raise ProviderUnavailableError("Provider not initialized")
431 return self._client
432
433 @property
434 def streaming(self) -> YandexMusicStreamingManager:
435 """Return the streaming manager."""
436 if self._streaming is None:
437 raise ProviderUnavailableError("Provider not initialized")
438 return self._streaming
439
440 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
441 """
442 Return Config entries to configure this provider.
443
444 Authentication runs in the interactive setup flow (see setup_flow.py). This
445 surface exposes a one-shot advanced token replacement alongside playback options
446 and the My Wave preset builder (whose actions use ``handle_config_action``).
447 """
448 return (
449 CONF_ENTRY_UNOFFICIAL_PROVIDER,
450 # Quality
451 ConfigEntry(
452 key=CONF_QUALITY,
453 type=ConfigEntryType.STRING,
454 options=[
455 ConfigValueOption(QUALITY_EFFICIENT),
456 ConfigValueOption(QUALITY_BALANCED),
457 ConfigValueOption(QUALITY_HIGH),
458 ConfigValueOption(QUALITY_SUPERB),
459 ],
460 default_value=QUALITY_BALANCED,
461 ),
462 # My Wave maximum tracks (advanced)
463 ConfigEntry(
464 key=CONF_MY_WAVE_MAX_TRACKS,
465 type=ConfigEntryType.INTEGER,
466 range=(10, 1000),
467 default_value=150,
468 required=False,
469 advanced=True,
470 ),
471 # User-defined wave presets: builder + save/delete actions (dynamic list)
472 *_wave_preset_config_entries(
473 self.get_config_value(CONF_WAVE_PRESETS_DATA, return_type=str)
474 ),
475 # Liked Tracks maximum tracks (advanced)
476 ConfigEntry(
477 key=CONF_LIKED_TRACKS_MAX_TRACKS,
478 type=ConfigEntryType.INTEGER,
479 range=(50, 2000),
480 default_value=200,
481 required=False,
482 advanced=True,
483 ),
484 # API Base URL (advanced)
485 ConfigEntry(
486 key=CONF_BASE_URL,
487 type=ConfigEntryType.STRING,
488 translation_params=[DEFAULT_BASE_URL],
489 default_value=DEFAULT_BASE_URL,
490 required=False,
491 advanced=True,
492 ),
493 # Restrictive rate limits (advanced)
494 ConfigEntry(
495 key=CONF_RESTRICTIVE_RATE_LIMITS,
496 type=ConfigEntryType.BOOLEAN,
497 default_value=False,
498 required=False,
499 advanced=True,
500 ),
501 # One-shot manual token replacement (advanced)
502 ConfigEntry(
503 key=CONF_MANUAL_TOKEN,
504 type=ConfigEntryType.SECURE_STRING,
505 required=False,
506 advanced=True,
507 requires_reload=True,
508 ),
509 )
510
511 async def handle_config_action(
512 self, action: str
513 ) -> tuple[ConfigEntry, ...] | ConfigActionResult | None:
514 """
515 Handle a wave-preset save/delete button press and re-render the entries.
516
517 Both actions mutate the hidden JSON store and clear the draft / selection
518 fields so the UI re-renders in a clean state. Draft values are read from
519 stored config (no form values are passed) and persisted immediately.
520
521 :param action: The action id of the pressed button.
522 """
523 if action == CONF_ACTION_SAVE_WAVE_PRESET:
524 new_presets = _merge_wave_preset(
525 self.get_config_value(CONF_WAVE_PRESET_DRAFT_NAME, return_type=str),
526 self.get_config_value(CONF_WAVE_PRESET_DRAFT_DIVERSITY, return_type=str),
527 self.get_config_value(CONF_WAVE_PRESET_DRAFT_MOOD, return_type=str),
528 self.get_config_value(CONF_WAVE_PRESET_DRAFT_LANGUAGE, return_type=str),
529 self.get_config_value(CONF_WAVE_PRESETS_DATA, return_type=str),
530 )
531 self._update_config_value(CONF_WAVE_PRESETS_DATA, new_presets, immediate=True)
532 # Clear draft so the UI is ready for the next preset
533 self._update_config_value(CONF_WAVE_PRESET_DRAFT_NAME, None, immediate=True)
534 self._update_config_value(CONF_WAVE_PRESET_DRAFT_DIVERSITY, "", immediate=True)
535 self._update_config_value(CONF_WAVE_PRESET_DRAFT_MOOD, "", immediate=True)
536 self._update_config_value(CONF_WAVE_PRESET_DRAFT_LANGUAGE, "", immediate=True)
537 return await self.get_config_entries()
538 if action == CONF_ACTION_DELETE_WAVE_PRESET:
539 new_presets = _remove_wave_preset(
540 self.get_config_value(CONF_WAVE_PRESET_TO_DELETE, return_type=str),
541 self.get_config_value(CONF_WAVE_PRESETS_DATA, return_type=str),
542 )
543 self._update_config_value(CONF_WAVE_PRESETS_DATA, new_presets, immediate=True)
544 self._update_config_value(CONF_WAVE_PRESET_TO_DELETE, "", immediate=True)
545 return await self.get_config_entries()
546 return await super().handle_config_action(action)
547
548 async def handle_async_init(self) -> None: # noqa: PLR0915
549 """Handle async initialization of the provider."""
550 manual_token = self.config.get_value(CONF_MANUAL_TOKEN)
551 token = self.get_setup_value(CONF_TOKEN)
552 x_token = self.get_setup_value(CONF_X_TOKEN)
553 refresh_token = self.get_setup_value(CONF_REFRESH_TOKEN)
554 base_url = self.config.get_value(CONF_BASE_URL, DEFAULT_BASE_URL)
555 restrictive = bool(self.config.get_value(CONF_RESTRICTIVE_RATE_LIMITS, False))
556 replacing_token = bool(manual_token)
557
558 if replacing_token:
559 token = str(manual_token)
560 x_token = None
561 refresh_token = None
562
563 if not token and not x_token:
564 raise LoginFailed("No Yandex Music token provided. Please authenticate.")
565
566 # Try existing music token first (fast path)
567 if token:
568 try:
569 self._client = YandexMusicClient(
570 SecretStr(str(token)),
571 base_url=str(base_url),
572 restrictive_rate_limits=restrictive,
573 )
574 await self._client.connect()
575 except LoginFailed:
576 if replacing_token:
577 self.logger.warning("Manually supplied music token was rejected")
578 self._client = None
579 self._update_config_value(CONF_MANUAL_TOKEN, None, immediate=True)
580 raise
581 self.logger.warning("Music token is invalid or expired")
582 # Clear the dead token so restarts go straight to refresh
583 self._update_setup_data(CONF_TOKEN, None)
584 if x_token:
585 self.logger.info("Attempting to refresh from session token")
586 token = None
587 self._client = None
588 else:
589 raise
590
591 if replacing_token:
592 self._update_setup_data(CONF_TOKEN, str(token))
593 self._update_setup_data(CONF_X_TOKEN, None)
594 self._update_setup_data(CONF_REFRESH_TOKEN, None)
595 self._update_config_value(CONF_MANUAL_TOKEN, None, immediate=True)
596
597 # Refresh from x_token if music token absent or failed
598 if not token and x_token:
599 try:
600 new_music_token = await refresh_music_token(SecretStr(str(x_token)))
601 self._update_setup_data(CONF_TOKEN, new_music_token.get_secret())
602 self._client = YandexMusicClient(
603 new_music_token,
604 base_url=str(base_url),
605 restrictive_rate_limits=restrictive,
606 )
607 await self._client.connect()
608 self.logger.info("Refreshed music token from session token")
609 except LoginFailed as err:
610 # x_token refresh failed. If a refresh_token is available
611 # (device-flow accounts), try silent re-issue of the full
612 # credential triple before giving up.
613 if refresh_token:
614 await self._reauth_via_refresh_token(
615 str(x_token), str(refresh_token), str(base_url), err
616 )
617 else:
618 # Definitive auth failure â clear dead credentials
619 self.logger.warning("Session token is invalid or expired")
620 self._update_setup_data(CONF_TOKEN, None)
621 self._update_setup_data(CONF_X_TOKEN, None)
622 raise LoginFailed("Session token expired. Please re-authenticate.") from err
623 except asyncio.CancelledError:
624 raise
625 except Exception as err:
626 # Transient/network failure â keep credentials for retry
627 self.logger.warning(
628 "Session token refresh failed (network): %s",
629 type(err).__name__,
630 )
631 raise ProviderUnavailableError(
632 "Unable to refresh music token right now. Please try again later."
633 ) from err
634
635 # Suppress yandex_music library DEBUG dumps (full API request/response JSON)
636 logging.getLogger("yandex_music").setLevel(self.logger.level + 10)
637 # Propagate the MA instance log level to our per-module loggers
638 # (api_client, streaming, parsers, auth) so DEBUG hooks there actually
639 # print when MA is set to DEBUG for this provider.
640 logging.getLogger("music_assistant.providers.yandex_music").setLevel(self.logger.level)
641 self._streaming = YandexMusicStreamingManager(self)
642 # Per-station wave state (incl. My Wave under ROTOR_STATION_MY_WAVE).
643 # Entries are created lazily by _get_wave_state() on first access.
644 self._wave_states = {}
645 self._wave_bg_colors = {}
646 self._liked_albums_lock, self._liked_albums_cache = asyncio.Lock(), None
647 self._audiobook_chapter_cache, self._audiobook_play_ids = {}, {}
648 self.logger.info("Successfully connected to Yandex Music")
649
650 async def unload(self, is_removed: bool = False) -> None:
651 """
652 Handle unload/close of the provider.
653
654 :param is_removed: Whether the provider is being removed.
655 """
656 if self._client:
657 await self._client.disconnect()
658 self._client = None
659 self._streaming = None
660 self._wave_states.clear()
661 self._wave_bg_colors.clear()
662 self._liked_albums_cache = None
663 self._audiobook_chapter_cache.clear()
664 self._audiobook_play_ids.clear()
665 await super().unload(is_removed)
666
667 def get_item_mapping(self, media_type: MediaType | str, key: str, name: str) -> ItemMapping:
668 """
669 Create a generic item mapping.
670
671 :param media_type: The media type.
672 :param key: The item ID.
673 :param name: The item name.
674 :return: An ItemMapping instance.
675 """
676 if isinstance(media_type, str):
677 media_type = MediaType(media_type)
678 return ItemMapping(
679 media_type=media_type,
680 item_id=key,
681 provider=self.instance_id,
682 name=name,
683 )
684
685 async def browse( # noqa: PLR0911, PLR0915
686 self, path: str
687 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
688 """
689 Browse provider items with locale-based folder names and My Wave.
690
691 Root level shows My Wave, artists, albums, liked tracks, playlists. Names
692 are in Russian when MA locale is ru_*, otherwise in English. My Wave
693 tracks use item_id format track_id@station_id for rotor feedback.
694
695 :param path: The path to browse (e.g. provider_id:// or provider_id://artists).
696 """
697 if ProviderFeature.BROWSE not in self.supported_features:
698 raise NotImplementedError
699
700 path_parts = path.split("://")[1].split("/") if "://" in path else []
701 subpath = path_parts[0] if len(path_parts) > 0 else None
702 sub_subpath = path_parts[1] if len(path_parts) > 1 else None
703
704 if subpath == MY_WAVE_PLAYLIST_ID:
705 async with self._get_wave_state(ROTOR_STATION_MY_WAVE).lock:
706 return await self._browse_my_wave(path, sub_subpath)
707
708 # Wave modes â accept two equivalent URL forms so both browse
709 # navigation (slash form "my_wave_modes/<preset>", emitted by our
710 # listing) and MA's play-time reconstruction (underscore form
711 # "my_wave_modes_<preset>", built as "<instance>://<item_id>") work.
712 mode_preset: str | None = None
713 if subpath == MY_WAVE_MODES_FOLDER_ID and sub_subpath is None:
714 return self._browse_my_wave_modes_list(path)
715 if subpath == MY_WAVE_MODES_FOLDER_ID and sub_subpath is not None:
716 mode_preset = sub_subpath if sub_subpath != "next" else None
717 if mode_preset is None:
718 return []
719 load_more_modes = len(path_parts) > 2 and path_parts[2] == "next"
720 elif subpath and subpath.startswith(f"{MY_WAVE_MODES_FOLDER_ID}_"):
721 mode_preset = subpath[len(MY_WAVE_MODES_FOLDER_ID) + 1 :]
722 load_more_modes = sub_subpath == "next"
723 if mode_preset is not None:
724 if mode_preset not in WAVE_MODE_PRESETS:
725 return []
726 station_key = f"{ROTOR_STATION_MY_WAVE}{WAVE_MODE_SEP}{mode_preset}"
727 async with self._get_wave_state(station_key).lock:
728 return await self._browse_my_wave_mode(path, station_key, load_more_modes)
729
730 # User-saved wave presets â same dual-form handling.
731 preset_idx: int | None = None
732 load_more_presets = False
733 if subpath == MY_WAVE_PRESETS_FOLDER_ID and sub_subpath is None:
734 return self._browse_user_presets_list(path, self._get_user_wave_presets())
735 if subpath == MY_WAVE_PRESETS_FOLDER_ID and sub_subpath is not None:
736 try:
737 preset_idx = int(sub_subpath)
738 except ValueError:
739 return []
740 load_more_presets = len(path_parts) > 2 and path_parts[2] == "next"
741 elif subpath and subpath.startswith(f"{MY_WAVE_PRESETS_FOLDER_ID}_"):
742 try:
743 preset_idx = int(subpath[len(MY_WAVE_PRESETS_FOLDER_ID) + 1 :])
744 except ValueError:
745 return []
746 load_more_presets = sub_subpath == "next"
747 if preset_idx is not None:
748 user_presets = self._get_user_wave_presets()
749 if not 0 <= preset_idx < len(user_presets):
750 return []
751 preset_data = user_presets[preset_idx]
752 station_key = f"{ROTOR_STATION_MY_WAVE}{WAVE_MODE_SEP}preset_{preset_idx}"
753 wave = self._get_wave_state(station_key)
754 # Stash user-chosen settings so _fetch_rotor_session_batch sends them
755 wave.settings = {
756 k: v
757 for k, v in preset_data.items()
758 if k in ("diversity", "moodEnergy", "language") and v
759 }
760 async with wave.lock:
761 return await self._browse_my_wave_mode(path, station_key, load_more_presets)
762
763 # For You folder (picks + mixes)
764 if subpath == FOR_YOU_FOLDER_ID:
765 return await self._browse_for_you(path, path_parts)
766
767 # Collection folder (library items). Two shapes:
768 # <prov>://collection â listing of library sub-folders
769 # <prov>://collection/<sub> â delegate to MA's library handler
770 # The nested form is what lets MA's "back" button return here (strip
771 # last /-segment) instead of dumping the user at the provider root.
772 if subpath == COLLECTION_FOLDER_ID:
773 if sub_subpath in _COLLECTION_SUB_FOLDERS:
774 return await super().browse(f"{self.instance_id}://{sub_subpath}")
775 return await self._browse_collection(path)
776
777 # Handle picks/ path (mood, activity, era, genres)
778 if subpath == "picks":
779 return await self._browse_picks(path, path_parts)
780
781 # Handle mixes/ path (seasonal collections)
782 if subpath == "mixes":
783 return await self._browse_mixes(path, path_parts)
784
785 # Handle waves/ and radio/ paths (rotor stations by genre/mood/activity)
786 if subpath in (WAVES_FOLDER_ID, RADIO_FOLDER_ID):
787 return await self._browse_waves(path, path_parts)
788
789 # Handle my_waves_set/ path (AI Wave Sets from /landing-blocks/mixes-waves)
790 if subpath == MY_WAVES_SET_FOLDER_ID:
791 return await self._browse_vibe_sets(path, path_parts)
792
793 # Pinned items folder
794 if subpath == PINNED_ITEMS_FOLDER_ID:
795 return await self._browse_pins()
796
797 # Listening history folder
798 if subpath == LISTENING_HISTORY_FOLDER_ID:
799 return await self._browse_history()
800
801 # Handle waves_landing/ path (Featured Waves from /landing-blocks/waves)
802 if subpath == WAVES_LANDING_FOLDER_ID:
803 return await self._browse_waves_landing(path, path_parts)
804
805 # Handle direct tag subpath (when folder is played by URI, the full path
806 # "picks/category/tag" is lost and only the tag slug arrives as subpath).
807 # Skip the API call for standard top-level folders that are never tag slugs.
808 _known_folders = {
809 "artists",
810 "albums",
811 "tracks",
812 "playlists",
813 "audiobooks",
814 "podcasts",
815 LIKED_TRACKS_PLAYLIST_ID,
816 WAVES_FOLDER_ID,
817 RADIO_FOLDER_ID,
818 MY_WAVES_FOLDER_ID,
819 MY_WAVES_SET_FOLDER_ID,
820 WAVES_LANDING_FOLDER_ID,
821 FOR_YOU_FOLDER_ID,
822 COLLECTION_FOLDER_ID,
823 PINNED_ITEMS_FOLDER_ID,
824 LISTENING_HISTORY_FOLDER_ID,
825 }
826 if subpath and subpath not in _known_folders:
827 # Handle direct wave station_id (e.g. "activity:workout") passed when
828 # MA plays a wave station folder using its item_id as the path subpath.
829 # Station IDs have format "category:tag" where category is non-numeric.
830 if ":" in subpath:
831 cat_part = subpath.split(":", 1)[0]
832 if not cat_part.isdigit():
833 return await self._browse_wave_station(subpath)
834
835 discovered_tags = await self._get_discovered_tag_slugs()
836 if subpath in discovered_tags:
837 return await self._get_tag_playlists_as_browse(subpath)
838
839 if subpath:
840 return await super().browse(path)
841
842 # The English name on each folder doubles as the fallback; translation_key localizes
843 # it for the connection locale at serialization (the server is the single source).
844 items: list[MediaItemType | ItemMapping | BrowseFolder] = []
845 base = path if path.endswith("//") else path.rstrip("/") + "/"
846 # My Wave is a dynamic playlist so the queue can request refills.
847 items.append(await self.get_playlist(MY_WAVE_PLAYLIST_ID))
848 # Wave modes folder (P4): discover / calm / active / language presets
849 items.append(
850 BrowseFolder(
851 item_id=MY_WAVE_MODES_FOLDER_ID,
852 provider=self.instance_id,
853 path=f"{base}{MY_WAVE_MODES_FOLDER_ID}",
854 name="Wave Modes",
855 translation_key=MY_WAVE_MODES_FOLDER_ID,
856 is_playable=False,
857 )
858 )
859 # User-defined wave presets (P8) â shown only when any configured.
860 if self._get_user_wave_presets():
861 items.append(
862 BrowseFolder(
863 item_id=MY_WAVE_PRESETS_FOLDER_ID,
864 provider=self.instance_id,
865 path=f"{base}{MY_WAVE_PRESETS_FOLDER_ID}",
866 name="My Presets",
867 translation_key=MY_WAVE_PRESETS_FOLDER_ID,
868 is_playable=False,
869 )
870 )
871 # For You folder â Picks + Mixes (Ð¯Ð½Ð´ÐµÐºÑ Â«ÐÐ»Ñ Ð²Ð°Ñ»)
872 items.append(
873 BrowseFolder(
874 item_id=FOR_YOU_FOLDER_ID,
875 provider=self.instance_id,
876 path=f"{base}{FOR_YOU_FOLDER_ID}",
877 name="For You",
878 translation_key=FOR_YOU_FOLDER_ID,
879 is_playable=False,
880 )
881 )
882 # Collection folder â library items (Ð¯Ð½Ð´ÐµÐºÑ Â«ÐоллекÑиÑ»)
883 has_library = any(
884 f in self.supported_features
885 for f in (
886 ProviderFeature.LIBRARY_ARTISTS,
887 ProviderFeature.LIBRARY_ALBUMS,
888 ProviderFeature.LIBRARY_TRACKS,
889 ProviderFeature.LIBRARY_PLAYLISTS,
890 )
891 )
892 if has_library:
893 items.append(
894 BrowseFolder(
895 item_id=COLLECTION_FOLDER_ID,
896 provider=self.instance_id,
897 path=f"{base}{COLLECTION_FOLDER_ID}",
898 name="Collection",
899 translation_key=COLLECTION_FOLDER_ID,
900 is_playable=False,
901 )
902 )
903 # Radio folder â rotor stations (Ð¯Ð½Ð´ÐµÐºÑ Ð²Ð¾Ð»Ð½Ñ, shown as Radio)
904 items.append(
905 BrowseFolder(
906 item_id=RADIO_FOLDER_ID,
907 provider=self.instance_id,
908 path=f"{base}{RADIO_FOLDER_ID}",
909 name="Radio",
910 translation_key=RADIO_FOLDER_ID,
911 is_playable=False,
912 )
913 )
914 # AI Wave Sets â parametric stations from /landing-blocks/mixes-waves
915 items.append(
916 BrowseFolder(
917 item_id=MY_WAVES_SET_FOLDER_ID,
918 provider=self.instance_id,
919 path=f"{base}{MY_WAVES_SET_FOLDER_ID}",
920 name="AI Wave Sets",
921 translation_key=MY_WAVES_SET_FOLDER_ID,
922 is_playable=False,
923 )
924 )
925 # Pinned items â user-pinned artists/albums/playlists/waves
926 items.append(
927 BrowseFolder(
928 item_id=PINNED_ITEMS_FOLDER_ID,
929 provider=self.instance_id,
930 path=f"{base}{PINNED_ITEMS_FOLDER_ID}",
931 name="Pinned",
932 translation_key=PINNED_ITEMS_FOLDER_ID,
933 is_playable=False,
934 )
935 )
936 # Listening history â recently played tracks/albums
937 items.append(
938 BrowseFolder(
939 item_id=LISTENING_HISTORY_FOLDER_ID,
940 provider=self.instance_id,
941 path=f"{base}{LISTENING_HISTORY_FOLDER_ID}",
942 name="Listening History",
943 translation_key=LISTENING_HISTORY_FOLDER_ID,
944 is_playable=False,
945 )
946 )
947 if len(items) == 1 and isinstance(items[0], BrowseFolder):
948 return await self.browse(items[0].path)
949 return items
950
951 # Search
952
953 @use_cache(3600 * 24, allow_expired_cache=True)
954 async def search(
955 self, search_query: str, media_types: list[MediaType], limit: int = 5
956 ) -> SearchResults:
957 """
958 Perform search on Yandex Music.
959
960 :param search_query: The search query.
961 :param media_types: List of media types to search for.
962 :param limit: Maximum number of results per type.
963 :return: SearchResults with found items.
964 """
965 result = SearchResults()
966
967 # Determine search type based on requested media types
968 # Map MediaType to Yandex API search type. AUDIOBOOK has no dedicated
969 # Yandex type â it maps to "album" and is filtered by classify_album below.
970 type_mapping = {
971 MediaType.TRACK: "track",
972 MediaType.ALBUM: "album",
973 MediaType.AUDIOBOOK: "album",
974 MediaType.ARTIST: "artist",
975 MediaType.PLAYLIST: "playlist",
976 MediaType.PODCAST: "podcast",
977 }
978 requested_types = list(
979 dict.fromkeys(type_mapping[mt] for mt in media_types if mt in type_mapping)
980 )
981
982 # Use specific type if only one requested, otherwise search all
983 search_type = requested_types[0] if len(requested_types) == 1 else "all"
984
985 search_result = await self.client.search(search_query, search_type=search_type)
986 if not search_result:
987 return result
988
989 # Parse tracks
990 if MediaType.TRACK in media_types and search_result.tracks:
991 for track in search_result.tracks.results[:limit]:
992 try:
993 result.tracks = [*result.tracks, parse_track(self, track)]
994 except InvalidDataError as err:
995 self.logger.debug("Error parsing track: %s", err)
996
997 # Parse albums â audiobooks are split into the audiobooks bucket via
998 # classify_album. Yandex-returned podcast albums are handled separately
999 # through the dedicated `.podcasts` node below. ``limit`` is applied per
1000 # bucket AFTER classification â slicing first would drop audiobooks when
1001 # the first ``limit`` results happen to be music albums (or vice versa).
1002 want_album = MediaType.ALBUM in media_types
1003 want_audiobook = MediaType.AUDIOBOOK in media_types
1004 if (want_album or want_audiobook) and search_result.albums:
1005 album_count = 0
1006 audiobook_count = 0
1007 for album in search_result.albums.results:
1008 album_full = not want_album or album_count >= limit
1009 audiobook_full = not want_audiobook or audiobook_count >= limit
1010 if album_full and audiobook_full:
1011 break
1012 kind = classify_album(album)
1013 try:
1014 if kind == "audiobook" and want_audiobook and not audiobook_full:
1015 result.audiobooks = [
1016 *result.audiobooks,
1017 parse_audiobook(self, album),
1018 ]
1019 audiobook_count += 1
1020 elif kind == "music" and want_album and not album_full:
1021 result.albums = [*result.albums, parse_album(self, album)]
1022 album_count += 1
1023 except InvalidDataError as err:
1024 self.logger.debug("Error parsing %s album: %s", kind, err)
1025
1026 # Parse artists
1027 if MediaType.ARTIST in media_types and search_result.artists:
1028 for artist in search_result.artists.results[:limit]:
1029 try:
1030 result.artists = [*result.artists, parse_artist(self, artist)]
1031 except InvalidDataError as err:
1032 self.logger.debug("Error parsing artist: %s", err)
1033
1034 # Parse playlists
1035 if MediaType.PLAYLIST in media_types and search_result.playlists:
1036 for playlist in search_result.playlists.results[:limit]:
1037 try:
1038 result.playlists = [*result.playlists, parse_playlist(self, playlist)]
1039 except InvalidDataError as err:
1040 self.logger.debug("Error parsing playlist: %s", err)
1041
1042 # Parse podcasts (Yandex returns them as albums under .podcasts)
1043 podcasts_node = getattr(search_result, "podcasts", None)
1044 if MediaType.PODCAST in media_types and podcasts_node:
1045 for album in podcasts_node.results[:limit]:
1046 try:
1047 result.podcasts = [*result.podcasts, parse_podcast(self, album)]
1048 except InvalidDataError as err:
1049 self.logger.debug("Error parsing podcast: %s", err)
1050
1051 return result
1052
1053 # Get single items
1054
1055 @use_cache(3600 * 24 * 30, allow_expired_cache=True)
1056 async def get_artist(self, prov_artist_id: str) -> Artist:
1057 """
1058 Get artist details by ID, enriched with description and listener stats.
1059
1060 :param prov_artist_id: The provider artist ID.
1061 :return: Artist object.
1062 :raises MediaNotFoundError: If artist not found.
1063 """
1064 artist, about = await asyncio.gather(
1065 self.client.get_artist(prov_artist_id),
1066 self.client.get_artist_about(prov_artist_id),
1067 )
1068 if not artist:
1069 raise MediaNotFoundError(f"Artist {prov_artist_id} not found")
1070 return parse_artist(self, artist, about=about)
1071
1072 @use_cache(3600 * 24 * 30, allow_expired_cache=True)
1073 async def get_album(self, prov_album_id: str) -> Album:
1074 """
1075 Get album details by ID.
1076
1077 :param prov_album_id: The provider album ID.
1078 :return: Album object.
1079 :raises MediaNotFoundError: If album not found.
1080 """
1081 album = await self.client.get_album(prov_album_id)
1082 if not album:
1083 raise MediaNotFoundError(f"Album {prov_album_id} not found")
1084 return parse_album(self, album)
1085
1086 @use_cache(3600 * 24, allow_expired_cache=True)
1087 async def get_podcast(self, prov_podcast_id: str) -> Podcast:
1088 """
1089 Get podcast details by ID (backed by a Yandex album).
1090
1091 :param prov_podcast_id: The provider podcast (album) ID.
1092 :return: Podcast object.
1093 :raises MediaNotFoundError: If not found.
1094 """
1095 album = await self.client.get_album(prov_podcast_id)
1096 if not album:
1097 raise MediaNotFoundError(f"Podcast {prov_podcast_id} not found")
1098 return parse_podcast(self, album)
1099
1100 async def get_podcast_episodes(self, prov_podcast_id: str) -> AsyncGenerator[PodcastEpisode]:
1101 """Iterate podcast episodes for a given podcast (album) ID."""
1102 album = await self.client.get_album_with_tracks(prov_podcast_id)
1103 if not album:
1104 raise MediaNotFoundError(f"Podcast {prov_podcast_id} not found")
1105 podcast = parse_podcast(self, album)
1106 position = 1
1107 for disc in album.volumes or []:
1108 for track_obj in disc:
1109 try:
1110 yield parse_podcast_episode(self, track_obj, podcast, position=position)
1111 except InvalidDataError as err:
1112 self.logger.debug("Error parsing podcast episode: %s", err)
1113 position += 1
1114
1115 async def get_podcast_episode(self, prov_episode_id: str) -> PodcastEpisode:
1116 """
1117 Get a single podcast episode by ID.
1118
1119 The parent Podcast is reconstructed from the track's parent album. If
1120 the album isn't present on the track, the episode cannot be converted
1121 into a valid MA model and InvalidDataError is raised.
1122 """
1123 tracks = await self.client.get_tracks([prov_episode_id])
1124 if not tracks:
1125 raise MediaNotFoundError(f"Podcast episode {prov_episode_id} not found")
1126 track_obj = tracks[0]
1127 if not track_obj.albums:
1128 raise InvalidDataError(
1129 f"Podcast episode {prov_episode_id} is missing parent podcast album data"
1130 )
1131 podcast = parse_podcast(self, track_obj.albums[0])
1132 return parse_podcast_episode(self, track_obj, podcast, position=0)
1133
1134 @use_cache(3600 * 24, allow_expired_cache=True)
1135 async def get_audiobook(self, prov_audiobook_id: str) -> Audiobook:
1136 """
1137 Get audiobook details by ID, including chapters built from tracks.
1138
1139 :param prov_audiobook_id: The provider audiobook (album) ID.
1140 :return: Audiobook object.
1141 :raises MediaNotFoundError: If not found.
1142 """
1143 album = await self.client.get_album_with_tracks(prov_audiobook_id)
1144 if not album:
1145 raise MediaNotFoundError(f"Audiobook {prov_audiobook_id} not found")
1146 audiobook = parse_audiobook(self, album)
1147
1148 chapters: list[MediaItemChapter] = []
1149 start = 0.0
1150 pos = 1
1151 for disc in album.volumes or []:
1152 for track_obj in disc:
1153 dur_s = (track_obj.duration_ms or 0) / 1000.0
1154 chapters.append(
1155 MediaItemChapter(
1156 position=pos,
1157 name=track_obj.title or f"Chapter {pos}",
1158 start=start,
1159 end=start + dur_s,
1160 )
1161 )
1162 start += dur_s
1163 pos += 1
1164 audiobook.metadata.chapters = chapters
1165 audiobook.duration = int(start)
1166 return audiobook
1167
1168 async def get_track(self, prov_track_id: str) -> Track:
1169 """
1170 Get track details by ID.
1171
1172 Supports composite item_id (track_id@station_id) for My Wave tracks;
1173 only the track_id part is used for the API. Normalizes the ID before
1174 caching to avoid duplicate cache entries.
1175
1176 :param prov_track_id: The provider track ID (or track_id@station_id).
1177 :return: Track object.
1178 :raises MediaNotFoundError: If track not found.
1179 """
1180 track_id, _ = _parse_radio_item_id(prov_track_id)
1181 return await self._get_track_cached(track_id)
1182
1183 async def get_playlist(self, prov_playlist_id: str) -> Playlist:
1184 """
1185 Get playlist details by ID.
1186
1187 Supports virtual playlists MY_WAVE_PLAYLIST_ID (My Wave) and
1188 LIKED_TRACKS_PLAYLIST_ID (Liked Tracks). Real playlists use format "owner_id:kind".
1189
1190 :param prov_playlist_id: The provider playlist ID (format: "owner_id:kind",
1191 my_wave, or liked_tracks).
1192 :return: Playlist object.
1193 :raises MediaNotFoundError: If playlist not found.
1194 """
1195 # Virtual playlists - constructed locally (no API call); translation_key localizes
1196 # the name for the connection locale at serialization.
1197 if prov_playlist_id == MY_WAVE_PLAYLIST_ID:
1198 return Playlist(
1199 item_id=MY_WAVE_PLAYLIST_ID,
1200 provider=self.instance_id,
1201 name="My Wave",
1202 translation_key=MY_WAVE_PLAYLIST_ID,
1203 owner=get_canonical_provider_name(self),
1204 provider_mappings={
1205 ProviderMapping(
1206 item_id=MY_WAVE_PLAYLIST_ID,
1207 provider_domain=self.domain,
1208 provider_instance=self.instance_id,
1209 is_unique=True,
1210 )
1211 },
1212 is_editable=False,
1213 is_dynamic=True,
1214 )
1215
1216 if prov_playlist_id == LIKED_TRACKS_PLAYLIST_ID:
1217 return Playlist(
1218 item_id=LIKED_TRACKS_PLAYLIST_ID,
1219 provider=self.instance_id,
1220 name="My Favorites",
1221 translation_key=LIKED_TRACKS_PLAYLIST_ID,
1222 owner=get_canonical_provider_name(self),
1223 provider_mappings={
1224 ProviderMapping(
1225 item_id=LIKED_TRACKS_PLAYLIST_ID,
1226 provider_domain=self.domain,
1227 provider_instance=self.instance_id,
1228 is_unique=True,
1229 )
1230 },
1231 is_editable=False,
1232 )
1233
1234 # Real playlists - use cached method
1235 return await self._get_real_playlist(prov_playlist_id)
1236
1237 # Get related items
1238
1239 @use_cache(3600 * 24 * 30, allow_expired_cache=True)
1240 async def get_album_tracks(self, prov_album_id: str) -> list[Track]:
1241 """
1242 Get album tracks.
1243
1244 :param prov_album_id: The provider album ID.
1245 :return: List of Track objects.
1246 """
1247 album = await self.client.get_album_with_tracks(prov_album_id)
1248 if not album or not album.volumes:
1249 return []
1250
1251 tracks = []
1252 for volume_index, volume in enumerate(album.volumes):
1253 for track_index, track in enumerate(volume):
1254 try:
1255 parsed_track = parse_track(self, track)
1256 parsed_track.disc_number = volume_index + 1
1257 parsed_track.track_number = track_index + 1
1258 tracks.append(parsed_track)
1259 except InvalidDataError as err:
1260 self.logger.debug("Error parsing album track: %s", err)
1261 return tracks
1262
1263 async def get_similar_tracks(self, prov_track_id: str, limit: int = 25) -> list[Track]:
1264 """
1265 Get similar tracks, preferring pre-fetched wave tracks when available.
1266
1267 Split in two paths with different caching policies:
1268
1269 - **Wave-drain path** (the seed carries a station suffix and
1270 ``wave.prefetched`` is non-empty). Uncached by design: it mutates
1271 state, a cache hit would replay the same drained tracks forever and
1272 the prefetch buffer would never advance.
1273 - **Fallback path** (plain track_id, no active wave, or empty buffer).
1274 Creates a per-seed rotor session under ``track:{id}`` and is cached
1275 for 3 hours â this is pure and safe to memoise.
1276
1277 :param prov_track_id: Provider track ID (plain or track_id@station_id).
1278 :param limit: Maximum number of tracks to return.
1279 :return: List of similar Track objects.
1280 """
1281 track_id, station_key = _parse_radio_item_id(prov_track_id)
1282
1283 if station_key:
1284 drained = await self._drain_prefetched_wave_tracks(station_key, limit)
1285 if drained:
1286 return drained
1287
1288 return await self._fetch_similar_tracks_for_seed(track_id, limit)
1289
1290 @use_cache(3600 * 3, allow_expired_cache=True)
1291 async def get_similar_artists(self, prov_artist_id: str, limit: int = 25) -> list[Artist]:
1292 """
1293 Get artists similar to the given one via Yandex artists/similar endpoint.
1294
1295 :param prov_artist_id: Provider artist ID.
1296 :param limit: Maximum number of artists to return.
1297 :return: List of similar Artist objects.
1298 """
1299 yandex_artists = await self.client.get_similar_artists(prov_artist_id, limit=limit)
1300 artists: list[Artist] = []
1301 for ya in yandex_artists:
1302 try:
1303 artists.append(parse_artist(self, ya))
1304 except InvalidDataError as err:
1305 self.logger.debug("Error parsing similar artist: %s", err)
1306 return artists
1307
1308 async def get_recommendations(self) -> list[RecommendationFolder]:
1309 """Return static recommendation row descriptors without backend calls."""
1310 seasonal_tag = TAG_SEASONAL_MAP.get(utc().month, "autumn")
1311 seasonal_name, _ = self._media_label(
1312 "folder", _media_label_key(seasonal_tag), seasonal_tag.title()
1313 )
1314 return [
1315 RecommendationFolder(
1316 item_id=MY_WAVE_PLAYLIST_ID,
1317 provider=self.instance_id,
1318 name="My Wave",
1319 translation_key=MY_WAVE_PLAYLIST_ID,
1320 icon="mdi-waveform",
1321 ),
1322 RecommendationFolder(
1323 item_id="feed",
1324 provider=self.instance_id,
1325 name="Made for You",
1326 translation_key="feed",
1327 icon="mdi-account-music",
1328 ),
1329 RecommendationFolder(
1330 item_id="chart",
1331 provider=self.instance_id,
1332 name="Chart",
1333 translation_key="chart",
1334 icon="mdi-chart-line",
1335 ),
1336 RecommendationFolder(
1337 item_id="new_releases",
1338 provider=self.instance_id,
1339 name="New Releases",
1340 translation_key="new_releases",
1341 icon="mdi-new-box",
1342 ),
1343 RecommendationFolder(
1344 item_id="new_playlists",
1345 provider=self.instance_id,
1346 name="New Playlists",
1347 translation_key="new_playlists",
1348 icon="mdi-playlist-star",
1349 ),
1350 RecommendationFolder(
1351 item_id="top_picks",
1352 provider=self.instance_id,
1353 name="Top Picks",
1354 translation_key="top_picks",
1355 icon="mdi-star",
1356 ),
1357 RecommendationFolder(
1358 item_id="mood_mix",
1359 provider=self.instance_id,
1360 name="Mood Mix",
1361 translation_key="mood_mix",
1362 subtitle=await self._rotating_row_tag_subtitle("mood"),
1363 icon="mdi-emoticon-outline",
1364 ),
1365 RecommendationFolder(
1366 item_id="activity_mix",
1367 provider=self.instance_id,
1368 name="Activity Mix",
1369 translation_key="activity_mix",
1370 subtitle=await self._rotating_row_tag_subtitle("activity"),
1371 icon="mdi-run",
1372 ),
1373 RecommendationFolder(
1374 item_id="seasonal_mix",
1375 provider=self.instance_id,
1376 name=f"Seasonal: {seasonal_name}",
1377 translation_key="seasonal_mix",
1378 translation_params=[seasonal_name],
1379 icon="mdi-weather-sunny",
1380 ),
1381 ]
1382
1383 async def get_recommendation_items(
1384 self, item_id: str
1385 ) -> UniqueList[MediaItemType | ItemMapping | BrowseFolder]:
1386 """Load items for one recommendation row."""
1387 folder: RecommendationFolder | None = None
1388 if item_id == MY_WAVE_PLAYLIST_ID:
1389 folder = await self._get_my_wave_recommendations()
1390 elif item_id == "feed":
1391 folder = await self._get_feed_recommendations()
1392 elif item_id == "chart":
1393 folder = await self._get_chart_recommendations()
1394 elif item_id == "new_releases":
1395 folder = await self._get_new_releases_recommendations()
1396 elif item_id == "new_playlists":
1397 folder = await self._get_new_playlists_recommendations()
1398 elif item_id == "top_picks":
1399 folder = await self._get_top_picks_recommendations()
1400 elif item_id == "mood_mix":
1401 if tags := await self._get_valid_tags_for_category("mood"):
1402 folder = await self._get_mood_mix_recommendations(
1403 self._rotating_row_tag("mood", tags)
1404 )
1405 elif item_id == "activity_mix":
1406 if tags := await self._get_valid_tags_for_category("activity"):
1407 folder = await self._get_activity_mix_recommendations(
1408 self._rotating_row_tag("activity", tags)
1409 )
1410 elif item_id == "seasonal_mix":
1411 folder = await self._get_seasonal_mix_recommendations()
1412 return folder.items if folder else UniqueList()
1413
1414 async def get_playlist_tracks(self, prov_playlist_id: str, page: int = 0) -> list[Track]:
1415 """
1416 Get playlist tracks.
1417
1418 :param prov_playlist_id: The provider playlist ID (format: "owner_id:kind",
1419 my_wave, or liked_tracks).
1420 :param page: Page number for pagination.
1421 :return: List of Track objects.
1422 """
1423 self.logger.debug(
1424 "get_playlist_tracks called: prov_playlist_id=%s, page=%s", prov_playlist_id, page
1425 )
1426
1427 if prov_playlist_id == MY_WAVE_PLAYLIST_ID:
1428 self.logger.debug("Fetching My Wave tracks")
1429 return await self._get_my_wave_playlist_tracks(page)
1430
1431 if prov_playlist_id == LIKED_TRACKS_PLAYLIST_ID:
1432 self.logger.debug("Fetching Liked Tracks for virtual playlist")
1433 result = await self._get_liked_tracks_playlist_tracks(page)
1434 self.logger.debug("Liked Tracks playlist returned %s tracks", len(result))
1435 return result
1436
1437 return await self._get_regular_playlist_tracks(prov_playlist_id, page)
1438
1439 @use_cache(3600 * 24 * 7, allow_expired_cache=True)
1440 async def get_artist_albums(self, prov_artist_id: str) -> list[Album]:
1441 """
1442 Get artist's albums.
1443
1444 :param prov_artist_id: The provider artist ID.
1445 :return: List of Album objects.
1446 """
1447 albums = await self.client.get_artist_albums(prov_artist_id)
1448 result = []
1449 for album in albums:
1450 try:
1451 result.append(parse_album(self, album))
1452 except InvalidDataError as err:
1453 self.logger.debug("Error parsing artist album: %s", err)
1454 return result
1455
1456 @use_cache(3600 * 24 * 7, allow_expired_cache=True)
1457 async def get_artist_toptracks(self, prov_artist_id: str) -> list[Track]:
1458 """
1459 Get artist's top tracks.
1460
1461 :param prov_artist_id: The provider artist ID.
1462 :return: List of Track objects.
1463 """
1464 tracks = await self.client.get_artist_tracks(prov_artist_id)
1465 result = []
1466 for track in tracks:
1467 try:
1468 result.append(parse_track(self, track))
1469 except InvalidDataError as err:
1470 self.logger.debug("Error parsing artist track: %s", err)
1471 return result
1472
1473 # Library methods
1474
1475 async def get_library_artists(self) -> AsyncGenerator[Artist]:
1476 """Retrieve library artists from Yandex Music."""
1477 artists = await self.client.get_liked_artists()
1478 for artist in artists:
1479 try:
1480 yield parse_artist(self, artist)
1481 except InvalidDataError as err:
1482 self.logger.debug("Error parsing library artist: %s", err)
1483
1484 async def get_library_albums(self) -> AsyncGenerator[Album]:
1485 """
1486 Retrieve library albums from Yandex Music.
1487
1488 Excludes entries classified as podcasts or audiobooks so they don't
1489 duplicate into the Albums library view.
1490 """
1491 for album in await self._get_liked_albums_cached():
1492 if classify_album(album) != "music":
1493 continue
1494 try:
1495 yield parse_album(self, album)
1496 except InvalidDataError as err:
1497 self.logger.debug("Error parsing library album: %s", err)
1498
1499 async def get_library_podcasts(self) -> AsyncGenerator[Podcast]:
1500 """Retrieve library podcasts from Yandex Music (filtered liked albums)."""
1501 for album in await self._get_liked_albums_cached():
1502 if classify_album(album) != "podcast":
1503 continue
1504 try:
1505 yield parse_podcast(self, album)
1506 except InvalidDataError as err:
1507 self.logger.debug("Error parsing library podcast: %s", err)
1508
1509 async def get_library_audiobooks(self) -> AsyncGenerator[Audiobook]:
1510 """Retrieve library audiobooks from Yandex Music (filtered liked albums)."""
1511 for album in await self._get_liked_albums_cached():
1512 if classify_album(album) != "audiobook":
1513 continue
1514 try:
1515 yield parse_audiobook(self, album)
1516 except InvalidDataError as err:
1517 self.logger.debug("Error parsing library audiobook: %s", err)
1518
1519 async def get_library_tracks(self) -> AsyncGenerator[Track]:
1520 """Retrieve library tracks from Yandex Music."""
1521 track_shorts = await self.client.get_liked_tracks()
1522 if not track_shorts:
1523 return
1524
1525 # Fetch full track details in batches
1526 track_ids = [str(ts.track_id) for ts in track_shorts if ts.track_id]
1527 batch_size = TRACK_BATCH_SIZE
1528 for i in range(0, len(track_ids), batch_size):
1529 batch_ids = track_ids[i : i + batch_size]
1530 full_tracks = await self.client.get_tracks(batch_ids)
1531 for track in full_tracks:
1532 try:
1533 yield parse_track(self, track)
1534 except InvalidDataError as err:
1535 self.logger.debug("Error parsing library track: %s", err)
1536
1537 async def get_library_playlists(self) -> AsyncGenerator[Playlist]:
1538 """
1539 Retrieve library playlists from Yandex Music.
1540
1541 Includes virtual playlists (My Wave and Liked Tracks if enabled), user-created playlists,
1542 and user-liked editorial playlists (returned by a separate API endpoint).
1543 """
1544 yield await self.get_playlist(MY_WAVE_PLAYLIST_ID)
1545 yield await self.get_playlist(LIKED_TRACKS_PLAYLIST_ID)
1546 seen_ids: set[str] = set()
1547 # User-created playlists
1548 playlists = await self.client.get_user_playlists()
1549 for playlist in playlists:
1550 try:
1551 parsed = parse_playlist(self, playlist)
1552 seen_ids.add(parsed.item_id)
1553 yield parsed
1554 except InvalidDataError as err:
1555 self.logger.debug("Error parsing library playlist: %s", err)
1556 # User-liked editorial playlists (not in users_playlists_list)
1557 liked_playlists = await self.client.get_liked_playlists()
1558 for playlist in liked_playlists:
1559 try:
1560 parsed = parse_playlist(self, playlist)
1561 if parsed.item_id not in seen_ids:
1562 yield parsed
1563 except InvalidDataError as err:
1564 self.logger.debug("Error parsing liked playlist: %s", err)
1565
1566 # Library edit methods
1567
1568 async def library_add(self, item: MediaItemType) -> bool:
1569 """
1570 Add item to library.
1571
1572 For tracks carrying a wave station context in the item_id (e.g. when
1573 the user adds a My Wave track to favourites during playback), also
1574 fires a rotor ``like`` feedback on the active session so the wave
1575 algorithm biases toward similar tracks immediately.
1576
1577 :param item: The media item to add.
1578 :return: True if successful.
1579 """
1580 prov_item_id = self._get_provider_item_id(item)
1581 if not prov_item_id:
1582 return False
1583 track_id, station_key = _parse_radio_item_id(prov_item_id)
1584
1585 if item.media_type == MediaType.TRACK:
1586 ok = await self.client.like_track(track_id)
1587 if ok and station_key:
1588 wave = self._wave_states.get(station_key)
1589 if wave and wave.session_id:
1590 await self._send_wave_feedback(wave, station_key, "like", track_id=track_id)
1591 return ok
1592 if item.media_type in (MediaType.ALBUM, MediaType.PODCAST, MediaType.AUDIOBOOK):
1593 return await self.client.like_album(prov_item_id)
1594 if item.media_type == MediaType.ARTIST:
1595 return await self.client.like_artist(prov_item_id)
1596 return False
1597
1598 async def library_remove(self, prov_item_id: str, media_type: MediaType) -> bool:
1599 """
1600 Remove item from library.
1601
1602 :param prov_item_id: The provider item ID (may be track_id@station_id for tracks).
1603 :param media_type: The media type.
1604 :return: True if successful.
1605 """
1606 track_id, _ = _parse_radio_item_id(prov_item_id)
1607 if media_type == MediaType.TRACK:
1608 return await self.client.unlike_track(track_id)
1609 if media_type in (MediaType.ALBUM, MediaType.PODCAST, MediaType.AUDIOBOOK):
1610 return await self.client.unlike_album(prov_item_id)
1611 if media_type == MediaType.ARTIST:
1612 return await self.client.unlike_artist(prov_item_id)
1613 return False
1614
1615 # Streaming
1616
1617 async def get_stream_details(
1618 self, item_id: str, media_type: MediaType = MediaType.TRACK
1619 ) -> StreamDetails:
1620 """
1621 Get stream details for a track, podcast episode, or audiobook.
1622
1623 A podcast episode is a track underneath the Yandex API, so it flows
1624 through the same per-track streaming path. An audiobook is an album
1625 with multiple tracks (chapters) â returned as a CUSTOM stream whose
1626 generator concatenates each chapter's bytes in order.
1627
1628 :param item_id: The track / episode ID (or track_id@station_id for My Wave),
1629 or the audiobook (album) ID when ``media_type`` is AUDIOBOOK.
1630 :param media_type: The media type.
1631 :return: StreamDetails for the item.
1632 """
1633 if media_type == MediaType.AUDIOBOOK:
1634 return await self._get_audiobook_stream_details(item_id)
1635 return await self.streaming.get_stream_details(item_id)
1636
1637 async def get_audio_stream(
1638 self, streamdetails: StreamDetails, seek_position: int = 0
1639 ) -> AsyncGenerator[bytes]:
1640 """
1641 Return the audio stream for the provider item.
1642
1643 For tracks and podcast episodes, streams via windowed Range requests
1644 (raw or AES-CTR encrypted). For audiobooks, iterates chapters: each
1645 chapter's bytes are streamed through the per-track path and concatenated.
1646
1647 :param streamdetails: Stream details with URL and optional decryption key.
1648 :param seek_position: Seek position in seconds (handled by provider for raw transport).
1649 :return: Async generator yielding audio chunks.
1650 """
1651 data = streamdetails.data if isinstance(streamdetails.data, dict) else None
1652 if streamdetails.media_type == MediaType.AUDIOBOOK and data and "chapter_ids" in data:
1653 async for chunk in self._stream_audiobook_chapters(data, seek_position):
1654 yield chunk
1655 return
1656 async for chunk in self.streaming.get_audio_stream(streamdetails, seek_position):
1657 yield chunk
1658
1659 async def get_rotor_station_tracks(
1660 self, station_id: str, queue: str | int | None = None
1661 ) -> tuple[list[Any], str | None]:
1662 """
1663 Fetch tracks from a rotor station using the session API.
1664
1665 Public surface â pinned by the ynison plugin
1666 (`YandexMusicProviderLike.get_rotor_station_tracks`). The
1667 ``(tracks, batch_id)`` return contract is kept for that caller even
1668 though batch_id is now a session-scoped identifier.
1669
1670 Routes to ``_fetch_rotor_session_batch`` so the wave session state
1671 (`session_id`, seen tracks, prefetch) is shared with our own Browse /
1672 on_played / on_streamed flows. ``queue`` is the most recently played
1673 track ID the external caller observed â we record it as the
1674 pagination cursor before calling through.
1675
1676 :param station_id: Rotor station ID (e.g. "user:onyourwave",
1677 "genre:rock", "mood:calm", "track:1234").
1678 :param queue: Last-played track ID for pagination. Ignored on the
1679 very first call (no session yet) but still recorded.
1680 :return: Tuple of (list of yandex tracks, batch_id or None).
1681 """
1682 wave = self._get_wave_state(station_id)
1683 # Cursor update + batch fetch run under the station's lock, matching
1684 # the discipline in browse / recommendations / prefetch. Without it,
1685 # ynison replenish racing with a concurrent MA browse could interleave
1686 # last_track_id writes and leave session_id / batch_id out of sync.
1687 async with wave.lock:
1688 if queue is not None:
1689 wave.last_track_id = str(queue)
1690 return await self._fetch_rotor_session_batch(wave, station_id)
1691
1692 def get_quality(self) -> str:
1693 """Return the configured audio quality tier (e.g. 'balanced', 'superb')."""
1694 quality = str(self.config.get_value(CONF_QUALITY) or QUALITY_BALANCED).strip().lower()
1695 if quality == "lossless":
1696 quality = QUALITY_SUPERB
1697 return quality
1698
1699 async def resolve_image(self, path: str) -> str | bytes:
1700 """
1701 Resolve wave cover image with background color fill for transparent PNGs.
1702
1703 If the image URL has an associated background color (stored in _wave_bg_colors),
1704 downloads the PNG from Yandex CDN and composites it on a solid color background
1705 using Pillow, returning JPEG bytes. Falls back to the original URL on any error.
1706
1707 :param path: Image URL (may include #rrggbb fragment used as cache key).
1708 :return: Composited JPEG bytes, or original path string as fallback.
1709 """
1710 bg_color = self._wave_bg_colors.get(path)
1711 if not bg_color:
1712 return path
1713
1714 # Strip the #color fragment before fetching the actual image
1715 fetch_url = path.split("#", maxsplit=1)[0] if "#" in path else path
1716 try:
1717 async with self.mass.http_session.get(fetch_url) as resp:
1718 resp.raise_for_status()
1719 raw = await resp.read()
1720 except Exception as err:
1721 self.logger.debug("Failed to fetch wave cover %s: %s", fetch_url, err)
1722 return fetch_url
1723
1724 def _composite() -> bytes:
1725 bg_clean = bg_color.lstrip("#")
1726 try:
1727 r = int(bg_clean[0:2], 16)
1728 g = int(bg_clean[2:4], 16)
1729 b = int(bg_clean[4:6], 16)
1730 except ValueError, IndexError:
1731 return raw
1732 fg = PilImage.open(BytesIO(raw)).convert("RGBA")
1733 bg = PilImage.new("RGBA", fg.size, (r, g, b, 255))
1734 bg.paste(fg, mask=fg)
1735 out = BytesIO()
1736 bg.convert("RGB").save(out, "JPEG", quality=92)
1737 return out.getvalue()
1738
1739 try:
1740 return await asyncio.to_thread(_composite)
1741 except Exception as err:
1742 self.logger.debug("Wave cover composite failed for %s: %s", fetch_url, err)
1743 return fetch_url
1744
1745 async def on_played(
1746 self,
1747 media_type: MediaType,
1748 prov_item_id: str,
1749 fully_played: bool,
1750 position: int,
1751 media_item: MediaItemType,
1752 is_playing: bool = False,
1753 ) -> None:
1754 """
1755 Report periodic playback updates.
1756
1757 - Audiobooks: persist chapter progress via play_audio so Yandex's
1758 own clients resume at the right point.
1759 - Wave tracks: send rotor ``trackStarted`` while actively playing and
1760 kick off a background prefetch so DSTM refill serves wave-curated
1761 tracks with no extra round-trip. DSTM itself is the user's toggle â
1762 the provider does not flip it.
1763
1764 Generic track history reporting is not attempted here â the only
1765 known channel Yandex writes into ``/handlers/music-history`` is a
1766 long-lived Ynison WebSocket session, which lives in the sibling
1767 yandex_ynison plugin. Regular tracks played through MA are therefore
1768 invisible to Listening History unless that plugin is also active.
1769 """
1770 if media_type == MediaType.AUDIOBOOK:
1771 await self._report_audiobook_progress(prov_item_id, position)
1772 return
1773 if media_type != MediaType.TRACK:
1774 return
1775 _, station_id = _parse_radio_item_id(prov_item_id)
1776 if station_id and is_playing:
1777 track_id, _ = _parse_radio_item_id(prov_item_id)
1778 wave = self._wave_states.get(station_id) or self._get_wave_state(station_id)
1779 await self._send_wave_feedback(wave, station_id, "trackStarted", track_id=track_id)
1780 self.mass.create_task(self._prefetch_rotor_session(station_id))
1781
1782 async def on_streamed(self, streamdetails: StreamDetails) -> None:
1783 """
1784 Report stream completion to Yandex.
1785
1786 - Audiobooks: a final ``play_audio`` with the absolute stream
1787 position so the last listening point is preserved across Yandex
1788 clients. Cleans up session state even when ``data`` was stripped.
1789 - Wave tracks (composite item_id carries a station suffix): a rotor
1790 ``trackFinished`` or ``skip`` event with the actual seconds streamed
1791 so Yandex can improve recommendations.
1792 """
1793 data = streamdetails.data if isinstance(streamdetails.data, dict) else None
1794 if streamdetails.media_type == MediaType.AUDIOBOOK:
1795 await self._report_audiobook_final(streamdetails, data or {})
1796 return
1797 if streamdetails.media_type != MediaType.TRACK:
1798 return
1799 track_id, station_id = _parse_radio_item_id(streamdetails.item_id)
1800 if not station_id:
1801 return
1802 seconds = int(streamdetails.seconds_streamed or 0)
1803 duration = int(streamdetails.duration or 0)
1804 feedback_type = "trackFinished" if duration and seconds >= max(0, duration - 10) else "skip"
1805 wave = self._wave_states.get(station_id) or self._get_wave_state(station_id)
1806 await self._send_wave_feedback(
1807 wave, station_id, feedback_type, track_id=track_id, total_played_seconds=seconds
1808 )
1809
1810 def _media_label(self, group: str, key: str, fallback: str) -> tuple[str, str | None]:
1811 """
1812 Resolve a media label to its English ``name`` and ``translation_key``.
1813
1814 The English source string lives in the provider's ``strings.json`` (the single source
1815 of truth) and is localized for the connection locale at serialization via the returned
1816 key. An unauthored key â e.g. a tag discovered from Yandex's landing API â returns
1817 ``(fallback, None)`` so its already-localized name is kept verbatim.
1818
1819 :param group: Media translation group (``folder``, ``recommendations`` or ``playlist``).
1820 :param key: Authoring key within the group; also the item's ``translation_key``.
1821 :param fallback: English name to use when no string is authored for *key*.
1822 """
1823 authored = self.mass.translations.get_translation(
1824 f"provider.{self.domain}.media.{group}.{key}.name"
1825 )
1826 if authored is None:
1827 return fallback, None
1828 return authored, key
1829
1830 async def _reauth_via_refresh_token(
1831 self, x_token: str, refresh_token: str, base_url: str, original_err: Exception
1832 ) -> None:
1833 """
1834 Silently re-issue full credentials when x_token refresh fails.
1835
1836 Device-flow accounts have a refresh_token that can mint a new
1837 x_token + refresh_token + music_token without any user interaction.
1838 Persists the rotated triple and connects the client. Any failure
1839 here is terminal â clears all credentials and forces re-auth.
1840 """
1841 try:
1842 new_creds = await refresh_credentials_via_passport(
1843 SecretStr(x_token), SecretStr(refresh_token)
1844 )
1845 except ResourceTemporarilyUnavailable as err2:
1846 # Transient Passport failure â keep creds, let MA retry later
1847 self.logger.warning(
1848 "Credential refresh temporarily unavailable: %s", type(err2).__name__
1849 )
1850 raise ProviderUnavailableError(
1851 "Unable to refresh credentials right now. Please try again later."
1852 ) from err2
1853 except LoginFailed as err2:
1854 self.logger.warning("Session and refresh tokens are both expired")
1855 self._update_setup_data(CONF_TOKEN, None)
1856 self._update_setup_data(CONF_X_TOKEN, None)
1857 self._update_setup_data(CONF_REFRESH_TOKEN, None)
1858 raise LoginFailed("Session expired. Please re-authenticate.") from err2
1859
1860 new_music_token = new_creds.music_token
1861 new_refresh_token = new_creds.refresh_token
1862 if new_music_token is None or new_refresh_token is None:
1863 self._update_setup_data(CONF_TOKEN, None)
1864 self._update_setup_data(CONF_X_TOKEN, None)
1865 self._update_setup_data(CONF_REFRESH_TOKEN, None)
1866 raise LoginFailed(
1867 "Credential refresh returned an incomplete response."
1868 ) from original_err
1869
1870 self._update_setup_data(CONF_TOKEN, new_music_token.get_secret())
1871 self._update_setup_data(CONF_X_TOKEN, new_creds.x_token.get_secret())
1872 self._update_setup_data(CONF_REFRESH_TOKEN, new_refresh_token.get_secret())
1873 restrictive = bool(self.config.get_value(CONF_RESTRICTIVE_RATE_LIMITS, False))
1874 self._client = YandexMusicClient(
1875 new_music_token, base_url=base_url, restrictive_rate_limits=restrictive
1876 )
1877 await self._client.connect()
1878 self.logger.info("Re-issued credentials silently from refresh token")
1879
1880 async def _browse_my_wave(
1881 self, path: str, sub_subpath: str | None
1882 ) -> list[Track | BrowseFolder]:
1883 """
1884 Browse My Wave tracks (must be called under the My Wave state lock).
1885
1886 :param path: Full browse path.
1887 :param sub_subpath: Sub-path part ('next' for load more, or track_id cursor).
1888 :return: List of Track and optional BrowseFolder for "Load more".
1889 """
1890 wave = self._get_wave_state(ROTOR_STATION_MY_WAVE)
1891 max_tracks_config = int(
1892 self.config.get_value(CONF_MY_WAVE_MAX_TRACKS) or 150 # type: ignore[arg-type]
1893 )
1894 batch_size_config = MY_WAVE_BATCH_SIZE
1895
1896 # Effective limit on tracks to collect for this call:
1897 # initial browse is capped to BROWSE_INITIAL_TRACKS to avoid marking
1898 # extra tracks as "seen" that are never shown to the user.
1899 effective_limit = min(
1900 BROWSE_INITIAL_TRACKS if sub_subpath != "next" else max_tracks_config,
1901 max_tracks_config,
1902 )
1903
1904 # Root my_wave: fetch up to batch_size_config batches so Play adds more tracks.
1905 # "Load more" always uses single next batch.
1906 max_batches = batch_size_config if sub_subpath != "next" else 1
1907
1908 # Reset seen tracks on fresh browse (not "load more")
1909 if sub_subpath != "next":
1910 wave.seen_track_ids = set()
1911
1912 queue: str | int | None = None
1913 if sub_subpath == "next":
1914 queue = wave.last_track_id
1915 elif sub_subpath:
1916 queue = sub_subpath
1917
1918 all_tracks: list[Track | BrowseFolder] = []
1919 last_batch_id: str | None = None
1920 first_track_id_this_batch: str | None = None
1921 total_track_count = 0
1922
1923 for _ in range(max_batches):
1924 if total_track_count >= effective_limit:
1925 break
1926
1927 # On a fresh browse (non-"next"), honour any sub_subpath cursor override
1928 # by seeding wave.last_track_id so the helper picks it up.
1929 if queue is not None:
1930 wave.last_track_id = str(queue)
1931 yandex_tracks, batch_id = await self._fetch_rotor_session_batch(
1932 wave, ROTOR_STATION_MY_WAVE
1933 )
1934 if batch_id:
1935 last_batch_id = batch_id
1936 if not wave.radio_started_sent and yandex_tracks:
1937 sent = await self._send_wave_feedback(wave, ROTOR_STATION_MY_WAVE, "radioStarted")
1938 if sent:
1939 wave.radio_started_sent = True
1940 first_track_id_this_batch = None
1941 for yt in yandex_tracks:
1942 if total_track_count >= effective_limit:
1943 break
1944
1945 track = self._parse_my_wave_track(yt, wave.seen_track_ids)
1946 if track is None:
1947 continue
1948 all_tracks.append(track)
1949 total_track_count += 1
1950
1951 track_id = track.item_id.split(RADIO_TRACK_ID_SEP, 1)[0]
1952 if first_track_id_this_batch is None:
1953 first_track_id_this_batch = track_id
1954
1955 if first_track_id_this_batch is not None:
1956 wave.last_track_id = first_track_id_this_batch
1957 if (
1958 first_track_id_this_batch is None
1959 or not batch_id
1960 or not yandex_tracks
1961 or total_track_count >= effective_limit
1962 ):
1963 break
1964 queue = first_track_id_this_batch
1965
1966 # Only show "Load more" if we haven't reached the limit and there's more data
1967 if last_batch_id and total_track_count < max_tracks_config:
1968 all_tracks.append(
1969 BrowseFolder(
1970 item_id="next",
1971 provider=self.instance_id,
1972 path=f"{path.rstrip('/')}/next",
1973 name="Load more",
1974 translation_key="load_more",
1975 is_playable=False,
1976 )
1977 )
1978 return all_tracks
1979
1980 def _get_user_wave_presets(self) -> list[dict[str, str]]:
1981 """
1982 Decode user-defined wave presets from the hidden JSON config key.
1983
1984 Thin wrapper around :func:`presets.parse_stored_presets` so browse
1985 code and settings actions use the exact same parsing â avoids schema
1986 drift when preset fields are added or renamed.
1987 """
1988 return parse_stored_presets(self.config.get_value(CONF_WAVE_PRESETS_DATA))
1989
1990 def _browse_user_presets_list(
1991 self, path: str, presets: list[dict[str, str]]
1992 ) -> list[BrowseFolder]:
1993 """
1994 Return one playable BrowseFolder per configured user preset.
1995
1996 ``path`` is nested (``my_wave_presets/<idx>``) so MA's back-nav â
1997 which strips the last ``/``-segment â returns the user to the
1998 listing instead of the provider root. ``item_id`` uses the
1999 underscore form (``my_wave_presets_<idx>``) because MA rebuilds a
2000 playable folder's path from its item_id at play time. The browse
2001 dispatcher accepts both forms.
2002
2003 :param path: Current browse path.
2004 :param presets: Sanitized presets from ``_get_user_wave_presets``.
2005 :return: List of playable BrowseFolder entries.
2006 """
2007 base = path if path.endswith("/") else f"{path}/"
2008 folders: list[BrowseFolder] = []
2009 for idx, preset in enumerate(presets):
2010 folders.append(
2011 BrowseFolder(
2012 item_id=f"{MY_WAVE_PRESETS_FOLDER_ID}_{idx}",
2013 provider=self.instance_id,
2014 path=f"{base}{idx}",
2015 name=preset.get("name", f"Preset {idx + 1}"),
2016 is_playable=True,
2017 )
2018 )
2019 return folders
2020
2021 def _browse_my_wave_modes_list(self, path: str) -> list[BrowseFolder]:
2022 """
2023 Return the 11 wave-mode entries as playable browse folders.
2024
2025 Same dual-form contract as user presets: nested ``path`` keeps
2026 back-navigation intact, underscore ``item_id`` survives MA's
2027 play-time reconstruction.
2028
2029 :param path: Browse path the user navigated into.
2030 :return: Ordered list of BrowseFolder entries, one per preset.
2031 """
2032 base = path if path.endswith("/") else f"{path}/"
2033 folders: list[BrowseFolder] = []
2034 for preset in WAVE_MODE_ORDER:
2035 name, translation_key = self._media_label(
2036 "folder", f"wave_mode_{preset}", preset.replace("_", " ").title()
2037 )
2038 folders.append(
2039 BrowseFolder(
2040 item_id=f"{MY_WAVE_MODES_FOLDER_ID}_{preset}",
2041 provider=self.instance_id,
2042 path=f"{base}{preset}",
2043 name=name,
2044 translation_key=translation_key,
2045 is_playable=True,
2046 )
2047 )
2048 return folders
2049
2050 async def _browse_my_wave_mode(
2051 self, path: str, station_key: str, load_more: bool
2052 ) -> list[Track | BrowseFolder]:
2053 """
2054 Fetch a batch of tracks for a specific wave-mode preset.
2055
2056 Reuses the session-API machinery: tracks live in
2057 ``_wave_states[station_key]`` where station_key is
2058 ``user:onyourwave#{preset}``. Tracks carry composite item_ids that
2059 route feedback back to this state.
2060
2061 :param path: Full browse path to this preset.
2062 :param station_key: Station key with a ``#preset`` suffix.
2063 :param load_more: True when called for ``.../next`` pagination.
2064 :return: Tracks + optional "Load more" folder.
2065 """
2066 wave = self._get_wave_state(station_key)
2067 max_tracks_config = int(
2068 self.config.get_value(CONF_MY_WAVE_MAX_TRACKS) or 150 # type: ignore[arg-type]
2069 )
2070 batch_size_config = MY_WAVE_BATCH_SIZE
2071 effective_limit = min(
2072 BROWSE_INITIAL_TRACKS if not load_more else max_tracks_config,
2073 max_tracks_config,
2074 )
2075 max_batches = batch_size_config if not load_more else 1
2076
2077 if not load_more:
2078 wave.seen_track_ids = set()
2079
2080 all_tracks: list[Track | BrowseFolder] = []
2081 last_batch_id: str | None = None
2082 total_track_count = 0
2083
2084 for _ in range(max_batches):
2085 if total_track_count >= effective_limit:
2086 break
2087 yandex_tracks, batch_id = await self._fetch_rotor_session_batch(wave, station_key)
2088 if batch_id:
2089 last_batch_id = batch_id
2090 if not wave.radio_started_sent and yandex_tracks:
2091 sent = await self._send_wave_feedback(wave, station_key, "radioStarted")
2092 if sent:
2093 wave.radio_started_sent = True
2094 first_track_id_this_batch: str | None = None
2095 for yt in yandex_tracks:
2096 if total_track_count >= effective_limit:
2097 break
2098 track = self._parse_my_wave_track(yt, wave.seen_track_ids, station_key=station_key)
2099 if track is None:
2100 continue
2101 all_tracks.append(track)
2102 total_track_count += 1
2103 track_id = track.item_id.split(RADIO_TRACK_ID_SEP, 1)[0]
2104 if first_track_id_this_batch is None:
2105 first_track_id_this_batch = track_id
2106 if first_track_id_this_batch is not None:
2107 wave.last_track_id = first_track_id_this_batch
2108 if (
2109 first_track_id_this_batch is None
2110 or not batch_id
2111 or not yandex_tracks
2112 or total_track_count >= effective_limit
2113 ):
2114 break
2115
2116 if last_batch_id and total_track_count < max_tracks_config:
2117 all_tracks.append(
2118 BrowseFolder(
2119 item_id="next",
2120 provider=self.instance_id,
2121 path=f"{path.rstrip('/')}/next",
2122 name="Load more",
2123 translation_key="load_more",
2124 is_playable=False,
2125 )
2126 )
2127 return all_tracks
2128
2129 def _parse_my_wave_track(
2130 self,
2131 yt: Any,
2132 seen_ids: set[str],
2133 *,
2134 station_key: str = ROTOR_STATION_MY_WAVE,
2135 ) -> Track | None:
2136 """
2137 Parse a Yandex track into a My Wave Track with composite item_id.
2138
2139 Extracts the track_id, checks for duplicates in the seen_ids set,
2140 sets composite item_id (track_id@station_key) and updates
2141 provider_mappings. `station_key` is the key in `_wave_states` under
2142 which the matching session lives; for preset modes it carries a
2143 `#preset` suffix so `on_played`/`on_streamed` find the right session.
2144
2145 Callers using shared state must hold the My Wave state lock.
2146
2147 :param yt: Yandex track object from rotor station response.
2148 :param seen_ids: Set of already-seen track IDs to check and update.
2149 :param station_key: Station key to embed in the composite item_id.
2150 Defaults to the plain My Wave station.
2151 :return: Parsed Track with composite item_id, or None if duplicate/invalid.
2152 """
2153 try:
2154 t = parse_track(self, yt)
2155 except InvalidDataError as err:
2156 self.logger.debug("Error parsing My Wave track: %s", err)
2157 return None
2158
2159 track_id = str(yt.id) if hasattr(yt, "id") and yt.id else getattr(yt, "track_id", None)
2160 if not track_id:
2161 return t
2162
2163 if track_id in seen_ids:
2164 self.logger.debug("Skipping duplicate My Wave track: %s", track_id)
2165 return None
2166
2167 seen_ids.add(track_id)
2168 t.item_id = f"{track_id}{RADIO_TRACK_ID_SEP}{station_key}"
2169 for pm in t.provider_mappings:
2170 if pm.provider_instance == self.instance_id:
2171 pm.item_id = t.item_id
2172 break
2173 return t
2174
2175 @use_cache(3600, allow_expired_cache=True)
2176 async def _get_valid_tags_for_category(self, category: str) -> list[str]:
2177 """
2178 Return tags for a category by combining hardcoded + landing-discovered.
2179
2180 Trusts the hardcoded ``TAG_CATEGORY_*`` lists (evergreen Yandex
2181 categories) and the landing API output (Yandex returns landing tags
2182 only when they have playlists). No per-tag runtime validation: that
2183 machinery was a parallel ``asyncio.gather`` over
2184 ``get_tag_playlists`` for every tag and tripped Yandex's edge
2185 per-endpoint concurrency limit on first browse â captcha within
2186 ~460ms of the burst. If a tag turns out to be empty at click time,
2187 ``_get_tag_playlists_as_browse`` already renders an empty folder.
2188
2189 :param category: Category name ('mood', 'activity', 'era', 'genres').
2190 :return: List of tag slugs (hardcoded order preserved, landing tags appended).
2191 """
2192 category_lists: dict[str, list[str]] = {
2193 "mood": list(TAG_CATEGORY_MOOD),
2194 "activity": list(TAG_CATEGORY_ACTIVITY),
2195 "era": list(TAG_CATEGORY_ERA),
2196 "genres": list(TAG_CATEGORY_GENRES),
2197 }
2198 tags = category_lists.get(category, [])
2199 try:
2200 landing_tags = await self.client.get_landing_tags()
2201 for slug, _title in landing_tags:
2202 cat = TAG_SLUG_CATEGORY.get(slug, "mood")
2203 if cat == category and slug not in tags:
2204 tags.append(slug)
2205 except Exception as err:
2206 self.logger.debug("Landing tag discovery failed: %s", err)
2207 return tags
2208
2209 @use_cache(3600, allow_expired_cache=True)
2210 async def _get_discovered_tags(self, locale: str) -> list[tuple[str, str, str | None]]:
2211 """
2212 Return all browse-able tags: hardcoded (non-seasonal) + landing-discovered.
2213
2214 Same rationale as :meth:`_get_valid_tags_for_category` â runtime
2215 validation removed to avoid the per-endpoint concurrency burst that
2216 triggered Yandex captcha. The locale parameter is part of the cache
2217 key so locale changes invalidate the cached landing titles.
2218
2219 :param locale: Current metadata locale (used as part of cache key).
2220 :return: List of (slug, English name, translation_key) tuples in
2221 hardcoded-then-discovered order. Landing-discovered tags carry their
2222 (already localized) API title and no translation_key.
2223 """
2224 all_tags: dict[str, tuple[str, str | None]] = {}
2225 for slug, cat in TAG_SLUG_CATEGORY.items():
2226 if cat != "seasonal":
2227 all_tags[slug] = self._media_label("folder", _media_label_key(slug), slug.title())
2228 try:
2229 landing_tags = await self.client.get_landing_tags()
2230 for slug, title in landing_tags:
2231 if slug not in all_tags:
2232 all_tags[slug] = (title, None)
2233 except Exception as err:
2234 self.logger.debug("Failed to discover tags from landing API: %s", err)
2235 return [(slug, name, translation_key) for slug, (name, translation_key) in all_tags.items()]
2236
2237 async def _get_discovered_tag_slugs(self) -> set[str]:
2238 """
2239 Get set of all valid tag slugs (cached).
2240
2241 :return: Set of tag slug strings that have playlists.
2242 """
2243 discovered = await self._get_discovered_tags(self.mass.metadata.locale or "en_US")
2244 return {slug for slug, _name, _key in discovered}
2245
2246 async def _browse_for_you(
2247 self, path: str, path_parts: list[str]
2248 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2249 """
2250 Browse «For You» folder â shows Picks and Mixes sub-folders.
2251
2252 :param path: Full browse path.
2253 :param path_parts: Split path parts after ://.
2254 :return: List of sub-folders (Picks, Mixes).
2255 """
2256 # Strip the for_you segment to build child paths that route to picks/mixes
2257 # Path format: ...//for_you â child paths should be ...//picks, ...//mixes
2258 # We build base from the root (before for_you) by dropping the last segment.
2259 base_parts = path.split("//", 1)
2260 root_base = (base_parts[0] + "//") if len(base_parts) > 1 else path.rstrip("/") + "/"
2261
2262 if len(path_parts) == 1:
2263 return [
2264 BrowseFolder(
2265 item_id="picks",
2266 provider=self.instance_id,
2267 path=f"{root_base}picks",
2268 name="Picks",
2269 translation_key="picks",
2270 is_playable=False,
2271 ),
2272 BrowseFolder(
2273 item_id="mixes",
2274 provider=self.instance_id,
2275 path=f"{root_base}mixes",
2276 name="Mixes",
2277 translation_key="mixes",
2278 is_playable=False,
2279 ),
2280 ]
2281 # Deeper path: delegate to picks or mixes handler via canonical paths
2282 return await super().browse(path)
2283
2284 async def _browse_collection(
2285 self, path: str
2286 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2287 """
2288 Browse «Collection» folder â shows library sub-folders (tracks/artists/albums/playlists).
2289
2290 Child ``path`` is nested (``â¦/collection/tracks``) so MA's "back"
2291 button lands on this listing instead of the provider root. The
2292 dispatcher then strips the ``collection/`` prefix and hands off to
2293 core's default library handler.
2294
2295 :param path: Full browse path.
2296 :return: List of library sub-folders.
2297 """
2298 base = path if path.endswith("/") else f"{path}/"
2299
2300 folders: list[BrowseFolder] = []
2301 for feature, sub_id, label_key, is_playable in _COLLECTION_SUBFOLDERS:
2302 if feature not in self.supported_features:
2303 continue
2304 name, translation_key = self._media_label(
2305 "folder", label_key, label_key.replace("_", " ").title()
2306 )
2307 folders.append(
2308 BrowseFolder(
2309 item_id=sub_id,
2310 provider=self.instance_id,
2311 path=f"{base}{sub_id}",
2312 name=name,
2313 translation_key=translation_key,
2314 is_playable=is_playable,
2315 )
2316 )
2317 return folders
2318
2319 async def _browse_pins(self) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2320 """
2321 Browse user's pinned items (artists/albums/playlists from Yandex Pins).
2322
2323 Resolves each pin to its full media item via existing single-item lookups.
2324 Wave pins are skipped â MA has no native concept for them.
2325
2326 :return: List of resolved media items.
2327 """
2328 pins_list = await self.client.get_pins()
2329 pins = getattr(pins_list, "pins", None) if pins_list else None
2330 if not pins:
2331 return []
2332
2333 items: list[MediaItemType] = []
2334 for pin in pins:
2335 pin_type = getattr(pin, "type", None)
2336 data = getattr(pin, "data", None)
2337 if data is None:
2338 continue
2339 try:
2340 if pin_type == "artist_item" and getattr(data, "id", None) is not None:
2341 items.append(await self.get_artist(str(data.id)))
2342 elif pin_type == "album_item" and getattr(data, "id", None) is not None:
2343 items.append(await self.get_album(str(data.id)))
2344 elif pin_type == "playlist_item":
2345 uid = getattr(data, "uid", None)
2346 kind = getattr(data, "kind", None)
2347 if uid is not None and kind is not None:
2348 items.append(await self.get_playlist(f"{uid}:{kind}"))
2349 except (MediaNotFoundError, InvalidDataError) as err:
2350 self.logger.debug("Skipping pin %s: %s", pin_type, err)
2351 return items
2352
2353 async def _browse_history(self) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2354 """
2355 Browse user's recent listening history (flattened across days).
2356
2357 Collects ``track_id`` values from each history entry's ``item_id``
2358 sub-object (``full_model`` is not populated by the current API
2359 response â MarshalX exposes the IDs separately), dedupes, and
2360 batch-resolves them via ``get_tracks`` so the returned Track objects
2361 carry full artist/album/cover metadata.
2362
2363 Entries without a resolvable ``track_id`` (e.g. album-only context
2364 rows) are skipped silently. Order is preserved â most recent first â
2365 by collecting unique IDs in response order into ``ordered_ids``,
2366 then rebuilding the final list by iterating ``ordered_ids`` and
2367 looking up each batch-fetched track in an idâtrack map.
2368
2369 :return: List of recently played Track items.
2370 """
2371 history = await self.client.get_music_history()
2372 tabs = getattr(history, "history_tabs", None) if history else None
2373 if not tabs:
2374 return []
2375
2376 seen_track_ids: set[str] = set()
2377 ordered_ids: list[str] = []
2378 for tab in tabs:
2379 for group in getattr(tab, "items", None) or []:
2380 for hist_item in getattr(group, "tracks", None) or []:
2381 if getattr(hist_item, "type", None) != "track":
2382 continue
2383 item_id_obj = getattr(getattr(hist_item, "data", None), "item_id", None)
2384 track_key: str | None = None
2385 if isinstance(item_id_obj, dict):
2386 track_key = item_id_obj.get("track_id") or item_id_obj.get("id")
2387 else:
2388 track_key = getattr(item_id_obj, "track_id", None) or getattr(
2389 item_id_obj, "id", None
2390 )
2391 if not track_key:
2392 continue
2393 track_key = str(track_key)
2394 if track_key in seen_track_ids:
2395 continue
2396 seen_track_ids.add(track_key)
2397 ordered_ids.append(track_key)
2398
2399 if not ordered_ids:
2400 return []
2401
2402 try:
2403 fetched = await self.client.get_tracks(ordered_ids)
2404 except ResourceTemporarilyUnavailable as err:
2405 self.logger.warning("Failed to hydrate history tracks: %s", err)
2406 return []
2407
2408 by_id = {str(t.id): t for t in fetched if getattr(t, "id", None) is not None}
2409 tracks: list[Track] = []
2410 for tid in ordered_ids:
2411 yt = by_id.get(tid)
2412 if yt is None:
2413 continue
2414 try:
2415 tracks.append(parse_track(self, yt))
2416 except InvalidDataError as err:
2417 self.logger.debug("Skipping history track %s: %s", tid, err)
2418 return tracks
2419
2420 async def _browse_picks(
2421 self, path: str, path_parts: list[str]
2422 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2423 """
2424 Browse picks folder using hardcoded tags validated against the API.
2425
2426 Tags are sourced from hardcoded category lists and landing API discovery,
2427 then validated via client.tags() to ensure they have playlists.
2428 Only categories with at least one valid tag are shown.
2429
2430 :param path: Full browse path.
2431 :param path_parts: Split path parts after ://.
2432 :return: List of folders or playlists.
2433 """
2434 base = path.rstrip("/") + "/"
2435
2436 # Get validated tags
2437 discovered = await self._get_discovered_tags(self.mass.metadata.locale or "en_US")
2438
2439 # Categorize valid tags, carrying each tag's (slug, English name, translation_key)
2440 categorized: dict[str, list[tuple[str, str, str | None]]] = {}
2441 for slug, name, translation_key in discovered:
2442 cat = TAG_SLUG_CATEGORY.get(slug, "mood")
2443 # Skip seasonal tags â they belong in mixes, not picks
2444 if cat == "seasonal":
2445 continue
2446 categorized.setdefault(cat, []).append((slug, name, translation_key))
2447
2448 # Sort tags within each category by preferred order
2449 for cat, cat_tags in categorized.items():
2450 order = TAG_CATEGORY_ORDER.get(cat, [])
2451 order_map = {s: i for i, s in enumerate(order)}
2452 cat_tags.sort(key=lambda t: order_map.get(t[0], len(order)))
2453
2454 # picks/ - show category folders (only those with valid tags)
2455 if len(path_parts) == 1:
2456 category_display_order = ["mood", "activity", "era", "genres"]
2457 folders: list[BrowseFolder] = []
2458 for cat in category_display_order:
2459 if cat in categorized:
2460 name, translation_key = self._media_label("folder", cat, cat.title())
2461 folders.append(
2462 BrowseFolder(
2463 item_id=cat,
2464 provider=self.instance_id,
2465 path=f"{base}{cat}",
2466 name=name,
2467 translation_key=translation_key,
2468 is_playable=False,
2469 )
2470 )
2471 # Show any extra categories not in the standard order
2472 for cat in categorized:
2473 if cat not in category_display_order:
2474 name, translation_key = self._media_label("folder", cat, cat.title())
2475 folders.append(
2476 BrowseFolder(
2477 item_id=cat,
2478 provider=self.instance_id,
2479 path=f"{base}{cat}",
2480 name=name,
2481 translation_key=translation_key,
2482 is_playable=False,
2483 )
2484 )
2485 return folders
2486
2487 category: str | None = path_parts[1] if len(path_parts) > 1 else None
2488 tag: str | None = path_parts[2] if len(path_parts) > 2 else None
2489
2490 self.logger.debug(
2491 "Browse picks: path=%s, category=%s, tag=%s",
2492 path,
2493 category,
2494 tag,
2495 )
2496
2497 # picks/category/ - show valid tag folders for this category
2498 if category and not tag:
2499 category_tags = categorized.get(category, [])
2500 folders = []
2501 for slug, name, translation_key in category_tags:
2502 folders.append(
2503 BrowseFolder(
2504 item_id=slug,
2505 provider=self.instance_id,
2506 path=f"{base}{slug}",
2507 name=name,
2508 translation_key=translation_key,
2509 is_playable=False,
2510 )
2511 )
2512 self.logger.debug("Returning %d tag folders for category %s", len(folders), category)
2513 return folders
2514
2515 # picks/category/tag - show playlists for the tag
2516 if tag:
2517 discovered_slugs = {slug for slug, _name, _key in discovered}
2518 if tag in discovered_slugs:
2519 self.logger.debug("Fetching playlists for tag: %s", tag)
2520 return await self._get_tag_playlists_as_browse(tag)
2521
2522 self.logger.debug("No match found, returning empty list")
2523 return []
2524
2525 async def _browse_mixes(
2526 self, path: str, path_parts: list[str]
2527 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2528 """
2529 Browse mixes folder (seasonal collections) using hardcoded tags.
2530
2531 Renders every seasonal tag from ``TAG_MIXES`` unconditionally. The
2532 old per-tag validation fired a ``Semaphore(5)+gather`` of
2533 ``get_tag_playlists`` calls and tripped Yandex's per-endpoint
2534 concurrency limit on first browse. If a season ends up empty at
2535 click time, ``_get_tag_playlists_as_browse`` already returns an
2536 empty folder.
2537
2538 :param path: Full browse path.
2539 :param path_parts: Split path parts after ://.
2540 :return: List of folders or playlists.
2541 """
2542 base = path.rstrip("/") + "/"
2543
2544 # mixes/ - show seasonal folders
2545 if len(path_parts) == 1:
2546 folders: list[BrowseFolder] = []
2547 for t in TAG_MIXES:
2548 name, translation_key = self._media_label("folder", t, t.title())
2549 folders.append(
2550 BrowseFolder(
2551 item_id=t,
2552 provider=self.instance_id,
2553 path=f"{base}{t}",
2554 name=name,
2555 translation_key=translation_key,
2556 is_playable=False,
2557 )
2558 )
2559 return folders
2560
2561 # mixes/tag - show playlists for the tag
2562 tag = path_parts[1] if len(path_parts) > 1 else None
2563 if tag and tag in TAG_MIXES:
2564 return await self._get_tag_playlists_as_browse(tag)
2565
2566 return []
2567
2568 def _get_wave_state(self, station_id: str) -> _WaveState:
2569 """
2570 Get or create per-station wave state.
2571
2572 :param station_id: Rotor station ID (e.g. 'genre:rock', 'mood:chill').
2573 :return: _WaveState instance for this station.
2574 """
2575 return self._wave_states.setdefault(station_id, _WaveState())
2576
2577 async def _send_wave_feedback(
2578 self,
2579 wave: _WaveState,
2580 station_id: str,
2581 event_type: str,
2582 *,
2583 track_id: str | None = None,
2584 total_played_seconds: int | None = None,
2585 ) -> bool:
2586 """
2587 Route rotor feedback to the session endpoint.
2588
2589 Requires an active ``wave.session_id`` â rotor feedback is only
2590 meaningful inside the session it originated from. The legacy
2591 stations-based endpoint (``/rotor/station/{id}/feedback``) is no
2592 longer reachable (returns 404 "not-found"), so when there's no
2593 session we skip silently rather than spamming the log.
2594
2595 This happens when the track's composite item_id was parsed in a
2596 previous provider run (e.g. loaded from MA's library cache) and
2597 the corresponding session_id is not in memory any more. History
2598 reporting via ``play_audio`` still works in that case â only the
2599 rotor recommendation signal is lost.
2600
2601 :param wave: Station state carrying session_id + batch_id.
2602 :param station_id: Rotor station ID (used only for logging here).
2603 :param event_type: Rotor event type (radioStarted, trackStarted, â¦).
2604 :param track_id: Yandex track ID the event refers to.
2605 :param total_played_seconds: Seconds played (trackFinished / skip only).
2606 :return: True if the feedback POST succeeded, False when skipped.
2607 """
2608 if not wave.session_id:
2609 self.logger.debug(
2610 "Skipping rotor feedback %s for %s: no active session",
2611 event_type,
2612 station_id,
2613 )
2614 return False
2615 return await self.client.rotor_session_feedback(
2616 wave.session_id,
2617 event_type,
2618 track_id=track_id,
2619 total_played_seconds=total_played_seconds,
2620 batch_id=wave.batch_id,
2621 )
2622
2623 async def _prefetch_rotor_session(self, station_key: str) -> None:
2624 """
2625 Fire-and-forget: fetch the next batch for an active wave session.
2626
2627 Called from ``on_played`` while a wave track starts playing, so by the
2628 time Music Assistant's DSTM asks for more via ``get_similar_tracks``,
2629 we already have Yandex-curated wave tracks sitting in
2630 ``wave.prefetched`` ready to serve (no extra round-trip).
2631
2632 No-op when the station has no active session yet (prefetch cannot
2633 safely create one â that requires holding the lock across the
2634 network call and would stall readers), or when the buffer already
2635 has items (avoids burning rate limit).
2636
2637 Three-phase lock discipline so the network round-trip does not
2638 block browse / drain paths that share the lock:
2639
2640 1. Acquire, verify session + empty buffer, snapshot
2641 ``session_id`` and ``last_track_id``, release.
2642 2. Call ``client.rotor_session_tracks`` **directly** (no
2643 ``_fetch_rotor_session_batch``) â that helper mutates shared
2644 state (session creation, batch_id write) and would race with
2645 other callers now that we hold no lock. The raw client call
2646 only reads the arguments we pass in.
2647 3. Re-acquire, verify the session hasn't been recycled and the
2648 buffer is still empty, then ``extend``.
2649
2650 :param station_key: Station key whose state to top up.
2651 """
2652 wave = self._wave_states.get(station_key)
2653 if wave is None:
2654 return
2655
2656 async with wave.lock:
2657 if wave.session_id is None or wave.prefetched:
2658 return
2659 session_id = wave.session_id
2660 cursor = wave.last_track_id
2661
2662 if not cursor:
2663 return # No anchor for the next batch yet; try again later.
2664
2665 tracks, _ = await self.client.rotor_session_tracks(session_id, current_track_id=str(cursor))
2666 if not tracks:
2667 return
2668
2669 async with wave.lock:
2670 # Another task could have restarted the session or filled the
2671 # buffer while we were awaiting the network call; bail in both
2672 # cases to avoid stale extends.
2673 if wave.session_id != session_id or wave.prefetched:
2674 return
2675 wave.prefetched.extend(tracks)
2676
2677 async def _fetch_rotor_session_batch(
2678 self, wave: _WaveState, station_id: str
2679 ) -> tuple[list[YandexTrack], str | None]:
2680 """
2681 Fetch the next rotor-session batch for any station.
2682
2683 On first call (wave.session_id is None), starts a new rotor session
2684 and records session_id + batch_id on the wave state. On subsequent
2685 calls, paginates via rotor_session_tracks using wave.last_track_id.
2686
2687 If station_id carries a wave-mode suffix (e.g. "user:onyourwave#discover"),
2688 the suffix maps to a preset in WAVE_MODE_PRESETS and its settings are
2689 merged with wave.settings (wave.settings wins on key conflict). The
2690 base station ID (before "#") is what actually goes to Yandex.
2691
2692 :param wave: The _WaveState for this station (persists across calls).
2693 :param station_id: Rotor station key (may include a "#preset" suffix).
2694 :return: Tuple of (list of yandex tracks, batch_id or None).
2695 """
2696 # Session-creation path: no session yet, or we have a session but no
2697 # cursor yet (`tracks` with an empty queue returns a hard-to-debug
2698 # empty batch â starting a fresh session is the same latency but
2699 # actually yields tracks).
2700 if wave.session_id is None or not wave.last_track_id:
2701 base_station, preset_settings = _split_wave_mode(station_id)
2702 merged = {**preset_settings, **wave.settings}
2703 session_id, tracks, batch_id = await self.client.rotor_session_new(
2704 base_station, settings=merged or None
2705 )
2706 if session_id:
2707 wave.session_id = session_id
2708 else:
2709 tracks, batch_id = await self.client.rotor_session_tracks(
2710 wave.session_id, current_track_id=str(wave.last_track_id)
2711 )
2712 if batch_id:
2713 wave.batch_id = batch_id
2714 return (tracks, batch_id)
2715
2716 async def _browse_waves(
2717 self, path: str, path_parts: list[str]
2718 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
2719 """
2720 Browse waves folder (rotor stations by genre/mood/activity/epoch/local).
2721
2722 Fetches available stations from the Yandex rotor API and groups them by category.
2723
2724 :param path: Full browse path.
2725 :param path_parts: Split path parts after ://.
2726 :return: List of folders or tracks.
2727 """
2728 base = path.rstrip("/") + "/"
2729
2730 locale = (self.mass.metadata.locale or "en_US").lower()
2731 language = "ru" if locale.startswith("ru") else "en"
2732
2733 all_stations = await self.client.get_wave_stations(language)
2734
2735 # Group stations by category, preserving image_url
2736 categorized: dict[str, list[tuple[str, str, str | None]]] = {}
2737 for station_id, cat_key, station_name, image_url in all_stations:
2738 categorized.setdefault(cat_key, []).append((station_id, station_name, image_url))
2739
2740 # waves/ â show category folders
2741 if len(path_parts) == 1:
2742 folders: list[BrowseFolder] = []
2743 # Personalized "My Waves" first â only show if dashboard returns stations
2744 dashboard_stations = await self._get_dashboard_stations_cached()
2745 if dashboard_stations:
2746 name, translation_key = self._media_label("folder", MY_WAVES_FOLDER_ID, "Personal")
2747 folders.append(
2748 BrowseFolder(
2749 item_id=MY_WAVES_FOLDER_ID,
2750 provider=self.instance_id,
2751 path=f"{base}{MY_WAVES_FOLDER_ID}",
2752 name=name,
2753 translation_key=translation_key,
2754 is_playable=False,
2755 )
2756 )
2757 # Featured Waves â only show if landing-blocks/waves returns data
2758 waves_landing = await self._get_waves_landing_cached()
2759 if waves_landing:
2760 name, translation_key = self._media_label(
2761 "folder", WAVES_LANDING_FOLDER_ID, "Featured Waves"
2762 )
2763 folders.append(
2764 BrowseFolder(
2765 item_id=WAVES_LANDING_FOLDER_ID,
2766 provider=self.instance_id,
2767 path=f"{base}{WAVES_LANDING_FOLDER_ID}",
2768 name=name,
2769 translation_key=translation_key,
2770 is_playable=False,
2771 )
2772 )
2773 for cat in WAVE_CATEGORY_DISPLAY_ORDER:
2774 if cat in categorized:
2775 name, translation_key = self._media_label("folder", cat, cat.title())
2776 folders.append(
2777 BrowseFolder(
2778 item_id=cat,
2779 provider=self.instance_id,
2780 path=f"{base}{cat}",
2781 name=name,
2782 translation_key=translation_key,
2783 is_playable=False,
2784 )
2785 )
2786 # Append any categories returned by API that aren't in the predefined order
2787 for cat in categorized:
2788 if cat not in WAVE_CATEGORY_DISPLAY_ORDER:
2789 name, translation_key = self._media_label("folder", cat, cat.title())
2790 folders.append(
2791 BrowseFolder(
2792 item_id=cat,
2793 provider=self.instance_id,
2794 path=f"{base}{cat}",
2795 name=name,
2796 translation_key=translation_key,
2797 is_playable=False,
2798 )
2799 )
2800 return folders
2801
2802 category: str | None = path_parts[1] if len(path_parts) > 1 else None
2803 tag: str | None = path_parts[2] if len(path_parts) > 2 else None
2804
2805 # waves/my_waves/ â show personalized stations from dashboard
2806 if category == MY_WAVES_FOLDER_ID and not tag:
2807 return await self._browse_my_waves_stations(path)
2808
2809 # waves/waves_landing/... â redirect to Featured Waves browse
2810 if category == WAVES_LANDING_FOLDER_ID:
2811 return await self._browse_waves_landing(path, path_parts[1:])
2812
2813 # waves/my_waves/<tag>[/next] â play a specific personal station
2814 # The full station_id has format "genre:allrock", not "my_waves:allrock".
2815 # Resolve by matching against dashboard stations cache.
2816 if category == MY_WAVES_FOLDER_ID and tag:
2817 dashboard_stations = await self._get_dashboard_stations_cached()
2818 for sid, _, _ in dashboard_stations:
2819 sid_tag = sid.split(":", 1)[1] if ":" in sid else sid
2820 if sid_tag == tag:
2821 return await self._browse_wave_station(sid, path=path)
2822 # Fallback: try tag as direct station_id (e.g. "genre:allrock" passed verbatim)
2823 if ":" in tag:
2824 return await self._browse_wave_station(tag, path=path)
2825 return []
2826
2827 # waves/<category>/ â show station folders with artwork
2828 if category and not tag:
2829 cat_stations = categorized.get(category, [])
2830 folders = []
2831 for station_id, station_name, image_url in cat_stations:
2832 tag_part = station_id.split(":", 1)[1] if ":" in station_id else station_id
2833 station_image: MediaItemImage | None = None
2834 if image_url:
2835 station_image = MediaItemImage(
2836 type=ImageType.THUMB,
2837 path=image_url,
2838 provider=self.instance_id,
2839 remotely_accessible=True,
2840 )
2841 folders.append(
2842 BrowseFolder(
2843 item_id=station_id,
2844 provider=self.instance_id,
2845 path=f"{base}{tag_part}",
2846 name=station_name,
2847 is_playable=True,
2848 image=station_image,
2849 )
2850 )
2851 return folders
2852
2853 # waves/<category>/<tag>[/next] â stream tracks from rotor station
2854 if category and tag:
2855 station_id = f"{category}:{tag}"
2856 return await self._browse_wave_station(station_id, path=path)
2857
2858 return []
2859
2860 @use_cache(600, allow_expired_cache=True)
2861 async def _get_dashboard_stations_cached(self) -> list[tuple[str, str, str | None]]:
2862 """
2863 Get personalized dashboard stations, cached for 10 minutes.
2864
2865 :return: List of (station_id, name, image_url) tuples.
2866 """
2867 return await self.client.get_dashboard_stations()
2868
2869 async def _browse_my_waves_stations(self, path: str) -> list[BrowseFolder]:
2870 """
2871 Browse personalized wave stations from rotor/stations/dashboard.
2872
2873 Names are resolved from the non-personalized station list so that
2874 stations show their actual genre/mood name (e.g. "Рок") rather than
2875 the generic "ÐÐ¾Ñ Ð²Ð¾Ð»Ð½Ð°" label that the dashboard API returns.
2876
2877 :param path: Full browse path (used to build sub-paths).
2878 :return: List of playable BrowseFolder items, one per station.
2879 """
2880 stations = await self._get_dashboard_stations_cached()
2881
2882 # Build a name map from the non-personalized list for proper localized names.
2883 locale = (self.mass.metadata.locale or "en_US").lower()
2884 language = "ru" if locale.startswith("ru") else "en"
2885 all_stations = await self.client.get_wave_stations(language)
2886 station_name_map: dict[str, str] = {sid: name for sid, _, name, _ in all_stations}
2887
2888 base = path.rstrip("/") + "/"
2889 folders: list[BrowseFolder] = []
2890 for station_id, fallback_name, image_url in stations:
2891 # Use full station_id (e.g. "genre:rock") in path to avoid collisions
2892 # when two stations share the same tag but differ by category.
2893 # The routing fallback (if ":" in tag) handles this correctly.
2894 name = station_name_map.get(station_id, fallback_name)
2895 station_image: MediaItemImage | None = None
2896 if image_url:
2897 station_image = MediaItemImage(
2898 type=ImageType.THUMB,
2899 path=image_url,
2900 provider=self.instance_id,
2901 remotely_accessible=True,
2902 )
2903 folders.append(
2904 BrowseFolder(
2905 item_id=station_id,
2906 provider=self.instance_id,
2907 path=f"{base}{station_id}",
2908 name=name,
2909 is_playable=True,
2910 image=station_image,
2911 )
2912 )
2913 return folders
2914
2915 async def _browse_wave_station(
2916 self, station_id: str, path: str = ""
2917 ) -> list[Track | BrowseFolder]:
2918 """
2919 Browse a rotor wave station and return tracks.
2920
2921 Fetches tracks from the rotor station, deduplicates within the current session,
2922 and sends radioStarted feedback on first call. Appends a "Load more" BrowseFolder
2923 at the end so MA can continue fetching the next batch automatically (radio mode).
2924
2925 :param station_id: Rotor station ID (e.g. 'genre:rock', 'mood:chill').
2926 :param path: Current browse path, used to construct the "Load more" next path.
2927 :return: List of Track objects with composite item_id (track_id@station_id),
2928 followed by a "Load more" BrowseFolder if more tracks are available.
2929 """
2930 state = self._get_wave_state(station_id)
2931 async with state.lock:
2932 max_tracks = int(
2933 self.config.get_value(CONF_MY_WAVE_MAX_TRACKS) or 150 # type: ignore[arg-type]
2934 )
2935
2936 self.logger.debug(
2937 "Browse wave station: station_id=%s path=%s last_track_id=%s session=%s",
2938 station_id,
2939 path,
2940 state.last_track_id,
2941 state.session_id,
2942 )
2943 # Tagged stations (genre:*, mood:*, activity:*, epoch:*) accept the
2944 # same /rotor/session/* endpoint as user:onyourwave / track:{id},
2945 # verified against the live Yandex API. Reuse the session helper so
2946 # batch_id + session_id stay anchored across browse/play/feedback.
2947 yandex_tracks, _ = await self._fetch_rotor_session_batch(state, station_id)
2948
2949 if not state.radio_started_sent and yandex_tracks:
2950 sent = await self._send_wave_feedback(state, station_id, "radioStarted")
2951 if sent:
2952 state.radio_started_sent = True
2953
2954 tracks: list[Track] = []
2955 first_track_id: str | None = None
2956 for yt in yandex_tracks:
2957 if len(state.seen_track_ids) >= max_tracks:
2958 break
2959 track = self._parse_my_wave_track(yt, state.seen_track_ids)
2960 if track is None:
2961 continue
2962 # Override station_id in composite item_id to reflect this specific station
2963 old_item_id = track.item_id
2964 track_id = old_item_id.split(RADIO_TRACK_ID_SEP, 1)[0]
2965 track.item_id = f"{track_id}{RADIO_TRACK_ID_SEP}{station_id}"
2966 # Keep provider mappings in sync with the new item_id
2967 for pm in getattr(track, "provider_mappings", []):
2968 if (
2969 getattr(pm, "item_id", None) == old_item_id
2970 and getattr(pm, "provider_instance", None) == self.instance_id
2971 ):
2972 pm.item_id = track.item_id
2973 if first_track_id is None:
2974 first_track_id = track_id
2975 tracks.append(track)
2976
2977 if first_track_id is not None:
2978 state.last_track_id = first_track_id
2979
2980 self.logger.debug(
2981 "Wave station %s returned %d tracks: %s",
2982 station_id,
2983 len(tracks),
2984 [t.item_id.split(RADIO_TRACK_ID_SEP, 1)[0] for t in tracks[:5]],
2985 )
2986 result: list[Track | BrowseFolder] = list(tracks)
2987
2988 # Append "Load more" sentinel so MA knows to call browse again for next batch.
2989 # This mirrors the My Wave mechanism and enables continuous radio playback.
2990 if tracks and len(state.seen_track_ids) < max_tracks and path:
2991 # Append /next to the current path (same pattern as _browse_my_wave).
2992 # This makes each "Load more" path unique (e.g. /next/next/next...)
2993 # so MA never serves a cached result for subsequent presses.
2994 result.append(
2995 BrowseFolder(
2996 item_id="next",
2997 provider=self.instance_id,
2998 path=f"{path.rstrip('/')}/next",
2999 name="Load more",
3000 translation_key="load_more",
3001 is_playable=False,
3002 )
3003 )
3004
3005 return result
3006
3007 @staticmethod
3008 def _extract_wave_item_cover(item: dict[str, Any]) -> tuple[str | None, str | None]:
3009 """
3010 Extract cover URI and background color from a wave/mix item.
3011
3012 Accepts both camelCase (``compactImageUrl`` â what /landing-blocks/
3013 actually returns) and snake_case (``compact_image_url`` â retained
3014 for safety if MarshalX ever normalises the payload).
3015
3016 :param item: Wave or mix item dict from the API.
3017 :return: (cover_uri, bg_color) tuple where bg_color is a hex string or None.
3018 """
3019 agent_uri = item.get("agent", {}).get("cover", {}).get("uri", "")
3020 cover_uri = agent_uri or item.get("compactImageUrl") or item.get("compact_image_url")
3021 bg_color = item.get("colors", {}).get("average")
3022 return cover_uri, bg_color
3023
3024 @use_cache(3600, allow_expired_cache=True)
3025 async def _get_mixes_waves_cached(self) -> list[dict[str, Any]] | None:
3026 """
3027 Get AI Wave Set data from /landing-blocks/mixes-waves, cached for 1 hour.
3028
3029 :return: List of mix category dicts from the API, or None on error.
3030 """
3031 return await self.client.get_mixes_waves()
3032
3033 @use_cache(3600, allow_expired_cache=True)
3034 async def _get_waves_landing_cached(self) -> list[dict[str, Any]] | None:
3035 """
3036 Get Featured Waves data from /landing-blocks/waves, cached for 1 hour.
3037
3038 :return: List of wave category dicts from the API, or None on error.
3039 """
3040 return await self.client.get_waves_landing()
3041
3042 async def _browse_waves_landing(
3043 self, path: str, path_parts: list[str]
3044 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
3045 """
3046 Browse Featured Waves (from /landing-blocks/waves).
3047
3048 :param path: Full browse path.
3049 :param path_parts: Split path parts after ://.
3050 :return: List of folders or tracks.
3051 """
3052 waves_data = await self._get_waves_landing_cached()
3053 return await self._browse_wave_categories(
3054 path, path_parts, waves_data or [], WAVES_LANDING_FOLDER_ID
3055 )
3056
3057 async def _browse_wave_categories(
3058 self,
3059 path: str,
3060 path_parts: list[str],
3061 categories_data: list[dict[str, Any]],
3062 id_prefix: str,
3063 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
3064 """
3065 Browse wave-like category folders and their station items.
3066
3067 Shared logic for both 'my_waves_set' browse trees:
3068 - Level 1 (e.g. my_waves_set/): category folders
3069 - Level 2 (e.g. my_waves_set/ai-sets/): playable station folders with artwork
3070 - Level 3+ (e.g. my_waves_set/ai-sets/genre:rock[/next]): track listing
3071
3072 :param path: Full browse path.
3073 :param path_parts: Split path parts after ://.
3074 :param categories_data: List of category dicts from the API.
3075 :param id_prefix: Prefix for BrowseFolder item_id (e.g. 'my_waves_set').
3076 :return: List of folders or tracks.
3077 """
3078 base = path.rstrip("/") + "/"
3079
3080 if not categories_data:
3081 return []
3082
3083 # Level 1 â category folders
3084 if len(path_parts) == 1:
3085 folders: list[BrowseFolder] = []
3086 for wave_category in categories_data:
3087 cat_id = wave_category.get("id", "")
3088 cat_title = wave_category.get("title", "")
3089 items = wave_category.get("items", [])
3090 if not items or not cat_id:
3091 continue
3092 display_name = cat_title.capitalize() if cat_title else cat_id.capitalize()
3093 folders.append(
3094 BrowseFolder(
3095 item_id=f"{id_prefix}_{cat_id}",
3096 provider=self.instance_id,
3097 path=f"{base}{cat_id}",
3098 name=display_name,
3099 is_playable=False,
3100 )
3101 )
3102 return folders
3103
3104 category_id = path_parts[1] if len(path_parts) > 1 else None
3105 if not category_id:
3106 return []
3107
3108 # Level 3+ â stream tracks from rotor station
3109 if len(path_parts) > 2:
3110 station_id = path_parts[2]
3111 return await self._browse_wave_station(station_id, path=path)
3112
3113 # Level 2 â playable station folders with artwork
3114 for wave_category in categories_data:
3115 if wave_category.get("id") == category_id:
3116 items = wave_category.get("items", [])
3117 result: list[BrowseFolder] = []
3118 for item in items:
3119 # API returns camelCase (`stationId`); keep snake_case as a
3120 # safety net if the payload is ever normalised upstream.
3121 station_id = item.get("stationId") or item.get("station_id") or ""
3122 title = item.get("title", "")
3123 if not station_id or not title:
3124 continue
3125 cover_uri, bg_color = self._extract_wave_item_cover(item)
3126 image: MediaItemImage | None = None
3127 if cover_uri:
3128 if cover_uri.startswith("http"):
3129 img_url: str = cover_uri.replace("%%", IMAGE_SIZE_MEDIUM)
3130 else:
3131 raw = get_image_url(cover_uri)
3132 img_url = "" if raw is None else raw
3133 if img_url:
3134 if bg_color:
3135 # Append bg_color as URL fragment for cache-key uniqueness.
3136 # MA will call resolve_image() to composite the transparent PNG.
3137 if len(self._wave_bg_colors) > 200:
3138 self._wave_bg_colors.clear()
3139 img_url = f"{img_url}#{bg_color.lstrip('#')}"
3140 self._wave_bg_colors[img_url] = bg_color
3141 image = MediaItemImage(
3142 type=ImageType.THUMB,
3143 path=img_url,
3144 provider=self.instance_id,
3145 remotely_accessible=bg_color is None,
3146 )
3147 result.append(
3148 BrowseFolder(
3149 item_id=station_id,
3150 provider=self.instance_id,
3151 path=f"{base}{station_id}",
3152 name=title,
3153 is_playable=True,
3154 image=image,
3155 )
3156 )
3157 return result
3158
3159 return []
3160
3161 async def _browse_vibe_sets(
3162 self, path: str, path_parts: list[str]
3163 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
3164 """
3165 Browse AI Wave Sets (from /landing-blocks/mixes-waves).
3166
3167 :param path: Full browse path.
3168 :param path_parts: Split path parts after ://.
3169 :return: List of folders or tracks.
3170 """
3171 mixes_data = await self._get_mixes_waves_cached()
3172 return await self._browse_wave_categories(
3173 path, path_parts, mixes_data or [], MY_WAVES_SET_FOLDER_ID
3174 )
3175
3176 @use_cache(600, allow_expired_cache=True)
3177 async def _get_tag_playlists_as_browse(
3178 self, tag_id: str
3179 ) -> Sequence[MediaItemType | ItemMapping | BrowseFolder]:
3180 """
3181 Get playlists for a tag and return as browse items.
3182
3183 :param tag_id: Tag identifier (e.g. 'chill', '80s').
3184 :return: List of Playlist objects.
3185 """
3186 self.logger.debug("Fetching playlists for tag: %s", tag_id)
3187 playlists = await self.client.get_tag_playlists(tag_id)
3188 self.logger.debug("Got %d playlists for tag %s", len(playlists), tag_id)
3189 result: list[Playlist] = []
3190 for playlist in playlists:
3191 try:
3192 result.append(parse_playlist(self, playlist))
3193 except InvalidDataError as err:
3194 self.logger.debug("Error parsing tag playlist: %s", err)
3195 self.logger.debug("Parsed %d playlists for tag %s", len(result), tag_id)
3196 return result
3197
3198 @use_cache(3600 * 24 * 30, allow_expired_cache=True)
3199 async def _get_track_cached(self, track_id: str) -> Track:
3200 """
3201 Get track details by normalized ID (cached).
3202
3203 :param track_id: Normalized track ID (without station suffix).
3204 :return: Track object.
3205 :raises MediaNotFoundError: If track not found.
3206 """
3207 yandex_track = await self.client.get_track(track_id)
3208 if not yandex_track:
3209 raise MediaNotFoundError(f"Track {track_id} not found")
3210
3211 # Use the already-fetched track object to avoid a duplicate API call
3212 lyrics, lyrics_synced = await self.client.get_track_lyrics_from_track(yandex_track)
3213
3214 return parse_track(self, yandex_track, lyrics=lyrics, lyrics_synced=lyrics_synced)
3215
3216 @use_cache(3600 * 24 * 30, allow_expired_cache=True)
3217 async def _get_real_playlist(self, prov_playlist_id: str) -> Playlist:
3218 """
3219 Get real playlist details by ID (cached).
3220
3221 :param prov_playlist_id: The provider playlist ID (format: "owner_id:kind").
3222 :return: Playlist object.
3223 :raises MediaNotFoundError: If playlist not found.
3224 """
3225 # Parse the playlist ID (format: owner_id:kind)
3226 if PLAYLIST_ID_SPLITTER in prov_playlist_id:
3227 owner_id, kind = prov_playlist_id.split(PLAYLIST_ID_SPLITTER, 1)
3228 else:
3229 owner_id = str(self.client.user_id)
3230 kind = prov_playlist_id
3231
3232 playlist = await self.client.get_playlist(owner_id, kind)
3233 if not playlist:
3234 raise MediaNotFoundError(f"Playlist {prov_playlist_id} not found")
3235 return parse_playlist(self, playlist)
3236
3237 @use_cache(3600 * 3, allow_expired_cache=True)
3238 async def _get_my_wave_playlist_tracks(self, page: int) -> list[Track]:
3239 """
3240 Get My Wave tracks for virtual playlist (uses cursor for page > 0).
3241
3242 Fetches MY_WAVE_BATCH_SIZE Rotor API batches per page call to reduce
3243 the number of round-trips when the player controller paginates through pages.
3244
3245 :param page: Page number (0 = first batch, 1+ = next batches via queue cursor).
3246 :return: List of Track objects for this page.
3247 """
3248 wave = self._get_wave_state(ROTOR_STATION_MY_WAVE)
3249 async with wave.lock:
3250 max_tracks_config = int(
3251 self.config.get_value(CONF_MY_WAVE_MAX_TRACKS) or 150 # type: ignore[arg-type]
3252 )
3253
3254 # Reset seen tracks on first page
3255 if page == 0:
3256 wave.seen_track_ids = set()
3257
3258 queue: str | int | None = None
3259 if page > 0:
3260 queue = wave.playlist_next_cursor
3261 if not queue:
3262 return []
3263
3264 # Check if we've already reached the limit
3265 if len(wave.seen_track_ids) >= max_tracks_config:
3266 return []
3267
3268 tracks: list[Track] = []
3269 next_cursor: str | None = None
3270
3271 # Fetch MY_WAVE_BATCH_SIZE Rotor API batches per page to reduce API round-trips
3272 for _ in range(MY_WAVE_BATCH_SIZE):
3273 if len(wave.seen_track_ids) >= max_tracks_config:
3274 break
3275
3276 if queue is not None:
3277 wave.last_track_id = str(queue)
3278 yandex_tracks, _ = await self._fetch_rotor_session_batch(
3279 wave, ROTOR_STATION_MY_WAVE
3280 )
3281 if not wave.radio_started_sent and yandex_tracks:
3282 sent = await self._send_wave_feedback(
3283 wave, ROTOR_STATION_MY_WAVE, "radioStarted"
3284 )
3285 if sent:
3286 wave.radio_started_sent = True
3287
3288 if not yandex_tracks:
3289 break
3290
3291 first_track_id_this_batch = None
3292 for yt in yandex_tracks:
3293 if len(wave.seen_track_ids) >= max_tracks_config:
3294 break
3295
3296 track = self._parse_my_wave_track(yt, wave.seen_track_ids)
3297 if track is None:
3298 continue
3299
3300 tracks.append(track)
3301 track_id = track.item_id.split(RADIO_TRACK_ID_SEP, 1)[0]
3302 if first_track_id_this_batch is None:
3303 first_track_id_this_batch = track_id
3304
3305 if first_track_id_this_batch is not None:
3306 next_cursor = first_track_id_this_batch
3307 queue = first_track_id_this_batch
3308 else:
3309 # All tracks in this batch were duplicates or failed to parse
3310 break
3311
3312 # Store cursor for next page call (None clears pagination so next call returns [])
3313 wave.playlist_next_cursor = next_cursor
3314 return tracks
3315
3316 @use_cache(3600 * 3, allow_expired_cache=True)
3317 async def _get_liked_tracks_playlist_tracks(self, page: int) -> list[Track]:
3318 """
3319 Get liked tracks for virtual playlist (sorted in reverse chronological order).
3320
3321 :param page: Page number (0 = all tracks limited by config, >0 = empty for pagination).
3322 :return: List of Track objects.
3323 """
3324 # Liked tracks API returns all tracks at once, so only return tracks on page 0
3325 if page > 0:
3326 return []
3327
3328 max_tracks_config = int(
3329 self.config.get_value(CONF_LIKED_TRACKS_MAX_TRACKS) or 200 # type: ignore[arg-type]
3330 )
3331
3332 # Fetch liked tracks (already sorted in reverse chronological order by api_client)
3333 track_shorts = await self.client.get_liked_tracks()
3334 if not track_shorts:
3335 self.logger.debug("No liked tracks found")
3336 return []
3337
3338 # Apply max tracks limit
3339 track_shorts = track_shorts[:max_tracks_config]
3340
3341 # Fetch full track details in batches
3342 track_ids = [str(ts.track_id) for ts in track_shorts if ts.track_id]
3343
3344 batch_size = TRACK_BATCH_SIZE
3345 full_tracks = []
3346 for i in range(0, len(track_ids), batch_size):
3347 batch_ids = track_ids[i : i + batch_size]
3348 batch_result = await self.client.get_tracks(batch_ids)
3349 full_tracks.extend(batch_result)
3350 # Spread bursts: insert a small jittered pause between batches so
3351 # a 500-track hydration doesn't look like a bot to Yandex's
3352 # smart-captcha. Skipped after the last batch.
3353 if i + batch_size < len(track_ids):
3354 await asyncio.sleep(
3355 LIKED_BATCH_JITTER_MIN_S + random.random() * LIKED_BATCH_JITTER_SPAN_S
3356 )
3357
3358 # Create track ID to full track mapping by track ID directly
3359 track_map = {}
3360 for t in full_tracks:
3361 if hasattr(t, "id") and t.id:
3362 track_map[str(t.id)] = t
3363
3364 # Parse tracks in the original order (reverse chronological)
3365 tracks = []
3366 for track_id in track_ids:
3367 # track_id may be compound "trackId:albumId", extract base ID for lookup
3368 base_id = track_id.split(":")[0] if ":" in track_id else track_id
3369 found = track_map.get(track_id) or track_map.get(base_id)
3370 if found:
3371 try:
3372 tracks.append(parse_track(self, found))
3373 except InvalidDataError as err:
3374 self.logger.debug("Error parsing liked track %s: %s", track_id, err)
3375
3376 self.logger.debug("Liked tracks: fetched %s, parsed %s", len(track_shorts), len(tracks))
3377 return tracks
3378
3379 @use_cache(3600 * 3, allow_expired_cache=True)
3380 async def _get_regular_playlist_tracks(self, prov_playlist_id: str, page: int) -> list[Track]:
3381 """
3382 Get the tracks of a regular (non-virtual) playlist.
3383
3384 :param prov_playlist_id: The provider playlist ID (format: "owner_id:kind").
3385 :param page: Page number for pagination.
3386 :return: List of Track objects.
3387 """
3388 # Yandex Music API returns all playlist tracks in one call (no server-side pagination).
3389 # Return empty list for page > 0 so the controller pagination loop terminates.
3390 if page > 0:
3391 return []
3392
3393 # Parse the playlist ID (format: owner_id:kind)
3394 if PLAYLIST_ID_SPLITTER in prov_playlist_id:
3395 owner_id, kind = prov_playlist_id.split(PLAYLIST_ID_SPLITTER, 1)
3396 else:
3397 owner_id = str(self.client.user_id)
3398 kind = prov_playlist_id
3399
3400 playlist = await self.client.get_playlist(owner_id, kind)
3401 if not playlist:
3402 return []
3403
3404 # API sometimes returns playlist without tracks; fetch them explicitly if needed
3405 tracks_list = playlist.tracks or []
3406 track_count = getattr(playlist, "track_count", None) or 0
3407 if not tracks_list and track_count > 0:
3408 self.logger.debug(
3409 "Playlist %s/%s: track_count=%s but no tracks in response, "
3410 "calling fetch_tracks_async",
3411 owner_id,
3412 kind,
3413 track_count,
3414 )
3415 try:
3416 tracks_list = await playlist.fetch_tracks_async()
3417 except Exception as err:
3418 self.logger.warning("fetch_tracks_async failed for %s/%s: %s", owner_id, kind, err)
3419 if not tracks_list:
3420 raise ResourceTemporarilyUnavailable(
3421 "Playlist tracks not available; try again later"
3422 )
3423
3424 if not tracks_list:
3425 return []
3426
3427 # Yandex returns TrackShort objects, we need to fetch full track info
3428 track_ids = [
3429 str(track.track_id) if hasattr(track, "track_id") else str(track.id)
3430 for track in tracks_list
3431 if track
3432 ]
3433 if not track_ids:
3434 return []
3435
3436 # Fetch full track details in batches to avoid timeouts
3437 batch_size = TRACK_BATCH_SIZE
3438 full_tracks = []
3439 for i in range(0, len(track_ids), batch_size):
3440 batch = track_ids[i : i + batch_size]
3441 batch_result = await self.client.get_tracks(batch)
3442 if not batch_result:
3443 # Skip this batch but keep going â the terminal guard below
3444 # raises if every batch comes back empty. Aborting on a single
3445 # empty batch threw away tracks already fetched from earlier
3446 # batches and forced a full retry hours later (under the
3447 # @use_cache TTL above).
3448 self.logger.warning(
3449 "Empty batch %s-%s for playlist %s, skipping",
3450 i,
3451 i + len(batch) - 1,
3452 prov_playlist_id,
3453 )
3454 continue
3455 full_tracks.extend(batch_result)
3456
3457 if track_ids and not full_tracks:
3458 raise ResourceTemporarilyUnavailable("Failed to load track details; try again later")
3459
3460 tracks = []
3461 for track in full_tracks:
3462 try:
3463 tracks.append(parse_track(self, track))
3464 except InvalidDataError as err:
3465 self.logger.debug("Error parsing playlist track: %s", err)
3466 return tracks
3467
3468 async def _drain_prefetched_wave_tracks(self, station_key: str, limit: int) -> list[Track]:
3469 """
3470 Pop up to ``limit`` prefetched tracks off the wave state.
3471
3472 Runs under ``wave.lock`` so it doesn't race with
3473 ``_prefetch_rotor_session`` which extends the same list under the
3474 same lock. Returns an empty list when there's no active session or
3475 nothing prefetched; callers then fall through to the cached fetch.
3476
3477 This method is intentionally not cached â it mutates wave state.
3478 """
3479 wave = self._wave_states.get(station_key)
3480 if not (wave and wave.session_id and wave.prefetched):
3481 return []
3482 async with wave.lock:
3483 if not wave.prefetched:
3484 return []
3485 drained_yt = wave.prefetched[:limit]
3486 wave.prefetched = wave.prefetched[limit:]
3487 tracks: list[Track] = []
3488 for yt in drained_yt:
3489 try:
3490 tracks.append(parse_track(self, yt))
3491 except InvalidDataError as err:
3492 self.logger.debug("Error parsing prefetched wave track: %s", err)
3493 return tracks
3494
3495 @use_cache(3600 * 3, allow_expired_cache=True)
3496 async def _fetch_similar_tracks_for_seed(self, track_id: str, limit: int) -> list[Track]:
3497 """
3498 Create a one-off rotor session for ``track:{id}`` and return up to ``limit`` tracks.
3499
3500 Stateless by design: similar-tracks results don't participate in
3501 playback feedback or prefetch, so there is no need to keep a
3502 ``_WaveState`` entry around. Going through ``_fetch_rotor_session_batch``
3503 would create one per unique seed and grow ``_wave_states`` without
3504 bound under normal DSTM usage; call ``rotor_session_new`` directly
3505 instead.
3506
3507 Pure function of ``track_id`` / ``limit``, hence safe to memoise
3508 via ``@use_cache``.
3509 """
3510 _, yandex_tracks, _ = await self.client.rotor_session_new(f"track:{track_id}")
3511 similar_tracks: list[Track] = []
3512 for yt in yandex_tracks[:limit]:
3513 try:
3514 similar_tracks.append(parse_track(self, yt))
3515 except InvalidDataError as err:
3516 self.logger.debug("Error parsing similar track: %s", err)
3517 return similar_tracks
3518
3519 @use_cache(600, allow_expired_cache=True)
3520 async def _get_my_wave_recommendations(self) -> RecommendationFolder | None:
3521 """
3522 Get My Wave recommendation folder with personalized tracks.
3523
3524 Shares the same `_WaveState(ROTOR_STATION_MY_WAVE)` with browse and
3525 virtual-playlist flows, so session_id + batch_id established here
3526 carry into `on_played`/`on_streamed` feedback even when the user
3527 starts playback from this discovery card.
3528
3529 :return: RecommendationFolder with My Wave tracks, or None if empty.
3530 """
3531 max_tracks_config = int(
3532 self.config.get_value(CONF_MY_WAVE_MAX_TRACKS) or 150 # type: ignore[arg-type]
3533 )
3534 batch_size_config = MY_WAVE_BATCH_SIZE
3535
3536 wave = self._get_wave_state(ROTOR_STATION_MY_WAVE)
3537 # Local dedup so the recommendations card stays independent from the
3538 # browse/virtual-playlist dedup set (which may be larger and stale).
3539 # Only session_id + batch_id + last_track_id are shared with `wave`.
3540 seen_track_ids: set[str] = set()
3541 items: list[Track] = []
3542
3543 # Hold the wave lock across the whole fetch chain â we mutate shared
3544 # session_id/batch_id/last_track_id via _fetch_rotor_session_batch,
3545 # and other call sites (browse, virtual-playlist) guard the same
3546 # state with this lock. Concurrent calls without the lock would
3547 # interleave cursor updates and leave the session inconsistent.
3548 async with wave.lock:
3549 for _ in range(batch_size_config):
3550 if len(seen_track_ids) >= max_tracks_config:
3551 break
3552
3553 yandex_tracks, _ = await self._fetch_rotor_session_batch(
3554 wave, ROTOR_STATION_MY_WAVE
3555 )
3556 if not yandex_tracks:
3557 break
3558
3559 first_track_id_this_batch: str | None = None
3560 for yt in yandex_tracks:
3561 if len(seen_track_ids) >= max_tracks_config:
3562 break
3563
3564 track = self._parse_my_wave_track(yt, seen_ids=seen_track_ids)
3565 if track is None:
3566 continue
3567
3568 items.append(track)
3569 track_id = track.item_id.split(RADIO_TRACK_ID_SEP, 1)[0]
3570 if first_track_id_this_batch is None:
3571 first_track_id_this_batch = track_id
3572
3573 if first_track_id_this_batch is None:
3574 break
3575 wave.last_track_id = first_track_id_this_batch
3576
3577 if not items:
3578 return None
3579
3580 initial_tracks_limit = DISCOVERY_INITIAL_TRACKS
3581 if len(items) > initial_tracks_limit:
3582 items = items[:initial_tracks_limit]
3583
3584 return RecommendationFolder(
3585 item_id=MY_WAVE_PLAYLIST_ID,
3586 provider=self.instance_id,
3587 name="My Wave",
3588 translation_key=MY_WAVE_PLAYLIST_ID,
3589 items=UniqueList(items),
3590 icon="mdi-waveform",
3591 )
3592
3593 @use_cache(1800, allow_expired_cache=True)
3594 async def _get_feed_recommendations(self) -> RecommendationFolder | None:
3595 """
3596 Get personalized feed playlists (Playlist of the Day, DejaVu, etc.).
3597
3598 :return: RecommendationFolder with generated playlists, or None if unavailable.
3599 """
3600 feed = await self.client.get_feed()
3601 if not feed or not feed.generated_playlists:
3602 return None
3603 items: list[Playlist] = []
3604 for gen_playlist in feed.generated_playlists:
3605 if gen_playlist.data and gen_playlist.ready:
3606 try:
3607 # Mark feed-generated playlists (Playlist of the Day, DejaVu,
3608 # Premiere, Missed Likes) as dynamic â Yandex regenerates them
3609 # on a schedule so MA must not long-cache the track list.
3610 items.append(parse_playlist(self, gen_playlist.data, is_dynamic=True))
3611 except InvalidDataError as err:
3612 self.logger.debug("Error parsing feed playlist: %s", err)
3613 if not items:
3614 return None
3615 return RecommendationFolder(
3616 item_id="feed",
3617 provider=self.instance_id,
3618 name="Made for You",
3619 translation_key="feed",
3620 items=UniqueList(items),
3621 icon="mdi-account-music",
3622 )
3623
3624 @use_cache(3600, allow_expired_cache=True)
3625 async def _get_chart_recommendations(self) -> RecommendationFolder | None:
3626 """
3627 Get chart tracks (hot tracks of the month).
3628
3629 :return: RecommendationFolder with chart tracks, or None if unavailable.
3630 """
3631 chart_info = await self.client.get_chart()
3632 if not chart_info or not chart_info.chart:
3633 return None
3634 playlist = chart_info.chart
3635 if not playlist.tracks:
3636 return None
3637 # TrackShort objects in chart context have .track (full Track) and .chart (position)
3638 tracks: list[Track] = []
3639 for track_short in playlist.tracks[:20]:
3640 track_obj = getattr(track_short, "track", None)
3641 if not track_obj:
3642 continue
3643 try:
3644 tracks.append(parse_track(self, track_obj))
3645 except InvalidDataError as err:
3646 self.logger.debug("Error parsing chart track: %s", err)
3647 if not tracks:
3648 return None
3649 return RecommendationFolder(
3650 item_id="chart",
3651 provider=self.instance_id,
3652 name="Chart",
3653 translation_key="chart",
3654 items=UniqueList(tracks),
3655 icon="mdi-chart-line",
3656 )
3657
3658 @use_cache(3600, allow_expired_cache=True)
3659 async def _get_new_releases_recommendations(self) -> RecommendationFolder | None:
3660 """
3661 Get new album releases.
3662
3663 :return: RecommendationFolder with new albums, or None if unavailable.
3664 """
3665 releases = await self.client.get_new_releases()
3666 if not releases or not releases.new_releases:
3667 return None
3668 # new_releases is a list of album IDs (int) â need to batch-fetch full details
3669 album_ids = [str(aid) for aid in releases.new_releases[:20]]
3670 if not album_ids:
3671 return None
3672 full_albums = await self.client.get_albums(album_ids)
3673 if not full_albums:
3674 return None
3675 albums: list[Album] = []
3676 for album in full_albums:
3677 try:
3678 albums.append(parse_album(self, album))
3679 except InvalidDataError as err:
3680 self.logger.debug("Error parsing new release album: %s", err)
3681 if not albums:
3682 return None
3683 return RecommendationFolder(
3684 item_id="new_releases",
3685 provider=self.instance_id,
3686 name="New Releases",
3687 translation_key="new_releases",
3688 items=UniqueList(albums),
3689 icon="mdi-new-box",
3690 )
3691
3692 @use_cache(3600, allow_expired_cache=True)
3693 async def _get_new_playlists_recommendations(self) -> RecommendationFolder | None:
3694 """
3695 Get new editorial playlists.
3696
3697 :return: RecommendationFolder with new playlists, or None if unavailable.
3698 """
3699 result = await self.client.get_new_playlists()
3700 if not result or not result.new_playlists:
3701 return None
3702 # new_playlists is a list of PlaylistId objects (uid, kind) â fetch full details
3703 playlist_ids = [
3704 f"{pid.uid}:{pid.kind}"
3705 for pid in result.new_playlists[:20]
3706 if hasattr(pid, "uid") and hasattr(pid, "kind")
3707 ]
3708 if not playlist_ids:
3709 return None
3710 full_playlists = await self.client.get_playlists(playlist_ids)
3711 if not full_playlists:
3712 return None
3713 playlists: list[Playlist] = []
3714 for playlist in full_playlists:
3715 try:
3716 playlists.append(parse_playlist(self, playlist))
3717 except InvalidDataError as err:
3718 self.logger.debug("Error parsing new playlist: %s", err)
3719 if not playlists:
3720 return None
3721 return RecommendationFolder(
3722 item_id="new_playlists",
3723 provider=self.instance_id,
3724 name="New Playlists",
3725 translation_key="new_playlists",
3726 items=UniqueList(playlists),
3727 icon="mdi-playlist-star",
3728 )
3729
3730 @use_cache(3600, allow_expired_cache=True)
3731 async def _get_top_picks_recommendations(self) -> RecommendationFolder | None:
3732 """
3733 Get Top Picks recommendation folder (tag: top).
3734
3735 :return: RecommendationFolder with top playlists, or None if unavailable.
3736 """
3737 playlists = await self.client.get_tag_playlists("top")
3738 if not playlists:
3739 return None
3740 items: list[Playlist] = []
3741 for playlist in playlists[:10]:
3742 try:
3743 items.append(parse_playlist(self, playlist))
3744 except InvalidDataError as err:
3745 self.logger.debug("Error parsing top picks playlist: %s", err)
3746 if not items:
3747 return None
3748 return RecommendationFolder(
3749 item_id="top_picks",
3750 provider=self.instance_id,
3751 name="Top Picks",
3752 translation_key="top_picks",
3753 items=UniqueList(items),
3754 icon="mdi-star",
3755 )
3756
3757 @use_cache(1800, allow_expired_cache=True)
3758 async def _get_mood_mix_recommendations(self, mood_tag: str) -> RecommendationFolder | None:
3759 """
3760 Get Mood Mix recommendation folder for a specific tag.
3761
3762 :param mood_tag: Preselected mood tag slug.
3763 :return: RecommendationFolder with mood playlists, or None if unavailable.
3764 """
3765 playlists = await self.client.get_tag_playlists(mood_tag)
3766 if not playlists:
3767 self.logger.debug("No playlists for mood tag %s, skipping recommendation", mood_tag)
3768 return None
3769 items: list[Playlist] = []
3770 for playlist in playlists[:8]:
3771 try:
3772 items.append(parse_playlist(self, playlist))
3773 except InvalidDataError as err:
3774 self.logger.debug("Error parsing mood playlist: %s", err)
3775 if not items:
3776 return None
3777 tag_name, _ = self._media_label("folder", _media_label_key(mood_tag), mood_tag.title())
3778 return RecommendationFolder(
3779 item_id="mood_mix",
3780 provider=self.instance_id,
3781 name=f"Mood Mix: {tag_name}",
3782 translation_key="mood_mix",
3783 translation_params=[tag_name],
3784 items=UniqueList(items),
3785 icon="mdi-emoticon-outline",
3786 )
3787
3788 @use_cache(1800, allow_expired_cache=True)
3789 async def _get_activity_mix_recommendations(
3790 self, activity_tag: str
3791 ) -> RecommendationFolder | None:
3792 """
3793 Get Activity Mix recommendation folder for a specific tag.
3794
3795 :param activity_tag: Preselected activity tag slug.
3796 :return: RecommendationFolder with activity playlists, or None if unavailable.
3797 """
3798 playlists = await self.client.get_tag_playlists(activity_tag)
3799 if not playlists:
3800 self.logger.debug(
3801 "No playlists for activity tag %s, skipping recommendation", activity_tag
3802 )
3803 return None
3804 items: list[Playlist] = []
3805 for playlist in playlists[:8]:
3806 try:
3807 items.append(parse_playlist(self, playlist))
3808 except InvalidDataError as err:
3809 self.logger.debug("Error parsing activity playlist: %s", err)
3810 if not items:
3811 return None
3812 tag_name, _ = self._media_label(
3813 "folder", _media_label_key(activity_tag), activity_tag.title()
3814 )
3815 return RecommendationFolder(
3816 item_id="activity_mix",
3817 provider=self.instance_id,
3818 name=f"Activity Mix: {tag_name}",
3819 translation_key="activity_mix",
3820 translation_params=[tag_name],
3821 items=UniqueList(items),
3822 icon="mdi-run",
3823 )
3824
3825 @use_cache(3600 * 6, allow_expired_cache=True)
3826 async def _get_seasonal_mix_recommendations(self) -> RecommendationFolder | None:
3827 """
3828 Get Seasonal Mix recommendation folder (based on current month).
3829
3830 :return: RecommendationFolder with seasonal playlists, or None if unavailable.
3831 """
3832 # Determine current season tag; fall back to autumn if the seasonal
3833 # endpoint returns nothing (e.g. spring/autumn handover gap).
3834 current_month = utc().month
3835 seasonal_tag = TAG_SEASONAL_MAP.get(current_month, "autumn")
3836 playlists = await self.client.get_tag_playlists(seasonal_tag)
3837 if not playlists and seasonal_tag != "autumn":
3838 seasonal_tag = "autumn"
3839 playlists = await self.client.get_tag_playlists(seasonal_tag)
3840 if not playlists:
3841 return None
3842 items: list[Playlist] = []
3843 for playlist in playlists[:8]:
3844 try:
3845 items.append(parse_playlist(self, playlist))
3846 except InvalidDataError as err:
3847 self.logger.debug("Error parsing seasonal playlist: %s", err)
3848 if not items:
3849 return None
3850 tag_name, _ = self._media_label(
3851 "folder", _media_label_key(seasonal_tag), seasonal_tag.title()
3852 )
3853 return RecommendationFolder(
3854 item_id="seasonal_mix",
3855 provider=self.instance_id,
3856 name=f"Seasonal: {tag_name}",
3857 translation_key="seasonal_mix",
3858 translation_params=[tag_name],
3859 items=UniqueList(items),
3860 icon="mdi-weather-sunny",
3861 )
3862
3863 async def _get_liked_albums_cached(self, ttl: float = 30.0) -> list[YandexAlbum]:
3864 """
3865 Return liked albums with a short in-process TTL cache + lock.
3866
3867 Albums, podcasts and audiobooks are all derived from the same
3868 ``users/{uid}/likes/albums`` endpoint, so a full library sync would
3869 otherwise trigger three sequential (or concurrent) identical calls.
3870 The lock serializes refreshes so only one request hits the API when
3871 multiple library syncs start together.
3872 """
3873 async with self._liked_albums_lock:
3874 now = asyncio.get_running_loop().time()
3875 if self._liked_albums_cache is not None:
3876 cached_at, cached = self._liked_albums_cache
3877 if now - cached_at < ttl:
3878 return cached
3879 albums = await self.client.get_liked_albums(batch_size=TRACK_BATCH_SIZE)
3880 self._liked_albums_cache = (now, albums)
3881 return albums
3882
3883 def _get_provider_item_id(self, item: MediaItemType) -> str | None:
3884 """Get provider item ID from media item."""
3885 for mapping in item.provider_mappings:
3886 if mapping.provider_instance == self.instance_id:
3887 return mapping.item_id
3888 return item.item_id if item.provider == self.instance_id else None
3889
3890 async def _get_audiobook_stream_details(self, audiobook_id: str) -> StreamDetails:
3891 """
3892 Build StreamDetails for an audiobook as a chapter-concatenated CUSTOM stream.
3893
3894 Loads the album's tracks, uses the first chapter to establish the audio
3895 format, and stores the per-chapter track-IDs + durations in ``data`` so
3896 ``get_audio_stream`` can iterate them. ``can_seek=True`` so MA routes
3897 ``seek_position`` into ``get_audio_stream``, where the provider translates
3898 it into ``(start_chapter, in_chapter_offset)``. In-chapter precision
3899 requires a byte-seekable chapter codec (raw MP3); otherwise the chapter
3900 is restarted from its beginning.
3901 """
3902 album = await self.client.get_album_with_tracks(audiobook_id)
3903 if not album or not (album.volumes or []):
3904 raise MediaNotFoundError(f"Audiobook {audiobook_id} has no chapters")
3905
3906 chapter_ids, chapter_durations_ms = _extract_chapter_map_from_album(album)
3907 if not chapter_ids:
3908 raise MediaNotFoundError(f"Audiobook {audiobook_id} has no chapters")
3909
3910 self._audiobook_chapter_cache[audiobook_id] = (chapter_ids, chapter_durations_ms)
3911
3912 # Resolve first-chapter format so MA/ffmpeg know what it's decoding
3913 first = await self.streaming.get_stream_details(chapter_ids[0])
3914 total_duration = sum(chapter_durations_ms) // 1000
3915
3916 return StreamDetails(
3917 item_id=audiobook_id,
3918 provider=self.instance_id,
3919 media_type=MediaType.AUDIOBOOK,
3920 audio_format=first.audio_format,
3921 stream_type=StreamType.CUSTOM,
3922 duration=total_duration,
3923 data={
3924 "chapter_ids": chapter_ids,
3925 "chapter_durations_ms": chapter_durations_ms,
3926 },
3927 can_seek=True,
3928 allow_seek=True,
3929 )
3930
3931 def _resolve_audiobook_seek(
3932 self, chapter_durations_ms: list[int], seek_position: int, n_chapters: int
3933 ) -> tuple[int, int]:
3934 """Map an audiobook ``seek_position`` (seconds) to (start_idx, chapter_seek)."""
3935 if seek_position <= 0 or not chapter_durations_ms:
3936 return 0, 0
3937 accumulated_ms = 0
3938 seek_ms = seek_position * 1000
3939 for idx, dur_ms in enumerate(chapter_durations_ms):
3940 if accumulated_ms + dur_ms > seek_ms:
3941 return idx, (seek_ms - accumulated_ms) // 1000
3942 accumulated_ms += dur_ms
3943 # Seek past end â start at last chapter from 0
3944 return max(n_chapters - 1, 0), 0
3945
3946 async def _resolve_audiobook_chapter_map(
3947 self, audiobook_id: str
3948 ) -> tuple[list[str], list[int]]:
3949 """
3950 Return (chapter_track_ids, chapter_durations_ms) for an audiobook.
3951
3952 Served from an in-memory cache populated by ``_get_audiobook_stream_details``.
3953 On a miss (e.g. ``on_played`` fires before streaming has started), falls back
3954 to a fresh ``get_album_with_tracks`` call and refills the cache.
3955 """
3956 cached = self._audiobook_chapter_cache.get(audiobook_id)
3957 if cached is not None:
3958 return cached
3959 album = await self.client.get_album_with_tracks(audiobook_id)
3960 if not album or not (album.volumes or []):
3961 return [], []
3962 chapter_ids, chapter_durations_ms = _extract_chapter_map_from_album(album)
3963 self._audiobook_chapter_cache[audiobook_id] = (chapter_ids, chapter_durations_ms)
3964 return chapter_ids, chapter_durations_ms
3965
3966 async def _stream_audiobook_chapters(
3967 self, data: dict[str, Any], seek_position: int
3968 ) -> AsyncGenerator[bytes]:
3969 """
3970 Concatenate per-chapter streams of an audiobook.
3971
3972 Translates ``seek_position`` into (start_chapter, in_chapter_offset) and
3973 delegates each chapter to the per-track streaming path. In-chapter offset
3974 is only applied when the chapter codec is byte-seekable (``can_seek``);
3975 otherwise the chapter is restarted from its beginning. Tracks consecutive
3976 chapter failures and raises ``MediaNotFoundError`` once the threshold is
3977 exceeded, so playback never silently truncates.
3978 """
3979 chapter_ids: list[str] = list(data.get("chapter_ids") or [])
3980 chapter_durations_ms: list[int] = list(data.get("chapter_durations_ms") or [])
3981 if not chapter_ids:
3982 return
3983
3984 start_idx, chapter_seek = self._resolve_audiobook_seek(
3985 chapter_durations_ms, seek_position, len(chapter_ids)
3986 )
3987
3988 max_consecutive_failures = 3
3989 consecutive_failures = 0
3990 has_yielded_audio = False
3991 last_error: Exception | None = None
3992
3993 for idx in range(start_idx, len(chapter_ids)):
3994 chapter_id = chapter_ids[idx]
3995 requested_offset = chapter_seek if idx == start_idx else 0
3996 chapter_details: StreamDetails | None = None
3997 try:
3998 chapter_details = await self.streaming.get_stream_details(chapter_id)
3999 except asyncio.CancelledError:
4000 raise
4001 except Exception as err:
4002 last_error = err
4003 self.logger.warning(
4004 "Audiobook chapter %d (%s) stream-details failed: %s",
4005 idx + 1,
4006 chapter_id,
4007 err,
4008 )
4009
4010 if chapter_details is None:
4011 consecutive_failures += 1
4012 if consecutive_failures >= max_consecutive_failures:
4013 raise MediaNotFoundError(
4014 "Unable to stream audiobook: too many consecutive chapter failures"
4015 ) from last_error
4016 continue
4017
4018 # Apply the in-chapter offset only when the chapter codec supports
4019 # byte-offset seeking; otherwise restart the chapter from 0 to avoid
4020 # decoding garbled bytes from mid-file of a container format.
4021 offset = requested_offset if chapter_details.can_seek else 0
4022 chapter_had_audio = False
4023 try:
4024 async for chunk in self.streaming.get_audio_stream(chapter_details, offset):
4025 chapter_had_audio = True
4026 has_yielded_audio = True
4027 yield chunk
4028 except asyncio.CancelledError:
4029 raise
4030 except Exception as err:
4031 last_error = err
4032 self.logger.warning(
4033 "Audiobook chapter %d (%s) stream failed mid-play: %s",
4034 idx + 1,
4035 chapter_id,
4036 err,
4037 )
4038
4039 if chapter_had_audio:
4040 consecutive_failures = 0
4041 last_error = None
4042 else:
4043 consecutive_failures += 1
4044 if consecutive_failures >= max_consecutive_failures:
4045 raise MediaNotFoundError(
4046 "Unable to stream audiobook: too many consecutive chapter failures"
4047 ) from last_error
4048
4049 if not has_yielded_audio:
4050 raise MediaNotFoundError(
4051 "Unable to stream audiobook: no playable chapters found"
4052 ) from last_error
4053
4054 def _audiobook_progress_point(
4055 self,
4056 chapter_durations_ms: list[int],
4057 n_chapters: int,
4058 absolute_sec: int,
4059 ) -> tuple[int, int, int]:
4060 """
4061 Resolve an absolute book position into a play_audio-ready tuple.
4062
4063 Returns ``(chapter_idx, track_length_seconds, offset_seconds)``, applying
4064 two invariants Yandex cares about and that ``_resolve_audiobook_seek``
4065 alone doesn't guarantee:
4066
4067 - At/beyond end-of-book, map to end of the last chapter (not start),
4068 so Yandex's resume point doesn't rewind to the start of the final
4069 chapter on natural completion.
4070 - ``track_length_seconds`` is clamped to at least 1 and ``offset`` to
4071 ``[0, track_length_seconds]`` â a chapter with ``duration_ms=None``
4072 (coerced to 0 by the chapter-map builder) would otherwise send
4073 ``track_length_seconds=0`` and block progress from syncing.
4074 """
4075 absolute_sec = max(0, absolute_sec)
4076 total_duration_sec = sum(chapter_durations_ms) // 1000
4077 last_idx = max(n_chapters - 1, 0)
4078 if absolute_sec >= total_duration_sec > 0:
4079 idx = last_idx
4080 track_length_sec = max(1, chapter_durations_ms[idx] // 1000)
4081 offset = track_length_sec
4082 else:
4083 idx, offset_raw = self._resolve_audiobook_seek(
4084 chapter_durations_ms, absolute_sec, n_chapters
4085 )
4086 track_length_sec = max(1, chapter_durations_ms[idx] // 1000)
4087 offset = max(0, min(int(offset_raw), track_length_sec))
4088 return idx, track_length_sec, offset
4089
4090 async def _report_audiobook_progress(self, audiobook_id: str, position_sec: int) -> None:
4091 """
4092 Push current listening position of an audiobook to Yandex.
4093
4094 Resolves the playing chapter + offset from the cached chapter map, then
4095 calls play_audio so Yandex persists the position for cross-client resume.
4096
4097 Best-effort: any non-cancellation failure while resolving the chapter
4098 map (rate-limit, network blip, auth edge case bubbling out of
4099 ``_call_with_retry``) must never break pause/stop, so it is swallowed
4100 here in addition to the errors already absorbed inside
4101 ``api_client.play_audio``.
4102 """
4103 try:
4104 chapter_ids, chapter_durations_ms = await self._resolve_audiobook_chapter_map(
4105 audiobook_id
4106 )
4107 except asyncio.CancelledError:
4108 raise
4109 except Exception as err:
4110 self.logger.debug(
4111 "Skipping audiobook progress report for %s (chapter map resolution failed): %s",
4112 audiobook_id,
4113 err,
4114 )
4115 return
4116 if not chapter_ids:
4117 self.logger.debug(
4118 "Audiobook %s has no chapter map; skipping progress report", audiobook_id
4119 )
4120 return
4121 idx, track_length_sec, offset = self._audiobook_progress_point(
4122 chapter_durations_ms, len(chapter_ids), int(position_sec)
4123 )
4124 play_id = self._audiobook_play_ids.setdefault(audiobook_id, uuid.uuid4().hex)
4125 await self.client.play_audio(
4126 track_id=chapter_ids[idx],
4127 album_id=audiobook_id,
4128 play_id=play_id,
4129 track_length_seconds=track_length_sec,
4130 total_played_seconds=offset,
4131 end_position_seconds=offset,
4132 )
4133
4134 async def _report_audiobook_final(
4135 self, streamdetails: StreamDetails, data: dict[str, Any]
4136 ) -> None:
4137 """
4138 Send a closing play_audio for an audiobook stream.
4139
4140 Uses the streamdetails' own ``chapter_ids`` / ``chapter_durations_ms``
4141 (populated when the StreamDetails was created) to stay consistent with
4142 what was actually played, then clears the session play_id and drops
4143 the chapter-map cache entry so long-running instances can't grow the
4144 cache without bound as users play more audiobooks.
4145 """
4146 audiobook_id = streamdetails.item_id
4147 chapter_ids = data.get("chapter_ids") or []
4148 chapter_durations_ms = data.get("chapter_durations_ms") or []
4149 play_id = self._audiobook_play_ids.pop(audiobook_id, None) or uuid.uuid4().hex
4150 self._audiobook_chapter_cache.pop(audiobook_id, None)
4151 if not chapter_ids or not chapter_durations_ms:
4152 return
4153 absolute_sec = int(streamdetails.seek_position + (streamdetails.seconds_streamed or 0))
4154 idx, track_length_sec, offset = self._audiobook_progress_point(
4155 chapter_durations_ms, len(chapter_ids), absolute_sec
4156 )
4157 await self.client.play_audio(
4158 track_id=chapter_ids[idx],
4159 album_id=audiobook_id,
4160 play_id=play_id,
4161 track_length_seconds=track_length_sec,
4162 total_played_seconds=offset,
4163 end_position_seconds=offset,
4164 )
4165
4166 async def _rotating_row_tag_subtitle(self, category: str) -> str | None:
4167 """
4168 Return the current rotating tag label from cache without backend I/O.
4169
4170 :param category: Tag category, such as ``mood`` or ``activity``.
4171 :return: Localized tag label, or ``None`` while the cache is cold.
4172 """
4173 tags, _, found = await self.mass.cache.get_with_freshness(
4174 f"_get_valid_tags_for_category.{category}",
4175 provider=self.instance_id,
4176 include_expired=True,
4177 )
4178 if not found or not isinstance(tags, list):
4179 return None
4180 valid_tags = [tag for tag in tags if isinstance(tag, str)]
4181 if not valid_tags:
4182 return None
4183 tag = self._rotating_row_tag(category, valid_tags)
4184 return self._media_label("folder", _media_label_key(tag), tag.title())[0]
4185
4186 def _rotating_row_tag(self, category: str, valid_tags: list[str]) -> str:
4187 """
4188 Deterministically select a tag for this provider and UTC hour.
4189
4190 :param category: Tag category the values belong to.
4191 :param valid_tags: Non-empty ordered tag list.
4192 :return: The selected tag slug.
4193 """
4194 hour_bucket = int(utc().timestamp()) // 3600
4195 seed = f"{self.instance_id}.{category}.{hour_bucket}".encode()
4196 index = int.from_bytes(hashlib.sha256(seed).digest()[:8], "big") % len(valid_tags)
4197 return valid_tags[index]
4198