sync.py 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282
  1. # -*- coding: utf-8 -*-
  2. """
  3. Katalog-Sync-Manager fuer den Streaming-Katalog.
  4. Fuehrt woechentliche Vollsynchronisierungen und taegliche
  5. Schnellaktualisierungen durch. Der Sync ist Resume-faehig:
  6. Nach jedem erfolgreich synchronisierten Dienst wird der Status
  7. gespeichert, sodass ein unterbrochener Sync fortgesetzt werden kann.
  8. """
  9. from __future__ import annotations
  10. import json
  11. from datetime import datetime, timezone
  12. from pathlib import Path
  13. from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
  14. from .cache import StreamingCache
  15. from .tmdb_discover import TMDBDiscoverClient
  16. class CatalogSync:
  17. """
  18. Woechentlicher Katalog-Sync mit Resume-Faehigkeit.
  19. Verwaltet den Synchronisierungszustand und koordiniert
  20. den Discover-Client mit dem StreamingCache. Nach jedem
  21. synchronisierten Dienst wird der Status persistiert.
  22. """
  23. STATE_FILE = "streaming_sync_state.json"
  24. # Standard-Sync-Intervall: 7 Tage
  25. _SYNC_INTERVAL_DAYS: int = 7
  26. def __init__(
  27. self,
  28. config: dict,
  29. cache: StreamingCache,
  30. discover: TMDBDiscoverClient,
  31. ) -> None:
  32. """
  33. Args:
  34. config: Plugin-Konfiguration mit ``cache_dir``, ``services``, etc.
  35. cache: StreamingCache-Instanz.
  36. discover: TMDBDiscoverClient-Instanz.
  37. """
  38. self._config = config
  39. self._cache = cache
  40. self._discover = discover
  41. cache_dir = Path(config.get("cache_dir", "cache/streaming"))
  42. self._state_path = cache_dir / self.STATE_FILE
  43. self._state: dict = {}
  44. # Verzeichnis sicherstellen
  45. cache_dir.mkdir(parents=True, exist_ok=True)
  46. # --- Status-Verwaltung ---
  47. def load_state(self) -> None:
  48. """
  49. Laedt den Sync-Status aus der JSON-Datei.
  50. Toleriert fehlende oder beschaedigte Dateien.
  51. """
  52. if not self._state_path.exists():
  53. pdebug("CatalogSync: Keine Status-Datei vorhanden, starte frisch")
  54. self._state = {}
  55. return
  56. try:
  57. raw = self._state_path.read_text(encoding="utf-8")
  58. self._state = json.loads(raw)
  59. pdebug(f"CatalogSync: Status geladen — letzter Sync: {self._state.get('last_sync', 'nie')}")
  60. except (json.JSONDecodeError, OSError) as e:
  61. perror(f"CatalogSync: Fehler beim Laden der Status-Datei: {e}")
  62. self._state = {}
  63. def save_state(self) -> None:
  64. """
  65. Speichert den Sync-Status atomar in die JSON-Datei.
  66. Nutzt tmp+rename fuer Atomizitaet.
  67. """
  68. tmp_path = self._state_path.with_suffix(".tmp")
  69. try:
  70. tmp_path.write_text(
  71. json.dumps(self._state, ensure_ascii=False, indent=2),
  72. encoding="utf-8",
  73. )
  74. tmp_path.replace(self._state_path)
  75. pdebug("CatalogSync: Status gespeichert")
  76. except OSError as e:
  77. perror(f"CatalogSync: Fehler beim Speichern der Status-Datei: {e}")
  78. # --- Sync-Steuerung ---
  79. def needs_sync(self) -> bool:
  80. """
  81. Prueft, ob ein neuer Vollsync noetig ist.
  82. Returns:
  83. True, wenn der letzte Sync aelter als das Sync-Intervall ist
  84. oder noch nie synchronisiert wurde.
  85. """
  86. last_sync = self._state.get("last_sync")
  87. if not last_sync:
  88. return True
  89. try:
  90. last_dt = datetime.fromisoformat(last_sync)
  91. now = datetime.now(timezone.utc)
  92. age_days = (now - last_dt).days
  93. return age_days >= self._SYNC_INTERVAL_DAYS
  94. except (ValueError, TypeError):
  95. pwarn("CatalogSync: Ungueltiges Datum in Status-Datei, Sync noetig")
  96. return True
  97. async def run_full_sync(self) -> dict:
  98. """
  99. Fuehrt einen vollstaendigen Katalog-Sync fuer alle aktivierten Dienste durch.
  100. Resume-sicher: Verfolgt, welche Dienste bereits synchronisiert wurden.
  101. Bei Unterbrechung werden nur die fehlenden Dienste nachgeholt.
  102. Returns:
  103. Stats-Dict mit service_id -> Titelanzahl.
  104. """
  105. services = self._config.get("services", {})
  106. if not services:
  107. pwarn("CatalogSync: Keine Dienste konfiguriert")
  108. return {}
  109. # Bereits synchronisierte Dienste aus dem aktuellen Lauf
  110. synced = set(self._state.get("services_synced", []))
  111. stats: dict[str, int] = {}
  112. pinfo(f"CatalogSync: Starte Vollsync fuer {len(services)} Dienste")
  113. for service_id, service_config in services.items():
  114. # Deaktivierte Dienste ueberspringen
  115. if not service_config.get("enabled", True):
  116. pdebug(f"CatalogSync: '{service_id}' deaktiviert, ueberspringe")
  117. continue
  118. # Bereits synchronisierte Dienste ueberspringen (Resume)
  119. if service_id in synced:
  120. pdebug(f"CatalogSync: '{service_id}' bereits synchronisiert, ueberspringe")
  121. continue
  122. try:
  123. count = await self._sync_service(service_id, service_config)
  124. stats[service_id] = count
  125. # Status nach jedem Dienst speichern (Resume-Punkt)
  126. synced.add(service_id)
  127. self._state["services_synced"] = list(synced)
  128. self.save_state()
  129. # Cache nach jedem Dienst speichern
  130. await self._cache.save_async()
  131. except Exception as exc:
  132. perror(f"CatalogSync: Fehler beim Sync von '{service_id}': {exc}")
  133. stats[service_id] = -1
  134. # Vollsync abgeschlossen — Status zuruecksetzen
  135. self._state["last_sync"] = datetime.now(timezone.utc).isoformat(timespec="seconds")
  136. self._state["services_synced"] = []
  137. self._state["last_full_sync_stats"] = stats
  138. self.save_state()
  139. total = sum(c for c in stats.values() if c > 0)
  140. pinfo(f"CatalogSync: Vollsync abgeschlossen — {total} Titel in {len(stats)} Diensten")
  141. return stats
  142. async def run_daily_update(self) -> dict:
  143. """
  144. Fuehrt eine schnelle taegliche Aktualisierung durch.
  145. Prueft nur die erste Seite pro Dienst (kuerzlich hinzugefuegte Titel).
  146. Speichert die neuen Titel als "recent" im Cache.
  147. Returns:
  148. Stats-Dict mit service_id -> Anzahl neuer Titel.
  149. """
  150. services = self._config.get("services", {})
  151. if not services:
  152. return {}
  153. stats: dict[str, int] = {}
  154. days = self._config.get("recent_days", 7)
  155. pdebug(f"CatalogSync: Starte taegliches Update (letzte {days} Tage)")
  156. for service_id, service_config in services.items():
  157. if not service_config.get("enabled", True):
  158. continue
  159. try:
  160. titles = await self._discover.discover_recent(service_id, days=days)
  161. recent_slugs: list[str] = []
  162. for title_data in titles:
  163. slug = StreamingCache._generate_slug(title_data.get("t", ""))
  164. if slug:
  165. self._cache.update_title(slug, title_data)
  166. recent_slugs.append(slug)
  167. self._cache.set_recent(service_id, recent_slugs)
  168. stats[service_id] = len(recent_slugs)
  169. except Exception as exc:
  170. perror(f"CatalogSync: Fehler beim Daily-Update von '{service_id}': {exc}")
  171. stats[service_id] = -1
  172. # Cache speichern
  173. await self._cache.save_async()
  174. # Status aktualisieren
  175. self._state["last_daily_update"] = datetime.now(timezone.utc).isoformat(timespec="seconds")
  176. self._state["last_daily_stats"] = stats
  177. self.save_state()
  178. total = sum(c for c in stats.values() if c > 0)
  179. pdebug(f"CatalogSync: Taegliches Update abgeschlossen — {total} neue Titel")
  180. return stats
  181. async def _sync_service(self, service_id: str, service_config: dict) -> int:
  182. """
  183. Synchronisiert einen einzelnen Streaming-Dienst.
  184. Ruft alle Titel ueber den Discover-Client ab, aktualisiert den Cache
  185. und entfernt Titel, die nicht mehr verfuegbar sind (Pruning).
  186. Args:
  187. service_id: ID des Streaming-Dienstes.
  188. service_config: Konfiguration des Dienstes.
  189. Returns:
  190. Anzahl der synchronisierten Titel.
  191. """
  192. max_titles = service_config.get("max_titles", 500)
  193. pinfo(f"CatalogSync: Synchronisiere '{service_id}' (max {max_titles} Titel)")
  194. titles = await self._discover.discover_service(service_id, max_titles=max_titles)
  195. if not titles:
  196. pwarn(f"CatalogSync: Keine Titel fuer '{service_id}' gefunden")
  197. return 0
  198. # Titel in den Cache eintragen
  199. current_slugs: set[str] = set()
  200. for title_data in titles:
  201. slug = StreamingCache._generate_slug(title_data.get("t", ""))
  202. if slug:
  203. self._cache.update_title(slug, title_data)
  204. current_slugs.add(slug)
  205. # Pruning: Titel entfernen, die nicht mehr auf dem Dienst sind
  206. pruned = self._cache.prune_missing(service_id, current_slugs)
  207. if pruned > 0:
  208. pdebug(f"CatalogSync: {pruned} veraltete Eintraege fuer '{service_id}' entfernt")
  209. # Dienst-Info aktualisieren
  210. self._cache.update_services({
  211. service_id: {
  212. "name": service_config.get("name", service_id),
  213. "last_sync": datetime.now(timezone.utc).isoformat(timespec="seconds"),
  214. "title_count": len(current_slugs),
  215. },
  216. })
  217. # Indizes neu aufbauen nach allen Aenderungen
  218. self._cache._rebuild_indices()
  219. pinfo(f"CatalogSync: '{service_id}' synchronisiert — {len(current_slugs)} Titel")
  220. return len(current_slugs)