From 7ac93befc40640a2cb0ed760768286246b6e7197 Mon Sep 17 00:00:00 2001 From: Mikei386 <44135113+Mikei386@users.noreply.github.com> Date: Sat, 22 Aug 2026 17:40:32 +0200 Subject: [PATCH] Add RTX 5080 FLUX hot-swap worker --- README.md | 7 +- compose.yaml | 40 ++++++- config/install.env.example | 2 + dev/test_profile_controller.py | 45 +++++++- docs/CURRENT_REFERENCE.md | 20 ++-- docs/STANDARD_PROFILE_MATRIX.md | 6 +- install.sh | 13 +++ platform/docker/flux-worker/Dockerfile | 18 +++ platform/docker/flux-worker/flux_worker.py | 109 ++++++++++++++++++ .../profile-controller/profile_controller.py | 68 +++++++++++ platform/models/manifest.example.yaml | 6 +- router/ai_profile_router.py | 105 +++++++++++++---- 12 files changed, 398 insertions(+), 41 deletions(-) create mode 100644 platform/docker/flux-worker/Dockerfile create mode 100644 platform/docker/flux-worker/flux_worker.py diff --git a/README.md b/README.md index 037f443..bc3c049 100644 --- a/README.md +++ b/README.md @@ -49,12 +49,13 @@ Neustart an; danach wird derselbe Befehl erneut ausgeführt. | Open WebUI | `:8080` | Chat und Administration | | Profile Router | `:8081` | OpenAI-kompatible API, Profilwahl | | llama.cpp | nur Docker-intern | Inferenz und integrierte Vision | -| Profile Controller | nur Docker-intern | eng begrenzter Containerwechsel | +| Profile Controller | nur Docker-intern | eng begrenzter Profil-/FLUX-Hot-Swap | +| FLUX Worker | nur Docker-intern, normalerweise gestoppt | Bildgenerierung auf RTX 5080 | | Piper | nur Docker-intern | lokale deutsche Sprachausgabe | | MCP-Tool-Stack | nur Docker-intern | Web, Home Assistant, ARR und Unraid | -Piper-TTS ist ein reproduzierbarer Kerndienst; STT und Bildgenerierung bleiben -optionale Dienste. Web-, Home-Assistant-, +Piper-TTS und der FLUX.2-Klein-Hot-Swap sind reproduzierbare Kerndienste; STT +bleibt optional. Web-, Home-Assistant-, ARR- und Unraid-Werkzeuge besitzen dagegen bereits getrennte Container unter `platform/mcp/`. Open WebUI erreicht sie ausschließlich über das interne `mike-ai-tools`-Netz; llama.cpp erhält keine MCP-Konfiguration und keine diff --git a/compose.yaml b/compose.yaml index 396e201..0a2f247 100644 --- a/compose.yaml +++ b/compose.yaml @@ -370,6 +370,7 @@ services: environment: CONTROLLER_TOKEN: "${CONTROLLER_TOKEN:?CONTROLLER_TOKEN is required}" ALLOWED_PROFILES: fast,medium,large,ultra,experimental + IMAGE_WORKER: flux networks: [control] security_opt: ["no-new-privileges:true"] healthcheck: @@ -405,8 +406,10 @@ services: SWITCH_TIMEOUT: "600" REQUEST_TIMEOUT: "600" IMAGE_DIR: /data/images + IMAGE_WORKER_URL: http://flux-worker:8086 + IMAGE_WORKER_TOKEN: "${CONTROLLER_TOKEN:?CONTROLLER_TOKEN is required}" CHAT_IMAGE_ALLOW_REMOTE_URLS: "false" - ENABLE_IMAGE_GENERATION: "false" + ENABLE_IMAGE_GENERATION: "true" ENABLE_TTS: "true" TTS_WORKER_URL: http://piper:8085 TTS_MODEL: piper @@ -435,6 +438,41 @@ services: piper: condition: service_healthy + flux-worker: + build: + context: platform/docker/flux-worker + args: + DIFFUSERS_VERSION: ${DIFFUSERS_VERSION:-0.40.0} + TRANSFORMERS_VERSION: ${TRANSFORMERS_VERSION:-5.15.1} + ACCELERATE_VERSION: ${ACCELERATE_VERSION:-1.14.0} + HF_HUB_VERSION: ${HF_HUB_VERSION:-1.28.0} + image: mike-ai/flux-worker:local + container_name: mike-ai-flux-worker + restart: "no" + profiles: [image] + labels: + com.mike-ai.image-worker: flux + gpus: all + read_only: true + tmpfs: ["/tmp:size=1g,mode=1777"] + volumes: + - "${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B}:/models/FLUX.2-klein-4B:ro" + - router-images:/data/images + environment: + NVIDIA_VISIBLE_DEVICES: ${IMAGE_GPU_DEVICES:-1} + NVIDIA_DRIVER_CAPABILITIES: compute,utility + WORKER_TOKEN: "${CONTROLLER_TOKEN:?CONTROLLER_TOKEN is required}" + FLUX_MODEL_DIR: /models/FLUX.2-klein-4B + IMAGE_DIR: /data/images + networks: [inference] + security_opt: ["no-new-privileges:true"] + cap_drop: [ALL] + healthcheck: + test: [CMD, python, -c, "import urllib.request; urllib.request.urlopen('http://127.0.0.1:8086/health', timeout=2)"] + interval: 5s + timeout: 3s + retries: 12 + piper: build: context: platform/docker/piper diff --git a/config/install.env.example b/config/install.env.example index 63aef27..4e00f4f 100644 --- a/config/install.env.example +++ b/config/install.env.example @@ -15,6 +15,8 @@ NVIDIA_DRIVER_BRANCH= NVIDIA_MIN_DRIVER_MAJOR=570 TEXT_GPU_DEVICES=0 SECONDARY_GPU_DEVICES=1 +IMAGE_GPU_DEVICES=0 +FLUX_MODEL_DIR=/data/models/FLUX.2-klein-4B # WireGuard client. The home peer must route 10.77.0.2/32 back to this host. WIREGUARD_ENABLE=true diff --git a/dev/test_profile_controller.py b/dev/test_profile_controller.py index 4915ba0..421afb0 100644 --- a/dev/test_profile_controller.py +++ b/dev/test_profile_controller.py @@ -21,6 +21,11 @@ def item(profile, state="exited"): } +def image_item(state="exited"): + return {"Id": "id-flux", "State": state, + "Labels": {controller.IMAGE_LABEL_KEY: controller.IMAGE_WORKER}} + + class ProfileControllerTests(unittest.TestCase): def test_rejects_unknown_profile_before_docker_call(self): with patch.object(controller, "docker_request") as request: @@ -38,6 +43,7 @@ class ProfileControllerTests(unittest.TestCase): return 204, b"" with patch.object(controller, "containers", return_value=profiles), \ + patch.object(controller, "image_container", return_value=image_item()), \ patch.object(controller, "docker_request", side_effect=request): result = controller.activate("medium") @@ -49,10 +55,47 @@ class ProfileControllerTests(unittest.TestCase): def test_fails_if_profile_container_is_missing(self): profiles = {name: item(name) for name in controller.ALLOWED[:-1]} - with patch.object(controller, "containers", return_value=profiles): + with patch.object(controller, "containers", return_value=profiles), \ + patch.object(controller, "image_container", return_value=image_item()): with self.assertRaisesRegex(RuntimeError, "missing"): controller.activate("fast") + def test_image_start_stops_inference_first(self): + profiles = {name: item(name) for name in controller.ALLOWED} + profiles["medium"] = item("medium", "running") + calls = [] + + def request(method, path): + calls.append((method, path)) + return 204, b"" + + with patch.object(controller, "containers", return_value=profiles), \ + patch.object(controller, "image_container", return_value=image_item()), \ + patch.object(controller, "docker_request", side_effect=request): + controller.set_image_worker(True) + self.assertEqual(calls, [ + ("POST", "/containers/id-medium/stop?t=120"), + ("POST", "/containers/id-flux/start"), + ]) + + def test_profile_activation_stops_image_worker_first(self): + profiles = {name: item(name) for name in controller.ALLOWED} + calls = [] + + def request(method, path): + calls.append((method, path)) + return 204, b"" + + with patch.object(controller, "containers", return_value=profiles), \ + patch.object(controller, "image_container", + return_value=image_item("running")), \ + patch.object(controller, "docker_request", side_effect=request): + controller.activate("fast") + self.assertEqual(calls, [ + ("POST", "/containers/id-flux/stop?t=120"), + ("POST", "/containers/id-fast/start"), + ]) + if __name__ == "__main__": unittest.main() diff --git a/docs/CURRENT_REFERENCE.md b/docs/CURRENT_REFERENCE.md index 1797cb9..e501aa4 100644 --- a/docs/CURRENT_REFERENCE.md +++ b/docs/CURRENT_REFERENCE.md @@ -54,9 +54,9 @@ Zielplattform. | Profil | Virtuelles Modell | Kontext | Besonderheit | |---|---|---:|---| -| Fast | `qwen-fast` | 76.800 | IQ4-MIX, MTP2, vollständig GPU, CPU-mmproj | -| Medium **(Standard)** | `qwen-medium` | 160.000 | IQ4_XS Pure, MTP3, beide GPUs 90:10, CPU-mmproj | -| Large | `qwen-large` | 192.000 | IQ4_XS Pure, MTP3, beide GPUs 86:14, CPU-mmproj | +| Fast | `qwen-fast` | 76.800 | IQ4-MIX, MTP2, Text auf RTX 5080, mmproj auf RTX 3060 | +| Medium **(Standard)** | `qwen-medium` | 160.000 | IQ4_XS Pure, MTP3, beide GPUs 90:10, mmproj auf RTX 3060 | +| Large | `qwen-large` | 192.000 | IQ4_XS Pure, MTP3, beide GPUs 86:14, mmproj auf RTX 3060 | | Ultra | `qwen-ultra` | 262.144 | IQ4_XS Pure, MTP2, beide GPUs 80:20, text-only; 68,2 Tok/s und 220K-Fülltest bestanden | ## Router @@ -85,7 +85,7 @@ Der Router übernimmt: |---|---| | Text-/Visionmodell | jeweils aktives Qwen3.8-27B-Profil | | Projektor | BF16-mmproj | -| Speicherort des Projektors | System-RAM (`--no-mmproj-offload`) | +| Speicherort des Projektors | RTX 3060 (`MTMD_BACKEND_DEVICE=CUDA1`) | | Kontext | entspricht Fast/Medium/Large; Ultra ist bewusst text-only | Vision ist Bestandteil von Fast, Medium und Large. Der Router prüft Bildgröße und URL, @@ -93,13 +93,13 @@ leitet das Bild dann direkt weiter und führt keinen Modellwechsel mehr aus. ## Bildgenerierung -- Modell: FLUX.2 klein Base 4B -- Runtime: PyTorch/Diffusers -- CPU-Offload aktiviert -- Standard: 30 Schritte -- High: 50 Schritte +- Modell: FLUX.2 Klein 4B Distilled, Apache-2.0 +- Runtime: eigener PyTorch-2.11/CUDA-12.8-/Diffusers-0.40-Container +- fest auf vier Schritte und Guidance 1,0 destilliert +- Worker läuft ausschließlich auf der RTX 5080 und ist im Normalbetrieb gestoppt +- der Controller beendet Qwen vor dem Job; Bildprompts bleiben im internen Netz - Worker wird nach jedem Job vollständig beendet -- Qwen wird anschließend mit dem vorherigen Profil wiederhergestellt +- Qwen wird anschließend mit exakt dem vorherigen Profil wiederhergestellt ## Sprache diff --git a/docs/STANDARD_PROFILE_MATRIX.md b/docs/STANDARD_PROFILE_MATRIX.md index e177cfd..61fb81c 100644 --- a/docs/STANDARD_PROFILE_MATRIX.md +++ b/docs/STANDARD_PROFILE_MATRIX.md @@ -6,9 +6,9 @@ und einer bewussten Aktualisierung dieser Datei. | Profil | Virtuelles Modell | GGUF | Kontext | GPUs / Split | MTP | Vision | gemessene kurze Ausgabe | |---|---|---|---:|---|---:|---|---:| -| Fast | `qwen-fast` | IQ4-MIX | 76.800 | RTX 5080 | 2 | ja, Projektor auf CPU | 85,5 Tok/s | -| **Medium (Default)** | `qwen-medium` | IQ4_XS Pure | 160.000 | RTX 5080 + RTX 3060, 90:10 | 3 | ja, Projektor auf CPU | 77,2 Tok/s | -| Large | `qwen-large` | IQ4_XS Pure | 192.000 | RTX 5080 + RTX 3060, 86:14 | 3 | ja, Projektor auf CPU | 75,3 Tok/s | +| Fast | `qwen-fast` | IQ4-MIX | 76.800 | RTX 5080 | 2 | ja, Projektor auf RTX 3060 | 85,5 Tok/s | +| **Medium (Default)** | `qwen-medium` | IQ4_XS Pure | 160.000 | RTX 5080 + RTX 3060, 90:10 | 3 | ja, Projektor auf RTX 3060 | 77,2 Tok/s | +| Large | `qwen-large` | IQ4_XS Pure | 192.000 | RTX 5080 + RTX 3060, 86:14 | 3 | ja, Projektor auf RTX 3060 | 75,3 Tok/s | | Ultra | `qwen-ultra` | IQ4_XS Pure | 262.144 | RTX 5080 + RTX 3060, 80:20 | 2 | nein, text-only | 68,2 Tok/s | ## Standardverhalten diff --git a/install.sh b/install.sh index 258c3fa..109891a 100755 --- a/install.sh +++ b/install.sh @@ -226,6 +226,8 @@ LARGE_TENSOR_SPLIT=${LARGE_TENSOR_SPLIT:-86,14} ULTRA_GPU_DEVICES=${TEXT_GPU_DEVICES:-0}${SECONDARY_GPU_DEVICES:+,$SECONDARY_GPU_DEVICES} ULTRA_TENSOR_SPLIT=${ULTRA_TENSOR_SPLIT:-80,20} EXPERIMENTAL_GPU_DEVICES=${TEXT_GPU_DEVICES:-0} +IMAGE_GPU_DEVICES=${IMAGE_GPU_DEVICES:-${TEXT_GPU_DEVICES:-0}} +FLUX_MODEL_DIR=${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B} LLAMA_THREADS=${LLAMA_THREADS:-6} LLAMA_THREADS_BATCH=${LLAMA_THREADS_BATCH:-6} EOF @@ -326,12 +328,23 @@ build_and_start() { cd "$STACK_DIR" docker build --progress=plain --build-arg LLAMA_CPP_COMMIT="$commit" \ -f platform/docker/llama-cpp/Dockerfile -t mike-ai/llama.cpp:local . + docker compose --env-file "$SECRETS_DIR/stack.env" --profile image build flux-worker + if [[ ! -s ${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B}/model_index.json ]]; then + log "FLUX.2 Klein Distilled laden" + install -d -m 0755 "${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B}" + docker run --rm --entrypoint python \ + -v "${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B}:/download" \ + mike-ai/flux-worker:local -c \ + "from huggingface_hub import snapshot_download; snapshot_download('black-forest-labs/FLUX.2-klein-4B', revision='303481f0390afb112393f9d77e8f0be72fcefeb7', local_dir='/download')" + chmod -R a-w "${FLUX_MODEL_DIR:-/data/models/FLUX.2-klein-4B}" + fi # Creates the shared internal tools network before Open WebUI is created. # Web search always starts; HA/ARR/Unraid only start when their root-only # secret files and required local artifacts are present. "$STACK_DIR/platform/mcp/install-tools.sh" docker compose --env-file "$SECRETS_DIR/stack.env" --profile inference create \ llama-fast llama-medium llama-large llama-ultra llama-experimental + docker compose --env-file "$SECRETS_DIR/stack.env" --profile image create flux-worker docker compose --env-file "$SECRETS_DIR/stack.env" up -d --build \ profile-controller router open-webui diff --git a/platform/docker/flux-worker/Dockerfile b/platform/docker/flux-worker/Dockerfile new file mode 100644 index 0000000..39d40f5 --- /dev/null +++ b/platform/docker/flux-worker/Dockerfile @@ -0,0 +1,18 @@ +FROM pytorch/pytorch:2.11.0-cuda12.8-cudnn9-runtime + +ARG DIFFUSERS_VERSION=0.40.0 +ARG TRANSFORMERS_VERSION=5.15.1 +ARG ACCELERATE_VERSION=1.14.0 +ARG HF_HUB_VERSION=1.28.0 + +RUN pip install --no-cache-dir \ + "diffusers==${DIFFUSERS_VERSION}" \ + "transformers==${TRANSFORMERS_VERSION}" \ + "accelerate==${ACCELERATE_VERSION}" \ + "huggingface-hub==${HF_HUB_VERSION}" \ + sentencepiece protobuf safetensors pillow && \ + useradd --system --uid 10002 --home /nonexistent --shell /usr/sbin/nologin flux + +COPY flux_worker.py /app/flux_worker.py +USER 10002:10002 +ENTRYPOINT ["python", "/app/flux_worker.py"] diff --git a/platform/docker/flux-worker/flux_worker.py b/platform/docker/flux-worker/flux_worker.py new file mode 100644 index 0000000..8ac92ee --- /dev/null +++ b/platform/docker/flux-worker/flux_worker.py @@ -0,0 +1,109 @@ +#!/usr/bin/env python3 +"""Private FLUX.2 Klein Distilled worker used only during a GPU hot swap.""" + +from __future__ import annotations + +import gc +import json +import os +import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path + +HOST = os.environ.get("WORKER_HOST", "0.0.0.0") +PORT = int(os.environ.get("WORKER_PORT", "8086")) +TOKEN = os.environ.get("WORKER_TOKEN", "").strip() +MODEL_DIR = os.environ.get("FLUX_MODEL_DIR", "/models/FLUX.2-klein-4B") +OUTPUT_DIR = Path(os.environ.get("IMAGE_DIR", "/data/images")).resolve() +PIPE = None +LOAD_SECONDS = 0.0 + +if len(TOKEN) < 32: + raise RuntimeError("WORKER_TOKEN is missing or too short") + + +def load_pipeline() -> None: + global PIPE, LOAD_SECONDS + if PIPE is not None: + return + import torch + from diffusers import DiffusionPipeline + started = time.monotonic() + PIPE = DiffusionPipeline.from_pretrained( + MODEL_DIR, torch_dtype=torch.bfloat16, device_map="cuda") + LOAD_SECONDS = time.monotonic() - started + + +def generate(data: dict) -> dict: + import torch + prompt = data.get("prompt") + filename = data.get("filename") + if not isinstance(prompt, str) or not prompt.strip() or len(prompt) > 8000: + raise ValueError("invalid prompt") + if (not isinstance(filename, str) or Path(filename).name != filename + or not filename.endswith(".png")): + raise ValueError("invalid filename") + width, height = int(data.get("width", 1024)), int(data.get("height", 1024)) + if (width, height) not in {(1024, 1024), (1536, 1024), (1024, 1536), + (1920, 1088), (1088, 1920)}: + raise ValueError("unsupported image size") + steps = int(data.get("steps", 4)) + guidance = float(data.get("guidance", 1.0)) + if steps != 4 or guidance != 1.0: + raise ValueError("distilled FLUX.2 Klein requires steps=4 and guidance=1.0") + seed = data.get("seed") + generator = None if seed is None else torch.Generator(device="cuda").manual_seed(int(seed)) + load_pipeline() + started = time.monotonic() + image = PIPE(prompt=prompt, height=height, width=width, + num_inference_steps=4, guidance_scale=1.0, + generator=generator).images[0] + OUTPUT_DIR.mkdir(parents=True, exist_ok=True) + output = OUTPUT_DIR / filename + image.save(output) + return {"status": "ok", "filename": filename, + "seconds": round(time.monotonic() - started, 3), + "load_seconds": round(LOAD_SECONDS, 3)} + + +class Handler(BaseHTTPRequestHandler): + def log_message(self, fmt: str, *args: object) -> None: + # Never log request bodies/prompts. + print(f"[flux-worker] {self.client_address[0]} {fmt % args}", flush=True) + + def reply(self, status: int, payload: dict) -> None: + body = json.dumps(payload, separators=(",", ":")).encode() + self.send_response(status) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def do_GET(self) -> None: # noqa: N802 + if self.path == "/health": + self.reply(200, {"status": "ok", "model_loaded": PIPE is not None}) + else: + self.reply(404, {"error": "not found"}) + + def do_POST(self) -> None: # noqa: N802 + if self.headers.get("Authorization", "") != f"Bearer {TOKEN}": + self.reply(401, {"error": "unauthorized"}) + return + if self.path != "/generate": + self.reply(404, {"error": "not found"}) + return + try: + length = int(self.headers.get("Content-Length", "0")) + if length < 2 or length > 16384: + raise ValueError("invalid request size") + self.reply(200, generate(json.loads(self.rfile.read(length)))) + except Exception as exc: + self.reply(400, {"status": "error", "message": str(exc)}) + + +try: + ThreadingHTTPServer((HOST, PORT), Handler).serve_forever() +finally: + if PIPE is not None: + del PIPE + gc.collect() diff --git a/platform/docker/profile-controller/profile_controller.py b/platform/docker/profile-controller/profile_controller.py index 7d97d2d..aca4756 100644 --- a/platform/docker/profile-controller/profile_controller.py +++ b/platform/docker/profile-controller/profile_controller.py @@ -23,6 +23,8 @@ TOKEN_FILE = os.environ.get("CONTROLLER_TOKEN_FILE", "/run/secrets/controller-to ALLOWED = tuple(x.strip() for x in os.environ.get( "ALLOWED_PROFILES", "fast,medium,large,ultra,experimental").split(",") if x.strip()) LABEL_KEY = "com.mike-ai.llama-profile" +IMAGE_LABEL_KEY = "com.mike-ai.image-worker" +IMAGE_WORKER = os.environ.get("IMAGE_WORKER", "flux") LOCK = threading.Lock() log = logging.getLogger("profile-controller") @@ -58,6 +60,56 @@ def containers() -> dict[str, dict]: return result +def labelled_containers(label: str) -> list[dict]: + filters = urllib.parse.quote(json.dumps({"label": [label]})) + status, body = docker_request("GET", f"/containers/json?all=1&filters={filters}") + if status != 200: + raise RuntimeError(f"Docker list failed with HTTP {status}") + return json.loads(body) + + +def image_container() -> dict: + matches = [item for item in labelled_containers(IMAGE_LABEL_KEY) + if item.get("Labels", {}).get(IMAGE_LABEL_KEY) == IMAGE_WORKER] + if len(matches) != 1: + raise RuntimeError( + f"expected exactly one image worker {IMAGE_WORKER!r}, found {len(matches)}") + return matches[0] + + +def stop_container(item: dict, timeout: int = 120) -> None: + if item.get("State") != "running": + return + status, _ = docker_request("POST", f"/containers/{item['Id']}/stop?t={timeout}") + if status not in (204, 304): + raise RuntimeError(f"failed to stop container: HTTP {status}") + + +def stop_inference() -> dict: + with LOCK: + items = containers() + previous = active_profile(items) + for item in items.values(): + stop_container(item) + return {"active_profile": None, "previous_profile": previous} + + +def set_image_worker(running: bool) -> dict: + with LOCK: + item = image_container() + if running: + # A FLUX worker may never overlap a llama profile on the 5080. + for profile_item in containers().values(): + stop_container(profile_item) + if item.get("State") != "running": + status, _ = docker_request("POST", f"/containers/{item['Id']}/start") + if status not in (204, 304): + raise RuntimeError(f"failed to start image worker: HTTP {status}") + else: + stop_container(item) + return {"image_worker": "running" if running else "stopped"} + + def active_profile(items: dict[str, dict] | None = None) -> str | None: items = items or containers() active = [name for name, item in items.items() if item.get("State") == "running"] @@ -70,6 +122,8 @@ def activate(profile: str) -> dict: if profile not in ALLOWED: raise ValueError("profile is not allowlisted") with LOCK: + # Defensive mutual exclusion even if a caller bypasses the router. + stop_container(image_container()) items = containers() missing = [name for name in ALLOWED if name not in items] if missing: @@ -144,6 +198,20 @@ class Handler(BaseHTTPRequestHandler): if not self.authenticated(): self.reply(401, {"error": "unauthorized"}) return + if self.path == "/inference/stop": + try: + self.reply(200, stop_inference()) + except Exception as exc: + log.exception("stopping inference failed") + self.reply(503, {"error": str(exc)}) + return + if self.path in ("/workers/image/start", "/workers/image/stop"): + try: + self.reply(200, set_image_worker(self.path.endswith("/start"))) + except Exception as exc: + log.exception("image worker transition failed") + self.reply(503, {"error": str(exc)}) + return prefix, suffix = "/profiles/", "/activate" if not self.path.startswith(prefix) or not self.path.endswith(suffix): self.reply(404, {"error": "not found"}) diff --git a/platform/models/manifest.example.yaml b/platform/models/manifest.example.yaml index a9bb4ec..6c6cea4 100644 --- a/platform/models/manifest.example.yaml +++ b/platform/models/manifest.example.yaml @@ -28,9 +28,9 @@ models: sha256: "REPLACE_AFTER_VERIFICATION" flux: role: image-generation - source: black-forest-labs/FLUX.2-klein-base-4B - target: /opt/mike-ai/models/FLUX.2-klein-base-4B - revision: "PIN_EXACT_REVISION" + source: black-forest-labs/FLUX.2-klein-4B + target: /data/models/FLUX.2-klein-4B + revision: "303481f0390afb112393f9d77e8f0be72fcefeb7" xtts: role: text-to-speech source: coqui/XTTS-v2 diff --git a/router/ai_profile_router.py b/router/ai_profile_router.py index 33b1ca9..becff7a 100755 --- a/router/ai_profile_router.py +++ b/router/ai_profile_router.py @@ -119,6 +119,8 @@ 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_WORKER_URL = os.environ.get("IMAGE_WORKER_URL", "").rstrip("/") +IMAGE_WORKER_TOKEN = os.environ.get("IMAGE_WORKER_TOKEN", "").strip() IMAGE_DIR = os.environ.get( "IMAGE_DIR", "/opt/mike-ai/ai-profile-router/images") IMAGE_WORKER_LOG = os.environ.get( @@ -147,10 +149,10 @@ IMAGE_SIZES = { "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} +# FLUX.2 Klein Distilled ist fest auf vier Schritte und Guidance 1.0 +# destilliert. Qualitätsstufen bleiben aus OpenAI-Kompatibilitätsgründen +# akzeptiert, ändern aber bewusst nicht die offiziellen Sampling-Werte. +IMAGE_QUALITY = {"standard": 4, "high": 4} IMAGE_DEFAULT_QUALITY = "standard" IMAGE_MAX_N = 4 @@ -693,11 +695,66 @@ def _worker() -> _Worker: if not img.worker or not img.worker.alive(): if img.worker: img.worker.stop() - img.worker = _Worker() + img.worker = _RemoteWorker() if IMAGE_WORKER_URL else _Worker() img.worker.start() return img.worker +class _RemoteWorker: + """Docker-Worker, dessen Lebenszyklus nur der Controller steuert.""" + + model_loaded = False + + def __init__(self) -> None: + self.running = False + + def alive(self) -> bool: + return self.running + + def _request(self, method: str, path: str, payload: dict | None = None, + timeout: float = 120) -> dict: + body = None if payload is None else json.dumps(payload).encode() + headers = {"Authorization": f"Bearer {IMAGE_WORKER_TOKEN}"} + if body is not None: + headers["Content-Type"] = "application/json" + req = urllib.request.Request(IMAGE_WORKER_URL + path, data=body, + method=method, headers=headers) + with urllib.request.urlopen(req, timeout=timeout) as response: + return json.load(response) + + def start(self) -> None: + if not IMAGE_WORKER_TOKEN or len(IMAGE_WORKER_TOKEN) < 32: + raise RuntimeError("Bild-Worker-Token fehlt oder ist zu kurz") + _profile_controller_request("POST", "/workers/image/start") + deadline = time.monotonic() + IMAGE_START_TIMEOUT + while time.monotonic() < deadline: + try: + self._request("GET", "/health", timeout=3) + self.running = True + return + except (OSError, urllib.error.URLError, TimeoutError): + time.sleep(1) + self.stop() + raise RuntimeError("Bild-Worker hat nicht gestartet") + + def request(self, payload: dict, timeout: float) -> dict: + if payload.get("cmd") != "generate": + raise RuntimeError("Remote-Bild-Worker erlaubt nur generate") + clean = dict(payload) + clean.pop("cmd", None) + output = clean.pop("output", "") + clean["filename"] = os.path.basename(output) + return self._request("POST", "/generate", clean, timeout) + + def stop(self) -> None: + try: + _profile_controller_request("POST", "/workers/image/stop") + finally: + self.running = False + self.model_loaded = False + RUNTIME.clear_worker("image") + + def _wait_upstream_down(deadline: float) -> None: """Wartet, bis llama.cpp den Port freigegeben hat (VRAM frei).""" while time.monotonic() < deadline: @@ -753,6 +810,11 @@ def _wait_vram_free(threshold_mib: int = 1000, 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) + if PROFILE_CONTROL_URL: + _profile_controller_request("POST", f"/profiles/{profile}/activate") + _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) + RUNTIME.save(last_profile=profile, phase="idle") + return try: proc = subprocess.run([SYSTEMCTL_BIN, "start", LLAMA_SERVICE], stdin=subprocess.DEVNULL, @@ -799,15 +861,18 @@ def generate_image(prompt: str, width: int, height: int, steps: int, # 1) Qwen stoppen (VRAM freigeben). img.phase = "stopping-qwen" 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:]}") + if PROFILE_CONTROL_URL: + _profile_controller_request("POST", "/inference/stop") + else: + 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). @@ -848,7 +913,7 @@ def generate_image(prompt: str, width: int, height: int, steps: int, "guidance": guidance, "quality": quality, "seconds": resp.get("seconds"), - "model": "FLUX.2-klein-base-4B", + "model": "FLUX.2-klein-4B", "created": time.strftime("%Y-%m-%dT%H:%M:%S"), } meta_path = os.path.join(IMAGE_DIR, filename[:-4] + ".json") @@ -1363,19 +1428,19 @@ class Handler(BaseHTTPRequestHandler): "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", + if not isinstance(steps, int) or isinstance(steps, bool) or steps != 4: + self._send_error(400, "FLUX.2 Klein Distilled erfordert 'steps'=4", "invalid_request_error", "invalid_steps") return - guidance = data.get("guidance", 4.0) + guidance = data.get("guidance", 1.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", + if guidance != 1.0: + self._send_error(400, "FLUX.2 Klein Distilled erfordert 'guidance'=1.0", "invalid_request_error", "invalid_guidance") return