feat: add local server fallback to AI analyzer, include_paths support, config_path in reports
- 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
This commit is contained in:
parent
989f5a933f
commit
013cf794bc
13 changed files with 479 additions and 60 deletions
|
|
@ -14,13 +14,17 @@ 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
|
||||
|
||||
|
|
@ -29,7 +33,16 @@ 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 = [
|
||||
|
|
@ -71,7 +84,7 @@ _RESPONSE_SCHEMA = {
|
|||
}
|
||||
|
||||
|
||||
_MODELS_URL = "https://openrouter.ai/api/v1/models"
|
||||
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -80,9 +93,9 @@ _MODELS_URL = "https://openrouter.ai/api/v1/models"
|
|||
|
||||
class ModelRouter:
|
||||
"""Wählt programmatisch das beste Modell aus einer Kette: gesunde zuerst,
|
||||
nach Latenz; rate-limited/ausgefallene Modelle bekommen einen Cooldown
|
||||
(Circuit-Breaker) und werden so lange übersprungen. Tote Slugs werden über
|
||||
den /models-Katalog failsafe aussortiert. Thread-sicher (für Parallelität)."""
|
||||
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()
|
||||
|
|
@ -91,7 +104,10 @@ class ModelRouter:
|
|||
self._refresh = ai_cfg.get("refresh_models", True)
|
||||
state = state or {}
|
||||
self._models: dict = state.get("models", {}) or {}
|
||||
self._catalog: dict = state.get("catalog", {"slugs": [], "fetched_at": None})
|
||||
# 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:
|
||||
|
|
@ -136,29 +152,59 @@ class ModelRouter:
|
|||
return ordered or list(models)
|
||||
|
||||
def _prune_locked(self, models: list[str]) -> list[str]:
|
||||
slugs = self._catalog.get("slugs") or []
|
||||
if not slugs:
|
||||
# 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 slugs]
|
||||
keep = [m for m in models if m in all_slugs]
|
||||
return keep or list(models) # nie alles wegfiltern
|
||||
|
||||
def refresh_catalog(self) -> None:
|
||||
"""Einmal pro Lauf /models ziehen (24 h gecacht). Failsafe: blockiert nie."""
|
||||
"""OpenRouter /models + lokale Server /models ziehen (24 h gecacht).
|
||||
Failsafe: blockiert nie."""
|
||||
if not self._refresh:
|
||||
return
|
||||
fetched_at = self._catalog.get("fetched_at")
|
||||
if fetched_at and self._catalog.get("slugs") and (time.time() - fetched_at) < 86400:
|
||||
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=15)
|
||||
resp = requests.get(_MODELS_URL, timeout=10)
|
||||
if resp.status_code == 200:
|
||||
slugs = [m["id"] for m in resp.json().get("data", []) if m.get("id")]
|
||||
if slugs:
|
||||
with self._lock:
|
||||
self._catalog = {"slugs": slugs, "fetched_at": time.time()}
|
||||
logger.info("KI-Router: /models-Katalog aktualisiert (%d Modelle).", len(slugs))
|
||||
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.warning("KI-Router: /models nicht abrufbar (failsafe, ungefiltert): %s", 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)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -209,7 +255,12 @@ def run_ai_analysis(cfg: dict, bm, snap: dict, diff: dict) -> dict:
|
|||
if ai_cfg.get("image", {}).get("enabled", True):
|
||||
dirty |= _analyze_images(cfg, ai_cfg, api_key, snap, diff,
|
||||
entries, result, router, executor, url_hashes)
|
||||
# Audio/Video: vorbereitet, default aus (siehe _collect_audio/video_candidates).
|
||||
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)
|
||||
|
|
@ -453,19 +504,178 @@ def _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result, router, e
|
|||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Audio / Video — vorbereitete Erweiterungspunkte (default deaktiviert)
|
||||
# Audio analysis
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _collect_audio_candidates(snap: dict, cfg: dict) -> list[str]:
|
||||
"""TODO: Audio-Quellen sammeln. Erfordert <audio>/<source>-Extraktion in
|
||||
extractor.py (heute nicht extrahiert). Aktivierbar über ai_analysis.audio.enabled."""
|
||||
return []
|
||||
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
|
||||
|
||||
|
||||
def _collect_video_candidates(snap: dict, cfg: dict) -> list[str]:
|
||||
"""TODO: Video-Quellen sammeln. Erfordert <video>/<source>-Extraktion in
|
||||
extractor.py + Frame-Sampling. Aktivierbar über ai_analysis.video.enabled."""
|
||||
return []
|
||||
# ---------------------------------------------------------------------------
|
||||
# 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
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -475,7 +685,6 @@ def _collect_video_candidates(snap: dict, cfg: dict) -> list[str]:
|
|||
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)."""
|
||||
import hashlib
|
||||
norm = normalize_text(text, cfg)
|
||||
return hashlib.sha256(norm.encode("utf-8")).hexdigest()
|
||||
|
||||
|
|
@ -521,10 +730,10 @@ def _maybe_finding(result, entry, fp, url, ai_cfg, asset_url=None) -> None:
|
|||
|
||||
|
||||
def _finding_points(finding: dict, sc: dict) -> int:
|
||||
"""Punkte für einen Fund. Bilder mindestens ai_suspicious_image."""
|
||||
"""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") == "image":
|
||||
if finding.get("kind") in ("image", "audio", "video"):
|
||||
pts = max(pts, sc.get("ai_suspicious_image", 40))
|
||||
return pts
|
||||
|
||||
|
|
@ -552,6 +761,29 @@ def _image_url_sources(snap: dict) -> dict[str, str]:
|
|||
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
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -638,6 +870,101 @@ 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
|
||||
|
|
@ -656,7 +983,7 @@ def _classify_batch(tasks, executor) -> dict:
|
|||
return out
|
||||
|
||||
|
||||
def _post_with_deadline(headers, payload, timeout):
|
||||
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
|
||||
|
|
@ -665,7 +992,7 @@ def _post_with_deadline(headers, payload, timeout):
|
|||
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, _OPENROUTER_URL, headers=headers,
|
||||
fut = ex.submit(requests.post, url, headers=headers,
|
||||
json=payload, timeout=timeout)
|
||||
try:
|
||||
resp = fut.result(timeout=timeout)
|
||||
|
|
@ -677,8 +1004,9 @@ def _post_with_deadline(headers, payload, timeout):
|
|||
|
||||
|
||||
def _openrouter_chat(model, messages, ai_cfg, api_key, timeout=None) -> dict | None:
|
||||
"""Einzelner OpenRouter-Call mit strukturierter JSON-Antwort.
|
||||
Returns das Verdikt-Dict oder None bei jedem Fehler/Timeout (graceful)."""
|
||||
"""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 = {
|
||||
|
|
@ -687,23 +1015,56 @@ def _openrouter_chat(model, messages, ai_cfg, api_key, timeout=None) -> dict | N
|
|||
"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(headers, payload, timeout)
|
||||
if resp.status_code != 200:
|
||||
logger.warning("OpenRouter HTTP %s (%s): %s", resp.status_code, model, resp.text[:200])
|
||||
return None
|
||||
content = resp.json()["choices"][0]["message"]["content"]
|
||||
verdict = json.loads(content)
|
||||
# Mindest-Validierung
|
||||
if verdict.get("category") not in _CATEGORIES:
|
||||
logger.warning("OpenRouter (%s): unbekannte Kategorie %r", model, verdict.get("category"))
|
||||
return None
|
||||
return verdict
|
||||
except (requests.exceptions.RequestException, KeyError, ValueError, json.JSONDecodeError) as exc:
|
||||
logger.warning("OpenRouter-Call fehlgeschlagen (%s): %s", model, exc)
|
||||
return None
|
||||
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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue