integrity_scanner_fuer_stat.../scanner/ai_analyzer.py
Dieter Schlüter de5be23d23 fix: Mi2-Mi4 + Mi1 partial — version pins, file-lock, catalog lock, __main__ tests
- Mi2: requirements.txt — Obergrenzen hinzugefügt (<3.0, <5.0, etc.)
- Mi3: cmd_scan() — File-Lock via fcntl (data_dir/.scan.lock)
- Mi4: refresh_catalog() — Cache-Check unter self._lock (thread-safe)
- Mi1: tests/test_main.py — 17 neue Tests (__main__ Coverage 10% → 18%)
- 294/294 Tests grün
2026-06-14 21:36:24 +02:00

1072 lines
46 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
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 (12 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