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