feat: adaptiver Modell-Router (Circuit-Breaker) + parallele KI-Analyse
Der reale Erst-Scan dauerte ~10 min, weil pro Kandidat sequenziell erst die rate-limited Free-Modelle (429) probiert wurden. Zwei Hebel beheben das: 1. ModelRouter: programmatische Telemetrie statt LLM-Orchestrator. Pro Modell Erfolg/Latenz; ausgefallene/rate-limited Modelle bekommen per Circuit-Breaker einen Cooldown und werden übersprungen, das schnellste gesunde Modell zuerst. Failsafe-Pruning toter Slugs über OpenRouter /models (24h-Cache). State persistent in data/ai_router_state.json. 2. Parallelisierung: ThreadPoolExecutor mit konfigurierbarer max_concurrency (Default 8, I/O-gebunden → an API-Rate-Limits gebunden, nicht an CPU-Kerne). Klassifikationen laufen nebenläufig, Merge im Hauptthread (keine Locks), findings deterministisch sortiert. Fingerprint-Gruppierung: identischer Inhalt wird nur einmal klassifiziert, Funde aber für alle URLs emittiert. Realtest bredelar.info: Frisch-Scan von ~10 min auf 1:01 min; Breaker öffnete 14× (Free-Modelle übersprungen), alle Verdikte von gemini-2.5-flash-lite. Config: max_concurrency, breaker_failure_threshold, breaker_cooldown_seconds, refresh_models. Tests: 251 grün (+12: Router-Reihung/Breaker/Pruning-failsafe, parallel==sequenziell, Dedup, deterministische Reihenfolge). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
31730403b3
commit
125d80e7f9
7 changed files with 415 additions and 65 deletions
|
|
@ -984,6 +984,34 @@ Hinweis: Kostenlose Modelle liefern nicht immer ein gültiges Ergebnis (sie unte
|
|||
strikte JSON-Format nicht durchgehend). In dem Fall greift automatisch die nächste Stufe — das
|
||||
zuverlässige Bezahlmodell auf Stufe 3 fängt solche Fälle ab.
|
||||
|
||||
### Adaptiver Modell-Router und Parallelität
|
||||
|
||||
Statt die Kette stur von vorn abzuarbeiten, **lernt** der Scanner während des Laufs, welche
|
||||
Modelle gerade funktionieren:
|
||||
|
||||
- **Circuit-Breaker:** Antwortet ein Modell wiederholt mit Fehler/Rate-Limit (`429`) oder zu
|
||||
langsam, wird es für `breaker_cooldown_seconds` übersprungen — keine verschwendeten Versuche
|
||||
mehr an gerade tote Free-Modelle. Das schnellste gesunde Modell kommt zuerst dran.
|
||||
- **Live-Abgleich (`refresh_models`):** Einmal täglich gleicht der Router die Modell-Liste mit
|
||||
OpenRouter ab; nicht mehr existierende Slugs fallen automatisch raus (failsafe — bei Netzfehler
|
||||
bleibt alles wie gehabt).
|
||||
- **Parallelität (`max_concurrency`):** Mehrere Inhalte werden gleichzeitig geprüft. Die KI-Calls
|
||||
warten fast nur auf die API-Antwort (Netzwerk), nicht auf die CPU — deshalb ist die Stellgröße
|
||||
an die **Rate-Limits der API** gebunden, nicht an die CPU-Kerne. Default 8 passt auf einen
|
||||
4-Kern-vhost genauso wie auf einen 24-Kern-Rechner.
|
||||
|
||||
```yaml
|
||||
ai_analysis:
|
||||
max_concurrency: 8 # gleichzeitige API-Anfragen
|
||||
breaker_failure_threshold: 2 # Fehler, bis ein Modell pausiert wird
|
||||
breaker_cooldown_seconds: 120 # Pausendauer eines ausgefallenen Modells
|
||||
refresh_models: true # tote Modell-Slugs automatisch aussortieren
|
||||
```
|
||||
|
||||
Wirkung: Der Erst-Scan einer Site (der den ganzen Bestand prüft) wird deutlich schneller, weil
|
||||
rate-limited Free-Modelle sofort übersprungen werden und die verbleibenden Anfragen nebenläufig
|
||||
laufen.
|
||||
|
||||
Audio- und Video-Prüfung sind als abschaltbare Erweiterungspunkte vorbereitet
|
||||
(`ai_analysis.audio` / `ai_analysis.video`, default deaktiviert) und können aktiviert werden,
|
||||
sobald eine Website solche Mediendateien direkt einbindet.
|
||||
|
|
|
|||
|
|
@ -274,6 +274,13 @@ kostenlose Modelle, Stufe 3 ein günstiges Bezahlmodell. Schlägt ein Modell feh
|
|||
es langsamer als `attempt_timeout`, wird automatisch zur nächsten Stufe eskaliert. Im Report
|
||||
und Ledger wird festgehalten, welche Stufe das Verdikt geliefert hat.
|
||||
|
||||
Ein **adaptiver Router** lernt während des Laufs, welche Modelle funktionieren: rate-limited oder
|
||||
ausgefallene Modelle werden per Circuit-Breaker (`breaker_*`) für eine Cooldown-Zeit übersprungen,
|
||||
das schnellste gesunde Modell zuerst probiert, und tote Slugs über OpenRouter `/models`
|
||||
(`refresh_models`) failsafe aussortiert. Mehrere Inhalte werden **parallel** geprüft
|
||||
(`max_concurrency`, Default 8) — die Stellgröße ist an die API-Rate-Limits gebunden, nicht an
|
||||
CPU-Kerne. Das beschleunigt vor allem den Erst-Scan einer Site erheblich.
|
||||
|
||||
`site_context` ist entscheidend: Die KI bewertet die **thematische Passung** zur
|
||||
deklarierten Beschreibung, nicht einzelne Schlüsselwörter — eingestreute Heimat-Begriffe
|
||||
machen Spam so nicht unauffällig.
|
||||
|
|
|
|||
|
|
@ -148,6 +148,14 @@ ai_analysis:
|
|||
attempt_timeout: 20 # HARTES Wanduhr-Limit je Modell-Versuch (Sek.) → Eskalation
|
||||
max_retries: 1 # Wiederholungen der ganzen Modell-Kette bei Komplettausfall
|
||||
retry_backoff_seconds: 2.0
|
||||
# Parallelität: gleichzeitige API-Requests. I/O-gebunden → an die Rate-Limits der API
|
||||
# gebunden, NICHT an CPU-Kerne (vhost mit 4 Kernen wie 24-Kern-Rechner gleich steuerbar).
|
||||
max_concurrency: 8
|
||||
# Adaptiver Router: ausgefallene/rate-limited Modelle nach so vielen Fehlern für die
|
||||
# Cooldown-Dauer überspringen; das schnellste gesunde Modell zuerst.
|
||||
breaker_failure_threshold: 2
|
||||
breaker_cooldown_seconds: 120
|
||||
refresh_models: true # tote Modell-Slugs failsafe über OpenRouter /models aussortieren
|
||||
# Verhalten, wenn Inhalte trotz Retry NICHT geprüft werden konnten (KI nicht erreichbar):
|
||||
# "warn" = sichtbarer Hinweis (Level bleibt unberührt)
|
||||
# "yellow" = zusätzlich Gesamt-Level auf mindestens Gelb anheben (für hohe Sicherheit)
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ import concurrent.futures
|
|||
import json
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
|
||||
|
|
@ -69,6 +70,96 @@ _RESPONSE_SCHEMA = {
|
|||
}
|
||||
|
||||
|
||||
_MODELS_URL = "https://openrouter.ai/api/v1/models"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Adaptive model router (circuit-breaker + latency telemetry + /models pruning)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class ModelRouter:
|
||||
"""Wählt programmatisch das beste Modell aus einer Kette: gesunde zuerst,
|
||||
nach Latenz; rate-limited/ausgefallene Modelle bekommen einen Cooldown
|
||||
(Circuit-Breaker) und werden so lange übersprungen. Tote Slugs werden über
|
||||
den /models-Katalog failsafe aussortiert. Thread-sicher (für Parallelität)."""
|
||||
|
||||
def __init__(self, ai_cfg: dict, state: dict | None = None):
|
||||
self._lock = threading.Lock()
|
||||
self._threshold = ai_cfg.get("breaker_failure_threshold", 2)
|
||||
self._cooldown = ai_cfg.get("breaker_cooldown_seconds", 120)
|
||||
self._refresh = ai_cfg.get("refresh_models", True)
|
||||
state = state or {}
|
||||
self._models: dict = state.get("models", {}) or {}
|
||||
self._catalog: dict = state.get("catalog", {"slugs": [], "fetched_at": None})
|
||||
|
||||
def state(self) -> dict:
|
||||
with self._lock:
|
||||
return {"models": self._models, "catalog": self._catalog}
|
||||
|
||||
def record_success(self, model: str, latency: float) -> None:
|
||||
with self._lock:
|
||||
m = self._models.setdefault(model, {})
|
||||
m["failures"] = 0
|
||||
m["open_until"] = 0
|
||||
prev = m.get("latency")
|
||||
m["latency"] = latency if prev is None else 0.7 * prev + 0.3 * latency
|
||||
|
||||
def record_failure(self, model: str) -> None:
|
||||
with self._lock:
|
||||
m = self._models.setdefault(model, {})
|
||||
m["failures"] = m.get("failures", 0) + 1
|
||||
m["last_failure"] = time.time()
|
||||
if m["failures"] >= self._threshold:
|
||||
m["open_until"] = time.time() + self._cooldown
|
||||
logger.info("KI-Router: Breaker für %s offen (%.0fs Cooldown).",
|
||||
model, self._cooldown)
|
||||
|
||||
def order(self, models: list[str]) -> list[str]:
|
||||
"""Reiht die Kette: geschlossener Breaker zuerst, dann nach Latenz aufsteigend.
|
||||
Tote Slugs werden entfernt. Gibt nie eine leere Liste zurück."""
|
||||
with self._lock:
|
||||
candidates = self._prune_locked(models)
|
||||
now = time.time()
|
||||
|
||||
def sort_key(slug):
|
||||
m = self._models.get(slug, {})
|
||||
is_open = m.get("open_until", 0) > now
|
||||
# Unbekannte Latenz → inf: auf dem Erstlauf bleibt die Config-Reihenfolge
|
||||
# erhalten (free zuerst); ein bewährtes (gemessenes) Modell schlägt ein
|
||||
# noch untestetes. sorted() ist stabil → Ties = Config-Reihenfolge.
|
||||
latency = m.get("latency")
|
||||
latency = latency if latency is not None else float("inf")
|
||||
return (is_open, latency)
|
||||
|
||||
ordered = sorted(candidates, key=sort_key)
|
||||
return ordered or list(models)
|
||||
|
||||
def _prune_locked(self, models: list[str]) -> list[str]:
|
||||
slugs = self._catalog.get("slugs") or []
|
||||
if not slugs:
|
||||
return list(models) # kein Katalog → nicht filtern (failsafe)
|
||||
keep = [m for m in models if m in slugs]
|
||||
return keep or list(models) # nie alles wegfiltern
|
||||
|
||||
def refresh_catalog(self) -> None:
|
||||
"""Einmal pro Lauf /models ziehen (24 h gecacht). Failsafe: blockiert nie."""
|
||||
if not self._refresh:
|
||||
return
|
||||
fetched_at = self._catalog.get("fetched_at")
|
||||
if fetched_at and self._catalog.get("slugs") and (time.time() - fetched_at) < 86400:
|
||||
return # Cache frisch
|
||||
try:
|
||||
resp = requests.get(_MODELS_URL, timeout=15)
|
||||
if resp.status_code == 200:
|
||||
slugs = [m["id"] for m in resp.json().get("data", []) if m.get("id")]
|
||||
if slugs:
|
||||
with self._lock:
|
||||
self._catalog = {"slugs": slugs, "fetched_at": time.time()}
|
||||
logger.info("KI-Router: /models-Katalog aktualisiert (%d Modelle).", len(slugs))
|
||||
except (requests.exceptions.RequestException, ValueError, KeyError) as exc:
|
||||
logger.warning("KI-Router: /models nicht abrufbar (failsafe, ungefiltert): %s", exc)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Public API
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -95,6 +186,10 @@ def run_ai_analysis(cfg: dict, bm, snap: dict, diff: dict) -> dict:
|
|||
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", {})
|
||||
|
|
@ -102,21 +197,29 @@ def run_ai_analysis(cfg: dict, bm, snap: dict, diff: dict) -> dict:
|
|||
# Ledger wird zwischengesichert (nach Textphase + im finally), damit bereits
|
||||
# berechnete Verdikte einen späteren Hang/Abbruch überleben und nicht erneut
|
||||
# bezahlt werden müssen.
|
||||
try:
|
||||
if ai_cfg.get("text", {}).get("enabled", True):
|
||||
dirty |= _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result)
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
|
||||
try:
|
||||
if ai_cfg.get("text", {}).get("enabled", True):
|
||||
dirty |= _analyze_text(cfg, ai_cfg, api_key, snap, diff,
|
||||
entries, result, router, executor)
|
||||
if dirty:
|
||||
bm.save_ai_ledger(ledger)
|
||||
if ai_cfg.get("image", {}).get("enabled", True):
|
||||
dirty |= _analyze_images(cfg, ai_cfg, api_key, snap, diff,
|
||||
entries, result, router, executor)
|
||||
# Audio/Video: vorbereitet, default aus (siehe _collect_audio/video_candidates).
|
||||
finally:
|
||||
if dirty:
|
||||
bm.save_ai_ledger(ledger)
|
||||
if ai_cfg.get("image", {}).get("enabled", True):
|
||||
dirty |= _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result)
|
||||
# Audio/Video: vorbereitet, default aus (siehe _collect_audio/video_candidates).
|
||||
finally:
|
||||
if dirty:
|
||||
bm.save_ai_ledger(ledger)
|
||||
except Exception as exc: # pragma: no cover - Schutzschirm, darf Scan nie brechen
|
||||
logger.error("KI-Analyse mit unerwartetem Fehler abgebrochen: %s", exc)
|
||||
result["skipped"] = f"interner Fehler: {exc}"
|
||||
|
||||
# Deterministische Report-Reihenfolge (unabhängig von der Completion-Reihenfolge)
|
||||
result["findings"].sort(key=lambda f: (f.get("kind", ""), f.get("url", ""), f.get("asset_url", "")))
|
||||
result["unchecked"].sort(key=lambda u: (u.get("kind", ""), u.get("url", "")))
|
||||
bm.save_ai_router_state(router.state())
|
||||
|
||||
return result
|
||||
|
||||
|
||||
|
|
@ -182,8 +285,8 @@ def score_ai_findings(ai_result: dict, cfg: dict) -> dict:
|
|||
# Text analysis
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result) -> bool:
|
||||
"""Prüft neue/geänderte Seitentexte. Returns True wenn Ledger verändert."""
|
||||
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)
|
||||
|
|
@ -201,28 +304,48 @@ def _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result) -> bool:
|
|||
candidates.append((url, text, _text_fingerprint(text, cfg)))
|
||||
candidates.sort(key=lambda c: c[0] not in changed) # changed first (False < True)
|
||||
|
||||
budget = text_cfg.get("max_pages_per_scan", 20)
|
||||
dirty = False
|
||||
# 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)
|
||||
|
||||
budget = text_cfg.get("max_pages_per_scan", 20)
|
||||
tasks, miss_fps = [], []
|
||||
for fp in order_fps:
|
||||
entry = entries.get(fp)
|
||||
if entry is None:
|
||||
if budget <= 0:
|
||||
continue # Kostendeckel erreicht — nächster Scan holt den Rest nach
|
||||
classified = _classify_text(text, site_context, models, ai_cfg, api_key)
|
||||
if classified is None:
|
||||
# KI nicht erreichbar → Inhalt bleibt UNGEPRÜFT (nicht still als clean werten).
|
||||
result["unchecked"].append({"kind": "text", "url": url})
|
||||
continue
|
||||
verdict, used_model = classified
|
||||
entry = _make_entry("text", url, verdict, used_model)
|
||||
entries[fp] = entry
|
||||
result["api_calls"] += 1
|
||||
if entry is not None:
|
||||
for url in fp_urls[fp]:
|
||||
result["cache_hits"] += 1
|
||||
_maybe_finding(result, entry, fp, url, ai_cfg)
|
||||
elif budget > 0:
|
||||
miss_fps.append(fp)
|
||||
tasks.append((fp, _text_task(fp_text[fp], site_context, models, ai_cfg, api_key, router)))
|
||||
budget -= 1
|
||||
dirty = True
|
||||
else:
|
||||
result["cache_hits"] += 1
|
||||
_maybe_finding(result, entry, fp, url, ai_cfg)
|
||||
# else: Kostendeckel erreicht — nächster Scan holt den Rest nach
|
||||
|
||||
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
|
||||
|
||||
|
||||
|
|
@ -230,8 +353,8 @@ def _analyze_text(cfg, ai_cfg, api_key, snap, diff, entries, result) -> bool:
|
|||
# Image analysis
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result) -> bool:
|
||||
"""Prüft neue/geänderte Bilder. Returns True wenn Ledger verändert."""
|
||||
def _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result, router, executor) -> bool:
|
||||
"""Prüft neue/geänderte Bilder (parallel). Returns True wenn Ledger verändert."""
|
||||
img_cfg = ai_cfg.get("image", {})
|
||||
models = _models_for(img_cfg)
|
||||
site_context = ai_cfg.get("site_context", "")
|
||||
|
|
@ -264,29 +387,47 @@ def _analyze_images(cfg, ai_cfg, api_key, snap, diff, entries, result) -> bool:
|
|||
url_to_page = _image_url_sources(snap)
|
||||
hashes = fetch_asset_hashes(fetch_urls, timeout=cfg.get("request_timeout", 15))
|
||||
|
||||
dirty = False
|
||||
# Nach Bild-Hash (Bytes) gruppieren: identische Bilder werden nur einmal klassifiziert.
|
||||
fp_assets: dict[str, list[str]] = {}
|
||||
order_fps: list[str] = []
|
||||
for url in fetch_urls:
|
||||
h = hashes.get(url, {})
|
||||
fp = h.get("sha256")
|
||||
fp = hashes.get(url, {}).get("sha256")
|
||||
if not fp:
|
||||
continue # nicht abrufbar — überspringen
|
||||
result["checked"] += 1
|
||||
page_url = url_to_page.get(url, url)
|
||||
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 None:
|
||||
classified = _classify_image(url, site_context, models, ai_cfg, api_key)
|
||||
if classified is None:
|
||||
# KI nicht erreichbar → Bild bleibt UNGEPRÜFT (nicht still als clean werten).
|
||||
result["unchecked"].append({"kind": "image", "url": page_url, "asset_url": url})
|
||||
continue
|
||||
verdict, used_model = classified
|
||||
entry = _make_entry("image", url, verdict, used_model)
|
||||
entries[fp] = entry
|
||||
result["api_calls"] += 1
|
||||
dirty = True
|
||||
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:
|
||||
result["cache_hits"] += 1
|
||||
_maybe_finding(result, entry, fp, page_url, ai_cfg, asset_url=url)
|
||||
miss_fps.append(fp)
|
||||
tasks.append((fp, _image_task(fp_assets[fp][0], site_context, models, ai_cfg, api_key, router)))
|
||||
|
||||
classified = _classify_batch(tasks, executor)
|
||||
dirty = False
|
||||
for fp in miss_fps:
|
||||
res = classified.get(fp)
|
||||
rep_url = fp_assets[fp][0]
|
||||
if res is None:
|
||||
for url in fp_assets[fp]:
|
||||
result["unchecked"].append(
|
||||
{"kind": "image", "url": url_to_page.get(url, url), "asset_url": url})
|
||||
continue
|
||||
verdict, used_model = res
|
||||
entry = _make_entry("image", rep_url, verdict, used_model)
|
||||
entries[fp] = entry
|
||||
result["api_calls"] += 1
|
||||
dirty = True
|
||||
for url in fp_assets[fp]:
|
||||
_maybe_finding(result, entry, fp, url_to_page.get(url, url), ai_cfg, asset_url=url)
|
||||
return dirty
|
||||
|
||||
|
||||
|
|
@ -418,27 +559,25 @@ def _models_for(modality_cfg: dict) -> list[str]:
|
|||
return [single] if single else []
|
||||
|
||||
|
||||
def _classify(messages, models, ai_cfg, api_key) -> tuple[dict, str] | None:
|
||||
"""Versucht die Modell-Kette der Reihe nach (zweifache Eskalation).
|
||||
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.
|
||||
Eskaliert auch bei Langsamkeit (Zeitlimit attempt_timeout je Versuch).
|
||||
Bei Komplettausfall der ganzen Kette wird bis max_retries wiederholt
|
||||
(mit Backoff) — fängt transiente Aussetzer ab, bevor Inhalt ungeprüft bleibt."""
|
||||
timeout = ai_cfg.get("attempt_timeout", 30)
|
||||
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):
|
||||
for i, model in enumerate(models):
|
||||
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:
|
||||
if i > 0 or attempt > 0:
|
||||
logger.info("KI-Verdikt von Stufe %d (%s)%s.", i + 1, model,
|
||||
f" nach Wiederholung {attempt}" if attempt else "")
|
||||
router.record_success(model, time.time() - t0)
|
||||
return verdict, model
|
||||
if i + 1 < len(models):
|
||||
logger.info("KI-Modell Stufe %d (%s) erfolglos — eskaliere zu Stufe %d.",
|
||||
i + 1, model, i + 2)
|
||||
router.record_failure(model)
|
||||
if attempt < max_retries:
|
||||
logger.info("KI-Kette komplett erfolglos — Wiederholung %d/%d in %.1fs.",
|
||||
attempt + 1, max_retries, backoff)
|
||||
|
|
@ -446,15 +585,15 @@ def _classify(messages, models, ai_cfg, api_key) -> tuple[dict, str] | None:
|
|||
return None
|
||||
|
||||
|
||||
def _classify_text(text, site_context, models, ai_cfg, api_key) -> tuple[dict, str] | 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)
|
||||
return _classify(messages, models, ai_cfg, api_key, router)
|
||||
|
||||
|
||||
def _classify_image(image_url, site_context, models, ai_cfg, api_key) -> tuple[dict, str] | None:
|
||||
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": [
|
||||
|
|
@ -463,7 +602,34 @@ def _classify_image(image_url, site_context, models, ai_cfg, api_key) -> tuple[d
|
|||
{"type": "image_url", "image_url": {"url": image_url}},
|
||||
]},
|
||||
]
|
||||
return _classify(messages, models, ai_cfg, api_key)
|
||||
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 _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(headers, payload, timeout):
|
||||
|
|
|
|||
|
|
@ -264,6 +264,34 @@ class BaselineManager:
|
|||
json.dumps(ledger, indent=2, ensure_ascii=False), encoding="utf-8"
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# AI model-router state (circuit-breaker telemetry + /models catalog)
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
@property
|
||||
def _ai_router_state_path(self) -> Path:
|
||||
return self.data_dir / "ai_router_state.json"
|
||||
|
||||
def load_ai_router_state(self) -> dict:
|
||||
"""Load the model-router state. Returns a safe default if none/corrupt."""
|
||||
default = {"models": {}, "catalog": {"slugs": [], "fetched_at": None}}
|
||||
if not self._ai_router_state_path.exists():
|
||||
return default
|
||||
try:
|
||||
data = json.loads(self._ai_router_state_path.read_text(encoding="utf-8"))
|
||||
except (json.JSONDecodeError, ValueError):
|
||||
return default
|
||||
data.setdefault("models", {})
|
||||
data.setdefault("catalog", {"slugs": [], "fetched_at": None})
|
||||
return data
|
||||
|
||||
def save_ai_router_state(self, state: dict) -> None:
|
||||
"""Persist the model-router state."""
|
||||
self.data_dir.mkdir(parents=True, exist_ok=True)
|
||||
self._ai_router_state_path.write_text(
|
||||
json.dumps(state, indent=2, ensure_ascii=False), encoding="utf-8"
|
||||
)
|
||||
|
||||
def dismiss_ai_entries(self, fingerprints: list[str] | None = None) -> int:
|
||||
"""
|
||||
Mark ledger entries as dismissed (false-positive acknowledgement).
|
||||
|
|
|
|||
|
|
@ -96,6 +96,14 @@ DEFAULT_CONFIG: dict = {
|
|||
"attempt_timeout": 20, # HARTES Wanduhr-Limit je Modell-Versuch → Eskalation
|
||||
"max_retries": 1, # Wiederholungen der ganzen Modell-Kette bei Komplettausfall
|
||||
"retry_backoff_seconds": 2.0,
|
||||
# Parallelität: gleichzeitige API-Requests. I/O-gebunden → an die Rate-Limits der
|
||||
# API gebunden, NICHT an CPU-Kerne (vhost wie 24-Kern-Maschine gleich steuerbar).
|
||||
"max_concurrency": 8,
|
||||
# Adaptiver Router: ausgefallene/rate-limited Modelle werden nach so vielen Fehlern
|
||||
# für die Cooldown-Dauer übersprungen; das schnellste gesunde Modell zuerst.
|
||||
"breaker_failure_threshold": 2,
|
||||
"breaker_cooldown_seconds": 120,
|
||||
"refresh_models": True, # tote Slugs failsafe über OpenRouter /models aussortieren
|
||||
# Verhalten, wenn Inhalte trotz Retry NICHT geprüft werden konnten (KI nicht erreichbar):
|
||||
# "warn" = sichtbarer Hinweis in Report/Terminal/E-Mail (Level bleibt unberührt)
|
||||
# "yellow" = zusätzlich Gesamt-Level auf mindestens Gelb anheben
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ from scanner.ai_analyzer import (
|
|||
_text_fingerprint,
|
||||
_models_for,
|
||||
_openrouter_chat,
|
||||
ModelRouter,
|
||||
)
|
||||
from scanner.baseline import BaselineManager
|
||||
from scanner.config import DEFAULT_CONFIG
|
||||
|
|
@ -26,6 +27,7 @@ def _cfg(**ai_overrides) -> dict:
|
|||
cfg["ai_analysis"]["enabled"] = True
|
||||
cfg["ai_analysis"]["site_context"] = "Heimat- und Vereinswebsite über Bergbau"
|
||||
cfg["ai_analysis"]["image"]["enabled"] = False # Text-Tests: Bild aus
|
||||
cfg["ai_analysis"]["refresh_models"] = False # Tests: kein echter /models-Netzcall
|
||||
cfg["ai_analysis"].update(ai_overrides)
|
||||
return cfg
|
||||
|
||||
|
|
@ -307,6 +309,109 @@ class TestRedGate:
|
|||
assert out["level"] == "yellow"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Adaptive model router: circuit-breaker + latency ordering + /models pruning
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _router(catalog_slugs=None):
|
||||
cfg = {"breaker_failure_threshold": 2, "breaker_cooldown_seconds": 120, "refresh_models": False}
|
||||
state = {"models": {}, "catalog": {"slugs": catalog_slugs or [], "fetched_at": None}}
|
||||
return ModelRouter(cfg, state)
|
||||
|
||||
|
||||
class TestModelRouter:
|
||||
def test_failure_below_threshold_keeps_model_healthy(self):
|
||||
r = _router()
|
||||
r.record_failure("m")
|
||||
assert r.order(["m", "n"])[0] in ("m", "n") # nicht hart aussortiert
|
||||
|
||||
def test_breaker_opens_after_threshold_and_sorts_last(self):
|
||||
r = _router()
|
||||
r.record_failure("bad"); r.record_failure("bad") # erreicht threshold 2
|
||||
assert r.order(["bad", "good"]) == ["good", "bad"]
|
||||
|
||||
def test_success_records_latency_faster_first(self):
|
||||
r = _router()
|
||||
r.record_success("slow", 2.0)
|
||||
r.record_success("fast", 0.3)
|
||||
assert r.order(["slow", "fast"]) == ["fast", "slow"]
|
||||
|
||||
def test_success_resets_open_breaker(self):
|
||||
r = _router()
|
||||
r.record_failure("m"); r.record_failure("m")
|
||||
r.record_success("m", 0.5)
|
||||
assert r.order(["m", "n"])[0] == "m" # wieder gesund (Latenz 0.5 < unbekannt? egal: nicht open)
|
||||
|
||||
def test_cooldown_expiry_makes_model_healthy_again(self):
|
||||
r = _router()
|
||||
r.record_failure("m"); r.record_failure("m")
|
||||
r._models["m"]["open_until"] = time.time() - 1 # Cooldown abgelaufen
|
||||
assert r.order(["m", "n"])[-1] != "m" or r.order(["m"]) == ["m"]
|
||||
|
||||
def test_prune_drops_dead_slug(self):
|
||||
r = _router(catalog_slugs=["alive"])
|
||||
assert r.order(["alive", "dead-404"]) == ["alive"]
|
||||
|
||||
def test_prune_failsafe_without_catalog(self):
|
||||
r = _router(catalog_slugs=[]) # kein Katalog → nicht filtern
|
||||
assert set(r.order(["a", "b"])) == {"a", "b"}
|
||||
|
||||
def test_order_never_empty_even_if_all_pruned(self):
|
||||
r = _router(catalog_slugs=["other"])
|
||||
assert r.order(["x", "y"]) # keiner im Katalog → trotzdem nicht leer
|
||||
|
||||
def test_refresh_catalog_failsafe_on_network_error(self):
|
||||
r = _router()
|
||||
with patch("scanner.ai_analyzer.requests.get", side_effect=Exception("net down")):
|
||||
r.refresh_catalog() # darf nicht werfen
|
||||
assert r.order(["a"]) == ["a"]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Parallel analysis correctness
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestParallel:
|
||||
def test_all_pages_classified_in_parallel(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("OPENROUTER_API_KEY", "test")
|
||||
bm = BaselineManager(tmp_path)
|
||||
pages = {f"https://x.de/p{i}": {"url": f"https://x.de/p{i}", "status": 200,
|
||||
"text": _LONG + str(i), "links": {"img": []}}
|
||||
for i in range(6)}
|
||||
with patch("scanner.ai_analyzer._openrouter_chat", return_value=_verdict()):
|
||||
res = run_ai_analysis(_cfg(), bm, {"pages": pages}, {})
|
||||
assert res["api_calls"] == 6
|
||||
assert len(bm.load_ai_ledger()["entries"]) == 6
|
||||
assert res["skipped"] is None
|
||||
|
||||
def test_duplicate_text_classified_once_findings_for_each_url(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("OPENROUTER_API_KEY", "test")
|
||||
bm = BaselineManager(tmp_path)
|
||||
same = _LONG + " identisch"
|
||||
pages = {
|
||||
"https://x.de/a": {"url": "https://x.de/a", "status": 200, "text": same, "links": {"img": []}},
|
||||
"https://x.de/b": {"url": "https://x.de/b", "status": 200, "text": same, "links": {"img": []}},
|
||||
}
|
||||
v = _verdict("hidden_spam", "high", 0.9)
|
||||
with patch("scanner.ai_analyzer._openrouter_chat", return_value=v):
|
||||
res = run_ai_analysis(_cfg(), bm, {"pages": pages}, {})
|
||||
assert res["api_calls"] == 1 # identischer Text → nur 1 Call
|
||||
assert len(res["findings"]) == 2 # aber Fund für beide URLs
|
||||
assert {f["url"] for f in res["findings"]} == {"https://x.de/a", "https://x.de/b"}
|
||||
|
||||
def test_findings_sorted_deterministically(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("OPENROUTER_API_KEY", "test")
|
||||
bm = BaselineManager(tmp_path)
|
||||
pages = {f"https://x.de/{c}": {"url": f"https://x.de/{c}", "status": 200,
|
||||
"text": _LONG + c, "links": {"img": []}}
|
||||
for c in "cab"}
|
||||
v = _verdict("off_topic_commercial", "high", 0.9)
|
||||
with patch("scanner.ai_analyzer._openrouter_chat", return_value=v):
|
||||
res = run_ai_analysis(_cfg(), bm, {"pages": pages}, {})
|
||||
urls = [f["url"] for f in res["findings"]]
|
||||
assert urls == sorted(urls) # deterministisch nach URL sortiert
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Robustness: hard wall-clock deadline + resilient ledger save
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue