From 2d39824f0e0b537c84af001709a39320efa7ddeb Mon Sep 17 00:00:00 2001 From: Nirvana Date: Wed, 23 Sep 2026 20:50:53 +0200 Subject: [PATCH] magentaeu: Add VOD --- .../providers/magentaeu/auth.py | 617 +++++++--- .../providers/magentaeu/provider.py | 146 ++- .../providers/magentaeu/vod_errors.py | 152 +++ .../providers/magentaeu/vod_manager.py | 1036 +++++++++++++++++ 4 files changed, 1745 insertions(+), 206 deletions(-) create mode 100644 lib/streaming_providers/providers/magentaeu/vod_errors.py create mode 100644 lib/streaming_providers/providers/magentaeu/vod_manager.py diff --git a/lib/streaming_providers/providers/magentaeu/auth.py b/lib/streaming_providers/providers/magentaeu/auth.py index 697c7b9..f38de5a 100644 --- a/lib/streaming_providers/providers/magentaeu/auth.py +++ b/lib/streaming_providers/providers/magentaeu/auth.py @@ -1,11 +1,13 @@ # streaming_providers/providers/magentaeu/auth.py # -*- coding: utf-8 -*- +from __future__ import annotations + import base64 import json import time import uuid from dataclasses import dataclass, field -from typing import Any, Dict, Optional +from typing import Any, Dict, Optional, Tuple # Updated imports for pycryptodome try: @@ -52,9 +54,12 @@ from .constants import ( ) +# --------------------------------------------------------------------------- +# JWT helpers +# --------------------------------------------------------------------------- + class InvalidTokenError(Exception): """Exception for invalid JWT tokens""" - pass @@ -65,7 +70,9 @@ def base64url_decode(input_str: str) -> bytes: def decode_jwt(token: str, verify: bool = True) -> Dict[str, Any]: - """Decode JWT token""" + """ + Decode a JWT payload. If verify=True, raise InvalidTokenError on expiry. + """ try: header_b64, payload_b64, signature = token.split(".") payload_json = base64url_decode(payload_b64).decode("utf-8") @@ -73,24 +80,37 @@ def decode_jwt(token: str, verify: bool = True) -> Dict[str, Any]: if verify and "exp" in payload: if payload["exp"] < time.time(): - raise InvalidTokenError("Token has expired") + raise InvalidTokenError( + f"Token expired at {payload['exp']} (now {int(time.time())})" + ) return payload - except (ValueError, json.JSONDecodeError, UnicodeDecodeError): - raise InvalidTokenError("Invalid token format") + except InvalidTokenError: + raise + except (ValueError, json.JSONDecodeError, UnicodeDecodeError) as exc: + raise InvalidTokenError(f"Invalid token format: {exc}") def is_token_valid(token: str) -> bool: - """Check if token is valid""" + """ + Check if token is valid. Logs the specific rejection reason, since a + silent False is undiagnosable in the field. + """ if not token: + logger.debug("is_token_valid: empty token") return False try: - decode_jwt(token) + decode_jwt(token, verify=True) return True - except InvalidTokenError: + except InvalidTokenError as exc: + logger.debug(f"is_token_valid: rejected -- {exc}") return False +# --------------------------------------------------------------------------- +# Token +# --------------------------------------------------------------------------- + @dataclass class MagentaAuthToken(BaseAuthToken): """Magenta TV authentication token""" @@ -99,14 +119,27 @@ class MagentaAuthToken(BaseAuthToken): device_id: Optional[str] = field(default="") session_id: Optional[str] = field(default="") channel_map_id: Optional[str] = field(default="") + # Epoch time device_id/session_id were last confirmed via real cookies # from the startup page. 0 means "never validated" -- treated as stale # regardless of GUEST_SESSION_TTL_SECONDS. session_id_updated_at: float = field(default=0.0) + # Access-token lifetime is on BaseAuthToken.expires_in. The refresh-token + # lifetime is surfaced separately here because the HR login response + # reports `refreshExpiresIn` (camelCase) and the base class's + # needs_refresh() consults it -- so it must never be None. See + # _create_token_from_response for the coercion. + refresh_expires_in: int = field(default=0) + + # --- /user/account payload, cached so VOD + entitlement checks don't + # each re-fetch it. Populated lazily by get_user_account(). --- + account_info: Optional[Dict[str, Any]] = field(default=None) + account_info_fetched_at: float = field(default=0.0) + def to_dict(self) -> Dict[str, Any]: """Convert token to dictionary""" - data = { + data: Dict[str, Any] = { "access_token": self.access_token, "refresh_token": self.refresh_token or "", "token_type": self.token_type, @@ -117,7 +150,6 @@ class MagentaAuthToken(BaseAuthToken): ), "credential_type": self.credential_type or "", } - # Include session data if self.device_id: data["device_id"] = self.device_id if self.session_id: @@ -126,18 +158,133 @@ class MagentaAuthToken(BaseAuthToken): data["channel_map_id"] = self.channel_map_id if self.session_id_updated_at: data["session_id_updated_at"] = self.session_id_updated_at + if self.refresh_expires_in: + data["refresh_expires_in"] = self.refresh_expires_in + if self.account_info is not None: + data["account_info"] = self.account_info + data["account_info_fetched_at"] = self.account_info_fetched_at return data def get_jwt_claims(self) -> Optional[Dict[str, Any]]: - """Extract JWT claims from access token""" + """Extract JWT claims from access token (no expiry verification).""" try: if not self.access_token: return None return decode_jwt(self.access_token, verify=False) - except Exception as e: - logger.debug(f"Failed to extract JWT claims: {e}") + except Exception as exc: + logger.debug(f"Failed to extract JWT claims: {exc}") return None + # ------------------------------------------------------------------ + # Composite-JWT claim accessors. + # + # The HR bifrost `accessToken` is a composite JWT. Its payload embeds + # the HAL / CTS / Persona tokens that theplatform-side services + # (licence server, concurrency service) require: + # + # dc_cts_accountId -> CTS (theplatform) account number + # dc_cts_personaToken -> RS512 JWT used as Widevine Basic-auth password + # dc_cts_personaId -> persona uuid (also inside personaToken.sub) + # dc_tvAccountId -> operator-side account number (e.g. HR + # 6000014999), distinct from dc_cts_accountId + # + # account_url / account_identifier come from /user/account, not the JWT. + # ------------------------------------------------------------------ + + @property + def tv_account_id(self) -> Optional[str]: + c = self.get_jwt_claims() or {} + return c.get("dc_tvAccountId") + + @property + def cts_account_id(self) -> Optional[str]: + c = self.get_jwt_claims() or {} + return c.get("dc_cts_accountId") + + @property + def persona_id(self) -> Optional[str]: + c = self.get_jwt_claims() or {} + return c.get("dc_cts_personaId") + + @property + def persona_jwt(self) -> Optional[str]: + """Raw RS512 persona token (Widevine Basic-auth password).""" + c = self.get_jwt_claims() or {} + return c.get("dc_cts_personaToken") + + @property + def account_uri(self) -> Optional[str]: + """ + MPX account URI, e.g. + http://access.auth.theplatform.com/data/Account/2709375564 + + Prefers the /user/account value (`account_url`); falls back to a + reconstruction from dc_cts_accountId. The reconstruction matches + the shape the live web app uses, but /user/account is authoritative. + """ + if self.account_info and self.account_info.get("account_url"): + return self.account_info["account_url"] + acct = self.cts_account_id + if not acct: + return None + return f"http://access.auth.theplatform.com/data/Account/{acct}" + + @property + def account_identifier(self) -> Optional[str]: + """Bare account uuid from /user/account (`account_identifier`).""" + if self.account_info: + return self.account_info.get("account_identifier") + return None + + # ------------------------------------------------------------------ + # VOD entitlement + # ------------------------------------------------------------------ + + @property + def vod_enabled(self) -> Optional[bool]: + """ + Whether this account is allowed to play VOD at all. + + Returns None when account_info has not been fetched yet -- callers + MUST distinguish "unknown" from "disabled", because acting on the + wrong one produces a false negative. The authoritative switch is + `managed_settings["TVSOA-setting-VodEnabled"]`; individual titles + have their own entitlement (HBO, Nova Plus, ...) which is separate. + """ + if not self.account_info: + return None + ms = self.account_info.get("managed_settings") or {} + return str(ms.get("TVSOA-setting-VodEnabled", "")).lower() == "true" + + @property + def vod_enabled_raw(self) -> Optional[str]: + """Raw value of the VOD-enabled managed setting, for diagnostics.""" + if not self.account_info: + return None + ms = self.account_info.get("managed_settings") or {} + return ms.get("TVSOA-setting-VodEnabled") + + @property + def entitlement_bouquets(self) -> list: + """ + Managed-setting keys whose value is "true" and that look like + package entitlements (e.g. "HR-package-basic-ftv"). Used by the + provider when the actions API reports `subscribe` rather than + `watch`, to decide whether a specific premium title is playable + for this subscriber. + """ + if not self.account_info: + return [] + ms = self.account_info.get("managed_settings") or {} + return sorted( + k for k, v in ms.items() + if str(v).lower() == "true" and "package" in k.lower() + ) + + +# --------------------------------------------------------------------------- +# Auth config +# --------------------------------------------------------------------------- class MagentaAuthConfig: """Configuration for Magenta TV authentication""" @@ -147,7 +294,6 @@ class MagentaAuthConfig: self.http_manager = http_manager self.country_config = COUNTRY_CONFIG[country] - # Application configuration self.app_version = APP_VERSION self.device_name = DEVICE_NAME self.user_agent = USER_AGENT @@ -159,10 +305,10 @@ class MagentaAuthConfig: call_type: str = CALL_TYPES["GUEST_USER"], flow: str = AUTH_FLOWS["START_UP"], step: str = AUTH_STEPS["GET_ACCESS_TOKEN"], - device_id: str = None, - session_id: str = None, - tracking_id: str = None, - call_time: str = None, + device_id: Optional[str] = None, + session_id: Optional[str] = None, + tracking_id: Optional[str] = None, + call_time: Optional[str] = None, ) -> Dict[str, str]: """Get authentication headers, including x-txn-id""" return build_auth_headers( @@ -177,25 +323,42 @@ class MagentaAuthConfig: ) def encrypt_password(self, password: str) -> str: - """Encrypt password using RSA public key""" - try: - rsa_key = self.country_config["rsa_key"] - if not rsa_key: - logger.error(f"No RSA public key configured for country: {self.country}") - return password + """ + Encrypt password using RSA public key. + Raises rather than returning the plaintext on failure. Sending a + plaintext credential in a login payload -- even over HTTPS -- is a + worse outcome than failing loudly. + """ + rsa_key = self.country_config["rsa_key"] + if not rsa_key: + raise RuntimeError( + f"No RSA public key configured for country: {self.country} -- " + f"refusing to send plaintext credentials" + ) + try: key = RSA.import_key(rsa_key) cipher = PKCS1_OAEP.new(key) ciphertext = cipher.encrypt(password.encode("utf-8")) return base64.b64encode(ciphertext).decode() - except Exception as e: - logger.error(f"Error encrypting password: {e}") - return password + except Exception as exc: + raise RuntimeError(f"Failed to encrypt password: {exc}") from exc +# --------------------------------------------------------------------------- +# Authenticator +# --------------------------------------------------------------------------- + class MagentaAuthenticator(BaseAuthenticator): """Magenta TV authenticator - directly extends BaseAuthenticator""" + # How long a /user/account response is reused before re-fetching. Long + # enough that the VOD manager and entitlement checks within one session + # hit the network at most once, short enough that a mid-session + # entitlement change (package added/removed) is picked up before the + # next playback attempt. + ACCOUNT_INFO_TTL_SECONDS = 15 * 60 # 15 minutes + def __init__( self, country: str = DEFAULT_COUNTRY, @@ -204,10 +367,9 @@ class MagentaAuthenticator(BaseAuthenticator): config_dir: Optional[str] = None, http_manager=None, proxy_config: Optional[ProxyConfig] = None, - device_id: Optional[str] = None, # New parameter + device_id: Optional[str] = None, session_id: Optional[str] = None, - ): # New parameter - + ): logger.info(f"=== MagentaAuthenticator.__init__ START ===") if country not in SUPPORTED_COUNTRIES: @@ -222,11 +384,8 @@ class MagentaAuthenticator(BaseAuthenticator): self._http_manager = http_manager self._proxy_config = proxy_config - # Setup config self._config = MagentaAuthConfig(self.country, self._http_manager) - # Call parent init (this will load existing session if available) - # Call parent init (this will load existing session if available) super().__init__( provider_name="magentaeu", settings_manager=settings_manager, @@ -238,13 +397,8 @@ class MagentaAuthenticator(BaseAuthenticator): logger.info(f"=== MagentaAuthenticator.__init__ AFTER super().__init__ ===") - # device_id/session_id are no longer decided once here and then - # trusted for the token's entire lifetime -- a persisted pair could - # be months old with no way to know it ever went invalid server-side. - # We just make sure a token object exists to read/write into; actual - # validation and refresh happens lazily in get_guest_session_ids(), - # called uniformly by every guest-flow call site (channel list, EPG, - # etc), independent of whether the user is authenticated. + # device_id/session_id are validated lazily by get_guest_session_ids(); + # we only make sure a token object exists to read/write into. if not self._current_token or not isinstance(self._current_token, MagentaAuthToken): self._current_token = MagentaAuthToken( access_token="", @@ -254,10 +408,10 @@ class MagentaAuthenticator(BaseAuthenticator): issued_at=time.time(), ) - # Explicit device_id/session_id constructor args (if the caller - # already knows good values) still take precedence, but they don't - # get treated as pre-validated -- session_id_updated_at is left at - # 0 so the first get_guest_session_ids() call still checks them. + # Explicit constructor args (if the caller already knows good values) + # take precedence, but they don't get treated as pre-validated -- + # session_id_updated_at stays at 0 so the first guest request still + # checks them. if device_id: self._current_token.device_id = device_id if session_id: @@ -268,6 +422,10 @@ class MagentaAuthenticator(BaseAuthenticator): f"will be validated on first guest request ===" ) + # ------------------------------------------------------------------ + # Properties + # ------------------------------------------------------------------ + @property def auth_endpoint(self) -> str: """Authentication endpoint - required by BaseAuthenticator""" @@ -283,31 +441,29 @@ class MagentaAuthenticator(BaseAuthenticator): return self._current_token.channel_map_id return "" + @property + def http_manager(self): + """Public access to HTTP manager""" + return self._http_manager + def get_auth_headers(self, call_type: str, flow: str, step: str) -> Dict[str, str]: return self._config.get_auth_headers(call_type, flow, step) def get_epg_headers(self) -> Dict[str, str]: return self.get_auth_headers("GUEST_USER", "START_UP", "EPG_CHANNEL") - @property - def http_manager(self): - """Public access to HTTP manager""" - return self._http_manager + # ------------------------------------------------------------------ + # Guest session + # ------------------------------------------------------------------ - def _initialize_guest_session(self) -> tuple[str, str, bool]: + def _initialize_guest_session(self) -> Tuple[str, str, bool]: """ Visit the provider's startup page to obtain a real deviceId/sessionId pair from Set-Cookie, the same way a browser would. - Uses API_ENDPOINTS["STARTUP_PAGE"] rather than a hardcoded "/epg" - path, since not every MagentaEU natco is guaranteed to serve the - startup page at that path -- a country override only needs to change - constants.py, not this method. - Returns (device_id, session_id, obtained) where obtained=False means we had to fall back to random UUIDs. Callers must NOT treat an - obtained=False result as a validated, cacheable session -- that's - exactly how a permanently-invalid session got persisted before. + obtained=False result as a validated, cacheable session. """ try: startup_url = API_ENDPOINTS["STARTUP_PAGE"].format( @@ -322,6 +478,7 @@ class MagentaAuthenticator(BaseAuthenticator): timeout=DEFAULT_REQUEST_TIMEOUT, ) + status = getattr(response, "status_code", None) device_id = "" session_id = "" @@ -332,34 +489,28 @@ class MagentaAuthenticator(BaseAuthenticator): if device_id and session_id: logger.debug( - f"[{self.country}] Guest session established from cookies - " - f"device_id: {device_id}, session_id: {session_id}" + f"[{self.country}] Guest session established from cookies " + f"(status={status}) - device_id: {device_id}, session_id: {session_id}" ) return device_id, session_id, True logger.warning( - f"[{self.country}] Startup page returned no deviceId/sessionId " - f"cookies; using unverified random fallback" + f"[{self.country}] Startup page (status={status}) returned no " + f"deviceId/sessionId cookies; using unverified random fallback" ) return str(uuid.uuid4()), str(uuid.uuid4()), False - except Exception as e: - logger.warning(f"[{self.country}] Guest session initialization failed: {e}") + except Exception as exc: + logger.warning(f"[{self.country}] Guest session initialization failed: {exc}") return str(uuid.uuid4()), str(uuid.uuid4()), False - def get_guest_session_ids(self, force_refresh: bool = False) -> tuple[str, str]: + def get_guest_session_ids(self, force_refresh: bool = False) -> Tuple[str, str]: """ Return (device_id, session_id) for guest-flow requests (channel - list, EPG, etc). This is the single source of truth for every guest - call site -- previously each call site (get_channels(), the EPG - manager) read current_token.device_id/session_id directly, which - were set once at first-ever init and never re-validated, so a pair - that went stale server-side stayed stale forever. + list, EPG, etc). Single source of truth for every guest call site. Values are re-validated whenever older than GUEST_SESSION_TTL_SECONDS - (or never validated at all -- session_id_updated_at == 0), regardless - of whether the user is authenticated; the bearer-token lifecycle - (_refresh_token) is separate from this guest session. + (or never validated at all -- session_id_updated_at == 0). """ token = self._current_token device_id = "" @@ -400,10 +551,6 @@ class MagentaAuthenticator(BaseAuthenticator): self._current_token.session_id_updated_at = time.time() self._save_session() else: - # Don't stamp session_id_updated_at -- an unverified fallback - # must not be cached as if it were a real, validated session. - # The next call will retry instead of trusting a guess for - # GUEST_SESSION_TTL_SECONDS. logger.warning( f"[{self.country}] Guest session unverified; will retry on " f"next call rather than caching this pair" @@ -411,12 +558,15 @@ class MagentaAuthenticator(BaseAuthenticator): return device_id, session_id + # ------------------------------------------------------------------ + # Header / payload builders + # ------------------------------------------------------------------ + def _get_auth_headers(self) -> Dict[str, str]: """Get headers for authentication request - required by BaseAuthenticator""" device_id = "" session_id = "" - # Get session data from current token if available if self._current_token and isinstance(self._current_token, MagentaAuthToken): device_id = self._current_token.device_id or "" session_id = self._current_token.session_id or "" @@ -441,16 +591,13 @@ class MagentaAuthenticator(BaseAuthenticator): if not self.credentials or not isinstance(self.credentials, UserPasswordCredentials): raise Exception("No valid credentials available") - # Enhanced validation if not self.credentials.username or not self.credentials.password: raise Exception("Username and password cannot be empty") - # Get device_id from current token or expect it to be provided via other means device_id = "" if self._current_token and isinstance(self._current_token, MagentaAuthToken): device_id = self._current_token.device_id or "" - # If no device_id, we need to get it from the provider if not device_id: device_id = str(uuid.uuid4()) @@ -477,68 +624,142 @@ class MagentaAuthenticator(BaseAuthenticator): "broadcastingStreamLimitationApplies": BROADCASTING_STREAM_LIMITATION_APPLIES, }, "telekomLogin": { - "username": self.credentials.username, # Works for both types! - "password": encrypted_password, # Works for both types! + "username": self.credentials.username, + "password": encrypted_password, }, } + # ------------------------------------------------------------------ + # Token construction + # ------------------------------------------------------------------ + + @staticmethod + def _read_token_field( + data: Dict[str, Any], + camel: str, + snake: str, + default: Any = None, + ) -> Any: + """ + Read a field that may appear in either camelCase (login/refresh + bodies) or snake_case (stored session format). is-not-None checks + rather than truthiness so a legitimate 0 is not treated as missing. + """ + if camel in data and data[camel] is not None: + return data[camel] + if snake in data and data[snake] is not None: + return data[snake] + return default + def _create_token_from_response(self, response_data: Dict[str, Any]) -> BaseAuthToken: """Create token from API response - required by BaseAuthenticator""" - # PRESERVE the existing session IDs (which follow the correct priority) + + # --- Preserve existing session data --------------------------------- device_id = "" session_id = "" channel_map_id = "" session_id_updated_at = 0.0 + existing_account_info: Optional[Dict[str, Any]] = None + existing_account_info_fetched_at = 0.0 - # Try to get session IDs from multiple sources in priority order: - - # 1. First from the response_data itself (when loading from stored session) + # Priority 1: fields stored directly on the token blob (restored + # sessions come through this path -- _load_session() calls + # _create_token_from_response with the persisted dict). if "device_id" in response_data: - device_id = response_data.get("device_id", "") + device_id = response_data.get("device_id", "") or "" if "session_id" in response_data: - session_id = response_data.get("session_id", "") + session_id = response_data.get("session_id", "") or "" if "session_id_updated_at" in response_data: session_id_updated_at = response_data.get("session_id_updated_at", 0.0) or 0.0 - # 2. Then from current token (for new authentications) - if ( - (not device_id or not session_id) - and self._current_token - and isinstance(self._current_token, MagentaAuthToken) - ): - device_id = self._current_token.device_id or "" - session_id = self._current_token.session_id or "" - channel_map_id = self._current_token.channel_map_id or "" - session_id_updated_at = self._current_token.session_id_updated_at or 0.0 + # account_info restore -- TTL-checked. A stale snapshot must not + # override a real entitlement change that happened while the app + # was closed. + if "account_info" in response_data and response_data["account_info"] is not None: + fetched_at = response_data.get("account_info_fetched_at", 0.0) or 0.0 + age = time.time() - fetched_at + if age < self.ACCOUNT_INFO_TTL_SECONDS: + existing_account_info = response_data["account_info"] + existing_account_info_fetched_at = fetched_at + logger.debug( + f"Restored persisted account_info (age={age:.0f}s)" + ) + else: + logger.debug( + f"Dropping persisted account_info: age {age:.0f}s " + f"exceeds TTL {self.ACCOUNT_INFO_TTL_SECONDS}s" + ) + + # Priority 2: carry forward from the current token (refresh / upgrade + # paths -- same user, so the entitlements remain valid). + if (not device_id or not session_id) and isinstance( + self._current_token, MagentaAuthToken + ): + device_id = device_id or (self._current_token.device_id or "") + session_id = session_id or (self._current_token.session_id or "") + channel_map_id = self._current_token.channel_map_id or "" + session_id_updated_at = ( + session_id_updated_at + or self._current_token.session_id_updated_at + or 0.0 + ) + if existing_account_info is None: + existing_account_info = self._current_token.account_info + existing_account_info_fetched_at = ( + self._current_token.account_info_fetched_at or 0.0 + ) - # 3. If we found session IDs, log it if device_id and session_id: logger.debug( - f"Creating new token with session IDs - device_id: {device_id}, session_id: {session_id}" + f"Creating new token with session IDs - device_id: {device_id}, " + f"session_id: {session_id}" ) else: logger.warning( - f"No session IDs found in response_data or current_token during token creation" + "No session IDs found in response_data or current_token during " + "token creation" ) - # DUAL KEY SUPPORT: Handle both camelCase (API responses) and snake_case (stored sessions) - # Access token - access_token = response_data.get("accessToken") or response_data.get("access_token") + # --- Access token (required) ---------------------------------------- + access_token = self._read_token_field( + response_data, "accessToken", "access_token" + ) if not access_token: - logger.error(f"CRITICAL: No access token found in response data") + logger.error("CRITICAL: No access token found in response data") logger.error(f"Available keys: {list(response_data.keys())}") raise Exception("No access token found in response data") - # Refresh token - refresh_token = response_data.get("refreshToken") or response_data.get("refresh_token", "") + # --- Refresh token -------------------------------------------------- + refresh_token = self._read_token_field( + response_data, "refreshToken", "refresh_token", "" + ) - # Expires in - expires_in = response_data.get("expiresIn") or response_data.get("expires_in", 3600) + # --- Expiries ------------------------------------------------------- + # The HR login response field is `accessExpiresIn`, NOT `expiresIn`. + # Reading the wrong key previously caused a fallback to 3600s and a + # refresh attempt every hour against 7-day access tokens. + expires_in_raw = self._read_token_field( + response_data, "accessExpiresIn", "expires_in", None + ) + if expires_in_raw is None: + # Older / other natco variants may use `expiresIn`. + expires_in_raw = self._read_token_field( + response_data, "expiresIn", "expires_in", 3600 + ) + expires_in = int(expires_in_raw) - # Token type - token_type = response_data.get("tokenType") or response_data.get("token_type", "Bearer") + # refresh_expires_in must never be None -- BaseAuthToken.needs_refresh + # does `if self.refresh_expires_in > 0:` and would raise TypeError. + refresh_expires_in_raw = self._read_token_field( + response_data, "refreshExpiresIn", "refresh_expires_in", None + ) + refresh_expires_in = ( + int(refresh_expires_in_raw) if refresh_expires_in_raw is not None else 0 + ) - # For stored sessions, issued_at might be in the data, otherwise use current time + token_type = self._read_token_field( + response_data, "tokenType", "token_type", "Bearer" + ) issued_at = response_data.get("issued_at", time.time()) token = MagentaAuthToken( @@ -551,12 +772,17 @@ class MagentaAuthenticator(BaseAuthenticator): session_id=session_id, channel_map_id=channel_map_id, session_id_updated_at=session_id_updated_at, + refresh_expires_in=refresh_expires_in, + account_info=existing_account_info, + account_info_fetched_at=existing_account_info_fetched_at, ) - # Classify token token.auth_level = self._classify_token(token) - logger.debug(f"Token created successfully from {len(response_data)} data fields") + logger.info( + f"Token created: access_expires_in={expires_in}s, " + f"refresh_expires_in={refresh_expires_in}s" + ) return token def get_fallback_credentials(self): @@ -565,32 +791,46 @@ class MagentaAuthenticator(BaseAuthenticator): return UserPasswordCredentials(username="", password="") + # ------------------------------------------------------------------ + # Authentication + # ------------------------------------------------------------------ + def _perform_authentication(self) -> BaseAuthToken: """Perform Magenta TV authentication - required by BaseAuthenticator""" - # Enhanced credential validation - FIXED VERSION if not self.credentials: raise Exception("No credentials available for authentication") - # Accept both MagentaCredentials AND base UserPasswordCredentials from ...base.auth.credentials import UserPasswordCredentials if not isinstance(self.credentials, UserPasswordCredentials): raise Exception( - f"Invalid credential type: {type(self.credentials)}. Expected UserPasswordCredentials or MagentaCredentials" + f"Invalid credential type: {type(self.credentials)}. " + f"Expected UserPasswordCredentials or MagentaCredentials" ) - # Validate credential content if not self.credentials.username or not self.credentials.password: raise Exception("Username and password are required for authentication") logger.info(f"Performing Magenta TV authentication for country: {self.country}") + # Clear any previously cached /user/account data before logging in. + # A full login may be for a *different* user than the token we + # currently hold -- carrying the old account_info forward would + # serve the previous user's vod_enabled/entitlement_bouquets for up + # to ACCOUNT_INFO_TTL_SECONDS. Refresh/upgrade paths (which reuse + # _create_token_from_response without going through this method) + # deliberately do NOT clear it, since those are same-user. + if isinstance(self._current_token, MagentaAuthToken): + self._current_token.account_info = None + self._current_token.account_info_fetched_at = 0.0 + try: - # Perform login headers = self._get_auth_headers() payload = self._build_auth_payload() - logger.debug(f"Authentication payload prepared for user: {self.credentials.username}") + logger.debug( + f"Authentication payload prepared for user: {self.credentials.username}" + ) response = self._http_manager.post( self.auth_endpoint, @@ -603,15 +843,16 @@ class MagentaAuthenticator(BaseAuthenticator): response.raise_for_status() token_data = response.json() - # Handle device limit exceeded if token_data.get("deviceLimitExceed", False): logger.info("Device limit exceeded, attempting token upgrade") token_data = self._upgrade_token(token_data["refreshToken"]) return self._create_token_from_response(token_data) - except Exception as e: - logger.error(f"Authentication failed for user {self.credentials.username}: {e}") + except Exception as exc: + logger.error( + f"Authentication failed for user {self.credentials.username}: {exc}" + ) raise def _upgrade_token(self, refresh_token: str) -> Dict[str, Any]: @@ -647,7 +888,7 @@ class MagentaAuthenticator(BaseAuthenticator): response.raise_for_status() return response.json() - def _get_session_data(self) -> tuple[str, str, str, float]: + def _get_session_data(self) -> Tuple[str, str, str, float]: """Safely get session data from current token""" if isinstance(self._current_token, MagentaAuthToken): return ( @@ -659,7 +900,13 @@ class MagentaAuthenticator(BaseAuthenticator): return "", "", "", 0.0 def _refresh_token(self) -> Optional[BaseAuthToken]: - """Refresh Magenta TV token - override base method""" + """ + Refresh Magenta TV token - override base method. + + Routes the response through _create_token_from_response() so the + camelCase/snake_case handling, the accessExpiresIn/refreshExpiresIn + reading, and the account-info carry-forward are done in one place. + """ if not self._current_token or not self._current_token.refresh_token: logger.debug("No valid refresh token available") return None @@ -669,12 +916,11 @@ class MagentaAuthenticator(BaseAuthenticator): refresh_url = API_ENDPOINTS["REFRESH_TOKEN"].format(natco=self.country) - device_id, session_id, channel_map_id, session_id_updated_at = self._get_session_data() + device_id, session_id, _, _ = self._get_session_data() tracking_id = str(uuid.uuid4()) call_time = str(int(time.time() * 1000)) - # Build headers according to your working example headers = self._config.get_auth_headers( call_type=CALL_TYPES["AUTH_USER"], flow=AUTH_FLOWS["START_UP"], @@ -685,7 +931,6 @@ class MagentaAuthenticator(BaseAuthenticator): call_time=call_time, ) - # Add the specific headers from your working example headers.update( { "Refresh_token": self._current_token.refresh_token, @@ -693,18 +938,13 @@ class MagentaAuthenticator(BaseAuthenticator): } ) - # Build payload matching your working example payload = { - "clientVersion": APP_VERSION, # Use current APP_VERSION + "clientVersion": APP_VERSION, "deviceId": device_id, "concurrencyLimitParam": DEVICE_CONCURRENCY_PARAM, } logger.debug(f"Refresh request - URL: {refresh_url}") - logger.debug( - f"Refresh request - Headers: { {k: v for k, v in headers.items() if k not in ['Authorization', 'Refresh_token']} }" - ) - logger.debug(f"Refresh request - Payload: {payload}") response = self._http_manager.post( refresh_url, @@ -717,69 +957,81 @@ class MagentaAuthenticator(BaseAuthenticator): response.raise_for_status() token_data = response.json() - # Create new token with updated data but preserve session IDs - new_token = MagentaAuthToken( - access_token=token_data["accessToken"], - refresh_token=token_data.get("refreshToken", self._current_token.refresh_token), - token_type="Bearer", - expires_in=token_data.get("expiresIn", 3600), - issued_at=time.time(), - device_id=device_id, - session_id=session_id, - channel_map_id=channel_map_id, - session_id_updated_at=session_id_updated_at, - ) + return self._create_token_from_response(token_data) - # Classify token - new_token.auth_level = self._classify_token(new_token) - logger.info("Token refresh successful") - return new_token - - except Exception as e: - logger.warning(f"Token refresh failed: {e}") - if hasattr(e, "response") and hasattr(e.response, "text"): - logger.error(f"Refresh response content: {e.response.text}") + except Exception as exc: + logger.warning(f"Token refresh failed: {exc}") + if hasattr(exc, "response") and hasattr(exc.response, "text"): + logger.error(f"Refresh response content: {exc.response.text}") return None + # ------------------------------------------------------------------ + # Token classification + # ------------------------------------------------------------------ + def _classify_token(self, token: BaseAuthToken) -> TokenAuthLevel: """Classify Magenta TV token - required by BaseAuthenticator""" try: if not token or not token.access_token: return TokenAuthLevel.UNKNOWN - # Only MagentaAuthToken has get_jwt_claims method if isinstance(token, MagentaAuthToken): claims = token.get_jwt_claims() else: - # For BaseAuthToken, try to decode JWT manually try: claims = decode_jwt(token.access_token, verify=False) except InvalidTokenError: return TokenAuthLevel.UNKNOWN - except (ValueError, json.JSONDecodeError, UnicodeDecodeError) as e: - logger.debug(f"Error decoding JWT token: {e}") - return TokenAuthLevel.UNKNOWN if not claims: return TokenAuthLevel.UNKNOWN - # Magenta tokens with username in claims indicate user authentication if "username" in claims or "preferred_username" in claims: return TokenAuthLevel.USER_AUTHENTICATED - # Anonymous tokens typically have limited claims - if len(claims) <= 3: # Basic claims like exp, iat, iss + if len(claims) <= 3: return TokenAuthLevel.ANONYMOUS return TokenAuthLevel.USER_AUTHENTICATED - except Exception as e: - logger.debug(f"Error classifying token: {e}") + except Exception as exc: + logger.debug(f"Error classifying token: {exc}") return TokenAuthLevel.UNKNOWN - def get_user_account(self) -> Dict[str, Any]: - """Get user account information""" + # ------------------------------------------------------------------ + # /user/account + # ------------------------------------------------------------------ + + def get_user_account(self, force_refresh: bool = False) -> Dict[str, Any]: + """ + Get user account information, with in-memory caching on the token. + + The response is cached for ACCOUNT_INFO_TTL_SECONDS and survives + restarts (persisted via _save_session / restored in + _create_token_from_response, with the TTL re-checked on restore). + Callers that need to pick up a mid-session entitlement change can + pass force_refresh=True. + + Side effects: + * populates current_token.channel_map_id + * populates current_token.account_info / account_info_fetched_at + * persists the token + """ + token = self._current_token + if not isinstance(token, MagentaAuthToken): + raise Exception("No token available to hold account info") + + if not force_refresh and token.account_info is not None: + age = time.time() - (token.account_info_fetched_at or 0.0) + if age < self.ACCOUNT_INFO_TTL_SECONDS: + logger.debug( + f"get_user_account: serving cached account info (age={age:.0f}s)" + ) + return token.account_info + access_token = self.get_bearer_token() + if access_token.startswith("Bearer "): + access_token = access_token[7:] account_url = API_ENDPOINTS["USER_ACCOUNT"].format( bifrost_url=get_bifrost_url(self.country) @@ -791,11 +1043,7 @@ class MagentaAuthenticator(BaseAuthenticator): "natco_code": self.country, } - device_id = "" - session_id = "" - if self._current_token and isinstance(self._current_token, MagentaAuthToken): - device_id = self._current_token.device_id or "" - session_id = self._current_token.session_id or "" + device_id, session_id, _, _ = self._get_session_data() tracking_id = str(uuid.uuid4()) call_time = str(int(time.time() * 1000)) @@ -822,14 +1070,21 @@ class MagentaAuthenticator(BaseAuthenticator): response.raise_for_status() account_data = response.json() - # Save channel map ID to current token - if ( - "channelMap_id" in account_data - and self._current_token - and isinstance(self._current_token, MagentaAuthToken) - ): - self._current_token.channel_map_id = account_data["channelMap_id"] - # Save the updated token with channel_map_id - self._save_session() + if "channelMap_id" in account_data: + token.channel_map_id = account_data["channelMap_id"] + token.account_info = account_data + token.account_info_fetched_at = time.time() + + ms = account_data.get("managed_settings") or {} + logger.info( + f"[{self.country}] account loaded: " + f"tvAccountId={account_data.get('tvAccountId')}, " + f"vod_enabled={ms.get('TVSOA-setting-VodEnabled')!r}, " + f"catchup_enabled={account_data.get('catchup_enabled')}, " + f"entitlement_bouquets=" + f"{sorted(k for k, v in ms.items() if 'package' in k.lower() and str(v).lower() == 'true')}" + ) + + self._save_session() return account_data \ No newline at end of file diff --git a/lib/streaming_providers/providers/magentaeu/provider.py b/lib/streaming_providers/providers/magentaeu/provider.py index c99475c..21a83d0 100644 --- a/lib/streaming_providers/providers/magentaeu/provider.py +++ b/lib/streaming_providers/providers/magentaeu/provider.py @@ -17,8 +17,11 @@ from ..lib_theplatform import ( build_widevine_drm_config, parse_bifrost_epg_channel, ) +import threading from .auth import MagentaAuthenticator from .epg_manager import MagentaEUEpgManager +from .vod_manager import MagentaEUVodManager +from .vod_errors import VodCatchupRequiredError, VodNotFoundError from .constants import ( API_ENDPOINTS, CONTENT_TYPE_LIVE, @@ -102,6 +105,11 @@ class MagentaEUProvider(StreamingProvider): http_manager=self.http_manager, authenticator=self.authenticator, ) + + # VOD manager — lazy, same reasoning as epg_manager but VOD is + # opt-in per account (see MagentaEUVodManager.is_vod_enabled). + self._vod_manager: Optional[MagentaEUVodManager] = None + self._vod_manager_lock = threading.Lock() # NEW logger.info(f"=== MagentaProvider.__init__ COMPLETE ===") def _load_proxy_from_manager(self, config_dir: Optional[str]) -> Optional[ProxyConfig]: @@ -145,6 +153,39 @@ class MagentaEUProvider(StreamingProvider): return 336 return 168 + @property + def vod_manager(self) -> MagentaEUVodManager: + if self._vod_manager is None: + with self._vod_manager_lock: + if self._vod_manager is None: # re-check inside the lock + self._vod_manager = MagentaEUVodManager( + country=self.country, + http_manager=self.http_manager, + authenticator=self.authenticator, + ) + return self._vod_manager + + @property + def implements_vod(self) -> bool: + return True + + def get_vod_category(self, content_id: str = "", **kwargs) -> List: + return self.vod_manager.get_category_children(content_id) + + def search_vod( + self, + query: str, + cursor: Optional[str] = None, + page_size: int = 24, + **kwargs, + ) -> List: + if cursor is not None: + logger.warning( + f"[{self.country}] search_vod: cursor={cursor!r} ignored — " + f"search pagination is not yet implemented" + ) + return self.vod_manager.search(query, size=page_size).entries + @property def supported_auth_types(self) -> List[str]: return ["user_credentials"] @@ -403,16 +444,63 @@ class MagentaEUProvider(StreamingProvider): return True def get_manifest(self, content_id: str, **kwargs) -> Optional[str]: - """Get manifest URL for a channel by ID""" - if not self._ensure_channels_cache(): + """ + Get manifest URL for a channel OR a VOD asset by ID. + + Tries the live-channel cache first (unchanged fast path), then + falls back to the VOD manager. + + Only VodNotFoundError and VodCatchupRequiredError are caught + and swallowed to None here — both are legitimate "there's + nothing to play at this content_id via this path" outcomes, + same as the pre-existing "channel not found" case below them. + Every other VodError (VodEntitlementError, + VodAccountVodDisabledError, VodGeoBlockError, VodAuthError + after its internal retry already failed, VodRateLimitError, + VodServerError, a bare VodError) is a real, user-relevant + failure and is NOT caught here — it propagates to the caller, + which is expected to already have error handling for + get_manifest() failing (this mirrors how the original + get_manifest() let arbitrary exceptions from + _ensure_channels_cache's callees propagate). Swallowing those + into a silent None would hide exactly the failures a person + needs to see (e.g. "you need HBO for this" vs. "nothing here"). + + A VodCatchupRequiredError means the "VOD" entry is actually a + linear catch-up item — see vod_errors.VodCatchupRequiredError + docstring, and vod_manager.py's note that the check producing + this exception is currently unconfirmed against real capture + data (see _resolve_playable_video_id). Full bridging into + get_catchup_manifest() is not implemented here: not because of + any epoch/ISO conversion difficulty (build_catchup_url() already + takes plain epoch ints, and the captured catchup_start_utc / + catchup_end_utc strings are ordinary ISO-8601 — trivial to parse + with stdlib datetime, no TimestampConverter reverse-direction + needed) — but because the exception that would carry the + catchup_schedules[] data isn't confirmed reachable from the + endpoint this code path actually calls. Wire the bridge once + that's confirmed; don't invent a conversion before there's data + to convert. + """ + if self._ensure_channels_cache(): + for channel in self._channels_cache: + if channel.channel_id == content_id: + return channel.manifest + + try: + return self.vod_manager.get_manifest(content_id, **kwargs) + except VodCatchupRequiredError as exc: + logger.warning( + f"[{self.country}] {content_id} is catch-up-only, not VOD. " + f"station_id={exc.station_id!r}, " + f"{len(exc.catchup_schedules)} catchup window(s) available. " + f"Bridging to get_catchup_manifest() is not yet implemented." + ) return None - - for channel in self._channels_cache: - if channel.channel_id == content_id: - return channel.manifest - - logger.warning(f"Channel {content_id} not found in available channels") - return None + except VodNotFoundError: + logger.warning(f"Content {content_id} not found in channels or VOD") + return None + # Deliberately no broader `except VodError` here — see docstring. def get_catchup_manifest( self, content_id: str, start_time: int, end_time: int, drm_variant: Optional[str] = "auto", **kwargs @@ -456,25 +544,33 @@ class MagentaEUProvider(StreamingProvider): return base_manifest def get_drm(self, content_id: str, **kwargs) -> List[DRMConfig]: - """Get DRM configurations for channel by ID""" - logger.info(f"=== get_drm_configs_by_id CALLED for channel_id: {content_id} ===") + """ + Get DRM configurations for a channel OR a VOD asset by ID. - if not self._ensure_channels_cache(): - logger.warning(f"Cannot get DRM for {content_id}, channels cache unavailable") + Same exception-narrowing rule as get_manifest() above: only + VodNotFoundError and VodCatchupRequiredError are swallowed to + an empty list. Every other VodError propagates. + """ + logger.info(f"=== get_drm CALLED for content_id: {content_id} ===") + + if self._ensure_channels_cache(): + for cached_channel in self._channels_cache: + if cached_channel.channel_id == content_id: + drm_config = self.get_drm_config(cached_channel, **kwargs) + return [drm_config] if drm_config else [] + + try: + return self.vod_manager.get_drm(content_id, **kwargs) + except VodCatchupRequiredError as exc: + logger.warning( + f"[{self.country}] {content_id} is catch-up-only, not VOD — " + f"no DRM config to resolve here. station_id={exc.station_id!r}" + ) return [] - - channel = None - for cached_channel in self._channels_cache: - if cached_channel.channel_id == content_id: - channel = cached_channel - break - - if not channel: - logger.warning(f"Channel with ID {content_id} not found in cache") + except VodNotFoundError: + logger.warning(f"Content {content_id} not found in channels or VOD") return [] - - drm_config = self.get_drm_config(channel, **kwargs) - return [drm_config] if drm_config else [] + # Deliberately no broader `except VodError` here — see docstring. def get_drm_configs(self, channel: StreamingChannel, **kwargs) -> List[DRMConfig]: """Get DRM configurations for channel""" diff --git a/lib/streaming_providers/providers/magentaeu/vod_errors.py b/lib/streaming_providers/providers/magentaeu/vod_errors.py new file mode 100644 index 0000000..0796f3b --- /dev/null +++ b/lib/streaming_providers/providers/magentaeu/vod_errors.py @@ -0,0 +1,152 @@ +# streaming_providers/providers/magentaeu/vod_errors.py +# -*- coding: utf-8 -*- +""" +Typed exceptions for the MagentaEU VOD path. + +Rationale: the bifrost API returns 4xx for several distinct conditions -- +expired token, geo-block, entitlement denial, content removed -- and they +need to be handled differently. `raise_for_status()` alone collapses them +into one opaque HTTPError. VODOperations and the Kodi player layer need +to distinguish "refresh and retry" from "tell the user they can't watch +this here" from "this title is gone" from "this is actually catch-up, +not VOD -- hand it to the channel/catchup pathway instead". + +Hierarchy: + + VodError (base, catch-all) + ├── VodAuthError (401 -- token expired/invalid) + ├── VodGeoBlockError (403 -- not available in your region) + ├── VodEntitlementError (403 -- not in your subscription) + │ └── VodAccountVodDisabledError (account-level VOD gate is off) + ├── VodNotFoundError (404 -- content removed) + ├── VodRateLimitError (429) + ├── VodServerError (5xx -- retryable) + ├── VodCatchupRequiredError (item has no watch/trailer action but + │ DOES have schedules/catchup_schedules + │ -- it's a linear-catchup item, not a + │ playable VOD asset; see docstring) + └── VodNotImplementedError (feature captured-but-not-yet-mapped) +""" + +from __future__ import annotations + +from typing import Any, Dict, List, Optional + + +class VodError(Exception): + """Base class for all MagentaEU VOD errors.""" + + def __init__(self, message: str, *, status: Optional[int] = None, + url: Optional[str] = None) -> None: + super().__init__(message) + self.status = status + self.url = url + + +class VodAuthError(VodError): + """ + Raised on 401 or on a locally-detected expired Bff_token. + + `MagentaEUVodManager._request()` catches this internally, refreshes + the token once via a forced re-authentication, and retries the + request exactly once. If the retry also fails, VodAuthError + propagates to the caller -- at that point re-auth is not going to + fix it (credentials are likely actually invalid). + """ + + +class VodGeoBlockError(VodError): + """Content is geo-blocked for the current network egress.""" + + +class VodEntitlementError(VodError): + """ + The subscriber is authenticated but not entitled to this content. + + Raised by the actions endpoint when `actions.watch` is empty and + `svod_subscription_message` is present, AND the item is not a + catchup-eligible linear item (see VodCatchupRequiredError -- that + takes priority when both watch/trailer are empty). + """ + + +class VodAccountVodDisabledError(VodEntitlementError): + """ + Account-level VOD gate is off. + + This is distinct from a per-title entitlement error: the account's + `managed_settings["TVSOA-setting-VodEnabled"]` is "false", which + means *no* VOD can be played, regardless of package. The captured + HR test account ships in this state (`vod_enabled: false` at + /user/account and echoed in every X-Account-Details header). + Callers should surface this differently than "you need to add the + HBO package". + """ + + +class VodNotFoundError(VodError): + """Content has been removed or never existed.""" + + +class VodRateLimitError(VodError): + """429 -- caller should back off.""" + + +class VodServerError(VodError): + """5xx -- retryable.""" + + +class VodCatchupRequiredError(VodError): + """ + This "VOD" list entry is actually a catch-up item from a linear + channel, not a TVOD/SVOD asset. + + Evidence from capture: some rail/search items marked as VOD series + (e.g. "Nakon poplave, serija") return `actions.watch: []` and + `actions.trailer: []` -- same shape as a genuine entitlement denial + -- but ALSO carry a non-empty `actions.schedules` / a `station_id` + /`channel_number`, and the episode list includes populated + `catchup_schedules[]` entries with their own `pid` and + `catchup_start_utc`/`catchup_end_utc` window. Treating this as a + plain VodEntitlementError is wrong: the content IS watchable, just + via the existing live-channel catchup path + (`provider.get_catchup_manifest()`), not via + `MagentaEUVodManager.get_manifest()`. + + This exception carries the raw schedule/catchup data so a caller + that wants to bridge into the catchup flow can do so without a + second round-trip. Bridging itself (mapping station_id + + catchup_start_utc/catchup_end_utc into a get_catchup_manifest() + call) is NOT implemented here -- it needs the channel/EPG manager, + which this module deliberately has no dependency on. Treat this as + a routing signal, not a playback failure. + """ + + def __init__( + self, + message: str, + *, + station_id: Optional[str] = None, + schedules: Optional[List[Dict[str, Any]]] = None, + catchup_schedules: Optional[List[Dict[str, Any]]] = None, + status: Optional[int] = None, + url: Optional[str] = None, + ) -> None: + super().__init__(message, status=status, url=url) + self.station_id = station_id + self.schedules = schedules or [] + self.catchup_schedules = catchup_schedules or [] + + +class VodNotImplementedError(VodError): + """ + Raised by methods whose endpoint shape has not been captured yet. + + Current cases: + * pricing resolution (needs a non-empty rent[]/purchase[] sample) + * recommendations rails (needs the recommendations endpoint) + * related-content parsing beyond the raw pass-through in + get_related_content() (field shapes only partially captured) + These raise rather than returning empty so callers can't silently + ship with a bug they don't know about. + """ \ No newline at end of file diff --git a/lib/streaming_providers/providers/magentaeu/vod_manager.py b/lib/streaming_providers/providers/magentaeu/vod_manager.py new file mode 100644 index 0000000..07c39b5 --- /dev/null +++ b/lib/streaming_providers/providers/magentaeu/vod_manager.py @@ -0,0 +1,1036 @@ +# streaming_providers/providers/magentaeu/vod_manager.py +# -*- coding: utf-8 -*- +""" +VOD manager for MagentaEU (HR/AT/PL/HU/ME). + +Changes vs. the original proposal (see inline comments tagged FIXED / +NEW for the specifics): + + FIXED DRM licence pid: get_drm() previously built the Widevine + licence URL from playinfo's title_id/program_id. Against the + capture, the real `releasePid` the client sends is + `video.pid` from the /media response -- a different value + that only exists after resolving the manifest. get_manifest() + and get_drm() now share one `_resolve_playback()` call so + both use the correct pid and neither re-fetches playinfo/media + redundantly. + FIXED 401 retry: `retry_on_auth` was accepted but never acted on. + `_request()` now actually catches VodAuthError, force-refreshes + the token once, and retries. + FIXED X-Tv-Step values: root categories capture shows step=DECK + (not CATEGORIES); series actions capture shows + step=SERIES_WATCH_ACTION (not SERIES_ACTIONS). + NEW VodCatchupRequiredError: some catalogue entries that look like + VOD (empty actions.watch/trailer) are actually catch-up-only + items from a linear channel (populated actions.schedules / + episode catchup_schedules[], a station_id). These now raise a + distinct, data-carrying exception instead of being reported as + a plain entitlement denial. + NEW get_program_metadata(): the real client fetches + `details/program/{id}` (no /actions/v2 suffix, + flow=SINGLE_PROGRAM_DETAIL, step=PROGRAM_METADATA) for + description/cast/genre fields the actions/v2 endpoint doesn't + carry. The response schema for this endpoint has not been + captured yet (only headers were observed, not a body), so this + returns the raw parsed JSON rather than a typed object -- do + not assume field names without a capture. + NEW get_related_content(): `relatedcontent/feed` is called by the + real client but was not implemented in the original proposal. + Field shapes are captured for the "Movie"/"Program" asset + entries seen so far; still un-captured shapes raise + VodNotImplementedError via the same discipline as search(). + +Everything else (browse, search, pagination) is unchanged from the +original proposal, which matched the capture well. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass +from typing import Any, Dict, List, Optional, Tuple, Union + +from ...base.models import ContentType, DRMConfig, StreamingMode +from ...base.models.vod import VodCategory, VodItem +from ...base.utils.logger import logger + +from ..lib_theplatform import ( + build_licence_url, + build_widevine_drm_config, +) +from .auth import MagentaAuthToken, MagentaAuthenticator +from .constants import ( + DEFAULT_REQUEST_TIMEOUT, + STREAMING_FORMAT_DASH, + USER_AGENT, + WV_URL, + build_auth_headers, + get_app_key, + get_base_url, + get_bifrost_url, + get_language, + get_natco_key, +) +from .vod_errors import ( + VodAccountVodDisabledError, + VodAuthError, + VodCatchupRequiredError, + VodEntitlementError, + VodError, + VodGeoBlockError, + VodNotFoundError, + VodRateLimitError, + VodServerError, +) + + +# --------------------------------------------------------------------------- +# Response wrappers +# --------------------------------------------------------------------------- + +@dataclass +class VodPage: + """A page of VOD results (see get_component_assets_page).""" + entries: List[Union[VodCategory, VodItem]] + next_offset: Optional[int] + + @property + def has_more(self) -> bool: + return self.next_offset is not None + + +@dataclass +class SearchResult: + """A search response (see search()).""" + entries: List[Union[VodCategory, VodItem]] + has_more_items: bool + + +@dataclass +class ResolvedPlayback: + """ + Everything needed to play one VOD asset, resolved once. + + `media_id` is the trailer/watch service item id (input to /media). + `release_pid` is /media's `video.pid` -- the value the real client + sends as `releasePid` to the Widevine licence endpoint. This is + NOT the same as playinfo's title_id/program_id; conflating those + was the original proposal's DRM bug. + """ + video_src: str + release_pid: str + content_type: str + media_id: str + + +# --------------------------------------------------------------------------- +# Manager +# --------------------------------------------------------------------------- + +class MagentaEUVodManager: + """VOD browse / search / playback-info manager for MagentaEU.""" + + DEFAULT_ASSET_PAGE_SIZE = 20 + DEFAULT_SEARCH_SIZE = 30 + + # How long a resolved-playback result stays valid for reuse between + # a get_manifest() call and the get_drm() call that typically + # follows it for the same title within one playback attempt. This + # is NOT a manifest-freshness guarantee -- it only avoids the + # provider re-hitting playinfo/media twice for one play action. + _PLAYBACK_CACHE_TTL_SECONDS = 60 + + def __init__( + self, + country: str, + http_manager, + authenticator: MagentaAuthenticator, + ) -> None: + self._country = country + self._http = http_manager + self._auth = authenticator + + self._bifrost_url = get_bifrost_url(country) + self._natco_key = get_natco_key(country) + self._app_language = get_language(country) + self._app_key = get_app_key(country) + self._origin = get_base_url(country) + + # (program_id, video_id) -> (ResolvedPlayback, resolved_at_epoch) + self._playback_cache: Dict[Tuple[str, str], Tuple[ResolvedPlayback, float]] = {} + + logger.info(f"[MagentaEUVodManager/{country}] initialised") + + # ================================================================== + # Public API -- VOD enablement + # ================================================================== + + def is_vod_enabled(self, force_refresh: bool = False) -> Optional[bool]: + try: + self._auth.get_user_account(force_refresh=force_refresh) + except Exception as exc: + logger.warning(f"[{self._country}] get_user_account failed: {exc}") + return None + return self._auth.current_token.vod_enabled + + def require_vod_enabled(self, force_refresh: bool = False) -> None: + enabled = self.is_vod_enabled(force_refresh=force_refresh) + if enabled is False: + raw = self._auth.current_token.vod_enabled_raw + raise VodAccountVodDisabledError( + f"VOD is disabled for this account " + f"(TVSOA-setting-VodEnabled={raw!r})" + ) + + # ================================================================== + # Public API -- browse + # ================================================================== + + def get_root_categories(self) -> List[VodCategory]: + """Top-level VOD navigation (POČETNA, SERIJE, FILMOVI, ZA DJECU, SPORT).""" + data = self._request( + "home/root/categories", + params={ + "device_type": "WEB", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="HOME", + step="DECK", # FIXED: capture shows DECK, not CATEGORIES + ) + + categories: List[VodCategory] = [] + for raw in data.get("categories") or []: + cat = self._parse_category_root(raw) + if cat is not None: + categories.append(cat) + + logger.debug(f"[{self._country}] root categories: {len(categories)} entries") + return categories + + def get_category_children(self, content_id: str) -> List[Union[VodCategory, VodItem]]: + if not content_id: + return self.get_root_categories() + + try: + return self._get_page_rails(content_id) + except VodNotFoundError: + pass + + return self._get_component_assets(content_id) + + def _get_page_rails(self, page_id: str) -> List[VodCategory]: + data = self._request( + f"home/page/{page_id}", + params={ + "page_size": 10, + "offset": 0, + "component_type": "all", + "is_opted_in": "true", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="HOME", + step="PAGE_COMPONENTS", # confirmed by capture (doc 7) + ) + + rails: List[VodCategory] = [] + for raw in data.get("components") or []: + template = raw.get("template_id") + if template not in ("RAIL", "HIGHLIGHT"): + continue + cat = self._parse_rail(raw) + if cat is not None: + rails.append(cat) + + logger.debug(f"[{self._country}] page {page_id}: {len(rails)} rails") + return rails + + def _get_component_assets( + self, + component_id: str, + offset: int = 0, + page_size: Optional[int] = None, + ) -> List[Union[VodCategory, VodItem]]: + # FIXED: previously filtered to `isinstance(e, VodItem)` only, + # which silently dropped every VodCategory (series) now + # returned by _parse_asset. A rail can legitimately mix + # playable movies (VodItem) and drill-down series + # (VodCategory); both need to reach the caller. + page = self.get_component_assets_page(component_id, offset=offset, page_size=page_size) + return page.entries + + def get_component_assets_page( + self, + component_id: str, + offset: int = 0, + page_size: Optional[int] = None, + ) -> VodPage: + size = page_size or self.DEFAULT_ASSET_PAGE_SIZE + + data = self._request( + f"home/component/{component_id}/assets", + params={ + "offset": offset, + "page_size": size, + "device_type": "WEB", + "store_id": "MHR Production - Main", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="HOME", + step="RAIL", + ) + + entries: List[VodItem] = [] + for raw in data.get("assets") or []: + item = self._parse_asset(raw) + if item is not None: + entries.append(item) + + next_offset = data.get("next_offset") + normalised_next = None if next_offset in (None, -1) else int(next_offset) + + logger.debug( + f"[{self._country}] component {component_id} offset={offset}: " + f"{len(entries)} items, next={normalised_next}" + ) + return VodPage(entries=entries, next_offset=normalised_next) + + # ================================================================== + # Public API -- details + # ================================================================== + + def get_series_detail(self, series_id: str) -> Dict[str, Any]: + """Raw series-actions response.""" + data = self._request( + f"details/series/{series_id}/actions/v2", + params={ + "interacted_with_nPVR": "false", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="SERIES_DETAIL", + step="SERIES_WATCH_ACTION", # FIXED: capture shows this, not SERIES_ACTIONS + ) + return data + + def get_season_episodes( + self, + series_id: str, + season_id: str, + season_number: int, + ) -> List[VodItem]: + """ + Episodes of one season, as VodItems. + + KNOWN LIMITATION: `season_id` must already be known (it comes + from the season list inside a series-actions response), but the + shape of that seasons list has not been captured -- every + series-actions capture so far has had exactly one season, so + it's unclear whether seasons come back as a `seasons: [...]` + array on the series-actions response, a separate endpoint, or + something else for a multi-season show. Callers currently have + no verified way to enumerate seasons for a series with more + than one season; this only works if the season_id is already + known some other way (e.g. surfaced elsewhere in the UI flow). + """ + data = self._request( + f"details/series/{series_id}/season/{season_id}/v2", + params={ + "interacted_with_nPVR": "false", + "season_number": season_number, + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="SERIES_DETAIL", + step="SEASON_EPISODE_LIST", + ) + + episodes: List[VodItem] = [] + for raw in data.get("episodes") or []: + item = self._parse_episode(raw, series_id=series_id) + if item is not None: + episodes.append(item) + return episodes + + def get_program_detail(self, program_id: str) -> Dict[str, Any]: + """Raw program-actions response (movies, one-off programmes).""" + data = self._request( + f"details/program/{program_id}/actions/v2", + params={ + "interacted_with_nPVR": "false", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="SERIES_DETAIL", + step="PROGRAM_ACTIONS", # NOTE: unconfirmed -- no headers were + # captured for this specific call, only its response body. + # Verify against a fresh capture before relying on this step + # value for anything that inspects X-Tv-Step server-side. + ) + return data + + def get_program_metadata(self, program_id: str) -> Dict[str, Any]: + """ + NEW. Raw `details/program/{id}` response (no /actions/v2 suffix). + + The real client calls this ahead of the /actions/v2 call for + single-program detail pages (flow=SINGLE_PROGRAM_DETAIL, + step=PROGRAM_METADATA) to get description/cast/genre fields + that /actions/v2 does not carry. Only the request (headers + + URL) was captured, not a response body, so this deliberately + returns the raw dict rather than a typed object -- do not + assume specific field names here until a body capture confirms + them. Callers that want description/genres for a program + should call this and merge defensively (e.g. `.get("description")`, + `.get("metadata")` following the same GENRES-tagged-list shape + used elsewhere in this file) rather than assuming it 1:1 + matches get_program_detail() or the episode `details` shape. + """ + return self._request( + f"details/program/{program_id}", + params={ + "interacted_with_nPVR": "false", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="SINGLE_PROGRAM_DETAIL", + step="PROGRAM_METADATA", + ) + + def get_related_content( + self, + program_id: str, + is_series: bool, + size: int = 16, + offset: int = 0, + ) -> List[VodItem]: + """ + NEW. `relatedcontent/feed` -- "more like this" rail on a detail page. + + Confirmed asset shape from capture: + {"id", "title", "type": "Movie", "content_type": "Program", + "thumbnail", "cta": {"deeplink"}, "ratings", "release_year"?} + `release_year` was present on some entries and absent on others + in the capture -- treated as optional here. + """ + data = self._request( + "relatedcontent/feed", + params={ + "program_id": program_id, + "live": "false", + "size": size, + "offset": offset, + "is_series": "true" if is_series else "false", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="SERIES_DETAIL" if is_series else "SINGLE_PROGRAM_DETAIL", + step="RELATED_CONTENT", + ) + + items: List[VodItem] = [] + for raw in data.get("assets") or []: + content_id = raw.get("id") + title = raw.get("title") + if not content_id or not title: + continue + items.append(VodItem( + name=title, + content_id=content_id, + provider="magentaeu", + logo_url=raw.get("thumbnail"), + mode=StreamingMode.VOD, + content_type=ContentType.MOVIE, + release_year=raw.get("release_year"), + rating=raw.get("ratings"), + streaming_format=STREAMING_FORMAT_DASH, + country=self._country.upper(), + language=self._app_language, + )) + return items + + def get_program_playinfo(self, program_id: str, video_id: str) -> Dict[str, Any]: + """Raw playinfo response for a program (movie).""" + data = self._request( + "player/playinfo/program", + params={ + "program_id": program_id, + "video_id": video_id, + "device_type": "WEB", + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="PLAYER", + step="VOD_PLAYBACK", + ) + return data + + # ================================================================== + # Public API -- search + # ================================================================== + + def search( + self, + query: str, + item_type: Optional[str] = None, + size: Optional[int] = None, + ) -> SearchResult: + if not query: + return SearchResult(entries=[], has_more_items=False) + + params = { + "device_type": "WEB", + "text_search": query, + "size": size or self.DEFAULT_SEARCH_SIZE, + "app_language": self._app_language, + "natco_code": self._country, + } + if item_type: + params["item_type"] = item_type + + data = self._request( + "search/items", + params=params, + flow="SEARCH", + step="SEARCH_ALL", + ) + + entries: List[Union[VodCategory, VodItem]] = [] + + for raw in data.get("tv_shows") or []: + cat = self._parse_search_series(raw) + if cat is not None: + entries.append(cat) + + for raw in data.get("movies") or []: + if raw.get("is_live"): + continue + item = self._parse_search_movie(raw) + if item is not None: + entries.append(item) + + streaming = data.get("streaming") or [] + if streaming: + logger.warning( + f"[{self._country}] search returned {len(streaming)} " + f"streaming[] entries -- parser not yet implemented, dropped." + ) + + return SearchResult( + entries=entries, + has_more_items=bool(data.get("has_more_items")), + ) + + # ================================================================== + # Public API -- playback + # ================================================================== + + def _resolve_playback(self, content_id: str, **kwargs) -> ResolvedPlayback: + """ + Resolve everything needed to play a VOD asset, exactly once per + (program_id, video_id) within the cache TTL. + + This is the single source of truth for both get_manifest() and + get_drm() -- previously each independently re-derived a video_id + via get_program_detail() and re-called get_program_playinfo(), + doubling the request count per playback and (for get_drm) + building the licence URL from the wrong field entirely. + + Raises VodCatchupRequiredError if the resolved program has no + playable watch/trailer action but does have schedule data, + signalling this is a linear-catchup item, not VOD. + """ + video_id = kwargs.get("video_id") + program_id = kwargs.get("program_id") or content_id + + if not video_id: + video_id = self._resolve_playable_video_id(program_id) + + cache_key = (program_id, video_id) + cached = self._playback_cache.get(cache_key) + if cached is not None: + resolved, resolved_at = cached + if time.time() - resolved_at < self._PLAYBACK_CACHE_TTL_SECONDS: + return resolved + + playinfo = self.get_program_playinfo(program_id, video_id) + + services = (playinfo.get("playBackInfoResponse") or {}).get("service_items") or [] + chosen = None + for item in services: + if not item.get("is_trailer"): + chosen = item + break + if chosen is None and services: + chosen = services[0] + + if chosen is None: + raise VodError( + f"No service_items in playinfo response for {program_id}/{video_id}" + ) + + media_id = chosen["media_id"] + content_type = (playinfo.get("playBackInfoResponse") or {}).get( + "content_type", "tvod" + ) + + media = self._request( + "media", + params={ + "client_id": self._auth.current_token.device_id or "", + "media_id": media_id, + "src_format": "MPEG-DASH", + "content_type": content_type, + "app_language": self._app_language, + "natco_code": self._country, + }, + flow="PLAYER", + step="MEDIA_CALL", + ) + + video = media.get("video") or {} + src = video.get("video_src") or video.get("ref_src") + if not src: + raise VodError(f"No video_src in /media response for media_id={media_id}") + + # FIXED: this is the actual releasePid the real client sends to + # the Widevine licence endpoint -- NOT playinfo's title_id or + # program_id (that was the original proposal's DRM bug). + release_pid = video.get("pid") + if not release_pid: + raise VodError(f"No pid in /media response for media_id={media_id}") + + resolved = ResolvedPlayback( + video_src=src, + release_pid=release_pid, + content_type=content_type, + media_id=media_id, + ) + self._playback_cache[cache_key] = (resolved, time.time()) + return resolved + + def get_manifest(self, content_id: str, **kwargs) -> Optional[str]: + resolved = self._resolve_playback(content_id, **kwargs) + logger.debug(f"[{self._country}] manifest for {content_id}: {resolved.video_src}") + return resolved.video_src + + def get_drm(self, content_id: str, **kwargs) -> List[DRMConfig]: + resolved = self._resolve_playback(content_id, **kwargs) + + token = self._auth.current_token + if not isinstance(token, MagentaAuthToken): + raise VodAuthError("No MagentaAuthToken available for DRM") + + persona_jwt = token.persona_jwt + account_uri = token.account_uri + if not persona_jwt or not account_uri: + self._auth.get_user_account(force_refresh=True) + token = self._auth.current_token + persona_jwt = token.persona_jwt + account_uri = token.account_uri + + if not persona_jwt: + raise VodAuthError("dc_cts_personaToken missing from access token") + if not account_uri: + raise VodAuthError("account_url missing from /user/account response") + + licence_url = build_licence_url( + widevine_endpoint=WV_URL, + release_pid=resolved.release_pid, # FIXED: was title_id/program_id + persona_jwt=persona_jwt, + account_uri=account_uri, + ) + + drm = build_widevine_drm_config( + licence_url=licence_url, + user_agent=USER_AGENT, + origin=self._origin, + ) + return [drm] + + # ================================================================== + # Internal -- HTTP + # ================================================================== + + def _build_headers(self, flow: str, step: str) -> Dict[str, str]: + device_id, session_id = self._auth.get_guest_session_ids() + return build_auth_headers( + country=self._country, + device_id=device_id, + session_id=session_id, + flow=flow, + step=step, + call_type="AUTH_USER", + ) + + def _request( + self, + path: str, + params: Dict[str, Any], + flow: str, + step: str, + *, + retry_on_auth: bool = True, + timeout: int = DEFAULT_REQUEST_TIMEOUT, + ) -> Dict[str, Any]: + """ + Single HTTP entry point for every VOD call. + + FIXED: retry_on_auth was previously accepted but never acted + on -- 401 always propagated straight out. It now actually + triggers one forced token refresh + one retry. + """ + token = self._auth.get_bearer_token() + if token.startswith("Bearer "): + token = token[7:] + + url = f"{self._bifrost_url}/{path.lstrip('/')}" + + full_params = dict(params) + full_params["natco_key"] = self._natco_key + + headers = self._build_headers(flow, step) + headers["Bff_token"] = token + headers["X-Account-Details"] = self._account_details_header() + + cmid = self._auth.current_token.channel_map_id + if cmid: + headers["X-Channel-Map-Id"] = str(cmid) + + try: + response = self._http.get( + url, + operation=f"vod_{path.replace('/', '_')}", + headers=headers, + params=full_params, + timeout=timeout, + ) + except Exception as exc: + raise VodError(f"Transport error on {path}: {exc}", url=url) from exc + + status = getattr(response, "status_code", None) + if status == 200: + return response.json() + + try: + self._raise_for_status(status, response, url, path) + except VodAuthError: + if not retry_on_auth: + raise + logger.info( + f"[{self._country}] 401 on {path}, refreshing token and retrying once" + ) + self._auth.get_bearer_token(force_refresh=True) + return self._request( + path, params, flow, step, + retry_on_auth=False, timeout=timeout, + ) + + # Defensive -- _raise_for_status always raises on non-200. + raise VodError(f"Unreachable: status {status}", status=status, url=url) + + @staticmethod + def _raise_for_status(status: int, response, url: str, path: str) -> None: + body_snippet = "" + try: + body_snippet = (response.text or "")[:512] + except Exception: + pass + + message = f"{status} from {path}: {body_snippet}" + + if status == 401: + raise VodAuthError(message, status=status, url=url) + if status == 403: + lower = body_snippet.lower() + if "geo" in lower or "region" in lower: + raise VodGeoBlockError(message, status=status, url=url) + raise VodEntitlementError(message, status=status, url=url) + if status == 404: + raise VodNotFoundError(message, status=status, url=url) + if status == 429: + raise VodRateLimitError(message, status=status, url=url) + if 500 <= status < 600: + raise VodServerError(message, status=status, url=url) + raise VodError(message, status=status, url=url) + + def _account_details_header(self) -> str: + token = self._auth.current_token + info = token.account_info or {} + + account_type = info.get("account_type") or "MHR_default" + user_id = info.get("tvAccountId") or info.get("user_id") or "" + channel_map_id = info.get("channelMap_id") or token.channel_map_id or "" + account_id = info.get("account_id") or "" + vod_enabled = str(token.vod_enabled).lower() if token.vod_enabled is not None else "false" + + import json + return json.dumps({ + "accountType": account_type, + "userId": user_id, + "recordingEnabledDVR": False, + "channelMapId": channel_map_id, + "rightsGroupIds": "", + "releaseSlot": "", + "accountId": account_id, + "vodEnabled": vod_enabled, + }) + + # ================================================================== + # Internal -- entitlement / catchup detection + # ================================================================== + + def _resolve_playable_video_id(self, program_id: str) -> str: + """ + Given a program id, return the media_id (video_id) of a + playable action. + + Order of preference: + 1. actions.watch[0].video_id (entitled -- UNVERIFIED, see + module docstring: every capture so far has an account + with vod_enabled=false, so this branch has never fired) + 2. actions.trailer[0].video_id (trailer -- confirmed working + against the capture, "Tom i Jerry" case) + 3. SPECULATIVE, likely dead: if watch/trailer are both + empty, check for schedules/catchup_schedules and raise + VodCatchupRequiredError instead of a plain entitlement + error. See the big comment at the check itself -- this + field pair is confirmed to exist on the *series*-actions + response (a different endpoint, + details/series/{id}/actions/v2, called by + get_series_detail()), but every *program*-actions + response captured so far (details/program/{id}/actions/v2 + -- the endpoint this method actually calls) has neither + key at all, whether for a movie or an episode. As + written this branch is reachable but has never actually + fired against real data. Left in as defensive coding + with this caveat rather than removed outright -- but do + not treat it as confirmed, and do not build further + logic on top of it until a capture shows the fields + present at THIS endpoint. + 4. raise VodEntitlementError + """ + detail = self.get_program_detail(program_id) + actions = detail.get("actions") or {} + + for key in ("watch", "trailer"): + bucket = actions.get(key) or [] + if bucket: + first = bucket[0] + vid = ( + first.get("video_id") + or (first.get("video") or {}).get("video_id") + ) + if vid: + return vid + + # SPECULATIVE (see docstring point 3 above): schedules / + # catchup_schedules are confirmed present on the *series*-level + # actions response (get_series_detail()'s endpoint), not on the + # *program*-level actions response this method calls. Kept as a + # defensive no-op-in-practice check rather than removed, since + # it's cheap and correct IF the fields ever do appear here -- + # but do not rely on it firing until a program/episode-actions + # capture confirms it. + schedules = actions.get("schedules") or [] + catchup_schedules = actions.get("catchup_schedules") or [] + if schedules or catchup_schedules: + station_id = None + if schedules: + station_id = schedules[0].get("station_id") + elif catchup_schedules: + station_id = catchup_schedules[0].get("station_id") + raise VodCatchupRequiredError( + f"{program_id} has no watch/trailer action but has " + f"schedule data -- this is a linear catch-up item, not " + f"VOD. Route to provider.get_catchup_manifest() instead.", + station_id=station_id, + schedules=schedules, + catchup_schedules=catchup_schedules, + ) + + msg = actions.get("svod_subscription_message") + if msg: + raise VodEntitlementError(f"No playable action for {program_id}: {msg}") + raise VodEntitlementError( + f"No playable action for {program_id} (actions.watch, " + f"actions.trailer, and schedule data all empty)" + ) + + # ================================================================== + # Internal -- parsers + # ================================================================== + + @staticmethod + def _parse_category_root(raw: Dict[str, Any]) -> Optional[VodCategory]: + page_id = raw.get("page_id") or raw.get("id") + title = raw.get("title") or raw.get("name") + if not page_id or not title: + return None + return VodCategory( + name=title.strip(), + content_id=page_id, + provider="magentaeu", + ) + + def _parse_rail(self, raw: Dict[str, Any]) -> Optional[VodCategory]: + comp_id = raw.get("id") + title = raw.get("title") + if not comp_id or not title: + return None + + content_details = raw.get("content_details") or {} + end_point = content_details.get("end_point") or "" + + if end_point.startswith("recommendations/"): + logger.debug( + f"[{self._country}] rail {comp_id} is a recommendations " + f"rail (end_point={end_point!r}); children not resolvable" + ) + + return VodCategory( + name=title.strip(), + content_id=comp_id, + provider="magentaeu", + child_count=None, + details_url=end_point or None, + fetch_url=None, + ) + + def _parse_asset(self, raw: Dict[str, Any]) -> Optional[Union[VodCategory, VodItem]]: + content_id = raw.get("id") + title = raw.get("title") + if not content_id or not title: + return None + + item_type = raw.get("type") or "" + content_type = raw.get("content_type") or "" + + # FIXED: series assets on a rail were previously dropped + # entirely (returned None). Against the capture, the entire + # "SERIJE" page's rail assets are content_type=="Series" -- + # dropping them broke top-level series navigation completely. + # Return a VodCategory drill-down node instead; the caller + # (get_category_children) already knows how to descend into a + # VodCategory via series_id. + if item_type == "TVShow" or content_type == "Series": + return VodCategory( + name=title, + content_id=content_id, + provider="magentaeu", + logo_url=raw.get("thumbnail"), + ) + + return VodItem( + name=title, + content_id=content_id, + provider="magentaeu", + logo_url=raw.get("thumbnail"), + mode=StreamingMode.VOD, + content_type=ContentType.MOVIE, + release_year=raw.get("release_year"), + rating=raw.get("ratings"), + streaming_format=STREAMING_FORMAT_DASH, + country=self._country.upper(), + language=self._app_language, + ) + + def _parse_episode( + self, + raw: Dict[str, Any], + series_id: Optional[str] = None, + ) -> Optional[VodItem]: + content_id = raw.get("id") + if not content_id: + return None + + episode_number: Optional[int] = None + try: + episode_number = int(raw.get("number")) + except (TypeError, ValueError): + pass + + details = raw.get("details") or {} + description = details.get("description") + + genres = None + for meta in details.get("metadata") or []: + if meta.get("type") == "GENRES" and meta.get("value"): + genres = [g.strip() for g in meta["value"].split(",") if g.strip()] or None + break + + runtime = raw.get("runtime_seconds") + try: + runtime = int(runtime) if runtime is not None else None + except (TypeError, ValueError): + runtime = None + + return VodItem( + name=raw.get("name") or f"Episode {raw.get('number')}", + content_id=content_id, + provider="magentaeu", + logo_url=raw.get("poster_image_url"), + mode=StreamingMode.VOD, + content_type=ContentType.SERIES, + description=description, + release_year=raw.get("release_year"), + rating=raw.get("ratings"), + genres=genres, + genre=genres[0] if genres else None, + duration_seconds=runtime, + episode_number=episode_number, + series_id=series_id, + streaming_format=STREAMING_FORMAT_DASH, + country=self._country.upper(), + language=self._app_language, + ) + + @staticmethod + def _parse_search_series(raw: Dict[str, Any]) -> Optional[VodCategory]: + series_id = raw.get("id") + title = raw.get("name") + if not series_id or not title: + return None + return VodCategory( + name=title, + content_id=series_id, + provider="magentaeu", + logo_url=raw.get("poster_image_url"), + ) + + def _parse_search_movie(self, raw: Dict[str, Any]) -> Optional[VodItem]: + content_id = raw.get("id") + title = raw.get("name") + if not content_id or not title: + return None + + runtime = raw.get("runtime_seconds") + try: + runtime = int(runtime) if runtime is not None else None + except (TypeError, ValueError): + runtime = None + + return VodItem( + name=title, + content_id=content_id, + provider="magentaeu", + logo_url=raw.get("poster_image_url"), + mode=StreamingMode.VOD, + content_type=ContentType.MOVIE, + release_year=raw.get("release_year"), + rating=raw.get("ratings"), + duration_seconds=runtime, + streaming_format=STREAMING_FORMAT_DASH, + country=self._country.upper(), + language=self._app_language, + ) + + # ================================================================== + # Internal -- serialisation helper (for debugging) + # ================================================================== + + @staticmethod + def as_dict(entries: List[Union[VodCategory, VodItem]]) -> List[Dict]: + return [e.to_dict() for e in entries] \ No newline at end of file