Plugin neuheiten: optionale LLM-Hilfsstufe für den Web-Quellen-Crawler

Neues Modul quellen/ki_hilfe.py (konsistent zum dedup-LLM-Muster):
- Struktur-Erkennung: bei verdächtig leerem Parser-Ergebnis schlägt das
  LLM Pagination-/Filter-Folge-URLs vor; hart gefiltert auf gleiche
  Domain, max. 20 je Seite, nicht erreichbare Vorschläge brechen den
  Lauf nicht ab.
- Feld-Extraktion: nur unsichere Treffer (fehlender Verlag, verdächtiger
  Titel); korrigiert ausschließlich titel/verlag/autor, niemals die URL.
- Env-Konfiguration SPIELE_NEUHEITEN_KI_* (AKTIV default 0,
  KONFIDENZ_MIN default 0.7), OpenAI-kompatible Chat-Completions via
  httpx mit Timeout und genau einem Retry.
- Fallback-Pflicht: ohne Konfiguration oder bei jedem Fehler läuft exakt
  der klassische Crawler; KI-Fehler blockieren den Sync nie.
- Audit: KI-Eingriffe je Quelle und Lauf ins audit-log (Modell,
  Konfidenz, korrigierte Felder, verfolgte Folge-URLs).

20 neue Tests (gemocktes HTTP/LLM): Fallback-Fälle, Konfidenz-Schwelle,
Domain-Filter, Max-20-Grenze, URL-Unveränderbarkeit, Audit.
This commit is contained in:
Flo Hartmann
2026-08-23 01:07:18 +00:00
parent c6f9b19f6b
commit f27940db8e
7 changed files with 1321 additions and 7 deletions

View File

@@ -340,12 +340,16 @@ class NeuheitenPlugin(BasePlugin):
def _quellen_sync_ausfuehren(self) -> list[QuellenLaufZeile]:
"""Ein Lauf über alle aktiven Web-Quellen (Fehler pro Quelle isoliert)."""
dedup = (
self.context.registry.get("dedup")
registry = (
self.context.registry
if self.context is not None and self.context.registry is not None
else None
)
service = WebQuellenSyncService(self.context.session_factory, dedup=dedup)
dedup = registry.get("dedup") if registry is not None else None
audit = registry.get("audit-log") if registry is not None else None
service = WebQuellenSyncService(
self.context.session_factory, dedup=dedup, audit=audit
)
return service.laufe()
def _quellen_sync_job(self) -> None:

View File

@@ -21,6 +21,7 @@ die Testsuite führt keine echten Netzwerkaufrufe durch.
"""
from __future__ import annotations
import logging
import re
import time
import xml.etree.ElementTree # noqa: F401 (dokumentiert: keine XML-Nutzung hier)
@@ -31,6 +32,10 @@ from typing import ClassVar
import httpx
from rapidfuzz import fuzz
from .ki_hilfe import KiHilfe
_logger = logging.getLogger(__name__)
#: Freundlicher User-Agent für alle Redaktions-Crawler (Courtesy-Regeln).
USER_AGENT = "SpieleRedaktionBot/1.0 (Redaktions-Tool; Kontakt siehe Repo)"
@@ -76,6 +81,12 @@ class SammelErgebnis:
#: True = mindestens eine Seite kam per 304 als unverändert zurück;
#: der Lauf wurde dann abgekürzt (Quelle übersprungen, kein Fehler).
unveraendert: bool = False
#: KI-Feldkorrekturen dieses Laufs (Protokoll fürs Audit-Log); leer bei
#: rein regelbasiertem Lauf.
ki_eingriffe: list[dict] = field(default_factory=list)
#: Von der KI vorgeschlagene und tatsächlich verfolgte Folge-URLs
#: (Struktur-Erkennung); leer bei rein regelbasiertem Lauf.
ki_folge_urls: list[str] = field(default_factory=list)
def jahr_aus_datumsangabe(angabe: str | None) -> int | None:
@@ -236,18 +247,40 @@ class QuellenAdapter:
def __init__(self, client: WebQuellenClient) -> None:
self.client = client
#: Optionale LLM-Hilfsstufe (`ki_hilfe.KiHilfe` oder gleiches Protokoll).
#: None (Standard) = rein regelbasiert; KI-Fehler blockieren den Lauf nie.
self.ki_hilfe: KiHilfe | None = None
def sammle(self) -> SammelErgebnis:
"""Crawlt die Quelle: Startseiten + gefundene Folge-Links (BFS)."""
"""Crawlt die Quelle: Startseiten + gefundene Folge-Links (BFS).
Mit gesetzter `ki_hilfe`: Bei verdächtig leer geparsten Seiten darf
das LLM Pagination-/Filter-Folge-URLs vorschlagen (nur gleiche
Domain, max. 20), unsichere Treffer werden per LLM nachgebessert
(nur titel/verlag/autor, nie die URL). Jeder KI-Fehler fällt auf den
klassischen Befund zurück.
"""
ergebnis = SammelErgebnis()
warteschlange = list(self.start_urls)
besucht: set[str] = set()
ki_vorschlaege: set[str] = set()
while warteschlange and len(besucht) < self.max_seiten:
url = warteschlange.pop(0)
if url in besucht:
continue
besucht.add(url)
antwort = self.client.hole(url)
try:
antwort = self.client.hole(url)
except Exception:
if url in ki_vorschlaege:
# Eine nicht erreichbare KI-Suggestion darf den klassischen
# Lauf nie abbrechen — überspringen statt Fehler.
_logger.warning(
"KI-vorgeschlagene URL nicht erreichbar (%s) — übersprungen.",
url,
)
continue
raise
if antwort.status_code == 304:
# Unverändert → Quelle überspringen (kein Fehler).
ergebnis.unveraendert = True
@@ -258,8 +291,48 @@ class QuellenAdapter:
for link in folge_links:
if link not in besucht:
warteschlange.append(link)
if self.ki_hilfe is not None:
for link in self._ki_folge_urls(treffer, antwort.content, url, besucht):
if link not in besucht and link not in warteschlange:
warteschlange.append(link)
ki_vorschlaege.add(link)
ergebnis.ki_folge_urls.append(link)
if self.ki_hilfe is not None:
self._ki_nachbessern(ergebnis)
return ergebnis
def _ki_folge_urls(
self,
treffer: list[QuellenTreffer],
inhalt: bytes,
url: str,
besucht: set[str],
) -> list[str]:
"""Struktur-Erkennung der Hilfsstufe (best effort — wirft nicht)."""
try:
return self.ki_hilfe.folge_urls_bei_bedarf(
len(treffer), inhalt, url, besucht
)
except Exception:
_logger.exception(
"KI-Struktur-Erkennung fehlgeschlagen (%s) — klassischer Lauf bleibt.",
url,
)
return []
def _ki_nachbessern(self, ergebnis: SammelErgebnis) -> None:
"""Feld-Korrektur unsicherer Treffer (best effort — wirft nicht)."""
try:
korrigiert, eingriffe = self.ki_hilfe.treffer_nachbessern(ergebnis.treffer)
except Exception:
_logger.exception(
"KI-Feldkorrektur fehlgeschlagen — klassischer Parserbefund bleibt."
)
return
if korrigiert is not None:
ergebnis.treffer = korrigiert
ergebnis.ki_eingriffe = eingriffe or []
def seite_verarbeiten(
self, inhalt: bytes, url: str
) -> tuple[list[QuellenTreffer], list[str]]:

View File

@@ -0,0 +1,575 @@
"""Optionale LLM-Hilfsstufe für den Web-Quellen-Crawler des Neuheiten-Plugins.
Der Crawler (`basis.py` + Adapter) bleibt maßgeblich und arbeitet rein
regelbasiert. Das LLM wird ausschließlich in zwei Situationen eingeschaltet:
1. **Struktur-Erkennung:** Liefert der Regel-Parser keine oder verdächtig
wenige Einträge (0 Treffer, oder wenige Treffer auf einer sehr link-
reichen Seite), darf das LLM Pagination-/Filter-Folge-URLs aus dem
gelieferten HTML vorschlagen. Der Vorschlag wird im Code hart gefiltert:
nur URLs derselben Domain wie die geparste Seite, maximal 20 je Seite —
Schutz vor Rate-Limit-Überlastung und Domain-Ausbruch.
2. **Feld-Extraktion:** Ein einzelner Treffer ist *unsicher* (fehlender
Verlag oder verdächtiger Titel). Dann darf das LLM höchstens `titel`,
`verlag` und `autor` korrigieren — niemals die URL.
Konfiguration konsistent zur dedup-KI über Umgebungsvariablen:
| Variable | Bedeutung |
|----------|-----------|
| `SPIELE_NEUHEITEN_KI_AKTIV` | `1` = Hilfsstufe an (Standard `0`) |
| `SPIELE_NEUHEITEN_KI_BASIS_URL` | Basis-URL, z. B. `https://api.openai.com/v1` |
| `SPIELE_NEUHEITEN_KI_API_KEY` | API-Key (wird als Bearer gesendet) |
| `SPIELE_NEUHEITEN_KI_MODELL` | Modellname, z. B. `gpt-4o-mini` |
| `SPIELE_NEUHEITEN_KI_KONFIDENZ_MIN` | Mindest-Konfidenz (Standard `0.7`) |
| `SPIELE_NEUHEITEN_KI_TIMEOUT_SEKUNDEN` | Timeout je Aufruf (Standard `20`) |
**Fallback-Pflicht:** Ohne Konfiguration oder bei jedem Fehler läuft EXAKT
der klassische Crawler weiter — ein KI-Fehler blockiert den Sync nie.
Ergebnisse unterhalb der Mindest-Konfidenz werden verworfen; der klassische
Parser-Befund bleibt dann unverändert bestehen.
Schnittstelle: OpenAI-kompatible Chat-Completions-API via httpx;
Netzwerk-/Server-Fehler werden genau einmal wiederholt (insgesamt zwei
Versuche), Antwort-Parsing-Fehler nicht.
"""
from __future__ import annotations
import dataclasses
import json
import logging
import os
import re
from dataclasses import dataclass
from typing import Any, Mapping
from urllib.parse import urljoin, urlsplit
import httpx
_logger = logging.getLogger("plugins.neuheiten")
#: Standard-Timeout je Chat-Completions-Aufruf in Sekunden.
STANDARD_TIMEOUT = 20.0
#: Gesamtzahl der Versuche (erster Versuch + genau ein Retry).
VERSUCHE = 2
#: HTTP-Statuscodes, die einen Retry rechtfertigen (zeitweilige Störungen).
RETRY_STATUS = frozenset({408, 429, 500, 502, 503, 504})
#: Mindest-Konfidenz, ab der ein KI-Befund übernommen wird.
STANDARD_KONFIDENZ_MIN = 0.7
#: Maximalzahl der KI-vorgeschlagenen Folge-URLs je Seite (Rate-Limit-Schutz).
MAX_FOLGE_URLS = 20
#: Obergrenze des HTML-Ausschnitts, der pro Struktur-Anfrage ans LLM geht.
HTML_MAX_ZEICHEN = 12000
# ---------------- Verdachts-Heuristiken (wann hilft das LLM?) ----------------
#: Weniger Treffer als das gilt als „verdächtig wenig“ …
STRUKTUR_MIN_TREFFER = 3
#: … aber nur, wenn die Seite zugleich sehr linkreich ist (Indiz für eine
#: nicht erkannte Liste/Pagination statt einer wirklich leeren Seite).
STRUKTUR_LINK_SCHWELLE = 30
_LINK_MUSTER = re.compile(r"<a\s", re.IGNORECASE)
#: Generische Navigations-/Platzhaltertexte, die als Spieltitel verdächtig sind.
GENERISCHE_TITEL = frozenset(
{
"hier", "mehr", "mehr infos", "details", "weiterlesen", "weiter",
"klick", "klick hier", "link", "produkt", "neuheit", "unbekannt",
"titel", "n/a", "-", "",
}
)
def _als_text(inhalt: bytes | str) -> str:
if isinstance(inhalt, bytes):
return inhalt.decode("utf-8", errors="replace")
return inhalt
def titel_ist_verdaechtig(titel: str | None) -> bool:
"""True bei fehlendem, zu kurzem oder generischem Navigationstext."""
if not titel or not titel.strip():
return True
gekuerzt = titel.strip()
if len(gekuerzt) < 3:
return True
return gekuerzt.rstrip(".!? ").lower() in GENERISCHE_TITEL
def treffer_ist_unsicher(treffer: Any) -> bool:
"""Ein Treffer ist unsicher, wenn der Verlag fehlt oder der Titel verdächtig ist.
Nur solche Treffer gehen in die KI-Feld-Extraktion; sichere Treffer
kosten kein Token und bleiben garantiert unangetastet.
"""
verlag = getattr(treffer, "verlag", None)
if not verlag or not str(verlag).strip():
return True
return titel_ist_verdaechtig(getattr(treffer, "titel", None))
def struktur_verdaechtig(anzahl_treffer: int, inhalt: bytes | str) -> bool:
"""True, wenn der Parser keine oder verdächtig wenige Einträge lieferte.
0 Treffer sind immer verdächtig. 12 Treffer nur dann, wenn die Seite
sehr viele Links enthält — sonst ist sie schlicht klein/leer.
"""
if anzahl_treffer <= 0:
return True
if anzahl_treffer >= STRUKTUR_MIN_TREFFER:
return False
links = len(_LINK_MUSTER.findall(_als_text(inhalt)))
return links > STRUKTUR_LINK_SCHWELLE
# ---------------- Konfiguration ----------------
@dataclass(frozen=True)
class KiKonfiguration:
"""Effektive Konfiguration der KI-Hilfsstufe (aus Env gelesen)."""
aktiv: bool = False
basis_url: str = ""
api_key: str = ""
modell: str = ""
konfidenz_min: float = STANDARD_KONFIDENZ_MIN
timeout: float = STANDARD_TIMEOUT
@property
def vollstaendig(self) -> bool:
"""True, wenn Aktiv-Schalter und Zugangsdaten vollständig sind."""
return bool(self.aktiv and self.basis_url and self.api_key and self.modell)
def lade_konfiguration(quelle: Mapping[str, str] | None = None) -> KiKonfiguration:
"""Liest die KI-Konfiguration aus der Umgebung (fehlertolerant).
Kaputte Zahlenwerte führen nicht zum Fehler, sondern zum Standardwert —
die Hilfsstufe darf den Crawler niemals blockieren.
"""
env = os.environ if quelle is None else quelle
aktiv = env.get("SPIELE_NEUHEITEN_KI_AKTIV", "0").strip() == "1"
try:
konfidenz_min = float(env.get("SPIELE_NEUHEITEN_KI_KONFIDENZ_MIN", ""))
except ValueError:
konfidenz_min = STANDARD_KONFIDENZ_MIN
konfidenz_min = min(1.0, max(0.0, konfidenz_min))
try:
timeout = float(env.get("SPIELE_NEUHEITEN_KI_TIMEOUT_SEKUNDEN", ""))
except ValueError:
timeout = STANDARD_TIMEOUT
timeout = max(1.0, timeout)
return KiKonfiguration(
aktiv=aktiv,
basis_url=env.get("SPIELE_NEUHEITEN_KI_BASIS_URL", "").strip().rstrip("/"),
api_key=env.get("SPIELE_NEUHEITEN_KI_API_KEY", "").strip(),
modell=env.get("SPIELE_NEUHEITEN_KI_MODELL", "").strip(),
konfidenz_min=konfidenz_min,
timeout=timeout,
)
# ---------------- Antwort-Parsing (robust gegen Code-Fences u. Ä.) ----------------
_FENCE_MUSTER = re.compile(r"```(?:json|JSON)?\s*(.*?)\s*```", re.DOTALL)
def extrahiere_json(text: str) -> dict | None:
"""Extrahiert das erste JSON-Objekt aus einer LLM-Antwort.
Toleriert Code-Fences und begleitenden Text; None, wenn nichts
Sinnvolles übrig bleibt.
"""
if not isinstance(text, str):
return None
text = text.strip()
if not text:
return None
versuche = [text]
versuche.extend(m.group(1).strip() for m in _FENCE_MUSTER.finditer(text))
erstes, letztes = text.find("{"), text.rfind("}")
if 0 <= erstes < letztes:
versuche.append(text[erstes : letztes + 1])
for kandidat in versuche:
try:
daten = json.loads(kandidat)
except (ValueError, TypeError):
continue
if isinstance(daten, dict):
return daten
return None
def _normalisiere_konfidenz(daten: Mapping[str, Any]) -> float | None:
"""Konfidenz aus der Antwort; Prozent-Skalen werden toleriert."""
try:
konfidenz = float(daten.get("konfidenz"))
except (TypeError, ValueError):
return None
if 10.0 <= konfidenz <= 100.0: # z. B. 92 → 0.92
konfidenz = konfidenz / 100.0
if not 0.0 <= konfidenz <= 1.0: # 1 < x < 10 ist auf keiner Skala plausibel
return None
return round(konfidenz, 4)
def _optional_text(wert: Any) -> str | None:
"""Normalisiert ein nullable String-Feld der Antwort."""
if isinstance(wert, str):
wert = wert.strip()
return wert or None
return None
def validiere_url_vorschlaege(daten: Any, seiten_url: str) -> dict | None:
"""Prüft die Struktur-Antwort und filtert hart auf gleiche Domain + max 20.
Erwartet `{folge_urls: [...], konfidenz: float}`. Relative URLs werden
gegen die Seiten-URL aufgelöst; fremde Domains, Nicht-HTTP-Schemata und
Duplikate (inkl. Fragment-Unterschiede) werden verworfen — auch wenn das
LLM sie vorschlägt.
"""
if not isinstance(daten, dict):
return None
roh = daten.get("folge_urls")
if not isinstance(roh, list):
return None
konfidenz = _normalisiere_konfidenz(daten)
if konfidenz is None:
return None
basis_domain = urlsplit(seiten_url).netloc.lower()
if not basis_domain:
return None
gefiltert: list[str] = []
gesehen: set[str] = set()
for eintrag in roh:
if not isinstance(eintrag, str) or not eintrag.strip():
continue
absolut = urljoin(seiten_url, eintrag.strip())
teile = urlsplit(absolut)
if teile.scheme not in ("http", "https"):
continue
if teile.netloc.lower() != basis_domain:
continue
ohne_fragment = absolut.split("#", 1)[0]
if ohne_fragment in gesehen:
continue
gesehen.add(ohne_fragment)
gefiltert.append(ohne_fragment)
if len(gefiltert) >= MAX_FOLGE_URLS:
break
return {"folge_urls": gefiltert, "konfidenz": konfidenz}
def validiere_feld_antwort(daten: Any) -> dict | None:
"""Prüft die Feld-Antwort gegen das Schema (titel/verlag/autor/konfidenz).
Die URL ist bewusst kein Teil des Schemas — eine Korrektur kann sie
strukturell nie ändern.
"""
if not isinstance(daten, dict):
return None
titel = daten.get("titel")
if not isinstance(titel, str) or not titel.strip():
return None
konfidenz = _normalisiere_konfidenz(daten)
if konfidenz is None:
return None
begruendung = daten.get("begruendung")
if not isinstance(begruendung, str):
begruendung = ""
return {
"titel": titel.strip(),
"verlag": _optional_text(daten.get("verlag")),
"autor": _optional_text(daten.get("autor")),
"konfidenz": konfidenz,
"begruendung": begruendung.strip(),
}
# ---------------- Prompt ----------------
_SYSTEM_STRUKTUR = (
"Du hilfst einem Crawler einer Spielemagazin-Redaktion, Brettspiel-"
"Neuheitenlisten im Web zu finden. Der Regel-Parser hat die folgende "
"Seite fast leer geparst — vermutlich nutzt die Liste Pagination, "
"Archiv- oder Filter-Links. Identifiziere aus dem HTML die URLs, unter "
"denen weitere Neuheiten-Einträge zu erwarten sind. Antworte "
"AUSSCHLIESSLICH mit einem JSON-Objekt nach exakt diesem Schema:\n"
'{"folge_urls": ["<absolute oder relative URL>", …], '
'"konfidenz": <float 0.0-1.0>}\n'
f"Höchstens {MAX_FOLGE_URLS} URLs und nur Links derselben Domain wie die "
"Seiten-URL. Wenn du nichts Sinnvolles findest, liefere eine leere Liste "
"mit niedriger Konfidenz. Kein weiterer Text, keine Code-Fences."
)
_SYSTEM_FELDER = (
"Du assistierst einer Spielemagazin-Redaktion beim Parsen von Brettspiel-"
"Neuheiten. Der folgende Treffer eines regelbasierten Crawlers ist "
"unsicher. Bereinige die Felder: trenne „Verlag: Titel“-Muster, entferne "
"Navigationstext und ergänze den Verlag, falls er sicher erkennbar ist. "
"Korrigiere ausschließlich titel, verlag und autor — ändere niemals URLs "
"und erfinde nichts. Ein Feld, das du nicht sicher bestimmen kannst, "
"bleibt null. Antworte AUSSCHLIESSLICH mit einem JSON-Objekt nach exakt "
"diesem Schema:\n"
'{"titel": "<bereinigter Titel>", "verlag": "<Verlagsname oder null>", '
'"autor": "<Autor oder null>", "konfidenz": <float 0.0-1.0>, '
'"begruendung": "<kurzer deutscher Satz>"}\n'
"Kein weiterer Text, keine Code-Fences."
)
def baue_struktur_prompt(seiten_url: str, html_text: str) -> str:
"""User-Nachricht für die Struktur-Erkennung (gekürzter HTML-Ausschnitt)."""
ausschnitt = html_text[:HTML_MAX_ZEICHEN]
return (
f"Seiten-URL: {seiten_url}\n\n"
"HTML-Ausschnitt der geparsten Seite:\n"
f"{ausschnitt}"
)
def baue_feld_prompt(titel: str, verlag: str | None, autor: str | None) -> str:
"""User-Nachricht mit dem unsicheren Treffer (ohne URL — bleibt fix)."""
zeilen = [f"Unsicherer Treffer:", f"titel: {titel}"]
zeilen.append(f"verlag: {verlag if verlag else '(leer)'}")
zeilen.append(f"autor: {autor if autor else '(leer)'}")
zeilen.append(
"Gib die bereinigten Felder zurück; die quellen_url selbst bleibt "
"nicht veränderbar und ist nicht Teil der Aufgabe."
)
return "\n".join(zeilen)
# ---------------- Client ----------------
class KiHilfe:
"""OpenAI-kompatibler Chat-Completions-Client für die Crawler-Hilfsstufe.
Bewusst genauso defensiv wie der dedup-KI-Client: Alle öffentlichen
Methoden werfen nicht (Ausnahmen werden intern geschluckt) — ohne
Konfiguration oder bei jedem Problem greift exakt der klassische
Crawler-Befund. Transport injizierbar, damit Tests netzwerkfrei bleiben.
"""
def __init__(
self,
konfiguration: KiKonfiguration,
*,
transport: httpx.BaseTransport | None = None,
) -> None:
self._konfiguration = konfiguration
self._client = httpx.Client(
base_url=konfiguration.basis_url,
timeout=konfiguration.timeout,
transport=transport,
headers={
"Authorization": f"Bearer {konfiguration.api_key}",
"Content-Type": "application/json",
},
)
@property
def modell(self) -> str:
return self._konfiguration.modell
def schliessen(self) -> None:
self._client.close()
# ---------- Struktur-Erkennung ----------
def folge_urls_bei_bedarf(
self,
anzahl_treffer: int,
inhalt: bytes | str,
seiten_url: str,
bereits_geprueft: set[str] | None = None,
) -> list[str]:
"""Pagination-/Filter-URLs bei verdächtig leerem Parser-Ergebnis.
Prüft selbst die Verdachts-Heuristik (sonst []), fragt dann das LLM
und filtert dessen Vorschläge hart: gleiche Domain, max. 20, bereits
geplante/besuchte URLs raus, Konfidenz unter der Schwelle → alles
verwerfen. Fehler → [] (klassischer Lauf bleibt).
"""
try:
if not struktur_verdaechtig(anzahl_treffer, inhalt):
return []
if not self._konfiguration.vollstaendig:
return []
nutzernachricht = baue_struktur_prompt(seiten_url, _als_text(inhalt))
daten = self._chat(_SYSTEM_STRUKTUR, nutzernachricht)
ergebnis = validiere_url_vorschlaege(daten, seiten_url)
if ergebnis is None:
return []
if ergebnis["konfidenz"] < self._konfiguration.konfidenz_min:
_logger.info(
"neuheiten/KI: URL-Vorschläge verworfen (Konfidenz %.2f < %.2f).",
ergebnis["konfidenz"], self._konfiguration.konfidenz_min,
)
return []
bekannt = bereits_geprueft or set()
return [
url for url in ergebnis["folge_urls"] if url not in bekannt
]
except Exception:
_logger.exception(
"neuheiten/KI: Struktur-Erkennung fehlgeschlagen — klassischer Lauf bleibt."
)
return []
# ---------- Feld-Extraktion ----------
def treffer_nachbessern(
self, treffer_liste: list
) -> tuple[list, list[dict]]:
"""Korrigiert ausschließlich unsichere Treffer (titel/verlag/autor).
Sichere Treffer werden nicht angefasst und kosten kein Token. Die
quellen_url bleibt immer unangetastet (kein Teil des Antwort-Schemas).
Rückgabe: (neue Trefferliste, Eingriffs-Protokoll fürs Audit-Log mit
Modell, Konfidenz und den jeweils korrigierten Feldern).
"""
if not treffer_liste:
return [], []
if not self._konfiguration.vollstaendig:
return list(treffer_liste), []
neue_liste: list = []
eingriffe: list[dict] = []
for einzel in treffer_liste:
if not treffer_ist_unsicher(einzel):
neue_liste.append(einzel)
continue
ergebnis = self._einzeln_korrigieren(einzel)
if ergebnis is None:
neue_liste.append(einzel) # klassischer Befund bleibt
continue
korrigiert, antwort = ergebnis
if korrigiert is einzel: # nichts tatsächlich geändert
neue_liste.append(einzel)
continue
neue_liste.append(korrigiert)
eingriffe.append(self._eingriff_protokoll(einzel, korrigiert, antwort))
return neue_liste, eingriffe
def _einzeln_korrigieren(self, treffer: Any) -> tuple[Any, dict] | None:
"""Eine Feld-Korrektur; Rückgabe (neuer Treffer, validierte Antwort).
None bzw. der unveränderte Treffer bedeutet: klassischer Befund bleibt.
"""
try:
nutzernachricht = baue_feld_prompt(
getattr(treffer, "titel", "") or "",
getattr(treffer, "verlag", None),
getattr(treffer, "autor", None),
)
daten = self._chat(_SYSTEM_FELDER, nutzernachricht)
antwort = validiere_feld_antwort(daten)
if antwort is None:
return None
if antwort["konfidenz"] < self._konfiguration.konfidenz_min:
_logger.info(
"neuheiten/KI: Feldkorrektur verworfen (Konfidenz %.2f < %.2f).",
antwort["konfidenz"], self._konfiguration.konfidenz_min,
)
return None
felder: dict[str, str | None] = {}
for feld in ("titel", "verlag", "autor"):
alter_wert = getattr(treffer, feld, None)
alter_norm = (alter_wert or "").strip() or None
neu_norm = (antwort[feld] or "").strip() or None
if neu_norm != alter_norm:
felder[feld] = neu_norm
if not felder:
return treffer, antwort # nichts tatsächlich korrigiert
try:
return dataclasses.replace(treffer, **felder), antwort
except Exception:
return None
except Exception:
_logger.exception(
"neuheiten/KI: Feldkorrektur fehlgeschlagen — klassischer Befund bleibt."
)
return None
def _eingriff_protokoll(self, alt: Any, neu: Any, antwort: dict) -> dict:
"""Audit-Zeile: was wurde korrigiert, mit welchem Modell/welcher Konfidenz."""
korrigiert: dict[str, str | None] = {}
for feld in ("titel", "verlag", "autor"):
vorher = getattr(alt, feld, None)
nachher = getattr(neu, feld, None)
if vorher != nachher:
korrigiert[feld] = nachher
return {
"quellen_url": getattr(alt, "quellen_url", ""),
"modell": self.modell,
"konfidenz": antwort["konfidenz"],
"begruendung": antwort.get("begruendung", ""),
"korrigiert": korrigiert,
"vorher": {f: getattr(alt, f, None) for f in korrigiert},
}
# ---------- Gemeinsamer Chat-Completions-Aufruf ----------
def _chat(self, system: str, nutzer: str) -> dict | None:
"""POST /chat/completions mit genau einem Retry; None bei jedem Problem.
Netzwerk-/Timeout-Fehler und zeitweilige Server-Störungen werden
genau einmal wiederholt (insgesamt zwei Versuche); andere 4xx und
kaputte Antworten nicht — dort hilft kein Retry.
"""
payload = {
"model": self._konfiguration.modell,
"temperature": 0,
"messages": [
{"role": "system", "content": system},
{"role": "user", "content": nutzer},
],
}
for versuch in range(1, VERSUCHE + 1):
try:
antwort = self._client.post("/chat/completions", json=payload)
except httpx.RequestError as exc:
_logger.warning(
"neuheiten/KI: Versuch %d/%d fehlgeschlagen (%s)",
versuch, VERSUCHE, exc,
)
continue
if antwort.status_code in RETRY_STATUS:
_logger.warning(
"neuheiten/KI: Versuch %d/%d mit Status %d — wiederholt.",
versuch, VERSUCHE, antwort.status_code,
)
continue
try:
antwort.raise_for_status() # andere 4xx: kein Retry hilft
except httpx.HTTPStatusError as exc:
_logger.warning("neuheiten/KI: Anfrage abgelehnt (%s).", exc)
return None
try:
inhalte = antwort.json()["choices"][0]["message"]["content"]
except Exception as exc: # kaputtes Antwort-Layout
_logger.warning("neuheiten/KI: unbrauchbare Antwort (%s)", exc)
return None
return extrahiere_json(inhalte)
_logger.warning(
"neuheiten/KI: alle %d Versuche fehlgeschlagen — klassischer Befund.", VERSUCHE
)
return None

View File

@@ -32,6 +32,7 @@ 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,
@@ -108,20 +109,27 @@ class WebQuellenSyncService:
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)
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(
@@ -141,7 +149,28 @@ class WebQuellenSyncService:
# ---------- Eine Quelle ----------
def _eine_quelle(self, klasse: type[QuellenAdapter]) -> QuellenLaufZeile:
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):
@@ -152,6 +181,7 @@ class WebQuellenSyncService:
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:
@@ -160,6 +190,7 @@ class WebQuellenSyncService:
zeile.fehlermeldung = str(exc)[:300]
return zeile
zeile.seiten = ergebnis.seiten
self._ki_auditieren(klasse.name, ergebnis)
if ergebnis.unveraendert:
zeile.unveraendert = True
else:
@@ -263,6 +294,37 @@ class WebQuellenSyncService:
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: