/
/
1"""Pandora music provider for Music Assistant."""
2
3from __future__ import annotations
4
5import json
6import time
7from typing import TYPE_CHECKING, Any
8
9import aiohttp
10from aiohttp import web
11from music_assistant_models.config_entries import (
12 ConfigActionResult,
13 ConfigEntry,
14 ConfigValueOption,
15)
16from music_assistant_models.enums import (
17 ConfigEntryType,
18 ContentType,
19 ImageType,
20 MediaType,
21 StreamType,
22)
23from music_assistant_models.errors import (
24 InvalidDataError,
25 LoginFailed,
26 MediaNotFoundError,
27 ProviderUnavailableError,
28)
29from music_assistant_models.media_items import (
30 AudioFormat,
31 MediaItemImage,
32 MediaItemMetadata,
33 ProviderMapping,
34 Radio,
35 SearchResults,
36 UniqueList,
37)
38from music_assistant_models.streamdetails import MultiPartPath, StreamDetails, StreamMetadata
39
40from music_assistant.constants import (
41 CONF_ENTRY_UNOFFICIAL_PROVIDER,
42 CONF_PASSWORD,
43 CONF_SOCKS_URL,
44 CONF_USERNAME,
45)
46from music_assistant.controllers.cache import use_cache
47from music_assistant.helpers.aiohttp_client import create_clientsession, get_socks5_url
48from music_assistant.helpers.compare import compare_strings
49from music_assistant.models.music_provider import MusicProvider
50
51from .constants import (
52 ACCOUNT_FLAG_HIGH_QUALITY,
53 CONF_QUALITY,
54 CONF_TAKEOVER_ACTION,
55 LOGIN_ENDPOINT,
56 PLAYBACK_RESUMED_ENDPOINT,
57 PLAYLIST_FRAGMENT_ENDPOINT,
58 QUALITY_HIGH,
59 QUALITY_STANDARD,
60 RETRY_REASON_AUTH,
61 RETRY_REASON_STREAM_VIOLATION,
62 STATIONS_ENDPOINT,
63)
64from .helpers import create_auth_headers, get_csrf_token, handle_pandora_error
65
66if TYPE_CHECKING:
67 from collections.abc import AsyncGenerator
68
69
70class PandoraStationSession:
71 """Manages streaming state for a single Pandora station."""
72
73 def __init__(self, station_id: str):
74 """
75 Initialize a new station streaming session.
76
77 Args:
78 station_id: The Pandora station ID.
79 """
80 self.station_id = station_id
81 self.fragments: list[dict[str, Any] | None] = []
82 self.track_map: list[tuple[int, int]] = []
83 self.cumulative_times: list[int] = []
84 self.last_accessed = time.time()
85
86 def get_track_duration(self, music_track_num: int) -> int:
87 """Calculate duration for a specific track index."""
88 if not (0 <= music_track_num < len(self.track_map)):
89 return 0
90 frag_idx, track_idx = self.track_map[music_track_num]
91 if frag_idx >= len(self.fragments) or not (frag := self.fragments[frag_idx]):
92 return 0
93 tracks = frag.get("tracks", [])
94 if track_idx >= len(tracks):
95 return 0
96 return int(tracks[track_idx].get("trackLength", 0))
97
98
99class StreamViolationError(InvalidDataError):
100 """Error raised when Pandora detects concurrent streaming on multiple devices."""
101
102
103class PandoraProvider(MusicProvider):
104 """Pandora Music Provider."""
105
106 _auth_token: str | None = None
107 _user_id: str | None = None
108 _csrf_token: str | None = None
109 _sessions: dict[str, PandoraStationSession]
110 _socks_proxy: bool = False
111 _high_quality_available: bool = False
112
113 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
114 """Return Config entries to configure this provider."""
115 return (
116 CONF_ENTRY_UNOFFICIAL_PROVIDER,
117 ConfigEntry(
118 key=CONF_QUALITY,
119 type=ConfigEntryType.STRING,
120 required=True,
121 default_value=QUALITY_STANDARD,
122 options=[
123 ConfigValueOption(QUALITY_STANDARD),
124 ConfigValueOption(QUALITY_HIGH),
125 ],
126 ),
127 ConfigEntry(
128 key=CONF_SOCKS_URL,
129 type=ConfigEntryType.STRING,
130 required=False,
131 default_value="",
132 advanced=True,
133 ),
134 ConfigEntry(
135 key=CONF_TAKEOVER_ACTION,
136 type=ConfigEntryType.ACTION,
137 action=CONF_TAKEOVER_ACTION,
138 required=False,
139 ),
140 )
141
142 async def handle_config_action(
143 self, action: str
144 ) -> tuple[ConfigEntry, ...] | ConfigActionResult | None:
145 """Handle a one-shot config action button press."""
146 if action == CONF_TAKEOVER_ACTION:
147 await self.takeover_stream()
148 return None
149 return await super().handle_config_action(action)
150
151 async def handle_async_init(self) -> None:
152 """Handle async initialization of the provider."""
153 self._on_unload_callbacks = []
154 self._sessions = {}
155
156 # Authenticate with Pandora
157 username = str(self.get_setup_value(CONF_USERNAME) or "")
158 password = str(self.get_setup_value(CONF_PASSWORD) or "")
159 if not username.strip() or not password.strip():
160 raise LoginFailed("Username and password are required")
161 socks_url = get_socks5_url(str(self.config.get_value(CONF_SOCKS_URL)))
162
163 if socks_url:
164 self.http_session = create_clientsession(
165 self.mass, verify_ssl=True, socks_url=socks_url
166 )
167 self._socks_proxy = True
168 else:
169 self.http_session = self.mass.http_session
170 await self._authenticate(username, password)
171
172 # Register dynamic stream route
173 self._on_unload_callbacks.append(
174 self.mass.streams.register_dynamic_route(
175 f"/{self.instance_id}_stream", self._handle_stream_request
176 )
177 )
178
179 async def unload(self, is_removed: bool = False) -> None:
180 """Handle unload/close of the provider."""
181 for callback in getattr(self, "_on_unload_callbacks", []):
182 callback()
183 await self.close()
184 await super().unload(is_removed)
185
186 async def _authenticate(self, username: str, password: str) -> None:
187 """Authenticate with Pandora and get auth token."""
188 try:
189 self._csrf_token = await get_csrf_token(self.http_session)
190
191 login_data = {
192 "username": username,
193 "password": password,
194 "keepLoggedIn": True,
195 "existingAuthToken": None,
196 }
197
198 headers = create_auth_headers(self._csrf_token)
199
200 async with self.http_session.post(
201 LOGIN_ENDPOINT,
202 headers=headers,
203 json=login_data,
204 timeout=aiohttp.ClientTimeout(total=30),
205 ) as response:
206 if response.status != 200:
207 await self.close()
208 raise LoginFailed(f"Login request failed with status {response.status}")
209
210 response_data = await response.json()
211 handle_pandora_error(response_data)
212
213 self._auth_token = response_data.get("authToken")
214 if not self._auth_token:
215 await self.close()
216 raise LoginFailed("No auth token received from Pandora")
217
218 self._user_id = response_data.get("listenerId")
219
220 # Check whether the account is eligible for high-quality streaming.
221 try:
222 flags: list[str] = response_data.get("config", {}).get("flags", [])
223 self._high_quality_available = ACCOUNT_FLAG_HIGH_QUALITY in flags
224 except AttributeError, TypeError:
225 self._high_quality_available = False
226
227 self.logger.info(
228 "Successfully authenticated with Pandora "
229 "(high-quality streaming available: %s)",
230 self._high_quality_available,
231 )
232
233 except aiohttp.ClientError as err:
234 await self.close()
235 self.logger.exception("Network error during authentication")
236 raise ProviderUnavailableError(
237 "Unable to connect to Pandora for authentication"
238 ) from err
239
240 async def _api_request(
241 self,
242 method: str,
243 url: str,
244 data: dict[str, Any] | None = None,
245 exhausted_retry_reasons: frozenset[str] = frozenset(),
246 ) -> dict[str, Any]:
247 """
248 Make an API request to Pandora.
249
250 :param method: HTTP method (GET, POST, etc.)
251 :param url: API endpoint URL
252 :param data: Optional JSON data to send
253 :param exhausted_retry_reasons: Set of retry reasons already attempted for this request.
254 Pass a pre-populated set to prevent specific retry strategies from being attempted.
255 """
256 if not self._csrf_token or not self._auth_token:
257 await self.close()
258 raise LoginFailed("Not authenticated with Pandora")
259
260 headers = create_auth_headers(self._csrf_token, self._auth_token)
261
262 try:
263 async with self.http_session.request(
264 method, url, json=data, headers=headers
265 ) as response:
266 # Check status BEFORE parsing JSON
267 if response.status == 401:
268 if RETRY_REASON_AUTH not in exhausted_retry_reasons:
269 # Auth token expired, re-authenticate and retry once
270 username = str(self.get_setup_value(CONF_USERNAME) or "")
271 password = str(self.get_setup_value(CONF_PASSWORD) or "")
272 await self._authenticate(username, password)
273 return await self._api_request(
274 method,
275 url,
276 data,
277 exhausted_retry_reasons=exhausted_retry_reasons | {RETRY_REASON_AUTH},
278 )
279 await self.close()
280 raise LoginFailed("Pandora authentication failed after retry")
281 if response.status == 404:
282 await self.close()
283 raise MediaNotFoundError("Resource not found")
284 if response.status == 429:
285 # Another device may already be streaming on this account.
286 # Parse the body to confirm it is a STREAM_VIOLATION.
287 try:
288 error_body: dict[str, Any] = await response.json()
289 except (aiohttp.ContentTypeError, json.JSONDecodeError) as err:
290 raise InvalidDataError(
291 "Unable to parse error 429 response body from Pandora"
292 ) from err
293 if error_body.get("errorString") == "STREAM_VIOLATION":
294 if RETRY_REASON_STREAM_VIOLATION not in exhausted_retry_reasons:
295 self.logger.warning(
296 "Pandora stream is already active on another device. "
297 "Automatically taking over the stream and retrying the request."
298 )
299 await self.takeover_stream()
300 return await self._api_request(
301 method,
302 url,
303 data,
304 exhausted_retry_reasons=exhausted_retry_reasons
305 | {RETRY_REASON_STREAM_VIOLATION},
306 )
307 raise StreamViolationError("STREAM_VIOLATION")
308 # This is some other, not concurrent streaming error kind of 429
309 raise ProviderUnavailableError(f"Pandora rate-limited (HTTP 429): {error_body}")
310 if response.status >= 500:
311 await self.close()
312 raise ProviderUnavailableError("Pandora server error")
313 if response.status >= 400:
314 await self.close()
315 raise InvalidDataError(f"Pandora API error: HTTP {response.status}")
316
317 result: dict[str, Any] = await response.json()
318 handle_pandora_error(result)
319 return result
320
321 except aiohttp.ClientError as err:
322 await self.close()
323 raise ProviderUnavailableError("Unable to connect to Pandora") from err
324 except (ValueError, KeyError) as err:
325 await self.close()
326 raise InvalidDataError("Invalid response from Pandora") from err
327
328 @use_cache(3600)
329 async def get_radio(self, prov_radio_id: str) -> Radio:
330 """Get single radio station details."""
331 return Radio(
332 item_id=prov_radio_id,
333 provider=self.instance_id,
334 name=f"Pandora Station {prov_radio_id}",
335 translation_key="pandora_station",
336 translation_params=[prov_radio_id],
337 provider_mappings={
338 ProviderMapping(
339 item_id=prov_radio_id,
340 provider_domain=self.domain,
341 provider_instance=self.instance_id,
342 )
343 },
344 )
345
346 async def get_library_radios(self) -> AsyncGenerator[Radio]:
347 """Retrieve library/subscribed radio stations from the provider."""
348 response = await self._api_request(
349 "POST",
350 STATIONS_ENDPOINT,
351 data={
352 "pageSize": 250,
353 },
354 )
355
356 stations = response.get("stations", [])
357 self.logger.debug("Retrieved %d stations from Pandora", len(stations))
358
359 for station in stations:
360 station_image = None
361 if art := station.get("art"):
362 art_url = next(
363 (item["url"] for item in art if item.get("size") == 500),
364 art[-1]["url"] if art else None,
365 )
366 if art_url:
367 station_image = MediaItemImage(
368 type=ImageType.THUMB,
369 path=art_url,
370 provider=self.instance_id,
371 remotely_accessible=True,
372 )
373 yield Radio(
374 item_id=station["stationId"],
375 provider=self.instance_id,
376 name=station["name"],
377 metadata=MediaItemMetadata(
378 images=UniqueList([station_image]) if station_image else None,
379 ),
380 provider_mappings={
381 ProviderMapping(
382 item_id=station["stationId"],
383 provider_domain=self.domain,
384 provider_instance=self.instance_id,
385 )
386 },
387 )
388
389 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
390 """Get streamdetails for a radio station."""
391 if media_type != MediaType.RADIO:
392 await self.close()
393 raise MediaNotFoundError(f"Unsupported media type: {media_type}")
394
395 # Clear any existing session so we get fresh tracks from the API
396 self._sessions.pop(item_id, None)
397
398 # Create playlist with 1000 track placeholders for continuous streaming
399 parts = [
400 MultiPartPath(
401 path=f"{self.mass.streams.base_url}/{self.instance_id}_stream?"
402 f"station_id={item_id}&track_num={i}"
403 )
404 for i in range(1000)
405 ]
406 return StreamDetails(
407 provider=self.instance_id,
408 item_id=item_id,
409 audio_format=AudioFormat(
410 content_type=ContentType.MP3 if self._use_high_quality() else ContentType.AAC,
411 ),
412 media_type=MediaType.RADIO,
413 stream_type=StreamType.HTTP,
414 path=parts,
415 can_seek=False,
416 allow_seek=False,
417 stream_metadata=StreamMetadata(
418 title="Pandora Radio",
419 ),
420 stream_metadata_update_callback=self._update_stream_metadata,
421 stream_metadata_update_interval=5, # Check every 5 seconds
422 )
423
424 async def _get_fragment_data(
425 self, session: PandoraStationSession, fragment_index: int
426 ) -> dict[str, Any]:
427 """Fetch fragment data from Pandora API."""
428 # Check if already cached in session
429 if fragment_index < len(session.fragments):
430 cached = session.fragments[fragment_index]
431 if cached is not None:
432 return cached
433
434 is_stream_start = fragment_index == 0
435
436 fragment_data = {
437 "stationId": session.station_id,
438 "isStationStart": is_stream_start,
439 "fragmentRequestReason": "Normal",
440 "audioFormat": "mp3-hifi" if self._use_high_quality() else "aacplus",
441 "startingAtTrackId": None,
442 "onDemandArtistMessageArtistUidHex": None,
443 "onDemandArtistMessageIdHex": None,
444 }
445
446 try:
447 result: dict[str, Any] = await self._api_request(
448 "POST",
449 PLAYLIST_FRAGMENT_ENDPOINT,
450 data=fragment_data,
451 # Mark stream violation retry as already exhausted for non-initial fragments
452 # this prevents us from fighting with the concurrent streaming limit
453 # if the user starts a stream on a different device while MA is already playing.
454 exhausted_retry_reasons=frozenset()
455 if is_stream_start
456 else frozenset({RETRY_REASON_STREAM_VIOLATION}),
457 )
458
459 # Store in session cache
460 while len(session.fragments) <= fragment_index:
461 session.fragments.append(None)
462 session.fragments[fragment_index] = result
463
464 tracks = result.get("tracks", [])
465
466 # Calculate starting cumulative time for this fragment
467 if session.cumulative_times:
468 # Get the last music track's end time
469 last_music_track_num = len(session.track_map) - 1
470 last_start = session.cumulative_times[-1]
471 last_duration = session.get_track_duration(last_music_track_num)
472 current_cumulative = last_start + last_duration
473 else:
474 current_cumulative = 0
475
476 for track_idx, track in enumerate(tracks):
477 title = track.get("songTitle", "")
478 # Skip curator messages from the mapping
479 if "Curator Message" not in title and "curator message" not in title.lower():
480 session.track_map.append((fragment_index, track_idx))
481 session.cumulative_times.append(current_cumulative)
482
483 duration = track.get("trackLength", 0)
484 current_cumulative += duration
485
486 return result
487
488 except MediaNotFoundError:
489 await self.close()
490 raise
491 except StreamViolationError:
492 self.logger.warning(
493 "Pandora stream is already active on another device. "
494 "To manually take over the stream on this device, use the "
495 "'Take over stream' button on the provider configuration page.",
496 )
497 raise
498 except InvalidDataError as err:
499 self.logger.error("Invalid fragment data for station %s: %s", session.station_id, err)
500 await self.close()
501 raise
502
503 async def _handle_stream_request(self, request: web.Request) -> web.Response:
504 """
505 Handle dynamic stream request.
506
507 Map track numbers to Pandora fragments and redirect to audio URLs.
508 """
509 if not (station_id := request.query.get("station_id")):
510 return web.Response(status=400, text="Missing station_id")
511 if not (track_num_str := request.query.get("track_num")):
512 return web.Response(status=400, text="Missing track_num")
513
514 try:
515 music_track_num = int(track_num_str)
516 except ValueError:
517 return web.Response(status=400, text="Invalid track_num")
518
519 # Get or create session with LRU eviction
520 session = self._get_or_create_session(station_id)
521
522 try:
523 # If we don't have this music track yet, fetch more fragments
524 while music_track_num >= len(session.track_map):
525 next_fragment_idx = len(session.fragments)
526 await self._get_fragment_data(session, next_fragment_idx)
527
528 # Look up the actual fragment/track position
529 fragment_idx, track_idx = session.track_map[music_track_num]
530
531 # Ensure fragment is loaded
532 if fragment_idx >= len(session.fragments) or not session.fragments[fragment_idx]:
533 await self._get_fragment_data(session, fragment_idx)
534
535 fragment = session.fragments[fragment_idx]
536 if not fragment:
537 return web.Response(status=404, text="Track unavailable")
538
539 # Get the track
540 tracks = fragment.get("tracks", [])
541 if track_idx >= len(tracks):
542 self.logger.error(
543 "Track index %d out of range (fragment has %d tracks)",
544 track_idx,
545 len(tracks),
546 )
547 return web.Response(status=404, text="Track unavailable")
548
549 track = tracks[track_idx]
550 audio_url = track.get("audioURL")
551
552 if not audio_url:
553 self.logger.error("No audio URL in track data")
554 return web.Response(status=404, text="Track unavailable")
555
556 # Redirect to the actual audio URL
557 return web.Response(status=302, headers={"Location": audio_url})
558
559 except (MediaNotFoundError, InvalidDataError) as err:
560 self.logger.error("Stream error: %s", err)
561 return web.Response(status=404, text="Stream unavailable")
562 except ProviderUnavailableError as err:
563 self.logger.error("Pandora service unavailable: %s", err)
564 return web.Response(status=503, text="Service temporarily unavailable")
565
566 def _get_or_create_session(self, station_id: str) -> PandoraStationSession:
567 """Get or create a session, with LRU eviction if needed."""
568 # Simple LRU: limit to 10 active sessions
569 if station_id not in self._sessions and len(self._sessions) >= 10:
570 # Remove oldest session
571 oldest = min(self._sessions.values(), key=lambda s: s.last_accessed)
572 self.logger.debug("Evicting session for station %s", oldest.station_id)
573 del self._sessions[oldest.station_id]
574
575 if station_id not in self._sessions:
576 self._sessions[station_id] = PandoraStationSession(station_id)
577
578 session = self._sessions[station_id]
579 session.last_accessed = time.time()
580 return session
581
582 async def search(
583 self,
584 search_query: str,
585 media_types: list[MediaType],
586 limit: int = 25,
587 ) -> SearchResults:
588 """Search library radio stations by name."""
589 # Search limited to library stations (API search requires legacy endpoints)
590 if MediaType.RADIO not in media_types:
591 return SearchResults()
592
593 results: list[Radio] = []
594
595 async for station in self.get_library_radios():
596 if compare_strings(station.name, search_query):
597 results.append(station)
598 if len(results) >= limit:
599 break
600
601 return SearchResults(radio=results)
602
603 async def _update_stream_metadata(
604 self, streamdetails: StreamDetails, elapsed_time: int
605 ) -> None:
606 """Update stream metadata based on elapsed playback time."""
607 station_id = streamdetails.item_id
608
609 # Get session if it exists
610 if station_id not in self._sessions:
611 return
612
613 session = self._sessions[station_id]
614 session.last_accessed = time.time()
615
616 if not session.track_map or not session.cumulative_times:
617 return
618
619 # Find the current track based on elapsed time
620 current_track_idx = None
621 for i, start_time in enumerate(session.cumulative_times):
622 # Calculate when this track ends
623 if i + 1 < len(session.cumulative_times):
624 end_time = session.cumulative_times[i + 1]
625 else:
626 end_time = start_time + session.get_track_duration(i)
627
628 if start_time <= elapsed_time < end_time:
629 current_track_idx = i
630 break
631
632 if current_track_idx is None:
633 return
634
635 # Get track data
636 frag_idx, track_idx = session.track_map[current_track_idx]
637 if frag_idx >= len(session.fragments):
638 return
639 fragment = session.fragments[frag_idx]
640 if not fragment:
641 return
642
643 tracks = fragment.get("tracks", [])
644 if track_idx >= len(tracks):
645 return
646
647 track = tracks[track_idx]
648
649 # Update metadata if title changed
650 if not streamdetails.stream_metadata or streamdetails.stream_metadata.title == track.get(
651 "songTitle"
652 ):
653 return
654
655 # Get album art
656 album_art_url = None
657 if album_art := track.get("albumArt"):
658 album_art_url = next(
659 (art["url"] for art in album_art if art.get("size") == 500),
660 album_art[-1]["url"] if album_art else None,
661 )
662
663 streamdetails.stream_metadata.title = track.get("songTitle", "Unknown Song")
664 streamdetails.stream_metadata.artist = track.get("artistName", "Unknown Artist")
665 streamdetails.stream_metadata.album = track.get("albumTitle")
666 streamdetails.stream_metadata.image_url = album_art_url
667 streamdetails.stream_metadata.duration = track.get("trackLength")
668 streamdetails.stream_metadata.uri = track.get("songDetailURL")
669
670 async def close(self) -> None:
671 """Handle closing of http session if using socks."""
672 if self._socks_proxy and self.http_session:
673 await self.http_session.close()
674
675 def _use_high_quality(self) -> bool:
676 """
677 Whether high quality audio should be requested from Pandora.
678
679 This allows a graceful fallback to standard quality if the account is not eligible for
680 high-quality streaming, while still respecting the user's preference if they are eligible.
681 """
682 return self._high_quality_available and self.config.get_value(CONF_QUALITY) == QUALITY_HIGH
683
684 async def takeover_stream(self) -> None:
685 """
686 Force Pandora to end any other active session and resume here.
687
688 This sends "forceActive=true" to the playbackResumed endpoint, which instructs Pandora to
689 terminate any conflicting stream on other devices. The user must manually restart playback
690 in MA after clicking the config button that triggers this call.
691 """
692 self.logger.debug("Sending playbackResumed request to Pandora to attempt stream takeover.")
693 await self._api_request(
694 "POST",
695 PLAYBACK_RESUMED_ENDPOINT,
696 data={"forceActive": True},
697 # This is called as part of handling a STREAM_VIOLATION 429, so mark that reason as
698 # already exhausted to prevent _api_request from retrying on another 429.
699 exhausted_retry_reasons=frozenset({RETRY_REASON_STREAM_VIOLATION}),
700 )
701