- ai_analyzer: OpenRouter primary + localhost:8001/:8002 as dev fallback - crawler: support include_paths for targeted crawling - __main__: pass config_path through report pipeline - tests: 266 tests passing
1070 lines
46 KiB
Python
1070 lines
46 KiB
Python
"""
|
||
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."""
|
||
if not self._refresh:
|
||
return
|
||
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
|