diff --git a/README.md b/README.md index 1495935..4a94bf4 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,7 @@ llama.cpp :8080 + profilabhängige MCP-Server - [Saubere Installation](docs/INSTALLATION.md) - [Betrieb und Profilwechsel](docs/OPERATIONS.md) - [Sicherheitsmodell](docs/SECURITY.md) +- [Router V2: Migration und Kompatibilität](docs/ROUTER_V2_MIGRATION.md) - [Migration vom bestehenden Host](docs/MIGRATION.md) - [MCP-Aufteilung](platform/mcp/README.md) - [llama.cpp-Build und Profile](platform/llama/README.md) @@ -92,11 +93,13 @@ Sprachausgabe bereit (XTTS-v2, CPU-only, OpenAI-kompatibel). | Endpunkt | Beschreibung | |---|---| | `GET /v1/models` | Die drei virtuellen Modelle inkl. `context_length`/`context_window` | +| `GET /health` | öffentliche Liveness-Prüfung des Routerprozesses | +| `GET /ready` | öffentliche Readiness-Prüfung von Router + Textmodell | | `GET /status` | Aktives Profil, Upstream-Zustand, Modell, Kontext, Uptime | -| `POST /fast` `/medium` `/long` | Profilwechsel (auch `GET` möglich) | +| `POST /fast` `/medium` `/long` | expliziter Profilwechsel | | `POST /v1/chat/completions` | Weiterleitung an llama.cpp (Streaming + Tool Calls) | | `POST /v1/images/generations` | Bildgenerierung (FLUX.2 [klein] 4B Base, OpenAI-kompatibel) | -| `GET /images` | Liste der gespeicherten Bilder (max. 200) | +| `GET /images` | Liste der aufbewahrten Bilder | | `GET /images/` | PNG-Download (nur `images/`-Verzeichnis, validiert) | | `POST /v1/audio/speech` | Sprachausgabe (XTTS-v2, OpenAI-kompatibel) | | `POST /v1/audio/transcriptions` | Deutsche Spracherkennung (whisper.cpp, OpenAI-kompatibel) | @@ -105,6 +108,12 @@ Sprachausgabe bereit (XTTS-v2, CPU-only, OpenAI-kompatibel). | `POST /vision/test` | Direkter Vision-Test (Bild + Frage → Q3-Analyse) | | alles andere | Transparente Weiterleitung an llama.cpp | +Bis auf `GET /health` und `GET /ready` benötigen alle Endpunkte einen +Router-API-Key als `Authorization: Bearer …` oder `X-API-Key`. Der Installer +erzeugt ihn einmalig in `/etc/mike-ai/router-api-key` und gibt ihn niemals im +Installationslog aus. Ein fehlender oder zu kurzer Schlüssel verhindert den +Produktivstart (fail closed). + ### Verhalten - **Virtuelles Modell** (`qwen-fast`/`qwen-medium`/`qwen-long` in @@ -159,6 +168,7 @@ Beispiel: ```bash curl -s http://AI_HOST:8081/v1/images/generations \ + -H "Authorization: Bearer $ROUTER_KEY" \ -H 'Content-Type: application/json' \ -d '{"prompt":"ein roter Würfel auf weißem Grund","size":"1024x1024","quality":"standard"}' ``` @@ -306,11 +316,13 @@ Beispiele: ```bash # MP3 (Default), Stimme claribel curl -s http://AI_HOST:8081/v1/audio/speech \ + -H "Authorization: Bearer $ROUTER_KEY" \ -H 'Content-Type: application/json' \ -d '{"input":"Hallo, dies ist ein Test.","voice":"claribel"}' -o out.mp3 # WAV, 1.5x Tempo curl -s http://AI_HOST:8081/v1/audio/speech \ + -H "Authorization: Bearer $ROUTER_KEY" \ -H 'Content-Type: application/json' \ -d '{"input":"Guten Tag.","voice":"claribel","speed":1.5,"response_format":"wav"}' -o out.wav ``` @@ -400,17 +412,20 @@ Beispiele: ```bash # WebM/Opus (z.B. aus Open WebUI-Mikrofon) curl -s http://AI_HOST:8081/v1/audio/transcriptions \ + -H "Authorization: Bearer $ROUTER_KEY" \ -F "file=@aufnahme.webm" \ -F "model=whisper-1" # WAV mit expliziter Sprache curl -s http://AI_HOST:8081/v1/audio/transcriptions \ + -H "Authorization: Bearer $ROUTER_KEY" \ -F "file=@aufnahme.wav" \ -F "model=whisper-1" \ -F "language=de" # Verbose-Format curl -s http://AI_HOST:8081/v1/audio/transcriptions \ + -H "Authorization: Bearer $ROUTER_KEY" \ -F "file=@aufnahme.wav" \ -F "model=whisper-1" \ -F "response_format=verbose_json" @@ -510,6 +525,11 @@ Die Installation ist idempotent (Update = erneut ausführen). |---|---|---| | `ROUTER_HOST` | `0.0.0.0` | Bind-Adresse | | `ROUTER_PORT` | `8081` | Port | +| `ROUTER_AUTH_MODE` | `required` | Authentifizierung; `off` nur für lokale Tests | +| `ROUTER_API_KEY_FILE` | `/etc/mike-ai/router-api-key` | Schlüsseldatei (0600) | +| `ROUTER_STATE_FILE` | `/var/lib/mike-ai-profile-router/state.json` | atomarer Crash-/Recovery-Zustand | +| `ROUTER_PROFILES_FILE` | `/etc/mike-ai/router-profiles.json` | Profile, Kontext und erwarteter Alias | +| `ROUTER_MAX_CONCURRENT_REQUESTS` | `16` | harte Grenze paralleler Requests | | `UPSTREAM_URL` | `http://127.0.0.1:8080` | llama.cpp | | `PROFILE_SCRIPT` | `/usr/local/bin/llama-profile` | Profil-Skript | | `PROFILE_DIR` | `/etc/systemd/system/mike-ai-llama-ui.service.d` | Ort der `override.conf` | @@ -538,6 +558,8 @@ Die Installation ist idempotent (Update = erneut ausführen). | `VISION_UNLOAD_TIMEOUT` | `120` | Warten auf VRAM-Freiheit (s) | | `VISION_MAX_TOKENS` | `4096` | Max. Tokens der Vision-Analyse | | `VISION_CACHE_MAX` | `64` | Größe des Analyse-Caches (LRU) | +| `VISION_MAX_IMAGE_BYTES` | `20971520` | maximales dekodiertes Bild (20 MiB) | +| `VISION_ALLOW_REMOTE_URLS` | `false` | externe Bild-URLs; standardmäßig SSRF-sicher aus | | `VISION_LOG` | `/opt/mike-ai/ai-profile-router/vision_server.log` | Vision-Server-Log | TTS-Worker (`mike-ai-xtts.service`): @@ -567,7 +589,8 @@ paralleler Chat während Bild-Job wartet statt 502) und **Sprachausgabe** (`/status` mit tts-Section, `POST /v1/audio/speech` wav/mp3, Validierung, Worker-Fehler→503, Worker down→503, Worker-Neustart→Recovery). -Aktuell: **43 Tests** (32 bestehende + 11 TTS-Assertions). +Aktuell: **63 Integrationsassertions** plus Chunked-, Multipart- und +Security/Recovery-Unit-Tests. ## Betrieb @@ -576,10 +599,12 @@ systemctl status mike-ai-profile-router systemctl status mike-ai-xtts journalctl -u mike-ai-profile-router -f journalctl -u mike-ai-xtts -f -curl -s http://AI_HOST:8081/status | python3 -m json.tool -curl -s -X POST http://AI_HOST:8081/medium +ROUTER_KEY='aus lokalem Secret-Store' +curl -s http://AI_HOST:8081/status -H "Authorization: Bearer $ROUTER_KEY" | python3 -m json.tool +curl -s -X POST http://AI_HOST:8081/medium -H "Authorization: Bearer $ROUTER_KEY" # TTS-Test curl -s http://AI_HOST:8081/v1/audio/speech \ + -H "Authorization: Bearer $ROUTER_KEY" \ -H 'Content-Type: application/json' \ -d '{"input":"Hallo","voice":"claribel"}' -o test.mp3 ``` @@ -589,7 +614,12 @@ curl -s http://AI_HOST:8081/v1/audio/speech \ - Keine Shell-Aufrufe: Profil-Skript wird mit `subprocess.run([script, profil])` aufgerufen, `profil` ist Whitelist-geprüft (`fast|medium|long`). - Keine Secrets/Tokens im Code oder in der Unit. -- systemd-Hardening: `NoNewPrivileges=true`, `PrivateTmp=true`. +- API-Key-Pflicht, konstante Schlüsselprüfung und keine Weitergabe des + Router-Schlüssels an llama.cpp. +- Remote-Bilder standardmäßig gesperrt; lokale Data-URLs sind typ- und + größenvalidiert. +- systemd-Hardening: `NoNewPrivileges`, `PrivateTmp`, `ProtectHome`, + Kernel-/Control-Group-Schutz und restriktive Dateirechte. - Die bestehende llama.cpp-/Profil-Konfiguration wird nicht verändert; der Router nutzt nur das vorhandene `llama-profile`-Skript und liest die `override.conf`. diff --git a/deploy/deploy.sh b/deploy/deploy.sh index e2685a9..751f405 100755 --- a/deploy/deploy.sh +++ b/deploy/deploy.sh @@ -9,7 +9,8 @@ SSH_KEY="${SSH_KEY:-$HOME/.ssh/id_ed25519}" STAGE="/tmp/ai-profile-router-$$" mkdir -p "$STAGE" -cp router/ai_profile_router.py router/image_worker.py router/xtts_worker.py \ +cp router/ai_profile_router.py router/router_support.py router/router_profiles.json \ + router/image_worker.py router/xtts_worker.py \ router/stt_worker.py \ deploy/install.sh deploy/mike-ai-profile-router.service \ deploy/mike-ai-xtts.service deploy/mike-ai-whisper.service \ diff --git a/deploy/install.sh b/deploy/install.sh index 8b1b454..f2d727c 100755 --- a/deploy/install.sh +++ b/deploy/install.sh @@ -38,7 +38,10 @@ fi # --- 2. Neue Dateien installieren ------------------------------------------- mkdir -p "$INSTALL_DIR" "$IMAGE_DIR" +install -d -m 0700 /etc/mike-ai /var/lib/mike-ai-profile-router install -m 0755 "$DIR/ai_profile_router.py" "$INSTALL_DIR/ai_profile_router.py" +install -m 0644 "$DIR/router_support.py" "$INSTALL_DIR/router_support.py" +install -m 0644 "$DIR/router_profiles.json" /etc/mike-ai/router-profiles.json install -m 0755 "$DIR/image_worker.py" "$INSTALL_DIR/image_worker.py" install -m 0755 "$DIR/xtts_worker.py" "$INSTALL_DIR/xtts_worker.py" install -m 0755 "$DIR/stt_worker.py" "$INSTALL_DIR/stt_worker.py" @@ -46,6 +49,14 @@ install -m 0644 "$DIR/${SERVICE}" "/etc/systemd/system/${SERVICE}" install -m 0644 "$DIR/${XTTS_SERVICE}" "/etc/systemd/system/${XTTS_SERVICE}" install -m 0644 "$DIR/${WHISPER_SERVICE}" "/etc/systemd/system/${WHISPER_SERVICE}" +# Eigener Router-Key. Er wird niemals in Git oder ins Installationslog geschrieben. +if [ ! -s /etc/mike-ai/router-api-key ]; then + umask 077 + python3 -c 'import secrets; print(secrets.token_urlsafe(48))' \ + > /etc/mike-ai/router-api-key +fi +chmod 0600 /etc/mike-ai/router-api-key + # --- 3. Python-Venv mit Bild-Abhängigkeiten --------------------------------- if [ ! -x "$VENV/bin/python" ]; then echo "-- Erstelle Python-Venv in $VENV" @@ -150,5 +161,17 @@ if ! systemctl is-active --quiet "$WHISPER_SERVICE"; then exit 1 fi echo "-- Whisper-Service läuft" -curl -sf "http://127.0.0.1:8081/status" | python3 -m json.tool +python3 - <<'PY' +import json +import pathlib +import urllib.request + +key = pathlib.Path("/etc/mike-ai/router-api-key").read_text().strip() +request = urllib.request.Request( + "http://127.0.0.1:8081/status", + headers={"Authorization": f"Bearer {key}"}) +with urllib.request.urlopen(request, timeout=10) as response: + print(json.dumps(json.load(response), indent=2)) +PY +echo "-- Router-API-Key liegt ausschließlich in /etc/mike-ai/router-api-key" echo "== Fertig ==" diff --git a/deploy/mike-ai-profile-router.service b/deploy/mike-ai-profile-router.service index a1d14a0..e897f3e 100644 --- a/deploy/mike-ai-profile-router.service +++ b/deploy/mike-ai-profile-router.service @@ -8,8 +8,14 @@ Type=simple ExecStart=/opt/mike-ai/ai-profile-router/venv/bin/python /opt/mike-ai/ai-profile-router/ai_profile_router.py Restart=on-failure RestartSec=3 +EnvironmentFile=-/etc/mike-ai/router.env Environment=ROUTER_HOST=0.0.0.0 Environment=ROUTER_PORT=8081 +Environment=ROUTER_AUTH_MODE=required +Environment=ROUTER_API_KEY_FILE=/etc/mike-ai/router-api-key +Environment=ROUTER_STATE_FILE=/var/lib/mike-ai-profile-router/state.json +Environment=ROUTER_PROFILES_FILE=/etc/mike-ai/router-profiles.json +Environment=ROUTER_MAX_CONCURRENT_REQUESTS=16 Environment=UPSTREAM_URL=http://127.0.0.1:8080 Environment=PROFILE_SCRIPT=/usr/local/bin/llama-profile Environment=PROFILE_DIR=/etc/systemd/system/mike-ai-llama-ui.service.d @@ -22,6 +28,9 @@ Environment=IMAGE_DIR=/opt/mike-ai/ai-profile-router/images Environment=IMAGE_WORKER_LOG=/opt/mike-ai/ai-profile-router/worker.log Environment=IMAGE_GEN_TIMEOUT=600 Environment=IMAGE_VRAM_FREE_TIMEOUT=120 +Environment=IMAGE_RETENTION_FILES=100 +Environment=IMAGE_RETENTION_BYTES=5368709120 +Environment=IMAGE_RETENTION_DAYS=30 Environment=LLAMA_SERVER_BIN=/opt/mike-ai/llama.cpp/build/bin/llama-server Environment=VISION_MODEL=/opt/mike-ai/models/qwen3.8-27b/Qwen3.8-27B-Q3_K_M.gguf Environment=VISION_MMPROJ=/opt/mike-ai/models/qwen3.8-27b-nvfp4/mmproj-BF16.gguf @@ -31,9 +40,20 @@ Environment=VISION_LOAD_TIMEOUT=300 Environment=VISION_INFER_TIMEOUT=300 Environment=VISION_UNLOAD_TIMEOUT=120 Environment=VISION_MAX_TOKENS=4096 +Environment=VISION_MAX_IMAGE_BYTES=20971520 +Environment=VISION_ALLOW_REMOTE_URLS=false Environment=VISION_LOG=/opt/mike-ai/ai-profile-router/vision_server.log NoNewPrivileges=true PrivateTmp=true +ProtectHome=true +ProtectKernelLogs=true +ProtectKernelModules=true +ProtectKernelTunables=true +ProtectControlGroups=true +LockPersonality=true +RestrictRealtime=true +RestrictSUIDSGID=true +UMask=0077 [Install] WantedBy=multi-user.target diff --git a/dev/fake-llama-profile.sh b/dev/fake-llama-profile.sh index 9217422..ec4580b 100755 --- a/dev/fake-llama-profile.sh +++ b/dev/fake-llama-profile.sh @@ -9,5 +9,11 @@ case "${1:-}" in *) echo "Usage: fake-llama-profile {fast|medium|long}"; exit 1 ;; esac +if [ -f /tmp/fake-profile-fail ] \ + && [ "$(cat /tmp/fake-profile-fail)" = "$1" ]; then + echo "fake: simulierter Profilfehler für '$1'" + exit 42 +fi + cp "$D/profile-$1.conf.disabled" "$D/override.conf" echo "fake: Profil '$1' gesetzt" diff --git a/dev/mock_upstream.py b/dev/mock_upstream.py index 1e0ff06..76f4ad9 100755 --- a/dev/mock_upstream.py +++ b/dev/mock_upstream.py @@ -58,6 +58,9 @@ class Handler(BaseHTTPRequestHandler): self._json(404, {"error": {"message": "not found"}}) def _completion(self, body: dict) -> dict: + delay = body.get("mock_delay", 0) + if isinstance(delay, (int, float)) and 0 < delay <= 10: + time.sleep(delay) model = body.get("model") if body.get("tools"): name = body["tools"][0]["function"]["name"] @@ -84,6 +87,8 @@ class Handler(BaseHTTPRequestHandler): "finish_reason": finish}], "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}, + "mock_ctx": current_ctx(), + "mock_authorization": self.headers.get("Authorization"), } def _stream(self, body: dict) -> None: diff --git a/dev/test_chunked.py b/dev/test_chunked.py index c2fca1f..a50e2fc 100644 --- a/dev/test_chunked.py +++ b/dev/test_chunked.py @@ -176,10 +176,17 @@ def cleanup() -> None: def main() -> None: global procs + try: + os.unlink("/tmp/test-chunked-router-state.json") + except FileNotFoundError: + pass env = os.environ.copy() env.update({ "ROUTER_HOST": "127.0.0.1", "ROUTER_PORT": str(PORTS["router"]), + "ROUTER_AUTH_MODE": "off", + "ROUTER_PROFILES_FILE": "", + "ROUTER_STATE_FILE": "/tmp/test-chunked-router-state.json", "UPSTREAM_URL": f"http://127.0.0.1:{PORTS['llama']}", "PROFILE_SCRIPT": FAKE_PROFILE, "PROFILE_DIR": FAKE_PROFILE_DIR, diff --git a/dev/test_local.sh b/dev/test_local.sh index 9f19069..2a9a81c 100755 --- a/dev/test_local.sh +++ b/dev/test_local.sh @@ -4,18 +4,27 @@ set -uo pipefail cd "$(dirname "$0")/.." -UP_PORT=18080 -RT_PORT=18081 -TTS_PORT=18082 -STT_PORT=18083 +UP_PORT="${UP_PORT:-18080}" +RT_PORT="${RT_PORT:-18081}" +TTS_PORT="${TTS_PORT:-18082}" +STT_PORT="${STT_PORT:-18083}" BASE="http://127.0.0.1:$RT_PORT" FAKE_DIR="$PWD/dev/fake-profile-dir" +TEST_ROUTER_KEY="test-router-key-0123456789-abcdefghijklmnopqrstuvwxyz" PASS=0 FAIL=0 +# Produktive Authentifizierung für alle Integrationstests. `command curl` +# umgeht diese Funktion bei den gezielten anonymen Negativtests. +curl() { command curl -H "Authorization: Bearer $TEST_ROUTER_KEY" "$@"; } + cleanup() { - kill "${MOCK_PID:-}" "${ROUTER_PID:-}" "${TTS_PID:-}" "${STT_PID:-}" 2>/dev/null || true - rm -f /tmp/mock_pid2 /tmp/mock_upstream_pid + CURRENT_MOCK_PID="$(cat /tmp/mock_upstream_pid 2>/dev/null || true)" + kill "${MOCK_PID:-}" "${CURRENT_MOCK_PID:-}" "${ROUTER_PID:-}" \ + "${TTS_PID:-}" "${STT_PID:-}" 2>/dev/null || true + rm -f /tmp/mock_pid2 /tmp/mock_upstream_pid /tmp/test-router-state.json \ + /tmp/fake-profile-fail + cp "$FAKE_DIR/profile-fast.conf.disabled" "$FAKE_DIR/override.conf" wait 2>/dev/null || true } trap cleanup EXIT @@ -23,8 +32,20 @@ trap cleanup EXIT ok() { echo " PASS: $1"; PASS=$((PASS+1)); } bad() { echo " FAIL: $1"; FAIL=$((FAIL+1)); } +wait_http() { + local url="$1" name="$2" + for _ in $(seq 1 50); do + curl -sf "$url" >/dev/null 2>&1 && return 0 + sleep 0.1 + done + echo "FEHLER: $name wurde nicht bereit: $url" >&2 + return 1 +} + # --- Mock-llama.cpp starten (über Fake-systemctl) ------------------------------ echo "== Starte Mock-llama.cpp (Port $UP_PORT)" +cp "$FAKE_DIR/profile-fast.conf.disabled" "$FAKE_DIR/override.conf" +rm -f /tmp/test-router-state.json /tmp/mock_upstream_pid FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid \ FAKE_SYSTEMD_PORT="$UP_PORT" \ FAKE_SYSTEMD_PROFILE_DIR="$FAKE_DIR" \ @@ -33,11 +54,15 @@ FAKE_SYSTEMD_LOG=/tmp/mock_upstream.log \ bash dev/fake-systemctl.sh start sleep 0.5 MOCK_PID=$(cat /tmp/mock_upstream_pid 2>/dev/null || echo "") +wait_http "http://127.0.0.1:$UP_PORT/health" "Mock-llama.cpp" || exit 1 # --- Router starten ----------------------------------------------------------- echo "== Starte Router (Port $RT_PORT)" rm -rf /tmp/test-images ROUTER_HOST=127.0.0.1 ROUTER_PORT="$RT_PORT" \ +ROUTER_AUTH_MODE=required ROUTER_API_KEY="$TEST_ROUTER_KEY" \ +ROUTER_PROFILES_FILE= \ +ROUTER_STATE_FILE=/tmp/test-router-state.json \ UPSTREAM_URL="http://127.0.0.1:$UP_PORT" \ PROFILE_SCRIPT="$PWD/dev/fake-llama-profile.sh" \ PROFILE_DIR="$FAKE_DIR" \ @@ -58,16 +83,27 @@ TTS_WORKER_URL="http://127.0.0.1:$TTS_PORT" \ STT_WORKER_URL="http://127.0.0.1:$STT_PORT" \ python3 router/ai_profile_router.py >/tmp/router_test.log 2>&1 & ROUTER_PID=$! -sleep 0.5 +wait_http "$BASE/health" "Router" || { + cat /tmp/router_test.log >&2 + exit 1 +} rm -f /tmp/test_worker_requests.jsonl +echo "== Test 0: Authentifizierung + Health/Readiness" +CODE=$(command curl -s -o /tmp/err0.json -w "%{http_code}" "$BASE/status") +[ "$CODE" = "401" ] && ok "Status ohne Key → 401" || bad "Status ohne Key: HTTP $CODE" +CODE=$(command curl -s -o /dev/null -w "%{http_code}" "$BASE/health") +[ "$CODE" = "200" ] && ok "öffentliche Liveness → 200" || bad "Liveness: HTTP $CODE" +CODE=$(command curl -s -o /dev/null -w "%{http_code}" "$BASE/ready") +[ "$CODE" = "200" ] && ok "öffentliche Readiness → 200" || bad "Readiness: HTTP $CODE" + # --- Mock-TTS-Worker starten ---------------------------------------------------- echo "== Starte Mock-TTS-Worker (Port $TTS_PORT)" MOCK_TTS_PORT="$TTS_PORT" MOCK_TTS_DELAY=0.1 \ MOCK_TTS_LOG=/tmp/test_tts_requests.jsonl \ python3 dev/mock_tts_worker.py >/tmp/mock_tts.log 2>&1 & TTS_PID=$! -sleep 0.5 +wait_http "http://127.0.0.1:$TTS_PORT/status" "Mock-TTS" || exit 1 rm -f /tmp/test_tts_requests.jsonl # --- Mock-STT-Worker starten ---------------------------------------------------- @@ -76,7 +112,7 @@ MOCK_STT_PORT="$STT_PORT" MOCK_STT_DELAY=0.1 \ MOCK_STT_LOG=/tmp/test_stt_requests.jsonl \ python3 dev/mock_stt_worker.py >/tmp/mock_stt.log 2>&1 & STT_PID=$! -sleep 0.5 +wait_http "http://127.0.0.1:$STT_PORT/status" "Mock-STT" || exit 1 rm -f /tmp/test_stt_requests.jsonl # --- 1. /v1/models ------------------------------------------------------------- @@ -115,6 +151,7 @@ import json,sys d=json.load(sys.stdin) assert d["model"].startswith("mock-model-"), d assert "Mock-Antwort" in d["choices"][0]["message"]["content"], d +assert d.get("mock_authorization") is None, d ' && ok "Request wurde weitergeleitet, Modell ersetzt" || bad "Forwarding" # --- 4. Streaming ---------------------------------------------------------------- @@ -176,17 +213,34 @@ d=json.load(sys.stdin) assert d["model"]=="mock-model-131072", d ' && ok "qwen-long hat Profil long aktiviert und weitergeleitet" || bad "virtuelles Modell" -# --- 9. Ungültiges Profil ------------------------------------------------------------------ -echo "== Test 9: Ungültiges Profil" -CODE=$(curl -s -o /tmp/err9.json -w "%{http_code}" -X POST "$BASE/huge") +# --- 9. Methoden und ungültiges virtuelles Modell ----------------------------------------- +echo "== Test 9: sichere Profilmethoden + ungültiges virtuelles Modell" +CODE=$(curl -s -o /tmp/err9.json -w "%{http_code}" "$BASE/long") cat /tmp/err9.json; echo -[ "$CODE" = "400" ] && ok "400 bei unbekanntem Profil (POST /huge)" || bad "erwartet 400, bekam $CODE" +[ "$CODE" = "405" ] && ok "GET /long verändert kein Profil" || bad "erwartet 405, bekam $CODE" CODE=$(curl -s -o /tmp/err9b.json -w "%{http_code}" -X POST "$BASE/v1/chat/completions" \ -H "Content-Type: application/json" -d '{"model":"qwen-huge","messages":[]}') cat /tmp/err9b.json; echo [ "$CODE" = "400" ] && ok "400 bei unbekanntem virtuellen Modell (qwen-huge)" || bad "erwartet 400, bekam $CODE" +# Unbekannte einteilige POST-Pfade gehören dem Upstream, nicht dem Profilrouter. +CODE=$(curl -s -o /tmp/err9c.json -w "%{http_code}" -X POST "$BASE/tokenize" \ + -H "Content-Type: application/json" -d '{"content":"Hallo"}') +[ "$CODE" = "404" ] && ok "POST /tokenize wurde transparent weitergeleitet" \ + || bad "erwartet Upstream-404, bekam $CODE" + +echo medium >/tmp/fake-profile-fail +CODE=$(curl -s -o /tmp/err9d.json -w "%{http_code}" -X POST "$BASE/medium") +[ "$CODE" = "503" ] && ok "fehlgeschlagenes Profilskript → 503" \ + || bad "Profilskript-Fehler: HTTP $CODE" +rm -f /tmp/fake-profile-fail +CODE=$(curl -s -o /tmp/chat9e.json -w "%{http_code}" "$BASE/v1/chat/completions" \ + -H "Content-Type: application/json" \ + -d '{"model":"qwen-fast","messages":[{"role":"user","content":"noch da?"}]}') +[ "$CODE" = "200" ] && ok "vorheriges Profil bleibt nach Skriptfehler verfügbar" \ + || bad "Qwen nach Profilskript-Fehler: HTTP $CODE" + # --- 10. llama.cpp down -> 502, danach Recovery --------------------------------------------------- echo "== Test 10: Upstream down -> 502, danach Recovery" # Profil auf fast setzen (aus Test 8 ist long aktiv) @@ -195,6 +249,10 @@ curl -sf -X POST "$BASE/fast" >/dev/null FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid FAKE_SYSTEMD_PORT="$UP_PORT" \ bash dev/fake-systemctl.sh stop sleep 0.5 +CODE=$(command curl -s -o /dev/null -w "%{http_code}" "$BASE/health") +[ "$CODE" = "200" ] && ok "Liveness bleibt bei downem Modell 200" || bad "Liveness down: HTTP $CODE" +CODE=$(command curl -s -o /dev/null -w "%{http_code}" "$BASE/ready") +[ "$CODE" = "503" ] && ok "Readiness zeigt downes Modell mit 503" || bad "Readiness down: HTTP $CODE" CODE=$(curl -s -o /tmp/err10.json -w "%{http_code}" -X POST "$BASE/v1/chat/completions" \ -H "Content-Type: application/json" -d '{"model":"qwen-fast","messages":[]}') cat /tmp/err10.json; echo @@ -641,6 +699,26 @@ wait $STT_PID43 [ "$CODE" = "200" ] && [ -s /tmp/tts43.mp3 ] \ && ok "STT + TTS parallel (beide 200)" || bad "STT + TTS parallel (TTS Code $CODE)" +# --- 44. Zwei konkurrierende Profilanfragen ------------------------------------------------ +echo "== Test 44: Profil-Lease verhindert Wechsel während eines Chats" +curl -sf "$BASE/v1/chat/completions" -H "Content-Type: application/json" \ + -d '{"model":"qwen-long","mock_delay":1.0,"messages":[{"role":"user","content":"Lang"}]}' \ + >/tmp/chat44-long.json & +CHAT44_PID=$! +sleep 0.2 +curl -sf "$BASE/v1/chat/completions" -H "Content-Type: application/json" \ + -d '{"model":"qwen-medium","messages":[{"role":"user","content":"Mittel"}]}' \ + >/tmp/chat44-medium.json +wait "$CHAT44_PID" +python3 -c ' +import json +long=json.load(open("/tmp/chat44-long.json")) +medium=json.load(open("/tmp/chat44-medium.json")) +assert long["mock_ctx"] == 131072, long +assert medium["mock_ctx"] == 94208, medium +' && ok "konkurrierende Chats behielten jeweils ihr Profil" \ + || bad "Profil-Lease bei konkurrierenden Chats" + # --- Ergebnis -------------------------------------------------------------------------------------------- echo echo "== Ergebnis: $PASS bestanden, $FAIL fehlgeschlagen ==" diff --git a/dev/test_router_support.py b/dev/test_router_support.py new file mode 100644 index 0000000..7fb0bd2 --- /dev/null +++ b/dev/test_router_support.py @@ -0,0 +1,109 @@ +#!/usr/bin/env python3 +"""Unit tests for security, runtime persistence and artifact retention.""" + +from __future__ import annotations + +import os +import tempfile +import threading +import time +import unittest +from pathlib import Path +from unittest.mock import patch + +import sys + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "router")) +os.environ.setdefault("ROUTER_PROFILES_FILE", "") + +from router_support import ( # noqa: E402 + AuthPolicy, + ConfigurationError, + RuntimeStore, + enforce_artifact_retention, + load_profile_registry, +) +from ai_profile_router import _normalize_vision_image # noqa: E402 + + +class AuthPolicyTests(unittest.TestCase): + def test_required_mode_fails_closed_without_key(self) -> None: + with patch.dict(os.environ, { + "ROUTER_AUTH_MODE": "required", + "ROUTER_API_KEY": "", + "ROUTER_API_KEY_FILE": "/definitely/missing"}, clear=False): + with self.assertRaises(ConfigurationError): + AuthPolicy.from_environment() + + def test_bearer_and_x_api_key(self) -> None: + key = "k" * 48 + policy = AuthPolicy("required", key) + self.assertTrue(policy.accepts(f"Bearer {key}", None)) + self.assertTrue(policy.accepts(None, key)) + self.assertFalse(policy.accepts("Bearer wrong", None)) + + +class RuntimeStoreTests(unittest.TestCase): + def test_atomic_roundtrip_and_delete(self) -> None: + with tempfile.TemporaryDirectory() as temp: + store = RuntimeStore(str(Path(temp) / "state.json")) + store.save(worker="vision", worker_pid=123, last_profile="fast") + self.assertEqual(store.load()["worker_pid"], 123) + store.clear_worker("vision") + state = store.load() + self.assertNotIn("worker", state) + self.assertEqual(state["last_profile"], "fast") + + def test_concurrent_updates_are_not_lost(self) -> None: + with tempfile.TemporaryDirectory() as temp: + store = RuntimeStore(str(Path(temp) / "state.json")) + threads = [threading.Thread(target=store.save, + kwargs={f"key_{index}": index}) + for index in range(20)] + for thread in threads: + thread.start() + for thread in threads: + thread.join() + state = store.load() + for index in range(20): + self.assertEqual(state[f"key_{index}"], index) + + +class ProfileRegistryTests(unittest.TestCase): + def test_explicit_missing_registry_fails_closed(self) -> None: + with self.assertRaises(ConfigurationError): + load_profile_registry("/definitely/missing/profiles.json") + + +class VisionInputTests(unittest.TestCase): + def test_small_png_data_url_is_accepted(self) -> None: + value = "data:image/png;base64,iVBORw0KGgo=" + self.assertEqual(_normalize_vision_image(value), value) + + def test_remote_url_is_denied_by_default(self) -> None: + with self.assertRaisesRegex(ValueError, "deaktiviert"): + _normalize_vision_image("https://example.com/private.png") + + def test_invalid_base64_is_rejected(self) -> None: + with self.assertRaisesRegex(ValueError, "Base64"): + _normalize_vision_image("data:image/png;base64,not!base64") + + +class RetentionTests(unittest.TestCase): + def test_oldest_pairs_are_removed(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + for index in range(3): + png = root / f"image-{index}.png" + png.write_bytes(b"x" * 10) + png.with_suffix(".json").write_text("{}", encoding="utf-8") + stamp = time.time() - (30 - index) + os.utime(png, (stamp, stamp)) + removed = enforce_artifact_retention( + temp, max_files=2, max_bytes=0, max_age_days=0) + self.assertEqual(removed, ["image-0.png"]) + self.assertFalse((root / "image-0.json").exists()) + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 0385a51..8dc5307 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -11,7 +11,7 @@ in diesem Repository oder im Modellmanifest beschrieben sind. | Komponente | Port | Ausführung | Aufgabe | |---|---:|---|---| -| AI Profile Router | 8081 | systemd, unprivilegiert empfohlen | zentrale Client-API und Orchestrierung | +| AI Profile Router | 8081 | systemd, gehärtet | zentrale Client-API und Orchestrierung | | llama.cpp | 8080 | systemd | Textmodell, Tool Calling und MCP | | Whisper | 8084, nur localhost | systemd | Speech-to-Text | | XTTS | 8085, nur localhost | systemd | Text-to-Speech | @@ -29,11 +29,20 @@ in diesem Repository oder im Modellmanifest beschrieben sind. 4. Der Router wartet auf Modellname und erwartete Kontextgröße. 5. Erst dann wird die Anfrage an Port 8080 weitergeleitet. +Profilwahl und Weiterleitung bilden dabei eine atomare Modell-Lease. Ein +zweiter Request kann das Profil nicht mehr zwischen Auswahl und Inferenz +wechseln. Laufende Requests werden vor einem GPU-Hotswap vollständig beendet; +bei Überschreiten des Drain-Timeouts wird der Wechsel abgebrochen, nicht der +Chat. + ## Profilprinzip Die Profile sind vollständige systemd-Overrides. Ein Profilwechsel kopiert die gewählte Datei atomar auf `override.conf`, lädt systemd neu und startet genau einen llama.cpp-Dienst neu. Es gibt niemals mehrere Textmodelle gleichzeitig. +Das unabhängige Register `/etc/mike-ai/router-profiles.json` definiert Kontext +und erwarteten Modellalias. Readiness gilt nur, wenn beides exakt passt; eine +abweichende oder fehlende Registry verhindert den Start. ## GPU-Hotswap @@ -47,6 +56,23 @@ Vision und Bildgenerierung teilen sich die RTX mit dem Textmodell. Der Router: 6. stellt das ursprüngliche Textprofil wieder her, 7. prüft Modell und Kontext vor der Freigabe. +Der zuletzt stabile Zustand und temporäre Worker-PIDs werden atomar unter +`/var/lib/mike-ai-profile-router/state.json` festgehalten. Beim Routerstart +werden ausschließlich dort erfasste Prozesse nach zusätzlicher +Kommandozeilenprüfung beendet und das letzte Profil wiederhergestellt. + +## Vertrauensgrenzen + +- Port 8081 verlangt einen eigenen API-Key; nur `/health` und `/ready` sind + absichtlich anonym und enthalten keine privaten Daten. +- Der Client-Key wird vor dem lokalen Upstream entfernt. +- Bild-Uploads sind begrenzt. Remote-Bild-URLs sind standardmäßig aus, damit + der Vision-Pfad nicht als Zugriff auf Intranet oder Metadatenendpunkte dient. +- Maximal 16 Requests werden gleichzeitig bearbeitet; weitere erhalten 429. +- Der Dienst benötigt derzeit wegen systemd-Profilwechsel und GPU-Hotswap noch + Root-Rechte. Die systemd-Sandbox begrenzt diese, ersetzt aber keine künftige + Aufteilung in unprivilegierten Proxy und eng begrenzten Root-Helper. + ## Verzeichnislayout auf dem Zielhost ```text diff --git a/docs/DISASTER_RECOVERY.md b/docs/DISASTER_RECOVERY.md index cdc92c7..a8a7380 100644 --- a/docs/DISASTER_RECOVERY.md +++ b/docs/DISASTER_RECOVERY.md @@ -31,6 +31,10 @@ laufen, sondern alle fachlichen Funktionen geprüft wurden. ## Phase C – Router +- [ ] Start ohne API-Key schlägt bewusst fehl +- [ ] `/health` bleibt bei Hotswap 200 und `/ready` wird vorübergehend 503 +- [ ] geschützte Endpunkte liefern ohne Key 401 +- [ ] Router-Key erscheint weder im Upstream noch im Journal - [ ] `/status` meldet den richtigen Upstream - [ ] `/v1/models` liefert drei virtuelle Modelle - [ ] `/fast`, `/medium` und `/long` wechseln zuverlässig @@ -40,6 +44,8 @@ laufen, sondern alle fachlichen Funktionen geprüft wurden. - [ ] Tool Calls funktionieren - [ ] Fehler sind OpenAI-kompatibel - [ ] ein abgebrochener Client hinterlässt keinen blockierten Job +- [ ] erzwungener Routerabbruch wird aus Zustandsdatei sauber rekonstruiert +- [ ] fehlende/falsche Profilregistry verhindert falsche Readiness ## Phase D – Web und MCP @@ -58,6 +64,7 @@ laufen, sondern alle fachlichen Funktionen geprüft wurden. - [ ] neues Bild löst genau einen Vision-Hotswap aus - [ ] Folgefrage verwendet Cache und keinen zweiten Hotswap - [ ] Bilddaten werden vor dem Textmodell sanitisiert +- [ ] Remote-/private Bild-URL wird abgewiesen und übergroße Data-URL blockiert - [ ] Qwen-Profil wird nach Vision wiederhergestellt - [ ] FLUX erzeugt Standard- und High-Bild - [ ] Qwen-Profil wird nach FLUX wiederhergestellt diff --git a/docs/INSTALLATION.md b/docs/INSTALLATION.md index 3956079..fddca7e 100644 --- a/docs/INSTALLATION.md +++ b/docs/INSTALLATION.md @@ -65,6 +65,17 @@ Das bestehende `deploy/install.sh` installiert Router, Vision/Bild-Worker, Whisper und XTTS. Vor produktiver Verwendung müssen Modellpfade in der systemd-Datei gegen das lokale Manifest geprüft werden. +Der Installer erzeugt `/etc/mike-ai/router-api-key` (0600), installiert das +Profilregister und legt die atomare Zustandsablage an. Danach wird der Key +einmal manuell in die Secret-Stores der erlaubten Clients übernommen. Er darf +nicht im Terminal-Log, in Screenshots oder in Git dokumentiert werden. + +Der Router läuft aktuell als gehärteter Root-Dienst, weil er den +llama.cpp-Systemdienst und temporäre GPU-Worker kontrolliert. Das ist eine +bewusste Restabweichung. Ein späterer V3-Schritt soll den HTTP-Proxy als eigenen +Benutzer ausführen und nur Profil-/Hotswap-Befehle an einen fest +parametrisierten Root-Helper delegieren. + ## 8. Websuche TinySearch und SearXNG bleiben als einziges Docker-Teilsystem isoliert. Die diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 08a01d5..f5335c0 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -20,12 +20,19 @@ Clients verbinden sich mit: http://HOST:8081/v1 ``` +Als API-Key verwenden sie den Inhalt von `/etc/mike-ai/router-api-key` über +Bearer-Authentifizierung. Der Key gehört in den Secret-Store des Clients, +nicht in Chat, Repository oder URL. Eine Rotation erfolgt atomar durch +Ersetzen der Datei und Neustart des Routerdienstes. + Sie sollen nicht direkt Port 8080 verwenden, weil sie sonst Profilumschaltung, Vision, Bildgenerierung, STT und TTS umgehen. ## Status -- `GET /status`: Router, Profil, Upstream, aktive Jobs +- `GET /health`: Routerprozess lebt; bleibt bei geplantem Hotswap grün +- `GET /ready`: Router und Textmodell sind einsatzbereit +- `GET /status`: authentifizierter Detailstatus, Profil, Upstream, aktive Jobs - `GET /v1/models`: virtuelle Modelle - llama.cpp-Metriken: Port 8080, nur im administrativen Netz freigeben - systemd-Journal: nur Metadaten und Fehler prüfen; keine Promptinhalte sammeln @@ -35,6 +42,22 @@ Vision, Bildgenerierung, STT und TTS umgehen. Niemals Build, Quantisierung und Profil gleichzeitig ändern. Immer genau eine Variable ändern und anschließend denselben Benchmark ausführen. +Profilkontext und Alias werden zusätzlich in +`/etc/mike-ai/router-profiles.json` gepflegt. Änderungen an Override und +Registry gehören in denselben getesteten Commit; andernfalls verweigert die +Readiness bewusst die Freigabe. + +## Fehler- und Recovery-Verhalten + +- Ein Profilwechsel bricht ab, wenn laufende Chats nicht innerhalb des + Drain-Timeouts enden. Er beendet niemals absichtlich einen Chat. +- Nach einem Routerabsturz wird das letzte stabile Profil aus der atomaren + Zustandsdatei rekonstruiert. +- `/health = 200`, aber `/ready = 503` bedeutet: Router lebt, Modell ist noch + nicht bereit oder wird gerade gewechselt. +- `429` bedeutet, dass die Parallelitätsgrenze erreicht ist; der Client soll + mit Backoff erneut versuchen. + ## Kapazitätsregeln - Systempartition dauerhaft unter 85 Prozent halten. diff --git a/docs/ROUTER_V2_MIGRATION.md b/docs/ROUTER_V2_MIGRATION.md new file mode 100644 index 0000000..6443edb --- /dev/null +++ b/docs/ROUTER_V2_MIGRATION.md @@ -0,0 +1,58 @@ +# Router V2 – Migration und Kompatibilität + +## Ergebnis + +V2 behält die OpenAI-kompatible Basis-URL und die virtuellen Modelle +`qwen-fast`, `qwen-medium` und `qwen-long`. Bestehende Chat-, Tool-, Audio-, +Vision- und Bildpfade bleiben erhalten. Die Änderungen betreffen absichtlich +die Stellen, an denen der alte Router unsicher oder nicht deterministisch war. + +## Bewusste Änderungen + +| Alt | V2 | +|---|---| +| alle LAN-Clients ohne Authentifizierung | API-Key für alle fachlichen Endpunkte | +| `GET /fast` schaltet ein Modell | nur noch `POST /fast` (analog medium/long) | +| `/health` hängt am Modellzustand | `/health` = Prozess, `/ready` = Textmodell | +| Profil durch Textvergleich erkannt | Kontext + Alias aus Profilregister geprüft | +| Profilwechsel und Request konnten sich überholen | atomare Modell-Lease | +| Drain-Timeout beendete trotzdem das Modell | Wechsel wird sicher abgebrochen | +| Workerzustand nur im RAM | atomare Zustandsdatei und Startup-Recovery | +| beliebige Bild-URL | Data-URL, Größenlimit; Remote standardmäßig aus | +| unbegrenzte Parallelität und Bildablage | Request- und Retention-Limits | +| Auth-Header potenziell am Upstream | Router-Credentials werden entfernt | + +## Client-Migration + +1. Router-Key aus `/etc/mike-ai/router-api-key` ohne Anzeige in einen lokalen + Secret-Store des Clients übernehmen. +2. Basis-URL unverändert auf `http://HOST:8081/v1` lassen. +3. Den Key als OpenAI-API-Key/Bearer-Token konfigurieren. +4. `GET /v1/models` testen und anschließend einen kurzen Chat über + `qwen-fast` senden. +5. Automationen, die Profile per GET schalten, auf POST umstellen. +6. Überwachung auf `/health` (Liveness) und `/ready` (Readiness) aufteilen. + +## Sicheres Rollout + +V2 wird nicht blind über einen laufenden Router kopiert: + +1. Repository-Commit und aktuelle produktive Konfiguration sichern. +2. Profilregister gegen alle drei systemd-Overrides prüfen. +3. API-Key erzeugen und Clients vorbereiten. +4. Router installieren und zuerst lokal mit Key prüfen. +5. Fast, Medium und Long jeweils einmal schalten und Alias/Kontext prüfen. +6. Streaming und einen Tool Call testen. +7. Erst danach normale Clients auf V2 freigeben. + +Ein Rollback stellt Routerdateien und Unit aus dem Installationsbackup wieder +her. Der neu erzeugte API-Key und die Zustandsdatei enthalten keine +Modelldateien oder Chatdaten. + +## Verbleibender Architekturpunkt + +Der Routerprozess läuft derzeit als root, weil er den systemweiten llama.cpp- +Dienst und temporäre GPU-Worker steuert. Die Unit ist stark gehärtet, dennoch +ist das nicht das langfristige Ideal. V3 soll HTTP/API und privilegierte +Orchestrierung trennen: unprivilegierter Proxy plus kleiner Root-Helper mit +festen, nicht frei parametrisierbaren Aktionen. diff --git a/docs/SECURITY.md b/docs/SECURITY.md index 1e79464..f32b077 100644 --- a/docs/SECURITY.md +++ b/docs/SECURITY.md @@ -42,6 +42,20 @@ Empfohlene Trennung: - Firewall erlaubt nur bekannte Quellnetze. - Externe Suche erhält nur die tatsächliche Suchanfrage, keine Chat-Historie. +## Router-Grenze + +- Alle fachlichen Endpunkte verlangen einen mindestens 32 Zeichen langen, + zufälligen Router-Key. Der Dienst startet ohne gültigen Key nicht. +- `/health` und `/ready` sind die einzigen anonymen Endpunkte und geben nur + groben Betriebszustand aus. +- Authentifizierungsheader werden niemals an llama.cpp weitergereicht. +- Remote-Bild-URLs sind standardmäßig gesperrt. Data-URLs werden auf MIME-Typ, + Base64-Gültigkeit und 20 MiB Maximalgröße geprüft. +- Die Zahl gleichzeitiger Requests ist begrenzt; große Uploads sind global + begrenzt und generierte Bilder werden nach Alter, Anzahl und Größe bereinigt. +- Crash-Recovery beendet keine PID nur aufgrund einer Zahl, sondern verlangt + zusätzlich einen erwarteten Prozessmarker in `/proc//cmdline`. + ## Schreibaktionen Jede destruktive oder persistente Aktion verwendet: diff --git a/router/ai_profile_router.py b/router/ai_profile_router.py index c141877..acc40b0 100755 --- a/router/ai_profile_router.py +++ b/router/ai_profile_router.py @@ -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() diff --git a/router/router_profiles.json b/router/router_profiles.json new file mode 100644 index 0000000..2521651 --- /dev/null +++ b/router/router_profiles.json @@ -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" + } + } +} diff --git a/router/router_support.py b/router/router_support.py new file mode 100644 index 0000000..23ec283 --- /dev/null +++ b/router/router_support.py @@ -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