Rebuild profile router as secure recoverable v2

This commit is contained in:
Mikei386
2026-08-20 13:41:00 +02:00
parent bdd6c08643
commit fd4a3055e8
18 changed files with 1162 additions and 129 deletions
+466 -106
View File
@@ -54,24 +54,39 @@ Nur Python-Standardbibliothek. Logging nach stdout (journald).
from __future__ import annotations
import base64
import binascii
import email
import hashlib
import ipaddress
import json
import logging
import os
import queue
import re
import socket
import subprocess
import sys
import threading
import time
import uuid
import http.client
import urllib.error
import urllib.parse
import urllib.request
from collections import OrderedDict
from email.parser import BytesParser
from email.policy import compat32
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from router_support import (
AuthPolicy,
ConfigurationError,
RuntimeStore,
enforce_artifact_retention,
load_profile_registry,
terminate_recorded_worker,
)
# ---------------------------------------------------------------------------
# Konfiguration (über Umgebungsvariablen, vgl. systemd-Unit)
# ---------------------------------------------------------------------------
@@ -102,6 +117,10 @@ IMAGE_WORKER_LOG = os.environ.get(
IMAGE_START_TIMEOUT = float(os.environ.get("IMAGE_START_TIMEOUT", "120")) # s, Worker-Start
IMAGE_GEN_TIMEOUT = float(os.environ.get("IMAGE_GEN_TIMEOUT", "1800")) # s, pro Bild
IMAGE_VRAM_FREE_TIMEOUT = float(os.environ.get("IMAGE_VRAM_FREE_TIMEOUT", "90")) # s, VRAM-Abgabe
IMAGE_RETENTION_FILES = int(os.environ.get("IMAGE_RETENTION_FILES", "100"))
IMAGE_RETENTION_BYTES = int(os.environ.get(
"IMAGE_RETENTION_BYTES", str(5 * 1024 * 1024 * 1024)))
IMAGE_RETENTION_DAYS = int(os.environ.get("IMAGE_RETENTION_DAYS", "30"))
# --- Vision-Orchestrierung (Q3 "Augen", temporär) ---
# Chat-Requests mit Bild im letzten User-Message lösen einen
@@ -121,6 +140,10 @@ VISION_INFER_TIMEOUT = float(os.environ.get("VISION_INFER_TIMEOUT", "300")) # s
VISION_UNLOAD_TIMEOUT = float(os.environ.get("VISION_UNLOAD_TIMEOUT", "120")) # s
VISION_MAX_TOKENS = int(os.environ.get("VISION_MAX_TOKENS", "4096"))
VISION_CACHE_MAX = int(os.environ.get("VISION_CACHE_MAX", "64"))
VISION_MAX_IMAGE_BYTES = int(os.environ.get(
"VISION_MAX_IMAGE_BYTES", str(20 * 1024 * 1024)))
VISION_ALLOW_REMOTE_URLS = os.environ.get(
"VISION_ALLOW_REMOTE_URLS", "false").lower() in {"1", "true", "yes"}
VISION_LOG = os.environ.get(
"VISION_LOG", "/opt/mike-ai/ai-profile-router/vision_server.log")
@@ -165,16 +188,32 @@ MAX_UPLOAD_SIZE = int(os.environ.get("MAX_UPLOAD_SIZE", 50 * 1024 * 1024))
CHAT_WAIT_TIMEOUT = float(os.environ.get("CHAT_WAIT_TIMEOUT", "300")) # s, max. Warten
CHAT_DRAIN_TIMEOUT = float(os.environ.get("CHAT_DRAIN_TIMEOUT", "60")) # s, max. Warten auf aktive Chats
PROFILES = {"fast": 73728, "medium": 94208, "long": 131072}
RUNTIME_STATE_FILE = os.environ.get(
"ROUTER_STATE_FILE", "/var/lib/mike-ai-profile-router/state.json")
PROFILE_REGISTRY_FILE = os.environ.get(
"ROUTER_PROFILES_FILE", "/etc/mike-ai/router-profiles.json")
ALLOW_LEGACY_GET_SWITCH = os.environ.get(
"ALLOW_LEGACY_GET_SWITCH", "false").lower() in {"1", "true", "yes"}
MAX_CONCURRENT_REQUESTS = int(os.environ.get(
"ROUTER_MAX_CONCURRENT_REQUESTS", "16"))
PROFILE_REGISTRY = load_profile_registry(PROFILE_REGISTRY_FILE)
PROFILES = {name: definition["context"]
for name, definition in PROFILE_REGISTRY.items()}
EXPECTED_MODELS = {name: definition.get("model_alias")
for name, definition in PROFILE_REGISTRY.items()}
VIRTUAL_MODELS = {f"qwen-{name}": name for name in PROFILES}
log = logging.getLogger("ai-profile-router")
AUTH: AuthPolicy | None = None
RUNTIME = RuntimeStore(RUNTIME_STATE_FILE)
REQUEST_SLOTS = threading.BoundedSemaphore(max(1, MAX_CONCURRENT_REQUESTS))
# Hop-by-hop-Header, die nicht an Upstream/Client weitergereicht werden.
HOP_BY_HOP = {
"host", "connection", "keep-alive", "proxy-authenticate",
"proxy-authorization", "te", "trailer", "transfer-encoding",
"upgrade", "content-length",
"upgrade", "content-length", "authorization", "x-api-key",
}
@@ -237,15 +276,17 @@ class _State:
gegenseitiger Ausschluss, kein Race zwischen beiden.
avail_lock : schützt qwen_unavailable + active_chats (Chat-Waiting).
"""
lock = threading.Lock() # GPU-/Model-Lock (Profilwechsel + Image + Vision)
switching: str | None = None # Profil, das gerade gewechselt wird
started = time.time()
image = _ImageState()
vision = _VisionState()
# Qwen-Verfügbarkeit für das Chat-Waiting:
qwen_unavailable = False # True, wenn Qwen down/neu geladen wird
active_chats = 0 # Anzahl laufender Chat-Requests
avail_lock = threading.Lock() # schützt die beiden Felder oben
def __init__(self) -> None:
# RLock erlaubt atomare Abläufe aus Profilwahl + Vision + Chat-Lease,
# während die darunterliegenden Funktionen denselben Lock verwenden.
self.lock = threading.RLock()
self.switching: str | None = None
self.started = time.time()
self.image = _ImageState()
self.vision = _VisionState()
self.qwen_unavailable = True
self.active_chats = 0
self.avail_lock = threading.Lock()
STATE = _State()
@@ -265,9 +306,9 @@ def _wait_chats_drained(timeout: float | None = None) -> None:
return
n = STATE.active_chats
if time.monotonic() > deadline:
log.warning("Chat-Drain-Timeout nach %.0f s (%d aktive Chats) – "
"fahre trotzdem fort", timeout, n)
return
raise RuntimeError(
f"Profil-/GPU-Wechsel nach {timeout:.0f} s abgebrochen: "
f"noch {n} aktive Chat-Anfrage(n)")
time.sleep(0.5)
@@ -426,17 +467,24 @@ def _read(path: str) -> str:
def current_profile() -> str | None:
"""Aktives Profil, ermittelt durch Vergleich der override.conf."""
"""Aktives Profil anhand semantischer Werte der override.conf.
Kommentare, Leerraum oder die Reihenfolge anderer llama.cpp-Optionen
beeinflussen die Erkennung nicht mehr.
"""
try:
override = _read(os.path.join(PROFILE_DIR, "override.conf"))
except OSError:
return None
for name in PROFILES:
try:
ref = _read(os.path.join(PROFILE_DIR, f"profile-{name}.conf.disabled"))
except OSError:
continue
if override == ref:
ctx_match = re.search(r"(?:^|\s)--ctx-size\s+(\d+)(?:\s|$)", override)
alias_match = re.search(r"(?:^|\s)--alias\s+([^\s]+)", override)
if not ctx_match:
return None
ctx = int(ctx_match.group(1))
alias = alias_match.group(1) if alias_match else None
for name, expected_ctx in PROFILES.items():
expected_alias = EXPECTED_MODELS.get(name)
if ctx == expected_ctx and (not expected_alias or alias == expected_alias):
return name
return None
@@ -444,20 +492,33 @@ def current_profile() -> str | None:
def _wait_ready(profile: str, deadline: float) -> None:
"""Wartet, bis llama.cpp das Profil geladen hat (Modell + ctx)."""
expected_ctx = PROFILES[profile]
expected_model = EXPECTED_MODELS.get(profile)
while True:
status = upstream_status()
if (status["reachable"] and status.get("model")
and status.get("ctx") == expected_ctx):
and status.get("ctx") == expected_ctx
and (not expected_model or status.get("model") == expected_model)):
log.info("llama.cpp bereit: Profil=%s Modell=%s ctx=%s",
profile, status.get("model"), status.get("ctx"))
return
if time.monotonic() > deadline:
raise RuntimeError(
f"llama.cpp nach {SWITCH_TIMEOUT:.0f} s nicht bereit "
f"(erwartet ctx {expected_ctx}, aktuell: {status.get('ctx')})")
f"(erwartet Modell {expected_model or '*'} / ctx "
f"{expected_ctx}, aktuell: {status.get('model')} / "
f"{status.get('ctx')})")
time.sleep(POLL_INTERVAL)
def _profile_is_ready(profile: str, status: dict | None = None) -> bool:
status = status or upstream_status()
expected_model = EXPECTED_MODELS.get(profile)
return bool(status.get("reachable") and status.get("model")
and status.get("ctx") == PROFILES[profile]
and (not expected_model
or status.get("model") == expected_model))
def switch_profile(profile: str, implicit: bool = False) -> None:
"""Stellt sicher, dass das Profil aktiv ist, und wartet bis es geladen ist.
@@ -479,8 +540,7 @@ def switch_profile(profile: str, implicit: bool = False) -> None:
try:
cur = current_profile()
up = upstream_status()
ready = (up["reachable"] and up.get("model")
and up.get("ctx") == PROFILES[profile])
ready = _profile_is_ready(profile, up)
if cur == profile and ready:
log.info("Profil %s ist bereits aktiv", profile)
return
@@ -510,19 +570,26 @@ def switch_profile(profile: str, implicit: bool = False) -> None:
if out:
log.info("llama-profile: %s", out[-500:])
if proc.returncode != 0:
# whiptail bricht das Skript ohne TTY ab – der Wechsel
# selbst (cp + systemctl restart) ist dann erledigt.
log.warning("llama-profile Exit-Code %d (ohne TTY "
"erwartet)", proc.returncode)
raise RuntimeError(
f"llama-profile fehlgeschlagen (Exit-Code "
f"{proc.returncode}): {out[-500:]}")
except subprocess.TimeoutExpired:
log.error("llama-profile hat 120 s überschritten")
raise RuntimeError("llama-profile hat 120 s überschritten")
if current_profile() != profile:
raise RuntimeError(
f"Profildatei wurde nicht gesetzt (erwartet: {profile})")
log.info("Warte, bis llama.cpp das Profil geladen hat ...")
_wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT)
RUNTIME.save(last_profile=profile, phase="idle")
finally:
_set_qwen_unavailable(False)
# Nach einem fehlgeschlagenen Skript/Timeout darf der Router
# Qwen nicht blind freigeben. Nur ein semantisch verifiziertes
# Profil (Alias + Kontext) wird wieder als verfügbar markiert.
active = current_profile()
up_after = upstream_status()
available = bool(active in PROFILES
and _profile_is_ready(active, up_after))
_set_qwen_unavailable(not available)
finally:
STATE.switching = None
@@ -539,6 +606,7 @@ class _Worker:
self.model_loaded = False
self._queue: queue.Queue[dict] = queue.Queue()
self._reader: threading.Thread | None = None
self._logf = None
def alive(self) -> bool:
return self.proc is not None and self.proc.poll() is None
@@ -547,15 +615,18 @@ class _Worker:
if self.alive():
return
log.info("starte Bild-Worker: %s %s", IMAGE_PYTHON, IMAGE_WORKER)
logf = open(IMAGE_WORKER_LOG, "ab")
self._logf = open(IMAGE_WORKER_LOG, "ab")
self.proc = subprocess.Popen(
[IMAGE_PYTHON, IMAGE_WORKER],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=logf,
stderr=self._logf,
text=True,
bufsize=1,
start_new_session=True,
)
RUNTIME.save(worker="image", worker_pid=self.proc.pid,
phase="loading-image")
self._reader = threading.Thread(target=self._read_loop, daemon=True)
self._reader.start()
try:
@@ -599,8 +670,19 @@ class _Worker:
self.proc.wait(timeout=10)
except subprocess.TimeoutExpired:
self.proc.kill()
try:
self.proc.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
if self._logf is not None:
try:
self._logf.close()
except OSError:
pass
self._logf = None
self.proc = None
self.model_loaded = False
RUNTIME.clear_worker("image")
def _worker() -> _Worker:
@@ -670,13 +752,19 @@ def _restore_qwen(profile: str) -> None:
"""Startet llama.cpp mit dem gemerkten Profil und wartet auf Readiness."""
log.info("stelle Qwen-Profil %s wieder her ...", profile)
try:
subprocess.run([SYSTEMCTL_BIN, "start", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=120)
proc = subprocess.run([SYSTEMCTL_BIN, "start", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=120)
if proc.returncode != 0:
out = proc.stdout.decode(errors="replace").strip()
raise RuntimeError(
f"systemctl start {LLAMA_SERVICE} fehlgeschlagen "
f"(Exit {proc.returncode}): {out[-500:]}")
except subprocess.TimeoutExpired:
log.error("systemctl start hat 120 s überschritten")
_wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT)
RUNTIME.save(last_profile=profile, phase="idle")
def generate_image(prompt: str, width: int, height: int, steps: int,
@@ -700,6 +788,7 @@ def generate_image(prompt: str, width: int, height: int, steps: int,
os.makedirs(IMAGE_DIR, exist_ok=True)
results: list[str] = []
warning: str | None = None
img.last_error = None
# Qwen wird gestoppt → für Chats nicht verfügbar (die warten).
_set_qwen_unavailable(True)
try:
@@ -707,10 +796,16 @@ def generate_image(prompt: str, width: int, height: int, steps: int,
# 1) Qwen stoppen (VRAM freigeben).
img.phase = "stopping-qwen"
subprocess.run([SYSTEMCTL_BIN, "stop", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=120)
RUNTIME.save(last_profile=profile, phase=img.phase)
proc = subprocess.run([SYSTEMCTL_BIN, "stop", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT, timeout=120)
if proc.returncode != 0:
out = proc.stdout.decode(errors="replace").strip()
raise RuntimeError(
f"systemctl stop {LLAMA_SERVICE} fehlgeschlagen "
f"(Exit {proc.returncode}): {out[-500:]}")
_wait_upstream_down(time.monotonic() + 60)
# 2) Worker starten (Modell wird beim ersten generate geladen).
@@ -763,6 +858,12 @@ def generate_image(prompt: str, width: int, height: int, steps: int,
log.info("Bild %d/%d: %s (%.1f s)", i + 1, n, filename,
resp.get("seconds", 0))
removed = enforce_artifact_retention(
IMAGE_DIR, IMAGE_RETENTION_FILES, IMAGE_RETENTION_BYTES,
IMAGE_RETENTION_DAYS, protected=results)
if removed:
log.info("Bild-Retention: %d alte Bilder entfernt", len(removed))
# 4) Worker vollständig beenden (VRAM + CUDA-Kontext freigeben).
img.phase = "unloading-image"
worker.stop()
@@ -804,8 +905,10 @@ def _image_filename_ok(name: str) -> bool:
VISION_ANALYST_PROMPT = (
"Du bist ein reiner Bild- und Screenshot-Analyst. Du beantwortest die "
"Benutzerfrage NICHT selbst. Du extrahierst aus dem Bild alle "
"Informationen, die für die Beantwortung relevant sein könnten.\n\n"
"Benutzerfrage NICHT selbst. Erstelle eine umfassende, von der aktuellen "
"Frage unabhängige Bestandsaufnahme. Lasse keine sichtbaren Details aus, "
"damit auch spätere Folgefragen allein mit dieser Analyse beantwortet "
"werden können.\n\n"
"Antworte NUR mit einer strukturierten Analyse in dieser Form:\n"
"1. SIEHTBARER TEXT: alle Texte wörtlich und vollständig, mit Anordnung\n"
"2. UI-ELEMENTE: Felder, Buttons, Menüs, Tabs, Dropdowns, Checkboxen – "
@@ -815,8 +918,7 @@ VISION_ANALYST_PROMPT = (
"links/rechts, Reihenfolge)\n"
"5. ZUSTÄNDE: Statusanzeigen, Farben (rot/grün/gelb), Ladezustände\n"
"6. OBJEKTE: relevante Objekte und Beziehungen zwischen Elementen\n"
"7. WEITERES: alles Weitere, was für die Benutzerfrage relevant sein "
"könnte\n\n"
"7. WEITERES: alle übrigen sichtbaren Details\n\n"
"Regeln: Bei Screenshots hat Text- und UI-Genauigkeit Vorrang vor "
"schöner Beschreibung. Keine Interpretation, keine Vermutungen – nur "
"was sichtbar ist. Unleserliches als [unleserlich] markieren."
@@ -869,7 +971,10 @@ class _VisionServer:
self._logf = open(VISION_LOG, "ab")
self.proc = subprocess.Popen(
cmd, stdin=subprocess.DEVNULL,
stdout=self._logf, stderr=subprocess.STDOUT)
stdout=self._logf, stderr=subprocess.STDOUT,
start_new_session=True)
RUNTIME.save(worker="vision", worker_pid=self.proc.pid,
phase="loading-vision")
log.info("Vision-Server gestartet (PID %d, Port %d, ctx %d)",
self.proc.pid, VISION_PORT, VISION_CTX)
@@ -917,6 +1022,7 @@ class _VisionServer:
pass
self._logf = None
self.proc = None
RUNTIME.clear_worker("vision")
def _extract_last_user_image(data: dict) -> tuple[str | None, str]:
@@ -955,6 +1061,74 @@ def _extract_last_user_image(data: dict) -> tuple[str | None, str]:
return None, ""
_VISION_DATA_TYPES = {
"image/jpeg", "image/png", "image/webp", "image/gif",
}
def _normalize_vision_image(image_url: str) -> str:
"""Validate an image and return a bounded data URL for llama.cpp.
Remote fetches are off by default. If explicitly enabled, the router
downloads the file itself after rejecting local/private destinations and
hands llama.cpp a data URL. The model server therefore never receives an
arbitrary user-controlled URL (SSRF protection).
"""
if image_url.startswith("data:"):
header, separator, payload = image_url.partition(",")
match = re.fullmatch(
r"data:([a-zA-Z0-9.+-]+/[a-zA-Z0-9.+-]+);base64", header)
if not separator or not match or match.group(1).lower() not in _VISION_DATA_TYPES:
raise ValueError("ungültige oder nicht unterstützte Bild-data-URL")
try:
raw = base64.b64decode(payload, validate=True)
except (ValueError, binascii.Error):
raise ValueError("ungültige Base64-Bilddaten") from None
if not raw or len(raw) > VISION_MAX_IMAGE_BYTES:
raise ValueError(
f"Bildgröße außerhalb des Limits (max {VISION_MAX_IMAGE_BYTES} Bytes)")
return image_url
parsed = urllib.parse.urlsplit(image_url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname:
raise ValueError("Bild muss eine data-URL oder eine gültige HTTP(S)-URL sein")
if not VISION_ALLOW_REMOTE_URLS:
raise ValueError(
"Remote-Bild-URLs sind deaktiviert; Bild bitte als data-URL hochladen")
try:
addresses = socket.getaddrinfo(
parsed.hostname,
parsed.port or (443 if parsed.scheme == "https" else 80))
except socket.gaierror as exc:
raise ValueError(f"Bild-Host kann nicht aufgelöst werden: {exc}") from None
for address in addresses:
try:
ip = ipaddress.ip_address(address[4][0])
except ValueError:
raise ValueError("Bild-Host liefert eine ungültige Adresse") from None
if not ip.is_global:
raise ValueError("private/lokale Bild-URLs sind nicht erlaubt")
request = urllib.request.Request(
image_url, headers={"User-Agent": "AI-Profile-Router/2.0"})
try:
with urllib.request.urlopen(request, timeout=15) as response:
final_url = urllib.parse.urlsplit(response.geturl())
if final_url.hostname != parsed.hostname:
raise ValueError("Weiterleitungen zu einem anderen Bild-Host sind nicht erlaubt")
content_type = response.headers.get_content_type().lower()
if content_type not in _VISION_DATA_TYPES:
raise ValueError(f"Remote-Inhalt ist kein unterstütztes Bild ({content_type})")
raw = response.read(VISION_MAX_IMAGE_BYTES + 1)
except urllib.error.URLError as exc:
raise ValueError(f"Remote-Bild kann nicht geladen werden: {exc}") from None
if not raw or len(raw) > VISION_MAX_IMAGE_BYTES:
raise ValueError(
f"Bildgröße außerhalb des Limits (max {VISION_MAX_IMAGE_BYTES} Bytes)")
return (f"data:{content_type};base64,"
+ base64.b64encode(raw).decode("ascii"))
def _image_hash(image_url: str) -> str:
"""Stabiler Hash für ein Bild (Base64-Payload oder URL).
@@ -1064,9 +1238,11 @@ def _vision_analyze(image_url: str, question: str) -> str:
{"role": "user", "content": [
{"type": "image_url", "image_url": {"url": image_url}},
{"type": "text",
"text": ("Benutzerfrage (nur zur Orientierung, NICHT "
"beantworten: " + (question or "(keine Frage)")
+ "\n\nErstelle jetzt die strukturierte Vision-Analyse.")},
"text": ("Aktuelle Benutzerfrage nur als zusätzlicher "
"Kontext, ohne die Analyse darauf zu beschränken: "
+ (question or "(keine Frage)")
+ "\n\nErstelle jetzt eine vollständige, "
"fragenunabhängige Vision-Analyse.")},
]},
],
"max_tokens": VISION_MAX_TOKENS,
@@ -1131,6 +1307,7 @@ def _vision_swap(data: dict) -> str:
if not image_url:
raise RuntimeError("kein Bild im Request")
img_hash = _image_hash(image_url)
safe_image_url = _normalize_vision_image(image_url)
vis.last_profile = profile
vis.last_error = None
t_total = time.monotonic()
@@ -1148,11 +1325,17 @@ def _vision_swap(data: dict) -> str:
# 1) Hauptprofil entladen (VRAM freigeben).
vis.phase = "stopping-main"
RUNTIME.save(last_profile=profile, phase=vis.phase)
t0 = time.monotonic()
subprocess.run([SYSTEMCTL_BIN, "stop", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=120)
proc = subprocess.run([SYSTEMCTL_BIN, "stop", LLAMA_SERVICE],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT, timeout=120)
if proc.returncode != 0:
out = proc.stdout.decode(errors="replace").strip()
raise RuntimeError(
f"systemctl stop {LLAMA_SERVICE} fehlgeschlagen "
f"(Exit {proc.returncode}): {out[-500:]}")
_wait_upstream_down(time.monotonic() + 60)
timings["main_unload"] = time.monotonic() - t0
log.info("Vision: Hauptprofil %s entladen (%.1f s)",
@@ -1171,7 +1354,7 @@ def _vision_swap(data: dict) -> str:
vis.phase = "analyzing"
t0 = time.monotonic()
log.info("Vision: Inferenz gestartet (127.0.0.1:%d)", VISION_PORT)
analysis = _vision_analyze(image_url, question)
analysis = _vision_analyze(safe_image_url, question)
timings["vision_infer"] = time.monotonic() - t0
log.info("Vision: Inferenz abgeschlossen (%.1f s, %d Zeichen)",
timings["vision_infer"], len(analysis))
@@ -1243,7 +1426,8 @@ def _vision_swap(data: dict) -> str:
# ---------------------------------------------------------------------------
class Handler(BaseHTTPRequestHandler):
server_version = "AIProfileRouter/1.0"
server_version = "AIProfileRouter/2.0"
sys_version = ""
timeout = 60 # Socket-Timeout für Client-Requests (s)
# ---------- Routing ----------
@@ -1254,13 +1438,51 @@ class Handler(BaseHTTPRequestHandler):
def do_POST(self):
self._route()
def do_PUT(self):
self._route()
def do_PATCH(self):
self._route()
def do_DELETE(self):
self._route()
def do_OPTIONS(self):
self._route()
def _route(self):
path = self.path.split("?", 1)[0]
started = time.monotonic()
slot_acquired = False
try:
if path == "/v1/models" and self.command == "GET":
if path == "/health" and self.command == "GET":
# Liveness: der Routerprozess lebt. Ein absichtlich entladenes
# Qwen (Vision/Bild) darf keinen Restart-Loop auslösen.
self._send_json(200, {"status": "ok", "router": "alive"})
return
if path == "/ready" and self.command == "GET":
up = upstream_status()
active = current_profile()
with STATE.avail_lock:
unavailable = STATE.qwen_unavailable
ready = bool(active in PROFILES and not unavailable
and _profile_is_ready(active, up))
self._send_json(200 if ready else 503, {
"status": "ok" if ready else "degraded",
"router": "alive",
"upstream": "ready" if ready else "unavailable",
})
return
slot_acquired = REQUEST_SLOTS.acquire(blocking=False)
if not slot_acquired:
self._send_error(429, "Router ist ausgelastet; bitte erneut versuchen",
"server_error", "too_many_requests")
return
if not self._authorized():
self._send_auth_required()
elif path == "/v1/models" and self.command == "GET":
self._send_json(200, self._models_payload())
elif path == "/status":
elif path == "/status" and self.command == "GET":
self._send_json(200, self._status_payload())
elif path == "/v1/audio/models" and self.command == "GET":
self._send_json(200, self._audio_models_payload())
@@ -1278,13 +1500,14 @@ class Handler(BaseHTTPRequestHandler):
self._images_list()
elif path.startswith("/images/") and self.command == "GET":
self._image_serve(path[len("/images/"):])
elif path in ("/fast", "/medium", "/long"):
elif (path in ("/fast", "/medium", "/long")
and (self.command == "POST"
or (self.command == "GET" and ALLOW_LEGACY_GET_SWITCH))):
self._switch(path[1:])
elif (self.command == "POST" and path.startswith("/")
and path.count("/") == 1):
# Kommandonamensraum: unbekanntes Profil
self._send_error(400, f"unbekanntes Profil: {path[1:]}",
"invalid_request_error", "invalid_profile")
elif (path in ("/fast", "/medium", "/long")
and self.command == "GET"):
self._send_error(405, "Profilwechsel erfordert POST",
"invalid_request_error", "method_not_allowed")
else:
self._forward()
except BrokenPipeError:
@@ -1293,9 +1516,31 @@ class Handler(BaseHTTPRequestHandler):
log.exception("Fehler bei %s %s", self.command, path)
self._safe_error(500, "interner Router-Fehler")
finally:
if slot_acquired:
REQUEST_SLOTS.release()
log.info("%s %s -> %s in %.3f s", self.command, path,
getattr(self, "_last_code", "-"), time.monotonic() - started)
def _authorized(self) -> bool:
assert AUTH is not None
return AUTH.accepts(self.headers.get("Authorization"),
self.headers.get("X-API-Key"))
def _send_auth_required(self) -> None:
body = json.dumps({"error": {
"message": "gültiger Router-API-Key erforderlich",
"type": "authentication_error",
"code": "invalid_api_key",
}}).encode()
self._last_code = 401
self.send_response(401)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.send_header("WWW-Authenticate", "Bearer")
self.send_header("Connection", "close")
self.end_headers()
self.wfile.write(body)
# ---------- Request-Body-Lesen (Content-Length + chunked) ----------
def _read_body(self) -> bytes:
@@ -1436,7 +1681,8 @@ class Handler(BaseHTTPRequestHandler):
"ctx": up.get("ctx"),
},
"qwen": {
"available": not qwen_unavailable,
"available": (not qwen_unavailable and up["reachable"]
and bool(up.get("model"))),
"active_chats": active_chats,
},
"image": {
@@ -1965,7 +2211,9 @@ class Handler(BaseHTTPRequestHandler):
return
data = None
# Virtuelles Modell? -> Profil sicherstellen, dann Modell ersetzen.
requested_profile: str | None = None
# Virtuelles Modell erkennen. Umschalten und Chat-Lease werden weiter
# unten atomar unter dem zentralen Orchestrierungs-Lock ausgeführt.
if body is not None and self.path.startswith("/v1/"):
try:
data = json.loads(body)
@@ -1973,57 +2221,98 @@ class Handler(BaseHTTPRequestHandler):
data = None
model = data.get("model") if isinstance(data, dict) else None
if isinstance(model, str) and model in VIRTUAL_MODELS:
profile = VIRTUAL_MODELS[model]
try:
switch_profile(profile, implicit=True)
except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error",
"upstream_unavailable")
return
up = upstream_status()
if not up["reachable"] or not up.get("model"):
self._send_error(502, "llama.cpp nicht erreichbar",
"server_error", "upstream_unavailable")
return
data["model"] = up["model"]
body = json.dumps(data).encode()
requested_profile = VIRTUAL_MODELS[model]
elif isinstance(model, str) and model.startswith("qwen-"):
# qwen-* ist der Namensraum des Routers
self._send_error(400, f"unbekanntes virtuelles Modell: {model}",
"invalid_request_error", "unknown_model")
return
# Vision: Bilder im Request → Q3-Vision-Analyse (nur für NEUE,
# noch nicht analysierte Bilder), danach erzeugt das
# (wiederhergestellte) Hauptmodell die Endantwort. Die gesamte
# History wird sanisiert: alle Bild-Parts werden durch ihre
# (gecachten) Vision-Analysen ersetzt, damit das Nicht-Vision-
# Hauptmodell (Fast/Medium/Long) keine Bilddaten bekommt. Das ist
# wichtig, weil Open WebUI bei Folgefragen den ursprünglichen
# multimodalen Verlauf erneut mitsendet.
# Alle modellbezogenen Requests erhalten eine atomare Lease. Damit
# kann kein zweiter Client zwischen Profilwahl und Upstream-Request das
# Modell austauschen. Vision-Vorbereitung gehört zur selben Transaktion.
if isinstance(data, dict) and path == "/v1/chat/completions":
image_url, _ = _extract_last_user_image(data)
if image_url:
if _vision_cached_analysis(_image_hash(image_url)) is None:
# Neues Bild → Vision-Hotswap (Q3 analysiert + cacht).
self.timeout = None # Vision-Swap kann Minuten dauern
try:
_vision_swap(data)
except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error",
"vision_failed")
return
else:
log.info("Vision: Bild bereits analysiert "
"(Cache-Treffer) – kein Hotswap")
if _request_has_image(data):
data = _sanitize_for_main_model(data)
body = json.dumps(data).encode()
log.info("Vision: finale Hauptmodell-Inferenz gestartet")
self._chat_proxy(body, data, requested_profile)
return
if requested_profile is not None:
self._profiled_proxy(body, data, requested_profile)
return
# An llama.cpp weiterleiten (mit Chat-Waiting, Streaming bleibt erhalten).
self._proxy_with_wait(body)
def _acquire_model_lease(self, profile: str | None = None) -> dict:
"""Atomar Profil sicherstellen und einen aktiven Request registrieren."""
with STATE.lock:
if profile is not None:
switch_profile(profile, implicit=True)
up = upstream_status()
if not up["reachable"] or not up.get("model"):
raise RuntimeError("llama.cpp nicht erreichbar")
with STATE.avail_lock:
if STATE.qwen_unavailable:
raise RuntimeError("Qwen wird gerade neu geladen")
STATE.active_chats += 1
return up
@staticmethod
def _release_model_lease() -> None:
with STATE.avail_lock:
STATE.active_chats = max(0, STATE.active_chats - 1)
def _profiled_proxy(self, body: bytes | None, data: dict,
profile: str) -> None:
try:
up = self._acquire_model_lease(profile)
except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error",
"upstream_unavailable")
return
try:
data["model"] = up["model"]
self._proxy(json.dumps(data).encode())
finally:
self._release_model_lease()
def _chat_proxy(self, body: bytes | None, data: dict,
profile: str | None) -> None:
"""Vision-Vorbereitung, Profilwahl und Chat-Lease als eine Transaktion."""
lease_acquired = False
try:
with STATE.lock:
if profile is not None:
switch_profile(profile, implicit=True)
image_url, _ = _extract_last_user_image(data)
if image_url:
if _vision_cached_analysis(_image_hash(image_url)) is None:
self.timeout = None
_vision_swap(data)
else:
log.info("Vision: Bild bereits analysiert "
"(Cache-Treffer) – kein Hotswap")
if _request_has_image(data):
data = _sanitize_for_main_model(data)
log.info("Vision: finale Hauptmodell-Inferenz gestartet")
up = upstream_status()
if not up["reachable"] or not up.get("model"):
raise RuntimeError("llama.cpp nicht erreichbar")
if profile is not None:
data["model"] = up["model"]
body = json.dumps(data).encode()
with STATE.avail_lock:
if STATE.qwen_unavailable:
raise RuntimeError("Qwen wird gerade neu geladen")
STATE.active_chats += 1
lease_acquired = True
self._proxy(body)
except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error", "upstream_unavailable")
finally:
if lease_acquired:
self._release_model_lease()
def _proxy_with_wait(self, body: bytes | None) -> None:
"""Leitet an llama.cpp weiter, wartet aber erst, bis Qwen verfügbar ist.
@@ -2128,15 +2417,86 @@ class _FlushHandler(logging.StreamHandler):
self.flush()
class RouterHTTPServer(ThreadingHTTPServer):
allow_reuse_address = True
def _startup_reconcile() -> None:
"""Reconcile persisted worker/model state before accepting requests."""
previous = RUNTIME.load()
worker = previous.get("worker")
if worker in {"image", "vision"}:
markers = ([IMAGE_WORKER] if worker == "image" else
[VISION_ALIAS, VISION_MODEL, str(VISION_PORT)])
terminated = terminate_recorded_worker(
previous, markers, lambda msg: log.warning("Recovery: %s", msg))
if terminated:
time.sleep(1)
RUNTIME.clear_worker(worker)
removed = enforce_artifact_retention(
IMAGE_DIR, IMAGE_RETENTION_FILES, IMAGE_RETENTION_BYTES,
IMAGE_RETENTION_DAYS)
if removed:
log.info("Startup-Retention: %d alte Bilder entfernt", len(removed))
profile = current_profile()
if profile is None:
saved = previous.get("last_profile")
if saved in PROFILES:
try:
switch_profile(saved)
log.info("Recovery: gespeichertes Profil %s neu angewendet", saved)
return
except Exception as exc:
log.error("Recovery: gespeichertes Profil %s konnte nicht "
"angewendet werden: %s", saved, exc)
profile = None
if profile is None:
log.error("Recovery: kein gültiges Profil gefunden; Router startet degraded")
_set_qwen_unavailable(True)
return
up = upstream_status()
if (up["reachable"] and up.get("model")
and up.get("ctx") == PROFILES[profile]
and (not EXPECTED_MODELS.get(profile)
or up.get("model") == EXPECTED_MODELS[profile])):
_set_qwen_unavailable(False)
RUNTIME.save(last_profile=profile, phase="idle")
log.info("Recovery: Profil %s ist bereits bereit", profile)
return
_set_qwen_unavailable(True)
try:
_restore_qwen(profile)
except Exception as exc:
log.error("Recovery: Profil %s konnte nicht gestartet werden: %s",
profile, exc)
return
_set_qwen_unavailable(False)
log.info("Recovery: Profil %s wurde wiederhergestellt", profile)
def main() -> None:
global AUTH
handler = _FlushHandler(sys.stdout)
handler.setFormatter(logging.Formatter(
"%(asctime)s %(levelname)s %(message)s"))
logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"),
handlers=[handler])
try:
AUTH = AuthPolicy.from_environment()
if not AUTH.enabled and HOST not in {"127.0.0.1", "::1", "localhost"}:
raise ConfigurationError(
"ROUTER_AUTH_MODE=off ist nur an einer Loopback-Adresse erlaubt")
except ConfigurationError as exc:
log.critical("Unsichere Router-Konfiguration: %s", exc)
raise SystemExit(2)
log.info("AI Profile Router startet: %s:%s -> %s (Profile: %s)",
HOST, PORT, UPSTREAM_URL, ", ".join(PROFILES))
server = ThreadingHTTPServer((HOST, PORT), Handler)
log.info("Authentifizierung: %s", "aktiv" if AUTH.enabled else "deaktiviert")
_startup_reconcile()
server = RouterHTTPServer((HOST, PORT), Handler)
server.daemon_threads = True
try:
server.serve_forever()
+16
View File
@@ -0,0 +1,16 @@
{
"profiles": {
"fast": {
"context": 73728,
"model_alias": "qwen38-27b-iq4mix-72k-mtp2"
},
"medium": {
"context": 94208,
"model_alias": "qwen38-27b-iq4xs-pure-92k"
},
"long": {
"context": 131072,
"model_alias": "qwen38-27b-iq4mix-128k-mtp2-ffn12"
}
}
}
+239
View File
@@ -0,0 +1,239 @@
#!/usr/bin/env python3
"""Security and runtime helpers for the AI Profile Router.
This module deliberately contains no model-specific logic. It provides the
small, testable building blocks that the HTTP gateway and the GPU orchestrator
share: API-key authentication, crash-state persistence, Linux child-process
cleanup and bounded artifact retention.
"""
from __future__ import annotations
import hmac
import json
import os
import signal
import threading
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Callable, Iterable
class ConfigurationError(RuntimeError):
"""Raised when the router would otherwise start in an unsafe state."""
def load_profile_registry(path: str | None) -> dict[str, dict]:
"""Load and validate the optional profile registry.
Passing no path keeps source-tree compatibility. An explicitly configured
but missing registry is a fatal configuration error: production must never
silently lose alias validation because a deployment forgot one file.
"""
fallback = {
"fast": {"context": 73728, "model_alias": None},
"medium": {"context": 94208, "model_alias": None},
"long": {"context": 131072, "model_alias": None},
}
if not path:
return fallback
try:
raw = json.loads(Path(path).read_text(encoding="utf-8"))
except (OSError, ValueError) as exc:
raise ConfigurationError(f"Profilregister kann nicht gelesen werden: {exc}")
profiles = raw.get("profiles") if isinstance(raw, dict) else None
if not isinstance(profiles, dict) or not profiles:
raise ConfigurationError("Profilregister enthält keine 'profiles'")
validated: dict[str, dict] = {}
for name, definition in profiles.items():
if not isinstance(name, str) or not re_full_profile_name(name):
raise ConfigurationError(f"ungültiger Profilname: {name!r}")
if not isinstance(definition, dict):
raise ConfigurationError(f"Profil {name!r} ist kein Objekt")
context = definition.get("context")
alias = definition.get("model_alias")
if not isinstance(context, int) or context < 1024:
raise ConfigurationError(f"ungültiger Kontext für Profil {name!r}")
if alias is not None and (not isinstance(alias, str) or not alias.strip()):
raise ConfigurationError(f"ungültiger Modellalias für Profil {name!r}")
validated[name] = {"context": context,
"model_alias": alias.strip() if alias else None}
return validated
def re_full_profile_name(value: str) -> bool:
return bool(value) and all(ch.isalnum() or ch in "-_" for ch in value)
@dataclass(frozen=True)
class AuthPolicy:
"""Bearer/X-API-Key authentication policy.
``mode`` is either ``required`` or ``off``. ``off`` is intended only for
loopback development tests. Production startup fails closed when the key
is missing or too short.
"""
mode: str
api_key: str | None
@classmethod
def from_environment(cls) -> "AuthPolicy":
mode = os.environ.get("ROUTER_AUTH_MODE", "required").strip().lower()
if mode not in {"required", "off"}:
raise ConfigurationError(
"ROUTER_AUTH_MODE muss 'required' oder 'off' sein")
if mode == "off":
return cls(mode=mode, api_key=None)
key = os.environ.get("ROUTER_API_KEY", "").strip()
key_file = os.environ.get(
"ROUTER_API_KEY_FILE", "/etc/mike-ai/router-api-key")
if not key and key_file:
try:
key = Path(key_file).read_text(encoding="utf-8").strip()
except OSError:
pass
if len(key) < 32:
raise ConfigurationError(
"Router-Authentifizierung ist aktiv, aber kein API-Key mit "
"mindestens 32 Zeichen vorhanden")
return cls(mode=mode, api_key=key)
@property
def enabled(self) -> bool:
return self.mode == "required"
def accepts(self, authorization: str | None,
x_api_key: str | None) -> bool:
if not self.enabled:
return True
candidate = (x_api_key or "").strip()
auth = (authorization or "").strip()
if not candidate and auth.lower().startswith("bearer "):
candidate = auth[7:].strip()
return bool(candidate and self.api_key
and hmac.compare_digest(candidate, self.api_key))
class RuntimeStore:
"""Tiny atomic JSON store used for crash reconciliation."""
def __init__(self, path: str) -> None:
self.path = Path(path)
self._lock = threading.RLock()
def load(self) -> dict:
with self._lock:
try:
value = json.loads(self.path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return {}
return value if isinstance(value, dict) else {}
def save(self, **updates) -> None:
with self._lock:
state = self.load()
for key, value in updates.items():
if value is None:
state.pop(key, None)
else:
state[key] = value
state["updated_at"] = int(time.time())
self.path.parent.mkdir(parents=True, exist_ok=True)
temp = self.path.with_name(
f".{self.path.name}.{os.getpid()}.{threading.get_ident()}.tmp")
try:
temp.write_text(
json.dumps(state, indent=2, sort_keys=True) + "\n",
encoding="utf-8")
os.chmod(temp, 0o600)
os.replace(temp, self.path)
finally:
try:
temp.unlink()
except FileNotFoundError:
pass
def clear_worker(self, worker: str) -> None:
with self._lock:
state = self.load()
if state.get("worker") == worker:
self.save(worker=None, worker_pid=None)
def terminate_recorded_worker(state: dict, allowed_markers: Iterable[str],
log: Callable[[str], None]) -> bool:
"""Terminate a previously recorded worker, but only after cmdline checks."""
pid = state.get("worker_pid")
if not isinstance(pid, int) or pid <= 1:
return False
cmdline_path = Path(f"/proc/{pid}/cmdline")
try:
cmdline = cmdline_path.read_bytes().replace(b"\0", b" ").decode(
errors="replace")
except OSError:
return False
if not any(marker in cmdline for marker in allowed_markers):
log(f"Recorded PID {pid} nicht beendet: Prozessprüfung fehlgeschlagen")
return False
try:
os.kill(pid, signal.SIGTERM)
except ProcessLookupError:
return False
log(f"Verwaister Router-Worker PID {pid} wurde beendet")
return True
def enforce_artifact_retention(directory: str, max_files: int, max_bytes: int,
max_age_days: int,
protected: Iterable[str] = ()) -> list[str]:
"""Delete oldest PNG + sidecar pairs until all retention limits hold."""
root = Path(directory)
if not root.is_dir():
return []
protected_set = set(protected)
now = time.time()
cutoff = now - max_age_days * 86400 if max_age_days > 0 else None
items: list[tuple[float, int, Path]] = []
for path in root.glob("*.png"):
if path.name in protected_set:
continue
try:
stat = path.stat()
except OSError:
continue
sidecar = path.with_suffix(".json")
size = stat.st_size
try:
size += sidecar.stat().st_size
except OSError:
pass
items.append((stat.st_mtime, size, path))
items.sort(key=lambda item: item[0])
total = sum(item[1] for item in items)
removed: list[str] = []
while items:
mtime, size, path = items[0]
too_old = cutoff is not None and mtime < cutoff
too_many = max_files > 0 and len(items) > max_files
too_large = max_bytes > 0 and total > max_bytes
if not (too_old or too_many or too_large):
break
items.pop(0)
try:
path.unlink()
removed.append(path.name)
except OSError:
continue
try:
path.with_suffix(".json").unlink()
except OSError:
pass
total -= size
return removed