From 125d80e7f9b499b53677cf07579f5ad7d7771438 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dieter=20Schl=C3=BCter?= Date: Sat, 13 Jun 2026 02:51:30 +0200 Subject: [PATCH] feat: adaptiver Modell-Router (Circuit-Breaker) + parallele KI-Analyse MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- BEDIENUNGSANLEITUNG.md | 28 ++++ README.md | 7 + meine-seite.de/config.yaml | 8 + scanner/ai_analyzer.py | 296 +++++++++++++++++++++++++++++-------- scanner/baseline.py | 28 ++++ scanner/config.py | 8 + tests/test_ai_analyzer.py | 105 +++++++++++++ 7 files changed, 415 insertions(+), 65 deletions(-) diff --git a/BEDIENUNGSANLEITUNG.md b/BEDIENUNGSANLEITUNG.md index dcc5acb..c21286d 100644 --- a/BEDIENUNGSANLEITUNG.md +++ b/BEDIENUNGSANLEITUNG.md @@ -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. diff --git a/README.md b/README.md index c7d4b98..095237f 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/meine-seite.de/config.yaml b/meine-seite.de/config.yaml index 5218bbc..a02eea1 100644 --- a/meine-seite.de/config.yaml +++ b/meine-seite.de/config.yaml @@ -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) diff --git a/scanner/ai_analyzer.py b/scanner/ai_analyzer.py index 636426d..7f0edd7 100644 --- a/scanner/ai_analyzer.py +++ b/scanner/ai_analyzer.py @@ -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): diff --git a/scanner/baseline.py b/scanner/baseline.py index dc0d824..22ca7b9 100644 --- a/scanner/baseline.py +++ b/scanner/baseline.py @@ -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). diff --git a/scanner/config.py b/scanner/config.py index 9ee8f74..921cb66 100644 --- a/scanner/config.py +++ b/scanner/config.py @@ -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 diff --git a/tests/test_ai_analyzer.py b/tests/test_ai_analyzer.py index 3412718..b185a30 100644 --- a/tests/test_ai_analyzer.py +++ b/tests/test_ai_analyzer.py @@ -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 # ---------------------------------------------------------------------------