/
/
1"""
2Overcast provider for Music Assistant.
3
4Imports podcast subscriptions and playback progress from an Overcast account
5using the account's extended OPML export. Synchronization is strictly one-way
6(Overcast -> Music Assistant): Overcast has no ingest API, so ``on_played`` is
7deliberately not implemented and no state is ever written back.
8
9The OPML export endpoint is rate limited by Overcast, therefore the export is
10cached aggressively and refreshed in the background, while the RSS feeds it
11refers to are fetched directly from the podcasts' own servers.
12"""
13
14from __future__ import annotations
15
16import asyncio
17import math
18import time
19from typing import TYPE_CHECKING, Any, cast
20
21import aiohttp
22from music_assistant_models.config_entries import ConfigEntry
23from music_assistant_models.enums import (
24 ConfigEntryType,
25 ContentType,
26 MediaType,
27 StreamType,
28)
29from music_assistant_models.errors import (
30 LoginFailed,
31 MediaNotFoundError,
32 ResourceTemporarilyUnavailable,
33)
34from music_assistant_models.media_items import AudioFormat, Podcast, PodcastEpisode
35from music_assistant_models.streamdetails import StreamDetails
36from yarl import URL
37
38from music_assistant.constants import CONF_ENTRY_UNOFFICIAL_PROVIDER, CONF_PASSWORD, CONF_USERNAME
39from music_assistant.controllers.cache import use_cache
40from music_assistant.helpers.aiohttp_client import create_clientsession
41from music_assistant.helpers.datetime import from_iso_string
42from music_assistant.helpers.podcast_parsers import (
43 enrich_episode_chapters,
44 find_episode_stream_url,
45 get_cached_podcast,
46 get_stream_url_and_guid_from_episode,
47 parse_podcast,
48 parse_podcast_episode,
49 refresh_cached_podcast,
50)
51from music_assistant.helpers.throttle_retry import parse_retry_after
52from music_assistant.models.music_provider import MusicProvider
53
54from .constants import (
55 AUTH_REJECT_STATUSES,
56 BASE_URL,
57 CACHE_CATEGORY_OPML,
58 CACHE_KEY_LAST_APPLIED,
59 CONF_MAX_NUM_EPISODES,
60 CONF_SESSION_COOKIE,
61 LOGIN_URL,
62 OPML_CACHE_EXPIRATION,
63 OPML_EXPORT_URL,
64 PODCASTS_URL,
65 RATE_LIMIT_FALLBACK_BACKOFF,
66 SESSION_COOKIE_NAME,
67)
68from .helpers import OvercastSubscription, match_episode_state, parse_extended_opml
69
70if TYPE_CHECKING:
71 from collections.abc import AsyncGenerator
72 from datetime import datetime
73
74
75class OvercastProvider(MusicProvider):
76 """Provider that imports podcast subscriptions from an Overcast account."""
77
78 http_session: aiohttp.ClientSession
79 max_episodes: int
80 # newest playback state applied per feed url: a feed that could not be retrieved
81 # keeps its own watermark, so its states are still applied once it recovers
82 _feed_watermarks: dict[str, datetime]
83 # the parsed export, kept alongside the raw text it was parsed from so a
84 # refreshed export is re-parsed while repeated lookups are not
85 _opml_cache: tuple[str, dict[str, OvercastSubscription]] | None = None
86 # monotonic deadline set from a 429's Retry-After, see _rate_limit_remaining
87 _rate_limited_until: float | None = None
88
89 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
90 """
91 Return the (options) config entries for the Overcast provider.
92
93 The account credentials are collected by the interactive setup flow
94 (see ``setup_flow.py``); only the max-episodes limit is configured here.
95 """
96 return (
97 CONF_ENTRY_UNOFFICIAL_PROVIDER,
98 ConfigEntry(
99 key=CONF_MAX_NUM_EPISODES,
100 type=ConfigEntryType.INTEGER,
101 required=False,
102 default_value=0,
103 ),
104 )
105
106 async def handle_async_init(self) -> None:
107 """Handle async initialization of the provider."""
108 self.max_episodes = cast("int", self.config.get_value(CONF_MAX_NUM_EPISODES, 0))
109 # Dedicated session with its own cookie jar to support multi-instance
110 # (each instance has its own Overcast login cookie)
111 self.http_session = create_clientsession(self.mass, cookie_jar=aiohttp.CookieJar())
112 stored_cookie = self.get_setup_value(CONF_SESSION_COOKIE)
113 cookie_valid = False
114 if isinstance(stored_cookie, str) and stored_cookie:
115 self.http_session.cookie_jar.update_cookies(
116 {SESSION_COOKIE_NAME: stored_cookie}, response_url=URL(BASE_URL)
117 )
118 cookie_valid = await self._session_cookie_valid()
119 if not cookie_valid:
120 await self._login()
121
122 raw_watermarks = await self.mass.cache.get(
123 key=CACHE_KEY_LAST_APPLIED,
124 provider=self.instance_id,
125 category=CACHE_CATEGORY_OPML,
126 default={},
127 )
128 self._feed_watermarks = {
129 feed_url: from_iso_string(raw) for feed_url, raw in raw_watermarks.items()
130 }
131
132 async def unload(self, is_removed: bool = False) -> None:
133 """Handle unload/close of the provider."""
134 if not self.http_session.closed:
135 await self.http_session.close()
136
137 @property
138 def is_streaming_provider(self) -> bool:
139 """Return False: the library mirrors the user's own Overcast subscriptions."""
140 return False
141
142 async def get_library_podcasts(self) -> AsyncGenerator[Podcast]:
143 """Retrieve the subscribed podcasts from the Overcast account."""
144 subscriptions = await self._get_opml_subscriptions()
145 for feed_url, subscription in subscriptions.items():
146 self.logger.debug("Adding podcast with feed %s to library", feed_url)
147 try:
148 parsed_podcast = await refresh_cached_podcast(
149 mass=self.mass,
150 provider_instance_id=self.instance_id,
151 feed_url=feed_url,
152 max_episodes=self.max_episodes,
153 )
154 except MediaNotFoundError as err:
155 self.report_skipped_sync_item(MediaType.PODCAST, feed_url, err)
156 continue
157 applied = await self._apply_playback_states(feed_url, subscription, parsed_podcast)
158 if applied is not None:
159 # stored right away: the caller may stop consuming this generator at
160 # any point, which would otherwise re-apply these states on the next sync
161 self._feed_watermarks[feed_url] = applied
162 await self._store_watermarks()
163 yield parse_podcast(
164 feed_url=feed_url,
165 parsed_feed=parsed_podcast,
166 instance_id=self.instance_id,
167 domain=self.domain,
168 )
169
170 async def get_podcast(self, prov_podcast_id: str) -> Podcast:
171 """Get the podcast for the given feed url."""
172 parsed_podcast = await self._cache_get_podcast(prov_podcast_id)
173 return parse_podcast(
174 feed_url=prov_podcast_id,
175 parsed_feed=parsed_podcast,
176 instance_id=self.instance_id,
177 domain=self.domain,
178 )
179
180 async def get_podcast_episodes(self, prov_podcast_id: str) -> AsyncGenerator[PodcastEpisode]:
181 """Get all episodes of a podcast, including their Overcast playback state."""
182 podcast = await self._cache_get_podcast(prov_podcast_id)
183 subscription = await self._get_subscription(prov_podcast_id)
184 podcast_cover = podcast.get("cover_url")
185 podcast_name = podcast.get("title")
186 for cnt, parsed_episode in enumerate(podcast.get("episodes", [])):
187 mass_episode = parse_podcast_episode(
188 episode=parsed_episode,
189 prov_podcast_id=prov_podcast_id,
190 episode_cnt=cnt,
191 podcast_cover=podcast_cover,
192 podcast_name=podcast_name,
193 instance_id=self.instance_id,
194 domain=self.domain,
195 )
196 if mass_episode is None:
197 # faulty episode
198 continue
199 try:
200 stream_url, _ = get_stream_url_and_guid_from_episode(episode=parsed_episode)
201 except ValueError:
202 # episode enclosure or stream url missing
203 continue
204 if subscription is not None:
205 state = match_episode_state(subscription, stream_url)
206 if state is not None and (state.played or state.progress_s):
207 mass_episode.resume_position_ms = (state.progress_s or 0) * 1000
208 mass_episode.fully_played = state.played
209 yield mass_episode
210
211 async def get_podcast_episode(self, prov_episode_id: str) -> PodcastEpisode:
212 """Get a single podcast episode."""
213 podcast_id, guid_or_stream_url = prov_episode_id.split(" ", 1)
214 async for mass_episode in self.get_podcast_episodes(podcast_id):
215 _, episode_key = mass_episode.item_id.split(" ", 1)
216 if episode_key == guid_or_stream_url:
217 await self._enrich_episode_chapters(podcast_id, guid_or_stream_url, mass_episode)
218 return mass_episode
219 raise MediaNotFoundError("Did not find episode.")
220
221 async def get_resume_position(
222 self, item_id: str, media_type: MediaType
223 ) -> tuple[bool, int, datetime | None]:
224 """Return fully_played, resume position (ms) and its timestamp from Overcast."""
225 if media_type != MediaType.PODCAST_EPISODE:
226 raise NotImplementedError
227 podcast_id, guid_or_stream_url = item_id.split(" ", 1)
228 stream_url = await self._get_episode_stream_url(podcast_id, guid_or_stream_url)
229 if stream_url is None:
230 raise NotImplementedError
231 subscription = await self._get_subscription(podcast_id)
232 state = match_episode_state(subscription, stream_url) if subscription else None
233 if state is None or (not state.played and state.progress_s is None):
234 # No known Overcast progress; raise NotImplementedError such that MA
235 # falls back to the resume position stored in its own playlog.
236 raise NotImplementedError
237 return state.played, max((state.progress_s or 0) * 1000, 0), state.user_updated_at
238
239 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
240 """Get streamdetails for item."""
241 podcast_id, guid_or_stream_url = item_id.split(" ", 1)
242 stream_url = await self._get_episode_stream_url(podcast_id, guid_or_stream_url)
243 if stream_url is None:
244 raise MediaNotFoundError
245 return StreamDetails(
246 provider=self.instance_id,
247 item_id=item_id,
248 audio_format=AudioFormat(
249 content_type=ContentType.try_parse(stream_url),
250 ),
251 media_type=MediaType.PODCAST_EPISODE,
252 stream_type=StreamType.HTTP,
253 path=stream_url,
254 can_seek=True,
255 allow_seek=True,
256 )
257
258 async def _login(self) -> None:
259 """Authenticate with Overcast and persist the session cookie."""
260 email = str(self.get_setup_value(CONF_USERNAME))
261 password = str(self.get_setup_value(CONF_PASSWORD))
262 try:
263 async with self.http_session.post(
264 LOGIN_URL,
265 data={"email": email, "password": password},
266 allow_redirects=False,
267 ) as response:
268 status = response.status
269 location = response.headers.get("Location", "")
270 morsel = response.cookies.get(SESSION_COOKIE_NAME)
271 except (TimeoutError, aiohttp.ClientError) as err:
272 raise ResourceTemporarilyUnavailable("Overcast is unreachable") from err
273 if status != 302 or "/podcasts" not in location or morsel is None:
274 raise LoginFailed("Overcast login failed, check your email and password")
275 self._update_setup_data(CONF_SESSION_COOKIE, morsel.value)
276
277 async def _session_cookie_valid(self) -> bool:
278 """Check whether the restored session cookie is still accepted by Overcast."""
279 try:
280 async with self.http_session.get(PODCASTS_URL, allow_redirects=False) as response:
281 return response.status == 200
282 except (TimeoutError, aiohttp.ClientError) as err:
283 raise ResourceTemporarilyUnavailable("Overcast is unreachable") from err
284
285 @use_cache(
286 expiration=OPML_CACHE_EXPIRATION,
287 category=CACHE_CATEGORY_OPML,
288 allow_expired_cache=True,
289 cache_none=False,
290 )
291 async def _fetch_opml_text(self) -> str:
292 """Fetch the account's extended OPML export."""
293 opml_text = await self._request_opml()
294 if opml_text is None:
295 # the session expired, log in again and retry once
296 await self._login()
297 opml_text = await self._request_opml()
298 if opml_text is None:
299 raise LoginFailed("Overcast rejected the session right after a fresh login")
300 return opml_text
301
302 async def _request_opml(self) -> str | None:
303 """Return the raw OPML document, or None if the session cookie was rejected."""
304 if remaining := self._rate_limit_remaining():
305 # spend no request while Overcast is still refusing them: the export allows
306 # only ~10 per day and every episode listing would otherwise cost one
307 raise ResourceTemporarilyUnavailable(
308 "Overcast OPML export is rate limited", backoff_time=remaining
309 )
310 try:
311 async with self.http_session.get(OPML_EXPORT_URL, allow_redirects=False) as response:
312 if response.status == 200:
313 return await response.text()
314 if response.status == 429:
315 backoff = (
316 parse_retry_after(response.headers.get("Retry-After"))
317 or RATE_LIMIT_FALLBACK_BACKOFF
318 )
319 self._rate_limited_until = time.monotonic() + backoff
320 raise ResourceTemporarilyUnavailable(
321 "Overcast OPML export is rate limited", backoff_time=backoff
322 )
323 if response.status in AUTH_REJECT_STATUSES:
324 return None
325 raise ResourceTemporarilyUnavailable(
326 f"Overcast OPML export failed with HTTP {response.status}"
327 )
328 except (TimeoutError, aiohttp.ClientError) as err:
329 raise ResourceTemporarilyUnavailable("Overcast is unreachable") from err
330
331 def _rate_limit_remaining(self) -> int:
332 """Return the seconds left of a known Overcast rate limit window, 0 if none."""
333 if self._rate_limited_until is None:
334 return 0
335 remaining = self._rate_limited_until - time.monotonic()
336 if remaining <= 0:
337 self._rate_limited_until = None
338 return 0
339 return math.ceil(remaining)
340
341 async def _get_opml_subscriptions(self) -> dict[str, OvercastSubscription]:
342 opml_text = await self._fetch_opml_text()
343 if self._opml_cache is None or self._opml_cache[0] != opml_text:
344 parsed = await asyncio.to_thread(parse_extended_opml, opml_text)
345 self._opml_cache = (opml_text, parsed)
346 return self._opml_cache[1]
347
348 async def _get_subscription(self, feed_url: str) -> OvercastSubscription | None:
349 """Return the Overcast subscription for a feed, or None if unavailable."""
350 try:
351 subscriptions = await self._get_opml_subscriptions()
352 except (ResourceTemporarilyUnavailable, LoginFailed) as err:
353 # episodes can still be listed without playback state
354 self.logger.debug("Could not obtain Overcast playback states: %s", err)
355 return None
356 return subscriptions.get(feed_url)
357
358 async def _apply_playback_states(
359 self,
360 feed_url: str,
361 subscription: OvercastSubscription,
362 parsed_podcast: dict[str, Any],
363 ) -> datetime | None:
364 """
365 Push a feed's new Overcast playback states to the playlog.
366
367 :param feed_url: The podcast's feed url (also the provider item id).
368 :param subscription: The Overcast subscription holding the episode states.
369 :param parsed_podcast: The podcastparser dict of the feed.
370 :return: The newest state timestamp that was applied, or None if none were.
371 """
372 watermark = self._feed_watermarks.get(feed_url)
373 newest_applied: datetime | None = None
374 podcast_cover = parsed_podcast.get("cover_url")
375 podcast_name = parsed_podcast.get("title")
376 for cnt, parsed_episode in enumerate(parsed_podcast.get("episodes", [])):
377 try:
378 stream_url, _ = get_stream_url_and_guid_from_episode(episode=parsed_episode)
379 except ValueError:
380 continue
381 state = match_episode_state(subscription, stream_url)
382 if state is None or state.user_updated_at is None:
383 continue
384 if not state.played and not state.progress_s:
385 # never mark items unplayed: an absent state cannot be told apart
386 # from an episode that simply was never touched in Overcast
387 continue
388 if watermark is not None and state.user_updated_at <= watermark:
389 # already applied in a previous sync; skipping it also makes sure
390 # local progress made since then is not overwritten
391 continue
392 mass_episode = parse_podcast_episode(
393 episode=parsed_episode,
394 prov_podcast_id=feed_url,
395 episode_cnt=cnt,
396 podcast_cover=podcast_cover,
397 podcast_name=podcast_name,
398 instance_id=self.instance_id,
399 domain=self.domain,
400 )
401 if mass_episode is None:
402 continue
403 if not state.played:
404 # never move the user backwards: the playlog write replaces whatever MA
405 # recorded itself, which may be a position further into the episode
406 _, local_position_ms = await self.mass.music.get_resume_position(mass_episode)
407 if local_position_ms > (state.progress_s or 0) * 1000:
408 continue
409 await self.mass.music.mark_item_played(
410 mass_episode,
411 fully_played=state.played,
412 seconds_played=state.progress_s or 0,
413 user_initiated=False,
414 )
415 if newest_applied is None or state.user_updated_at > newest_applied:
416 newest_applied = state.user_updated_at
417 return newest_applied
418
419 async def _store_watermarks(self) -> None:
420 """Persist the per-feed watermarks of the applied playback states."""
421 # watermarks of feeds that are no longer subscribed are kept on purpose,
422 # so re-subscribing does not re-apply the feed's entire playback history
423 await self.mass.cache.set(
424 key=CACHE_KEY_LAST_APPLIED,
425 provider=self.instance_id,
426 category=CACHE_CATEGORY_OPML,
427 data={feed_url: ts.isoformat() for feed_url, ts in self._feed_watermarks.items()},
428 )
429
430 async def _enrich_episode_chapters(
431 self, prov_podcast_id: str, guid_or_stream_url: str, mass_episode: PodcastEpisode
432 ) -> None:
433 """
434 Attach external ``podcast:chapters`` JSON to a resolved single episode, if any.
435
436 :param prov_podcast_id: Provider podcast id the episode belongs to.
437 :param guid_or_stream_url: Episode identifier used to locate the raw parsed episode.
438 :param mass_episode: The episode to enrich in place; left untouched on any failure.
439 """
440 if mass_episode.metadata.chapters:
441 return
442 podcast = await self._cache_get_podcast(prov_podcast_id)
443 for episode in podcast.get("episodes", []):
444 try:
445 stream_url, guid = get_stream_url_and_guid_from_episode(episode=episode)
446 except ValueError:
447 continue
448 if guid_or_stream_url in (guid, stream_url):
449 await enrich_episode_chapters(
450 session=self.mass.http_session,
451 chapters_json_url=episode.get("chapters_json_url"),
452 mass_episode=mass_episode,
453 )
454 return
455
456 async def _get_episode_stream_url(self, podcast_id: str, guid_or_stream_url: str) -> str | None:
457 parsed_podcast = await self._cache_get_podcast(podcast_id)
458 return find_episode_stream_url(
459 parsed_feed=parsed_podcast, guid_or_stream_url=guid_or_stream_url
460 )
461
462 async def _cache_get_podcast(self, prov_podcast_id: str) -> dict[str, Any]:
463 # raises MediaNotFoundError when the feed is gone
464 return await get_cached_podcast(
465 mass=self.mass,
466 provider_instance_id=self.instance_id,
467 feed_url=prov_podcast_id,
468 max_episodes=self.max_episodes,
469 )
470