Files
script.service.ultimate/lib/streaming_providers/base/managed_provider.py
T
2026-10-07 12:12:10 +02:00

376 lines
14 KiB
Python

# 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