From 04ededea69819a30b9be5135910d9de9be8e511d Mon Sep 17 00:00:00 2001 From: Nirvana Date: Wed, 7 Oct 2026 12:12:10 +0200 Subject: [PATCH] work towards provider v2 --- .../base/auth/base_oauth2_auth.py | 60 +- lib/streaming_providers/base/errors.py | 40 ++ .../base/managed_provider.py | 376 +++++++++++++ .../base/managers/__init__.py | 17 +- .../base/managers/bookmarks.py | 47 +- .../base/managers/catchup.py | 50 +- .../base/managers/channel.py | 40 +- lib/streaming_providers/base/managers/epg.py | 26 +- .../base/managers/favorites.py | 48 +- .../base/managers/recordings.py | 70 ++- lib/streaming_providers/base/managers/vod.py | 46 +- .../base/models/channel.py | 6 + .../base/models/content.py | 60 +- .../base/models/drm/drm_config.py | 68 ++- .../base/models/drm/exceptions.py | 20 +- .../base/models/drm/license_config.py | 101 +++- .../base/models/pricing.py | 20 +- lib/streaming_providers/base/models/vod.py | 10 +- .../base/provider_registry.py | 125 +++-- lib/streaming_providers/base/vod.py | 4 +- .../providers/_template/README.md | 36 +- .../providers/_template/provider.py | 448 ++++----------- .../providers/joyn/channel_manager.py | 515 ++++++++++++------ .../providers/joyn/constants.py | 59 +- 24 files changed, 1428 insertions(+), 864 deletions(-) create mode 100644 lib/streaming_providers/base/managed_provider.py diff --git a/lib/streaming_providers/base/auth/base_oauth2_auth.py b/lib/streaming_providers/base/auth/base_oauth2_auth.py index f977c2e..3e4c10b 100644 --- a/lib/streaming_providers/base/auth/base_oauth2_auth.py +++ b/lib/streaming_providers/base/auth/base_oauth2_auth.py @@ -88,42 +88,37 @@ class WafBlockedException(Exception): pass class SessionAwareHTTPManager: - """Wraps http_manager to provide session-like cookie handling while maintaining proxy support""" + """ + Thin wrapper over HTTPManager for OAuth flows: shared headers and + operation="auth" (so proxy scope honours ProxyScope.authentication). + + Cookies are handled entirely by the underlying requests.Session cookie + jar. Never set a 'Cookie' header manually here — doing so overrides + requests' own cookie handling and drops intermediate-redirect cookies + (e.g. Joyn's `status_id`), which breaks multi-step OAuth flows. + + The `headers` dict is merged into every request; per-call headers win + over session-level ones. + """ def __init__(self, http_manager): self.http_manager = http_manager - self.cookies: Dict[str, str] = {} self.headers: Dict[str, str] = {} + def _prepare(self, kwargs): + kwargs.setdefault("operation", "auth") + merged = {**self.headers, **(kwargs.get("headers") or {})} # per-call wins + if merged: + kwargs["headers"] = merged + else: + kwargs.pop("headers", None) + return kwargs + def get(self, url: str, **kwargs): - """GET request with cookie handling""" - headers = kwargs.get("headers", {}).copy() - headers.update(self.headers) - if self.cookies: - cookie_str = "; ".join([f"{k}={v}" for k, v in self.cookies.items()]) - headers["Cookie"] = cookie_str - kwargs["headers"] = headers - response = self.http_manager.get(url, operation="oauth", **kwargs) - self._update_cookies_from_response(response) - return response + return self.http_manager.get(url, **self._prepare(kwargs)) def post(self, url: str, **kwargs): - """POST request with cookie handling""" - headers = kwargs.get("headers", {}).copy() - headers.update(self.headers) - if self.cookies: - cookie_str = "; ".join([f"{k}={v}" for k, v in self.cookies.items()]) - headers["Cookie"] = cookie_str - kwargs["headers"] = headers - response = self.http_manager.post(url, operation="oauth", **kwargs) - self._update_cookies_from_response(response) - return response - - def _update_cookies_from_response(self, response): - """Extract and update cookies from response""" - if hasattr(response, "cookies"): - for cookie in response.cookies: - self.cookies[cookie.name] = cookie.value + return self.http_manager.post(url, **self._prepare(kwargs)) # OAuth2RemoteLoginMixin comes first so its methods take precedence over any @@ -519,11 +514,10 @@ class BaseOAuth2Authenticator(OAuth2RemoteLoginMixin, BaseAuthenticator): def _create_oauth_session(self) -> SessionAwareHTTPManager: """Create a session-aware HTTP manager for OAuth flows""" session = SessionAwareHTTPManager(self.http_manager) - session.headers.update({ - "User-Agent": self.config.user_agent, - "Referer": getattr(self.config, "base_website", ""), - "Origin": getattr(self.config, "base_website", ""), - }) + base = getattr(self.config, "base_website", "") + session.headers["User-Agent"] = self.config.user_agent + if base: + session.headers.update({"Referer": base, "Origin": base}) return session # ======================================================================== diff --git a/lib/streaming_providers/base/errors.py b/lib/streaming_providers/base/errors.py index d4cc7c4..b1b1dc9 100644 --- a/lib/streaming_providers/base/errors.py +++ b/lib/streaming_providers/base/errors.py @@ -182,6 +182,25 @@ class NotFoundError(ProviderError): """ +class ItemNotFoundError(NotFoundError, KeyError): + """A user-scoped item (recording, favorite, bookmark) does not exist. + + Raised by delete_recording / remove_favorite / delete_bookmark when the + id is not currently stored. It is deliberately BOTH a NotFoundError + (uniform ProviderError handling) and a KeyError (the pre-existing + convention documented on the Provider*Mixin classes), so callers can + migrate from ``except KeyError`` to ``except ItemNotFoundError`` at + their own pace. Do not confuse with "this manager doesn't handle that + content_id", which is a None / [] return value. + """ + + def __str__(self) -> str: + # KeyError.__str__ would repr() the message (extra quotes). + if len(self.args) == 1: + return str(self.args[0]) + return Exception.__str__(self) + + class BadRequestError(ProviderError): """400 -- the request was malformed for the endpoint called. @@ -220,10 +239,31 @@ class NotImplementedYetError(ProviderError): """Feature captured but not yet wired up.""" +class UnsupportedOperationError(NotImplementedYetError): + """The provider does not support this operation (and will not). + + Distinct in meaning from NotImplementedYetError ("planned, not done"). + Subclasses it so existing handlers keep working. Prefer checking the + capability flag (e.g. RecordingsManager.supports_scheduling) before + calling, and treat this as the safety net. + """ + + class ConfigurationError(ProviderError): """Provider is misconfigured.""" +class OperationFailedError(ProviderError, RuntimeError): + """A write operation was rejected or failed on the provider backend. + + Raised by add_favorite / remove_favorite / update_bookmark / + delete_bookmark / delete_recording on backend refusal or failure. Also a + RuntimeError so existing ``except RuntimeError`` handlers keep working. + Prefer a more specific subclass of ProviderError (ServerError, + EntitlementError, ...) when the cause is known. + """ + + # --------------------------------------------------------------------------- # HTTP status -> error class heuristic # --------------------------------------------------------------------------- diff --git a/lib/streaming_providers/base/managed_provider.py b/lib/streaming_providers/base/managed_provider.py new file mode 100644 index 0000000..752c052 --- /dev/null +++ b/lib/streaming_providers/base/managed_provider.py @@ -0,0 +1,376 @@ +# streaming_providers/base/managed_provider.py +""" +ManagedProvider -- StreamingProvider plus manager composition. DRAFT. + +Supersedes provider_v2.py. It SUBCLASSES StreamingProvider, so registry +metadata, isinstance checks in the backend and the legacy mixins keep +working. Only new/migrated providers use it; legacy providers are +untouched. + +It moves out of every provider (see simplitv/provider.py) what is +identical in all of them: + * the _build_*() hooks (default None = capability absent) + * the implements_* flags + * _route() and the default get_manifest() / get_drm() routing + * get_channels(), EPG + header delegation, DRM validation + +A provider still owns: http/auth/cache setup, the _build_*() bodies, and +anything with provider-specific grammar (simplitv: the catchup: branch +and get_restart()). + +Usage in a provider's __init__, after http_manager / auth / caches exist: + + self._init_managers() + +and for provider-specific prefixes: + + def get_manifest(self, content_id, **kw): + if content_id.startswith(CATCHUP_PREFIX): + ... + return super().get_manifest(content_id, **kw) + +Routing rules +------------- +1. handles_content_id() is a cheap pre-filter; False skips the manager. +2. A manager that passes may still return None/[]; the next is tried. +3. NotFoundError is remembered and re-raised only if nobody resolves. +4. BadRequestError is never caught. +5. The manager that resolved an id is remembered (bounded, locked), so + headers go to the same manager that produced the manifest. +6. content_type narrows only for "live" and "vod"; anything else widens. +""" + +from __future__ import annotations + +import threading +from collections import OrderedDict +from datetime import datetime +from typing import Any, Callable, ClassVar, Dict, List, Optional, Tuple + +from .errors import ConfigurationError, NotFoundError +from .managers import ( + BookmarksManager, + CatchupManager, + ChannelManager, + EpgManager, + FavoritesManager, + RecordingsManager, + VodManager, +) +from .models import StreamingChannel +# Module paths below assume models/drm/{exceptions,drm_config}.py (see TODO R-1). +from .models.drm.drm_config import validate_drm_set +from .models.drm.exceptions import DRMError +from .protocols import DrmManagerProtocol +from .provider import StreamingProvider + +_ROUTED = ("channels", "vod") + + +class ManagedProvider(StreamingProvider): + # Folded DRM: managers override get_channel_drm / get_vod_drm. + # Explicit on purpose -- no override sniffing via method identity. + DRM_IN_MANAGERS: ClassVar[bool] = False + + # True: header hooks on the managers are honoured (default is then + # auth.build_headers(), NOT the legacy {}). Set False to keep legacy {}. + HEADERS_FROM_MANAGERS: ClassVar[bool] = True + + ROUTE_CACHE_SIZE: ClassVar[int] = 512 + + # NOTE: `self.channels` holds the ChannelManager here (template + # convention) and shadows the legacy list set by StreamingProvider. + # to_output_format() is overridden below for that reason. + + # ------------------------------------------------------------------ + # Wiring + # ------------------------------------------------------------------ + + def _init_managers(self) -> None: + """Build managers in dependency order (catchup sees channels).""" + self._route_cache: "OrderedDict[str, str]" = OrderedDict() + self._route_lock = threading.Lock() + + self.channels = self._build_channels() + self.vod = self._build_vod() + self.epg = self._build_epg() + self.recordings = self._build_recordings() + self.favorites = self._build_favorites() + self.bookmarks = self._build_bookmarks() + self.catchup = self._build_catchup() + self.drm = self._build_drm() + self._check_wiring() + + def _build_channels(self) -> Optional[ChannelManager]: + return None + + def _build_vod(self) -> Optional[VodManager]: + return None + + def _build_epg(self) -> Optional[EpgManager]: + return None + + def _build_recordings(self) -> Optional[RecordingsManager]: + return None + + def _build_favorites(self) -> Optional[FavoritesManager]: + return None + + def _build_bookmarks(self) -> Optional[BookmarksManager]: + return None + + def _build_catchup(self) -> Optional[CatchupManager]: + return None + + def _build_drm(self) -> Optional[DrmManagerProtocol]: + return None + + def _check_wiring(self) -> None: + expected = { + "channels": ChannelManager, + "vod": VodManager, + "epg": EpgManager, + "recordings": RecordingsManager, + "favorites": FavoritesManager, + "bookmarks": BookmarksManager, + "catchup": CatchupManager, + } + for attr, abc_type in expected.items(): + mgr = getattr(self, attr) + if mgr is not None and not isinstance(mgr, abc_type): + raise ConfigurationError( + f"{type(self).__name__}.{attr} is {type(mgr).__name__}, " + f"expected a {abc_type.__name__}" + ) + if self.drm is not None and self.DRM_IN_MANAGERS: + raise ConfigurationError( + f"{type(self).__name__}: use _build_drm() OR " + f"DRM_IN_MANAGERS, not both" + ) + + # ------------------------------------------------------------------ + # Capability flags (one signal each) + # ------------------------------------------------------------------ + + @property + def implements_channels(self) -> bool: + return self.channels is not None + + @property + def implements_vod(self) -> bool: + return self.vod is not None + + @property + def implements_epg(self) -> bool: + # Manager present AND it reports a window (epg_window != (0, 0)). + return self.epg is not None and self.epg.implements_epg + + @property + def implements_recordings(self) -> bool: + return self.recordings is not None + + @property + def implements_favorites(self) -> bool: + return self.favorites is not None + + @property + def implements_bookmarks(self) -> bool: + return self.bookmarks is not None + + @property + def implements_catchup(self) -> bool: + return self.catchup is not None and self.catchup.supports_catchup + + @property + def implements_drm(self) -> bool: + return self.drm is not None or self.DRM_IN_MANAGERS + + @property + def capabilities(self) -> Dict[str, bool]: + """All derived flags in one dict (same keys the registry reports).""" + return { + name: getattr(self, f"implements_{name}") + for name in ( + "channels", "vod", "epg", "recordings", "favorites", + "bookmarks", "catchup", "drm", + ) + } + + # ------------------------------------------------------------------ + # Router + # ------------------------------------------------------------------ + + def _remember(self, content_id: str, name: str) -> None: + with self._route_lock: + self._route_cache[content_id] = name + self._route_cache.move_to_end(content_id) + while len(self._route_cache) > self.ROUTE_CACHE_SIZE: + self._route_cache.popitem(last=False) + + def _route( + self, + content_id: str, + attempts: List[Tuple[str, Callable[[Any], Any]]], + ) -> Any: + """ + attempts: [(manager attribute name, call(manager)), ...] + A falsy result (None, "", []) means "not mine". + """ + last_not_found: Optional[NotFoundError] = None + for name, call in attempts: + manager = getattr(self, name, None) + if manager is None or not manager.handles_content_id(content_id): + continue + try: + result = call(manager) + except NotFoundError as exc: + last_not_found = exc + continue + if result: + self._remember(content_id, name) + return result + if last_not_found is not None: + raise last_not_found + return None + + def _routed_manager(self, content_id: str) -> Tuple[Optional[str], Any]: + """Manager that resolved this id before, else the first that accepts it.""" + with self._route_lock: + name = self._route_cache.get(content_id) + if name is not None: + return name, getattr(self, name, None) + for candidate in _ROUTED: + mgr = getattr(self, candidate, None) + if mgr is not None and mgr.handles_content_id(content_id): + return candidate, mgr + return None, None + + # ------------------------------------------------------------------ + # Public surface + # ------------------------------------------------------------------ + + def get_channels(self, **kw: Any) -> List[StreamingChannel]: + if self.channels is None: + return [] + return self.channels.get_channels(**kw) + + # --- EPG ----------------------------------------------------------- + # The legacy ProviderEpgMixin never looks at self.epg: epg_window is + # (0, 0), get_epg() returns [] and get_epg_grid() returns {} unless the + # provider overrides them. Delegate once, here. + + @property + def epg_window(self) -> Tuple[int, int]: + return self.epg.epg_window if self.epg is not None else (0, 0) + + def get_epg( + self, + channel_id: str, + start_time: Optional[datetime] = None, + end_time: Optional[datetime] = None, + country: Optional[str] = None, # legacy arg; the manager has its own + **kw: Any, + ) -> List: + if not self.implements_epg or not self.epg.handles_channel_id(channel_id): + return [] + return self.epg.get_epg( + channel_id, start_time=start_time, end_time=end_time, **kw + ) + + def get_epg_grid( + self, + start_time: Optional[datetime] = None, + end_time: Optional[datetime] = None, + channel_ids: Optional[List[str]] = None, + country: Optional[str] = None, + **kw: Any, + ) -> Dict[str, List]: + # Legacy order is (start, end, channel_ids); the ABC's is + # (channel_ids, start, end). Always pass by keyword. + if not self.implements_epg: + return {} + if channel_ids is None: + # Assumes EPG channel_id == content_id; verify per provider. + channel_ids = [c.content_id for c in self.get_channels()] + return self.epg.get_epg_grid( + channel_ids, start_time=start_time, end_time=end_time, **kw + ) + + def get_program_details(self, program_id: str, **kw: Any): + if not self.implements_epg: + return None + return self.epg.get_program_details(program_id, **kw) + + def to_output_format(self, channels: Optional[List[StreamingChannel]] = None) -> Dict: + # Legacy implementation iterates self.channels, which is a manager here. + if channels is None: + channels = self.get_channels() + return super().to_output_format(channels) + + def get_manifest(self, content_id: str, **kw: Any) -> Optional[str]: + return self._route(content_id, [ + ("channels", lambda m: m.get_channel_manifest(content_id, **kw)), + ("vod", lambda m: m.get_vod_manifest(content_id, **kw)), + ]) + + def get_manifest_headers(self, content_id: str, **kw: Any) -> Dict[str, str]: + if not self.HEADERS_FROM_MANAGERS: + return super().get_manifest_headers(content_id, **kw) + name, mgr = self._routed_manager(content_id) + if name == "channels": + return mgr.get_channel_manifest_headers(content_id, **kw) + if name == "vod": + return mgr.get_vod_manifest_headers(content_id, **kw) + return super().get_manifest_headers(content_id, **kw) + + def get_segment_headers(self, content_id: str, **kw: Any) -> Dict[str, str]: + if self.HEADERS_FROM_MANAGERS: + name, mgr = self._routed_manager(content_id) + if name in _ROUTED: # both ChannelManager and VodManager have the hook + return mgr.get_segment_headers(content_id, **kw) + return self.get_manifest_headers(content_id, **kw) + + def get_drm( + self, + content_id: str, + drm_variant: Optional[str] = None, # same position as legacy + content_type: Optional[str] = None, # "live" | "vod" narrow; else widen + **kw: Any, + ) -> List: + if drm_variant is not None: + kw["drm_variant"] = drm_variant + + if self.drm is not None: + configs = self.drm.get_drm_configs( + content_id, content_type=content_type, **kw + ) + else: + attempts: List[Tuple[str, Callable[[Any], Any]]] = [] + if content_type != "vod": + attempts.append( + ("channels", lambda m: m.get_channel_drm(content_id, **kw)) + ) + if content_type != "live": + attempts.append( + ("vod", lambda m: m.get_vod_drm(content_id, **kw)) + ) + configs = self._route(content_id, attempts) or [] + return self._validate_drm(content_id, configs) + + def _validate_drm(self, content_id: str, configs: List) -> List: + """ + Reject configs ISA would mishandle: invalid, or equal priorities. + + Delegates to validate_drm_set (models/drm/drm_config.py). Note the + factories (create_widevine, create_playready, ...) all default to + priority=1, so a provider returning several DRMs must assign + priorities (or use merge_drm_configs(..., auto_priority=True)). + """ + try: + validate_drm_set(configs) + except DRMError as exc: + raise ConfigurationError( + f"{type(self).__name__}: invalid DRM configuration for " + f"{content_id}: {exc}" + ) from exc + return configs \ No newline at end of file diff --git a/lib/streaming_providers/base/managers/__init__.py b/lib/streaming_providers/base/managers/__init__.py index 7531c79..22859f9 100644 --- a/lib/streaming_providers/base/managers/__init__.py +++ b/lib/streaming_providers/base/managers/__init__.py @@ -8,6 +8,8 @@ methods (headers, DRM defaults, search no-ops) come for free. Design rules ------------ +* All managers share one constructor contract, implemented once in the + private ManagerBase (managers/_base.py). * Constructors take four required collaborators -- http_manager, auth, country, config -- plus keyword-only extras (caches, collaborators). * Managers never hold a reference to the provider. Shared state is passed @@ -30,12 +32,17 @@ conventions documented in the template. Manager list ------------ -Three required capabilities: - ChannelManager - VodManager (providers without a browseable catalogue return None) - EpgManager (providers without EPG return None) +All seven managers are OPTIONAL (see providers/_template/README.md): a +provider wires the ones its service offers and returns None from the other +_build_*() factories. A VOD-only provider has no ChannelManager; a +metadata-only provider may have just an EpgManager. -Four optional capabilities -- providers implement the ones they support: +Core capabilities: + ChannelManager (live channels, channel manifest / DRM) + VodManager (browseable catalogue) + EpgManager (guide data) + +Four further capabilities -- providers implement the ones they support: RecordingsManager (cloud / network PVR) FavoritesManager (user bookmarks on programs / channels) BookmarksManager (resume position) diff --git a/lib/streaming_providers/base/managers/bookmarks.py b/lib/streaming_providers/base/managers/bookmarks.py index 3417392..164d3c8 100644 --- a/lib/streaming_providers/base/managers/bookmarks.py +++ b/lib/streaming_providers/base/managers/bookmarks.py @@ -16,14 +16,19 @@ Return-value conventions ------------------------ get_bookmarks returns [] when the user has no bookmarks. Not an error. -update_bookmark raises RuntimeError if the provider rejects the write +update_bookmark raises OperationFailedError if the provider rejects the write (e.g. content inaccessible, backend error). It is called on every playback stop / pause, so providers should tolerate a write that overwrites the same position with a no-op rather than failing. -delete_bookmark raises KeyError if no bookmark exists for content_id, -so callers can distinguish "already gone" from "successfully deleted". -The base ProviderBookmarksMixin documents the same rule. +delete_bookmark raises ItemNotFoundError if no bookmark exists for +content_id, so callers can distinguish "already gone" from "successfully +deleted". The base ProviderBookmarksMixin documents the same rule (as +KeyError; ItemNotFoundError is a KeyError subclass). + +ItemNotFoundError is both a NotFoundError (ProviderError) and a KeyError, and +OperationFailedError is both a ProviderError and a RuntimeError, so handlers +written against the old KeyError / RuntimeError convention keep working. Caller guidance --------------- @@ -35,36 +40,16 @@ they are given. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, List, Optional from ..models.bookmark import Bookmark, ContentType -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class BookmarksManager(ABC): +class BookmarksManager(ManagerBase): """Abstract base for provider bookmarks managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Abstract # ------------------------------------------------------------------ @@ -94,7 +79,8 @@ class BookmarksManager(ABC): Called on playback stop / pause, so should be tolerant of repeated writes to the same position (a no-op write is fine). - Raises RuntimeError if the provider rejects the write. + Raises OperationFailedError (also a RuntimeError) if the provider + rejects the write. """ raise NotImplementedError @@ -104,7 +90,8 @@ class BookmarksManager(ABC): Delete a bookmark. Raises: - KeyError: if no bookmark exists for content_id. - RuntimeError: on backend failure. + ItemNotFoundError: if no bookmark exists for content_id + (also a KeyError). + OperationFailedError: on backend failure (also a RuntimeError). """ raise NotImplementedError \ No newline at end of file diff --git a/lib/streaming_providers/base/managers/catchup.py b/lib/streaming_providers/base/managers/catchup.py index ac58635..6a292b9 100644 --- a/lib/streaming_providers/base/managers/catchup.py +++ b/lib/streaming_providers/base/managers/catchup.py @@ -9,6 +9,7 @@ Public interface get_catchup_manifest(content_id, start_time, end_time=None, ...) -> Optional[str] [abstract] get_catchup_manifest_headers(...) -> Dict[str,str] [concrete] + get_catchup_segment_headers(...) -> Dict[str,str] [concrete] get_catchup_drm(...) -> List[DRMConfig] [concrete] Constructor contract @@ -30,8 +31,8 @@ get_catchup_drm returns [] when catchup shares DRM with live (the common case), or a provider-specific list when catchup uses different DRM. -Return the manifest of the live stream --------------------------------------- +What a catchup manifest is (and is not) +---------------------------------------- The catchup manifest is a *modified* live manifest URL (with time parameters) for providers like Magenta and MoveTV, or a distinct URL for providers whose catchup is served from a different origin. This @@ -60,36 +61,16 @@ when the ABC accepts None. Pass None. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, Dict, List, Optional from ..models import DRMConfig -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class CatchupManager(ABC): +class CatchupManager(ManagerBase): """Abstract base for provider catchup managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Capability # ------------------------------------------------------------------ @@ -165,6 +146,25 @@ class CatchupManager(ABC): """ return self.auth.build_headers() + def get_catchup_segment_headers( + self, + content_id: str, + start_time: int, + end_time: Optional[int] = None, + epg_id: Optional[str] = None, + **kw: Any, + ) -> Dict[str, str]: + """ + Headers for catchup segment requests. + + Default: the catchup manifest headers. Mirrors + ChannelManager.get_segment_headers; override for providers whose + catchup segments need different or token-bound headers. + """ + return self.get_catchup_manifest_headers( + content_id, start_time, end_time, epg_id, **kw + ) + def get_catchup_drm( self, content_id: str, diff --git a/lib/streaming_providers/base/managers/channel.py b/lib/streaming_providers/base/managers/channel.py index ba5736c..0e0ee7b 100644 --- a/lib/streaming_providers/base/managers/channel.py +++ b/lib/streaming_providers/base/managers/channel.py @@ -13,10 +13,10 @@ Public interface Constructor contract -------------------- -Four required keyword-only collaborators. Subclasses that need extra state -declare additional keyword-only args and store them on self AFTER calling -super().__init__. The base does NOT accept **kwargs -- a typo at a call -site becomes an immediate TypeError, which is what you want. +Four required keyword-only collaborators (see ManagerBase). Subclasses that +need extra state declare additional keyword-only args and store them on +self AFTER calling super().__init__. The base does NOT accept **kwargs -- a +typo at a call site becomes an immediate TypeError, which is what you want. Return-value conventions ------------------------ @@ -27,39 +27,16 @@ failures propagate as exceptions from base.errors. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, Dict, List, Optional from ..models import Channel, DRMConfig -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class ChannelManager(ABC): +class ChannelManager(ManagerBase): """Abstract base for provider channel managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - # Sanity check the Auth collaborator against the runtime_checkable - # protocol. isinstance on a runtime_checkable Protocol only verifies - # method presence, not signatures -- that's the intended check here. - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Routing # ------------------------------------------------------------------ @@ -71,7 +48,8 @@ class ChannelManager(ABC): Default: True. Override in providers whose channel content_ids have a distinguishable grammar (e.g. numeric-only for live channels). The orchestrator uses this to route get_manifest / get_drm without - a wasted request to the wrong manager. + a wasted request to the wrong manager. It is a cheap, I/O-free + PRE-FILTER: False means the manager is never asked. """ return True diff --git a/lib/streaming_providers/base/managers/epg.py b/lib/streaming_providers/base/managers/epg.py index 385f7e7..8b34a4c 100644 --- a/lib/streaming_providers/base/managers/epg.py +++ b/lib/streaming_providers/base/managers/epg.py @@ -22,37 +22,17 @@ epg_window returns (past_days, future_days). (0, 0) means no EPG support. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from datetime import datetime from typing import Any, Dict, List, Optional, Tuple from ..models.epg_models import EPGEntry, EPGProgramDetails -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class EpgManager(ABC): +class EpgManager(ManagerBase): """Abstract base for provider EPG managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Capability # ------------------------------------------------------------------ diff --git a/lib/streaming_providers/base/managers/favorites.py b/lib/streaming_providers/base/managers/favorites.py index 9cf5e66..97b6380 100644 --- a/lib/streaming_providers/base/managers/favorites.py +++ b/lib/streaming_providers/base/managers/favorites.py @@ -18,13 +18,17 @@ Return-value conventions get_favorites returns [] when the user has no favorites or the provider does not support favorites. Not an error. -add_favorite returns the created Favorite. Raises RuntimeError on +add_favorite returns the created Favorite. Raises OperationFailedError on rejection (e.g. provider's backend refuses the operation). -remove_favorite raises KeyError if the content_id is not currently -favorited, and RuntimeError on backend failure. Deleting a -non-existent favorite is a KeyError -- consistent with -FavoritesManager's counterpart in ProviderFavoritesMixin. +remove_favorite raises ItemNotFoundError if the content_id is not currently +favorited, and OperationFailedError on backend failure. Deleting a +non-existent favorite is an ItemNotFoundError -- which is also a KeyError, +consistent with ProviderFavoritesMixin. + +ItemNotFoundError is both a NotFoundError (ProviderError) and a KeyError, and +OperationFailedError is both a ProviderError and a RuntimeError, so handlers +written against the old KeyError / RuntimeError convention keep working. Provider guidance ----------------- @@ -36,36 +40,16 @@ them (PROGRAM, CHANNEL, CLIP, LIVE, EVENT). from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, List, Optional from ..models.favorite import Favorite, FavoriteType -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class FavoritesManager(ABC): +class FavoritesManager(ManagerBase): """Abstract base for provider favorites managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Abstract # ------------------------------------------------------------------ @@ -91,7 +75,8 @@ class FavoritesManager(ABC): """ Add a favorite. - Raises RuntimeError if the provider refuses the operation. + Raises OperationFailedError (also a RuntimeError) if the provider + refuses the operation. """ raise NotImplementedError @@ -101,8 +86,9 @@ class FavoritesManager(ABC): Remove a favorite. Raises: - KeyError: if content_id is not currently favorited. - RuntimeError: on backend failure. + ItemNotFoundError: if content_id is not currently favorited + (also a KeyError). + OperationFailedError: on backend failure (also a RuntimeError). """ raise NotImplementedError diff --git a/lib/streaming_providers/base/managers/recordings.py b/lib/streaming_providers/base/managers/recordings.py index 0c038c9..2115e9d 100644 --- a/lib/streaming_providers/base/managers/recordings.py +++ b/lib/streaming_providers/base/managers/recordings.py @@ -20,13 +20,14 @@ Return-value conventions get_recordings returns [] when the provider has no recordings. Callers should treat an empty list as "no recordings", not as an error. -delete_recording raises KeyError when the recording does not exist, and -ProviderError subclasses (usually ServerError / EntitlementError) on -transport or permission failures. Returning silently on a failed delete -would hide real errors. +delete_recording raises ItemNotFoundError (also a KeyError) when the +recording does not exist, and ProviderError subclasses (usually ServerError +/ EntitlementError / OperationFailedError) on transport or permission +failures. Returning silently on a failed delete would hide real errors. schedule_recording is optional -- the ABC provides a default that raises -NotImplementedYetError. Providers that support scheduling override it. +UnsupportedOperationError, and supports_scheduling is False. Providers that +support scheduling override both. Recording identity ------------------ @@ -38,16 +39,15 @@ that distinction is preserved. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, List -from ..errors import NotImplementedYetError +from ..errors import UnsupportedOperationError from ..models import Channel -from ..protocols import AuthProtocol -from ..utils.logger import logger +from ._base import ManagerBase -class RecordingsManager(ABC): +class RecordingsManager(ManagerBase): """ Abstract base for provider recordings managers. @@ -72,29 +72,10 @@ class RecordingsManager(ABC): Providers that need a recording's manifest should implement the routing in their provider's get_manifest, not by adding a - get_manifest method here. See providers/simpli/provider.py for - the pattern. + get_manifest method here. See the simpliTV provider's provider.py + for the pattern. """ - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Routing # ------------------------------------------------------------------ @@ -143,8 +124,9 @@ class RecordingsManager(ABC): Delete a recording. Raises: - KeyError: if the recording does not exist. - ProviderError: on transport / permission / server failure. + ItemNotFoundError: if the recording does not exist (also a + KeyError). + ProviderError: on transport / permission / server failure. """ raise NotImplementedError @@ -152,6 +134,17 @@ class RecordingsManager(ABC): # Concrete -- optional # ------------------------------------------------------------------ + @property + def supports_scheduling(self) -> bool: + """ + True when schedule_recording() is implemented. Default False. + + Providers that override schedule_recording MUST also override this + to return True, so callers can gate the UI instead of catching + UnsupportedOperationError. + """ + return False + def schedule_recording(self, content_id: str, **kw: Any) -> Any: """ Schedule a recording of content_id. @@ -161,10 +154,11 @@ class RecordingsManager(ABC): indicating success. The ABC does not constrain the shape because there is no shared one across providers. - Default raises NotImplementedYetError. Providers that support - scheduling override this. + Default raises UnsupportedOperationError (a NotImplementedYetError + subclass, so old handlers still match). "Not supported by this + provider" is not the same as "planned, not done yet". Providers + that support scheduling override this AND supports_scheduling. """ - raise NotImplementedYetError( - f"{self.__class__.__name__}.schedule_recording is not " - f"implemented" + raise UnsupportedOperationError( + f"{self.__class__.__name__} does not support schedule_recording" ) \ No newline at end of file diff --git a/lib/streaming_providers/base/managers/vod.py b/lib/streaming_providers/base/managers/vod.py index 85ba3a0..0231b19 100644 --- a/lib/streaming_providers/base/managers/vod.py +++ b/lib/streaming_providers/base/managers/vod.py @@ -10,6 +10,7 @@ Public interface get_vod_manifest(content_id, **kw) -> Optional[str] [abstract] search_vod(query, cursor=None, page_size=24, **kw) -> VodPage [concrete] get_vod_manifest_headers(content_id, **kw) -> Dict[str, str] [concrete] + get_segment_headers(content_id, **kw) -> Dict[str, str] [concrete] get_vod_drm(content_id, **kw) -> List[DRMConfig] [concrete] content_id grammar @@ -24,9 +25,9 @@ documents it here. The base class does not parse content_id. Examples: Constructor contract -------------------- -Four required keyword-only collaborators. No **kwargs. Subclasses accept -extra keyword-only args explicitly and call super().__init__ with only -the four required. A typo at a call site becomes an immediate TypeError. +Four required keyword-only collaborators (see ManagerBase). No **kwargs. +Subclasses accept extra keyword-only args explicitly and call +super().__init__ with only the four required. Return-value conventions ------------------------ @@ -37,37 +38,17 @@ propagate as exceptions from base.errors. from __future__ import annotations -from abc import ABC, abstractmethod +from abc import abstractmethod from typing import Any, Dict, List, Optional from ..models import DRMConfig -from ..protocols import AuthProtocol -from ..utils.logger import logger from ..vod import VodPage +from ._base import ManagerBase -class VodManager(ABC): +class VodManager(ManagerBase): """Abstract base for provider VOD managers.""" - def __init__( - self, - *, - http_manager: Any, - auth: AuthProtocol, - country: str, - config: Any, - ) -> None: - if not isinstance(auth, AuthProtocol): - logger.warning( - f"{self.__class__.__name__}: auth does not match AuthProtocol " - f"(missing one of get_access_token / build_headers / " - f"invalidate). Got {type(auth).__name__}." - ) - self.http_manager = http_manager - self.auth = auth - self.country = country - self.config = config - # ------------------------------------------------------------------ # Routing # ------------------------------------------------------------------ @@ -78,6 +59,7 @@ class VodManager(ABC): Default: True. Override in providers whose VOD content_ids have a distinguishable grammar (e.g. start with "details_" or "clip_"). + Cheap, I/O-free PRE-FILTER: False means the manager is never asked. """ return True @@ -123,6 +105,18 @@ class VodManager(ABC): """Headers for the manifest request. Default: auth headers.""" return self.auth.build_headers() + def get_segment_headers( + self, content_id: str, **kw: Any + ) -> Dict[str, str]: + """ + Headers for segment requests. Default: manifest headers. + + Mirrors ChannelManager.get_segment_headers so the orchestrator can + ask any routed manager for segment headers. Override for providers + with token-bound segment URLs. + """ + return self.get_vod_manifest_headers(content_id, **kw) + def get_vod_drm( self, content_id: str, **kw: Any ) -> List[DRMConfig]: diff --git a/lib/streaming_providers/base/models/channel.py b/lib/streaming_providers/base/models/channel.py index 89c7d2e..745d039 100644 --- a/lib/streaming_providers/base/models/channel.py +++ b/lib/streaming_providers/base/models/channel.py @@ -34,6 +34,12 @@ class Channel(Content): f"Channel {self.name} ({self.content_id}): " f"session_manifest=True but manifest URL is set - manifest will be ignored" ) + # Content.__post_init__ is deliberately NOT called: it raises on + # pricing/mode mismatches, and Channel's documented contract is + # "advisory warnings, the code runs" (one bad upstream channel must + # not kill the whole channel list). VodItem stays strict. + for problem in self._pricing_mode_problems(): + logger.warning(f"Channel {self.name} ({self.content_id}): {problem}") # Backward compatibility alias @property diff --git a/lib/streaming_providers/base/models/content.py b/lib/streaming_providers/base/models/content.py index 35265db..537d9f7 100644 --- a/lib/streaming_providers/base/models/content.py +++ b/lib/streaming_providers/base/models/content.py @@ -1,6 +1,6 @@ from dataclasses import dataclass from decimal import Decimal -from typing import Dict, List, Optional, Set +from typing import Dict, List, Optional from .drm import DRMConfig from .pricing import AccessType, PricePoint, Pricing @@ -68,34 +68,44 @@ class Content: # Pricing (None = UNKNOWN, NOT FREE) pricing: Optional[Pricing] = None - def __post_init__(self): - """Validate pricing and mode consistency.""" + def _pricing_mode_problems(self) -> List[str]: + """ + Pricing-vs-mode inconsistencies, as human-readable strings. + + Only combinations that are unambiguous are constrained: + PPV_LIVE -> mode must be "live" + TVOD_RENTAL / TVOD_PURCHASE -> mode must be "vod" + AVOD, PPV_REPLAY and SVOD_PPV used to be forced to "live". That made + every ad-supported on-demand item and every PPV replay (the VodItem + docstring's own "past sports event") fail at construction, because + VodItem normalises mode to "vod" before validating. Those access + types describe how content is MONETISED, not how it is delivered, so + they no longer constrain the mode. + """ if not self.pricing: - return - - # Define mode mappings - live_types: Set[AccessType] = { - AccessType.PPV_LIVE, - AccessType.PPV_REPLAY, - AccessType.SVOD_PPV, - AccessType.AVOD - } - vod_types: Set[AccessType] = { - AccessType.TVOD_RENTAL, - AccessType.TVOD_PURCHASE - } - - # Check mode consistency - if self.pricing.access_type in live_types and self.mode != StreamingMode.LIVE: - raise ValueError( - f"Content '{self.name}' has {self.pricing.access_type.value} pricing " + return [] + access = self.pricing.access_type + problems: List[str] = [] + if access == AccessType.PPV_LIVE and self.mode != StreamingMode.LIVE: + problems.append( + f"Content '{self.name}' has {access.value} pricing " f"but mode is '{self.mode}' (expected '{StreamingMode.LIVE}')" ) - elif self.pricing.access_type in vod_types and self.mode != StreamingMode.VOD: - raise ValueError( - f"Content '{self.name}' has {self.pricing.access_type.value} pricing " + elif ( + access in (AccessType.TVOD_RENTAL, AccessType.TVOD_PURCHASE) + and self.mode != StreamingMode.VOD + ): + problems.append( + f"Content '{self.name}' has {access.value} pricing " f"but mode is '{self.mode}' (expected '{StreamingMode.VOD}')" ) + return problems + + def __post_init__(self): + """Validate pricing and mode consistency (strict: raises).""" + problems = self._pricing_mode_problems() + if problems: + raise ValueError(problems[0]) # --- Pricing Properties --- @@ -226,6 +236,8 @@ class Content: "ContentType": self.content_type, "Country": self.country, "Language": self.language, + "Description": self.description, + "Genre": self.genre, "StreamingFormat": self.streaming_format, "LicenseUrl": self.license_url, "CertificateUrl": self.certificate_url, diff --git a/lib/streaming_providers/base/models/drm/drm_config.py b/lib/streaming_providers/base/models/drm/drm_config.py index a9f060d..69a4a06 100644 --- a/lib/streaming_providers/base/models/drm/drm_config.py +++ b/lib/streaming_providers/base/models/drm/drm_config.py @@ -4,8 +4,8 @@ DRM Configuration Model Main configuration class that combines DRM system, PSSH data, and license config. """ -from dataclasses import dataclass -from typing import Optional +from dataclasses import dataclass, replace +from typing import Dict, Optional, Sequence from .drm_systems import DRMSystem from .license_config import LicenseConfig @@ -213,4 +213,66 @@ class DRMConfig: has_license = "with license" if self.license else "no license" return ( f"" - ) \ No newline at end of file + ) + + +# --------------------------------------------------------------------------- +# Multi-DRM helpers +# --------------------------------------------------------------------------- + +def validate_drm_set(configs: Sequence[DRMConfig]) -> None: + """ + Validate a set of DRMConfig objects that will be sent to ISA together. + + Checks each config (DRMConfig.validate) and that no two configs share a + priority. The factories (create_widevine, create_playready, ...) all + default to priority=1, so a naive [widevine, playready] list is INVALID + until priorities are assigned -- see merge_drm_configs(auto_priority=True). + + Raises: + LicenseConfigError: on any invalid config or duplicate priority. + """ + seen: Dict[int, DRMConfig] = {} + for cfg in configs: + cfg.validate() + if cfg.priority in seen: + raise LicenseConfigError( + f"{cfg.system.name} and {seen[cfg.priority].system.name} " + f"share priority {cfg.priority}; each DRM needs a unique " + f"priority (lower number = higher priority)." + ) + seen[cfg.priority] = cfg + + +def merge_drm_configs( + configs: Sequence[DRMConfig], *, auto_priority: bool = False +) -> dict: + """ + Merge several DRMConfig objects into ONE inputstream.adaptive.drm dict. + + Args: + configs: DRM configs in preference order (first = most preferred). + auto_priority: Renumber priorities 1..n in the given order (on + copies; the inputs are not mutated). Without it the configs + must already carry unique priorities. + + Raises: + LicenseConfigError: invalid config, duplicate priority, or the same + DRM system listed twice. (A config with pre_init_data must keep + priority=1; auto_priority will trip that rule if it is not first.) + """ + items = list(configs) + if auto_priority: + items = [replace(c, priority=i) for i, c in enumerate(items, start=1)] + validate_drm_set(items) + + merged: dict = {} + for cfg in items: + entry = cfg.to_dict() + clash = set(entry) & set(merged) + if clash: + raise LicenseConfigError( + f"DRM system listed twice: {sorted(clash)}" + ) + merged.update(entry) + return merged \ No newline at end of file diff --git a/lib/streaming_providers/base/models/drm/exceptions.py b/lib/streaming_providers/base/models/drm/exceptions.py index 528c573..8d12205 100644 --- a/lib/streaming_providers/base/models/drm/exceptions.py +++ b/lib/streaming_providers/base/models/drm/exceptions.py @@ -2,12 +2,26 @@ DRM-specific Exception Classes Custom exceptions for better error handling and debugging in DRM operations. + +All of them descend from DRMError, which is a ConfigurationError and +therefore a ProviderError (base/errors.py). A caller that handles +``ProviderError`` now also catches DRM configuration problems; before this +change DRMError derived from plain Exception and slipped past every +``except ProviderError`` handler. """ +from ...errors import ConfigurationError -class DRMError(Exception): - """Base exception for all DRM-related errors""" - pass + +class DRMError(ConfigurationError): + """Base exception for all DRM-related errors. + + ``message`` is optional so call sites that raise a bare + ``InvalidPSSHError()`` keep working (ProviderError itself requires one). + """ + + def __init__(self, message: str = "", **kw) -> None: + super().__init__(message, **kw) class InvalidPSSHError(DRMError): diff --git a/lib/streaming_providers/base/models/drm/license_config.py b/lib/streaming_providers/base/models/drm/license_config.py index 19c4275..95d9b00 100644 --- a/lib/streaming_providers/base/models/drm/license_config.py +++ b/lib/streaming_providers/base/models/drm/license_config.py @@ -5,12 +5,13 @@ Data classes for DRM license configuration, including server URLs, certificates, headers, and unwrapper parameters. """ +import base64 import json import re from dataclasses import dataclass, field from enum import Enum from typing import Dict, Optional, Union -from urllib.parse import urlencode, unquote +from urllib.parse import quote, urlencode from .utils import safe_base64_decode, safe_base64_encode, normalize_key_id from .exceptions import LicenseConfigError @@ -149,15 +150,17 @@ class LicenseConfig: norm_kid = normalize_key_id(kid) norm_key = key.lower().replace("-", "") normalized_keyids[norm_kid] = norm_key - except Exception: - pass + except Exception as e: + # Previously swallowed: a bad pair vanished silently, + # keyids ended up empty and validate() still passed. + raise LicenseConfigError( + f"Invalid ClearKey entry (kid={kid!r}): {e}" + ) from e self.keyids = normalized_keyids @staticmethod def _is_base64(s: str) -> bool: """Check if string appears to be base64 encoded.""" - import base64 - import re # Base64 pattern: only valid chars, length multiple of 4 if not re.match(r'^[A-Za-z0-9+/]*={0,2}$', s): return False @@ -167,6 +170,21 @@ class LicenseConfig: except Exception: return False + @staticmethod + def _split_flags(value: Optional[str]) -> list: + """Split a comma-separated flag string ("base64,urlenc") into tokens.""" + return [t.strip().lower() for t in (value or "").split(",") if t.strip()] + + @classmethod + def _validate_flag_list(cls, name: str, value: Optional[str], enum_cls) -> None: + allowed = {m.value for m in enum_cls} + for token in cls._split_flags(value): + if token not in allowed: + raise LicenseConfigError( + f"Invalid {name} flag '{token}'. " + f"Allowed: {', '.join(sorted(allowed))}" + ) + def validate(self) -> None: """ Validate the license configuration. @@ -201,6 +219,17 @@ class LicenseConfig: "Construct via the normal dataclass constructor to ensure normalization." ) + self._validate_flag_list("wrapper", self.wrapper, WrapperType) + self._validate_flag_list("unwrapper", self.unwrapper, UnwrapperType) + if ( + set(self._split_flags(self.unwrapper)) & {"json", "xml"} + and self.unwrapper_params is None + ): + raise LicenseConfigError( + "unwrapper_params is required when unwrapper includes " + "'json' or 'xml'." + ) + if self.keyids: for kid, key in self.keyids.items(): if len(kid) != 32: @@ -272,7 +301,7 @@ class LicenseConfig: # Case 1: Already a dict if isinstance(headers, dict): - return urlencode(headers) + return self._encode_headers(headers) # Case 2: String input if isinstance(headers, str): @@ -285,21 +314,23 @@ class LicenseConfig: try: header_dict = json.loads(headers) if isinstance(header_dict, dict): - return urlencode(header_dict) + return self._encode_headers(header_dict) except json.JSONDecodeError as e: raise LicenseConfigError( f"Invalid JSON in req_headers: {e}" ) from e - # Already URL-encoded: contains '=' and at least one '&' or ';' - if "=" in headers and ("&" in headers or ";" in headers): + # "Key: Value" plain text: a ':' appears before any '='. + if self._looks_plain(headers): + return self._parse_plain_headers(headers) + + # URL-encoded "k=v[&k=v...]". A SINGLE pair is valid too (the + # class docstring's own example is one pair, which the old + # '&'/';' requirement rejected). + if "=" in headers: self._validate_urlencoded_headers(headers) return headers - # "Key: Value" plain-text format (no '=' present) - if ":" in headers and "=" not in headers: - return self._parse_plain_headers(headers) - raise LicenseConfigError( f"Invalid req_headers format: '{headers}'. " f"Must be a dict, JSON string, URL-encoded string, " @@ -310,6 +341,25 @@ class LicenseConfig: f"req_headers must be dict, string, or None, got {type(headers).__name__}" ) + @staticmethod + def _looks_plain(headers: str) -> bool: + """True when ``headers`` is "Key: Value" text rather than "k=v".""" + colon = headers.find(":") + equals = headers.find("=") + return colon != -1 and (equals == -1 or colon < equals) + + @staticmethod + def _encode_headers(header_dict: Dict[str, str]) -> str: + """ + URL-encode a header dict. + + quote_via=quote with safe="" gives "%20" for spaces and "%2F" for + "/" (as documented), where the default quote_plus wrote "+" for + spaces. TODO(device-verify): confirm ISA decodes "+" and "%20" alike; + "%20" is the unambiguous choice either way. + """ + return urlencode(header_dict, quote_via=quote, safe="") + @staticmethod def _validate_urlencoded_headers(headers: str) -> None: """ @@ -334,20 +384,22 @@ class LicenseConfig: key, value = pair.split("=", 1) if not key.strip(): raise LicenseConfigError(f"Empty header key in pair: '{pair}'") - try: - unquote(value) # Raises if the percent-encoding is broken - except Exception as e: + # unquote() never raises, so the old try/except validated nothing. + if re.search(r"%(?![0-9A-Fa-f]{2})", pair): raise LicenseConfigError( - f"Header value is not properly URL-encoded: '{value}'" - ) from e + f"Header value has a malformed percent-escape: '{pair}'" + ) @staticmethod def _parse_plain_headers(headers: str) -> str: """ Convert ``"Key: Value"`` lines to a URL-encoded string. - Lines may be separated by newlines, carriage-returns, semicolons, - or commas. Blank lines are silently skipped. + Headers are separated by newlines / carriage-returns, or by a ';' or + ',' that is followed by another "Name:" token. A comma or semicolon + inside a value ("Mozilla/5.0 (KHTML, like Gecko)", "text/html;q=0.9") + does NOT split. Blank lines are skipped; a header with an empty value + is kept (it used to be dropped silently). Example:: @@ -364,7 +416,7 @@ class LicenseConfig: LicenseConfigError: If any non-blank line lacks a ``':'``. """ header_dict: Dict[str, str] = {} - for line in re.split(r"[\n\r;,]+", headers): + for line in re.split(r"[\n\r]+|[;,]\s*(?=[A-Za-z0-9_\-]+\s*:)", headers): line = line.strip() if not line: continue @@ -375,9 +427,10 @@ class LicenseConfig: key, value = line.split(":", 1) key = key.strip() value = value.strip() - if key and value: - header_dict[key] = value - return urlencode(header_dict) + if not key: + raise LicenseConfigError(f"Empty header name in line: '{line}'") + header_dict[key] = value + return LicenseConfig._encode_headers(header_dict) def to_dict(self) -> dict: """ diff --git a/lib/streaming_providers/base/models/pricing.py b/lib/streaming_providers/base/models/pricing.py index e7c1339..336db16 100644 --- a/lib/streaming_providers/base/models/pricing.py +++ b/lib/streaming_providers/base/models/pricing.py @@ -19,6 +19,22 @@ class AccessType(Enum): SVOD_PPV = "svod_ppv" # Subscription + extra PPV surcharge +def _as_comparable(now: datetime, bound: datetime) -> datetime: + """ + Make ``now`` comparable with ``bound``. + + Python refuses to compare naive and timezone-aware datetimes. Upstream + APIs usually deliver aware valid_from/valid_until values, while the + default "now" was naive, so is_active() raised TypeError. A naive value + is interpreted as local wall-clock time. + """ + if (now.tzinfo is None) == (bound.tzinfo is None): + return now + if bound.tzinfo is not None: + return now.astimezone() # naive now -> aware (local) + return now.astimezone().replace(tzinfo=None) # aware now -> naive local + + @dataclass class PricePoint: """A specific price for a region/quality/time period.""" @@ -42,9 +58,9 @@ class PricePoint: def is_active(self, at: Optional[datetime] = None) -> bool: """Check if this price point is currently valid.""" now = at or datetime.now() - if self.valid_from and now < self.valid_from: + if self.valid_from and _as_comparable(now, self.valid_from) < self.valid_from: return False - if self.valid_until and now > self.valid_until: + if self.valid_until and _as_comparable(now, self.valid_until) > self.valid_until: return False return True diff --git a/lib/streaming_providers/base/models/vod.py b/lib/streaming_providers/base/models/vod.py index 574a740..f38c2ab 100644 --- a/lib/streaming_providers/base/models/vod.py +++ b/lib/streaming_providers/base/models/vod.py @@ -124,8 +124,10 @@ class VodCategory: # of reconstructing it from content_id, which loses query params. fetch_url: Optional[str] = None - # Cached slug (computed lazily if not set) - _slug: Optional[str] = field(default=None, repr=False) + # Cached slug (computed lazily). init=False / compare=False: it is a + # cache, not state -- two equal categories must compare equal whether or + # not .slug was accessed, and it must not be an __init__ parameter. + _slug: Optional[str] = field(default=None, init=False, repr=False, compare=False) def __post_init__(self): """Validate and clean up required fields.""" @@ -218,8 +220,8 @@ class VodItem(Content): # False → full broadcast recording (videoType == "STANDALONE_EVENT" etc.) is_highlight: bool = False - # Cached slug - _slug: Optional[str] = field(default=None, repr=False) + # Cached slug (see VodCategory._slug) + _slug: Optional[str] = field(default=None, init=False, repr=False, compare=False) def __post_init__(self): """Validate and clean up required fields.""" diff --git a/lib/streaming_providers/base/provider_registry.py b/lib/streaming_providers/base/provider_registry.py index 6ba46f1..b6b06a1 100644 --- a/lib/streaming_providers/base/provider_registry.py +++ b/lib/streaming_providers/base/provider_registry.py @@ -3,20 +3,62 @@ Core provider registry handling discovery, metadata, and lifecycle management. """ +import threading from typing import Any, Dict, List, Optional from .provider import StreamingProvider from .utils.logger import logger +_CAPABILITY_NAMES = ( + "channels", "vod", "epg", "recordings", "favorites", "bookmarks", + "catchup", "drm", +) + + +def _collect_capabilities(instance: Optional[StreamingProvider]) -> Optional[Dict[str, bool]]: + """ + Read the derived ``implements_*`` flags of a live provider instance. + + Capabilities are instance-level (they come from which managers the + provider wired), so they are only known once the provider exists, i.e. + for ENABLED providers. Returns None when there is no instance or the + provider exposes no boolean flags (e.g. an old-style provider). + """ + if instance is None: + return None + caps: Dict[str, bool] = {} + for name in _CAPABILITY_NAMES: + try: + value = getattr(instance, f"implements_{name}") + except Exception: # a property may raise on a half-built provider + continue + if isinstance(value, bool): + caps[name] = value + return caps or None + class ProviderMetadata: """Metadata for a provider instance with lazy initialization.""" - def __init__(self, plugin_class, country: str, enabled: bool = False): + def __init__( + self, + plugin_class, + country: str, + enabled: bool = False, + instance_name: Optional[str] = None, + ): self.plugin_class = plugin_class self.country = country.lower() self.enabled = enabled self.instance: Optional[StreamingProvider] = None + # The registry key this metadata is stored under. When given, it + # becomes ``name`` so to_dict()["name"] always resolves in + # get_provider(); otherwise the name is derived from the class name + # and the two can drift apart. + self._instance_name = instance_name + # Guards lazy create/destroy: two threads asking for the same + # provider must not build (and authenticate) two instances. + self._lock = threading.RLock() self._extract_metadata() def _extract_metadata(self): @@ -55,27 +97,36 @@ class ProviderMetadata: for auth_type in self.supported_auth_types ) + if self._instance_name: + self.name = self._instance_name + def create_instance(self) -> Optional[StreamingProvider]: """Lazily create provider instance if enabled""" if not self.enabled: return None - if self.instance is None: - try: - logger.info(f"Creating instance for provider: {self.name}") - self.instance = self.plugin_class(country=self.country) - logger.debug(f"Successfully created instance for {self.name}") - except Exception as e: - logger.error(f"Failed to create instance for {self.name}: {e}") - self.instance = None + with self._lock: + if self.instance is None: + try: + logger.info(f"Creating instance for provider: {self.name}") + self.instance = self.plugin_class(country=self.country) + logger.debug(f"Successfully created instance for {self.name}") + except Exception: + # logger.exception keeps the traceback: a missing + # provider_name, a bad manager wiring (ConfigurationError) + # or an import error otherwise shows up as one cryptic + # line and the provider silently vanishes from the UI. + logger.exception(f"Failed to create instance for {self.name}") + self.instance = None - return self.instance + return self.instance def destroy_instance(self): """Clean up provider instance""" - if self.instance: - logger.debug(f"Destroying instance for provider: {self.name}") - self.instance = None + with self._lock: + if self.instance: + logger.debug(f"Destroying instance for provider: {self.name}") + self.instance = None def set_enabled(self, enabled: bool): """Update enabled status and manage instance accordingly""" @@ -99,6 +150,7 @@ class ProviderMetadata: "logo": self.logo, "is_multi_country": self.is_multi_country, "supported_countries": self.supported_countries, + "capabilities": _collect_capabilities(self.instance), } @@ -125,6 +177,7 @@ class M3UGroupMetadata: self.country = country.lower() self.enabled = enabled self.instance: Optional[StreamingProvider] = None + self._lock = threading.RLock() # Extract metadata from class self._extract_metadata() @@ -161,28 +214,30 @@ class M3UGroupMetadata: if not self.enabled: return None - if self.instance is None: - try: - logger.info(f"Creating M3U group instance: {self.name} (group: {self.group})") + with self._lock: + if self.instance is None: + try: + logger.info(f"Creating M3U group instance: {self.name} (group: {self.group})") - # Create instance with group filter - __init__ will handle proxy setup - self.instance = self.plugin_class( - country=self.country, - group_filter=self.group - ) + # Create instance with group filter - __init__ will handle proxy setup + self.instance = self.plugin_class( + country=self.country, + group_filter=self.group + ) - logger.debug(f"Successfully created instance for M3U group: {self.name}") - except Exception as e: - logger.error(f"Failed to create M3U group instance {self.name}: {e}") - self.instance = None + logger.debug(f"Successfully created instance for M3U group: {self.name}") + except Exception: + logger.exception(f"Failed to create M3U group instance {self.name}") + self.instance = None - return self.instance + return self.instance def destroy_instance(self): """Clean up provider instance""" - if self.instance: - logger.debug(f"Destroying M3U group instance: {self.name}") - self.instance = None + with self._lock: + if self.instance: + logger.debug(f"Destroying M3U group instance: {self.name}") + self.instance = None def set_enabled(self, enabled: bool): """Update enabled status and manage instance accordingly""" @@ -206,6 +261,7 @@ class M3UGroupMetadata: "logo": self.logo, "is_multi_country": self.is_multi_country, "supported_countries": self.supported_countries, + "capabilities": _collect_capabilities(self.instance), # M3U-specific fields "group": self.group, "type": "m3u_group", @@ -260,7 +316,9 @@ class ProviderRegistry: instance_name = f"{plugin_name}_{country}" enabled = self._is_provider_enabled(plugin_name, country) - metadata = ProviderMetadata(plugin_class, country, enabled) + metadata = ProviderMetadata( + plugin_class, country, enabled, instance_name=instance_name + ) self.provider_metadata[instance_name] = metadata discovered.append(instance_name) @@ -272,7 +330,9 @@ class ProviderRegistry: instance_name = plugin_name enabled = self._is_provider_enabled(plugin_name) - metadata = ProviderMetadata(plugin_class, default_country, enabled) + metadata = ProviderMetadata( + plugin_class, default_country, enabled, instance_name=instance_name + ) self.provider_metadata[instance_name] = metadata discovered.append(instance_name) @@ -405,6 +465,9 @@ class ProviderRegistry: return False try: + # Drop the registry's reference first: if re-creation fails we + # must not keep serving the instance the metadata already forgot. + self.providers.pop(provider_name, None) metadata.destroy_instance() new_instance = metadata.create_instance() if new_instance: diff --git a/lib/streaming_providers/base/vod.py b/lib/streaming_providers/base/vod.py index 0566d63..9fa6115 100644 --- a/lib/streaming_providers/base/vod.py +++ b/lib/streaming_providers/base/vod.py @@ -22,6 +22,7 @@ from dataclasses import dataclass, field from typing import Iterator, List, Optional, Union from .models.vod import VodCategory, VodItem +from .utils.logger import logger VodEntry = Union[VodCategory, VodItem] @@ -102,8 +103,7 @@ def normalize_vod_result(result) -> VodPage: ) if isinstance(result, list): return VodPage(entries=result) - import logging - logging.getLogger(__name__).warning( + logger.warning( f"normalize_vod_result: unexpected type {type(result).__name__}; " f"returning empty page" ) diff --git a/lib/streaming_providers/providers/_template/README.md b/lib/streaming_providers/providers/_template/README.md index e9feb85..c497876 100644 --- a/lib/streaming_providers/providers/_template/README.md +++ b/lib/streaming_providers/providers/_template/README.md @@ -1,5 +1,34 @@ # Provider template +> **v2 NOTICE (ManagedProvider) -- read first.** +> New providers should subclass `ManagedProvider` +> (`base/managed_provider.py`, itself a `StreamingProvider`) instead of copying +> the v1 boilerplate. The sections below describe the v1 shape; where they +> disagree with the v2 template (`provider.py` in this directory), v2 wins: +> +> * Capability flags, `_route()`, default `get_manifest` / `get_drm`, header, +> channel and EPG delegation are inherited. Call `self._init_managers()` at +> the end of `__init__`. +> * `implements_epg` = EPG manager present AND `epg_window != (0, 0)`; +> `implements_catchup` = catchup manager present AND `supports_catchup`. +> (v1 said "manager is not None" for every flag.) +> * `get_drm(content_id, drm_variant=None, content_type=None, **kw)` -- the +> second positional argument is `drm_variant`, as on `StreamingProvider`. +> (v1 documented `get_drm(content_id, content_type=None)`.) +> * DRM architecture is declared: `_build_drm()` (dedicated) or +> `DRM_IN_MANAGERS = True` (folded). Not both. +> * The legacy mixins do NOT delegate to `self.epg` / `self.recordings` / +> ...: the v1 claim "no delegation needed" is wrong for EPG (verified) and +> unverified for the others. See docs/provider-v2/TODO.md item M-1. +> * Errors: `ItemNotFoundError` (a `KeyError`) and `OperationFailedError` +> (a `RuntimeError`) replace bare `KeyError` / `RuntimeError`; +> `UnsupportedOperationError` replaces `NotImplementedYetError` for +> "this provider will never support it". +> * Multiple DRMs need unique priorities: `merge_drm_configs(..., +> auto_priority=True)` in `models/drm/drm_config.py`. +> +> Open items and the full change log: `docs/provider-v2/`. + Copy this directory to `providers/{your_provider}/`, rename the classes, and fill in the stubs. Read this file first — it explains the contract. @@ -1071,9 +1100,10 @@ New providers using the dedicated-manager architecture implement folded architecture override `get_channel_drm` and/or `get_vod_drm` and leave `get_drm_configs` alone. -`StreamingProvider.get_drm(content_id, content_type=None)` is the public -method callers use; it dispatches to whichever architecture the provider -chose. +`get_drm(content_id, drm_variant=None, content_type=None, **kw)` is the +public method callers use (v1 providers: `get_drm(content_id, +content_type=None)` -- see the v2 notice at the top); it dispatches to +whichever architecture the provider chose. ### content_type hint semantics diff --git a/lib/streaming_providers/providers/_template/provider.py b/lib/streaming_providers/providers/_template/provider.py index 24e1929..610535b 100644 --- a/lib/streaming_providers/providers/_template/provider.py +++ b/lib/streaming_providers/providers/_template/provider.py @@ -1,40 +1,50 @@ -# streaming_providers/providers/_template/provider.py +# streaming_providers/providers/_template/provider.py (v2 template) """ {TODO: Provider name} orchestrator. -Owns shared resources (http_manager, caches, auth, managers) and exposes -the public StreamingProvider interface. +Owns shared resources (http_manager, caches, auth, managers) and declares +which managers exist. Everything generic -- capability flags, content_id +routing, DRM validation, header delegation, EPG delegation -- is inherited +from ManagedProvider (base/managed_provider.py), which itself subclasses +StreamingProvider, so the registry and the backend see an ordinary provider. -Subclasses the existing StreamingProvider. Does NOT subclass any new -base class. Authentication is lazy -- no network I/O in __init__. +Authentication is lazy -- no network I/O in __init__. -There are seven optional manager ABCs (channels, vod, epg, recordings, -favorites, bookmarks, catchup) plus an optional DRM manager (a protocol, -not an ABC). This template shows the shape for a provider that has -channels, VOD, and EPG. Delete the factories for capabilities you don't -have, or return None from them. +What you write +-------------- + * provider_name, class metadata (PROVIDER_LABEL, ...) + * __init__: http manager, auth, provider-owned caches, then + ``self._init_managers()`` + * the _build_*() factories for the capabilities you have (the rest + inherit "return None" = capability absent) + * get_manifest / get_drm ONLY if you have provider-specific prefixes that + need parsed arguments (catchup timestamps, ...); call super() for the rest + +What you no longer write (vs. the v1 template) +---------------------------------------------- + * the eight implements_* properties + * _route() + * default get_manifest / get_drm routing + * get_channels / get_epg / get_epg_grid / get_program_details delegation + * header delegation to the managers Plugin name: the registry derives it from the CLASS NAME via -`cls.__name__.lower().replace("provider", "")`. `YourProvider` becomes -"your". Name the class so that this matches your plugin directory, e.g. -`SimpliTVProvider` -> "simplitv" -> providers/simplitv/. +`cls.__name__.lower().replace("provider", "")`. Keep it consistent with +provider_name (see docs/provider-v2/TODO.md, item N-1). """ -from datetime import datetime -from typing import Any, Callable, ClassVar, Dict, List, Optional, Tuple +from typing import ClassVar, Dict, List, Optional -from ...base.errors import NotFoundError +from ...base.managed_provider import ManagedProvider from ...base.managers import ChannelManager, VodManager from ...base.models.proxy_models import ProxyConfig from ...base.protocols import DrmManagerProtocol -from ...base.provider import StreamingProvider from ...base.utils.logger import logger from ...base.vod import VodPage from .auth import YourProviderAuth from .channel_manager import YourChannelManager from .constants import YourConfig, YourDefaults -# from .channel_manager import parse_catchup_id # if you route catchup # from .vod_manager import YourVodManager # from .epg_manager import YourEpgManager # from .recordings_manager import YourRecordingsManager @@ -44,61 +54,36 @@ from .constants import YourConfig, YourDefaults # from .drm_manager import YourDrmManager -class YourProvider(StreamingProvider): +class YourProvider(ManagedProvider): """{TODO: provider name} streaming provider.""" - # ------------------------------------------------------------------ - # provider_name -- ABSTRACT, must be implemented - # ------------------------------------------------------------------ - # - # `provider_name` is declared as an @property @abstractmethod on - # StreamingProvider. If this class does not override it, Python - # raises TypeError at instantiation: - # - # Can't instantiate abstract class YourProvider with abstract - # method provider_name - # - # The registry catches that exception, logs it at ERROR level, and - # skips the provider. The symptom is a provider that is registered - # but never appears in the UI. - # - # The value is the machine identifier: lowercase, no spaces, matching - # the plugin directory name and the PROVIDER_NAME constant in - # constants.py. Used in settings keys, log lines, and the `provider` - # field on models. Return the constant so the two cannot drift. - # - # Do not delete this property. Override the return value; do not - # replace it with a class attribute. + # provider_name is ABSTRACT on StreamingProvider. If it is missing the + # registry logs "Can't instantiate abstract class ..." (with traceback + # since the registry fix) and the provider never shows up in the UI. + # Machine identifier: lowercase, no spaces, == plugin directory == + # PROVIDER_NAME in constants.py. @property def provider_name(self) -> str: return YourDefaults.PROVIDER_NAME - # ------------------------------------------------------------------ - # Class metadata (read by the registry BEFORE any instance exists) - # ------------------------------------------------------------------ - + # Class metadata, read by the registry BEFORE any instance exists. PROVIDER_LABEL: ClassVar[str] = "TODO: display label" PROVIDER_LOGO: ClassVar[str] = YourDefaults.PROVIDER_LOGO SUPPORTED_AUTH_TYPES: ClassVar[List[str]] = ["user_credentials"] - # ALWAYS set SUPPORTED_COUNTRIES. Never leave it at the base - # default (an empty list), which has a specific meaning: "no - # country concept at all." Even a single-country provider declares - # a one-element list. - # - # Single country: ["AT"] - # Multi-country: ["hr", "pl", "me", "at", "hu"] - # Wildcard: ["*"] (country discovered at runtime; - # also the right answer for "not sure yet") - # - # See the README's "SUPPORTED_COUNTRIES is not optional" section. + # ALWAYS set SUPPORTED_COUNTRIES (never the empty base default): + # ["AT"] single | ["hr", "pl"] multi (one instance per country) | ["*"] SUPPORTED_COUNTRIES: ClassVar[List[str]] = ["TODO"] - # Only "live" and "vod" narrow the folded DRM search; anything else - # (None, "event", "catchup", a typo) tries both domains. See the - # README's "content_type hint semantics". - _LIVE_ONLY_CONTENT_TYPES = frozenset({"live"}) - _VOD_ONLY_CONTENT_TYPES = frozenset({"vod"}) + # DRM architecture, declared explicitly (no override sniffing): + # dedicated manager -> return it from _build_drm() + # folded -> set True and override get_channel_drm / get_vod_drm + # Setting both raises ConfigurationError at construction. + DRM_IN_MANAGERS: ClassVar[bool] = False + + # True (default): header hooks on the managers are honoured and the + # default is auth.build_headers(). Set False to keep the legacy "{}". + HEADERS_FROM_MANAGERS: ClassVar[bool] = True def __init__( self, @@ -112,15 +97,13 @@ class YourProvider(StreamingProvider): super().__init__(country) # Unknown kwargs are tolerated (the registry may pass host-level - # extras) but never silently: a typo here is otherwise invisible. + # extras) but never silently. if kwargs: logger.debug( - f"{self.provider_name}: ignoring unknown kwargs " - f"{sorted(kwargs)}" + f"{self.provider_name}: ignoring unknown kwargs {sorted(kwargs)}" ) - config = config or {} - self.config = YourConfig(config, country=self.country) + self.config = YourConfig(config or {}, country=self.country) # 1. HTTP manager. self.http_manager = self._setup_http_manager( @@ -130,46 +113,29 @@ class YourProvider(StreamingProvider): timeout=self.config.timeout, ) - # 2. Auth (lazy -- no network call in __init__). - # - # Providers WITHOUT auth: _build_auth returns None (accept the - # manager base constructors' AuthProtocol warning) or a minimal - # stub with the three token methods. See the README's - # "Providers without auth" section. + # 2. Auth (lazy). Providers WITHOUT auth: return None or a minimal + # stub (README, "Providers without auth"). self._credentials = credentials self.auth = self._build_auth(settings_manager) - # 3. Provider-owned caches. Managers borrow these by reference. - # Add caches here as the provider needs them. + # 3. Provider-owned caches; managers borrow them by reference. + # (Plain dicts are not thread-safe for check-then-set; see + # TODO item C-3.) self._channels_cache: Dict = {} self._playback_cache: Dict = {} - # 4. Managers, in DEPENDENCY ORDER. Every factory returns a - # manager or None. Catchup (and a dedicated DRM manager) may - # need channels/epg/vod, so they are built after them. Do not - # introduce cycles between managers. - self.channels = self._build_channels() - self.vod = self._build_vod() - self.epg = self._build_epg() - self.recordings = self._build_recordings() - self.favorites = self._build_favorites() - self.bookmarks = self._build_bookmarks() - self.catchup = self._build_catchup() # sees channels + epg - self.drm = self._build_drm() + # 4. Build the managers in dependency order (catchup after + # channels/epg). NOTE: this sets self.channels etc. to MANAGERS, + # shadowing the legacy list attribute; ManagedProvider overrides + # to_output_format() for that reason. + self._init_managers() # ------------------------------------------------------------------ - # Factory methods + # Factories -- override only the capabilities you have # ------------------------------------------------------------------ def _build_auth(self, settings_manager): - """ - Return the provider's Auth instance, or None if the provider - needs no authentication. - - `credentials=` MUST be forwarded: it is credentials source #1 - (constructor argument). The auth class falls back to the - settings_manager and CredentialManager on its own. - """ + """`credentials=` MUST be forwarded (credentials source #1).""" return YourProviderAuth( http_manager=self.http_manager, country=self.country, @@ -179,13 +145,6 @@ class YourProvider(StreamingProvider): ) def _build_channels(self) -> Optional[ChannelManager]: - """ - Return a ChannelManager, or None if the provider has no live - channels. - - A VOD-only provider returns None here. A free linear-only - provider returns a manager here and None from _build_vod. - """ return YourChannelManager( http_manager=self.http_manager, auth=self.auth, @@ -195,194 +154,45 @@ class YourProvider(StreamingProvider): ) def _build_vod(self) -> Optional[VodManager]: - """ - Return a VodManager, or None if the provider has no browseable - VOD catalogue. - """ - # TODO: return YourVodManager( - # http_manager=self.http_manager, - # auth=self.auth, - # country=self.country, - # config=self.config, - # playback_cache=self._playback_cache, - # ) + # TODO: return YourVodManager(http_manager=..., auth=..., country=..., + # config=..., playback_cache=...) return None - def _build_epg(self): - """ - Return an EpgManager, or None if the provider has no EPG. - - Some providers have channels but no EPG; some have EPG but no - channels. The two capabilities are independent. - """ - return None - - def _build_recordings(self): - """Return a RecordingsManager, or None.""" - return None - - def _build_favorites(self): - """Return a FavoritesManager, or None.""" - return None - - def _build_bookmarks(self): - """Return a BookmarksManager, or None.""" - return None - - def _build_catchup(self): - """ - Return a CatchupManager, or None. - - Catchup usually needs the channel manager and/or the EPG manager - as collaborators. Pass them as explicit keyword-only arguments; - self.channels and self.epg are already built when this runs. - """ - # TODO: return YourCatchupManager( - # http_manager=self.http_manager, - # auth=self.auth, - # country=self.country, - # config=self.config, - # channels=self.channels, - # epg=self.epg, - # ) - return None + # _build_epg / _build_recordings / _build_favorites / _build_bookmarks / + # _build_catchup: inherited "return None". Override to enable, e.g. + # + # def _build_catchup(self): + # return YourCatchupManager( + # http_manager=self.http_manager, auth=self.auth, + # country=self.country, config=self.config, + # channels=self.channels, epg=self.epg, # already built + # ) def _build_drm(self) -> Optional[DrmManagerProtocol]: - """ - Return a dedicated DRM manager, or None. - - Two supported architectures (see the README's "DRM" section): - - * Dedicated manager: return a class matching - DrmManagerProtocol here; the provider's get_drm() delegates - to it. - - * Folded into managers: leave this returning None, and - override get_channel_drm() on your ChannelManager and/or - get_vod_drm() on your VodManager. The provider's get_drm() - falls back to routing to those. - - Rule: fold if DRM shares state with the manifest step; otherwise - use the dedicated manager. See providers/_template/drm_manager.py - for the four existing source patterns. - """ + """Dedicated DRM manager, or None (folded / no DRM). See README "DRM".""" return None # ------------------------------------------------------------------ - # Capability flags (derived from manager presence) + # Provider-specific routing (only if you need it) # ------------------------------------------------------------------ - - @property - def implements_channels(self) -> bool: - return self.channels is not None - - @property - def implements_vod(self) -> bool: - return self.vod is not None - - @property - def implements_epg(self) -> bool: - return self.epg is not None - - @property - def implements_recordings(self) -> bool: - return self.recordings is not None - - @property - def implements_favorites(self) -> bool: - return self.favorites is not None - - @property - def implements_bookmarks(self) -> bool: - return self.bookmarks is not None - - @property - def implements_catchup(self) -> bool: - return self.catchup is not None - - @property - def implements_drm(self) -> bool: - """ - True when this provider can produce DRM configurations. - - Derived from either source of DRM: - * a dedicated DRM manager (_build_drm returned a class), OR - * a ChannelManager that overrides get_channel_drm, OR - * a VodManager that overrides get_vod_drm. - - The base classes' defaults return []; we detect overrides by - comparing the bound method against the base class's method. - """ - if self.drm is not None: - return True - - folded_channels = ( - self.channels is not None - and type(self.channels).get_channel_drm - is not ChannelManager.get_channel_drm - ) - folded_vod = ( - self.vod is not None - and type(self.vod).get_vod_drm - is not VodManager.get_vod_drm - ) - return folded_channels or folded_vod - - # ------------------------------------------------------------------ - # Router - # ------------------------------------------------------------------ - - def _route(self, content_id: str, attempts: List[Tuple[Any, Callable]]): - """ - Try managers in order, using handles_content_id() to skip those - that declare they don't handle the id. - - Distinguishes three outcomes: - * manager returned a truthy result -> return it - * manager returned None / [] / falsy -> try next manager - * manager raised NotFoundError -> remember it, try next - - If nobody resolved and a NotFoundError was seen, re-raise it -- - that's "the content existed in some manager's domain but is gone", - distinct from "nobody handles this id at all" (which returns None). - - Any OTHER exception propagates. In particular BadRequestError is - deliberately NOT caught: if a manager uses a 400 as an endpoint - dispatch signal, override handles_content_id() on that manager - instead of relying on try-and-catch. When both channels and vod - exist, override handles_content_id() on both -- the default - (True) makes the first manager see every id. - """ - last_not_found: Optional[NotFoundError] = None - for manager, call in attempts: - if manager is None or not manager.handles_content_id(content_id): - continue - try: - result = call(manager) - except NotFoundError as e: - last_not_found = e - continue - if result: - return result - if last_not_found is not None: - raise last_not_found - return None - - # ------------------------------------------------------------------ - # Public delegations # - # The optional capabilities (recordings, favorites, bookmarks, - # catchup) are exposed by the base layer's Provider*Mixin classes, - # which correspond one-to-one to those managers -- no delegation is - # needed here. Channels, VOD and EPG are delegated explicitly below. - # VERIFY these three signatures against StreamingProvider when you - # copy the template; remove any the base class already provides. - # ------------------------------------------------------------------ + # def get_manifest(self, content_id: str, **kw) -> Optional[str]: + # if content_id.startswith("catchup:"): + # parsed = parse_catchup_id(content_id) # raises BadRequestError + # return self.catchup.get_catchup_manifest( + # parsed.content_id, parsed.start_time, parsed.end_time, **kw + # ) if self.catchup else None + # return super().get_manifest(content_id, **kw) + # + # Do NOT fall back to the live manifest when catchup fails. - def get_channels(self, **kw): - if self.channels is None: - return [] - return self.channels.get_channels(**kw) + # ------------------------------------------------------------------ + # VOD delegation + # ------------------------------------------------------------------ + # ManagedProvider delegates channels, EPG, manifest, headers and DRM. + # VOD navigation is NOT delegated yet because the legacy + # ProviderVodMixin signature has not been verified (TODO item M-2). + # Until then providers with VOD keep this delegation themselves. def get_vod_category( self, @@ -395,80 +205,4 @@ class YourProvider(StreamingProvider): return VodPage() return self.vod.get_vod_category( content_id, cursor=cursor, page_size=page_size, **kw - ) - - def get_epg( - self, - channel_id: str, - start_time: Optional[datetime] = None, - end_time: Optional[datetime] = None, - **kw, - ): - if self.epg is None: - return [] - return self.epg.get_epg( - channel_id, start_time=start_time, end_time=end_time, **kw - ) - - def get_manifest(self, content_id: str, **kw) -> Optional[str]: - """ - Return the manifest URL for the given content, routing by manager. - - Route through `_route` when the manager needs only the - content_id. Add an explicit prefix branch ABOVE `_route` when the - manager needs parsed arguments (a timestamp, an episode index). - One parser, in the router -- see the README's "Parsers vs. - dispatch". Catchup is the usual example: - - # if content_id.startswith("catchup:"): - # parsed = parse_catchup_id(content_id) # raises - # return self.catchup.get_catchup_manifest( # BadRequestError - # parsed.content_id, parsed.start_time, parsed.end_time, - # **kw, - # ) if self.catchup else None - - Do NOT fall back to the live manifest when catchup fails. - """ - return self._route(content_id, [ - (self.channels, lambda m: m.get_channel_manifest( - content_id, **kw - )), - (self.vod, lambda m: m.get_vod_manifest(content_id, **kw)), - ]) - - def get_drm( - self, - content_id: str, - content_type: Optional[str] = None, - **kw, - ) -> List: - """ - Return DRM configuration(s) for the given content. - - content_type is an optional hint. Only "live" and "vod" narrow - the search on the folded path; any other value (including None, - "event", "catchup", or an unrecognized string) tries both - domains. This is deliberate: a wrong narrowing produces a silent - [] for protected content, which is the hardest kind of bug to - trace. Widening on unknown input is always safe. - - If a dedicated DRM manager is configured, the hint is passed - through unchanged and the manager decides what to do with it. - """ - if self.drm is not None: - return self.drm.get_drm_configs( - content_id, content_type=content_type, **kw - ) - - attempts: List[Tuple[Any, Callable]] = [] - if content_type not in self._VOD_ONLY_CONTENT_TYPES: - attempts.append( - (self.channels, lambda m: m.get_channel_drm( - content_id, **kw - )) - ) - if content_type not in self._LIVE_ONLY_CONTENT_TYPES: - attempts.append( - (self.vod, lambda m: m.get_vod_drm(content_id, **kw)) - ) - return self._route(content_id, attempts) or [] \ No newline at end of file + ) \ No newline at end of file diff --git a/lib/streaming_providers/providers/joyn/channel_manager.py b/lib/streaming_providers/providers/joyn/channel_manager.py index 79be3af..67d4a37 100644 --- a/lib/streaming_providers/providers/joyn/channel_manager.py +++ b/lib/streaming_providers/providers/joyn/channel_manager.py @@ -9,7 +9,8 @@ import json import time import urllib.parse from base64 import b64decode -from typing import Dict, List, Optional +from typing import Dict, List, Optional, Tuple +from urllib.parse import urlparse from ...base.models import DRMConfig, DRMSystem, LicenseConfig, StreamingChannel from ...base.provider import AuthType @@ -41,7 +42,13 @@ from .constants import ( MODE_VOD, SIGNATURE_SECRET_KEY, ) -from .models import JoynChannel, JoynError, JoynEntitlementError, PlaybackRestrictedException, SubscriptionRequiredException +from .models import ( + JoynChannel, + JoynError, + JoynEntitlementError, + PlaybackRestrictedException, + SubscriptionRequiredException, +) def create_video_payload(config: Optional[Dict] = None, compact: bool = True) -> str: @@ -78,6 +85,10 @@ class JoynChannelManager: self._cache_timestamp: float = 0.0 self._cache_ttl: int = 300 # 5 minutes + # channel_id -> resolved_id (usually "-hd"). Populated lazily by + # get_channel_entitlement_token so we don't re-probe on every playback. + self._resolved_channel_variants: Dict[str, str] = {} + logger.info(f"[JoynChannelManager] Initialised for country={provider.country}") @property @@ -108,6 +119,10 @@ class JoynChannelManager: self._cache_timestamp = time.time() return self._channels_cache or [] + # ======================================================================== + # HEADERS + # ======================================================================== + def _get_graphql_headers(self) -> Dict[str, str]: return self.provider._build_provider_headers( base_headers=JOYN_GRAPHQL_BASE_HEADERS, @@ -133,11 +148,308 @@ class JoynChannelManager: "joyn-distribution-tenant": self.distribution_tenant, "joyn-platform": self.platform, "joyn-b2b-context": "UNKNOWN", - "joyn-client-os": "UNKNOWN", # Restored missing header + "joyn-client-os": "UNKNOWN", "origin": JOYN_DOMAINS.get(self.country, JOYN_DOMAINS["de"]), }, ) + def _get_entitlement_headers(self) -> Dict[str, str]: + """ + Headers for the entitlement host. + + Entitlement lives on a *separate* host from the Joyn GraphQL/streaming + APIs and does not accept the joyn-* header set. Sending a minimal + header set (Authorization + Content-Type + UA) matches the working + reference client. + + The bearer is fetched via the authenticator (not read from + provider.bearer_token), so a long-running session picks up refreshes + automatically instead of sending a stale token. + """ + token = "" + if self.authenticator is not None: + try: + token = self.authenticator.get_bearer_token() or "" + except Exception as e: + logger.warning(f"Could not obtain bearer for entitlement call: {e}") + + headers = { + "Content-Type": "application/json", + "Accept": "application/json", + "User-Agent": "Mozilla/5.0", + } + if token: + headers["Authorization"] = f"Bearer {token}" + return headers + + def get_manifest_headers(self, content_id: str, **kwargs) -> Dict[str, str]: + # The CDN-served manifest URL is self-authorizing; sending the Joyn + # provider bearer token to the CDN causes: + # 400 InvalidArgument: Unsupported Authorization Type + # so we deliberately omit Authorization here. + return self.provider._build_provider_headers( + base_headers={}, + auth_type=AuthType.NONE, + provider_headers={ + "User-Agent": JOYN_USER_AGENT, + "Origin": JOYN_DOMAINS.get(self.country, JOYN_DOMAINS["de"]), + }, + ) + + # ======================================================================== + # ENTITLEMENT + # ======================================================================== + + def get_entitlement_token(self, content_id: str, content_type: str = CONTENT_TYPE_LIVE) -> str: + headers = self._get_entitlement_headers() + payload = {"content_id": content_id, "content_type": content_type} + url = JOYN_STREAMING_ENDPOINTS["ENTITLEMENT"] + host = urlparse(url).netloc + + try: + response = self.http_manager.post( + url, + operation="auth", + headers=headers, + json_data=payload, + timeout=DEFAULT_REQUEST_TIMEOUT, + ) + + if response.status_code == 400: + try: + error_data = response.json() + if isinstance(error_data, list) and len(error_data) > 0: + error = error_data[0] + code = error.get("code", "UNKNOWN") + msg = error.get("msg", "No error message provided") + if code == ERROR_CODES["PLAYBACK_RESTRICTED"]: + raise PlaybackRestrictedException( + f"Playback restricted for {content_id}: {msg}" + ) + elif code == ERROR_CODES["BUSINESS_MODEL_NOT_SUITABLE"]: + raise SubscriptionRequiredException( + f"Subscription required for {content_id} ({code}): {msg}" + ) + else: + raise JoynEntitlementError( + f"Entitlement error for {content_id} ({code}): {msg}" + ) + except (json.JSONDecodeError, KeyError, IndexError) as e: + logger.warning( + f"Entitlement 400 for {content_id} from {host}: " + f"failed to parse error body: {e}" + ) + raise JoynEntitlementError( + f"Bad response for {content_id} (400), failed to parse error: {e}" + ) + + if response.status_code >= 400: + logger.warning( + f"Entitlement failed for {content_id} (type={content_type}): " + f"HTTP {response.status_code} from {host}" + ) + response.raise_for_status() + + data = response.json() + # Working reference accepts either key; keep both for compatibility. + token = data.get("entitlement_token") or data.get("token") + if not token: + logger.warning( + f"Entitlement response for {content_id} from {host} " + f"had no entitlement_token: keys={list(data.keys())}" + ) + raise JoynEntitlementError(f"No entitlement_token in response for {content_id}") + return token + + except PlaybackRestrictedException: + raise + except SubscriptionRequiredException: + raise + except JoynEntitlementError: + raise + except Exception as e: + raise JoynEntitlementError(f"Error getting entitlement token for {content_id}: {e}") + + def get_channel_entitlement_token(self, channel_id: str) -> Tuple[str, str]: + """ + Resolve entitlement for a live channel, trying the -hd variant first. + + Joyn's live channels are indexed with an "-hd" suffix at the + entitlement service; asking for the bare slug returns no token for + HD-only streams. Mirrors the working reference's retry order. + + Results are cached per original channel_id so the extra probe only + happens once per channel per process. + + Returns: + (resolved_channel_id, entitlement_token) + + Raises: + PlaybackRestrictedException / SubscriptionRequiredException for + account/rights errors (terminal — not retried on the other variant). + JoynEntitlementError if neither variant yields a token. + """ + # Fast path: we already know which variant works for this channel. + cached_resolved = self._resolved_channel_variants.get(channel_id) + if cached_resolved: + try: + token = self.get_entitlement_token( + content_id=cached_resolved, content_type=CONTENT_TYPE_LIVE + ) + if token: + return cached_resolved, token + except (PlaybackRestrictedException, SubscriptionRequiredException): + # Rights changed since we cached — let it propagate. + raise + except JoynEntitlementError: + # Cached variant no longer works; drop it and re-probe below. + logger.debug(f"Cached variant {cached_resolved} no longer resolves; re-probing") + self._resolved_channel_variants.pop(channel_id, None) + + candidates: List[str] = [] + if channel_id.endswith("-sd"): + candidates.append(channel_id[:-3] + "-hd") + elif not channel_id.endswith("-hd"): + candidates.append(channel_id + "-hd") + candidates.append(channel_id) + + last_error: Optional[Exception] = None + for cid in candidates: + try: + token = self.get_entitlement_token( + content_id=cid, content_type=CONTENT_TYPE_LIVE + ) + if token: + self._resolved_channel_variants[channel_id] = cid + return cid, token + except (PlaybackRestrictedException, SubscriptionRequiredException): + # Rights errors are terminal — do not retry the next candidate. + raise + except JoynEntitlementError as e: + last_error = e + continue + + raise last_error or JoynEntitlementError( + f"No entitlement token for {channel_id} (tried {candidates})" + ) + + # ======================================================================== + # PLAYLIST / MANIFEST / DRM + # ======================================================================== + + def get_channel_playlist( + self, + channel_id: str, + entitlement_token: str, + video_config: Optional[Dict] = None, + ) -> Dict: + video_payload = create_video_payload(video_config) + signature = build_signature(entitlement_token, video_payload) + + url = JOYN_STREAMING_ENDPOINTS["PLAYLIST"].format(channel_id=channel_id) + url += f"?signature={signature}" + + headers = JOYN_API_BASE_HEADERS.copy() + headers["Authorization"] = f"Bearer {entitlement_token}" + + try: + response = self.http_manager.post( + url, + operation="manifest", + headers=headers, + data=video_payload, + timeout=DEFAULT_REQUEST_TIMEOUT, + ) + response.raise_for_status() + return response.json() + except Exception as e: + # Let our custom entitlement exceptions bubble up untouched so callers + # (and the UI layer) can distinguish "needs subscription" / "not allowed + # here" from a generic network failure instead of seeing everything as + # a flat JoynError. + if isinstance( + e, + (PlaybackRestrictedException, SubscriptionRequiredException, JoynEntitlementError), + ): + raise + raise JoynError(f"Error getting playlist for {channel_id}: {e}") + + def get_manifest( + self, + content_id: str, + content_type: str = CONTENT_TYPE_LIVE, + video_config: Optional[Dict] = None, + **kwargs, + ) -> Optional[str]: + try: + if content_type == CONTENT_TYPE_LIVE: + resolved_id, entitlement_token = self.get_channel_entitlement_token(content_id) + else: + entitlement_token = self.get_entitlement_token( + content_id=content_id, content_type=content_type + ) + resolved_id = content_id + playlist_data = self.get_channel_playlist(resolved_id, entitlement_token, video_config) + return playlist_data.get("manifestUrl") + except Exception as e: + logger.error(f"Error getting manifest for channel {content_id}: {e}") + return None + + def _build_drm_config(self, playlist_data: Dict) -> Optional[DRMConfig]: + """Build a DRMConfig object from a playlist response. + + The license endpoint authenticates via the token embedded in the URL's + signature query param — no Authorization header is sent. We include + Origin and User-Agent to satisfy Cloudflare WAF requirements, matching + captured browser traffic. + """ + license_url = playlist_data.get("licenseUrl") + if not license_url: + return None + + return DRMConfig( + system=DRMSystem.WIDEVINE, + priority=1, + license=LicenseConfig( + server_url=license_url, + server_certificate=playlist_data.get("certificateUrl"), + req_headers=json.dumps({ + "User-Agent": JOYN_USER_AGENT, + "Origin": JOYN_DOMAINS.get(self.country, JOYN_DOMAINS["de"]), + "Content-Type": DRM_REQUEST_HEADERS["Content-Type"], + }), + req_data="{CHA-RAW}", + use_http_get_request=False, + ), + ) + + def get_drm( + self, + content_id: str, + content_type: str = CONTENT_TYPE_LIVE, + video_config: Optional[Dict] = None, + **kwargs, + ) -> List[DRMConfig]: + try: + if content_type == CONTENT_TYPE_LIVE: + resolved_id, entitlement_token = self.get_channel_entitlement_token(content_id) + else: + entitlement_token = self.get_entitlement_token( + content_id=content_id, content_type=content_type + ) + resolved_id = content_id + playlist_data = self.get_channel_playlist(resolved_id, entitlement_token, video_config) + + drm_config = self._build_drm_config(playlist_data) + return [drm_config] if drm_config else [] + except Exception as e: + logger.error(f"Error getting DRM configs for channel {content_id}: {e}") + return [] + + # ======================================================================== + # CHANNEL FETCHING + # ======================================================================== + def get_channels( self, time_window_hours: int = DEFAULT_EPG_WINDOW_HOURS, @@ -229,162 +541,17 @@ class JoynChannelManager: if stream_data.get("eventStream", False): joyn_channel.raw_data["is_event_stream"] = True - channels.append(joyn_channel.to_streaming_channel(provider_name=self.provider.provider_name)) + channels.append( + joyn_channel.to_streaming_channel(provider_name=self.provider.provider_name) + ) except Exception as e: logger.warning(f"Error processing channel data: {e}") return channels - def get_entitlement_token(self, content_id: str, content_type: str = CONTENT_TYPE_LIVE) -> str: - headers = self.get_api_headers() - payload = {"content_id": content_id, "content_type": content_type} - - try: - response = self.http_manager.post( - JOYN_STREAMING_ENDPOINTS["ENTITLEMENT"], - operation="auth", - headers=headers, - json_data=payload, - timeout=DEFAULT_REQUEST_TIMEOUT, - ) - - if response.status_code == 400: - try: - error_data = response.json() - if isinstance(error_data, list) and len(error_data) > 0: - error = error_data[0] - code = error.get("code", "UNKNOWN") - msg = error.get("msg", "No error message provided") - if code == ERROR_CODES["PLAYBACK_RESTRICTED"]: - raise PlaybackRestrictedException(f"Playback restricted for {content_id}: {msg}") - elif code == ERROR_CODES["BUSINESS_MODEL_NOT_SUITABLE"]: - raise SubscriptionRequiredException( - f"Subscription required for {content_id} ({code}): {msg}") - else: - raise JoynEntitlementError(f"Entitlement error for {content_id} ({code}): {msg}") - except (json.JSONDecodeError, KeyError, IndexError) as e: - raise JoynEntitlementError(f"Bad response for {content_id} (400), failed to parse error: {e}") - - response.raise_for_status() - data = response.json() - return data["entitlement_token"] - - except PlaybackRestrictedException: - raise - except JoynEntitlementError: - raise - except KeyError: - raise JoynEntitlementError(f"No entitlement_token in response for {content_id}") - except Exception as e: - raise JoynEntitlementError(f"Error getting entitlement token for {content_id}: {e}") - - def get_channel_playlist( - self, - channel_id: str, - entitlement_token: str, - video_config: Optional[Dict] = None, - ) -> Dict: - video_payload = create_video_payload(video_config) - signature = build_signature(entitlement_token, video_payload) - - url = JOYN_STREAMING_ENDPOINTS["PLAYLIST"].format(channel_id=channel_id) - url += f"?signature={signature}" - - headers = JOYN_API_BASE_HEADERS.copy() - headers["Authorization"] = f"Bearer {entitlement_token}" - - try: - response = self.http_manager.post( - url, - operation="manifest", - headers=headers, - data=video_payload, - timeout=DEFAULT_REQUEST_TIMEOUT, - ) - response.raise_for_status() - return response.json() - except Exception as e: - # Let our custom entitlement exceptions bubble up untouched so callers - # (and the UI layer) can distinguish "needs subscription" / "not allowed - # here" from a generic network failure instead of seeing everything as - # a flat JoynError. - if isinstance(e, (PlaybackRestrictedException, SubscriptionRequiredException, JoynEntitlementError)): - raise - raise JoynError(f"Error getting playlist for {channel_id}: {e}") - - def get_manifest( - self, - content_id: str, - content_type: str = CONTENT_TYPE_LIVE, - video_config: Optional[Dict] = None, - **kwargs, - ) -> Optional[str]: - try: - entitlement_token = self.get_entitlement_token(content_id=content_id, content_type=content_type) - playlist_data = self.get_channel_playlist(content_id, entitlement_token, video_config) - return playlist_data.get("manifestUrl") - except Exception as e: - logger.error(f"Error getting manifest for channel {content_id}: {e}") - return None - - def get_manifest_headers(self, content_id: str, **kwargs) -> Dict[str, str]: - # The CDN-served manifest URL is self-authorizing; sending the Joyn - # provider bearer token to the CDN causes: - # 400 InvalidArgument: Unsupported Authorization Type - # so we deliberately omit Authorization here. - return self.provider._build_provider_headers( - base_headers={}, - auth_type=AuthType.NONE, - provider_headers={ - "User-Agent": JOYN_USER_AGENT, - "Origin": JOYN_DOMAINS.get(self.country, JOYN_DOMAINS["de"]), - }, - ) - - def _build_drm_config(self, playlist_data: Dict) -> Optional[DRMConfig]: - """Build a DRMConfig object from a playlist response. - - The license endpoint authenticates via the token embedded in the URL's - signature query param — no Authorization header is sent. We include - Origin and User-Agent to satisfy Cloudflare WAF requirements, matching - captured browser traffic. - """ - license_url = playlist_data.get("licenseUrl") - if not license_url: - return None - - return DRMConfig( - system=DRMSystem.WIDEVINE, - priority=1, - license=LicenseConfig( - server_url=license_url, - server_certificate=playlist_data.get("certificateUrl"), - req_headers=json.dumps({ - "User-Agent": JOYN_USER_AGENT, - "Origin": JOYN_DOMAINS.get(self.country, JOYN_DOMAINS["de"]), - "Content-Type": DRM_REQUEST_HEADERS["Content-Type"], - }), - req_data="{CHA-RAW}", - use_http_get_request=False, - ), - ) - - def get_drm( - self, - content_id: str, - content_type: str = CONTENT_TYPE_LIVE, - video_config: Optional[Dict] = None, - **kwargs, - ) -> List[DRMConfig]: - try: - entitlement_token = self.get_entitlement_token(content_id=content_id, content_type=content_type) - playlist_data = self.get_channel_playlist(content_id, entitlement_token, video_config) - - drm_config = self._build_drm_config(playlist_data) - return [drm_config] if drm_config else [] - except Exception as e: - logger.error(f"Error getting DRM configs for channel {content_id}: {e}") - return [] + # ======================================================================== + # ENRICHMENT + # ======================================================================== def enrich_channel_data( self, @@ -394,8 +561,18 @@ class JoynChannelManager: ) -> Optional[StreamingChannel]: try: content_type = kwargs.get("content_type", channel.content_type) - entitlement_token = self.get_entitlement_token(content_id=channel.channel_id, content_type=content_type) - playlist_data = self.get_channel_playlist(channel.channel_id, entitlement_token, video_config) + if content_type == CONTENT_TYPE_LIVE: + resolved_id, entitlement_token = self.get_channel_entitlement_token( + channel.channel_id + ) + else: + entitlement_token = self.get_entitlement_token( + content_id=channel.channel_id, content_type=content_type + ) + resolved_id = channel.channel_id + playlist_data = self.get_channel_playlist( + resolved_id, entitlement_token, video_config + ) manifest_url = playlist_data.get("manifestUrl") if not manifest_url: @@ -430,11 +607,17 @@ class JoynChannelManager: while retries < max_retries and not success and not is_restricted: try: - entitlement_token = self.get_entitlement_token( - content_id=channel.channel_id, content_type=channel.content_type - ) + if channel.content_type == CONTENT_TYPE_LIVE: + resolved_id, entitlement_token = self.get_channel_entitlement_token( + channel.channel_id + ) + else: + entitlement_token = self.get_entitlement_token( + content_id=channel.channel_id, content_type=channel.content_type + ) + resolved_id = channel.channel_id playlist_data = self.get_channel_playlist( - channel.channel_id, entitlement_token, video_config + resolved_id, entitlement_token, video_config ) manifest_url = playlist_data.get("manifestUrl") @@ -450,8 +633,11 @@ class JoynChannelManager: success = True else: raise JoynError("No manifestUrl in response") - except PlaybackRestrictedException as e: - logger.warning(f"Playback restricted for {channel.name}: {e}") + except (PlaybackRestrictedException, SubscriptionRequiredException) as e: + # Both are terminal, per-account/rights failures — do not retry. + logger.warning( + f"Playback restricted/subscription-required for {channel.name}: {e}" + ) is_restricted = True except Exception as e: retries += 1 @@ -460,5 +646,8 @@ class JoynChannelManager: else: logger.error(f"Failed to get streaming data for {channel.name}: {e}") - logger.info(f"Streaming data population complete: {len(successful_channels)}/{len(channels)}") + logger.info( + f"Streaming data population complete: " + f"{len(successful_channels)}/{len(channels)}" + ) return successful_channels \ No newline at end of file diff --git a/lib/streaming_providers/providers/joyn/constants.py b/lib/streaming_providers/providers/joyn/constants.py index 63aed26..ec166e1 100644 --- a/lib/streaming_providers/providers/joyn/constants.py +++ b/lib/streaming_providers/providers/joyn/constants.py @@ -4,6 +4,8 @@ Joyn provider constants - Cleaned and organized """ +import os + # ============================================================================ # Provider Metadata # ============================================================================ @@ -49,8 +51,10 @@ DEVICE_IDS = { # HTTP Headers & User Agent # ============================================================================ -JOYN_USER_AGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/153.0.0.0 Safari/537.36" -JOYN_CLIENT_VERSION = "5.1587.0" +# Bumped to match the working reference client (Chrome 154). +JOYN_USER_AGENT = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/154.0.0.0 Safari/537.36" +# Bumped to match the working reference client. +JOYN_CLIENT_VERSION = "5.1592.3" DEFAULT_PLATFORM = "web" # Base authentication headers (without dynamic values) @@ -105,12 +109,27 @@ GRAPHQL_OFFSET = 0 # Streaming Configuration # ============================================================================ +# Entitlement host is overridable so a future Joyn migration (they already +# moved once, from entitlement.p7s1.io) does not require a code change. +# Set JOYN_ENTITLEMENT_URL in the environment to override. +_DEFAULT_ENTITLEMENT_URL = ( + "https://entitlements-service-alb.prd.platform.s.joyn.de/api/user/entitlement-token" +) +_ENTITLEMENT_URL = os.environ.get("JOYN_ENTITLEMENT_URL", _DEFAULT_ENTITLEMENT_URL).strip() \ + or _DEFAULT_ENTITLEMENT_URL + JOYN_STREAMING_ENDPOINTS = { - "ENTITLEMENT": "https://entitlement.p7s1.io/api/user/entitlement-token", + "ENTITLEMENT": _ENTITLEMENT_URL, "PLAYLIST": "https://api.vod-prd.s.joyn.de/v1/channel/{channel_id}/playlist", } -# Default video configuration for playlist requests +# Default video configuration for playlist requests. +# +# IMPORTANT: the playlist request is signed over the *exact* JSON string built +# from this dict (see create_video_payload). If you change any key or value +# here, the signature sent to api.vod-prd.s.joyn.de will change with it. If +# the server validates the signature against its own expected payload shape, +# a mismatch here produces 403 on every live and VOD playback request. DEFAULT_VIDEO_CONFIG = { "enableDolbyAtmos": True, "enableSubtitles": True, @@ -125,8 +144,16 @@ DEFAULT_VIDEO_CONFIG = { "maxSecurityLevel": 5, } -# Signature secret key (base64 encoded) -SIGNATURE_SECRET_KEY = "MzU0MzM3MzgzMzM4MzMzNjM1NDMzNzM4MzYzNDM2MzYzNTQzMzczODM2MzYzMzM4MzIzNjM1NDMzNzM4MzMzMDM2MzQzNTM5MzU0MzM3MzgzMzM5MzMzNTMyMzQzNTQzMzczODM2MzUzMzM5MzU0MzM3MzgzMzM4MzMzMjMzNDYzNTQzMzczODM2MzYzMzMzMzM0NDMzNDIzNTQzMzczODMzMzgzNjM2MzMzNQ==" +# Signature secret key (base64 encoded). +# +# Decodes to a digit-string secret used as-is. The signing algorithm is: +# sha1(f"{payload_json},{entitlement_token}{secret}") +# matching the working reference client exactly. +# +# This constant is intentionally NOT configurable: it must byte-match the +# value embedded in Joyn's own web client. If Joyn rotates it, the fix is a +# new constant shipped in an update, not a user setting. +SIGNATURE_SECRET_KEY = "MzU0MzM3MzgzMzM4MzMzNjM1NDMzNzM4MzYzNDM2MzYzNTQzMzk3MzgzNjM2MzMzODMyMzYzNTQzMzc3MzgzMzMwMzYzNDM1MzkzNTQzMzc3MzgzMzM5MzMzNTMyMzQzNTQzMzc3MzgzNjM1MzMzOTM1NDMzNzM4MzMzODMzMjMzNDYzNTQzMzc4MzYzNjMzMzMzMzQ0MzM0NDMyNzA2NTQzMzczODMzMzgzNjM2MzMzMw==" # ============================================================================ # Content Types & Modes @@ -165,22 +192,42 @@ ERROR_CODES = { SUPPORTED_COUNTRIES = ["de", "at", "ch"] DEFAULT_COUNTRY = "de" +# GraphQL / streaming API tenant values. +# Used by provider.py and channel_manager._get_graphql_headers. COUNTRY_TENANT_MAPPING = { "de": "JOYN", "at": "JOYN_AT", "ch": "JOYN_CH", } +# Auth (7pass / auth.joyn.de) tenant values. +# +# NOTE: Joyn sends a *different* tenant on auth calls than on GraphQL calls — +# Germany is "JOYN_DE" on auth, but plain "JOYN" on GraphQL. These must stay +# separate; using the GraphQL map on auth calls silently downgrades the +# session to anonymous. +# +# NOT YET WIRED UP: auth.py still reads COUNTRY_TENANT_MAPPING. The switch is +# part of the auth rework batch. Kept here so the two maps stay visible +# together. +AUTH_TENANT_MAPPING = { + "de": "JOYN_DE", + "at": "JOYN_AT", + "ch": "JOYN_CH", +} + JOYN_DOMAINS = { "de": "https://www.joyn.de", "at": "https://www.joyn.at", "ch": "https://www.joyn.ch", } + def get_oauth_redirect_uri(country: str) -> str: """Get country-specific OAuth redirect URI""" return f"https://www.joyn.{country}/oauth" + # ============================================================================ # DRM Configuration # ============================================================================