diff --git a/lib/streaming_providers/providers/lib_theplatform.py b/lib/streaming_providers/providers/lib_theplatform.py index 229580a..01520a3 100644 --- a/lib/streaming_providers/providers/lib_theplatform.py +++ b/lib/streaming_providers/providers/lib_theplatform.py @@ -41,16 +41,94 @@ Usage (magentaeu): account_uri), USER_AGENT) """ +from __future__ import annotations import base64 from dataclasses import dataclass, field -from typing import List, Optional +from typing import Any, Callable, Dict, List, Optional from urllib.parse import quote from ..base.models import DRMConfig, DRMSystem, LicenseConfig from ..base.utils.logger import logger from ..base.utils.timestamp_converter import TimestampConverter +class PaginationError(Exception): + """ + Raised when a page fetch fails outright. Carries whatever entries were + successfully collected before the failure, so callers can choose to use + partial data — but they have to choose explicitly, rather than a failed + fetch silently masquerading as "end of data". + """ + + def __init__(self, message: str, partial_entries: Optional[List[Dict[str, Any]]] = None): + super().__init__(message) + self.partial_entries = partial_entries or [] + + +def paginate_feed( + fetch_page: Callable[[int, int], Optional[Dict[str, Any]]], + items_per_page: int = 100, + max_pages: int = 20, + start_index: int = 1, + feed_name: str = "feed", +) -> List[Dict[str, Any]]: + """ + Args: + fetch_page: callable(start_index, items_per_page) -> response dict + with an "entries" key, or None on failure. This callable owns + headers, retries, the `operation` tag, and raise_for_status — + pagination here is HTTP-agnostic by design. + items_per_page: page size requested via `range=`. + max_pages: safety cap. Hitting it is logged at ERROR (not silently + swallowed) since it usually means termination isn't firing or + the feed grew beyond expectations. + feed_name: for logging only. + + Returns: + Combined entries across all pages. + + Raises: + PaginationError if a page fetch fails. exc.partial_entries holds + whatever was collected before the failure. + """ + all_entries: List[Dict[str, Any]] = [] + index = start_index + page = 0 + + while page < max_pages: + page += 1 + response = fetch_page(index, items_per_page) + + if response is None: + raise PaginationError( + f"{feed_name}: page {page} (range={index}-{index + items_per_page - 1}) " + f"fetch failed", + partial_entries=all_entries, + ) + + entries = response.get("entries", []) + all_entries.extend(entries) + + logger.debug( + f"{feed_name}: page {page} range={index}-{index + items_per_page - 1} -> " + f"{len(entries)} entries (entryCount={response.get('entryCount')!r})" + ) + + # Termination based on what THIS page actually returned, never on + # entryCount — its meaning varies by feed and isn't worth trusting. + if len(entries) < items_per_page: + break + + index += items_per_page + else: + logger.error( + f"{feed_name}: hit max_pages={max_pages} safety cap with " + f"{len(all_entries)} entries collected — feed may be larger than " + f"expected, or termination isn't triggering. Returning partial data." + ) + + logger.info(f"{feed_name}: {len(all_entries)} entries across {page} page(s)") + return all_entries # --------------------------------------------------------------------------- # Canonical channel dataclass @@ -265,67 +343,48 @@ def fetch_entitled_channels_feed( user_agent: str, timeout: int = 30, ) -> List[TheplatformChannel]: - """ - Fetch the entitled-channels feed for the given distribution rights and - return parsed TheplatformChannel objects. - - Handles pagination automatically (range=1-100, 101-200, …). - - Args: - http_manager: HTTPManager instance. - feed_url: Resolved feed URL (MpxAccountPid already substituted). - e.g. "https://feed.entertainment.tv.theplatform.eu/f/mdeprod/mdeprod-entitled-channels" - distribution_rights: List of DistributionRight URL strings. - cid: Correlation ID string. - user_agent: Platform user-agent string. - timeout: HTTP timeout in seconds. - - Returns: - List of TheplatformChannel, one per DASH-capable channel entry. - """ if not distribution_rights: logger.warning("lib_theplatform: no distribution rights — cannot fetch entitled channels") return [] rights_param = "|".join(distribution_rights) - channels: List[TheplatformChannel] = [] - page_size = 100 - start = 1 + headers = {"User-Agent": user_agent, "Accept": "application/json"} - while True: + def fetch_page(start_index: int, page_size: int) -> Optional[Dict[str, Any]]: + params = { + "byDistributionRightId": rights_param, + "range": f"{start_index}-{start_index + page_size - 1}", + "cid": cid, + } try: - params = { - "byDistributionRightId": rights_param, - "range": f"{start}-{start + page_size - 1}", - "cid": cid, - } response = http_manager.get( feed_url, operation="entitled_channels_feed", - headers={"User-Agent": user_agent, "Accept": "application/json"}, + headers=headers, params=params, timeout=timeout, ) - if response.status_code != 200: - logger.error( - f"lib_theplatform: entitled channels feed failed [{response.status_code}]" - ) - break - - data = response.json() - page_channels = parse_entitled_channels_feed(data) - channels.extend(page_channels) - - entry_count = data.get("entryCount", 0) - if len(page_channels) < page_size or entry_count < page_size: - break - - start += page_size - + response.raise_for_status() + return response.json() except Exception as exc: - logger.error(f"lib_theplatform: error fetching entitled channels page {start}: {exc}") - break + logger.warning(f"lib_theplatform: entitled channels page fetch failed: {exc}") + return None + try: + entries = paginate_feed( + fetch_page, items_per_page=100, max_pages=20, feed_name="entitled_channels" + ) + except PaginationError as exc: + if not exc.partial_entries: + logger.error(f"lib_theplatform: {exc} — no entitled channels fetched") + raise + logger.error( + f"lib_theplatform: {exc} — channel lineup INCOMPLETE, " + f"using {len(exc.partial_entries)} raw entries fetched before failure" + ) + entries = exc.partial_entries + + channels = parse_entitled_channels_feed({"entries": entries}) logger.info(f"lib_theplatform: fetched {len(channels)} entitled channels total") return channels diff --git a/lib/streaming_providers/providers/magenta2/channel_manager.py b/lib/streaming_providers/providers/magenta2/channel_manager.py index 409cc29..7d17ed9 100644 --- a/lib/streaming_providers/providers/magenta2/channel_manager.py +++ b/lib/streaming_providers/providers/magenta2/channel_manager.py @@ -11,6 +11,7 @@ from typing import Callable, Dict, List, Optional from ...base.models import StreamingChannel from ...base.utils.logger import logger +from ..lib_theplatform import paginate_feed, PaginationError from .constants import ( CONTENT_TYPE_LIVE, DEFAULT_EPG_WINDOW_HOURS, @@ -493,12 +494,6 @@ class ChannelManager: return None def _fetch_station_metadata(self) -> Dict[str, Dict]: - """ - Fetch the unauthenticated channel-stations feed and build metadata. - - Returns a dict keyed by station_id (numeric), each value containing: - title, logo_url, quality, channel_number, playback_id - """ metadata: Dict[str, Dict] = {} self._station_to_playback_id.clear() @@ -506,23 +501,42 @@ class ChannelManager: url = None if self._endpoint_manager: url = ( - self._endpoint_manager.get_endpoint("channel_stations") - or self._endpoint_manager.get_endpoint("channel_list") + self._endpoint_manager.get_endpoint("channel_stations") + or self._endpoint_manager.get_endpoint("channel_list") ) - url = ( - url - or "https://feed.entertainment.tv.theplatform.eu/f/mdeprod/mdeprod-channel-stations-main" + base_url = ( + url + or "https://feed.entertainment.tv.theplatform.eu/f/mdeprod/mdeprod-channel-stations-main" ) - url += "?lang=short-de&sort=dt%24displayChannelNumber&range=1-1000" - headers = self._get_api_headers(require_auth=False) - response = self._http.get( - url, operation="api", headers=headers, timeout=DEFAULT_REQUEST_TIMEOUT - ) - response.raise_for_status() - data = response.json() - for entry in data.get("entries", []): + def _fetch_page(start_index: int, page_size: int) -> Optional[Dict]: + page_url = ( + f"{base_url}?lang=short-de&sort=dt%24displayChannelNumber" + f"&range={start_index}-{start_index + page_size - 1}" + ) + try: + response = self._http.get( + page_url, operation="api", headers=headers, timeout=DEFAULT_REQUEST_TIMEOUT + ) + response.raise_for_status() + return response.json() + except Exception as exc: + logger.warning(f"_fetch_station_metadata: page fetch failed: {exc}") + return None + + try: + entries = paginate_feed( + _fetch_page, items_per_page=1000, max_pages=10, feed_name="channel_stations" + ) + except PaginationError as exc: + logger.warning( + f"_fetch_station_metadata: {exc} — using " + f"{len(exc.partial_entries)} entries fetched before failure" + ) + entries = exc.partial_entries + + for entry in entries: try: stations = entry.get("stations", {}) if not stations: @@ -573,9 +587,7 @@ class ChannelManager: f"{len(self._station_to_playback_id)} playback mappings" ) except Exception as exc: - logger.warning( - f"_fetch_station_metadata: failed, channels will have no names: {exc}" - ) + logger.warning(f"_fetch_station_metadata: failed, channels will have no names: {exc}") return metadata diff --git a/lib/streaming_providers/providers/magenta2/epg_manager.py b/lib/streaming_providers/providers/magenta2/epg_manager.py index 5d0ad88..eb09a66 100644 --- a/lib/streaming_providers/providers/magenta2/epg_manager.py +++ b/lib/streaming_providers/providers/magenta2/epg_manager.py @@ -23,6 +23,7 @@ from typing import Any, Dict, List, Optional, Set, Tuple from ...base.models.epg_models import EPGEntry, EPGProgramDetails from ...base.utils.logger import logger +from ..lib_theplatform import paginate_feed, PaginationError from .constants import DEFAULT_REQUEST_TIMEOUT @@ -296,53 +297,46 @@ class Magenta2EpgManager: # Schedule fetching # ------------------------------------------------------------------ - def _fetch_schedule( - self, - date_from: datetime, - date_to: datetime, - ) -> Dict[str, Any]: - """ - Fetch the full schedule grid for the [date_from, date_to] window in - a single ThePlatform API call using byListingTime=ISO~ISO. - - Returns the raw decoded JSON (with an "entries" list) or an empty - dict on failure / missing config. - """ + def _fetch_schedule(self, date_from: datetime, date_to: datetime) -> Dict[str, Any]: if not self._schedule_feed_url or not self._location_id: return {} utc_from = self._ensure_tz(date_from).astimezone(timezone.utc) utc_to = self._ensure_tz(date_to).astimezone(timezone.utc) - - # ThePlatform expects: YYYY-MM-DDTHH:MM:SS.000Z~YYYY-MM-DDTHH:MM:SS.000Z start_iso = utc_from.strftime("%Y-%m-%dT%H:%M:%S.000Z") end_iso = utc_to.strftime("%Y-%m-%dT%H:%M:%S.000Z") - logger.debug( - f"[Magenta2EpgManager] Fetching schedule window " - f"{start_iso} ~ {end_iso} (UTC)" - ) + logger.debug(f"[Magenta2EpgManager] Fetching schedule window {start_iso} ~ {end_iso} (UTC)") - # Schedule only needs: time, title, guid, stationId. - # Everything else is fetched via get_program_details() on demand. fields = ( "listings.stationId," "listings.program.guid,listings.program.title," "listings.startTime,listings.endTime" ) - url = ( - f"{self._schedule_feed_url}" - f"?byListingTime={start_iso}~{end_iso}" - f"&byLocationId={self._location_id}" - f"&fields={fields}" - f"&cid={uuid.uuid4()}" - ) + def _fetch_page(start_index: int, page_size: int) -> Optional[Dict[str, Any]]: + url = ( + f"{self._schedule_feed_url}" + f"?byListingTime={start_iso}~{end_iso}" + f"&byLocationId={self._location_id}" + f"&fields={fields}" + f"&range={start_index}-{start_index + page_size - 1}" + f"&cid={uuid.uuid4()}" + ) + return self._get_with_retry(url, operation="epg_schedule_window") - data = self._get_with_retry(url, operation="epg_schedule_window") - if data and data.get("entries"): - return data - return {} + try: + entries = paginate_feed( + _fetch_page, items_per_page=200, max_pages=10, feed_name="epg_schedule_window" + ) + except PaginationError as exc: + logger.error( + f"[Magenta2EpgManager] {exc} — EPG grid will be INCOMPLETE " + f"({len(exc.partial_entries)} stations fetched before failure)" + ) + entries = exc.partial_entries + + return {"entries": entries} if entries else {} # ------------------------------------------------------------------ # Listing -> EPGEntry