| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201 |
- # -*- coding: utf-8 -*-
- """
- EPG-Daten-Collector.
- Laedt XMLTV-Daten von konfigurierten Quellen via aiohttp herunter,
- cached die Rohdaten und dedupliziert die geparsten Sendungen.
- """
- from __future__ import annotations
- import asyncio
- import gzip
- import hashlib
- from pathlib import Path
- from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
- from plugins.tv_program.epg.channels import ChannelMapper
- from plugins.tv_program.epg.xmltv_parser import XMLTVParser
- # Optionale aiohttp-Abhaengigkeit
- _AIOHTTP_AVAILABLE = False
- try:
- import aiohttp
- _AIOHTTP_AVAILABLE = True
- except ImportError:
- pass
- # Download-Timeout in Sekunden
- _FETCH_TIMEOUT = 60
- class EPGCollector:
- """
- Sammelt EPG-Daten von konfigurierten XMLTV-Quellen.
- Laedt XMLTV-Dateien asynchron herunter, parst sie und
- dedupliziert die Ergebnisse. Rohdaten werden lokal gecached.
- """
- def __init__(
- self,
- sources: list[dict],
- cache_dir: Path,
- channel_mapper: ChannelMapper,
- ) -> None:
- """
- Args:
- sources: Liste von Quellen-Dicts mit mindestens "url" und "name".
- Optional "enabled" (default True).
- cache_dir: Verzeichnis fuer den Download-Cache.
- channel_mapper: Mapper zur XMLTV-ID-Aufloesung.
- """
- self._sources = sources
- self._cache_dir = cache_dir
- self._raw_dir = cache_dir / "raw"
- self._mapper = channel_mapper
- self._parser = XMLTVParser(channel_mapper)
- # Cache-Verzeichnisse erstellen
- self._raw_dir.mkdir(parents=True, exist_ok=True)
- async def fetch_all(self) -> list[dict]:
- """
- Laedt alle aktivierten Quellen herunter und gibt deduplizierte
- Programm-Dicts zurueck.
- Returns:
- Zusammengefuehrte, deduplizierte Liste aller Sendungen.
- """
- if not _AIOHTTP_AVAILABLE:
- perror("EPGCollector: aiohttp nicht installiert, Download nicht moeglich")
- return []
- all_programs: list[dict] = []
- enabled_sources = [
- s for s in self._sources if s.get("enabled", True)
- ]
- if not enabled_sources:
- pwarn("EPGCollector: Keine aktiven EPG-Quellen konfiguriert")
- return []
- pinfo(f"EPGCollector: Lade {len(enabled_sources)} EPG-Quelle(n)...")
- for source in enabled_sources:
- name = source.get("name", "unbekannt")
- try:
- raw_data = await self._fetch_source(source)
- if raw_data:
- # XML-Parsen in Thread auslagern (CPU-intensiv, ~500ms fuer 11 MB)
- # damit der Event-Loop nicht blockiert wird
- programs = await asyncio.to_thread(self._parser.parse_bytes, raw_data)
- pinfo(f"EPGCollector: {len(programs)} Sendungen von '{name}' geparst")
- all_programs.extend(programs)
- else:
- pwarn(f"EPGCollector: Keine Daten von '{name}' erhalten")
- except Exception as e:
- perror(f"EPGCollector: Fehler bei Quelle '{name}': {e}")
- # Deduplizieren
- deduplicated = self._deduplicate(all_programs)
- pinfo(
- f"EPGCollector: {len(deduplicated)} Sendungen nach Deduplizierung "
- f"(von {len(all_programs)} gesamt)"
- )
- return deduplicated
- async def _fetch_source(self, source: dict) -> bytes:
- """
- Laedt eine einzelne XMLTV-Quelle herunter.
- Speichert die Rohdaten zusaetzlich im Cache-Verzeichnis
- fuer Debugging-Zwecke.
- Args:
- source: Quellen-Dict mit "url" und "name".
- Returns:
- Rohe XML-Bytes der heruntergeladenen Datei.
- Raises:
- aiohttp.ClientError: Bei Netzwerkfehlern.
- asyncio.TimeoutError: Bei Timeout.
- """
- url = source["url"]
- name = source.get("name", "unbekannt")
- pdebug(f"EPGCollector: Lade '{name}' von {url}")
- timeout = aiohttp.ClientTimeout(total=_FETCH_TIMEOUT)
- async with aiohttp.ClientSession(timeout=timeout) as session:
- async with session.get(url) as response:
- response.raise_for_status()
- data = await response.read()
- # gzip-komprimierte Dateien entpacken
- if url.endswith(".gz"):
- try:
- data = gzip.decompress(data)
- pdebug(f"EPGCollector: gzip entpackt: {len(data)} Bytes")
- except (gzip.BadGzipFile, OSError) as e:
- perror(f"EPGCollector: gzip-Entpacken fehlgeschlagen: {e}")
- return b""
- # Rohdaten im Cache speichern
- cache_file = self._raw_dir / f"{name}.xml"
- try:
- cache_file.write_bytes(data)
- pdebug(f"EPGCollector: Rohdaten gespeichert: {cache_file}")
- except OSError as e:
- pwarn(f"EPGCollector: Konnte Rohdaten nicht cachen: {e}")
- pdebug(f"EPGCollector: {len(data)} Bytes von '{name}' heruntergeladen")
- return data
- @staticmethod
- def _deduplicate(programs: list[dict]) -> list[dict]:
- """
- Entfernt Duplikate basierend auf Programm-ID (Hash aus Sender + Startzeit).
- Bei Duplikaten wird der Eintrag mit mehr Informationen bevorzugt.
- Args:
- programs: Liste aller geparsten Sendungen.
- Returns:
- Deduplizierte Liste.
- """
- seen: dict[str, dict] = {}
- for program in programs:
- prog_id = EPGCollector._generate_id(program)
- if prog_id not in seen:
- seen[prog_id] = program
- else:
- # Eintrag mit mehr Feldern bevorzugen
- existing = seen[prog_id]
- if len(program) > len(existing):
- seen[prog_id] = program
- return list(seen.values())
- @staticmethod
- def _generate_id(program: dict) -> str:
- """
- Generiert eine eindeutige ID fuer eine Sendung.
- Basiert auf MD5-Hash von Sender-ID und Startzeit.
- Args:
- program: Programm-Dict mit "ch" und "st" Feldern.
- Returns:
- 12-stelliger Hex-Hash als Programm-ID.
- """
- key = f"{program.get('ch', '')}_{program.get('s', '')}"
- return hashlib.md5(key.encode()).hexdigest()[:12]
|