| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282 |
- # -*- 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)
|