/
/
1"""
2gPodder provider for Music Assistant.
3
4Tested against opodsync, https://github.com/kd2org/opodsync
5and nextcloud-gpodder, https://github.com/thrillfall/nextcloud-gpodder
6gpodder.net is not supported due to responsiveness/ frequent downtimes of domain.
7
8Note:
9 - it can happen, that we have the guid and use that for identification, but the sync state
10 provider, eg. opodsync might use only the stream url. So always make sure, to compare both
11 when relying on an external service
12 - The service calls have a timestamp (int, unix epoch s), which give the changes since then.
13"""
14
15from __future__ import annotations
16
17import time
18from collections.abc import AsyncGenerator
19from datetime import datetime
20from typing import TYPE_CHECKING, Any, cast
21
22from music_assistant_models.config_entries import ConfigEntry, ProviderConfig
23from music_assistant_models.enums import (
24 ConfigEntryType,
25 ContentType,
26 MediaType,
27 ProviderFeature,
28 StreamType,
29)
30from music_assistant_models.errors import (
31 LoginFailed,
32 MediaNotFoundError,
33 ResourceTemporarilyUnavailable,
34)
35from music_assistant_models.media_items import AudioFormat, MediaItemType, Podcast, PodcastEpisode
36from music_assistant_models.streamdetails import StreamDetails
37
38from music_assistant.helpers.datetime import from_utc_timestamp
39from music_assistant.helpers.podcast_parsers import (
40 enrich_episode_chapters,
41 find_episode_stream_url,
42 get_cached_podcast,
43 get_stream_url_and_guid_from_episode,
44 parse_podcast,
45 parse_podcast_episode,
46 refresh_cached_podcast,
47)
48from music_assistant.models.music_provider import MusicProvider
49
50from .client import EpisodeActionDelete, EpisodeActionNew, EpisodeActionPlay, GPodderClient
51
52if TYPE_CHECKING:
53 from music_assistant_models.provider import ProviderManifest
54
55 from music_assistant.mass import MusicAssistant
56 from music_assistant.models import ProviderInstanceType
57
58# Config for "classic" gpodder api
59CONF_URL = "url"
60CONF_USERNAME = "username"
61CONF_PASSWORD = "password"
62CONF_DEVICE_ID = "device_id"
63
64# Config for nextcloud
65CONF_TOKEN_NC = "token"
66CONF_URL_NC = "url_nc"
67
68# General config
69CONF_VERIFY_SSL = "verify_ssl"
70CONF_MAX_NUM_EPISODES = "max_num_episodes"
71
72
73# category 0 holds the individual parsed podcasts, see CACHE_CATEGORY_PODCAST_FEED
74CACHE_CATEGORY_OTHER = 1
75CACHE_KEY_TIMESTAMP = (
76 "timestamp" # tuple of two ints, timestamp_subscriptions and timestamp_actions
77)
78CACHE_KEY_FEEDS = "feeds" # list[str] : all available rss feed urls
79
80SUPPORTED_FEATURES = {
81 ProviderFeature.LIBRARY_PODCASTS,
82 ProviderFeature.BROWSE,
83}
84
85
86async def setup(
87 mass: MusicAssistant, manifest: ProviderManifest, config: ProviderConfig
88) -> ProviderInstanceType:
89 """Initialize provider(instance) with given configuration."""
90 return GPodder(mass, manifest, config, SUPPORTED_FEATURES)
91
92
93class GPodder(MusicProvider):
94 """gPodder MusicProvider."""
95
96 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
97 """
98 Return the (options) config entries for the gPodder provider.
99
100 The server/account connection (gpodder API or Nextcloud) is set up by the interactive
101 setup flow (see ``setup_flow.py``); only the max-episodes limit is configured here.
102 """
103 return (
104 ConfigEntry(
105 key=CONF_MAX_NUM_EPISODES,
106 type=ConfigEntryType.INTEGER,
107 required=False,
108 default_value=0,
109 ),
110 )
111
112 async def handle_async_init(self) -> None:
113 """Pass config values to client and initialize."""
114 base_url = str(self.get_setup_value(CONF_URL))
115 _username = self.get_setup_value(CONF_USERNAME)
116 _password = self.get_setup_value(CONF_PASSWORD)
117 _device_id = self.get_setup_value(CONF_DEVICE_ID)
118 nc_url = str(self.get_setup_value(CONF_URL_NC))
119 nc_token = self.get_setup_value(CONF_TOKEN_NC)
120 verify_ssl = bool(self.get_setup_value(CONF_VERIFY_SSL, True))
121
122 self.max_episodes = cast("int", self.config.get_value(CONF_MAX_NUM_EPISODES, 0))
123
124 self._client = GPodderClient(
125 session=self.mass.http_session, logger=self.logger, verify_ssl=verify_ssl
126 )
127
128 if nc_token is not None:
129 assert nc_url is not None
130 self._client.init_nc(base_url=nc_url, nc_token=str(nc_token))
131 else:
132 if _username is None or _password is None or _device_id is None:
133 raise LoginFailed("Must provide username, password and device_id.")
134 username = str(_username)
135 password = str(_password)
136 device_id = str(_device_id)
137
138 if base_url.rstrip("/") == "https://gpodder.net":
139 raise LoginFailed("Do not use gpodder.net. See docs for explanation.")
140 try:
141 await self._client.init_gpodder(
142 username=username, password=password, base_url=base_url, device=device_id
143 )
144 except RuntimeError as exc:
145 raise LoginFailed("Login failed.") from exc
146
147 timestamps = await self.mass.cache.get(
148 key=CACHE_KEY_TIMESTAMP,
149 provider=self.instance_id,
150 category=CACHE_CATEGORY_OTHER,
151 default=None,
152 )
153 if timestamps is None:
154 self.timestamp_subscriptions: int = 0
155 self.timestamp_actions: int = 0
156 else:
157 self.timestamp_subscriptions, self.timestamp_actions = timestamps
158
159 self.logger.debug(
160 "Our timestamps are (subscriptions, actions) (%s, %s)",
161 self.timestamp_subscriptions,
162 self.timestamp_actions,
163 )
164
165 feeds = await self.mass.cache.get(
166 key=CACHE_KEY_FEEDS,
167 provider=self.instance_id,
168 category=CACHE_CATEGORY_OTHER,
169 default=None,
170 )
171 if feeds is None:
172 self.feeds: set[str] = set()
173 else:
174 self.feeds = set(feeds) # feeds is a list here
175
176 # we are syncing the playlog, but not event based. A simple check in on_played,
177 # should be sufficient
178 self.progress_guard_timestamp = 0.0
179
180 @property
181 def is_streaming_provider(self) -> bool:
182 """Return True if the provider is a streaming provider."""
183 # For streaming providers return True here but for local file based providers return False.
184 # While the streams are remote, the user controls what is added.
185 return False
186
187 async def get_library_podcasts(self) -> AsyncGenerator[Podcast]:
188 """Retrieve library/subscribed podcasts from the provider."""
189 try:
190 subscriptions = await self._client.get_subscriptions()
191 except RuntimeError:
192 raise ResourceTemporarilyUnavailable(backoff_time=30)
193 if subscriptions is None:
194 return
195
196 for feed_url in subscriptions.add:
197 self.feeds.add(feed_url)
198 for feed_url in subscriptions.remove:
199 try:
200 self.feeds.remove(feed_url)
201 except KeyError:
202 # a podcast might have been added and removed in our absence...
203 continue
204
205 episode_actions, timestamp_action = await self._client.get_episode_actions()
206 for feed_url in self.feeds:
207 self.logger.debug("Adding podcast with feed %s to library", feed_url)
208 # parse podcast
209 try:
210 parsed_podcast = await refresh_cached_podcast(
211 mass=self.mass,
212 provider_instance_id=self.instance_id,
213 feed_url=feed_url,
214 max_episodes=self.max_episodes,
215 )
216 except MediaNotFoundError as err:
217 self.report_skipped_sync_item(MediaType.PODCAST, feed_url, err)
218 continue
219
220 # playlog
221 # be safe, if there should be multiple episodeactions. client already sorts
222 # progresses in descending order.
223 _already_processed = set()
224 _episode_actions = [x for x in episode_actions if x.podcast == feed_url]
225 for _action in _episode_actions:
226 if _action.episode not in _already_processed:
227 _already_processed.add(_action.episode)
228 # we do not have to add the progress, these would make calls twice,
229 # and we only use the object to propagate to playlog
230 self.progress_guard_timestamp = time.time()
231 _episode_ids: list[str] = []
232 if _action.guid is not None:
233 _episode_ids.append(f"{feed_url} {_action.guid}")
234 _episode_ids.append(f"{feed_url} {_action.episode}")
235 mass_episode: PodcastEpisode | None = None
236 for _episode_id in _episode_ids:
237 try:
238 mass_episode = await self.get_podcast_episode(
239 _episode_id, add_progress=False
240 )
241 break
242 except MediaNotFoundError:
243 continue
244 if mass_episode is None:
245 self.logger.debug(
246 f"Was unable to use progress for episode {_action.episode}."
247 )
248 continue
249 match _action:
250 case EpisodeActionNew():
251 await self.mass.music.mark_item_unplayed(mass_episode)
252 case EpisodeActionPlay():
253 await self.mass.music.mark_item_played(
254 mass_episode,
255 fully_played=_action.position >= _action.total,
256 seconds_played=_action.position,
257 user_initiated=False,
258 )
259
260 # cache
261 yield parse_podcast(
262 feed_url=feed_url,
263 parsed_feed=parsed_podcast,
264 instance_id=self.instance_id,
265 domain=self.domain,
266 )
267
268 self.timestamp_subscriptions = subscriptions.timestamp
269 if timestamp_action is not None:
270 self.timestamp_actions = timestamp_action
271 await self._cache_set_timestamps()
272 await self._cache_set_feeds()
273
274 async def get_podcast(self, prov_podcast_id: str) -> Podcast:
275 """Get Podcast."""
276 parsed_podcast = await self._cache_get_podcast(prov_podcast_id)
277
278 return parse_podcast(
279 feed_url=prov_podcast_id,
280 parsed_feed=parsed_podcast,
281 instance_id=self.instance_id,
282 domain=self.domain,
283 )
284
285 async def get_podcast_episodes(
286 self, prov_podcast_id: str, add_progress: bool = True
287 ) -> AsyncGenerator[PodcastEpisode]:
288 """Get Podcast episodes. Add progress information."""
289 if add_progress:
290 episode_actions, timestamp = await self._client.get_episode_actions()
291 else:
292 episode_actions, timestamp = [], None
293
294 podcast = await self._cache_get_podcast(prov_podcast_id)
295 podcast_cover = podcast.get("cover_url")
296 parsed_episodes = podcast.get("episodes", [])
297
298 if timestamp is not None:
299 self.timestamp_actions = timestamp
300 await self._cache_set_timestamps()
301
302 for cnt, parsed_episode in enumerate(parsed_episodes):
303 mass_episode = parse_podcast_episode(
304 episode=parsed_episode,
305 prov_podcast_id=prov_podcast_id,
306 episode_cnt=cnt,
307 podcast_cover=podcast_cover,
308 podcast_name=podcast.get("title"),
309 domain=self.domain,
310 instance_id=self.instance_id,
311 )
312 if mass_episode is None:
313 # faulty episode
314 continue
315 try:
316 stream_url, guid = get_stream_url_and_guid_from_episode(episode=parsed_episode)
317 except ValueError:
318 # episode enclosure or stream url missing
319 continue
320
321 for action in episode_actions:
322 # we have to test both, as we are comparing to external input.
323 _test = [action.guid, action.episode]
324 if prov_podcast_id == action.podcast and (guid in _test or stream_url in _test):
325 self.progress_guard_timestamp = time.time()
326 if isinstance(action, EpisodeActionNew):
327 mass_episode.resume_position_ms = 0
328 mass_episode.fully_played = False
329
330 # propagate to playlog
331 await self.mass.music.mark_item_unplayed(
332 mass_episode,
333 )
334 elif isinstance(action, EpisodeActionPlay):
335 fully_played = action.position >= action.total
336 resume_position_s = action.position
337 mass_episode.resume_position_ms = resume_position_s * 1000
338 mass_episode.fully_played = fully_played
339
340 # propagate progress to playlog
341 await self.mass.music.mark_item_played(
342 mass_episode,
343 fully_played=fully_played,
344 seconds_played=resume_position_s,
345 user_initiated=False,
346 )
347 elif isinstance(action, EpisodeActionDelete):
348 for mapping in mass_episode.provider_mappings:
349 mapping.available = False
350 break
351 yield mass_episode
352
353 async def get_podcast_episode(
354 self, prov_episode_id: str, add_progress: bool = True
355 ) -> PodcastEpisode:
356 """Get Podcast Episode. Add progress information."""
357 podcast_id, guid_or_stream_url = prov_episode_id.split(" ")
358 async for mass_episode in self.get_podcast_episodes(podcast_id, add_progress=add_progress):
359 _, _guid_or_stream_url = mass_episode.item_id.split(" ")
360 # this is enough, as internal
361 if guid_or_stream_url == _guid_or_stream_url:
362 await self._enrich_episode_chapters(podcast_id, guid_or_stream_url, mass_episode)
363 return mass_episode
364 raise MediaNotFoundError("Did not find episode.")
365
366 async def get_resume_position(
367 self, item_id: str, media_type: MediaType
368 ) -> tuple[bool, int, datetime | None]:
369 """Return: finished, position_ms."""
370 assert media_type == MediaType.PODCAST_EPISODE
371 podcast_id, guid_or_stream_url = item_id.split(" ")
372 stream_url = await self._get_episode_stream_url(podcast_id, guid_or_stream_url)
373 try:
374 progresses, timestamp = await self._client.get_episode_actions(
375 since=self.timestamp_actions
376 )
377 except RuntimeError:
378 self.logger.warning("Was unable to obtain progresses.")
379 raise NotImplementedError # fallback to internal position.
380 for action in progresses:
381 _test = [action.guid, action.episode]
382 # progress is external, compare guid and stream_url
383 if action.podcast == podcast_id and (
384 guid_or_stream_url in _test or stream_url in _test
385 ):
386 dt_timestamp: datetime | None = None
387 if timestamp is not None:
388 self.timestamp_actions = timestamp
389 await self._cache_set_timestamps()
390 dt_timestamp = from_utc_timestamp(timestamp)
391 if isinstance(action, EpisodeActionNew | EpisodeActionDelete):
392 # no progress, it might have been actively reset
393 # in case of delete, we start from start.
394 return False, 0, None
395 _progress = (action.position >= action.total, max(action.position * 1000, 0))
396 self.logger.debug("Found an updated external resume position.")
397 return action.position >= action.total, max(action.position * 1000, 0), dt_timestamp
398 self.logger.debug("Did not find an updated resume position, falling back to stored.")
399 # If we did not find a resume position, nothing changed since our last timestamp
400 # we raise NotImplementedError, such that MA falls back to the already stored
401 # resume_position in its playlog.
402 raise NotImplementedError
403
404 async def on_played(
405 self,
406 media_type: MediaType,
407 prov_item_id: str,
408 fully_played: bool,
409 position: int,
410 media_item: MediaItemType,
411 is_playing: bool = False,
412 ) -> None:
413 """Update progress."""
414 if media_item is None or not isinstance(media_item, PodcastEpisode):
415 return
416 if media_type != MediaType.PODCAST_EPISODE:
417 return
418 if time.time() - self.progress_guard_timestamp <= 5:
419 return
420 podcast_id, guid_or_stream_url = prov_item_id.split(" ")
421 stream_url = await self._get_episode_stream_url(podcast_id, guid_or_stream_url)
422 assert stream_url is not None
423 duration = media_item.duration
424 try:
425 await self._client.update_progress(
426 podcast_id=podcast_id,
427 episode_id=stream_url,
428 guid=guid_or_stream_url,
429 position_s=position,
430 duration_s=duration,
431 )
432 self.logger.debug(f"Updated progress to {position / duration * 100:.2f}%")
433 except RuntimeError as exc:
434 self.logger.debug(exc)
435 self.logger.debug("Failed to update progress.")
436
437 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
438 """Get streamdetails for item."""
439 podcast_id, guid_or_stream_url = item_id.split(" ")
440 stream_url = await self._get_episode_stream_url(podcast_id, guid_or_stream_url)
441 if stream_url is None:
442 raise MediaNotFoundError
443 return StreamDetails(
444 provider=self.instance_id,
445 item_id=item_id,
446 audio_format=AudioFormat(
447 content_type=ContentType.try_parse(stream_url),
448 ),
449 media_type=MediaType.PODCAST_EPISODE,
450 stream_type=StreamType.HTTP,
451 path=stream_url,
452 can_seek=True,
453 allow_seek=True,
454 )
455
456 async def _enrich_episode_chapters(
457 self, prov_podcast_id: str, guid_or_stream_url: str, mass_episode: PodcastEpisode
458 ) -> None:
459 """
460 Attach external ``podcast:chapters`` JSON to a resolved single episode, if any.
461
462 :param prov_podcast_id: Provider podcast id the episode belongs to.
463 :param guid_or_stream_url: Episode identifier used to locate the raw parsed episode.
464 :param mass_episode: The episode to enrich in place; left untouched on any failure.
465 """
466 if mass_episode.metadata.chapters:
467 return
468 podcast = await self._cache_get_podcast(prov_podcast_id)
469 for episode in podcast.get("episodes", []):
470 try:
471 stream_url, guid = get_stream_url_and_guid_from_episode(episode=episode)
472 except ValueError:
473 continue
474 if guid_or_stream_url in (guid, stream_url):
475 await enrich_episode_chapters(
476 session=self.mass.http_session,
477 chapters_json_url=episode.get("chapters_json_url"),
478 mass_episode=mass_episode,
479 )
480 return
481
482 async def _get_episode_stream_url(self, podcast_id: str, guid_or_stream_url: str) -> str | None:
483 parsed_podcast = await self._cache_get_podcast(podcast_id)
484 return find_episode_stream_url(
485 parsed_feed=parsed_podcast, guid_or_stream_url=guid_or_stream_url
486 )
487
488 async def _cache_get_podcast(self, prov_podcast_id: str) -> dict[str, Any]:
489 # raises MediaNotFoundError when the feed is gone
490 return await get_cached_podcast(
491 mass=self.mass,
492 provider_instance_id=self.instance_id,
493 feed_url=prov_podcast_id,
494 max_episodes=self.max_episodes,
495 )
496
497 async def _cache_set_timestamps(self) -> None:
498 # seven days default
499 await self.mass.cache.set(
500 key=CACHE_KEY_TIMESTAMP,
501 provider=self.instance_id,
502 category=CACHE_CATEGORY_OTHER,
503 data=[self.timestamp_subscriptions, self.timestamp_actions],
504 )
505
506 async def _cache_set_feeds(self) -> None:
507 # seven days default
508 await self.mass.cache.set(
509 key=CACHE_KEY_FEEDS,
510 provider=self.instance_id,
511 category=CACHE_CATEGORY_OTHER,
512 data=list(self.feeds),
513 )
514