"""Sync-Orchestrierung für die Web-Quellen des Neuheiten-Plugins. Pro Quelle: - Env-Schalter SPIELE_NEUHEITEN_QUELLE__AKTIV (Standard „1“ = an). - Adapter-Lauf mit konditionalen Requests (ETag/Last-Modified werden in der Statuszeile persistiert; 304 → Quelle übersprungen, kein Fehler). - Fehler-Isolation: eine kaputte Quelle bricht den Gesamtdurchlauf nicht ab; ihr Ergebnis wird als Fehler protokolliert, die anderen Quellen laufen normal weiter. - Duplikatvermeidung: Web-Quellen haben keine BGG-ID, daher wird jeder Treffer vor dem Speichern fuzzy gegen die bestehende Liste geprüft (Titel-Match, rapidfuzz) — Treffer gegen BGG-Einträge werden übersprungen. Zusätzlich läuft jeder neue Treffer durch die öffentliche dedup-Prüfung (via Plugin-Registry, sofern das Plugin geladen ist). - Upsert: bereits vorhandene Treffer derselben Quelle (gleicher Titel) werden aktualisiert statt doppelt angelegt; neue Einträge bekommen quelle = '' statt 'boardgamegeek' und keine BGG-ID. Das Ergebnis jedes Laufes wird pro Quelle in `neuheiten_quellen_status` protokolliert und im Plugin-UI unter /neuheiten angezeigt. """ from __future__ import annotations import asyncio import logging import os import time from dataclasses import dataclass from datetime import datetime, timezone from typing import Callable from sqlalchemy import select from sqlalchemy.orm import sessionmaker from .quellen.ki_hilfe import KiHilfe, lade_konfiguration as lade_ki_konfiguration from .models import STATUS_NEUHEIT, Neuheit, QuellenStatus from .quellen import ( ADAPTER_KLASSEN, QuellenAdapter, QuellenTreffer, WebQuellenClient, jahr_aus_datumsangabe, titel_aehnlichkeit, titel_normalisieren, ) _logger = logging.getLogger("plugins.neuheiten") #: Schwelle für den Fuzzy-Titel-Match gegen bestehende Einträge (0–100). TITEL_FUZZY_SCHWELLE = 92.0 #: Zweite Regel: etwas lockererer Titel-Match, dafür gleicher/ähnlicher Verlag. TITEL_FUZZY_MIT_VERLAG_SCHWELLE = 85.0 VERLAG_FUZZY_SCHWELLE = 85.0 ENV_SCHABLONE_AKTIV = "SPIELE_NEUHEITEN_QUELLE_{name}_AKTIV" def quelle_aktiv(name: str, umgebung=None) -> bool: """Liest den Env-Schalter einer Quelle (Standard: an).""" env = umgebung if umgebung is not None else os.environ wert = env.get(ENV_SCHABLONE_AKTIV.format(name=name.upper()), "1") return str(wert).strip() != "0" @dataclass class QuellenLaufZeile: """Ergebnis einer Quelle in einem Sync-Lauf (Grundlage der UI-Übersicht).""" quelle: str anzeigename: str aktiv: bool = True neu: int = 0 aktualisiert: int = 0 gefiltert: int = 0 fehleranzahl: int = 0 fehlermeldung: str | None = None unveraendert: bool = False seiten: int = 0 dauer_sekunden: float = 0.0 url: str | None = None def als_text(self) -> str: if not self.aktiv: return f"{self.anzeigename}: ausgeschaltet" if self.fehlermeldung: return f"{self.anzeigename}: Fehler — {self.fehlermeldung}" if self.unveraendert: return f"{self.anzeigename}: unverändert (304) — übersprungen" return ( f"{self.anzeigename}: {self.neu} neu, {self.aktualisiert} aktualisiert, " f"{self.gefiltert} gefiltert (Duplikate)" ) class WebQuellenSyncService: """Führt einen Sync-Lauf über alle Web-Quellen aus. `dedup` ist die öffentliche API des dedup-Plugins (via Plugin-Registry) oder None, wenn das Plugin nicht geladen ist — die Prüfung ist dann optional und wird übersprungen. `client_fabrik` baut den HTTP-Client (in Tests mit Fake-Transport überschrieben). """ def __init__( self, session_factory: sessionmaker, *, adapter_klassen: tuple[type[QuellenAdapter], ...] = ADAPTER_KLASSEN, dedup=None, client_fabrik: Callable[[dict], WebQuellenClient] | None = None, umgebung=None, ki_fabrik: Callable[[], KiHilfe] | None = None, audit=None, ) -> None: self.session_factory = session_factory self.adapter_klassen = adapter_klassen self.dedup = dedup self.audit = audit self.umgebung = umgebung if umgebung is not None else os.environ self._client_fabrik = client_fabrik or self._standard_client_fabrik # KI-Hilfsstufe: Fabrik für Tests; ohne vollständige Konfiguration wird # nie ein Client gebaut (und damit nie ein LLM-Aufruf getätigt). self._ki_fabrik = ki_fabrik # ---------- Gesamtdurchlauf ---------- def laufe(self) -> list[QuellenLaufZeile]: zeilen: list[QuellenLaufZeile] = [] ki = self._ki_hilfe() for klasse in self.adapter_klassen: try: zeile = self._eine_quelle(klasse, ki) except Exception as exc: # letzte Verteidigungslinie pro Quelle _logger.exception("Quellen-Sync fehlgeschlagen für %s", klasse.name) zeile = QuellenLaufZeile( klasse.name, klasse.anzeigename, fehleranzahl=1, fehlermeldung=str(exc)[:300], ) self._status_speichern(zeile) zeilen.append(zeile) return zeilen def gesamt_text(self, zeilen: list[QuellenLaufZeile]) -> str: """Kompakte deutschsprachige Zusammenfassung für Banner/Protokoll.""" teile = [zeile.als_text() for zeile in zeilen] return " · ".join(teile) if teile else "Keine Web-Quellen konfiguriert." # ---------- Eine Quelle ---------- def _ki_hilfe(self) -> KiHilfe | None: """KI-Hilfsstufe nur bei vollständiger Konfiguration (sonst None). Die Aktivierungsprüfung gilt vor jeder injizierten Fabrik — ist die Hilfsstufe aus oder unvollständig konfiguriert, wird nie ein Client gebaut und damit nie ein LLM-Aufruf getätigt. """ try: konfiguration = lade_ki_konfiguration(self.umgebung) if not konfiguration.vollstaendig: return None if self._ki_fabrik is not None: return self._ki_fabrik() return KiHilfe(konfiguration) except Exception: # Auch ein kaputtes KI-Setup darf den Sync nie blockieren. _logger.exception("neuheiten/KI: Hilfsstufe nicht verfügbar — klassischer Lauf.") return None def _eine_quelle( self, klasse: type[QuellenAdapter], ki: KiHilfe | None = None ) -> QuellenLaufZeile: zeile = QuellenLaufZeile(klasse.name, klasse.anzeigename) zeile.url = klasse.start_urls[0] if klasse.start_urls else None if not quelle_aktiv(klasse.name, self.umgebung): zeile.aktiv = False return zeile start = time.monotonic() client = self._client_fabrik(self._validatoren_laden(klasse.name)) try: adapter = klasse(client) adapter.ki_hilfe = ki # None = rein regelbasiert (Fallback-Pflicht) try: ergebnis = adapter.sammle() except Exception as exc: _logger.warning("Quelle %s fehlgeschlagen: %s", klasse.name, exc) zeile.fehleranzahl = 1 zeile.fehlermeldung = str(exc)[:300] return zeile zeile.seiten = ergebnis.seiten self._ki_auditieren(klasse.name, ergebnis) if ergebnis.unveraendert: zeile.unveraendert = True else: self._treffer_verarbeiten(zeile, ergebnis.treffer, klasse.name) finally: client.schliessen() zeile.dauer_sekunden = round(time.monotonic() - start, 2) self._validatoren_speichern(klasse.name, zeile.url, client) return zeile def _treffer_verarbeiten( self, zeile: QuellenLaufZeile, treffer: list[QuellenTreffer], quelle: str ) -> None: asyncio.run(self._verarbeite_async(zeile, treffer, quelle)) async def _verarbeite_async( self, zeile: QuellenLaufZeile, treffer: list[QuellenTreffer], quelle: str ) -> None: with self.session_factory() as db: try: for einzel in treffer: if not einzel.titel or not einzel.titel.strip(): zeile.gefiltert += 1 continue if self._eigenen_eintrag_aktualisieren(db, einzel, quelle): zeile.aktualisiert += 1 continue if self._ist_duplikat(db, einzel): zeile.gefiltert += 1 continue if not await self._dedup_erklaert_kein_konflikt(einzel): zeile.gefiltert += 1 continue db.add( Neuheit( titel=einzel.titel.strip()[:300], verlag=einzel.verlag, autor=einzel.autor, erscheinungsjahr=jahr_aus_datumsangabe( einzel.erscheinungsdatum_oder_quartal ), bgg_id=None, bild_url=einzel.bild_url, quellen_url=einzel.quellen_url, status=STATUS_NEUHEIT, quelle=quelle, ) ) zeile.neu += 1 db.commit() except Exception: db.rollback() raise # ---------- Duplikate / Upsert ---------- def _eigenen_eintrag_aktualisieren( self, db, treffer: QuellenTreffer, quelle: str ) -> bool: """Update statt Duplikat: gleicher Titel + gleiche Quelle.""" schluessel = titel_normalisieren(treffer.titel) vorhanden = db.scalars(select(Neuheit).where(Neuheit.quelle == quelle)).all() for eintrag in vorhanden: if titel_normalisieren(eintrag.titel) == schluessel: eintrag.verlag = treffer.verlag eintrag.autor = treffer.autor eintrag.erscheinungsjahr = jahr_aus_datumsangabe( treffer.erscheinungsdatum_oder_quartal ) # Coverbild + Quellen-URL mitpflegen; fehlen sie im Treffer, # bleiben die bereits gespeicherten Werte erhalten. if treffer.bild_url: eintrag.bild_url = treffer.bild_url if treffer.quellen_url: eintrag.quellen_url = treffer.quellen_url return True return False def _ist_duplikat(self, db, treffer: QuellenTreffer) -> bool: """Titel-Match gegen bestehende Einträge (z. B. aus BGG). BGG-IDs haben Web-Quellen nicht — gematcht wird fuzzy auf den Titel; bei Grenzfällen entscheidet ein ähnlicher Verlag mit. """ for eintrag in db.scalars(select(Neuheit)).all(): if titel_aehnlichkeit(treffer.titel, eintrag.titel) >= TITEL_FUZZY_SCHWELLE: return True if ( treffer.verlag and eintrag.verlag and titel_aehnlichkeit(treffer.titel, eintrag.titel) >= TITEL_FUZZY_MIT_VERLAG_SCHWELLE and titel_aehnlichkeit(treffer.verlag, eintrag.verlag) >= VERLAG_FUZZY_SCHWELLE ): return True return False async def _dedup_erklaert_kein_konflikt(self, treffer: QuellenTreffer) -> bool: """Öffentliche dedup-Prüfung (sofern Plugin geladen); Fehler tolerieren.""" if self.dedup is None: return True try: ergebnis = await self.dedup.check_titel(treffer.titel, treffer.verlag, None) except Exception as exc: _logger.warning("dedup-Prüfung fehlgeschlagen (%s): %s", treffer.titel, exc) return True if ergebnis is None: return True return not ergebnis.hat_konflikte # ---------- KI-Audit ---------- def _ki_auditieren(self, quelle: str, ergebnis) -> None: """KI-Eingriffe des Laufs ins audit-log schreiben (best effort). Ein Eintrag je Quelle und Lauf mit Modell, Konfidenz, korrigierten Feldern und den verfolgten KI-Folge-URLs. Ohne audit-log (Plugin nicht geladen) oder bei Fehlern wird nur protokolliert — der Sync wird nie blockiert. """ eingriffe = list(getattr(ergebnis, "ki_eingriffe", []) or []) folge_urls = list(getattr(ergebnis, "ki_folge_urls", []) or []) if not eingriffe and not folge_urls: return if self.audit is None: _logger.info( "Quelle %s: %d KI-Korrekturen, %d KI-Folge-URLs (kein audit-log geladen).", quelle, len(eingriffe), len(folge_urls), ) return details = { "quelle": quelle, "modell": eingriffe[0].get("modell", "") if eingriffe else "", "eingriffe": eingriffe, "folge_urls": folge_urls, } try: self.audit.log_sync(None, "korrigiert", "neuheiten_ki_crawler", quelle, details) except Exception: _logger.exception("neuheiten/KI: Audit-Eintrag konnte nicht geschrieben werden.") # ---------- Statusprotokoll ---------- def _status_speichern(self, zeile: QuellenLaufZeile) -> None: with self.session_factory() as db: eintrag = db.scalar( select(QuellenStatus).where(QuellenStatus.quelle == zeile.quelle) ) if eintrag is None: eintrag = QuellenStatus(quelle=zeile.quelle, url=zeile.url or "") db.add(eintrag) if zeile.url: eintrag.url = zeile.url eintrag.letzte_laufzeit = datetime.now(timezone.utc).replace(tzinfo=None) eintrag.dauer_sekunden = zeile.dauer_sekunden eintrag.anzahl_neu = zeile.neu eintrag.anzahl_aktualisiert = zeile.aktualisiert eintrag.anzahl_gefiltert = zeile.gefiltert eintrag.anzahl_fehler = zeile.fehleranzahl eintrag.fehlermeldung = zeile.fehlermeldung eintrag.unveraendert = zeile.unveraendert db.commit() def _validatoren_laden(self, quelle: str) -> dict: with self.session_factory() as db: eintrag = db.scalar( select(QuellenStatus).where(QuellenStatus.quelle == quelle) ) return dict(eintrag.validatoren) if eintrag is not None else {} def _validatoren_speichern( self, quelle: str, url: str | None, client: WebQuellenClient ) -> None: if not client.validatoren: return with self.session_factory() as db: eintrag = db.scalar( select(QuellenStatus).where(QuellenStatus.quelle == quelle) ) if eintrag is None: eintrag = QuellenStatus(quelle=quelle, url=url or "") db.add(eintrag) eintrag.validatoren = client.validatoren db.commit() # ---------- Hilfen ---------- @staticmethod def _standard_client_fabrik(validatoren) -> WebQuellenClient: return WebQuellenClient(validatoren=validatoren)