""" KI-gestützte semantische Inhaltsanalyse via OpenRouter. Ergänzung zum Integritäts-Kern, niemals dessen Ersatz. Erkennt problematische *Inhalte ohne Link-Signal*: Pornografie, Propaganda, diffamierende/strafbare Texte, versteckten Spam, thematisch unpassende Werbung, widersprüchliche Aussagen — in Text und in Bildern (inkl. eingebettetem Text via OCR). Kostengate: Jeder Kandidat bekommt einen Fingerprint (SHA-256). Ein bereits geprüfter Fingerprint liegt mit seinem Verdikt im Ledger (data/ai_ledger.json) → Cache-Treffer → KEIN API-Call. Nur neue/geänderte Inhalte kosten etwas. Robustheit: Fehlt der API-Key oder schlägt ein Call fehl, wird die Analyse übersprungen und der Scan läuft unverändert weiter (graceful degradation). KI darf den Kern-Scan nie brechen. """ import concurrent.futures import json import logging import os import threading import time from datetime import datetime, timezone import requests from .checker import fetch_asset_hashes from .differ import normalize_text logger = logging.getLogger(__name__) _OPENROUTER_URL = "https://openrouter.ai/api/v1/chat/completions" # Erlaubte Kategorien, die das Modell zurückgeben darf. _CATEGORIES = [ "clean", "pornography", "propaganda", "defamation_illegal", "hidden_spam", "off_topic_commercial", "contradiction", ] # Kategorie → Scoring-Schlüssel (Fallback-Punkte, falls cfg sie nicht liefert). _CATEGORY_SCORE_KEY = { "pornography": ("ai_pornography", 50), "defamation_illegal": ("ai_defamation_illegal", 50), "propaganda": ("ai_propaganda", 40), "hidden_spam": ("ai_hidden_spam", 40), "off_topic_commercial": ("ai_off_topic_commercial", 30), "contradiction": ("ai_contradiction", 20), } _SEVERITY_ORDER = {"none": 0, "low": 1, "medium": 2, "high": 3} # JSON-Schema für die strukturierte Modell-Antwort. _RESPONSE_SCHEMA = { "type": "json_schema", "json_schema": { "name": "content_verdict", "strict": True, "schema": { "type": "object", "properties": { "category": {"type": "string", "enum": _CATEGORIES}, "severity": {"type": "string", "enum": ["none", "low", "medium", "high"]}, "confidence": {"type": "number"}, "explanation": {"type": "string"}, }, "required": ["category", "severity", "confidence", "explanation"], "additionalProperties": False, }, }, } _MODELS_URL = "https://openrouter.ai/api/v1/models" # --------------------------------------------------------------------------- # Adaptive model router (circuit-breaker + latency telemetry + /models pruning) # --------------------------------------------------------------------------- class ModelRouter: """Wählt programmatisch das beste Modell aus einer Kette: gesunde zuerst, nach Latenz; rate-limited/ausgefallene Modelle bekommen einen Cooldown (Circuit-Breaker) und werden so lange übersprungen. Tote Slugs werden über den /models-Katalog failsafe aussortiert. Thread-sicher (für Parallelität).""" def __init__(self, ai_cfg: dict, state: dict | None = None): self._lock = threading.Lock() self._threshold = ai_cfg.get("breaker_failure_threshold", 2) self._cooldown = ai_cfg.get("breaker_cooldown_seconds", 120) self._refresh = ai_cfg.get("refresh_models", True) state = state or {} self._models: dict = state.get("models", {}) or {} self._catalog: dict = state.get("catalog", {"slugs": [], "fetched_at": None}) def state(self) -> dict: with self._lock: return {"models": self._models, "catalog": self._catalog} def record_success(self, model: str, latency: float) -> None: with self._lock: m = self._models.setdefault(model, {}) m["failures"] = 0 m["open_until"] = 0 prev = m.get("latency") m["latency"] = latency if prev is None else 0.7 * prev + 0.3 * latency def record_failure(self, model: str) -> None: with self._lock: m = self._models.setdefault(model, {}) m["failures"] = m.get("failures", 0) + 1 m["last_failure"] = time.time() if m["failures"] >= self._threshold: m["open_until"] = time.time() + self._cooldown logger.info("KI-Router: Breaker für %s offen (%.0fs Cooldown).", model, self._cooldown) def order(self, models: list[str]) -> list[str]: """Reiht die Kette: geschlossener Breaker zuerst, dann nach Latenz aufsteigend. Tote Slugs werden entfernt. Gibt nie eine leere Liste zurück.""" with self._lock: candidates = self._prune_locked(models) now = time.time() def sort_key(slug): m = self._models.get(slug, {}) is_open = m.get("open_until", 0) > now # Unbekannte Latenz → inf: auf dem Erstlauf bleibt die Config-Reihenfolge # erhalten (free zuerst); ein bewährtes (gemessenes) Modell schlägt ein # noch untestetes. sorted() ist stabil → Ties = Config-Reihenfolge. latency = m.get("latency") latency = latency if latency is not None else float("inf") return (is_open, latency) ordered = sorted(candidates, key=sort_key) return ordered or list(models) def _prune_locked(self, models: list[str]) -> list[str]: slugs = self._catalog.get("slugs") or [] if not slugs: return list(models) # kein Katalog → nicht filtern (failsafe) keep = [m for m in models if m in slugs] return keep or list(models) # nie alles wegfiltern def refresh_catalog(self) -> None: """Einmal pro Lauf /models ziehen (24 h gecacht). Failsafe: blockiert nie.""" if not self._refresh: return fetched_at = self._catalog.get("fetched_at") if fetched_at and self._catalog.get("slugs") and (time.time() - fetched_at) < 86400: return # Cache frisch try: resp = requests.get(_MODELS_URL, timeout=15) if resp.status_code == 200: slugs = [m["id"] for m in resp.json().get("data", []) if m.get("id")] if slugs: with self._lock: self._catalog = {"slugs": slugs, "fetched_at": time.time()} logger.info("KI-Router: /models-Katalog aktualisiert (%d Modelle).", len(slugs)) except (requests.exceptions.RequestException, ValueError, KeyError) as exc: logger.warning("KI-Router: /models nicht abrufbar (failsafe, ungefiltert): %s", exc) # --------------------------------------------------------------------------- # Public API # --------------------------------------------------------------------------- def run_ai_analysis(cfg: dict, bm, snap: dict, diff: dict) -> dict: """ Hash-gegate KI-Analyse über neue/geänderte Texte und Bilder. Returns {"findings": [...], "checked": int, "cache_hits": int, "api_calls": int, "skipped": str|None}. Bricht nie mit einer Exception — bei Problemen wird "skipped" gesetzt. """ result = {"findings": [], "unchecked": [], "checked": 0, "cache_hits": 0, "api_calls": 0, "skipped": None} ai_cfg = cfg.get("ai_analysis", {}) if not ai_cfg.get("enabled"): result["skipped"] = "deaktiviert" return result api_key = os.environ.get(ai_cfg.get("api_key_env", "OPENROUTER_API_KEY"), "") if not api_key: result["skipped"] = f"kein API-Key ({ai_cfg.get('api_key_env', 'OPENROUTER_API_KEY')})" logger.warning("KI-Analyse übersprungen: %s", result["skipped"]) return result router = ModelRouter(ai_cfg, bm.load_ai_router_state()) router.refresh_catalog() # einmal /models (24 h gecacht), failsafe max_workers = max(1, int(ai_cfg.get("max_concurrency", 8))) try: ledger = bm.load_ai_ledger() entries = ledger.setdefault("entries", {}) dirty = False # Ledger wird zwischengesichert (nach Textphase + im finally), damit bereits # berechnete Verdikte einen späteren Hang/Abbruch überleben und nicht erneut # bezahlt werden müssen. with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: try: if ai_cfg.get("text", {}).get("enabled", True): dirty |= _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor) if dirty: bm.save_ai_ledger(ledger) if ai_cfg.get("image", {}).get("enabled", True): dirty |= _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor) # Audio/Video: vorbereitet, default aus (siehe _collect_audio/video_candidates). finally: if dirty: bm.save_ai_ledger(ledger) except Exception as exc: # pragma: no cover - Schutzschirm, darf Scan nie brechen logger.error("KI-Analyse mit unerwartetem Fehler abgebrochen: %s", exc) result["skipped"] = f"interner Fehler: {exc}" # Deterministische Report-Reihenfolge (unabhängig von der Completion-Reihenfolge) result["findings"].sort(key=lambda f: (f.get("kind", ""), f.get("url", ""), f.get("asset_url", ""))) result["unchecked"].sort(key=lambda u: (u.get("kind", ""), u.get("url", ""))) bm.save_ai_router_state(router.state()) return result def score_ai_findings(ai_result: dict, cfg: dict) -> dict: """ Bewertet KI-Funde additiv. Das Level ist auf 'yellow' gedeckelt — AUSSER schwerwiegende Kategorien (red_categories, default Pornografie/strafbar) lösen bei hoher Schwere + hoher Konfidenz ROT aus. Die Punktsumme selbst ergibt nie Rot; nur dieser kategoriebasierte Override tut das. Returns {score, level, reasons, exit_code}. """ sc = cfg.get("scoring", {}) thr = cfg.get("thresholds", {"yellow": 20, "red": 60}) yellow = thr.get("yellow", 20) score = 0 reasons: list[str] = [] for f in ai_result.get("findings", []): score += (pts := _finding_points(f, sc)) conf = f.get("confidence", 0.0) reasons.append( f"KI [{f.get('kind', '?')}] {f.get('url', '?')}: " f"{f.get('category', '?')} ({f.get('severity', '?')}, {conf:.0%}) (+{pts})" ) # Deckelung auf gelb: auch bei Score ≥ rot-Schwelle bleibt es gelb. if score >= yellow: level, exit_code = "yellow", 1 else: level, exit_code = "green", 0 # Ungeprüfte Inhalte: "konnte nicht prüfen" ≠ "sauber". unchecked = ai_result.get("unchecked", []) if unchecked: reasons.append( f"{len(unchecked)} Inhalt(e) konnten von der KI nicht geprüft werden " "(Dienst nicht erreichbar) — bitte später erneut prüfen." ) if cfg.get("ai_analysis", {}).get("unchecked_level") == "yellow" and level == "green": level, exit_code = "yellow", 1 # Grave-Override: nur die schwersten Kategorien lösen ROT aus (Vorrang vor gelb). ai_cfg = cfg.get("ai_analysis", {}) red_cats = set(ai_cfg.get("red_categories", [])) red_min_sev = _SEVERITY_ORDER.get(ai_cfg.get("red_min_severity", "high"), 3) red_min_conf = ai_cfg.get("red_min_confidence", 0.9) for f in ai_result.get("findings", []): if (f.get("category") in red_cats and _SEVERITY_ORDER.get(f.get("severity", "none"), 0) >= red_min_sev and f.get("confidence", 0.0) >= red_min_conf): level, exit_code = "red", 2 reasons.append( f"Schwerwiegender KI-Fund ({f['category']}, {f.get('confidence', 0):.0%}) " f"auf {f.get('url', '?')} → ROT" ) break return {"score": score, "level": level, "reasons": reasons, "exit_code": exit_code} # --------------------------------------------------------------------------- # Text analysis # --------------------------------------------------------------------------- def _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor) -> bool: """Prüft neue/geänderte Seitentexte (parallel). Returns True wenn Ledger verändert.""" text_cfg = ai_cfg.get("text", {}) min_chars = text_cfg.get("min_chars", 200) models = _models_for(text_cfg) site_context = ai_cfg.get("site_context", "") # Fingerprint je Seite (lokal, gratis). Geänderte/neue Seiten zuerst (Risiko). changed = {pd["url"] for pd in diff.get("page_diffs", [])} | set(diff.get("new_internal_urls", [])) candidates: list[tuple[str, str, str]] = [] for url, page in snap.get("pages", {}).items(): if page.get("status") != 200: continue text = page.get("text", "") or "" if len(text) < min_chars: continue candidates.append((url, text, _text_fingerprint(text, cfg))) candidates.sort(key=lambda c: c[0] not in changed) # changed first (False < True) # Nach Fingerprint gruppieren: identischer Text wird nur EINMAL klassifiziert, # die Findings aber für alle betroffenen URLs emittiert (wie sequenziell). fp_urls: dict[str, list[str]] = {} fp_text: dict[str, str] = {} order_fps: list[str] = [] for url, text, fp in candidates: result["checked"] += 1 if fp not in fp_urls: fp_urls[fp] = [] fp_text[fp] = text order_fps.append(fp) fp_urls[fp].append(url) # Drossel NUR für den historischen Bestand; 0 = unbegrenzt. Geänderte/neue Seiten # werden IMMER geprüft (Echtzeit-Schutz), ungeachtet des Limits. backlog_limit = text_cfg.get("max_pages_per_scan", 0) backlog_used = 0 tasks, miss_fps = [], [] for fp in order_fps: entry = entries.get(fp) if entry is not None: for url in fp_urls[fp]: result["cache_hits"] += 1 _maybe_finding(result, entry, fp, url, ai_cfg) continue is_changed = any(u in changed for u in fp_urls[fp]) if not is_changed and backlog_limit and backlog_used >= backlog_limit: continue # Bestands-Drossel erreicht — nächster Scan holt den Rest nach miss_fps.append(fp) tasks.append((fp, _text_task(fp_text[fp], site_context, models, ai_cfg, api_key, router))) if not is_changed: backlog_used += 1 classified = _classify_batch(tasks, executor) dirty = False for fp in miss_fps: res = classified.get(fp) if res is None: for url in fp_urls[fp]: result["unchecked"].append({"kind": "text", "url": url}) continue verdict, used_model = res entry = _make_entry("text", fp_urls[fp][0], verdict, used_model) entries[fp] = entry result["api_calls"] += 1 dirty = True for url in fp_urls[fp]: _maybe_finding(result, entry, fp, url, ai_cfg) return dirty # --------------------------------------------------------------------------- # Image analysis # --------------------------------------------------------------------------- def _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor) -> bool: """Prüft neue/geänderte Bilder (parallel). Returns True wenn Ledger verändert.""" img_cfg = ai_cfg.get("image", {}) models = _models_for(img_cfg) site_context = ai_cfg.get("site_context", "") current_img_urls = _all_image_urls(snap) if not current_img_urls: return False changed = {pd["url"] for pd in diff.get("page_diffs", [])} | set(diff.get("new_internal_urls", [])) changed_img_urls = _all_image_urls(snap, only_pages=changed) # Markierte Bilder erneut holen: ersetztes Bild → neue Bytes → Neubewertung # (verhindert Dauer-Fehlalarm nach Bereinigung); unverändert → Cache-Treffer. flagged_img_urls = { e["url"] for e in entries.values() if e.get("kind") == "image" and e.get("category", "clean") != "clean" and not e.get("dismissed") and e.get("url") in current_img_urls } # Noch nie analysierte Bilder (URL nicht im Ledger): bauen über mehrere Scans die # VOLLE Abdeckung auf — sonst bliebe jedes Bild jenseits des ersten Budget-Fensters # ein dauerhafter blinder Fleck (z. B. ein eingeschleustes strafbares Logo auf einer # ansonsten stabilen Seite). Bytes-Änderungen bekannter Bilder deckt check-assets ab. known_img_urls = {e["url"] for e in entries.values() if e.get("kind") == "image"} coverage_urls = current_img_urls - known_img_urls # Geänderte/geflaggte Bilder: IMMER prüfen (Echtzeit-Schutz, nie gedeckelt). # Bestands-Abdeckung: Drossel optional; 0 = unbegrenzt (volle Abdeckung im ersten Scan). priority = sorted((changed_img_urls | flagged_img_urls) & current_img_urls) coverage = sorted(coverage_urls - set(priority)) backlog_limit = img_cfg.get("max_images_per_scan", 0) if backlog_limit and backlog_limit > 0: coverage = coverage[:backlog_limit] fetch_urls = priority + coverage if not fetch_urls: return False # url → Seite(n) für den Report url_to_page = _image_url_sources(snap) hashes = fetch_asset_hashes(fetch_urls, timeout=cfg.get("request_timeout", 15)) # Nach Bild-Hash (Bytes) gruppieren: identische Bilder werden nur einmal klassifiziert. fp_assets: dict[str, list[str]] = {} order_fps: list[str] = [] for url in fetch_urls: fp = hashes.get(url, {}).get("sha256") if not fp: continue # nicht abrufbar — überspringen result["checked"] += 1 if fp not in fp_assets: fp_assets[fp] = [] order_fps.append(fp) fp_assets[fp].append(url) tasks, miss_fps = [], [] for fp in order_fps: entry = entries.get(fp) if entry is not None: for url in fp_assets[fp]: result["cache_hits"] += 1 _maybe_finding(result, entry, fp, url_to_page.get(url, url), ai_cfg, asset_url=url) else: miss_fps.append(fp) tasks.append((fp, _image_task(fp_assets[fp][0], site_context, models, ai_cfg, api_key, router))) classified = _classify_batch(tasks, executor) dirty = False for fp in miss_fps: res = classified.get(fp) rep_url = fp_assets[fp][0] if res is None: for url in fp_assets[fp]: result["unchecked"].append( {"kind": "image", "url": url_to_page.get(url, url), "asset_url": url}) continue verdict, used_model = res entry = _make_entry("image", rep_url, verdict, used_model) entries[fp] = entry result["api_calls"] += 1 dirty = True for url in fp_assets[fp]: _maybe_finding(result, entry, fp, url_to_page.get(url, url), ai_cfg, asset_url=url) return dirty # --------------------------------------------------------------------------- # Audio / Video — vorbereitete Erweiterungspunkte (default deaktiviert) # --------------------------------------------------------------------------- def _collect_audio_candidates(snap: dict, cfg: dict) -> list[str]: """TODO: Audio-Quellen sammeln. Erfordert