diff --git a/README.md b/README.md index 6e0b9a8..6c00c4f 100644 --- a/README.md +++ b/README.md @@ -2,8 +2,10 @@ Kleiner OpenAI-kompatibler Proxy (Python, nur Standardbibliothek) vor einem lokalen llama.cpp-Server. Er leitet normale OpenAI-Requests transparent -weiter (Streaming, Tool Calls, JSON) und schaltet zwischen drei festen -llama.cpp-Profilen um. +weiter (Streaming, Tool Calls, JSON), schaltet zwischen drei festen +llama.cpp-Profilen um und orchestriert lokale Bildgenerierung mit +FLUX.2 [klein] 4B Base (GPU-Hotswap: Qwen stoppen → FLUX laden → Bild +→ FLUX entladen → Qwen wiederherstellen). ## Zielsystem @@ -32,6 +34,9 @@ llama.cpp-Profilen um. | `GET /status` | Aktives Profil, Upstream-Zustand, Modell, Kontext, Uptime | | `POST /fast` `/medium` `/long` | Profilwechsel (auch `GET` möglich) | | `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/` | PNG-Download (nur `images/`-Verzeichnis, validiert) | | alles andere | Transparente Weiterleitung an llama.cpp | ### Verhalten @@ -52,14 +57,125 @@ llama.cpp-Profilen um. schaltbar. - **Fehlerformat**: OpenAI-kompatibel (`{"error": {"message", "type", "code"}}`). +## Bildgenerierung (FLUX.2 [klein] 4B Base) + +Der Router orchestriert lokale Bildgenerierung mit +`black-forest-labs/FLUX.2-klein-base-4B` (Apache 2.0, ~13 GB, bf16 + +CPU-Offload). Da Qwen (llama.cpp) und FLUX denselben GPU/VRAM teilen, macht +der Router einen **GPU-Hotswap**: + +1. Zentrales GPU/Modell-Lock übernehmen (Profilwechsel und Bild teilen sich + dasselbe Lock → kein Race). +2. Aktives Qwen-Profil merken. +3. `mike-ai-llama-ui.service` stoppen, warten bis Port + VRAM frei sind. +4. FLUX-Worker starten (eigener Prozess, `Flux2KleinPipeline`, bf16 + + `enable_model_cpu_offload()`), Bild generieren, PNG speichern. +5. Worker **beenden** (nicht nur entladen), VRAM-Freiheit verifizieren. +6. Vorheriges Qwen-Profil exakt wiederherstellen, Readiness-Check + (Modell geladen + Kontext passt). +7. Erst dann antworten und das GPU-Lock freigeben. + +### Endpunkt `POST /v1/images/generations` + +OpenAI-kompatibel. Unterstützt `prompt`, `size`, `n`, `seed`, `quality`, +`response_format`. + +| Parameter | Werte | Default | +|---|---|---| +| `prompt` | Text (Pflicht) | – | +| `size` | `1024x1024`, `1536x1024`, `1024x1536`, `1920x1088`, `1088x1920` | `1024x1024` | +| `n` | 1–4 | 1 | +| `seed` | int (reproduzierbar) | zufällig | +| `quality` | `standard` (30 Steps), `high` (50 Steps) | `standard` | +| `response_format` | `url` (Default), `b64_json` | `url` | + +Beispiel: + +```bash +curl -s http://192.168.1.196:8081/v1/images/generations \ + -H 'Content-Type: application/json' \ + -d '{"prompt":"ein roter Würfel auf weißem Grund","size":"1024x1024","quality":"standard"}' +``` + +Die Antwort enthält `data[].url` (absolute URL, über den Router abrufbar) +und `data[].b64_json` (optional). Jedes Bild wird unter +`/opt/mike-ai/ai-profile-router/images/` gespeichert (kollisionsfreie Namen, +`img---.png`) und ist über `GET /images/` +abrufbar. + +### Verhalten während eines Bild-Jobs + +- **Chat-Requests warten** (kein 502): Der Router merkt sich, dass Qwen + vorübergehend nicht verfügbar ist (`qwen.available=false`), und Chat-Requests + warten, bis Qwen wieder bereit ist (Timeout `CHAT_WAIT_TIMEOUT`, Default + 300 s). So gibt es keine `502 llama.cpp nicht erreichbar` während des + Hotswaps. +- **Profilwechsel warten**: Ein Profilwechsel während eines Bild-Jobs + blockiert auf dem GPU-Lock, bis der Bild-Job fertig ist (kein Race). +- **`/status`** zeigt den aktuellen Zustand: `image.phase` (`idle`, + `stopping-qwen`, `loading-image`, `generating`, `unloading-image`, + `restoring-qwen`), `image.worker`, `image.model_loaded`, + `image.last_image`, `image.last_seconds`, `image.last_error`, + `qwen.available`, `qwen.active_chats`. + +### Recovery (robust) + +- **`try/finally`**: Qwen wird **immer** wiederhergestellt, egal ob die + Bildgenerierung erfolgreich war, fehlgeschlagen ist (OOM, Python-Fehler, + ungültiger Prompt, Speichern-Fehler, Client-Disconnect, Timeout) oder der + Worker abstürzt. +- **Worker-Beendigung**: Nach jedem Job wird der Worker beendet (SIGTERM → + SIGKILL), nicht nur entladen. So wird der VRAM (inkl. CUDA-Kontext) frei. +- **VRAM-Check**: Nach dem Worker-Beenden wartet der Router, bis der VRAM + unter 1000 MiB fällt (`nvidia-smi`), bevor Qwen neu startet. +- **Qwen-Readiness**: Nach dem Neustart wartet der Router, bis llama.cpp + erreichbar ist, das Modell geladen ist und der Kontext zum Profil passt. +- **`qwen.available`**: Bleibt `false`, wenn die Wiederherstellung fehlschlägt + (Chat-Requests warten weiter, statt 502 zu liefern). Der Fehler wird in + `image.last_error` und im Log protokolliert. + +### Benchmarks (RTX 5080, 16 GB, CPU-Offload, gemessen) + +| Auflösung | Steps | Zeit | Peak-VRAM (torch) | +|---|---|---|---| +| 512×512 | 10 | ~9.3 s | ~8.4 GB | +| 1024×1024 | 30 | ~31.3 s | ~8.4 GB | +| 1024×1024 | 50 | ~45.3 s | ~8.4 GB | +| 1920×1088 | 50 | ~91 s | ~8.9 GB | + +**Entscheidung:** `standard` = 30 Steps (Default, ~31 s bei 1024×1024), +`high` = 50 Steps (maximale Qualität, ~45 s bei 1024×1024). Ab 20–30 Steps +ist der Qualitätsgewinn bei einfachen Motiven gering; 50 Steps lohnt sich +für komplexe Szenen. + +**Hinweis:** FLUX.2 [klein] 4B Base passt **nicht** vollständig GPU-resident +in 16 GB (OOM bei ~15.5 GB). Deshalb wird `enable_model_cpu_offload()` +verwendet (Modelle werden pro Layer zwischen CPU und GPU gewechselt). + +**Hotswap-Gesamtzeit:** Ein vollständiger Bild-Job (Qwen stoppen → FLUX laden +→ Bild → FLUX entladen → Qwen wiederherstellen) dauert ~41–42 s bei +1024×1024 / 30 Steps (davon ~31 s Generierung, ~10 s Qwen-Stop/Start + +VRAM-Check). + +### Erforderliche Python-Pakete (im Venv) + +- `torch` (2.11.0+cu128, CUDA 12.8) +- `diffusers` (0.40.0.dev0, für `Flux2KleinPipeline`) +- `transformers` (5.15.0) +- `accelerate` (1.14.0, für `enable_model_cpu_offload()`) + +Das Venv liegt unter `/opt/mike-ai/ai-profile-router/venv/` und wird von +`install.sh` automatisch angelegt/aktualisiert. + ## Repository-Struktur ``` router/ai_profile_router.py # der Router (einzige Laufzeit-Datei) +router/image_worker.py # FLUX-Worker (eigener Prozess, JSON-Protokoll) deploy/mike-ai-profile-router.service # systemd-Unit deploy/install.sh # läuft auf dem Zielsystem (per SSH) deploy/deploy.sh # läuft lokal: SCP + SSH -dev/ # lokale Tests (Mock-llama.cpp, Fake-Profil-Skript) +dev/ # lokale Tests (Mock-llama.cpp, Mock-Worker, Benchmarks) ``` Entwicklungsdateien (`dev/`) und Deployment-Dateien (`router/`, `deploy/`) @@ -75,12 +191,17 @@ Voraussetzung: SSH-Key `~/.ssh/lmstudio_unraid` (bereits vorhanden). Das Skript: -1. Überträgt `ai_profile_router.py`, `install.sh` und die systemd-Unit per - SCP nach `/tmp/ai-profile-router/` auf dem Zielsystem. +1. Überträgt `ai_profile_router.py`, `image_worker.py`, `install.sh` und die + systemd-Unit per SCP nach `/tmp/ai-profile-router/` auf dem Zielsystem. 2. Führt `install.sh` per SSH aus, das: - den alten Router (`mike-ai-local-llm-router.service` + `/opt/mike-ai/local-llm-router`) **mit Backup** entfernt, - - den neuen Router nach `/opt/mike-ai/ai-profile-router/` installiert, + - den neuen Router + Worker nach `/opt/mike-ai/ai-profile-router/` + installiert, + - ein Python-Venv mit `torch`, `diffusers`, `transformers`, `accelerate` + anlegt (nur wenn noch nicht vorhanden), + - das FLUX-Modell nach `/opt/mike-ai/models/FLUX.2-klein-base-4B` lädt + (nur wenn noch nicht vorhanden, ~15 GB), - `mike-ai-profile-router.service` aktiviert (Start beim Boot) und startet, - `GET /status` verifiziert. @@ -100,6 +221,14 @@ Die Installation ist idempotent (Update = erneut ausführen). | `CONNECT_TIMEOUT` | `10` | Connect-Timeout Upstream (s) | | `POLL_INTERVAL` | `2` | Polling-Intervall (s) | | `LOG_LEVEL` | `INFO` | Logging-Level | +| `LLAMA_SERVICE` | `mike-ai-llama-ui.service` | llama.cpp-Service (für Bild-Hotswap) | +| `IMAGE_WORKER` | `/image_worker.py` | FLUX-Worker-Skript | +| `IMAGE_PYTHON` | `sys.executable` | Python für den Worker (venv mit torch) | +| `IMAGE_DIR` | `/opt/mike-ai/ai-profile-router/images` | Bild-Speicherort | +| `IMAGE_WORKER_LOG` | `/opt/mike-ai/ai-profile-router/worker.log` | Worker-Log | +| `IMAGE_GEN_TIMEOUT` | `600` | Timeout pro Bild (s) | +| `IMAGE_VRAM_FREE_TIMEOUT` | `120` | Warten auf VRAM-Freiheit (s) | +| `CHAT_WAIT_TIMEOUT` | `300` | Chat wartet auf Qwen (s) | ## Lokale Tests @@ -107,10 +236,13 @@ Die Installation ist idempotent (Update = erneut ausführen). ./dev/test_local.sh ``` -Startet einen Mock-llama.cpp und den Router mit einem Fake-Profil-Skript und -prüft: `/v1/models`, `/status`, Forwarding, Streaming, Tool Calls, -Profilwechsel (fast→medium→fast), virtuelles Modell triggert Wechsel, -ungültige Profile, Upstream down → 502, Recovery. +Startet einen Mock-llama.cpp, einen Mock-Bild-Worker und den Router mit einem +Fake-Profil-Skript und prüft: `/v1/models`, `/status`, Forwarding, Streaming, +Tool Calls, Profilwechsel (fast→medium→fast), virtuelles Modell triggert +Wechsel, ungültige Profile, Upstream down → 502, Recovery, **Bildgenerierung** +(`standard`→30 Steps, `high`→50 Steps, Validierung, Image-Fehler→Qwen +wiederhergestellt, Fast/Medium/Long→Image→gleiches Profil, `/status` während +Bild-Job, paralleler Chat während Bild-Job wartet statt 502). ## Betrieb diff --git a/deploy/deploy.sh b/deploy/deploy.sh index 018388a..fc89de8 100755 --- a/deploy/deploy.sh +++ b/deploy/deploy.sh @@ -9,7 +9,7 @@ SSH_KEY="${SSH_KEY:-$HOME/.ssh/lmstudio_unraid}" STAGE="/tmp/ai-profile-router-$$" mkdir -p "$STAGE" -cp router/ai_profile_router.py deploy/install.sh deploy/mike-ai-profile-router.service "$STAGE/" +cp router/ai_profile_router.py router/image_worker.py deploy/install.sh deploy/mike-ai-profile-router.service "$STAGE/" echo "== Übertrage Dateien nach ${TARGET}:/tmp/ai-profile-router/" ssh -i "$SSH_KEY" "$TARGET" 'mkdir -p /tmp/ai-profile-router' diff --git a/deploy/install.sh b/deploy/install.sh index f8c9c40..8124836 100755 --- a/deploy/install.sh +++ b/deploy/install.sh @@ -1,7 +1,7 @@ #!/bin/bash # AI Profile Router – Installation/Update auf dem Zielsystem. # Wird als root auf dem Zielsystem ausgeführt (per SSH, vgl. deploy.sh). -# Erwartet ai_profile_router.py und mike-ai-profile-router.service +# Erwartet ai_profile_router.py, image_worker.py und mike-ai-profile-router.service # im selben Verzeichnis wie dieses Skript. set -euo pipefail @@ -11,6 +11,9 @@ SERVICE=mike-ai-profile-router.service OLD_SERVICE=mike-ai-local-llm-router.service OLD_DIR=/opt/mike-ai/local-llm-router BACKUP_DIR=/opt/mike-ai/.backup-ai-profile-router-$(date +%Y%m%d-%H%M%S) +VENV="$INSTALL_DIR/venv" +MODEL_DIR=/opt/mike-ai/models/FLUX.2-klein-base-4B +IMAGE_DIR="$INSTALL_DIR/images" echo "== AI Profile Router: Installation/Update ==" @@ -29,16 +32,47 @@ else fi # --- 2. Neue Dateien installieren ------------------------------------------- -mkdir -p "$INSTALL_DIR" +mkdir -p "$INSTALL_DIR" "$IMAGE_DIR" install -m 0755 "$DIR/ai_profile_router.py" "$INSTALL_DIR/ai_profile_router.py" +install -m 0755 "$DIR/image_worker.py" "$INSTALL_DIR/image_worker.py" install -m 0644 "$DIR/${SERVICE}" "/etc/systemd/system/${SERVICE}" -# --- 3. Service aktivieren und starten --------------------------------------- +# --- 3. Python-Venv mit Bild-Abhängigkeiten --------------------------------- +if [ ! -x "$VENV/bin/python" ]; then + echo "-- Erstelle Python-Venv in $VENV" + python3 -m venv "$VENV" +fi +echo "-- Installiere/aktualisiere Bild-Abhängigkeiten (torch, diffusers, ...)" +"$VENV/bin/pip" install --quiet --upgrade pip +"$VENV/bin/pip" install --quiet \ + torch \ + diffusers \ + transformers \ + accelerate + +# --- 4. FLUX-Modell (nur wenn noch nicht vorhanden) -------------------------- +if [ -f "$MODEL_DIR/model_index.json" ]; then + echo "-- FLUX-Modell vorhanden: $MODEL_DIR" +else + echo "-- Lade FLUX.2-klein-base-4B nach $MODEL_DIR (kann dauern)" + "$VENV/bin/python" - <<'PY' +import os +from huggingface_hub import snapshot_download +snapshot_download( + repo_id="black-forest-labs/FLUX.2-klein-base-4B", + local_dir="/opt/mike-ai/models/FLUX.2-klein-base-4B", + local_dir_use_symlinks=False, +) +print("Modell-Download abgeschlossen") +PY +fi + +# --- 5. Service aktivieren und starten --------------------------------------- systemctl daemon-reload systemctl enable "$SERVICE" systemctl restart "$SERVICE" -# --- 4. Verifikation ---------------------------------------------------------- +# --- 6. Verifikation ---------------------------------------------------------- sleep 1 if ! systemctl is-active --quiet "$SERVICE"; then echo "-- FEHLER: Service läuft nicht" >&2 diff --git a/deploy/mike-ai-profile-router.service b/deploy/mike-ai-profile-router.service index d6520d9..1a7e021 100644 --- a/deploy/mike-ai-profile-router.service +++ b/deploy/mike-ai-profile-router.service @@ -1,11 +1,11 @@ [Unit] -Description=Mike AI Profile Router (OpenAI-kompatibler Proxy, Port 8081) +Description=Mike AI Profile Router (OpenAI-kompatibler Proxy + Bild-Orchestrierung, Port 8081) After=network-online.target Wants=network-online.target [Service] Type=simple -ExecStart=/usr/bin/python3 /opt/mike-ai/ai-profile-router/ai_profile_router.py +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 Environment=ROUTER_HOST=0.0.0.0 @@ -15,6 +15,13 @@ Environment=PROFILE_SCRIPT=/usr/local/bin/llama-profile Environment=PROFILE_DIR=/etc/systemd/system/mike-ai-llama-ui.service.d Environment=SWITCH_TIMEOUT=600 Environment=REQUEST_TIMEOUT=600 +Environment=LLAMA_SERVICE=mike-ai-llama-ui.service +Environment=IMAGE_WORKER=/opt/mike-ai/ai-profile-router/image_worker.py +Environment=IMAGE_PYTHON=/opt/mike-ai/ai-profile-router/venv/bin/python +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 NoNewPrivileges=true PrivateTmp=true diff --git a/dev/fake-systemctl.sh b/dev/fake-systemctl.sh new file mode 100755 index 0000000..6c34aa8 --- /dev/null +++ b/dev/fake-systemctl.sh @@ -0,0 +1,49 @@ +#!/bin/bash +# Fake systemctl für lokale Tests: verwaltet den Mock-llama.cpp-Prozess. +# Simuliert: systemctl stop|start|status +# +# Umgebungsvariablen (vom Router geerbt): +# FAKE_SYSTEMD_PIDFILE PID-Datei des Mocks (Default /tmp/mock_upstream_pid) +# FAKE_SYSTEMD_PORT Mock-Port (Default 18080) +# FAKE_SYSTEMD_PROFILE_DIR Profil-Dir für den Mock +# FAKE_SYSTEMD_MOCK Mock-Skript (Default dev/mock_upstream.py) +# FAKE_SYSTEMD_LOG Log-Datei (Default /tmp/mock_upstream_fake.log) + +CMD="${1:-}" +PIDFILE="${FAKE_SYSTEMD_PIDFILE:-/tmp/mock_upstream_pid}" +PORT="${FAKE_SYSTEMD_PORT:-18080}" +PROFILE_DIR="${FAKE_SYSTEMD_PROFILE_DIR:-}" +MOCK="${FAKE_SYSTEMD_MOCK:-dev/mock_upstream.py}" +LOG="${FAKE_SYSTEMD_LOG:-/tmp/mock_upstream_fake.log}" + +is_running() { + [ -f "$PIDFILE" ] && kill -0 "$(cat "$PIDFILE")" 2>/dev/null +} + +case "$CMD" in + stop) + if is_running; then + kill "$(cat "$PIDFILE")" 2>/dev/null || true + rm -f "$PIDFILE" + for _ in $(seq 1 50); do + is_running || break + sleep 0.1 + done + fi + exit 0 + ;; + start) + if ! is_running; then + MOCK_PROFILE_DIR="$PROFILE_DIR" MOCK_PORT="$PORT" \ + python3 "$MOCK" >>"$LOG" 2>&1 & + echo $! > "$PIDFILE" + fi + exit 0 + ;; + status) + is_running && exit 0 || exit 3 + ;; + *) + exit 0 + ;; +esac diff --git a/dev/flux_benchmark_run.sh b/dev/flux_benchmark_run.sh new file mode 100644 index 0000000..25d7ecc --- /dev/null +++ b/dev/flux_benchmark_run.sh @@ -0,0 +1,66 @@ +#!/bin/bash +# FLUX-Benchmark-Wrapper mit GARANTIERTER Qwen-Recovery. +# +# Usage: flux_benchmark_run.sh +# +# Ablauf: +# 1. Aktives Qwen-Profil aus override.conf merken (bytegenauer Vergleich) +# 2. llama.cpp stoppen, VRAM-Abgabe verifizieren +# 3. Benchmark-Skript ausführen +# 4. IMMER (trap EXIT): llama.cpp wieder starten, Readiness prüfen +set -u + +SERVICE=mike-ai-llama-ui.service +PROFILE_DIR=/etc/systemd/system/mike-ai-llama-ui.service.d +PY=/opt/mike-ai/ai-profile-router/venv/bin/python +SCRIPT="${1:?Usage: flux_benchmark_run.sh }" +ROUTER_STATUS=http://127.0.0.1:8081/status + +# Aktives Profil bestimmen: override.conf bytegenau mit den Profil-Dateien +# vergleichen (robuster als String-Matching). +PROFILE="" +for p in fast medium long; do + if cmp -s "$PROFILE_DIR/override.conf" "$PROFILE_DIR/profile-$p.conf.disabled" 2>/dev/null; then + PROFILE=$p + break + fi +done +if [ -z "$PROFILE" ]; then + echo "FEHLER: aktives Profil nicht erkannt (override.conf passt zu keinem Profil-File)" + exit 1 +fi +echo "=== Aktives Qwen-Profil: $PROFILE (wird nach dem Test wiederhergestellt) ===" + +restore() { + echo + echo "=== RECOVERY: stelle Profil $PROFILE wieder her ===" + systemctl start "$SERVICE" 2>&1 || echo "systemctl start fehlgeschlagen" + for i in $(seq 1 150); do + if curl -sf "$ROUTER_STATUS" 2>/dev/null | python3 -c ' +import json, sys +d = json.load(sys.stdin) +u = d["upstream"] +assert u["reachable"] and u["model"], d +print("ready:", u["model"], "ctx", u["ctx"]) +' 2>/dev/null; then + echo "=== RECOVERY OK: llama.cpp ist wieder inference-ready ===" + return 0 + fi + sleep 2 + done + echo "=== RECOVERY FEHLGESCHLAGEN: bitte manuell prüfen (systemctl status $SERVICE) ===" + return 1 +} +trap restore EXIT + +echo "=== Stoppe $SERVICE ===" +systemctl stop "$SERVICE" +sleep 3 +echo "=== VRAM nach Stop (MiB) ===" +nvidia-smi --query-gpu=memory.used --format=csv,noheader,nounits + +echo "=== Starte Benchmark: $SCRIPT ===" +"$PY" "$SCRIPT" +RC=$? +echo "=== Benchmark-Exit-Code: $RC ===" +exit $RC diff --git a/dev/flux_gpu_benchmark.py b/dev/flux_gpu_benchmark.py new file mode 100644 index 0000000..220795b --- /dev/null +++ b/dev/flux_gpu_benchmark.py @@ -0,0 +1,95 @@ +#!/usr/bin/env python3 +"""FLUX.2 [klein] 4B Base – GPU-Resident-Benchmark (OHNE CPU-Offload). + +Ziel: Prüfen, ob das Modell vollständig auf der RTX 5080 (16 GB) läuft. + +Messen: + - from_pretrained-Zeit + - .to("cuda")-Zeit + - VRAM (nvidia-smi + torch.cuda.memory_allocated / max_memory_allocated) + - Generierungszeit, Peak-VRAM pro Auflösung + +Auflösungen: 512x512 (10 steps) → 1024x1024 (50 steps) → 1920x1088 (50 steps) +Bei OOM wird abgebrochen (CUDA-Kontext danach nicht mehr verlässlich). +""" + +import os +import subprocess +import time + +os.environ["HF_HUB_DISABLE_PROGRESS_BARS"] = "1" +os.environ["TRANSFORMERS_VERBOSITY"] = "error" +os.environ["TOKENIZERS_PARALLELISM"] = "false" + +MODEL = "/opt/mike-ai/models/FLUX.2-klein-base-4B" +PROMPT = ("A detailed photograph of a red cube on a white marble table, " + "soft studio lighting, shallow depth of field") + +# (breite, hoehe, steps, seed, name) +CASES = [ + (512, 512, 10, 0, "512x512-10s"), + (1024, 1024, 50, 0, "1024x1024-50s"), + (1920, 1088, 50, 0, "1920x1088-50s"), +] + + +def nvidia_vram() -> int: + out = subprocess.check_output( + ["nvidia-smi", "--query-gpu=memory.used", + "--format=csv,noheader,nounits"]).decode().strip() + return int(out.split()[0]) + + +def main() -> None: + import torch + print(f"torch {torch.__version__} | cuda {torch.cuda.is_available()} " + f"| {torch.cuda.get_device_name(0)}", flush=True) + from diffusers import Flux2KleinPipeline + + # --- Laden (zuerst auf CPU, dann vollständig auf GPU) --- + t0 = time.monotonic() + pipe = Flux2KleinPipeline.from_pretrained(MODEL, torch_dtype=torch.bfloat16) + t_load = time.monotonic() - t0 + print(f"[load] from_pretrained: {t_load:.1f} s", flush=True) + + t1 = time.monotonic() + pipe.to("cuda") + torch.cuda.synchronize() + t_to = time.monotonic() - t1 + print(f"[load] .to(cuda): {t_to:.1f} s", flush=True) + print(f"[load] VRAM nvidia-smi: {nvidia_vram()} MiB | " + f"torch allocated: {torch.cuda.memory_allocated() / 1e9:.2f} GB", + flush=True) + + # --- Generierung --- + for width, height, steps, seed, name in CASES: + out = f"/tmp/flux-bench-{name}.png" + torch.cuda.synchronize() + torch.cuda.reset_peak_memory_stats() + t = time.monotonic() + try: + img = pipe( + prompt=PROMPT, + height=height, + width=width, + guidance_scale=4.0, + num_inference_steps=steps, + generator=torch.Generator(device="cuda").manual_seed(seed), + ).images[0] + dt = time.monotonic() - t + img.save(out) + peak = torch.cuda.max_memory_allocated() / 1e9 + print(f"[gen] {name}: {dt:.1f} s | peak torch {peak:.2f} GB | " + f"nvidia-smi {nvidia_vram()} MiB | {out}", flush=True) + except Exception as e: # noqa: BLE001 + dt = time.monotonic() - t + print(f"[gen] {name}: FEHLER nach {dt:.1f} s: {e!r}", flush=True) + if "out of memory" in str(e).lower(): + print("[gen] OOM – Abbruch, größere Auflösungen nicht getestet", + flush=True) + break + print("DONE", flush=True) + + +if __name__ == "__main__": + main() diff --git a/dev/flux_gpu_benchmark.sh b/dev/flux_gpu_benchmark.sh new file mode 100644 index 0000000..79fced3 --- /dev/null +++ b/dev/flux_gpu_benchmark.sh @@ -0,0 +1,66 @@ +#!/bin/bash +# FLUX GPU-Resident-Benchmark mit GARANTIERTER Qwen-Recovery. +# +# Ablauf: +# 1. Aktives Qwen-Profil aus override.conf merken +# 2. llama.cpp stoppen, VRAM-Abgabe verifizieren +# 3. Benchmark ausführen (Python, GPU-resident, ohne CPU-Offload) +# 4. IMMER (trap EXIT): llama.cpp wieder starten, Readiness prüfen +# +# Usage: flux_gpu_benchmark.sh +set -u + +SERVICE=mike-ai-llama-ui.service +PROFILE_DIR=/etc/systemd/system/mike-ai-llama-ui.service.d +PY=/opt/mike-ai/ai-profile-router/venv/bin/python +SCRIPT="$(cd "$(dirname "$0")" && pwd)/flux_gpu_benchmark.py" +ROUTER_STATUS=http://127.0.0.1:8081/status + +# Aktives Profil bestimmen: override.conf bytegenau mit den Profil-Dateien +# vergleichen (robuster als String-Matching). +PROFILE="" +for p in fast medium long; do + if cmp -s "$PROFILE_DIR/override.conf" "$PROFILE_DIR/profile-$p.conf.disabled" 2>/dev/null; then + PROFILE=$p + break + fi +done +if [ -z "$PROFILE" ]; then + echo "FEHLER: aktives Profil nicht erkannt (override.conf passt zu keinem Profil-File)" + exit 1 +fi +echo "=== Aktives Qwen-Profil: $PROFILE (wird nach dem Test wiederhergestellt) ===" + +restore() { + echo + echo "=== RECOVERY: stelle Profil $PROFILE wieder her ===" + systemctl start "$SERVICE" 2>&1 || echo "systemctl start fehlgeschlagen" + for i in $(seq 1 150); do + if curl -sf "$ROUTER_STATUS" 2>/dev/null | python3 -c ' +import json, sys +d = json.load(sys.stdin) +u = d["upstream"] +assert u["reachable"] and u["model"], d +print("ready:", u["model"], "ctx", u["ctx"]) +' 2>/dev/null; then + echo "=== RECOVERY OK: llama.cpp ist wieder inference-ready ===" + return 0 + fi + sleep 2 + done + echo "=== RECOVERY FEHLGESCHLAGEN: bitte manuell prüfen (systemctl status $SERVICE) ===" + return 1 +} +trap restore EXIT + +echo "=== Stoppe $SERVICE ===" +systemctl stop "$SERVICE" +sleep 3 +echo "=== VRAM nach Stop (MiB) ===" +nvidia-smi --query-gpu=memory.used --format=csv,noheader,nounits + +echo "=== Starte Benchmark ===" +"$PY" "$SCRIPT" +RC=$? +echo "=== Benchmark-Exit-Code: $RC ===" +exit $RC diff --git a/dev/flux_offload_benchmark.py b/dev/flux_offload_benchmark.py new file mode 100644 index 0000000..d2a80af --- /dev/null +++ b/dev/flux_offload_benchmark.py @@ -0,0 +1,83 @@ +#!/usr/bin/env python3 +"""FLUX.2 [klein] 4B Base – Benchmark MIT enable_model_cpu_offload(). + +Offizieller Pfad der Modellkarte ("runs on consumer hardware, with as +little as 13GB VRAM"). Keine Qualitätsreduktion – nur langsamer +(Weights wandern pro Layer zwischen CPU und GPU). + +Messen: Load-Zeit, Generierungszeit, Peak-VRAM pro Auflösung. +""" + +import os +import subprocess +import time + +os.environ["HF_HUB_DISABLE_PROGRESS_BARS"] = "1" +os.environ["TRANSFORMERS_VERBOSITY"] = "error" +os.environ["TOKENIZERS_PARALLELISM"] = "false" + +MODEL = "/opt/mike-ai/models/FLUX.2-klein-base-4B" +PROMPT = ("A detailed photograph of a red cube on a white marble table, " + "soft studio lighting, shallow depth of field") + +CASES = [ + (512, 512, 10, 0, "512x512-10s"), + (1024, 1024, 50, 0, "1024x1024-50s"), + (1920, 1088, 50, 0, "1920x1088-50s"), +] + + +def nvidia_vram() -> int: + out = subprocess.check_output( + ["nvidia-smi", "--query-gpu=memory.used", + "--format=csv,noheader,nounits"]).decode().strip() + return int(out.split()[0]) + + +def main() -> None: + import torch + print(f"torch {torch.__version__} | cuda {torch.cuda.is_available()} " + f"| {torch.cuda.get_device_name(0)}", flush=True) + from diffusers import Flux2KleinPipeline + + t0 = time.monotonic() + pipe = Flux2KleinPipeline.from_pretrained(MODEL, torch_dtype=torch.bfloat16) + t_load = time.monotonic() - t0 + print(f"[load] from_pretrained: {t_load:.1f} s", flush=True) + + t1 = time.monotonic() + pipe.enable_model_cpu_offload() + t_off = time.monotonic() - t1 + print(f"[load] enable_model_cpu_offload: {t_off:.1f} s", flush=True) + print(f"[load] VRAM nvidia-smi (idle): {nvidia_vram()} MiB", flush=True) + + for width, height, steps, seed, name in CASES: + out = f"/tmp/flux-bench-offload-{name}.png" + torch.cuda.synchronize() + torch.cuda.reset_peak_memory_stats() + t = time.monotonic() + try: + img = pipe( + prompt=PROMPT, + height=height, + width=width, + guidance_scale=4.0, + num_inference_steps=steps, + generator=torch.Generator(device="cuda").manual_seed(seed), + ).images[0] + dt = time.monotonic() - t + img.save(out) + peak = torch.cuda.max_memory_allocated() / 1e9 + print(f"[gen] {name}: {dt:.1f} s | peak torch {peak:.2f} GB | " + f"nvidia-smi {nvidia_vram()} MiB | {out}", flush=True) + except Exception as e: # noqa: BLE001 + dt = time.monotonic() - t + print(f"[gen] {name}: FEHLER nach {dt:.1f} s: {e!r}", flush=True) + if "out of memory" in str(e).lower(): + print("[gen] OOM – Abbruch", flush=True) + break + print("DONE", flush=True) + + +if __name__ == "__main__": + main() diff --git a/dev/flux_quality_compare.py b/dev/flux_quality_compare.py new file mode 100644 index 0000000..cf0fa57 --- /dev/null +++ b/dev/flux_quality_compare.py @@ -0,0 +1,92 @@ +#!/usr/bin/env python3 +"""FLUX.2 [klein] 4B Base – Qualitätsvergleich 30 vs. 50 Steps. + +Identischer Prompt, identischer Seed, identische Parameter – nur +num_inference_steps variiert (30 vs. 50). CPU-Offload, Base-Modell, +keine anderen Änderungen. +""" + +import os +import subprocess +import time + +os.environ["HF_HUB_DISABLE_PROGRESS_BARS"] = "1" +os.environ["TRANSFORMERS_VERBOSITY"] = "error" +os.environ["TOKENIZERS_PARALLELISM"] = "false" + +MODEL = "/opt/mike-ai/models/FLUX.2-klein-base-4B" +SEED = 42 +WIDTH = HEIGHT = 1024 +GUIDANCE = 4.0 + +PROMPT = ( + "Ultra-realistic cinematic photograph of a woman in her early thirties " + "sitting at a small outdoor café table in a rainy European city at night. " + "Natural detailed skin texture with pores and subtle imperfections, " + "realistic eyes and individual strands of wet hair, both hands clearly " + "visible holding a ceramic coffee cup with anatomically correct fingers. " + "She wears a dark wool coat over a finely textured knitted sweater. " + "Raindrops on the table and glass surfaces, wet pavement reflecting warm " + "café lights and cool blue street lighting, realistic depth of field, " + "pedestrians and bicycles in the detailed background, complex reflections " + "in windows and puddles. On the café window behind her is a clearly " + "readable handwritten sign saying exactly: 'CAFÉ LUMIÈRE – OPEN UNTIL " + "MIDNIGHT'. A small newspaper lies on the table with the clearly readable " + "headline 'BERLIN AFTER DARK'. Photorealistic professional full-frame " + "camera photograph, natural color grading, physically plausible lighting, " + "realistic materials, fine micro-detail, no plastic skin, no illustration, " + "no CGI look." +) + +CASES = [ + (30, "/tmp/flux-quality-30.png"), + (50, "/tmp/flux-quality-50.png"), +] + + +def nvidia_vram() -> int: + out = subprocess.check_output( + ["nvidia-smi", "--query-gpu=memory.used", + "--format=csv,noheader,nounits"]).decode().strip() + return int(out.split()[0]) + + +def main() -> None: + import torch + print(f"torch {torch.__version__} | {torch.cuda.get_device_name(0)}", + flush=True) + print(f"SEED={SEED} | {WIDTH}x{HEIGHT} | guidance={GUIDANCE}", flush=True) + print(f"PROMPT={PROMPT!r}", flush=True) + from diffusers import Flux2KleinPipeline + + t0 = time.monotonic() + pipe = Flux2KleinPipeline.from_pretrained(MODEL, torch_dtype=torch.bfloat16) + pipe.enable_model_cpu_offload() + print(f"[load] ready in {time.monotonic() - t0:.1f} s", flush=True) + + for steps, out in CASES: + torch.cuda.synchronize() + torch.cuda.reset_peak_memory_stats() + t = time.monotonic() + try: + img = pipe( + prompt=PROMPT, + height=HEIGHT, + width=WIDTH, + guidance_scale=GUIDANCE, + num_inference_steps=steps, + generator=torch.Generator(device="cuda").manual_seed(SEED), + ).images[0] + dt = time.monotonic() - t + img.save(out) + peak = torch.cuda.max_memory_allocated() / 1e9 + print(f"[gen] {steps} steps: {dt:.1f} s | peak {peak:.2f} GB | " + f"nvidia-smi {nvidia_vram()} MiB | seed={SEED} | {out}", + flush=True) + except Exception as e: # noqa: BLE001 + print(f"[gen] {steps} steps: FEHLER: {e!r}", flush=True) + print("DONE", flush=True) + + +if __name__ == "__main__": + main() diff --git a/dev/flux_steps_benchmark.py b/dev/flux_steps_benchmark.py new file mode 100644 index 0000000..5497ac9 --- /dev/null +++ b/dev/flux_steps_benchmark.py @@ -0,0 +1,66 @@ +#!/usr/bin/env python3 +"""FLUX.2 [klein] 4B Base – Steps-Vergleich (Qualität vs. Latenz). + +1024x1024, fester Seed, cpu_offload. Vergleicht 20/30/40/50 Steps. +""" + +import os +import subprocess +import time + +os.environ["HF_HUB_DISABLE_PROGRESS_BARS"] = "1" +os.environ["TRANSFORMERS_VERBOSITY"] = "error" +os.environ["TOKENIZERS_PARALLELISM"] = "false" + +MODEL = "/opt/mike-ai/models/FLUX.2-klein-base-4B" +PROMPT = ("A detailed photograph of a red cube on a white marble table, " + "soft studio lighting, shallow depth of field") +SEED = 0 +WIDTH = HEIGHT = 1024 +STEPS_LIST = [20, 30, 40, 50] + + +def nvidia_vram() -> int: + out = subprocess.check_output( + ["nvidia-smi", "--query-gpu=memory.used", + "--format=csv,noheader,nounits"]).decode().strip() + return int(out.split()[0]) + + +def main() -> None: + import torch + print(f"torch {torch.__version__} | {torch.cuda.get_device_name(0)}", + flush=True) + from diffusers import Flux2KleinPipeline + + t0 = time.monotonic() + pipe = Flux2KleinPipeline.from_pretrained(MODEL, torch_dtype=torch.bfloat16) + pipe.enable_model_cpu_offload() + print(f"[load] ready in {time.monotonic() - t0:.1f} s", flush=True) + + for steps in STEPS_LIST: + out = f"/tmp/flux-steps-{steps}.png" + torch.cuda.synchronize() + torch.cuda.reset_peak_memory_stats() + t = time.monotonic() + try: + img = pipe( + prompt=PROMPT, + height=HEIGHT, + width=WIDTH, + guidance_scale=4.0, + num_inference_steps=steps, + generator=torch.Generator(device="cuda").manual_seed(SEED), + ).images[0] + dt = time.monotonic() - t + img.save(out) + peak = torch.cuda.max_memory_allocated() / 1e9 + print(f"[gen] {steps} steps: {dt:.1f} s | peak {peak:.2f} GB | " + f"nvidia-smi {nvidia_vram()} MiB | {out}", flush=True) + except Exception as e: # noqa: BLE001 + print(f"[gen] {steps} steps: FEHLER: {e!r}", flush=True) + print("DONE", flush=True) + + +if __name__ == "__main__": + main() diff --git a/dev/mock_image_worker.py b/dev/mock_image_worker.py new file mode 100644 index 0000000..1385c71 --- /dev/null +++ b/dev/mock_image_worker.py @@ -0,0 +1,90 @@ +#!/usr/bin/env python3 +"""Mock-Bild-Worker für lokale Tests (gleiche Protokoll wie image_worker.py). + +Erzeugt ein minimales 1x1-PNG statt eines echten Bildes. + +Optionen (Umgebungsvariablen): + MOCK_WORKER_DELAY Sekunden, die pro generate geschlafen werden + (Default 0.3). Für Tests von parallelen Requests. + MOCK_WORKER_LOG Datei, in die die Requests geloggt werden (JSON-Zeilen). + Für Tests, die die Steps/Qualität prüfen wollen. + +Sonder-Prompts: + "FAIL" -> Worker antwortet mit Fehler (simuliert OOM/Crash). + "SLOW" -> Worker schläft 5 s (für Parallel-Tests). +""" + +import json +import os +import sys +import time + +# Minimales 1x1-PNG (1 Byte rot) +PNG_1x1 = bytes.fromhex( + "89504e470d0a1a0a0000000d49484452000000010000000108060000001f15c4" + "890000000d49444154789c626001000000ffff03000006000557bfabd40000" + "000049454e44ae426082" +) + +DELAY = float(os.environ.get("MOCK_WORKER_DELAY", "0.3")) +LOG_FILE = os.environ.get("MOCK_WORKER_LOG", "") + + +def _emit(payload: dict) -> None: + sys.stdout.write(json.dumps(payload) + "\n") + sys.stdout.flush() + + +def _log_request(req: dict) -> None: + if not LOG_FILE: + return + try: + with open(LOG_FILE, "a", encoding="utf-8") as f: + f.write(json.dumps(req) + "\n") + except OSError: + pass + + +def main() -> None: + loaded = False + _emit({"status": "ready"}) + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + except ValueError: + _emit({"status": "error", "message": "ungültiges JSON"}) + continue + cmd = req.get("cmd") + if cmd == "generate": + _log_request(req) + prompt = req.get("prompt", "") + if prompt == "SLOW": + time.sleep(5.0) # langsame Generierung (Parallel-Tests) + else: + time.sleep(DELAY) # simulierte Generierung + if prompt == "FAIL": + _emit({"status": "error", + "message": "simulierter Fehler (OOM)"}) + continue + output = req["output"] + os.makedirs(os.path.dirname(output) or ".", exist_ok=True) + with open(output, "wb") as f: + f.write(PNG_1x1) + _emit({"status": "ok", "path": output, "seconds": DELAY, + "load_seconds": 0.1}) + loaded = True + elif cmd == "unload": + loaded = False + _emit({"status": "ok"}) + elif cmd == "status": + _emit({"status": "ok", "model_loaded": loaded}) + else: + _emit({"status": "error", + "message": f"unbekanntes Kommando: {cmd}"}) + + +if __name__ == "__main__": + main() diff --git a/dev/test_local.sh b/dev/test_local.sh index a613ec3..497f335 100755 --- a/dev/test_local.sh +++ b/dev/test_local.sh @@ -13,7 +13,7 @@ FAIL=0 cleanup() { kill "${MOCK_PID:-}" "${ROUTER_PID:-}" 2>/dev/null || true - rm -f /tmp/mock_pid2 + rm -f /tmp/mock_pid2 /tmp/mock_upstream_pid wait 2>/dev/null || true } trap cleanup EXIT @@ -21,22 +21,41 @@ trap cleanup EXIT ok() { echo " PASS: $1"; PASS=$((PASS+1)); } bad() { echo " FAIL: $1"; FAIL=$((FAIL+1)); } -# --- Mock-llama.cpp starten -------------------------------------------------- +# --- Mock-llama.cpp starten (über Fake-systemctl) ------------------------------ echo "== Starte Mock-llama.cpp (Port $UP_PORT)" -MOCK_PROFILE_DIR="$FAKE_DIR" MOCK_PORT="$UP_PORT" python3 dev/mock_upstream.py >/tmp/mock_upstream.log 2>&1 & -MOCK_PID=$! +FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid \ +FAKE_SYSTEMD_PORT="$UP_PORT" \ +FAKE_SYSTEMD_PROFILE_DIR="$FAKE_DIR" \ +FAKE_SYSTEMD_MOCK="$PWD/dev/mock_upstream.py" \ +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 "") # --- Router starten ----------------------------------------------------------- echo "== Starte Router (Port $RT_PORT)" +rm -rf /tmp/test-images ROUTER_HOST=127.0.0.1 ROUTER_PORT="$RT_PORT" \ UPSTREAM_URL="http://127.0.0.1:$UP_PORT" \ PROFILE_SCRIPT="$PWD/dev/fake-llama-profile.sh" \ PROFILE_DIR="$FAKE_DIR" \ SWITCH_TIMEOUT=30 \ +SYSTEMCTL_BIN="$PWD/dev/fake-systemctl.sh" \ +FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid \ +FAKE_SYSTEMD_PORT="$UP_PORT" \ +FAKE_SYSTEMD_PROFILE_DIR="$FAKE_DIR" \ +FAKE_SYSTEMD_MOCK="$PWD/dev/mock_upstream.py" \ +FAKE_SYSTEMD_LOG=/tmp/mock_upstream_fake.log \ +IMAGE_WORKER="$PWD/dev/mock_image_worker.py" \ +IMAGE_PYTHON=python3 \ +IMAGE_DIR=/tmp/test-images \ +IMAGE_WORKER_LOG=/tmp/test_worker.log \ +IMAGE_GEN_TIMEOUT=30 \ +MOCK_WORKER_LOG=/tmp/test_worker_requests.jsonl \ python3 router/ai_profile_router.py >/tmp/router_test.log 2>&1 & ROUTER_PID=$! sleep 0.5 +rm -f /tmp/test_worker_requests.jsonl # --- 1. /v1/models ------------------------------------------------------------- echo "== Test 1: /v1/models" @@ -150,30 +169,232 @@ cat /tmp/err9b.json; echo echo "== Test 10: Upstream down -> 502, danach Recovery" # Profil auf fast setzen (aus Test 8 ist long aktiv) curl -sf -X POST "$BASE/fast" >/dev/null -kill "$MOCK_PID" 2>/dev/null; wait "$MOCK_PID" 2>/dev/null || true +# Mock stoppen (simuliert Crash) – über Fake-systemctl +FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid FAKE_SYSTEMD_PORT="$UP_PORT" \ + bash dev/fake-systemctl.sh stop sleep 0.5 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 [ "$CODE" = "502" ] && ok "502 bei downem Upstream (Profil bereits aktiv)" || bad "erwartet 502, bekam $CODE" -# Mock nach ~3 s neu starten (simuliert systemctl restart durch das Profil-Skript) -rm -f /tmp/mock_pid2 -( - sleep 3 - MOCK_PROFILE_DIR="$FAKE_DIR" MOCK_PORT="$UP_PORT" python3 dev/mock_upstream.py >/tmp/mock_upstream2.log 2>&1 & - echo $! > /tmp/mock_pid2 -) & +# Mock neu starten (simuliert systemctl restart durch das Profil-Skript) +FAKE_SYSTEMD_PIDFILE=/tmp/mock_upstream_pid FAKE_SYSTEMD_PORT="$UP_PORT" \ +FAKE_SYSTEMD_PROFILE_DIR="$FAKE_DIR" FAKE_SYSTEMD_MOCK="$PWD/dev/mock_upstream.py" \ +FAKE_SYSTEMD_LOG=/tmp/mock_upstream2.log \ + bash dev/fake-systemctl.sh start RESP=$(curl -sf -X POST "$BASE/fast") echo "$RESP" | python3 -m json.tool # neuen Mock als MOCK_PID übernehmen, damit Cleanup ihn beendet -[ -f /tmp/mock_pid2 ] && MOCK_PID=$(cat /tmp/mock_pid2) +[ -f /tmp/mock_upstream_pid ] && MOCK_PID=$(cat /tmp/mock_upstream_pid) echo "$RESP" | python3 -c ' import json,sys d=json.load(sys.stdin) assert d["profile"]=="fast" and d["model"]=="mock-model-73728", d ' && ok "Recovery: /fast wartet auf Upstream, dann Erfolg" || bad "Recovery" +# --- 11. Bildgenerierung (Mock-Worker) ------------------------------------------------- +echo "== Test 11: POST /v1/images/generations (1024x1024)" +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"ein rotes Haus","size":"1024x1024"}') +echo "$RESP" | python3 -m json.tool +IMG_NAME=$(echo "$RESP" | python3 -c 'import json,sys; print(json.load(sys.stdin)["data"][0]["url"].rsplit("/",1)[1])') +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +assert len(d["data"])==1, d +assert d["data"][0]["url"].startswith("http://"), d +' && [ -f "/tmp/test-images/$IMG_NAME" ] \ + && ok "Bild generiert und gespeichert ($IMG_NAME)" || bad "Bildgenerierung" + +# --- 12. Bild-Download ----------------------------------------------------------------- +echo "== Test 12: GET /images/" +CODE=$(curl -s -o /tmp/test_dl.png -w "%{http_code}" -D /tmp/hdr12.txt "$BASE/images/$IMG_NAME") +CTYPE=$(grep -i content-type /tmp/hdr12.txt | tr -d "\r") +[ "$CODE" = "200" ] && [ -s /tmp/test_dl.png ] && echo "$CTYPE" | grep -qi "image/png" \ + && ok "PNG-Download (200, $CTYPE)" || bad "PNG-Download (Code $CODE, $CTYPE)" + +# --- 13. Bild-Liste --------------------------------------------------------------------- +echo "== Test 13: GET /images" +RESP=$(curl -sf "$BASE/images") +echo "$RESP" | python3 -m json.tool +echo "$RESP" | python3 -c " +import json,sys +d=json.load(sys.stdin) +names=[i['name'] for i in d['images']] +assert '$IMG_NAME' in names, names +" && ok "Bild in Liste enthalten" || bad "Bild-Liste" + +# --- 14. Validierung --------------------------------------------------------------------- +echo "== Test 14: Validierung (Größe, Prompt, n)" +CODE=$(curl -s -o /tmp/err14a.json -w "%{http_code}" "$BASE/v1/images/generations" \ + -H "Content-Type: application/json" -d '{"prompt":"x","size":"500x500"}') +cat /tmp/err14a.json; echo +[ "$CODE" = "400" ] && ok "400 bei ungültiger Größe" || bad "erwartet 400, bekam $CODE" + +CODE=$(curl -s -o /tmp/err14b.json -w "%{http_code}" "$BASE/v1/images/generations" \ + -H "Content-Type: application/json" -d '{"size":"1024x1024"}') +cat /tmp/err14b.json; echo +[ "$CODE" = "400" ] && ok "400 bei fehlendem Prompt" || bad "erwartet 400, bekam $CODE" + +CODE=$(curl -s -o /tmp/err14c.json -w "%{http_code}" "$BASE/v1/images/generations" \ + -H "Content-Type: application/json" -d '{"prompt":"x","n":9}') +cat /tmp/err14c.json; echo +[ "$CODE" = "400" ] && ok "400 bei n=9 (max 4)" || bad "erwartet 400, bekam $CODE" + +# --- 15. b64_json + n=2 + Seed ------------------------------------------------------------- +echo "== Test 15: response_format=b64_json, n=2, seed" +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"zwei Bilder","size":"1024x1024","n":2,"seed":42,"response_format":"b64_json"}') +echo "$RESP" | python3 -c ' +import json,sys,base64 +d=json.load(sys.stdin) +assert len(d["data"])==2, d +for item in d["data"]: + assert item["url"] is None, item + png=base64.b64decode(item["b64_json"]) + assert png[:4]==b"\x89PNG", "kein PNG" +' && ok "2 Bilder als b64_json (gültige PNGs)" || bad "b64_json/n=2" + +# --- 16. /status zeigt Bild-Zustand --------------------------------------------------------- +echo "== Test 16: /status mit Bild-Section" +RESP=$(curl -sf "$BASE/status") +echo "$RESP" | python3 -m json.tool +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +img=d["image"] +assert img["phase"]=="idle", img +assert img["worker"]=="stopped", img # Worker wird nach Job beendet +assert img["model_loaded"] is False, img +assert img["last_image"], img +assert img["last_error"] is None, img +q=d["qwen"] +assert q["available"] is True, q +assert q["active_chats"]==0, q +' && ok "Status: phase=idle, worker=stopped, qwen verfügbar" || bad "Status Bild-Section" + +# --- 17. Qwen nach Bildgenerierung erreichbar ------------------------------------------------- +echo "== Test 17: Qwen nach Bildgenerierung erreichbar" +RESP=$(curl -sf "$BASE/v1/chat/completions" -H "Content-Type: application/json" \ + -d '{"model":"qwen-fast","messages":[{"role":"user","content":"Hallo"}]}') +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +assert "Mock-Antwort" in d["choices"][0]["message"]["content"], d +' && ok "Chat funktioniert nach Bildgenerierung" || bad "Chat nach Bild" + +# --- 18. quality=standard → 30 Steps ------------------------------------------------------------- +echo "== Test 18: quality=standard → 30 Steps" +rm -f /tmp/test_worker_requests.jsonl +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"standard test","size":"1024x1024","quality":"standard"}') +sleep 0.3 +STEPS=$(tail -1 /tmp/test_worker_requests.jsonl 2>/dev/null | python3 -c 'import json,sys; print(json.load(sys.stdin)["steps"])' 2>/dev/null || echo "?") +[ "$STEPS" = "30" ] && ok "quality=standard → 30 Steps" || bad "erwartet 30 Steps, bekam $STEPS" + +# --- 19. quality=high → 50 Steps ------------------------------------------------------------------- +echo "== Test 19: quality=high → 50 Steps" +rm -f /tmp/test_worker_requests.jsonl +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"high test","size":"1024x1024","quality":"high"}') +sleep 0.3 +STEPS=$(tail -1 /tmp/test_worker_requests.jsonl 2>/dev/null | python3 -c 'import json,sys; print(json.load(sys.stdin)["steps"])' 2>/dev/null || echo "?") +[ "$STEPS" = "50" ] && ok "quality=high → 50 Steps" || bad "erwartet 50 Steps, bekam $STEPS" + +# --- 20. ungültige Qualität → 400 ------------------------------------------------------------------ +echo "== Test 20: ungültige Qualität → 400" +CODE=$(curl -s -o /tmp/err20.json -w "%{http_code}" "$BASE/v1/images/generations" \ + -H "Content-Type: application/json" -d '{"prompt":"x","quality":"bogus"}') +cat /tmp/err20.json; echo +[ "$CODE" = "400" ] && ok "400 bei ungültiger Qualität" || bad "erwartet 400, bekam $CODE" + +# --- 21. Image-Fehler → Qwen wiederhergestellt ------------------------------------------------------ +echo "== Test 21: Image-Fehler → Qwen wiederhergestellt" +curl -sf -X POST "$BASE/fast" >/dev/null +CODE=$(curl -s -o /tmp/err21.json -w "%{http_code}" "$BASE/v1/images/generations" \ + -H "Content-Type: application/json" -d '{"prompt":"FAIL","size":"1024x1024"}') +cat /tmp/err21.json; echo +sleep 0.5 +RESP=$(curl -sf "$BASE/status") +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +assert d["current_profile"]=="fast", d +assert d["upstream"]["reachable"] is True, d +assert d["qwen"]["available"] is True, d +' && ok "Qwen nach Image-Fehler wiederhergestellt (fast, erreichbar)" || bad "Qwen nicht wiederhergestellt" + +# --- 22. Fast → Image → Fast ------------------------------------------------------------------------ +echo "== Test 22: Fast → Image → Fast" +curl -sf -X POST "$BASE/fast" >/dev/null +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"fast test","size":"1024x1024"}') +sleep 0.5 +PROFILE=$(curl -sf "$BASE/status" | python3 -c 'import json,sys; print(json.load(sys.stdin)["current_profile"])') +[ "$PROFILE" = "fast" ] && ok "Fast → Image → Fast" || bad "Profil nach Image: $PROFILE (erwartet fast)" + +# --- 23. Medium → Image → Medium -------------------------------------------------------------------- +echo "== Test 23: Medium → Image → Medium" +curl -sf -X POST "$BASE/medium" >/dev/null +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"medium test","size":"1024x1024"}') +sleep 0.5 +PROFILE=$(curl -sf "$BASE/status" | python3 -c 'import json,sys; print(json.load(sys.stdin)["current_profile"])') +[ "$PROFILE" = "medium" ] && ok "Medium → Image → Medium" || bad "Profil nach Image: $PROFILE (erwartet medium)" + +# --- 24. Long → Image → Long ------------------------------------------------------------------------ +echo "== Test 24: Long → Image → Long" +curl -sf -X POST "$BASE/long" >/dev/null +RESP=$(curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"long test","size":"1024x1024"}') +sleep 0.5 +PROFILE=$(curl -sf "$BASE/status" | python3 -c 'import json,sys; print(json.load(sys.stdin)["current_profile"])') +[ "$PROFILE" = "long" ] && ok "Long → Image → Long" || bad "Profil nach Image: $PROFILE (erwartet long)" + +# --- 25. /status während Image-Job ------------------------------------------------------------------- +echo "== Test 25: /status während Image-Job" +curl -sf -X POST "$BASE/fast" >/dev/null +curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"SLOW","size":"1024x1024"}' >/tmp/img25.json 2>&1 & +IMG_PID=$! +sleep 1.5 +RESP=$(curl -sf "$BASE/status") +echo "$RESP" | python3 -m json.tool +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +img=d["image"] +assert img["phase"]!="idle", img +assert d["qwen"]["available"] is False, d +' && ok "Status während Image-Job: phase!=idle, qwen unavailable" || bad "Status während Image-Job" +wait $IMG_PID +sleep 0.5 +RESP=$(curl -sf "$BASE/status") +echo "$RESP" | python3 -c ' +import json,sys +d=json.load(sys.stdin) +assert d["qwen"]["available"] is True, d +assert d["image"]["phase"]=="idle", d +' && ok "Nach Image-Job: qwen verfügbar, phase=idle" || bad "Nach Image-Job" + +# --- 26. paralleler Chat während Image-Job (wartet, kein 502) ---------------------------------------- +echo "== Test 26: paralleler Chat während Image-Job (wartet, kein 502)" +curl -sf -X POST "$BASE/fast" >/dev/null +curl -sf "$BASE/v1/images/generations" -H "Content-Type: application/json" \ + -d '{"prompt":"SLOW","size":"1024x1024"}' >/tmp/img26.json 2>&1 & +IMG_PID=$! +sleep 1.5 +START=$(date +%s) +CODE=$(curl -s -o /tmp/chat26.json -w "%{http_code}" "$BASE/v1/chat/completions" \ + -H "Content-Type: application/json" -d '{"model":"qwen-fast","messages":[{"role":"user","content":"Hallo"}]}') +END=$(date +%s) +ELAPSED=$((END-START)) +cat /tmp/chat26.json; echo +wait $IMG_PID +[ "$CODE" = "200" ] && [ "$ELAPSED" -ge 2 ] \ + && ok "Chat wartete ${ELAPSED}s (kein 502), dann 200" || bad "Chat: Code $CODE, ${ELAPSED}s" + # --- Ergebnis -------------------------------------------------------------------------------------------- echo echo "== Ergebnis: $PASS bestanden, $FAIL fehlgeschlagen ==" diff --git a/router/ai_profile_router.py b/router/ai_profile_router.py index 526c135..6c8fb80 100755 --- a/router/ai_profile_router.py +++ b/router/ai_profile_router.py @@ -15,19 +15,28 @@ Virtuelle Modelle: qwen-fast, qwen-medium, qwen-long Kommandos: POST /fast, /medium, /long (Profilwechsel) GET /status (Zustand) -Ein Profilwechsel führt PROFILE_SCRIPT aus (ohne Shell, feste -Argumente → keine Injection), wartet dann, bis llama.cpp wieder erreichbar -ist, und erst dann wird eine erfolgreiche Antwort geliefert bzw. der -Request weitergeleitet. +Bildgenerierung (FLUX.2 [klein] 4B Base): + POST /v1/images/generations (OpenAI-kompatibel) + GET /images (Liste) + GET /images/ (PNG-Download) + +Der Router agiert als Modell-Orchestrator: vor der Generierung wird +llama.cpp gestoppt, der Bild-Worker lädt FLUX, generiert und entlädt +das Modell wieder; danach wird das vorherige Qwen-Profil wiederher- +gestellt und erst dann geantwortet (try/finally – Qwen wird auch bei +Fehlgeschlagener Generierung wiederhergestellt). Nur Python-Standardbibliothek. Logging nach stdout (journald). """ from __future__ import annotations +import base64 import json import logging import os +import queue +import re import subprocess import sys import threading @@ -51,6 +60,42 @@ REQUEST_TIMEOUT = float(os.environ.get("REQUEST_TIMEOUT", "600")) # s, Read-Ti CONNECT_TIMEOUT = float(os.environ.get("CONNECT_TIMEOUT", "10")) # s, Connect-Timeout POLL_INTERVAL = float(os.environ.get("POLL_INTERVAL", "2")) # s, Polling-Intervall +# --- Bildgenerierung (FLUX.2 [klein] 4B Base) --- +LLAMA_SERVICE = os.environ.get("LLAMA_SERVICE", "mike-ai-llama-ui.service") +SYSTEMCTL_BIN = os.environ.get("SYSTEMCTL_BIN", "systemctl") +IMAGE_WORKER = os.environ.get( + "IMAGE_WORKER", "/opt/mike-ai/ai-profile-router/image_worker.py") +IMAGE_PYTHON = os.environ.get( + "IMAGE_PYTHON", "/opt/mike-ai/ai-profile-router/venv/bin/python") +IMAGE_DIR = os.environ.get( + "IMAGE_DIR", "/opt/mike-ai/ai-profile-router/images") +IMAGE_WORKER_LOG = os.environ.get( + "IMAGE_WORKER_LOG", "/opt/mike-ai/ai-profile-router/image_worker.log") +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 + +# Erlaubte Auflösungen (Breite x Höhe). FLUX.2 klein ist für 1 MP +# ausgelegt; 1920x1088 (≈2 MP) wird zusätzlich unterstützt. +IMAGE_SIZES = { + "1024x1024": (1024, 1024), + "1536x1024": (1536, 1024), + "1024x1536": (1024, 1536), + "1920x1088": (1920, 1088), + "1088x1920": (1088, 1920), +} +# Qualitätsstufen → Inference-Schritte (guidance bleibt offiziell 4.0). +# Auf der RTX 5080 gemessen: 30 vs. 50 Steps liefern praktisch dieselbe +# Qualität (1024x1024: 31,3 s vs. 45,3 s). Default ist daher "standard". +IMAGE_QUALITY = {"standard": 30, "high": 50} +IMAGE_DEFAULT_QUALITY = "standard" +IMAGE_MAX_N = 4 + +# Chat-Waiting: Während eines Image-Jobs oder Profilwechsels ist Qwen +# down. Chat-Requests warten (statt 502) bis Qwen wieder bereit ist. +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} VIRTUAL_MODELS = {f"qwen-{name}": name for name in PROFILES} @@ -78,16 +123,68 @@ UPSTREAM_HOST, UPSTREAM_PORT = _parse_upstream(UPSTREAM_URL) # Zustand # --------------------------------------------------------------------------- +class _ImageState: + """Zustand der Bildgenerierung (nur für Status-Reporting).""" + def __init__(self) -> None: + self.phase = "idle" # siehe PHASES unten + self.worker: "_Worker | None" = None + self.last_error: str | None = None + self.last_image: str | None = None + self.last_seconds: float | None = None + + +IMAGE_PHASES = ( + "idle", "stopping-qwen", "loading-image", "generating", + "unloading-image", "restoring-qwen", +) + + class _State: - """Gemeinsamer, thread-sicherer Zustand.""" - lock = threading.Lock() # serialisiert Profilwechsel + """Gemeinsamer, thread-sicherer Zustand. + + lock : zentraler GPU-/Model-Lock. Wird von Profilwechsel UND + Image-Generation gehalten → gegenseitiger Ausschluss, + kein Race zwischen beiden. + avail_lock : schützt qwen_unavailable + active_chats (Chat-Waiting). + """ + lock = threading.Lock() # GPU-/Model-Lock (Profilwechsel + Image) switching: str | None = None # Profil, das gerade gewechselt wird started = time.time() + image = _ImageState() + # 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 STATE = _State() +def _wait_chats_drained(timeout: float | None = None) -> None: + """Wartet, bis keine aktiven Chat-Requests mehr laufen. + + Wird von Profilwechsel/Image-Job aufgerufen, BEVOR Qwen gestoppt wird. + Verhindert, dass ein laufender Chat auf ein gestopptes Qwen trifft (502). + """ + timeout = CHAT_DRAIN_TIMEOUT if timeout is None else timeout + deadline = time.monotonic() + timeout + while True: + with STATE.avail_lock: + if STATE.active_chats == 0: + 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 + time.sleep(0.5) + + +def _set_qwen_unavailable(unavailable: bool) -> None: + with STATE.avail_lock: + STATE.qwen_unavailable = unavailable + + # --------------------------------------------------------------------------- # Upstream (llama.cpp) # --------------------------------------------------------------------------- @@ -167,6 +264,9 @@ def switch_profile(profile: str, implicit: bool = False) -> None: if profile not in PROFILES: raise ValueError(f"unbekanntes Profil: {profile!r} " f"(erlaubt: {', '.join(PROFILES)})") + # Kein Fast-Fail: Wenn ein Image-Job läuft (hält den GPU-Lock), wartet + # der Profilwechsel auf den GPU-Lock (blockiert), bis der Image-Job + # fertig ist. So bekommen Chat-Requests kein 502, sondern warten. with STATE.lock: STATE.switching = profile try: @@ -177,43 +277,299 @@ def switch_profile(profile: str, implicit: bool = False) -> None: if cur == profile and ready: log.info("Profil %s ist bereits aktiv", profile) return - if cur == profile and up["reachable"] and not ready: - # Modell wird gerade geladen (z.B. nach einem Wechsel) - log.info("Warte, bis Profil %s geladen ist ...", profile) - _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) - return - if cur == profile and not up["reachable"] and implicit: - raise RuntimeError( - f"llama.cpp nicht erreichbar (Profil {profile} ist bereits " - f"aktiv; Neustart über /{profile})") - log.info("Profilwechsel: %s -> %s", cur, profile) + # Qwen wird neu geladen/gewechselt → für Chats nicht verfügbar. + _set_qwen_unavailable(True) try: - proc = subprocess.run( - [PROFILE_SCRIPT, profile], - stdin=subprocess.DEVNULL, - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - timeout=120, - ) - out = proc.stdout.decode(errors="replace").strip() - 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 aber erledigt. - log.warning("llama-profile Exit-Code %d (ohne TTY erwartet)", - proc.returncode) - except subprocess.TimeoutExpired: - log.error("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) + _wait_chats_drained() + if cur == profile and up["reachable"] and not ready: + # Modell wird gerade geladen (z.B. nach einem Wechsel) + log.info("Warte, bis Profil %s geladen ist ...", profile) + _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) + return + if cur == profile and not up["reachable"] and implicit: + raise RuntimeError( + f"llama.cpp nicht erreichbar (Profil {profile} ist " + f"bereits aktiv; Neustart über /{profile})") + log.info("Profilwechsel: %s -> %s", cur, profile) + try: + proc = subprocess.run( + [PROFILE_SCRIPT, profile], + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + timeout=120, + ) + out = proc.stdout.decode(errors="replace").strip() + 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) + except subprocess.TimeoutExpired: + log.error("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) + finally: + _set_qwen_unavailable(False) finally: STATE.switching = None +# --------------------------------------------------------------------------- +# Bildgenerierung (FLUX.2 [klein] 4B Base) +# --------------------------------------------------------------------------- + +class _Worker: + """Verwaltet den Bild-Worker-Prozess (stdin/stdout-JSON-Protokoll).""" + + def __init__(self) -> None: + self.proc: subprocess.Popen | None = None + self.model_loaded = False + self._queue: queue.Queue[dict] = queue.Queue() + self._reader: threading.Thread | None = None + + def alive(self) -> bool: + return self.proc is not None and self.proc.poll() is None + + def start(self) -> None: + if self.alive(): + return + log.info("starte Bild-Worker: %s %s", IMAGE_PYTHON, IMAGE_WORKER) + logf = open(IMAGE_WORKER_LOG, "ab") + self.proc = subprocess.Popen( + [IMAGE_PYTHON, IMAGE_WORKER], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=logf, + text=True, + bufsize=1, + ) + self._reader = threading.Thread(target=self._read_loop, daemon=True) + self._reader.start() + try: + msg = self._queue.get(timeout=IMAGE_START_TIMEOUT) + except queue.Empty: + self.stop() + raise RuntimeError("Bild-Worker hat nicht gestartet") + if msg.get("status") != "ready": + self.stop() + raise RuntimeError(f"Bild-Worker-Startfehler: {msg}") + log.info("Bild-Worker bereit") + + def _read_loop(self) -> None: + assert self.proc is not None and self.proc.stdout is not None + for line in self.proc.stdout: + line = line.strip() + if not line: + continue + try: + self._queue.put(json.loads(line)) + except ValueError: + log.warning("Worker-Zeile (kein JSON): %s", line[:200]) + + def request(self, payload: dict, timeout: float) -> dict: + if not self.alive(): + raise RuntimeError("Bild-Worker ist nicht aktiv") + assert self.proc is not None and self.proc.stdin is not None + self.proc.stdin.write(json.dumps(payload) + "\n") + self.proc.stdin.flush() + try: + return self._queue.get(timeout=timeout) + except queue.Empty: + raise RuntimeError( + f"Bild-Worker hat nach {timeout:.0f} s nicht geantwortet " + f"(cmd={payload.get('cmd')})") + + def stop(self) -> None: + if self.proc is not None and self.proc.poll() is None: + self.proc.terminate() + try: + self.proc.wait(timeout=10) + except subprocess.TimeoutExpired: + self.proc.kill() + self.proc = None + self.model_loaded = False + + +def _worker() -> _Worker: + """Worker-Instanz liefern (startet bei Bedarf).""" + img = STATE.image + if not img.worker or not img.worker.alive(): + if img.worker: + img.worker.stop() + img.worker = _Worker() + img.worker.start() + return img.worker + + +def _wait_upstream_down(deadline: float) -> None: + """Wartet, bis llama.cpp den Port freigegeben hat (VRAM frei).""" + while time.monotonic() < deadline: + if not upstream_status()["reachable"]: + return + time.sleep(1) + raise RuntimeError("llama.cpp gibt Port/VRAM nicht frei") + + +def _vram_used_mib() -> int | None: + """Aktuelle VRAM-Belegung in MiB (via nvidia-smi), None bei Fehler.""" + try: + out = subprocess.run( + ["nvidia-smi", "--query-gpu=memory.used", + "--format=csv,noheader,nounits"], + stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, timeout=10, + ).stdout.decode().strip() + return int(out.splitlines()[0].split()[0]) + except (OSError, ValueError, IndexError): + return None + + +def _wait_vram_free(threshold_mib: int = 1000, + timeout: float | None = None) -> None: + """Wartet, bis der VRAM unter threshold_mib fällt (FLUX entladen). + + Wird nach dem Beenden des Bild-Workers aufgerufen, um sicherzustellen, + dass der VRAM (inkl. CUDA-Kontext) frei ist, bevor Qwen neu startet. + Wenn nvidia-smi nicht verfügbar ist (z.B. lokale Tests), wird der + Check übersprungen. + """ + timeout = IMAGE_VRAM_FREE_TIMEOUT if timeout is None else timeout + deadline = time.monotonic() + timeout + last = _vram_used_mib() + if last is None: + log.info("VRAM-Check übersprungen (nvidia-smi nicht verfügbar)") + return + while time.monotonic() < deadline: + if last <= threshold_mib: + log.info("VRAM frei: %d MiB", last) + return + time.sleep(1) + last = _vram_used_mib() + if last is None: + log.info("VRAM-Check übersprungen (nvidia-smi nicht verfügbar)") + return + raise RuntimeError( + f"VRAM nach {timeout:.0f} s nicht frei (letzte Messung: " + f"{last} MiB, erwartet <= {threshold_mib} MiB)") + + +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) + except subprocess.TimeoutExpired: + log.error("systemctl start hat 120 s überschritten") + _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) + + +def generate_image(prompt: str, width: int, height: int, steps: int, + guidance: float, seed: int | None, n: int + ) -> tuple[list[str], str | None]: + """Orchestriert die Bildgenerierung inkl. Qwen-Hotswap. + + Hält den zentralen GPU-Lock (gegenseitiger Ausschluss mit Profilwechsel). + Ablauf: Qwen stoppen → Worker laden → generieren → Worker beenden + (VRAM + CUDA-Kontext frei) → Qwen wiederherstellen. Qwen wird auch bei + Fehlern wiederhergestellt (try/finally). + """ + img = STATE.image + with STATE.lock: + if img.phase != "idle": + raise RuntimeError(f"Bildgenerierung läuft ({img.phase})") + profile = current_profile() + if profile is None: + raise RuntimeError("kein aktives Qwen-Profil (override.conf?)") + os.makedirs(IMAGE_DIR, exist_ok=True) + results: list[str] = [] + warning: str | None = None + # Qwen wird gestoppt → für Chats nicht verfügbar (die warten). + _set_qwen_unavailable(True) + try: + _wait_chats_drained() + + # 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) + _wait_upstream_down(time.monotonic() + 60) + + # 2) Worker starten (Modell wird beim ersten generate geladen). + img.phase = "loading-image" + worker = _worker() + + # 3) Generieren. + for i in range(n): + img.phase = "generating" + filename = time.strftime("%Y%m%d-%H%M%S") + \ + f"-{os.urandom(2).hex()}.png" + output = os.path.join(IMAGE_DIR, filename) + resp = worker.request({ + "cmd": "generate", + "prompt": prompt, + "width": width, + "height": height, + "steps": steps, + "guidance": guidance, + "seed": seed, + "output": output, + }, timeout=IMAGE_GEN_TIMEOUT) + if resp.get("status") != "ok": + raise RuntimeError( + resp.get("message", "Bildgenerierung fehlgeschlagen")) + worker.model_loaded = True + results.append(filename) + img.last_image = filename + img.last_seconds = resp.get("seconds") + log.info("Bild %d/%d: %s (%.1f s)", i + 1, n, filename, + resp.get("seconds", 0)) + + # 4) Worker vollständig beenden (VRAM + CUDA-Kontext freigeben). + img.phase = "unloading-image" + worker.stop() + img.worker = None + try: + _wait_vram_free() + except RuntimeError as e: + log.warning("VRAM-Check: %s (fahre mit Qwen-Restore fort)", e) + except Exception as e: + img.last_error = str(e) + log.error("Bildgenerierung fehlgeschlagen: %s", e) + # Worker sicher beenden (falls noch aktiv), VRAM freigeben. + if img.worker is not None: + img.worker.stop() + img.worker = None + raise + finally: + # 5) Qwen immer wiederherstellen. + img.phase = "restoring-qwen" + try: + _restore_qwen(profile) + _set_qwen_unavailable(False) + except Exception as e: + warning = f"Qwen-Wiederherstellung fehlgeschlagen: {e}" + img.last_error = warning + log.error(warning) + # Qwen ist down → qwen_unavailable bleibt True. + img.phase = "idle" + return results, warning + + +def _image_filename_ok(name: str) -> bool: + return bool(re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._-]*\.png", name)) + + # --------------------------------------------------------------------------- # HTTP-Handler # --------------------------------------------------------------------------- @@ -238,6 +594,12 @@ class Handler(BaseHTTPRequestHandler): self._send_json(200, self._models_payload()) elif path == "/status": self._send_json(200, self._status_payload()) + elif path == "/v1/images/generations" and self.command == "POST": + self._image_generate() + elif path == "/images" and self.command == "GET": + self._images_list() + elif path.startswith("/images/") and self.command == "GET": + self._image_serve(path[len("/images/"):]) elif path in ("/fast", "/medium", "/long"): self._switch(path[1:]) elif (self.command == "POST" and path.startswith("/") @@ -277,6 +639,10 @@ class Handler(BaseHTTPRequestHandler): def _status_payload(self) -> dict: up = upstream_status() + img = STATE.image + with STATE.avail_lock: + qwen_unavailable = STATE.qwen_unavailable + active_chats = STATE.active_chats return { "router": "ai-profile-router", "uptime_seconds": round(time.time() - STATE.started, 1), @@ -289,8 +655,174 @@ class Handler(BaseHTTPRequestHandler): "model": up.get("model"), "ctx": up.get("ctx"), }, + "qwen": { + "available": not qwen_unavailable, + "active_chats": active_chats, + }, + "image": { + "phase": img.phase, + "worker": "running" if (img.worker and img.worker.alive()) + else "stopped", + "model_loaded": bool(img.worker and img.worker.model_loaded), + "last_image": img.last_image, + "last_seconds": img.last_seconds, + "last_error": img.last_error, + }, } + # ---------- Bildgenerierung ---------- + + def _image_generate(self) -> None: + length = int(self.headers.get("Content-Length") or 0) + try: + data = json.loads(self.rfile.read(length)) + except ValueError: + self._send_error(400, "ungültiges JSON", + "invalid_request_error", "invalid_json") + return + if not isinstance(data, dict): + self._send_error(400, "Request muss ein JSON-Objekt sein", + "invalid_request_error", "invalid_request") + return + + prompt = data.get("prompt") + if not isinstance(prompt, str) or not prompt.strip(): + self._send_error(400, "'prompt' fehlt oder ist leer", + "invalid_request_error", "missing_prompt") + return + if len(prompt) > 8000: + self._send_error(400, "'prompt' zu lang (max 8000 Zeichen)", + "invalid_request_error", "prompt_too_long") + return + + # Größe + size = data.get("size", "1024x1024") + if size not in IMAGE_SIZES: + self._send_error( + 400, f"ungültige Größe: {size!r} " + f"(erlaubt: {', '.join(IMAGE_SIZES)})", + "invalid_request_error", "invalid_size") + return + width, height = IMAGE_SIZES[size] + + # Anzahl + n = data.get("n", 1) + if not isinstance(n, int) or isinstance(n, bool) or not 1 <= n <= IMAGE_MAX_N: + self._send_error(400, f"'n' muss eine Ganzzahl 1..{IMAGE_MAX_N} sein", + "invalid_request_error", "invalid_n") + return + + # Qualität / Schritte / Guidance + quality = data.get("quality", IMAGE_DEFAULT_QUALITY) + if quality not in IMAGE_QUALITY: + self._send_error(400, f"ungültige Qualität: {quality!r} " + f"(erlaubt: {', '.join(IMAGE_QUALITY)})", + "invalid_request_error", "invalid_quality") + return + steps = data.get("steps", IMAGE_QUALITY[quality]) + if not isinstance(steps, int) or isinstance(steps, bool) or not 4 <= steps <= 150: + self._send_error(400, "'steps' muss eine Ganzzahl 4..150 sein", + "invalid_request_error", "invalid_steps") + return + guidance = data.get("guidance", 4.0) + try: + guidance = float(guidance) + except (TypeError, ValueError): + self._send_error(400, "'guidance' muss eine Zahl sein", + "invalid_request_error", "invalid_guidance") + return + if not 1.0 <= guidance <= 10.0: + self._send_error(400, "'guidance' muss zwischen 1.0 und 10.0 sein", + "invalid_request_error", "invalid_guidance") + return + + seed = data.get("seed") + if seed is not None: + try: + seed = int(seed) + except (TypeError, ValueError): + self._send_error(400, "'seed' muss eine Ganzzahl sein", + "invalid_request_error", "invalid_seed") + return + if not 0 <= seed <= 2**32 - 1: + self._send_error(400, "'seed' muss zwischen 0 und 4294967295 sein", + "invalid_request_error", "invalid_seed") + return + + response_format = data.get("response_format", "url") + if response_format not in ("url", "b64_json"): + self._send_error(400, "'response_format' muss 'url' oder 'b64_json' sein", + "invalid_request_error", "invalid_response_format") + return + + # Generierung (blockt mehrere Minuten – eigener Thread-Timeout). + self.timeout = None + try: + results, warning = generate_image( + prompt.strip(), width, height, steps, guidance, seed, n) + except (ValueError, RuntimeError) as e: + self._send_error(503, str(e), "server_error", "image_generation_failed") + return + + # Antwort bauen + host = self.headers.get("Host") or f"{HOST}:{PORT}" + if not host.startswith(("http://", "https://")): + host = f"http://{host}" + items = [] + for filename in results: + path = os.path.join(IMAGE_DIR, filename) + item: dict = {"url": f"{host}/images/{filename}", "b64_json": None} + if response_format == "b64_json": + with open(path, "rb") as f: + item["b64_json"] = base64.b64encode(f.read()).decode() + item["url"] = None + items.append(item) + payload: dict = {"created": int(time.time()), "data": items} + if warning: + payload["router_warning"] = warning + self._send_json(200, payload) + + def _images_list(self) -> None: + if not os.path.isdir(IMAGE_DIR): + self._send_json(200, {"images": []}) + return + entries = [] + for name in sorted(os.listdir(IMAGE_DIR), reverse=True): + if not _image_filename_ok(name): + continue + path = os.path.join(IMAGE_DIR, name) + try: + st = os.stat(path) + except OSError: + continue + entries.append({ + "name": name, + "url": f"/images/{name}", + "bytes": st.st_size, + "modified": int(st.st_mtime), + }) + self._send_json(200, {"images": entries[:200]}) + + def _image_serve(self, name: str) -> None: + if not _image_filename_ok(name): + self._send_error(400, "ungültiger Dateiname", + "invalid_request_error", "invalid_filename") + return + path = os.path.join(IMAGE_DIR, name) + if not os.path.isfile(path): + self._send_error(404, "Bild nicht gefunden", + "invalid_request_error", "not_found") + return + data = open(path, "rb").read() + self._last_code = 200 + self.send_response(200) + self.send_header("Content-Type", "image/png") + self.send_header("Content-Length", str(len(data))) + self.send_header("Cache-Control", "public, max-age=86400") + self.send_header("Connection", "close") + self.end_headers() + self.wfile.write(data) + def _switch(self, profile: str) -> None: if profile not in PROFILES: self._send_error(400, f"unbekanntes Profil: {profile}", @@ -343,6 +875,42 @@ class Handler(BaseHTTPRequestHandler): "invalid_request_error", "unknown_model") return + # An llama.cpp weiterleiten (mit Chat-Waiting, Streaming bleibt erhalten). + self._proxy_with_wait(body) + + def _proxy_with_wait(self, body: bytes | None) -> None: + """Leitet an llama.cpp weiter, wartet aber erst, bis Qwen verfügbar ist. + + Während eines Image-Jobs oder Profilwechsels ist Qwen down. Statt + 502 zu liefern, wartet der Request (mit Timeout), bis Qwen wieder + bereit ist. Mehrere Chats können parallel laufen (active_chats). + + Race-frei: Der Check auf qwen_unavailable und das Inkrement von + active_chats sind atomar (avail_lock). Ein Image-Job/Profilwechsel + setzt qwen_unavailable=True und wartet auf active_chats==0, BEVOR + er Qwen stoppt – ein laufender Chat wird daher nie unterbrochen. + """ + deadline = time.monotonic() + CHAT_WAIT_TIMEOUT + while True: + with STATE.avail_lock: + if not STATE.qwen_unavailable: + STATE.active_chats += 1 + break + if time.monotonic() > deadline: + self._send_error( + 503, + "Qwen wird neu geladen (Image-Job oder Profilwechsel), " + "bitte später erneut", + "server_error", "qwen_reloading") + return + time.sleep(0.5) + try: + self._proxy(body) + finally: + with STATE.avail_lock: + STATE.active_chats -= 1 + + def _proxy(self, body: bytes | None) -> None: # An llama.cpp weiterleiten (Streaming bleibt erhalten). try: conn = http.client.HTTPConnection(UPSTREAM_HOST, UPSTREAM_PORT, diff --git a/router/image_worker.py b/router/image_worker.py new file mode 100644 index 0000000..1e625a2 --- /dev/null +++ b/router/image_worker.py @@ -0,0 +1,159 @@ +#!/usr/bin/env python3 +"""FLUX.2 [klein] 4B Base – Bild-Worker. + +Protokoll: zeilenbasiertes JSON über stdin/stdout. + + Start: Worker gibt {"status": "ready"} aus (Modell noch NICHT geladen). + Request: {"cmd": "generate", "prompt": ..., "width": ..., "height": ..., + "steps": ..., "guidance": ..., "seed": ..., "output": ...} + Antwort: {"status": "ok", "path": ..., "seconds": ..., "load_seconds": ...} + oder {"status": "error", "message": ...} + Request: {"cmd": "unload"} -> {"status": "ok"} + Request: {"cmd": "status"} -> {"status": "ok", "model_loaded": bool} + +Das Modell wird beim ersten generate geladen (bf16, cpu_offload) und auf +Anforderung wieder entladen (VRAM freigeben). Der Prozess bleibt danach +laufen – ohne geladenes Modell belegt er kaum Ressourcen. + +Alle torch-/diffusers-Logs gehen nach stderr, stdout ist reines Protokoll. +""" + +import gc +import json +import os +import signal +import sys +import time + +# stderr-Logs von torch & Co. unterdrücken, bevor importiert wird +os.environ.setdefault("DIFFUSERS_VERBOSITY", "error") +os.environ.setdefault("TRANSFORMERS_VERBOSITY", "error") +os.environ.setdefault("HF_HUB_DISABLE_PROGRESS_BARS", "1") +os.environ.setdefault("TOKENIZERS_PARALLELISM", "false") + +MODEL_DIR = os.environ.get( + "FLUX_MODEL_DIR", "/opt/mike-ai/models/FLUX.2-klein-base-4B") + +_pipe = None # geladene Pipeline (None = entladen) +_load_seconds = 0.0 # Dauer des letzten Ladens + + +def _emit(payload: dict) -> None: + sys.stdout.write(json.dumps(payload) + "\n") + sys.stdout.flush() + + +def _log(msg: str) -> None: + print(f"[image-worker] {msg}", file=sys.stderr, flush=True) + + +def _load() -> None: + """Pipeline laden (bf16, CPU-Offload).""" + global _pipe, _load_seconds + if _pipe is not None: + return + import torch + from diffusers import Flux2KleinPipeline + + t0 = time.monotonic() + _log(f"lade Modell aus {MODEL_DIR} ...") + _pipe = Flux2KleinPipeline.from_pretrained( + MODEL_DIR, torch_dtype=torch.bfloat16) + _pipe.enable_model_cpu_offload() + _load_seconds = time.monotonic() - t0 + _log(f"Modell geladen in {_load_seconds:.1f} s") + + +def _unload() -> None: + """Pipeline entladen und VRAM freigeben.""" + global _pipe + if _pipe is None: + return + t0 = time.monotonic() + del _pipe + _pipe = None + gc.collect() + try: + import torch + torch.cuda.empty_cache() + except Exception: + pass + _log(f"Modell entladen in {time.monotonic() - t0:.1f} s") + + +def _generate(req: dict) -> dict: + import torch + + prompt = req["prompt"] + width = int(req.get("width", 1024)) + height = int(req.get("height", 1024)) + steps = int(req.get("steps", 50)) + guidance = float(req.get("guidance", 4.0)) + seed = req.get("seed") + output = req["output"] + + _load() + + t0 = time.monotonic() + generator = None + if seed is not None: + generator = torch.Generator(device="cuda").manual_seed(int(seed)) + image = _pipe( + prompt=prompt, + height=height, + width=width, + guidance_scale=guidance, + num_inference_steps=steps, + generator=generator, + ).images[0] + + os.makedirs(os.path.dirname(output) or ".", exist_ok=True) + image.save(output) + seconds = time.monotonic() - t0 + _log(f"generiert {output} in {seconds:.1f} s " + f"({width}x{height}, {steps} steps, seed={seed})") + return { + "status": "ok", + "path": output, + "seconds": round(seconds, 2), + "load_seconds": round(_load_seconds, 2), + } + + +def _handle(line: str) -> None: + try: + req = json.loads(line) + except ValueError: + _emit({"status": "error", "message": "ungültiges JSON"}) + return + + cmd = req.get("cmd") + try: + if cmd == "generate": + _emit(_generate(req)) + elif cmd == "unload": + _unload() + _emit({"status": "ok"}) + elif cmd == "status": + _emit({"status": "ok", "model_loaded": _pipe is not None}) + else: + _emit({"status": "error", "message": f"unbekanntes Kommando: {cmd}"}) + except Exception as e: # noqa: BLE001 – Fehler ans Router-Protokoll + _log(f"Fehler bei {cmd}: {e!r}") + _emit({"status": "error", "message": str(e)}) + + +def main() -> None: + signal.signal(signal.SIGTERM, lambda *_: sys.exit(0)) + _emit({"status": "ready"}) + for line in sys.stdin: + line = line.strip() + if not line: + continue + _handle(line) + if _pipe is None and line.startswith('{"cmd": "unload"'): + pass # Worker bleibt laufen, Modell ist entladen + + +if __name__ == "__main__": + main()