""" 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 base64 import concurrent.futures import hashlib import json import logging import os import subprocess import threading import time from datetime import datetime, timezone from pathlib import Path import requests from .checker import fetch_asset_hashes from .differ import normalize_text logger = logging.getLogger(__name__) # Primärer Endpunkt: OpenRouter (wie bisher). _OPENROUTER_URL = "https://openrouter.ai/api/v1/chat/completions" _MODELS_URL = "https://openrouter.ai/api/v1/models" # Lokale llama.cpp-Server als Fallback — werden nur genutzt, wenn OpenRouter # nicht erreichbar ist (Rate-Limit, Timeout, Ausfall). _LOCAL_SERVERS: list[str] = [ "http://localhost:8001", "http://localhost:8002", ] # 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": { "description": {"type": "string"}, "category": {"type": "string", "enum": _CATEGORIES}, "severity": {"type": "string", "enum": ["none", "low", "medium", "high"]}, "confidence": {"type": "number"}, "explanation": {"type": "string"}, }, "required": ["description", "category", "severity", "confidence", "explanation"], "additionalProperties": False, }, }, } # --------------------------------------------------------------------------- # Adaptive model router (circuit-breaker + latency telemetry + /models pruning) # --------------------------------------------------------------------------- class ModelRouter: """Wählt programmatisch das beste Modell aus einer Kette: gesunde zuerst, nach Latenz; ausgefallene Modelle bekommen einen Cooldown (Circuit-Breaker). Katalog-Pruning funktioniert für OpenRouter- und lokale Modell-Slugs. 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 {} # catalog: {"slugs": [...], "servers": {...}, "fetched_at": ...} # slugs = OpenRouter-Katalog (primär), servers = lokale Server-Modelle saved = state.get("catalog", {"slugs": [], "servers": {}, "fetched_at": None}) self._catalog: dict = saved if isinstance(saved, dict) else {"slugs": [], "servers": {}, "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]: # Kombiniere OpenRouter-Slugs + lokale Server-Modelle. all_slugs: set[str] = set(self._catalog.get("slugs", [])) for slugs in self._catalog.get("servers", {}).values(): all_slugs.update(slugs) if not all_slugs: return list(models) # kein Katalog → nicht filtern (failsafe) keep = [m for m in models if m in all_slugs] return keep or list(models) # nie alles wegfiltern def refresh_catalog(self) -> None: """OpenRouter /models + lokale Server /models ziehen (24 h gecacht). Failsafe: blockiert nie. Thread-safe via self._lock.""" if not self._refresh: return # Cache-Check unter Lock — vermeidet doppelte Fetches bei parallelen Scans with self._lock: fetched_at = self._catalog.get("fetched_at") slugs = self._catalog.get("slugs") or [] servers = self._catalog.get("servers") or {} if fetched_at and (slugs or servers) and (time.time() - fetched_at) < 86400: return # Cache frisch new_slugs, new_servers = [], {} # 1. OpenRouter-Katalog (primär) try: resp = requests.get(_MODELS_URL, timeout=10) if resp.status_code == 200: new_slugs = [m["id"] for m in resp.json().get("data", []) if m.get("id")] logger.info("KI-Router: OpenRouter-Katalog aktualisiert (%d Modelle).", len(new_slugs)) except (requests.exceptions.RequestException, ValueError, KeyError) as exc: logger.debug("KI-Router: OpenRouter /models nicht abrufbar: %s", exc) # 2. Lokale Server-Kataloge (Fallback) for base_url in _LOCAL_SERVERS: try: resp = requests.get(f"{base_url}/v1/models", timeout=5) if resp.status_code == 200: data = resp.json() raw = data.get("data") or data.get("models", []) local_slugs = [m["id"] for m in raw if isinstance(m, dict) and m.get("id")] if local_slugs: new_servers[base_url] = local_slugs logger.info("KI-Router: %s — %d Modelle gefunden.", base_url, len(local_slugs)) elif resp.status_code == 503: logger.debug("KI-Router: %s — Modell wird geladen (503).", base_url) except requests.exceptions.RequestException as exc: logger.debug("KI-Router: %s nicht erreichbar: %s", base_url, exc) if new_slugs or new_servers: with self._lock: self._catalog = {"slugs": new_slugs, "servers": new_servers, "fetched_at": time.time()} total = len(new_slugs) + sum(len(v) for v in new_servers.values()) logger.info("KI-Router: Katalog aktualisiert (%d Modelle gesamt).", total) # --------------------------------------------------------------------------- # 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", {}) url_hashes = ledger.setdefault("url_hashes", {}) # URL→Byte-Hash-Erinnerung (Bilder) 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, url_hashes) if ai_cfg.get("audio", {}).get("enabled", False): dirty |= _analyze_audio(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor, url_hashes) if ai_cfg.get("video", {}).get("enabled", False): dirty |= _analyze_video(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor, url_hashes) 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, url_hashes) -> bool: """Prüft neue/geänderte Bilder (parallel). Returns True wenn Ledger verändert. url_hashes ist die persistente URL→Byte-Hash-Erinnerung: eine bereits gesehene Bild-URL wird NICHT erneut heruntergeladen (kein redundanter Fetch, keine doppelte Analyse). Nur nie gesehene Bilder sowie geänderte/geflaggte werden geholt.""" 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 } # Geänderte/geflaggte Bilder: IMMER (erneut) prüfen (Echtzeit-Schutz, nie gedeckelt). # Abdeckung: nur noch NIE gesehene URLs (nicht in url_hashes) — schon bekannte Bilder # werden nicht erneut geladen (Hashwert gemerkt). Volle Abdeckung baut sich so EINMALIG # auf; Duplikate (gleiche Bytes an mehreren URLs) werden je URL höchstens einmal geholt # und nie doppelt analysiert (Ledger ist nach Byte-Hash indiziert). Drossel 0 = unbegrenzt. priority = sorted((changed_img_urls | flagged_img_urls) & current_img_urls) coverage = sorted((current_img_urls - set(url_hashes)) - 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. # Zugleich URL→Hash merken (auch für Duplikate), damit künftige Scans sie überspringen. dirty = False 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 (nicht merken → später erneut versuchen) if url_hashes.get(url) != fp: url_hashes[url] = fp # neu gesehen oder Bytes geändert dirty = True 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) 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 analysis # --------------------------------------------------------------------------- def _analyze_audio(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor, url_hashes) -> bool: """Prüft neue/geänderte Audio-Dateien (parallel). Returns True wenn Ledger verändert.""" aud_cfg = ai_cfg.get("audio", {}) models = _models_for(aud_cfg) site_context = ai_cfg.get("site_context", "") current_urls = _all_media_urls(snap, "audio") if not current_urls: return False changed = {pd["url"] for pd in diff.get("page_diffs", [])} | set(diff.get("new_internal_urls", [])) changed_urls = _all_media_urls(snap, "audio", only_pages=changed) flagged_urls = { e["url"] for e in entries.values() if e.get("kind") == "audio" and e.get("category", "clean") != "clean" and not e.get("dismissed") and e.get("url") in current_urls } priority = sorted((changed_urls | flagged_urls) & current_urls) coverage = sorted((current_urls - set(url_hashes)) - set(priority)) backlog_limit = aud_cfg.get("max_files_per_scan", 5) if backlog_limit and backlog_limit > 0: coverage = coverage[:backlog_limit] fetch_urls = priority + coverage if not fetch_urls: return False url_to_page = _media_url_sources(snap, "audio") # Sampling: für jede URL via ffmpeg einen 90-s-Ausschnitt holen (hash + bytes). # Ersetzt _local_file_hashes + fetch_asset_hashes — kein 40-60 MB Download mehr. samples: dict[str, tuple[bytes, str]] = {} # url → (mp3_bytes, sha256) dirty = False fp_assets: dict[str, list[str]] = {} order_fps: list[str] = [] for url in fetch_urls: sample = _audio_sample(url, timeout=ai_cfg.get("request_timeout", 60)) if sample is None: result["unchecked"].append( {"kind": "audio", "url": url_to_page.get(url, url), "asset_url": url} ) continue data, fp = sample samples[url] = (data, fp) if url_hashes.get(url) != fp: url_hashes[url] = fp dirty = True 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) rep_bytes = samples[fp_assets[fp][0]][0] tasks.append((fp, _audio_task(rep_bytes, site_context, models, ai_cfg, api_key, router))) classified = _classify_batch(tasks, executor) 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": "audio", "url": url_to_page.get(url, url), "asset_url": url}) continue verdict, used_model = res entry = _make_entry("audio", 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 # --------------------------------------------------------------------------- # Video analysis # --------------------------------------------------------------------------- def _analyze_video(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor, url_hashes) -> bool: """Prüft neue/geänderte Video-Dateien (parallel). Returns True wenn Ledger verändert.""" vid_cfg = ai_cfg.get("video", {}) models = _models_for(vid_cfg) site_context = ai_cfg.get("site_context", "") current_urls = _all_media_urls(snap, "video") if not current_urls: return False changed = {pd["url"] for pd in diff.get("page_diffs", [])} | set(diff.get("new_internal_urls", [])) changed_urls = _all_media_urls(snap, "video", only_pages=changed) flagged_urls = { e["url"] for e in entries.values() if e.get("kind") == "video" and e.get("category", "clean") != "clean" and not e.get("dismissed") and e.get("url") in current_urls } priority = sorted((changed_urls | flagged_urls) & current_urls) coverage = sorted((current_urls - set(url_hashes)) - set(priority)) backlog_limit = vid_cfg.get("max_files_per_scan", 3) if backlog_limit and backlog_limit > 0: coverage = coverage[:backlog_limit] fetch_urls = priority + coverage if not fetch_urls: return False url_to_page = _media_url_sources(snap, "video") local_urls = [u for u in fetch_urls if u.startswith("file://")] remote_urls = [u for u in fetch_urls if not u.startswith("file://")] hashes = _local_file_hashes(local_urls) if remote_urls: hashes.update(fetch_asset_hashes(remote_urls, timeout=cfg.get("request_timeout", 15))) dirty = False 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 if url_hashes.get(url) != fp: url_hashes[url] = fp dirty = True 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, _video_task(fp_assets[fp][0], site_context, models, ai_cfg, api_key, router))) classified = _classify_batch(tasks, executor) 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": "video", "url": url_to_page.get(url, url), "asset_url": url}) continue verdict, used_model = res entry = _make_entry("video", 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 # --------------------------------------------------------------------------- # Fingerprinting & helpers # --------------------------------------------------------------------------- def _text_fingerprint(text: str, cfg: dict) -> str: """SHA-256 des normalisierten Texts (Reuse von differ.normalize_text → Cache- Buster wie CF7-Platzhalter ändern den Hash nicht).""" norm = normalize_text(text, cfg) return hashlib.sha256(norm.encode("utf-8")).hexdigest() def _make_entry(kind: str, url: str, verdict: dict, model: str) -> dict: return { "kind": kind, "url": url, "description": verdict.get("description", ""), "category": verdict.get("category", "clean"), "severity": verdict.get("severity", "none"), "confidence": float(verdict.get("confidence", 0.0)), "explanation": verdict.get("explanation", ""), "checked_at": datetime.now(timezone.utc).isoformat(), "model": model, "dismissed": False, } def _maybe_finding(result, entry, fp, url, ai_cfg, asset_url=None) -> None: """Hängt einen Fund an, wenn das Verdikt gewertet werden soll (nicht clean, nicht quittiert, severity≥medium, confidence≥Schwelle).""" if entry.get("dismissed"): return if entry.get("category", "clean") == "clean": return if _SEVERITY_ORDER.get(entry.get("severity", "none"), 0) < _SEVERITY_ORDER["medium"]: return if entry.get("confidence", 0.0) < ai_cfg.get("ai_confidence_min", 0.7): return finding = { "kind": entry["kind"], "url": url, "fingerprint": fp, "category": entry["category"], "severity": entry["severity"], "confidence": entry["confidence"], "explanation": entry.get("explanation", ""), } if asset_url: finding["asset_url"] = asset_url result["findings"].append(finding) def _finding_points(finding: dict, sc: dict) -> int: """Punkte für einen Fund. Bilder/Audio/Video mindestens ai_suspicious_image.""" key, default = _CATEGORY_SCORE_KEY.get(finding.get("category", ""), (None, 30)) pts = sc.get(key, default) if key else 30 if finding.get("kind") in ("image", "audio", "video"): pts = max(pts, sc.get("ai_suspicious_image", 40)) return pts def _all_image_urls(snap: dict, only_pages: set | None = None) -> set[str]: urls: set[str] = set() for page_url, page in snap.get("pages", {}).items(): if only_pages is not None and page_url not in only_pages: continue if page.get("status") != 200: continue for link in page.get("links", {}).get("img", []): if isinstance(link, dict) and link.get("url"): urls.add(link["url"]) return urls def _image_url_sources(snap: dict) -> dict[str, str]: """Map image-URL → erste Seite, auf der es vorkommt (für den Report).""" sources: dict[str, str] = {} for page_url, page in snap.get("pages", {}).items(): for link in page.get("links", {}).get("img", []): if isinstance(link, dict) and link.get("url"): sources.setdefault(link["url"], page_url) return sources def _all_media_urls(snap: dict, kind: str, only_pages: set | None = None) -> set[str]: urls: set[str] = set() for page_url, page in snap.get("pages", {}).items(): if only_pages is not None and page_url not in only_pages: continue if page.get("status") != 200: continue for link in page.get("links", {}).get(kind, []): if isinstance(link, dict) and link.get("url"): urls.add(link["url"]) return urls def _media_url_sources(snap: dict, kind: str) -> dict[str, str]: """Map Medien-URL → erste Seite, auf der sie vorkommt (für den Report).""" sources: dict[str, str] = {} for page_url, page in snap.get("pages", {}).items(): for link in page.get("links", {}).get(kind, []): if isinstance(link, dict) and link.get("url"): sources.setdefault(link["url"], page_url) return sources # --------------------------------------------------------------------------- # OpenRouter classification # --------------------------------------------------------------------------- _SYSTEM_PROMPT = ( "Du bist ein Sicherheits-Prüfer für Website-Inhalte. Die zu prüfende Website " "hat folgendes Thema:\n\"{context}\"\n\n" "Beurteile, ob der folgende Inhalt problematisch ist. Achte auf: sexuell " "expliziten Text/Pornografie, politische Propaganda, diffamierende oder " "strafbare Aussagen, versteckten Spam, thematisch völlig unpassende " "kommerzielle Werbung, in sich widersprüchliche Aussagen.\n" "WICHTIG: Beurteile die thematische PASSUNG zum oben genannten Thema, nicht " "einzelne Schlüsselwörter. Heimat-, Orts- oder Vereinsbegriffe allein machen " "einen Werbe- oder Spam-Text nicht unproblematisch — ein Angreifer kann sie " "gezielt einstreuen. Bewerte den Gesamtcharakter des Inhalts.\n" "Gib 'clean' zurück, wenn der Inhalt zum Thema passt und unproblematisch ist. " "Gib im Feld 'description' eine knappe, neutrale Inhaltsangabe (1–2 Sätze, deutsch): " "was ist auf dem Bild zu sehen bzw. worum geht es im Text — sachlich, unabhängig vom " "Sicherheitsurteil. Antworte ausschließlich im vorgegebenen JSON-Format." ) def _models_for(modality_cfg: dict) -> list[str]: """Modell-Kette einer Modalität. Rückwärts-kompatibel: einzelnes `model` → Liste.""" models = modality_cfg.get("models") if models: return list(models) single = modality_cfg.get("model") return [single] if single else [] def _classify(messages, models, ai_cfg, api_key, router) -> tuple[dict, str] | None: """Klassifiziert über die vom Router gereihte Modell-Kette. Erste gültige Antwort → (verdict, model_used). Alle gescheitert → None. Der Router reiht gesunde/schnelle Modelle nach vorn und überspringt rate-limited/ausgefallene (Circuit-Breaker). Harter Deadline je Versuch; bei Komplettausfall bis max_retries Wiederholung (Backoff).""" timeout = ai_cfg.get("attempt_timeout", 20) max_retries = ai_cfg.get("max_retries", 1) backoff = ai_cfg.get("retry_backoff_seconds", 2.0) for attempt in range(max_retries + 1): ordered = router.order(models) for model in ordered: t0 = time.time() verdict = _openrouter_chat(model, messages, ai_cfg, api_key, timeout=timeout) if verdict is not None: router.record_success(model, time.time() - t0) return verdict, model router.record_failure(model) if attempt < max_retries: logger.info("KI-Kette komplett erfolglos — Wiederholung %d/%d in %.1fs.", attempt + 1, max_retries, backoff) time.sleep(backoff) return None def _classify_text(text, site_context, models, ai_cfg, api_key, router) -> tuple[dict, str] | None: messages = [ {"role": "system", "content": _SYSTEM_PROMPT.format(context=site_context or "(nicht angegeben)")}, {"role": "user", "content": f"Zu prüfender Seitentext:\n\n{text[:8000]}"}, ] return _classify(messages, models, ai_cfg, api_key, router) def _classify_image(image_url, site_context, models, ai_cfg, api_key, router) -> tuple[dict, str] | None: messages = [ {"role": "system", "content": _SYSTEM_PROMPT.format(context=site_context or "(nicht angegeben)")}, {"role": "user", "content": [ {"type": "text", "text": "Prüfe dieses Bild auf problematische Inhalte. " "Beziehe auch im Bild sichtbaren Text mit ein (OCR)."}, {"type": "image_url", "image_url": {"url": image_url}}, ]}, ] return _classify(messages, models, ai_cfg, api_key, router) def _text_task(text, site_context, models, ai_cfg, api_key, router): """Factory: gibt einen parameterlosen Aufruf für den Thread-Pool zurück.""" return lambda: _classify_text(text, site_context, models, ai_cfg, api_key, router) def _image_task(image_url, site_context, models, ai_cfg, api_key, router): return lambda: _classify_image(image_url, site_context, models, ai_cfg, api_key, router) def _local_file_hashes(urls: list[str]) -> dict[str, dict]: """Berechnet SHA-256 für file://-URLs direkt aus dem Dateisystem (kein HTTP).""" result: dict[str, dict] = {} for url in urls: if not url.startswith("file://"): continue path = Path(url[7:]) try: data = path.read_bytes() result[url] = {"sha256": hashlib.sha256(data).hexdigest(), "size": len(data)} except OSError as exc: logger.warning("Lokale Datei nicht lesbar (%s): %s", url, exc) return result def _audio_sample(url: str, timeout: int = 60) -> tuple[bytes, str] | None: """Extrahiert via ffmpeg einen 90-s-Ausschnitt ab Sekunde 30 (Intro überspringen). Arbeitet für file://-Pfade und https://-URLs gleich; lädt nur ~3 MB statt 40-60 MB. Gibt (mp3_bytes, sha256_hex) zurück oder None bei Fehler.""" input_path = url[7:] if url.startswith("file://") else url def _run(ss: str) -> subprocess.CompletedProcess: return subprocess.run( ["ffmpeg", "-hide_banner", "-loglevel", "error", "-i", input_path, "-ss", ss, "-t", "90", "-ar", "22050", "-ac", "1", "-b:a", "32k", "-f", "mp3", "pipe:1"], capture_output=True, timeout=timeout, ) try: r = _run("30") if r.returncode != 0 or not r.stdout: r = _run("0") # Fallback: Datei kürzer als 30 s if r.returncode != 0 or not r.stdout: logger.warning("ffmpeg-Sampling fehlgeschlagen für %s: %s", url, r.stderr.decode()[:200]) return None return r.stdout, hashlib.sha256(r.stdout).hexdigest() except (subprocess.TimeoutExpired, FileNotFoundError, OSError) as exc: logger.warning("ffmpeg nicht aufrufbar für %s: %s", url, exc) return None def _media_content(url: str, media_type: str) -> dict: """Gibt den OpenRouter-Message-Content für eine Medien-URL zurück. Bei file://-Pfaden: Datei einlesen und als inline-base64 senden (input_audio/input_video). Bei http(s)-URLs: URL-Referenz (audio_url/video_url).""" if url.startswith("file://"): path = Path(url[7:]) data = path.read_bytes() b64 = base64.b64encode(data).decode() fmt = path.suffix.lstrip(".").lower() or media_type if media_type == "audio": return {"type": "input_audio", "input_audio": {"data": b64, "format": fmt}} return {"type": "video_url", "video_url": {"url": f"data:video/{fmt};base64,{b64}"}} if media_type == "audio": return {"type": "audio_url", "audio_url": {"url": url}} return {"type": "video_url", "video_url": {"url": url}} def _classify_audio(audio_bytes: bytes, site_context, models, ai_cfg, api_key, router) -> tuple[dict, str] | None: b64 = base64.b64encode(audio_bytes).decode() messages = [ {"role": "system", "content": _SYSTEM_PROMPT.format(context=site_context or "(nicht angegeben)")}, {"role": "user", "content": [ {"type": "text", "text": ( "Prüfe diesen Audio-Ausschnitt (bis zu 90 Sekunden, ab Sekunde 30) " "auf problematische Inhalte (Sprache, Geräusche, Musik — soweit erkennbar)." )}, {"type": "input_audio", "input_audio": {"data": b64, "format": "mp3"}}, ]}, ] return _classify(messages, models, ai_cfg, api_key, router) def _classify_video(video_url, site_context, models, ai_cfg, api_key, router) -> tuple[dict, str] | None: messages = [ {"role": "system", "content": _SYSTEM_PROMPT.format(context=site_context or "(nicht angegeben)")}, {"role": "user", "content": [ {"type": "text", "text": "Prüfe dieses Video auf problematische Inhalte " "(Bild, Ton, eingebetteter Text — soweit erkennbar)."}, _media_content(video_url, "video"), ]}, ] return _classify(messages, models, ai_cfg, api_key, router) def _audio_task(audio_bytes: bytes, site_context, models, ai_cfg, api_key, router): return lambda: _classify_audio(audio_bytes, site_context, models, ai_cfg, api_key, router) def _video_task(video_url, site_context, models, ai_cfg, api_key, router): return lambda: _classify_video(video_url, site_context, models, ai_cfg, api_key, router) def _classify_batch(tasks, executor) -> dict: """Führt [(key, fn)] nebenläufig im Pool aus. Returns {key: fn()-Ergebnis|None}. Eine fehlgeschlagene Aufgabe blubbert nie hoch (→ None, der Aufrufer wertet das als 'ungeprüft').""" if not tasks: return {} futures = {executor.submit(fn): key for key, fn in tasks} out: dict = {} for fut in concurrent.futures.as_completed(futures): key = futures[fut] try: out[key] = fut.result() except Exception as exc: # pragma: no cover - Schutzschirm logger.warning("KI-Klassifikation fehlgeschlagen (%s): %s", key, exc) out[key] = None return out def _post_with_deadline(url: str, headers, payload, timeout): """requests.post in einem Worker-Thread mit HARTEM Wanduhr-Limit. Der requests-Timeout ist nur ein Inaktivitäts-Timeout: ein Server, der die Verbindung offen hält oder Tokens langsam tröpfeln lässt, kann ihn umgehen und den Aufruf unbegrenzt hängen lassen. Dieser Deadline kappt die Gesamtdauer hart; bei Überschreitung läuft der verwaiste Thread im Hintergrund aus und wir eskalieren. """ ex = concurrent.futures.ThreadPoolExecutor(max_workers=1) fut = ex.submit(requests.post, url, headers=headers, json=payload, timeout=timeout) try: resp = fut.result(timeout=timeout) ex.shutdown(wait=False) return resp except concurrent.futures.TimeoutError: ex.shutdown(wait=False) # Thread nicht abwarten — wir machen weiter raise requests.exceptions.Timeout(f"Hartes Zeitlimit {timeout}s überschritten") def _openrouter_chat(model, messages, ai_cfg, api_key, timeout=None) -> dict | None: """Chat-Call: zuerst OpenRouter (primär), dann lokale llama.cpp-Server als Fallback. Returns das Verdikt-Dict oder None wenn alle Endpunkte scheitern (graceful).""" if timeout is None: timeout = ai_cfg.get("attempt_timeout", ai_cfg.get("request_timeout", 20)) payload = { "model": model, "messages": messages, "response_format": _RESPONSE_SCHEMA, "temperature": 0, } # --- 1. OpenRouter (primär) --- headers = { "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", "X-Title": "integrity-scanner", } try: resp = _post_with_deadline(_OPENROUTER_URL, headers, payload, timeout) if resp.status_code == 200: content = resp.json()["choices"][0]["message"]["content"] verdict = json.loads(content) if verdict.get("category") not in _CATEGORIES: logger.warning("OpenRouter (%s): unbekannte Kategorie %r", model, verdict.get("category")) else: logger.debug("KI-Antwort von OpenRouter (Modell: %s).", model) return verdict else: logger.info("OpenRouter HTTP %s (%s) — Fallback zu lokalen Servern.", resp.status_code, model) except (requests.exceptions.RequestException, KeyError, ValueError, json.JSONDecodeError) as exc: logger.debug("OpenRouter-Call fehlgeschlagen (%s): %s", model, exc) # --- 2. Lokale Server (Fallback) --- headers = {"Content-Type": "application/json"} if api_key: headers["Authorization"] = f"Bearer {api_key}" for base_url in _LOCAL_SERVERS: url = f"{base_url}/v1/chat/completions" try: resp = _post_with_deadline(url, headers, payload, timeout) if resp.status_code == 200: content = resp.json()["choices"][0]["message"]["content"] verdict = json.loads(content) if verdict.get("category") not in _CATEGORIES: logger.warning("KI (%s@%s): unbekannte Kategorie %r", model, base_url, verdict.get("category")) continue logger.debug("KI-Antwort von %s (Modell: %s).", base_url, model) return verdict elif resp.status_code == 503: logger.info("KI (%s@%s): Modell wird geladen — nächster Server.", model, base_url) else: logger.warning("KI HTTP %s (%s@%s): %s", resp.status_code, model, base_url, resp.text[:200]) except (requests.exceptions.RequestException, KeyError, ValueError, json.JSONDecodeError) as exc: logger.debug("KI-Call fehlgeschlagen (%s@%s): %s", model, base_url, exc) logger.warning("KI: OpenRouter + alle lokalen Server erfolglos für Modell %s.", model) return None