collector.py 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. # -*- coding: utf-8 -*-
  2. """
  3. EPG-Daten-Collector.
  4. Laedt XMLTV-Daten von konfigurierten Quellen via aiohttp herunter,
  5. cached die Rohdaten und dedupliziert die geparsten Sendungen.
  6. """
  7. from __future__ import annotations
  8. import asyncio
  9. import gzip
  10. import hashlib
  11. from pathlib import Path
  12. from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
  13. from plugins.tv_program.epg.channels import ChannelMapper
  14. from plugins.tv_program.epg.xmltv_parser import XMLTVParser
  15. # Optionale aiohttp-Abhaengigkeit
  16. _AIOHTTP_AVAILABLE = False
  17. try:
  18. import aiohttp
  19. _AIOHTTP_AVAILABLE = True
  20. except ImportError:
  21. pass
  22. # Download-Timeout in Sekunden
  23. _FETCH_TIMEOUT = 60
  24. class EPGCollector:
  25. """
  26. Sammelt EPG-Daten von konfigurierten XMLTV-Quellen.
  27. Laedt XMLTV-Dateien asynchron herunter, parst sie und
  28. dedupliziert die Ergebnisse. Rohdaten werden lokal gecached.
  29. """
  30. def __init__(
  31. self,
  32. sources: list[dict],
  33. cache_dir: Path,
  34. channel_mapper: ChannelMapper,
  35. ) -> None:
  36. """
  37. Args:
  38. sources: Liste von Quellen-Dicts mit mindestens "url" und "name".
  39. Optional "enabled" (default True).
  40. cache_dir: Verzeichnis fuer den Download-Cache.
  41. channel_mapper: Mapper zur XMLTV-ID-Aufloesung.
  42. """
  43. self._sources = sources
  44. self._cache_dir = cache_dir
  45. self._raw_dir = cache_dir / "raw"
  46. self._mapper = channel_mapper
  47. self._parser = XMLTVParser(channel_mapper)
  48. # Cache-Verzeichnisse erstellen
  49. self._raw_dir.mkdir(parents=True, exist_ok=True)
  50. async def fetch_all(self) -> list[dict]:
  51. """
  52. Laedt alle aktivierten Quellen herunter und gibt deduplizierte
  53. Programm-Dicts zurueck.
  54. Returns:
  55. Zusammengefuehrte, deduplizierte Liste aller Sendungen.
  56. """
  57. if not _AIOHTTP_AVAILABLE:
  58. perror("EPGCollector: aiohttp nicht installiert, Download nicht moeglich")
  59. return []
  60. all_programs: list[dict] = []
  61. enabled_sources = [
  62. s for s in self._sources if s.get("enabled", True)
  63. ]
  64. if not enabled_sources:
  65. pwarn("EPGCollector: Keine aktiven EPG-Quellen konfiguriert")
  66. return []
  67. pinfo(f"EPGCollector: Lade {len(enabled_sources)} EPG-Quelle(n)...")
  68. for source in enabled_sources:
  69. name = source.get("name", "unbekannt")
  70. try:
  71. raw_data = await self._fetch_source(source)
  72. if raw_data:
  73. # XML-Parsen in Thread auslagern (CPU-intensiv, ~500ms fuer 11 MB)
  74. # damit der Event-Loop nicht blockiert wird
  75. programs = await asyncio.to_thread(self._parser.parse_bytes, raw_data)
  76. pinfo(f"EPGCollector: {len(programs)} Sendungen von '{name}' geparst")
  77. all_programs.extend(programs)
  78. else:
  79. pwarn(f"EPGCollector: Keine Daten von '{name}' erhalten")
  80. except Exception as e:
  81. perror(f"EPGCollector: Fehler bei Quelle '{name}': {e}")
  82. # Deduplizieren
  83. deduplicated = self._deduplicate(all_programs)
  84. pinfo(
  85. f"EPGCollector: {len(deduplicated)} Sendungen nach Deduplizierung "
  86. f"(von {len(all_programs)} gesamt)"
  87. )
  88. return deduplicated
  89. async def _fetch_source(self, source: dict) -> bytes:
  90. """
  91. Laedt eine einzelne XMLTV-Quelle herunter.
  92. Speichert die Rohdaten zusaetzlich im Cache-Verzeichnis
  93. fuer Debugging-Zwecke.
  94. Args:
  95. source: Quellen-Dict mit "url" und "name".
  96. Returns:
  97. Rohe XML-Bytes der heruntergeladenen Datei.
  98. Raises:
  99. aiohttp.ClientError: Bei Netzwerkfehlern.
  100. asyncio.TimeoutError: Bei Timeout.
  101. """
  102. url = source["url"]
  103. name = source.get("name", "unbekannt")
  104. pdebug(f"EPGCollector: Lade '{name}' von {url}")
  105. timeout = aiohttp.ClientTimeout(total=_FETCH_TIMEOUT)
  106. async with aiohttp.ClientSession(timeout=timeout) as session:
  107. async with session.get(url) as response:
  108. response.raise_for_status()
  109. data = await response.read()
  110. # gzip-komprimierte Dateien entpacken
  111. if url.endswith(".gz"):
  112. try:
  113. data = gzip.decompress(data)
  114. pdebug(f"EPGCollector: gzip entpackt: {len(data)} Bytes")
  115. except (gzip.BadGzipFile, OSError) as e:
  116. perror(f"EPGCollector: gzip-Entpacken fehlgeschlagen: {e}")
  117. return b""
  118. # Rohdaten im Cache speichern
  119. cache_file = self._raw_dir / f"{name}.xml"
  120. try:
  121. cache_file.write_bytes(data)
  122. pdebug(f"EPGCollector: Rohdaten gespeichert: {cache_file}")
  123. except OSError as e:
  124. pwarn(f"EPGCollector: Konnte Rohdaten nicht cachen: {e}")
  125. pdebug(f"EPGCollector: {len(data)} Bytes von '{name}' heruntergeladen")
  126. return data
  127. @staticmethod
  128. def _deduplicate(programs: list[dict]) -> list[dict]:
  129. """
  130. Entfernt Duplikate basierend auf Programm-ID (Hash aus Sender + Startzeit).
  131. Bei Duplikaten wird der Eintrag mit mehr Informationen bevorzugt.
  132. Args:
  133. programs: Liste aller geparsten Sendungen.
  134. Returns:
  135. Deduplizierte Liste.
  136. """
  137. seen: dict[str, dict] = {}
  138. for program in programs:
  139. prog_id = EPGCollector._generate_id(program)
  140. if prog_id not in seen:
  141. seen[prog_id] = program
  142. else:
  143. # Eintrag mit mehr Feldern bevorzugen
  144. existing = seen[prog_id]
  145. if len(program) > len(existing):
  146. seen[prog_id] = program
  147. return list(seen.values())
  148. @staticmethod
  149. def _generate_id(program: dict) -> str:
  150. """
  151. Generiert eine eindeutige ID fuer eine Sendung.
  152. Basiert auf MD5-Hash von Sender-ID und Startzeit.
  153. Args:
  154. program: Programm-Dict mit "ch" und "st" Feldern.
  155. Returns:
  156. 12-stelliger Hex-Hash als Programm-ID.
  157. """
  158. key = f"{program.get('ch', '')}_{program.get('s', '')}"
  159. return hashlib.md5(key.encode()).hexdigest()[:12]