# -*- coding: utf-8 -*- """ Katalog-Sync-Manager fuer den Streaming-Katalog. Fuehrt woechentliche Vollsynchronisierungen und taegliche Schnellaktualisierungen durch. Der Sync ist Resume-faehig: Nach jedem erfolgreich synchronisierten Dienst wird der Status gespeichert, sodass ein unterbrochener Sync fortgesetzt werden kann. """ from __future__ import annotations import json from datetime import datetime, timezone from pathlib import Path from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn from .cache import StreamingCache from .tmdb_discover import TMDBDiscoverClient class CatalogSync: """ Woechentlicher Katalog-Sync mit Resume-Faehigkeit. Verwaltet den Synchronisierungszustand und koordiniert den Discover-Client mit dem StreamingCache. Nach jedem synchronisierten Dienst wird der Status persistiert. """ STATE_FILE = "streaming_sync_state.json" # Standard-Sync-Intervall: 7 Tage _SYNC_INTERVAL_DAYS: int = 7 def __init__( self, config: dict, cache: StreamingCache, discover: TMDBDiscoverClient, ) -> None: """ Args: config: Plugin-Konfiguration mit ``cache_dir``, ``services``, etc. cache: StreamingCache-Instanz. discover: TMDBDiscoverClient-Instanz. """ self._config = config self._cache = cache self._discover = discover cache_dir = Path(config.get("cache_dir", "cache/streaming")) self._state_path = cache_dir / self.STATE_FILE self._state: dict = {} # Verzeichnis sicherstellen cache_dir.mkdir(parents=True, exist_ok=True) # --- Status-Verwaltung --- def load_state(self) -> None: """ Laedt den Sync-Status aus der JSON-Datei. Toleriert fehlende oder beschaedigte Dateien. """ if not self._state_path.exists(): pdebug("CatalogSync: Keine Status-Datei vorhanden, starte frisch") self._state = {} return try: raw = self._state_path.read_text(encoding="utf-8") self._state = json.loads(raw) pdebug(f"CatalogSync: Status geladen — letzter Sync: {self._state.get('last_sync', 'nie')}") except (json.JSONDecodeError, OSError) as e: perror(f"CatalogSync: Fehler beim Laden der Status-Datei: {e}") self._state = {} def save_state(self) -> None: """ Speichert den Sync-Status atomar in die JSON-Datei. Nutzt tmp+rename fuer Atomizitaet. """ tmp_path = self._state_path.with_suffix(".tmp") try: tmp_path.write_text( json.dumps(self._state, ensure_ascii=False, indent=2), encoding="utf-8", ) tmp_path.replace(self._state_path) pdebug("CatalogSync: Status gespeichert") except OSError as e: perror(f"CatalogSync: Fehler beim Speichern der Status-Datei: {e}") # --- Sync-Steuerung --- def needs_sync(self) -> bool: """ Prueft, ob ein neuer Vollsync noetig ist. Returns: True, wenn der letzte Sync aelter als das Sync-Intervall ist oder noch nie synchronisiert wurde. """ last_sync = self._state.get("last_sync") if not last_sync: return True try: last_dt = datetime.fromisoformat(last_sync) now = datetime.now(timezone.utc) age_days = (now - last_dt).days return age_days >= self._SYNC_INTERVAL_DAYS except (ValueError, TypeError): pwarn("CatalogSync: Ungueltiges Datum in Status-Datei, Sync noetig") return True async def run_full_sync(self) -> dict: """ Fuehrt einen vollstaendigen Katalog-Sync fuer alle aktivierten Dienste durch. Resume-sicher: Verfolgt, welche Dienste bereits synchronisiert wurden. Bei Unterbrechung werden nur die fehlenden Dienste nachgeholt. Returns: Stats-Dict mit service_id -> Titelanzahl. """ services = self._config.get("services", {}) if not services: pwarn("CatalogSync: Keine Dienste konfiguriert") return {} # Bereits synchronisierte Dienste aus dem aktuellen Lauf synced = set(self._state.get("services_synced", [])) stats: dict[str, int] = {} pinfo(f"CatalogSync: Starte Vollsync fuer {len(services)} Dienste") for service_id, service_config in services.items(): # Deaktivierte Dienste ueberspringen if not service_config.get("enabled", True): pdebug(f"CatalogSync: '{service_id}' deaktiviert, ueberspringe") continue # Bereits synchronisierte Dienste ueberspringen (Resume) if service_id in synced: pdebug(f"CatalogSync: '{service_id}' bereits synchronisiert, ueberspringe") continue try: count = await self._sync_service(service_id, service_config) stats[service_id] = count # Status nach jedem Dienst speichern (Resume-Punkt) synced.add(service_id) self._state["services_synced"] = list(synced) self.save_state() # Cache nach jedem Dienst speichern await self._cache.save_async() except Exception as exc: perror(f"CatalogSync: Fehler beim Sync von '{service_id}': {exc}") stats[service_id] = -1 # Vollsync abgeschlossen — Status zuruecksetzen self._state["last_sync"] = datetime.now(timezone.utc).isoformat(timespec="seconds") self._state["services_synced"] = [] self._state["last_full_sync_stats"] = stats self.save_state() total = sum(c for c in stats.values() if c > 0) pinfo(f"CatalogSync: Vollsync abgeschlossen — {total} Titel in {len(stats)} Diensten") return stats async def run_daily_update(self) -> dict: """ Fuehrt eine schnelle taegliche Aktualisierung durch. Prueft nur die erste Seite pro Dienst (kuerzlich hinzugefuegte Titel). Speichert die neuen Titel als "recent" im Cache. Returns: Stats-Dict mit service_id -> Anzahl neuer Titel. """ services = self._config.get("services", {}) if not services: return {} stats: dict[str, int] = {} days = self._config.get("recent_days", 7) pdebug(f"CatalogSync: Starte taegliches Update (letzte {days} Tage)") for service_id, service_config in services.items(): if not service_config.get("enabled", True): continue try: titles = await self._discover.discover_recent(service_id, days=days) recent_slugs: list[str] = [] for title_data in titles: slug = StreamingCache._generate_slug(title_data.get("t", "")) if slug: self._cache.update_title(slug, title_data) recent_slugs.append(slug) self._cache.set_recent(service_id, recent_slugs) stats[service_id] = len(recent_slugs) except Exception as exc: perror(f"CatalogSync: Fehler beim Daily-Update von '{service_id}': {exc}") stats[service_id] = -1 # Cache speichern await self._cache.save_async() # Status aktualisieren self._state["last_daily_update"] = datetime.now(timezone.utc).isoformat(timespec="seconds") self._state["last_daily_stats"] = stats self.save_state() total = sum(c for c in stats.values() if c > 0) pdebug(f"CatalogSync: Taegliches Update abgeschlossen — {total} neue Titel") return stats async def _sync_service(self, service_id: str, service_config: dict) -> int: """ Synchronisiert einen einzelnen Streaming-Dienst. Ruft alle Titel ueber den Discover-Client ab, aktualisiert den Cache und entfernt Titel, die nicht mehr verfuegbar sind (Pruning). Args: service_id: ID des Streaming-Dienstes. service_config: Konfiguration des Dienstes. Returns: Anzahl der synchronisierten Titel. """ max_titles = service_config.get("max_titles", 500) pinfo(f"CatalogSync: Synchronisiere '{service_id}' (max {max_titles} Titel)") titles = await self._discover.discover_service(service_id, max_titles=max_titles) if not titles: pwarn(f"CatalogSync: Keine Titel fuer '{service_id}' gefunden") return 0 # Titel in den Cache eintragen current_slugs: set[str] = set() for title_data in titles: slug = StreamingCache._generate_slug(title_data.get("t", "")) if slug: self._cache.update_title(slug, title_data) current_slugs.add(slug) # Pruning: Titel entfernen, die nicht mehr auf dem Dienst sind pruned = self._cache.prune_missing(service_id, current_slugs) if pruned > 0: pdebug(f"CatalogSync: {pruned} veraltete Eintraege fuer '{service_id}' entfernt") # Dienst-Info aktualisieren self._cache.update_services({ service_id: { "name": service_config.get("name", service_id), "last_sync": datetime.now(timezone.utc).isoformat(timespec="seconds"), "title_count": len(current_slugs), }, }) # Indizes neu aufbauen nach allen Aenderungen self._cache._rebuild_indices() pinfo(f"CatalogSync: '{service_id}' synchronisiert — {len(current_slugs)} Titel") return len(current_slugs)