#!/usr/bin/env python3 """AI Profile Router – OpenAI-kompatibler Proxy vor llama.cpp. Leitet OpenAI-kompatible Requests transparent an den lokalen llama.cpp-Server weiter (Streaming, Tool Calls, JSON) und schaltet zwischen fünf festen Profilen um: Profil Kontext ------ -------- fast 76800 medium 160000 large 192000 ultra 262144 uncensored 80000 Virtuelle Modelle: qwen-fast, qwen-medium, qwen-large, qwen-ultra, qwen-uncensored Kommandos: POST /fast, /medium, /large, /ultra, /uncensored GET /status (Zustand) Bildgenerierung und Editing (Qwen-Image-2.1 INT8): POST /v1/images/generations (OpenAI-kompatibel) POST /v1/images/edits (lokal, Referenzbilder) GET /images (Liste) GET /images/ (PNG-Download) Sprachausgabe (Qwen3-TTS auf RTX 3060): POST /v1/audio/speech (OpenAI-kompatibel) GET /v1/audio/voices (verfügbare Stimmen) Spracherkennung (whisper.cpp, deutsch, CPU-only): POST /v1/audio/transcriptions (OpenAI-kompatibel) GET /v1/audio/models (verfügbare Audio-Modelle) Der TTS-Worker (mike-ai-xtts.service) und der STT-Worker (mike-ai-whisper.service) laufen als separate, langlebige Prozesse. Der Router leitet /v1/audio/speech und /v1/audio/transcriptions per HTTP an die Worker weiter. Der Router agiert als Modell-Orchestrator: vor der Generierung wird llama.cpp und Qwen3-TTS gestoppt, der Bild-Worker lädt Qwen-Image und den Textencoder auf die RTX 5080, generiert/bearbeitet 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). Vision: POST /v1/chat/completions mit Bild wird direkt an das aktive multimodale Qwen-Profil weitergeleitet. Der Vision-Projektor ist Bestandteil des Profils; es findet kein Modellwechsel statt. Nur Python-Standardbibliothek. Logging nach stdout (journald). """ from __future__ import annotations import base64 import binascii import email import ipaddress import json import logging import os import queue import re import select import socket import subprocess import sys import threading import time import uuid import http.client from contextlib import contextmanager import urllib.error import urllib.parse import urllib.request from email.parser import BytesParser from email.policy import compat32 from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from router_support import ( AuthPolicy, ConfigurationError, RuntimeStore, enforce_artifact_retention, load_profile_registry, terminate_recorded_worker, ) # --------------------------------------------------------------------------- # Konfiguration (über Umgebungsvariablen, vgl. systemd-Unit) # --------------------------------------------------------------------------- HOST = os.environ.get("ROUTER_HOST", "0.0.0.0") PORT = int(os.environ.get("ROUTER_PORT", "8081")) UPSTREAM_URL = os.environ.get("UPSTREAM_URL", "http://127.0.0.1:8080").rstrip("/") REVIEW_UPSTREAM_URL = os.environ.get("REVIEW_UPSTREAM_URL", "").rstrip("/") REVIEW_MODEL_NAME = os.environ.get("REVIEW_MODEL_NAME", "qwen-review").strip() REVIEW_CONTEXT_LENGTH = int(os.environ.get("REVIEW_CONTEXT_LENGTH", "32768")) PROFILE_SCRIPT = os.environ.get("PROFILE_SCRIPT", "/usr/local/bin/llama-profile") PROFILE_DIR = os.environ.get( "PROFILE_DIR", "/etc/systemd/system/mike-ai-llama-ui.service.d") PROFILE_CONTROL_URL = os.environ.get("PROFILE_CONTROL_URL", "").rstrip("/") PROFILE_CONTROL_TOKEN_FILE = os.environ.get( "PROFILE_CONTROL_TOKEN_FILE", "/run/secrets/controller-token") ENABLE_MUSIC_MODE = os.environ.get( "ENABLE_MUSIC_MODE", "false").lower() in {"1", "true", "yes"} MUSIC_START_TIMEOUT = float(os.environ.get("MUSIC_START_TIMEOUT", "600")) YUE2_START_TIMEOUT = float(os.environ.get("YUE2_START_TIMEOUT", "600")) SEPARATOR_START_TIMEOUT = float(os.environ.get("SEPARATOR_START_TIMEOUT", "600")) VOICE_START_TIMEOUT = float(os.environ.get("VOICE_START_TIMEOUT", "600")) VOICE_CHANGE_START_TIMEOUT = float(os.environ.get("VOICE_CHANGE_START_TIMEOUT", "600")) APPLIO_START_TIMEOUT = float(os.environ.get("APPLIO_START_TIMEOUT", "900")) TRELLIS_START_TIMEOUT = float(os.environ.get("TRELLIS_START_TIMEOUT", "900")) VIDEO_START_TIMEOUT = float(os.environ.get("VIDEO_START_TIMEOUT", "900")) # Optional worker APIs. The clean Docker baseline deliberately ships only # text/multimodal chat; absent workers must fail explicitly instead of trying # legacy systemd paths inside the container. ENABLE_IMAGE_GENERATION = os.environ.get( "ENABLE_IMAGE_GENERATION", "true").lower() in {"1", "true", "yes"} ENABLE_TTS = os.environ.get( "ENABLE_TTS", "true").lower() in {"1", "true", "yes"} ENABLE_STT = os.environ.get( "ENABLE_STT", "true").lower() in {"1", "true", "yes"} SWITCH_TIMEOUT = float(os.environ.get("SWITCH_TIMEOUT", "600")) # s, Warten auf llama.cpp REQUEST_TIMEOUT = float(os.environ.get("REQUEST_TIMEOUT", "600")) # s, Read-Timeout Upstream CONNECT_TIMEOUT = float(os.environ.get("CONNECT_TIMEOUT", "10")) # s, Connect-Timeout POLL_INTERVAL = float(os.environ.get("POLL_INTERVAL", "2")) # s, Polling-Intervall MAX_GENERATION_TOKENS = int(os.environ.get("MAX_GENERATION_TOKENS", "8192")) DEFAULT_REASONING_EFFORT = os.environ.get( "DEFAULT_REASONING_EFFORT", "off").strip().lower() GLOBAL_SYSTEM_POLICY_FILE = os.environ.get( "GLOBAL_SYSTEM_POLICY_FILE", "").strip() # --- Bildgenerierung und Referenzbild-Bearbeitung (Qwen-Image-2.1 INT8) --- 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_WORKER_URL = os.environ.get("IMAGE_WORKER_URL", "").rstrip("/") IMAGE_WORKER_TOKEN = os.environ.get("IMAGE_WORKER_TOKEN", "").strip() IMAGE_MODEL_NAME = os.environ.get( "IMAGE_MODEL_NAME", "Qwen-Image-2.1-int8") IMAGE_INFERENCE_STEPS = int(os.environ.get("IMAGE_INFERENCE_STEPS", "25")) 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 IMAGE_RETENTION_FILES = int(os.environ.get("IMAGE_RETENTION_FILES", "100")) IMAGE_RETENTION_BYTES = int(os.environ.get( "IMAGE_RETENTION_BYTES", str(5 * 1024 * 1024 * 1024))) IMAGE_RETENTION_DAYS = int(os.environ.get("IMAGE_RETENTION_DAYS", "30")) # --- Multimodale Chat-Eingaben --- # Bilder werden validiert und direkt an das aktive Qwen-Profil weitergeleitet. CHAT_IMAGE_MAX_BYTES = int(os.environ.get( "CHAT_IMAGE_MAX_BYTES", str(20 * 1024 * 1024))) CHAT_IMAGE_ALLOW_REMOTE_URLS = os.environ.get( "CHAT_IMAGE_ALLOW_REMOTE_URLS", "false").lower() in {"1", "true", "yes"} # Erlaubte Auflösungen (Breite x Höhe). IMAGE_SIZES = { "1024x1024": (1024, 1024), "1536x1024": (1536, 1024), "1024x1536": (1024, 1536), "1920x1088": (1920, 1088), "1088x1920": (1088, 1920), } # Der produktive Qwen-Image-2.1-Workflow ist auf 25 Schritte festgelegt. IMAGE_QUALITY = {"standard": IMAGE_INFERENCE_STEPS, "high": IMAGE_INFERENCE_STEPS} IMAGE_DEFAULT_QUALITY = "standard" IMAGE_MAX_N = 4 # --- Sprachausgabe (Qwen3-TTS über das interne Normalisierungs-Gateway) --- TTS_WORKER_URL = os.environ.get("TTS_WORKER_URL", "http://127.0.0.1:8085") TTS_TIMEOUT = float(os.environ.get("TTS_TIMEOUT", "300")) # s, pro Synthese TTS_CONNECT_TIMEOUT = float(os.environ.get("TTS_CONNECT_TIMEOUT", "5")) TTS_MODEL = os.environ.get("TTS_MODEL", "qwen3-tts") TTS_VOICES = tuple(v.strip() for v in os.environ.get( "TTS_VOICES", "alloy").split(",") if v.strip()) TTS_DEFAULT_VOICE = os.environ.get( "TTS_DEFAULT_VOICE", TTS_VOICES[0] if TTS_VOICES else "alloy") TTS_FORMATS = ("mp3", "wav", "pcm") TTS_DEFAULT_FORMAT = "mp3" # --- Spracherkennung (whisper.cpp, deutsch, CPU-only) --- STT_WORKER_URL = os.environ.get("STT_WORKER_URL", "http://127.0.0.1:8084") STT_TIMEOUT = float(os.environ.get("STT_TIMEOUT", "120")) # s, pro Transkription STT_CONNECT_TIMEOUT = float(os.environ.get("STT_CONNECT_TIMEOUT", "5")) STT_MODEL = "whisper-1" # virtuelles Modell für /v1/audio/transcriptions # Maximale Upload-Größe (Bytes) – verhindert unbegrenzten RAM-Verbrauch. # 50 MB ist für Audio-Dateien (WebM/Opus, WAV, MP3) mehr als ausreichend. MAX_UPLOAD_SIZE = int(os.environ.get("MAX_UPLOAD_SIZE", 50 * 1024 * 1024)) # 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 RUNTIME_STATE_FILE = os.environ.get( "ROUTER_STATE_FILE", "/var/lib/mike-ai-profile-router/state.json") PROFILE_REGISTRY_FILE = os.environ.get( "ROUTER_PROFILES_FILE", "/etc/mike-ai/router-profiles.json") ALLOW_LEGACY_GET_SWITCH = os.environ.get( "ALLOW_LEGACY_GET_SWITCH", "false").lower() in {"1", "true", "yes"} MAX_CONCURRENT_REQUESTS = int(os.environ.get( "ROUTER_MAX_CONCURRENT_REQUESTS", "16")) PROFILE_REGISTRY = load_profile_registry(PROFILE_REGISTRY_FILE) PROFILES = {name: definition["context"] for name, definition in PROFILE_REGISTRY.items()} EXPECTED_MODELS = {name: definition.get("model_alias") for name, definition in PROFILE_REGISTRY.items()} VIRTUAL_MODELS = { (EXPECTED_MODELS.get(name) or f"qwen-{name}"): name for name in PROFILES } log = logging.getLogger("ai-profile-router") AUTH: AuthPolicy | None = None RUNTIME = RuntimeStore(RUNTIME_STATE_FILE) REQUEST_SLOTS = threading.BoundedSemaphore(max(1, MAX_CONCURRENT_REQUESTS)) # Hop-by-hop-Header, die nicht an Upstream/Client weitergereicht werden. HOP_BY_HOP = { "host", "connection", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailer", "transfer-encoding", "upgrade", "content-length", "authorization", "x-api-key", } def _parse_upstream(url: str) -> tuple[str, int]: """'http://127.0.0.1:8080' -> ('127.0.0.1', 8080)""" hostport = url.split("://", 1)[-1] host, _, port = hostport.partition(":") return host, int(port) if port else 80 UPSTREAM_HOST, UPSTREAM_PORT = _parse_upstream(UPSTREAM_URL) if REVIEW_UPSTREAM_URL: REVIEW_UPSTREAM_HOST, REVIEW_UPSTREAM_PORT = _parse_upstream( REVIEW_UPSTREAM_URL) else: REVIEW_UPSTREAM_HOST, REVIEW_UPSTREAM_PORT = "", 0 # --------------------------------------------------------------------------- # 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 self.current_model: str | None = None IMAGE_PHASES = ( "idle", "stopping-qwen", "loading-image", "generating", "unloading-image", "restoring-qwen", ) class _State: """Gemeinsamer, thread-sicherer Zustand. lock : zentraler GPU-/Model-Lock. Wird von Profilwechsel und Bildgenerierung gehalten. avail_lock : schützt qwen_unavailable + active_chats (Chat-Waiting). """ def __init__(self) -> None: # RLock erlaubt atomare Abläufe aus Profilwahl + Chat-Lease, # während die darunterliegenden Funktionen denselben Lock verwenden. self.lock = threading.RLock() self.switching: str | None = None self.started = time.time() self.image = _ImageState() self.qwen_unavailable = True self.active_chats = 0 self.avail_lock = threading.Lock() self.mode = "llm" self.mode_phase = "ready" self.mode_error: str | None = None STATE = _State() class ModelWaitTimeout(RuntimeError): """The request did not obtain the GPU coordinator within its wait budget.""" @contextmanager def _model_lock(): # Socket timeouts do not bound a Python lock acquisition. if not STATE.lock.acquire(timeout=max(0, CHAT_WAIT_TIMEOUT)): raise ModelWaitTimeout("Wartezeit auf das Modell überschritten; bitte erneut versuchen") try: yield finally: STATE.lock.release() 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: raise RuntimeError( f"Profil-/GPU-Wechsel nach {timeout:.0f} s abgebrochen: " f"noch {n} aktive Chat-Anfrage(n)") time.sleep(0.5) def _set_qwen_unavailable(unavailable: bool) -> None: with STATE.avail_lock: STATE.qwen_unavailable = unavailable def _music_worker_state() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "music_worker", "missing")) except Exception as exc: log.warning("Musik-Worker-Status nicht verfügbar: %s", exc) return "unknown" def _music_worker_health() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "music_health", "unknown")) except Exception: return "unknown" def _separator_worker_state() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "separator_worker", "missing")) except Exception as exc: log.warning("Stem-Separator-Status nicht verfügbar: %s", exc) return "unknown" def _separator_worker_health() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "separator_health", "unknown")) except Exception: return "unknown" def _voice_worker_state() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "voice_worker", "missing")) except Exception as exc: log.warning("Voice-Worker-Status nicht verfügbar: %s", exc) return "unknown" def _voice_worker_health() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "voice_health", "unknown")) except Exception: return "unknown" def _voice_change_worker_state() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "voice_change_worker", "missing")) except Exception as exc: log.warning("Voice-Change-Worker-Status nicht verfügbar: %s", exc) return "unknown" def _voice_change_worker_health() -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get( "voice_change_health", "unknown")) except Exception: return "unknown" def _worker_field(field: str) -> str: if not PROFILE_CONTROL_URL: return "unsupported" try: return str(_profile_controller_request("GET", "/status").get(field, "missing")) except Exception as exc: log.warning("Spezial-Worker-Status %s nicht verfügbar: %s", field, exc) return "unknown" def _wait_music_ready() -> None: deadline = time.monotonic() + MUSIC_START_TIMEOUT while time.monotonic() < deadline: status = _profile_controller_request("GET", "/status") if (status.get("music_worker") == "running" and status.get("music_health") == "healthy"): return if status.get("music_health") == "unhealthy": raise RuntimeError("ACE-Step-Container ist unhealthy") time.sleep(POLL_INTERVAL) raise RuntimeError( f"ACE-Step nach {MUSIC_START_TIMEOUT:.0f} s nicht bereit") def _wait_yue2_ready() -> None: _wait_aux_voice_ready("yue2_worker", "yue2_health", "YuE2", YUE2_START_TIMEOUT) def _wait_separator_ready() -> None: deadline = time.monotonic() + SEPARATOR_START_TIMEOUT while time.monotonic() < deadline: status = _profile_controller_request("GET", "/status") if (status.get("separator_worker") == "running" and status.get("separator_health") == "healthy"): return if status.get("separator_health") == "unhealthy": raise RuntimeError("BS-RoFormer-Container ist unhealthy") time.sleep(POLL_INTERVAL) raise RuntimeError( f"BS-RoFormer nach {SEPARATOR_START_TIMEOUT:.0f} s nicht bereit") def _wait_voice_ready() -> None: deadline = time.monotonic() + VOICE_START_TIMEOUT while time.monotonic() < deadline: status = _profile_controller_request("GET", "/status") if (status.get("voice_worker") == "running" and status.get("voice_health") == "healthy"): return if status.get("voice_health") == "unhealthy": raise RuntimeError("OmniVoice-Container ist unhealthy") time.sleep(POLL_INTERVAL) raise RuntimeError( f"OmniVoice nach {VOICE_START_TIMEOUT:.0f} s nicht bereit") def _wait_voice_change_ready() -> None: deadline = time.monotonic() + VOICE_CHANGE_START_TIMEOUT while time.monotonic() < deadline: status = _profile_controller_request("GET", "/status") if (status.get("voice_change_worker") == "running" and status.get("voice_change_health") == "healthy"): return if status.get("voice_change_health") == "unhealthy": raise RuntimeError("X-VC-Container ist unhealthy") time.sleep(POLL_INTERVAL) raise RuntimeError( f"X-VC nach {VOICE_CHANGE_START_TIMEOUT:.0f} s nicht bereit") def _wait_aux_voice_ready(worker_field: str, health_field: str, label: str, timeout: float) -> None: deadline = time.monotonic() + timeout while time.monotonic() < deadline: status = _profile_controller_request("GET", "/status") if (status.get(worker_field) == "running" and status.get(health_field) == "healthy"): return if status.get(health_field) == "unhealthy": raise RuntimeError(f"{label}-Container ist unhealthy") time.sleep(POLL_INTERVAL) raise RuntimeError(f"{label} nach {timeout:.0f} s nicht bereit") def _wait_applio_ready() -> None: _wait_aux_voice_ready("applio_worker", "applio_health", "Applio", APPLIO_START_TIMEOUT) def _wait_trellis_ready() -> None: _wait_aux_voice_ready("trellis_worker", "trellis_health", "TRELLIS.2", TRELLIS_START_TIMEOUT) def _wait_video_ready() -> None: _wait_aux_voice_ready("video_worker", "video_health", "LTX-2 Studio", VIDEO_START_TIMEOUT) def _special_worker(mode: str) -> tuple[str, str, callable]: if mode == "music": return "/workers/music/start", _music_worker_state(), _wait_music_ready if mode == "yue2": return ("/workers/yue2/start", _worker_field("yue2_worker"), _wait_yue2_ready) if mode == "separation": return "/workers/separator/start", _separator_worker_state(), _wait_separator_ready if mode == "voice": return "/workers/voice/start", _voice_worker_state(), _wait_voice_ready if mode == "voicechange": return ("/workers/voice-change/start", _voice_change_worker_state(), _wait_voice_change_ready) if mode == "applio": return ("/workers/applio/start", _worker_field("applio_worker"), _wait_applio_ready) if mode == "trellis": return ("/workers/trellis/start", _worker_field("trellis_worker"), _wait_trellis_ready) if mode == "video": return ("/workers/video/start", _worker_field("video_worker"), _wait_video_ready) raise ValueError(f"unbekannter Spezialmodus: {mode}") def set_operating_mode(mode: str) -> dict: """Atomarer Wechsel zwischen LLM und den exklusiven GPU-Werkzeugen.""" if not ENABLE_MUSIC_MODE or not PROFILE_CONTROL_URL: raise RuntimeError("Musikmodus ist nicht konfiguriert") special_modes = {"music", "yue2", "separation", "voice", "voicechange", "applio", "trellis", "video"} if mode not in {"llm", *special_modes}: raise ValueError("unbekannter Betriebsmodus") with STATE.lock: STATE.mode_error = None if mode in special_modes: path, worker_state, wait_ready = _special_worker(mode) if STATE.mode == mode and worker_state == "running": return {"status": "ok", "mode": mode, "changed": False} profile = current_profile() previous = RUNTIME.load() saved = previous.get("return_profile") or previous.get("last_profile") return_profile = profile if profile in PROFILES else saved if return_profile not in PROFILES: return_profile = next(iter(PROFILES)) STATE.mode_phase = f"starting-{mode}" _set_qwen_unavailable(True) try: # Persist intent before stopping anything so a router restart # during ACE-Step loading can resume the same transition. RUNTIME.save(mode=mode, return_profile=return_profile, last_profile=return_profile, phase=f"starting-{mode}") _wait_chats_drained() _profile_controller_request( "POST", path, timeout=120) wait_ready() STATE.mode = mode STATE.mode_phase = "ready" RUNTIME.save(mode=mode, return_profile=return_profile, last_profile=return_profile, phase=mode) return {"status": "ok", "mode": mode, "changed": True, "return_profile": return_profile} except Exception as exc: STATE.mode_error = str(exc) STATE.mode_phase = "error" raise previous = RUNTIME.load() profile = previous.get("return_profile") or previous.get("last_profile") if profile not in PROFILES: profile = next(iter(PROFILES)) STATE.mode_phase = "restoring-llm" _set_qwen_unavailable(True) try: _profile_controller_request("POST", "/workers/music/stop") _profile_controller_request("POST", "/workers/yue2/stop") _profile_controller_request("POST", "/workers/separator/stop") _profile_controller_request("POST", "/workers/voice/stop") _profile_controller_request("POST", "/workers/voice-change/stop") _profile_controller_request("POST", "/workers/applio/stop") _profile_controller_request("POST", "/workers/trellis/stop") _profile_controller_request("POST", "/workers/video/stop") _restore_qwen(profile) STATE.mode = "llm" STATE.mode_phase = "ready" _set_qwen_unavailable(False) RUNTIME.save(mode="llm", return_profile=None, last_profile=profile, phase="idle") return {"status": "ok", "mode": "llm", "changed": True, "profile": profile} except Exception as exc: STATE.mode_error = str(exc) STATE.mode_phase = "error" raise def schedule_operating_mode(mode: str) -> tuple[bool, str]: """Start a transition in the background so chat/UI acknowledgement is instant.""" if not ENABLE_MUSIC_MODE or not PROFILE_CONTROL_URL: raise RuntimeError("Musikmodus ist nicht konfiguriert") if mode not in {"llm", "music", "yue2", "separation", "voice", "voicechange", "applio", "trellis", "video"}: raise ValueError("unbekannter Betriebsmodus") with STATE.lock: if STATE.mode_phase not in {"ready", "error"}: return False, STATE.mode_phase if STATE.mode == mode and STATE.mode_phase == "ready": return False, "ready" STATE.mode_phase = f"starting-{mode}" if mode != "llm" else "restoring-llm" def transition() -> None: try: set_operating_mode(mode) log.info("Betriebsmodus ist jetzt %s", mode) except Exception: log.exception("Betriebsmodus-Wechsel zu %s fehlgeschlagen", mode) threading.Thread(target=transition, name=f"mode-{mode}", daemon=True).start() return True, STATE.mode_phase def _control_command(data: dict, path: str) -> str | None: """Recognise exact local commands without invoking an LLM.""" text: object = None if path == "/v1/chat/completions": messages = data.get("messages") if isinstance(messages, list): for item in reversed(messages): if isinstance(item, dict) and item.get("role") == "user": text = item.get("content") break elif path == "/v1/responses": text = data.get("input") if not isinstance(text, str): return None command = text.strip().casefold() return command if command in {"/athena music", "/athena yue2", "/athena stems", "/athena separation", "/athena llm", "/athena voice", "/athena voicechange", "/athena changer", "/athena applio", "/athena 3d", "/athena trellis", "/athena video", "/athena ltx", "/athena ltx2", "/athena status"} else None # --------------------------------------------------------------------------- # Upstream (llama.cpp) # --------------------------------------------------------------------------- def tts_status() -> dict: """Prüft den TTS-Worker: erreichbar? bereit? welche Stimmen?""" hostport = TTS_WORKER_URL.split("://", 1)[-1] host, _, port = hostport.partition(":") try: conn = http.client.HTTPConnection(host, int(port) if port else 80, timeout=TTS_CONNECT_TIMEOUT) conn.request("GET", "/status") resp = conn.getresponse() data = json.loads(resp.read()) conn.close() return {"reachable": True, **data} except (OSError, ValueError) as e: return {"reachable": False, "error": str(e)} def tts_synthesize(text: str, voice: str, speed: float, fmt: str) -> tuple[bytes, str]: """Synthetisiert Audio über den TTS-Worker. Liefert (audio_bytes, content_type). Wirft RuntimeError bei Fehler. """ hostport = TTS_WORKER_URL.split("://", 1)[-1] host, _, port = hostport.partition(":") payload = json.dumps({"text": text, "voice": voice, "speed": speed, "format": fmt}).encode() try: conn = http.client.HTTPConnection(host, int(port) if port else 80, timeout=TTS_CONNECT_TIMEOUT) conn.request("POST", "/tts", body=payload, headers={"Content-Type": "application/json"}) conn.sock.settimeout(TTS_TIMEOUT) resp = conn.getresponse() body = resp.read() conn.close() except (OSError, http.client.HTTPException) as e: raise RuntimeError(f"TTS-Worker nicht erreichbar: {e}") if resp.status != 200: try: err = json.loads(body) msg = err.get("error", str(err)) except ValueError: msg = body.decode(errors="replace")[:200] raise RuntimeError(f"TTS-Fehler ({resp.status}): {msg}") content_type = {"mp3": "audio/mpeg", "wav": "audio/wav", "flac": "audio/flac", "pcm": "application/octet-stream"}[fmt] return body, content_type def stt_status() -> dict: """Prüft den STT-Worker: erreichbar? bereit?""" hostport = STT_WORKER_URL.split("://", 1)[-1] host, _, port = hostport.partition(":") try: conn = http.client.HTTPConnection(host, int(port) if port else 80, timeout=STT_CONNECT_TIMEOUT) conn.request("GET", "/status") resp = conn.getresponse() data = json.loads(resp.read()) conn.close() return {"reachable": True, **data} except (OSError, ValueError) as e: return {"reachable": False, "error": str(e)} def stt_transcribe(file_data: bytes, filename: str, language: str | None = None, prompt: str | None = None, temperature: float | None = None) -> dict: """Transkribiert Audio über den STT-Worker. Liefert dict mit 'text'. Wirft RuntimeError bei Fehler. """ hostport = STT_WORKER_URL.split("://", 1)[-1] host, _, port = hostport.partition(":") # Multipart-Form-Data bauen boundary = "----STTBoundary" + uuid.uuid4().hex[:16] parts = [] parts.append( f"--{boundary}\r\n" f'Content-Disposition: form-data; name="file"; filename="{filename}"\r\n' f"Content-Type: application/octet-stream\r\n\r\n".encode("utf-8") ) parts.append(file_data) parts.append(b"\r\n") for key, value in [("language", language), ("prompt", prompt), ("temperature", temperature)]: if value is not None: parts.append( f"--{boundary}\r\n" f'Content-Disposition: form-data; name="{key}"\r\n\r\n' f"{value}\r\n".encode("utf-8") ) parts.append(f"--{boundary}--\r\n".encode("utf-8")) body = b"".join(parts) try: conn = http.client.HTTPConnection(host, int(port) if port else 80, timeout=STT_CONNECT_TIMEOUT) conn.request("POST", "/transcribe", body=body, headers={"Content-Type": f"multipart/form-data; boundary={boundary}"}) conn.sock.settimeout(STT_TIMEOUT) resp = conn.getresponse() data = json.loads(resp.read()) conn.close() except (OSError, http.client.HTTPException) as e: raise RuntimeError(f"STT-Worker nicht erreichbar: {e}") if resp.status != 200: msg = data.get("error", str(data)) if isinstance(data, dict) else str(data) raise RuntimeError(f"STT-Fehler ({resp.status}): {msg}") return data def upstream_status() -> dict: """Prüft llama.cpp: erreichbar? welches Modell? welcher Kontext?""" try: conn = http.client.HTTPConnection(UPSTREAM_HOST, UPSTREAM_PORT, timeout=CONNECT_TIMEOUT) conn.request("GET", "/v1/models") resp = conn.getresponse() data = json.loads(resp.read()) conn.close() except (OSError, ValueError) as e: return {"reachable": False, "error": str(e)} models = data.get("data") or [] if not models: return {"reachable": True, "model": None, "ctx": None} m = models[0] return {"reachable": True, "model": m.get("id"), "ctx": (m.get("meta") or {}).get("n_ctx")} _TELEMETRY_LOCK = threading.Lock() _TELEMETRY_AT = 0.0 _TELEMETRY_CACHE: dict = {} def _upstream_read(path: str, *, timeout: float = 1.5) -> tuple[int, bytes]: """Read a bounded, read-only llama.cpp telemetry endpoint.""" conn = http.client.HTTPConnection(UPSTREAM_HOST, UPSTREAM_PORT, timeout=timeout) try: conn.request("GET", path, headers={"Accept": "application/json,text/plain"}) resp = conn.getresponse() return resp.status, resp.read(2 * 1024 * 1024) finally: conn.close() def _parse_prometheus_metrics(raw: str) -> dict: wanted = { "llamacpp:prompt_tokens_total", "llamacpp:prompt_tokens_cached_total", "llamacpp:prompt_seconds_total", "llamacpp:tokens_predicted_total", "llamacpp:tokens_predicted_seconds_total", "llamacpp:n_decode_total", "llamacpp:n_tokens_max", "llamacpp:spec_decode_num_draft_tokens_total", "llamacpp:spec_decode_num_accepted_tokens_total", "llamacpp:spec_decode_num_drafts_total", "llamacpp:prompt_tokens_seconds", "llamacpp:predicted_tokens_seconds", "llamacpp:requests_processing", "llamacpp:requests_deferred", "llamacpp:n_busy_slots_per_decode", } result: dict[str, int | float] = {} for line in raw.splitlines(): if not line or line.startswith("#") or " " not in line: continue name, value = line.rsplit(None, 1) if "{" in name or name not in wanted: continue try: number = float(value) result[name.removeprefix("llamacpp:")] = ( int(number) if number.is_integer() else number ) except ValueError: continue return result def upstream_telemetry() -> dict: """Compact llama.cpp slots, rates, cache and MTP telemetry. The result is cached briefly because the dashboard refreshes every second. Failures never affect inference or the normal router status response. """ global _TELEMETRY_AT, _TELEMETRY_CACHE now = time.monotonic() with _TELEMETRY_LOCK: if now - _TELEMETRY_AT < 0.75 and _TELEMETRY_CACHE: return _TELEMETRY_CACHE result: dict = {"available": False, "slots": [], "metrics": {}} errors: dict[str, str] = {} try: status, body = _upstream_read("/slots") if status == 200: raw_slots = json.loads(body) for slot in raw_slots if isinstance(raw_slots, list) else []: next_token = (slot.get("next_token") or [{}])[0] params = slot.get("params") or {} prompt = int(slot.get("n_prompt_tokens") or 0) decoded = int(next_token.get("n_decoded") or 0) n_ctx = int(slot.get("n_ctx") or 0) result["slots"].append({ "id": slot.get("id"), "task_id": slot.get("id_task"), "processing": bool(slot.get("is_processing")), "speculative": bool(slot.get("speculative")), "n_ctx": n_ctx, "prompt_tokens": prompt, "prompt_processed": int(slot.get("n_prompt_tokens_processed") or 0), "prompt_cached": int(slot.get("n_prompt_tokens_cache") or 0), "decoded_tokens": decoded, "context_used": min(n_ctx, prompt + decoded) if n_ctx else prompt + decoded, "remaining_generation": next_token.get("n_remain"), "max_tokens": params.get("max_tokens", params.get("n_predict")), "temperature": params.get("temperature"), "stream": params.get("stream"), }) else: errors["slots"] = f"HTTP {status}" except (OSError, ValueError, KeyError, TypeError, http.client.HTTPException) as exc: errors["slots"] = str(exc) try: status, body = _upstream_read("/metrics") if status == 200: result["metrics"] = _parse_prometheus_metrics( body.decode("utf-8", "replace") ) else: errors["metrics"] = f"HTTP {status}" except (OSError, ValueError, http.client.HTTPException) as exc: errors["metrics"] = str(exc) try: status, body = _upstream_read("/props") if status == 200: props = json.loads(body) result["props"] = { "total_slots": props.get("total_slots"), "model_alias": props.get("model_alias"), "model_ftype": props.get("model_ftype"), "model_path": props.get("model_path"), "modalities": props.get("modalities") or {}, "default_context": ((props.get("default_generation_settings") or {}) .get("n_ctx")), } else: errors["props"] = f"HTTP {status}" except (OSError, ValueError, TypeError, http.client.HTTPException) as exc: errors["props"] = str(exc) result["available"] = bool(result["slots"] or result["metrics"]) if errors: result["errors"] = errors _TELEMETRY_CACHE = result _TELEMETRY_AT = now return result # --------------------------------------------------------------------------- # Profile # --------------------------------------------------------------------------- def _read(path: str) -> str: with open(path, encoding="utf-8") as f: return f.read().strip() def _profile_controller_request(method: str, path: str, timeout: float = 120) -> dict: token = os.environ.get("PROFILE_CONTROL_TOKEN", "").strip() if not token: token = _read(PROFILE_CONTROL_TOKEN_FILE) if len(token) < 32: raise RuntimeError("Profil-Controller-Token fehlt oder ist zu kurz") request = urllib.request.Request( PROFILE_CONTROL_URL + path, method=method, headers={"Authorization": f"Bearer {token}"}, ) try: with urllib.request.urlopen(request, timeout=timeout) as response: return json.load(response) except urllib.error.HTTPError as exc: body = exc.read(500).decode(errors="replace") raise RuntimeError( f"Profil-Controller HTTP {exc.code}: {body}") from exc def current_profile() -> str | None: """Aktives Profil anhand semantischer Werte der override.conf. Kommentare, Leerraum oder die Reihenfolge anderer llama.cpp-Optionen beeinflussen die Erkennung nicht mehr. """ if PROFILE_CONTROL_URL: try: profile = _profile_controller_request("GET", "/profiles/status", timeout=3).get( "active_profile") return profile if profile in PROFILES else None except Exception as exc: log.warning("Profil-Controller-Status nicht verfügbar: %s", exc) return None try: override = _read(os.path.join(PROFILE_DIR, "override.conf")) except OSError: return None ctx_match = re.search(r"(?:^|\s)--ctx-size\s+(\d+)(?:\s|$)", override) alias_match = re.search(r"(?:^|\s)--alias\s+([^\s]+)", override) if not ctx_match: return None ctx = int(ctx_match.group(1)) alias = alias_match.group(1) if alias_match else None for name, expected_ctx in PROFILES.items(): expected_alias = EXPECTED_MODELS.get(name) if ctx == expected_ctx and (not expected_alias or alias == expected_alias): return name return None def _wait_ready(profile: str, deadline: float) -> None: """Wartet, bis llama.cpp das Profil geladen hat (Modell + ctx).""" expected_ctx = PROFILES[profile] expected_model = EXPECTED_MODELS.get(profile) while True: status = upstream_status() if (status["reachable"] and status.get("model") and _context_matches(expected_ctx, status.get("ctx")) and (not expected_model or status.get("model") == expected_model)): log.info("llama.cpp bereit: Profil=%s Modell=%s ctx=%s", profile, status.get("model"), status.get("ctx")) return if time.monotonic() > deadline: raise RuntimeError( f"llama.cpp nach {SWITCH_TIMEOUT:.0f} s nicht bereit " f"(erwartet Modell {expected_model or '*'} / ctx " f"{expected_ctx}, aktuell: {status.get('model')} / " f"{status.get('ctx')})") time.sleep(POLL_INTERVAL) def _profile_is_ready(profile: str, status: dict | None = None) -> bool: status = status or upstream_status() expected_model = EXPECTED_MODELS.get(profile) return bool(status.get("reachable") and status.get("model") and _context_matches(PROFILES[profile], status.get("ctx")) and (not expected_model or status.get("model") == expected_model)) def _context_matches(expected: int, reported: object) -> bool: """Allow llama.cpp's small MTP/speculative context overhead. Recent builds can report a slot context slightly above the requested ``--ctx-size`` (currently 128 tokens with the tested Qwen MTP setup). The alias check and this narrow bound still reject neighbouring profiles. """ return (isinstance(reported, int) and expected <= reported <= expected + 1024) def switch_profile(profile: str, implicit: bool = False) -> dict: """Stellt sicher, dass das Profil aktiv ist, und wartet bis es geladen ist. Wirft RuntimeError, wenn das Profil nicht aktiviert werden konnte. implicit=True (ausgelöst durch ein virtuelles Modell in einem Chat-Request): Wenn das Profil bereits aktiv ist, aber llama.cpp down ist, wird sofort eine RuntimeError geworfen (kein stiller Neustart). Der Nutzer kann den Neustart explizit über / anstoßen. """ 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: cur = current_profile() up = upstream_status() ready = _profile_is_ready(profile, up) if cur == profile and ready: # A previous failed switch can leave the in-memory availability # flag set even though the controller and llama.cpp have since # recovered. The semantic readiness check above is authoritative, # so clear the stale flag before returning. _set_qwen_unavailable(False) log.info("Profil %s ist bereits aktiv", profile) return up # Qwen wird neu geladen/gewechselt → für Chats nicht verfügbar. _set_qwen_unavailable(True) try: _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 upstream_status() 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) if PROFILE_CONTROL_URL: _profile_controller_request( "POST", f"/profiles/{profile}/activate") else: 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: raise RuntimeError( f"llama-profile fehlgeschlagen (Exit-Code " f"{proc.returncode}): {out[-500:]}") except subprocess.TimeoutExpired: raise RuntimeError("llama-profile hat 120 s überschritten") if current_profile() != profile: raise RuntimeError( f"Profildatei wurde nicht gesetzt (erwartet: {profile})") log.info("Warte, bis llama.cpp das Profil geladen hat ...") _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) STATE.mode = "llm" STATE.mode_phase = "ready" RUNTIME.save(last_profile=profile, mode="llm", return_profile=None, phase="idle") finally: # Nach einem fehlgeschlagenen Skript/Timeout darf der Router # Qwen nicht blind freigeben. Nur ein semantisch verifiziertes # Profil (Alias + Kontext) wird wieder als verfügbar markiert. active = current_profile() up_after = upstream_status() available = bool(active in PROFILES and _profile_is_ready(active, up_after)) _set_qwen_unavailable(not available) finally: STATE.switching = None return upstream_status() # --------------------------------------------------------------------------- # Bildgenerierung und Editing (Qwen-Image-2.1 INT8) # --------------------------------------------------------------------------- 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 self._logf = 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) self._logf = open(IMAGE_WORKER_LOG, "ab") self.proc = subprocess.Popen( [IMAGE_PYTHON, IMAGE_WORKER], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=self._logf, text=True, bufsize=1, start_new_session=True, ) RUNTIME.save(worker="image", worker_pid=self.proc.pid, phase="loading-image") self._reader = threading.Thread(target=self._read_loop, daemon=True) self._reader.start() try: 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() try: self.proc.wait(timeout=5) except subprocess.TimeoutExpired: pass if self._logf is not None: try: self._logf.close() except OSError: pass self._logf = None self.proc = None self.model_loaded = False RUNTIME.clear_worker("image") def _worker() -> _Worker: """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 = (_RemoteWorker(kind="image", url=IMAGE_WORKER_URL, token=IMAGE_WORKER_TOKEN, endpoint="/generate") 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, *, kind: str, url: str, token: str, endpoint: str) -> None: self.kind = kind self.url = url self.token = token self.endpoint = endpoint 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 {self.token}"} if body is not None: headers["Content-Type"] = "application/json" req = urllib.request.Request(self.url + path, data=body, method=method, headers=headers) try: with urllib.request.urlopen(req, timeout=timeout) as response: return json.load(response) except urllib.error.HTTPError as exc: try: message = json.loads(exc.read(4096)).get("message") except Exception: message = None raise RuntimeError(message or f"Bild-Worker HTTP {exc.code}") from exc except (OSError, urllib.error.URLError, TimeoutError) as exc: raise RuntimeError(f"Bild-Worker nicht erreichbar: {exc}") from exc def start(self) -> None: if not self.token or len(self.token) < 32: raise RuntimeError("Bild-Worker-Token fehlt oder ist zu kurz") _profile_controller_request("POST", f"/workers/{self.kind}/start") deadline = time.monotonic() + IMAGE_START_TIMEOUT while time.monotonic() < deadline: try: health = self._request("GET", "/health", timeout=3) self.running = True self.model_loaded = bool(health.get("model_loaded")) return except (OSError, urllib.error.URLError, TimeoutError, RuntimeError): 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", self.endpoint, clean, timeout) def stop(self) -> None: try: _profile_controller_request("POST", f"/workers/{self.kind}/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: 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 (Bildmodell 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) 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, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, timeout=120) if proc.returncode != 0: out = proc.stdout.decode(errors="replace").strip() raise RuntimeError( f"systemctl start {LLAMA_SERVICE} fehlgeschlagen " f"(Exit {proc.returncode}): {out[-500:]}") except subprocess.TimeoutExpired: log.error("systemctl start hat 120 s überschritten") _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) RUNTIME.save(last_profile=profile, phase="idle") def generate_image(prompt: str, width: int, height: int, steps: int, guidance: float, seed: int | None, n: int, quality: str = "standard", source_files: list[str] | None = None, model: str = IMAGE_MODEL_NAME, ) -> 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 img.last_error = None img.current_model = model # 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" RUNTIME.save(last_profile=profile, phase=img.phase) 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). 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) worker_payload = { "cmd": "generate", "prompt": prompt, "width": width, "height": height, "steps": steps, "guidance": guidance, "seed": seed, "output": output, "source_files": source_files or [], } resp = worker.request(worker_payload, 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") # Metadaten speichern (Sidecar-JSON). meta = { "prompt": prompt, "seed": seed, "width": width, "height": height, "size": f"{width}x{height}", "steps": steps, "guidance": guidance, "quality": quality, "mode": ("image-edit" if source_files else "text-to-image"), "reference_images": len(source_files or []), "seconds": resp.get("seconds"), "model": model, "created": time.strftime("%Y-%m-%dT%H:%M:%S"), } meta_path = os.path.join(IMAGE_DIR, filename[:-4] + ".json") try: with open(meta_path, "w", encoding="utf-8") as f: json.dump(meta, f, ensure_ascii=False, indent=2) except OSError as e: log.warning("Metadaten-Speicherung fehlgeschlagen: %s", e) log.info("Bild %d/%d: %s (%.1f s)", i + 1, n, filename, resp.get("seconds", 0)) removed = enforce_artifact_retention( IMAGE_DIR, IMAGE_RETENTION_FILES, IMAGE_RETENTION_BYTES, IMAGE_RETENTION_DAYS, protected=results) if removed: log.info("Bild-Retention: %d alte Bilder entfernt", len(removed)) # 4) Worker vollständig beenden (VRAM + CUDA-Kontext freigeben). img.phase = "unloading-image" worker.stop() 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" img.current_model = None 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)) # --------------------------------------------------------------------------- # Multimodale Chat-Eingaben # --------------------------------------------------------------------------- _CHAT_IMAGE_DATA_TYPES = { "image/jpeg", "image/png", "image/webp", "image/gif", } def _normalize_chat_image(image_url: str) -> str: """Validiert ein Bild und liefert eine begrenzte data-URL. Remote-Downloads sind standardmäßig deaktiviert. Wenn sie ausdrücklich aktiviert werden, lädt der Router das Bild nach SSRF-Prüfung selbst und übergibt llama.cpp ausschließlich eine data-URL. """ if image_url.startswith("data:"): header, separator, payload = image_url.partition(",") match = re.fullmatch( r"data:([a-zA-Z0-9.+-]+/[a-zA-Z0-9.+-]+);base64", header) if (not separator or not match or match.group(1).lower() not in _CHAT_IMAGE_DATA_TYPES): raise ValueError("ungültige oder nicht unterstützte Bild-data-URL") try: raw = base64.b64decode(payload, validate=True) except (ValueError, binascii.Error): raise ValueError("ungültige Base64-Bilddaten") from None if not raw or len(raw) > CHAT_IMAGE_MAX_BYTES: raise ValueError( "Bildgröße außerhalb des Limits " f"(max {CHAT_IMAGE_MAX_BYTES} Bytes)") return image_url parsed = urllib.parse.urlsplit(image_url) if parsed.scheme not in {"http", "https"} or not parsed.hostname: raise ValueError( "Bild muss eine data-URL oder eine gültige HTTP(S)-URL sein") if not CHAT_IMAGE_ALLOW_REMOTE_URLS: raise ValueError( "Remote-Bild-URLs sind deaktiviert; Bild bitte als data-URL hochladen") try: addresses = socket.getaddrinfo( parsed.hostname, parsed.port or (443 if parsed.scheme == "https" else 80)) except socket.gaierror as exc: raise ValueError( f"Bild-Host kann nicht aufgelöst werden: {exc}") from None for address in addresses: try: ip = ipaddress.ip_address(address[4][0]) except ValueError: raise ValueError("Bild-Host liefert eine ungültige Adresse") from None if not ip.is_global: raise ValueError("private/lokale Bild-URLs sind nicht erlaubt") request = urllib.request.Request( image_url, headers={"User-Agent": "AI-Profile-Router/2.0"}) try: with urllib.request.urlopen(request, timeout=15) as response: final_url = urllib.parse.urlsplit(response.geturl()) if final_url.hostname != parsed.hostname: raise ValueError( "Weiterleitungen zu einem anderen Bild-Host sind nicht erlaubt") content_type = response.headers.get_content_type().lower() if content_type not in _CHAT_IMAGE_DATA_TYPES: raise ValueError( "Remote-Inhalt ist kein unterstütztes Bild " f"({content_type})") raw = response.read(CHAT_IMAGE_MAX_BYTES + 1) except urllib.error.URLError as exc: raise ValueError( f"Remote-Bild kann nicht geladen werden: {exc}") from None if not raw or len(raw) > CHAT_IMAGE_MAX_BYTES: raise ValueError( f"Bildgröße außerhalb des Limits (max {CHAT_IMAGE_MAX_BYTES} Bytes)") return (f"data:{content_type};base64," + base64.b64encode(raw).decode("ascii")) def _normalize_chat_images(data: dict) -> dict: """Validiert alle image_url-Parts, ohne sie aus dem Chat zu entfernen.""" out = json.loads(json.dumps(data)) messages = out.get("messages") if not isinstance(messages, list): return out for message in messages: if not isinstance(message, dict): continue content = message.get("content") if not isinstance(content, list): continue for part in content: if not isinstance(part, dict) or part.get("type") != "image_url": continue image = part.get("image_url") if isinstance(image, dict): url = image.get("url") if not isinstance(url, str) or not url: raise ValueError("image_url.url fehlt") image["url"] = _normalize_chat_image(url) elif isinstance(image, str) and image: part["image_url"] = _normalize_chat_image(image) else: raise ValueError("image_url fehlt") return out def _request_has_image(data: dict) -> bool: """True, wenn irgendwo im Request ein image_url-Part vorkommt.""" messages = data.get("messages") if not isinstance(messages, list): return False return any( isinstance(part, dict) and part.get("type") == "image_url" for message in messages if isinstance(message, dict) for part in (message.get("content") if isinstance(message.get("content"), list) else []) ) # --------------------------------------------------------------------------- # llama.cpp Chat-Template-Parameter # --------------------------------------------------------------------------- _REASONING_EFFORT_MAP = { "minimal": "low", "low": "low", "medium": "medium", "high": "xhigh", "xhigh": "xhigh", "max": "xhigh", "ultra": "xhigh", } _REASONING_OFF = {"", "none", "off", "disabled", "false"} _REASONING_BUDGET_TOKENS = { "minimal": 256, "low": 768, "medium": 2048, "high": 4096, "xhigh": 8192, "max": 8192, "ultra": 8192, } def _load_global_system_policy() -> str: """Load the static cross-client platform policy. The file is read for every request so operators can revise the policy without rebuilding or restarting the router. Its contents remain stable between edits and therefore remain friendly to upstream prompt caches. """ if not GLOBAL_SYSTEM_POLICY_FILE: return "" try: with open(GLOBAL_SYSTEM_POLICY_FILE, encoding="utf-8") as handle: return handle.read().strip() except OSError as exc: raise ValueError( f"globale Systemrichtlinie nicht lesbar: {exc}") from exc def _inject_global_system_policy(data: dict, path: str) -> dict: """Prepend the shared policy to OpenAI chat and Responses requests.""" policy = _load_global_system_policy() if not policy: return data if path == "/v1/chat/completions": messages = data.get("messages") if not isinstance(messages, list): return data if messages and isinstance(messages[0], dict) and ( messages[0].get("role") == "system" and isinstance(messages[0].get("content"), str)): existing = messages[0]["content"] if policy not in existing: messages[0]["content"] = f"{policy}\n\n{existing}" elif not any( isinstance(message, dict) and message.get("role") == "system" and message.get("content") == policy for message in messages): messages.insert(0, {"role": "system", "content": policy}) return data if path == "/v1/responses": instructions = data.get("instructions") if isinstance(instructions, str) and instructions: if policy not in instructions: data["instructions"] = f"{policy}\n\n{instructions}" elif instructions is None or instructions == "": data["instructions"] = policy return data def _normalize_llamacpp_reasoning(data: dict) -> dict: """Mappt OpenAI/Hermes-Reasoning auf llama.cpp-Template-Parameter. Hermes sendet ``reasoning_effort`` bei einem Custom Provider als Top-Level-Feld. llama.cpp akzeptiert das Feld zwar, reicht es dort aber nicht an das Jinja-Chat-Template weiter. Qwen3.8 erwartet stattdessen ``chat_template_kwargs.reasoning_effort`` bzw. ``enable_thinking=false``. Zusätzlich erhält llama.cpp mit ``thinking_budget_tokens`` eine echte, pro Request geltende Obergrenze. Dadurch sind die in Hermes sichtbaren Stufen nicht bloß unterschiedlich formulierte Template-Hinweise. Die Funktion verändert den übergebenen Request absichtlich in-place und entfernt das wirkungslose Top-Level-Feld. Andere Template-Argumente des Clients bleiben erhalten. """ explicit_effort = "reasoning_effort" in data if explicit_effort: raw_effort = data.pop("reasoning_effort") else: # A client may already speak llama.cpp's native template dialect. # Preserve that explicit choice; otherwise apply the platform-wide # default. Hermes deliberately omits reasoning_effort when its UI is # set to Off, so the safe default must remain Off; enabled levels are # sent explicitly by Hermes and other capable clients. existing_kwargs = data.get("chat_template_kwargs") if "thinking_budget_tokens" in data or ( isinstance(existing_kwargs, dict) and ( "reasoning_effort" in existing_kwargs or "enable_thinking" in existing_kwargs)): return data raw_effort = DEFAULT_REASONING_EFFORT effort = str(raw_effort).strip().lower() if raw_effort is not None else "" template_kwargs = data.get("chat_template_kwargs") if not isinstance(template_kwargs, dict): template_kwargs = {} data["chat_template_kwargs"] = template_kwargs if effort in _REASONING_OFF: template_kwargs.pop("reasoning_effort", None) template_kwargs["enable_thinking"] = False data["thinking_budget_tokens"] = 0 return data mapped = _REASONING_EFFORT_MAP.get(effort) if mapped is None: # Unbekannte OpenAI-Erweiterungen dürfen das Qwen-Template nicht mit # einem ungültigen Wert zum Abbruch bringen. Das Template verwendet # in diesem Fall seine eigene Voreinstellung. log.warning("Unbekanntes reasoning_effort=%r ignoriert", raw_effort) if not template_kwargs: data.pop("chat_template_kwargs", None) return data template_kwargs["enable_thinking"] = True template_kwargs["reasoning_effort"] = mapped data["thinking_budget_tokens"] = _REASONING_BUDGET_TOKENS[effort] return data def _cap_chat_generation(data: dict) -> dict: """Apply a client-independent upper bound to one chat generation. Some OpenAI-compatible clients omit both token-limit fields. llama.cpp interprets that as ``n_predict=-1`` and can remain in hidden reasoning until the context is exhausted. Smaller explicit limits are preserved. """ cap = MAX_GENERATION_TOKENS if cap <= 0: return data fields = ("max_tokens", "max_completion_tokens") present = False for field in fields: if field not in data: continue present = True value = data[field] if isinstance(value, bool): data[field] = cap continue try: parsed = int(value) except (TypeError, ValueError): data[field] = cap continue data[field] = min(parsed, cap) if parsed > 0 else cap if not present: data["max_tokens"] = cap return data # --------------------------------------------------------------------------- # HTTP-Handler # --------------------------------------------------------------------------- class Handler(BaseHTTPRequestHandler): server_version = "AIProfileRouter/2.0" sys_version = "" timeout = 60 # Socket-Timeout für Client-Requests (s) # ---------- Routing ---------- def do_GET(self): self._route() def do_POST(self): self._route() def do_PUT(self): self._route() def do_PATCH(self): self._route() def do_DELETE(self): self._route() def do_OPTIONS(self): self._route() def _route(self): path = self.path.split("?", 1)[0] started = time.monotonic() slot_acquired = False try: if path == "/health" and self.command == "GET": # Liveness: der Routerprozess lebt. Ein absichtlich entladenes # Qwen (Vision/Bild) darf keinen Restart-Loop auslösen. self._send_json(200, {"status": "ok", "router": "alive"}) return if path == "/ready" and self.command == "GET": up = upstream_status() active = current_profile() with STATE.avail_lock: unavailable = STATE.qwen_unavailable ready = bool(active in PROFILES and not unavailable and _profile_is_ready(active, up)) self._send_json(200 if ready else 503, { "status": "ok" if ready else "degraded", "router": "alive", "upstream": "ready" if ready else "unavailable", }) return slot_acquired = REQUEST_SLOTS.acquire(blocking=False) if not slot_acquired: self._send_error(429, "Router ist ausgelastet; bitte erneut versuchen", "server_error", "too_many_requests") return if not self._authorized(): self._send_auth_required() elif path == "/v1/models" and self.command == "GET": self._send_json(200, self._models_payload()) elif path == "/models" and self.command == "GET": # llama.cpp clients (notably OpenClaw's existing-server # provider) probe the native catalog before /v1/models. Do # not proxy this request to the one currently active profile, # otherwise the remaining switchable profiles disappear from # discovery. self._send_json(200, self._llamacpp_models_payload()) elif path == "/status" and self.command == "GET": self._send_json(200, self._status_payload()) elif path == "/mode" and self.command == "GET": self._send_json(200, self._mode_payload()) elif path == "/mode" and self.command == "POST": self._mode_change() elif path == "/v1/audio/models" and self.command == "GET": self._send_json(200, self._audio_models_payload()) elif path == "/v1/audio/voices" and self.command == "GET": self._send_json(200, self._audio_voices_payload()) elif path == "/v1/images/generations" and self.command == "POST": if ENABLE_IMAGE_GENERATION: self._image_generate() else: self._send_error(503, "Bildgenerierung ist nicht installiert", "server_error", "feature_disabled") elif path == "/v1/images/edits" and self.command == "POST": if ENABLE_IMAGE_GENERATION: self._image_edit() else: self._send_error(503, "Bildbearbeitung ist nicht installiert", "server_error", "feature_disabled") elif path == "/v1/audio/speech" and self.command == "POST": if ENABLE_TTS: self._speech() else: self._send_error(503, "Sprachausgabe ist nicht installiert", "server_error", "feature_disabled") elif (path == "/v1/audio/speech/pcm-stream" and self.command == "POST"): if ENABLE_TTS: self._speech_pcm_stream() else: self._send_error(503, "Sprachausgabe ist nicht installiert", "server_error", "feature_disabled") elif path == "/v1/audio/transcriptions" and self.command == "POST": if ENABLE_STT: self._transcribe() else: self._send_error(503, "Spracherkennung ist nicht installiert", "server_error", "feature_disabled") 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.removeprefix("/") in PROFILES and (self.command == "POST" or (self.command == "GET" and ALLOW_LEGACY_GET_SWITCH))): self._switch(path[1:]) elif (path.removeprefix("/") in PROFILES and self.command == "GET"): self._send_error(405, "Profilwechsel erfordert POST", "invalid_request_error", "method_not_allowed") else: self._forward() except BrokenPipeError: log.warning("Client getrennt: %s %s", self.command, path) except Exception: log.exception("Fehler bei %s %s", self.command, path) self._safe_error(500, "interner Router-Fehler") finally: if slot_acquired: REQUEST_SLOTS.release() log.info("%s %s -> %s in %.3f s", self.command, path, getattr(self, "_last_code", "-"), time.monotonic() - started) def _authorized(self) -> bool: assert AUTH is not None return AUTH.accepts(self.headers.get("Authorization"), self.headers.get("X-API-Key")) def _send_auth_required(self) -> None: body = json.dumps({"error": { "message": "gültiger Router-API-Key erforderlich", "type": "authentication_error", "code": "invalid_api_key", }}).encode() self._last_code = 401 self.send_response(401) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.send_header("WWW-Authenticate", "Bearer") self.send_header("Connection", "close") self.end_headers() self.wfile.write(body) # ---------- Request-Body-Lesen (Content-Length + chunked) ---------- def _read_body(self) -> bytes: """Liest den HTTP-Request-Body (Content-Length oder chunked). Liefert die Body-Bytes. Wirft ValueError bei: - malformed chunked encoding - Upload größer als MAX_UPLOAD_SIZE - unvollständiger Body """ te = self.headers.get("Transfer-Encoding", "").lower() if "chunked" in te: return self._read_chunked_body() length = int(self.headers.get("Content-Length") or 0) if length > MAX_UPLOAD_SIZE: raise ValueError( f"Upload zu groß: {length} bytes (max {MAX_UPLOAD_SIZE})") if length == 0: return b"" data = self.rfile.read(length) if len(data) != length: raise ValueError( f"Unvollständiger Body: {len(data)}/{length} bytes") return data def _read_chunked_body(self) -> bytes: """Liest und dekodiert einen HTTP/1.1 chunked-Transfer-Encoding Body. RFC 7230 §4.1: chunked-body = *chunk last-chunk trailer-part CRLF chunk = chunk-size [chunk-ext] CRLF chunk-data CRLF chunk-size = 1*HEXDIG last-chunk = 0 [chunk-ext] CRLF trailer-part = *( field-line CRLF ) - Chunk-Größen werden hexadezimal geparst. - Chunk Extensions (nach ';') werden toleriert/ignoriert. - 0-Chunk markiert das Ende. - Trailer werden konsumiert und ignoriert. - MAX_UPLOAD_SIZE wird durchgesetzt. """ chunks: list[bytes] = [] total_size = 0 while True: # Chunk-Size-zeile lesen: "hex-size [chunk-ext] CRLF" size_line = self.rfile.readline(65537) if not size_line: raise ValueError("Chunked Body: unerwartetes Ende") # CRLF/LF entfernen size_line = size_line.rstrip(b"\r\n") # Chunk Extension entfernen (alles nach dem ersten ';') if b";" in size_line: size_line = size_line.split(b";", 1)[0] # Hexadezimale Größe parsen size_str = size_line.strip() if not size_str: raise ValueError("Chunked Body: leere Chunk-Size") try: chunk_size = int(size_str, 16) except ValueError: raise ValueError( f"Malformed Chunk-Size: {size_str!r}") # 0-Chunk = Ende des chunked-body if chunk_size == 0: break # Uploadgrößenlimit prüfen total_size += chunk_size if total_size > MAX_UPLOAD_SIZE: raise ValueError( f"Upload zu groß: {total_size} bytes " f"(max {MAX_UPLOAD_SIZE})") # Chunk-Daten lesen chunk_data = self.rfile.read(chunk_size) if len(chunk_data) != chunk_size: raise ValueError( f"Unvollständiges Chunk: {len(chunk_data)}/{chunk_size} bytes") chunks.append(chunk_data) # CRLF nach Chunk-Daten lesen crlf = self.rfile.read(2) if crlf != b"\r\n": raise ValueError( f"Erwartet CRLF nach Chunk, erhalten: {crlf!r}") # Trailer lesen und ignorieren # trailer-part = *( field-line CRLF ), beendet durch leere Zeile while True: line = self.rfile.readline(65537) if not line or line in (b"\r\n", b"\n"): break # Trailer-Header ignorieren return b"".join(chunks) # ---------- Router-eigene Endpunkte ---------- @staticmethod def _models_payload() -> dict: models = [ { "id": EXPECTED_MODELS.get(name) or f"qwen-{name}", "object": "model", "created": 0, "owned_by": "ai-profile-router", "context_length": ctx, "context_window": ctx, } for name, ctx in PROFILES.items() ] if REVIEW_UPSTREAM_URL: models.append({ "id": REVIEW_MODEL_NAME, "object": "model", "created": 0, "owned_by": "ai-profile-router", "context_length": REVIEW_CONTEXT_LENGTH, "context_window": REVIEW_CONTEXT_LENGTH, }) return { "object": "list", "data": models, } @staticmethod def _llamacpp_models_payload() -> dict: """Return every switchable profile in llama.cpp's native catalog.""" active = current_profile() models = [] for name, ctx in PROFILES.items(): model_id = EXPECTED_MODELS.get(name) or f"qwen-{name}" models.append({ "id": model_id, "object": "model", "created": 0, "owned_by": "ai-profile-router", "context_length": ctx, "context_window": ctx, "meta": {"n_ctx_train": ctx}, "status": { "value": "loaded" if name == active else "unloaded", }, "architecture": { "input_modalities": (["text", "image"] if PROFILE_REGISTRY[name]["vision"] else ["text"]), "output_modalities": ["text"], }, }) return {"data": models} 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), "current_profile": current_profile(), "switching": STATE.switching, "profiles": PROFILES, "mode": self._mode_payload(), "upstream": { "url": UPSTREAM_URL, "reachable": up["reachable"], "model": up.get("model"), "ctx": up.get("ctx"), }, "qwen": { "available": (not qwen_unavailable and up["reachable"] and bool(up.get("model"))), "active_chats": active_chats, }, "llama_telemetry": (upstream_telemetry() if up["reachable"] else { "available": False, "slots": [], "metrics": {}, "errors": {"upstream": up.get("error", "not reachable")}, }), "image": { "phase": img.phase, "worker": "running" if (img.worker and img.worker.alive()) else "stopped", "model": img.current_model if img.phase != "idle" else None, "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, }, "tts": tts_status(), "stt": stt_status(), } @staticmethod def _mode_payload() -> dict: state = RUNTIME.load() snapshot = {} if PROFILE_CONTROL_URL: try: snapshot = _profile_controller_request("GET", "/status", timeout=3) except Exception as exc: log.warning("Controller-Status nicht verfügbar: %s", exc) fallback = "unknown" if PROFILE_CONTROL_URL else "unsupported" result = { "active": STATE.mode, "phase": STATE.mode_phase, "return_profile": state.get("return_profile"), "last_error": STATE.mode_error, "enabled": ENABLE_MUSIC_MODE, "worker_errors": snapshot.get("worker_errors", {}), } for worker in ("music", "yue2", "separator", "voice", "voice_change", "applio", "trellis", "video"): for suffix in ("worker", "health"): field = f"{worker}_{suffix}" result[field] = snapshot.get(field, fallback) return result def _mode_change(self) -> None: try: data = json.loads(self._read_body() or b"{}") mode = data.get("mode") if isinstance(data, dict) else None if mode not in {"llm", "music", "yue2", "separation", "voice", "voicechange", "applio", "trellis", "video"}: raise ValueError("Feld 'mode' enthält einen unbekannten Betriebsmodus") started, phase = schedule_operating_mode(mode) self._send_json(202 if started else 200, { "status": "accepted" if started else "ok", "requested_mode": mode, "phase": phase, }) except ValueError as exc: self._send_error(400, str(exc), "invalid_request_error", "invalid_mode") except RuntimeError as exc: self._send_error(503, str(exc), "server_error", "mode_unavailable") # ---------- Bildgenerierung ---------- def _image_generate(self) -> None: data = self._read_image_request() if data is not None: self._image_request(data, []) def _image_edit(self) -> None: """Edit with OpenAI multipart uploads or legacy JSON/base64 input.""" content_type = self.headers.get("Content-Type", "") if "multipart/form-data" in content_type: try: body = self._read_body() files, data = self._parse_multipart_parts(body, content_type) except ValueError as exc: self._send_error(400, str(exc), "invalid_request_error", "invalid_multipart") return images = [payload for name, _filename, payload in files if name in {"image", "image[]"}] if not images: self._send_error(400, "Referenzbild fehlt", "invalid_request_error", "missing_image") return self._image_edit_bytes(data, images) return data = self._read_image_request() if data is None: return encoded: list[str] = [] primary = data.pop("image_b64", None) if isinstance(primary, str) and primary: encoded.append(primary) references = data.pop("reference_images_b64", []) if references is None: references = [] if not isinstance(references, list) or any( not isinstance(item, str) for item in references): self._send_error(400, "'reference_images_b64' muss eine Liste sein", "invalid_request_error", "invalid_references") return encoded.extend(references) if not encoded: self._send_error(400, "Referenzbild fehlt", "invalid_request_error", "missing_image") return if len(encoded) > 4: self._send_error(400, "höchstens vier Referenzbilder erlaubt", "invalid_request_error", "too_many_images") return decoded: list[bytes] = [] try: for item in encoded: if item.startswith("data:"): header, separator, item = item.partition(",") if not separator or not header.lower().startswith("data:image/"): raise ValueError("ungültige Bild-Data-URI") try: raw = base64.b64decode(item, validate=True) except Exception as exc: raise ValueError("ungültige Base64-Bilddaten") from exc decoded.append(raw) self._image_edit_bytes(data, decoded) except ValueError as exc: self._send_error(400, str(exc), "invalid_request_error", "invalid_image") def _image_edit_bytes(self, data: dict, images: list[bytes]) -> None: if len(images) > 4: self._send_error(400, "höchstens vier Referenzbilder erlaubt", "invalid_request_error", "too_many_images") return if not images: self._send_error(400, "Referenzbild fehlt", "invalid_request_error", "missing_image") return source_files: list[str] = [] try: for raw in images: if not raw or len(raw) > CHAT_IMAGE_MAX_BYTES: raise ValueError( f"Referenzbild muss 1..{CHAT_IMAGE_MAX_BYTES} Bytes groß sein") name = f".edit-{os.urandom(12).hex()}.ref" os.makedirs(IMAGE_DIR, exist_ok=True) with open(os.path.join(IMAGE_DIR, name), "xb") as output: output.write(raw) source_files.append(name) self._image_request(data, source_files) except ValueError as exc: self._send_error(400, str(exc), "invalid_request_error", "invalid_image") finally: for name in source_files: try: os.unlink(os.path.join(IMAGE_DIR, name)) except FileNotFoundError: pass except OSError as exc: log.warning("temporäres Referenzbild nicht gelöscht: %s", exc) def _read_image_request(self) -> dict | None: try: body = self._read_body() except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_body") return None try: data = json.loads(body) except ValueError: self._send_error(400, "ungültiges JSON", "invalid_request_error", "invalid_json") return None if not isinstance(data, dict): self._send_error(400, "Request muss ein JSON-Objekt sein", "invalid_request_error", "invalid_request") return None return data def _image_request(self, data: dict, source_files: list[str]) -> None: 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 model = data.get("model", IMAGE_MODEL_NAME) if model != IMAGE_MODEL_NAME: self._send_error( 400, f"unbekanntes Bildmodell: {model!r}", "invalid_request_error", "invalid_model") 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 steps != IMAGE_INFERENCE_STEPS): self._send_error( 400, f"{IMAGE_MODEL_NAME} erfordert " f"'steps'={IMAGE_INFERENCE_STEPS}", "invalid_request_error", "invalid_steps") return 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 guidance != 1.0: self._send_error(400, f"{IMAGE_MODEL_NAME} erfordert 'guidance'=1.0", "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_explicit = "response_format" in data 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, quality, source_files, model) 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} # OpenClaw's OpenAI-compatible image parser requires b64_json but # does not send response_format. Keep the URL for existing clients # and add inline image data only when the format was omitted. if response_format == "b64_json" or not response_format_explicit: with open(path, "rb") as f: item["b64_json"] = base64.b64encode(f.read()).decode() if response_format == "b64_json": 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 entry = { "name": name, "url": f"/images/{name}", "bytes": st.st_size, "modified": int(st.st_mtime), } # Metadaten laden (Sidecar-JSON, falls vorhanden). meta_path = os.path.join(IMAGE_DIR, name[:-4] + ".json") if os.path.isfile(meta_path): try: with open(meta_path, encoding="utf-8") as f: entry["meta"] = json.load(f) except (OSError, ValueError): pass entries.append(entry) 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) # ---------- Sprachausgabe (Qwen3-TTS) ---------- def _speech(self) -> None: try: body = self._read_body() except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_body") return try: data = json.loads(body) 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 # input (OpenAI) – auch 'text' akzeptieren (bequemer für curl) text = data.get("input", data.get("text")) if not isinstance(text, str) or not text.strip(): self._send_error(400, "'input' fehlt oder ist leer", "invalid_request_error", "missing_input") return if len(text) > 8000: self._send_error(400, "'input' zu lang (max 8000 Zeichen)", "invalid_request_error", "input_too_long") return voice = data.get("voice", TTS_DEFAULT_VOICE) if voice not in TTS_VOICES: self._send_error( 400, f"ungültige Stimme: {voice!r} " f"(erlaubt: {', '.join(TTS_VOICES)})", "invalid_request_error", "invalid_voice") return fmt = data.get("response_format", TTS_DEFAULT_FORMAT) if fmt not in TTS_FORMATS: self._send_error( 400, f"ungültiges response_format: {fmt!r} " f"(erlaubt: {', '.join(TTS_FORMATS)})", "invalid_request_error", "invalid_format") return speed = data.get("speed", 1.0) try: speed = float(speed) except (TypeError, ValueError): self._send_error(400, "'speed' muss eine Zahl sein", "invalid_request_error", "invalid_speed") return if not 0.5 <= speed <= 2.0: self._send_error(400, "'speed' muss zwischen 0.5 und 2.0 sein", "invalid_request_error", "invalid_speed") return # Modell-Name optional; falls angegeben, muss es Qwen3-TTS sein. model = data.get("model") if model is not None and model != TTS_MODEL: self._send_error(400, f"unbekanntes Modell: {model!r} " f"(erwartet: {TTS_MODEL})", "invalid_request_error", "unknown_model") return self.timeout = None # Synthese kann dauern try: audio, content_type = tts_synthesize( text.strip(), voice, speed, fmt) except RuntimeError as e: self._send_error(503, str(e), "server_error", "tts_failed") return self._last_code = 200 self.send_response(200) self.send_header("Content-Type", content_type) self.send_header("Content-Length", str(len(audio))) self.send_header("Connection", "close") self.end_headers() self.wfile.write(audio) def _speech_pcm_stream(self) -> None: """Pass through Qwen's native 24 kHz PCM stream without buffering.""" try: body = self._read_body() data = json.loads(body) except ValueError as exc: self._send_error(400, str(exc) or "ungültiges JSON", "invalid_request_error", "invalid_body") return if not isinstance(data, dict): self._send_error(400, "Request muss ein JSON-Objekt sein", "invalid_request_error", "invalid_request") return text = data.get("input", data.get("text")) if not isinstance(text, str) or not text.strip() or len(text) > 8000: self._send_error(400, "'input' fehlt, ist leer oder zu lang", "invalid_request_error", "invalid_input") return voice = data.get("voice", TTS_DEFAULT_VOICE) if voice not in TTS_VOICES: self._send_error(400, f"ungültige Stimme: {voice!r}", "invalid_request_error", "invalid_voice") return hostport = TTS_WORKER_URL.split("://", 1)[-1] host, _, port = hostport.partition(":") connection = None headers_sent = False try: connection = http.client.HTTPConnection( host, int(port) if port else 80, timeout=TTS_CONNECT_TIMEOUT) connection.request( "POST", "/tts/pcm-stream", body=body, headers={"Content-Type": "application/json", "Accept": "application/octet-stream"}) connection.sock.settimeout(TTS_TIMEOUT) response = connection.getresponse() if response.status != 200: message = response.read(512).decode(errors="replace") self._send_error(503, f"TTS-Stream fehlgeschlagen: {message}", "server_error", "tts_failed") return self._last_code = 200 self.send_response(200) self.send_header("Content-Type", "application/octet-stream") self.send_header("Cache-Control", "no-store") self.send_header("Connection", "close") self.end_headers() headers_sent = True while True: chunk = response.read1(16384) if not chunk: break self.wfile.write(chunk) self.wfile.flush() except (OSError, http.client.HTTPException) as exc: if not headers_sent: try: self._send_error(502, f"TTS-Worker nicht erreichbar: {exc}", "server_error", "tts_unavailable") except (OSError, BrokenPipeError): pass finally: if connection is not None: connection.close() self.close_connection = True # ---------- Audio-Discovery ---------- def _audio_models_payload(self) -> dict: """Listet verfügbare Audio-Modelle (STT + TTS).""" tts = tts_status() stt = stt_status() models = [] if stt.get("ready"): models.append({ "id": STT_MODEL, "object": "model", "owned_by": "whisper.cpp", "type": "transcription", }) if tts.get("ready"): models.append({ "id": TTS_MODEL, "object": "model", "owned_by": "qwen", "type": "speech", }) return {"object": "list", "data": models} def _audio_voices_payload(self) -> dict: """Listet verfügbare TTS-Stimmen.""" tts = tts_status() voices = [] for v in tts.get("voices", []): voices.append({ "id": v, "object": "voice", "language": "de", }) return {"object": "list", "data": voices} # ---------- STT (Spracherkennung) ---------- def _parse_multipart_parts(self, data: bytes, content_type: str ) -> tuple[list[tuple[str, str, bytes]], dict]: """Parst alle Datei- und Textteile eines multipart/form-data-Body. Nutzt email.parser.BytesParser (Standardbibliothek) für robustes MIME-Parsing. Handhabt quoted und unquoted Boundaries, beliebige Feldreihenfolge, zusätzliche Header und binäre Payloads. """ # MIME-Message aus rohen Bytes + Content-Type-Header bauen raw = (f"Content-Type: {content_type}\r\n\r\n" ).encode("utf-8") + data msg = BytesParser(policy=compat32).parsebytes(raw) if not msg.is_multipart(): raise ValueError("Kein multipart/form-data") files: list[tuple[str, str, bytes]] = [] fields: dict[str, str] = {} for part in msg.get_payload(): disposition = part.get("Content-Disposition", "") name = None part_filename = None for kv in disposition.split(";"): kv = kv.strip() if kv.startswith("name="): name = kv[len("name="):].strip('"') elif kv.startswith("filename="): part_filename = kv[len("filename="):].strip('"') if name is None: continue payload = part.get_payload(decode=True) if payload is None: payload = b"" if part_filename is not None: # Dateifeld (binär, nicht dekodieren) files.append((name, part_filename or "", payload)) else: # Textfeld fields[name] = payload.decode("utf-8", errors="replace") return files, fields def _parse_multipart(self, data: bytes, content_type: str ) -> tuple[bytes, str, dict]: """Kompatibler Einzeldatei-Wrapper für den STT-Pfad.""" files, fields = self._parse_multipart_parts(data, content_type) if not files: return b"", "", fields _name, filename, file_data = files[-1] return file_data, filename, fields def _transcribe(self) -> None: """POST /v1/audio/transcriptions – STT (OpenAI-kompatibel).""" content_type = self.headers.get("Content-Type", "") if "multipart/form-data" not in content_type: self._send_error(400, "Content-Type muss multipart/form-data sein", "invalid_request_error", "invalid_content_type") return try: data = self._read_body() except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_body") return try: file_data, filename, fields = self._parse_multipart( data, content_type) except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_multipart") return if not file_data: self._send_error(400, "Keine Datei im Request", "invalid_request_error", "missing_file") return # Modell-Validierung model = fields.get("model", STT_MODEL) if model not in (STT_MODEL, "whisper"): self._send_error(400, f"unbekanntes Modell: {model!r} " f"(erwartet: {STT_MODEL})", "invalid_request_error", "unknown_model") return # Optionale Felder language = fields.get("language") prompt = fields.get("prompt") temperature = None if fields.get("temperature"): try: temperature = float(fields["temperature"]) except ValueError: self._send_error(400, "'temperature' muss eine Zahl sein", "invalid_request_error", "invalid_temperature") return response_format = fields.get("response_format", "json") self.timeout = None # Transkription kann dauern try: result = stt_transcribe( file_data, filename, language=language, prompt=prompt, temperature=temperature) except RuntimeError as e: self._send_error(503, str(e), "server_error", "stt_failed") return # OpenAI-kompatibles Antwort-Format if response_format == "verbose_json": resp = { "text": result.get("text", ""), "language": result.get("language", "de"), "duration": result.get("audio_duration_ms", 0) / 1000.0, } else: resp = {"text": result.get("text", "")} self._send_json(200, resp) def _switch(self, profile: str) -> None: if profile not in PROFILES: self._send_error(400, f"unbekanntes Profil: {profile}", "invalid_request_error", "invalid_profile") return try: switch_profile(profile) except (ValueError, RuntimeError) as e: self._send_error(503, str(e), "server_error", "profile_switch_failed") return up = upstream_status() self._send_json(200, { "status": "ok", "profile": profile, "context_length": PROFILES[profile], "model": up.get("model"), }) # ---------- Transparentes Forwarding ---------- def _forward(self) -> None: path = self.path.split("?", 1)[0] try: body = self._read_body() or None except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_body") return # /v1/streams/lookup (Open-WebUI-Stream-Recovery): darf NIEMALS auf # die Qwen-Wiederherstellung warten (während Profilwechsel oder # Image-Job ist Qwen down). Wenn Qwen down ist, # gibt es per Definition keine aktiven Streams → sofortige lokale # Antwort []. Ansonsten normal an llama.cpp weiterleiten. if path == "/v1/streams/lookup": with STATE.avail_lock: qwen_unavailable = STATE.qwen_unavailable if qwen_unavailable: self._send_json(200, []) return self._proxy(body) return data = None requested_profile: str | None = None requested_review = False # Virtuelles Modell erkennen. Umschalten und Chat-Lease werden weiter # unten atomar unter dem zentralen Orchestrierungs-Lock ausgeführt. if body is not None and self.path.startswith("/v1/"): try: data = json.loads(body) except ValueError: data = None if (isinstance(data, dict) and path in {"/v1/chat/completions", "/v1/responses"}): command = _control_command(data, path) if command is not None: self._control_response(command, data, path) return model = data.get("model") if isinstance(data, dict) else None if (isinstance(model, str) and REVIEW_UPSTREAM_URL and model == REVIEW_MODEL_NAME): requested_review = True elif isinstance(model, str) and model in VIRTUAL_MODELS: requested_profile = VIRTUAL_MODELS[model] elif isinstance(model, str) and model.startswith("qwen-"): # qwen-* ist der Namensraum des Routers self._send_error(400, f"unbekanntes virtuelles Modell: {model}", "invalid_request_error", "unknown_model") return if path in {"/v1/chat/completions", "/v1/responses"}: try: data = _inject_global_system_policy(data, path) except ValueError as exc: self._send_error(500, str(exc), "server_error", "system_policy_unavailable") return body = json.dumps(data).encode() # Alle modellbezogenen Requests erhalten eine atomare Lease. Damit # kann kein zweiter Client zwischen Profilwahl und Upstream-Request das # Modell austauschen. Vision-Vorbereitung gehört zur selben Transaktion. if isinstance(data, dict) and path == "/v1/chat/completions": if requested_review: self._review_chat_proxy(data) return self._chat_proxy(body, data, requested_profile) return if requested_profile is not None: self._profiled_proxy(body, data, requested_profile) return # An llama.cpp weiterleiten (mit Chat-Waiting, Streaming bleibt erhalten). self._proxy_with_wait(body) def _control_response(self, command: str, data: dict, path: str) -> None: """Return OpenAI-compatible local replies for Athena control commands.""" if command == "/athena status": mode = self._mode_payload() profile = current_profile() text = (f"Athena läuft im {mode['active'].upper()}-Modus. " f"Phase: {mode['phase']}. Musik-Worker: " f"{mode['music_worker']}. YuE2: " f"{mode['yue2_worker']}. Stem-Separator: " f"{mode['separator_worker']}. Voice Studio: " f"{mode['voice_worker']}. Voice Changer: " f"{mode['voice_change_worker']}. 3D Studio: " f"{mode['trellis_worker']}. LTX-2 Studio: " f"{mode['video_worker']}. LLM-Profil: {profile or 'entladen'}.") else: target = ("music" if command == "/athena music" else "yue2" if command == "/athena yue2" else "separation" if command in {"/athena stems", "/athena separation"} else "voice" if command == "/athena voice" else "voicechange" if command in {"/athena voicechange", "/athena changer"} else "applio" if command == "/athena applio" else "trellis" if command in {"/athena 3d", "/athena trellis"} else "video" if command in {"/athena video", "/athena ltx", "/athena ltx2"} else "llm") try: started, phase = schedule_operating_mode(target) if started: text = ("Musikstudio wird gestartet. LLM und TTS werden entladen." if target == "music" else "YuE2 Studio wird gestartet. LLM und TTS werden entladen." if target == "yue2" else "Stimmtrennung wird gestartet. LLM und TTS werden entladen." if target == "separation" else "Voice Studio wird gestartet. LLM und TTS werden entladen." if target == "voice" else "Voice Changer wird gestartet. LLM und TTS werden entladen." if target == "voicechange" else "Applio wird gestartet. LLM und TTS werden entladen." if target == "applio" else "3D Studio wird gestartet. LLM und TTS werden entladen." if target == "trellis" else "Spezialmodus wird beendet und das vorherige LLM-Profil wiederhergestellt.") else: text = (f"Athena ist bereits im {target.upper()}-Modus " f"oder wechselt gerade ({phase}).") except RuntimeError as exc: self._send_error(503, str(exc), "server_error", "mode_unavailable") return model = str(data.get("model") or "athena-control") created = int(time.time()) request_id = f"athena-mode-{uuid.uuid4().hex[:16]}" if path == "/v1/responses": self._send_json(200, { "id": request_id, "object": "response", "created_at": created, "status": "completed", "model": model, "output": [{"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": text}]}], "output_text": text, "usage": {"input_tokens": 0, "output_tokens": 0, "total_tokens": 0}, }) return if data.get("stream") is True: chunks = [ {"id": request_id, "object": "chat.completion.chunk", "created": created, "model": model, "choices": [{"index": 0, "delta": {"role": "assistant", "content": text}, "finish_reason": None}]}, {"id": request_id, "object": "chat.completion.chunk", "created": created, "model": model, "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}]}, ] body = "".join(f"data: {json.dumps(chunk)}\n\n" for chunk in chunks) body += "data: [DONE]\n\n" encoded = body.encode() self._last_code = 200 self.send_response(200) self.send_header("Content-Type", "text/event-stream") self.send_header("Content-Length", str(len(encoded))) self.send_header("Connection", "close") self.end_headers() self.wfile.write(encoded) return self._send_json(200, { "id": request_id, "object": "chat.completion", "created": created, "model": model, "choices": [{"index": 0, "message": {"role": "assistant", "content": text}, "finish_reason": "stop"}], "usage": {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0}, }) def _acquire_model_lease(self, profile: str | None = None) -> dict: """Atomar Profil sicherstellen und einen aktiven Request registrieren.""" with _model_lock(): up = (switch_profile(profile, implicit=True) if profile is not None else upstream_status()) if not up["reachable"] or not up.get("model"): raise RuntimeError("llama.cpp nicht erreichbar") with STATE.avail_lock: if STATE.qwen_unavailable: raise RuntimeError("Qwen wird gerade neu geladen") STATE.active_chats += 1 return up @staticmethod def _release_model_lease() -> None: with STATE.avail_lock: STATE.active_chats = max(0, STATE.active_chats - 1) def _profiled_proxy(self, body: bytes | None, data: dict, profile: str) -> None: try: up = self._acquire_model_lease(profile) except ModelWaitTimeout as e: self._send_error(503, str(e), "server_error", "model_wait_timeout") return except (ValueError, RuntimeError) as e: self._send_error(502, str(e), "server_error", "upstream_unavailable") return try: data["model"] = up["model"] self._proxy(json.dumps(data).encode()) finally: self._release_model_lease() def _chat_proxy(self, body: bytes | None, data: dict, profile: str | None) -> None: """Bildvalidierung, Profilwahl und Chat-Lease als eine Transaktion.""" lease_acquired = False try: # Validation and image decoding need no GPU ownership. data = _normalize_llamacpp_reasoning(data) data = _cap_chat_generation(data) has_image = _request_has_image(data) if has_image: data = _normalize_chat_images(data) with _model_lock(): selected = profile or (current_profile() if has_image else None) if has_image and selected in PROFILE_REGISTRY and not PROFILE_REGISTRY[selected]["vision"]: raise ValueError(f"Profil {selected} unterstützt keine Bilder") up = (switch_profile(profile, implicit=True) if profile is not None else upstream_status()) if has_image: log.info("Vision: Bild wird direkt an das aktive " "multimodale Qwen-Profil weitergeleitet") if not up["reachable"] or not up.get("model"): raise RuntimeError("llama.cpp nicht erreichbar") if profile is not None: data["model"] = up["model"] body = json.dumps(data).encode() with STATE.avail_lock: if STATE.qwen_unavailable: raise RuntimeError("Qwen wird gerade neu geladen") STATE.active_chats += 1 lease_acquired = True self._proxy(body) except ModelWaitTimeout as e: self._send_error(503, str(e), "server_error", "model_wait_timeout") except ValueError as e: self._send_error(400, str(e), "invalid_request_error", "invalid_chat_request") except RuntimeError as e: self._send_error(502, str(e), "server_error", "upstream_unavailable") finally: if lease_acquired: self._release_model_lease() def _review_chat_proxy(self, data: dict) -> None: """Leitet kompakte Hermes-Hintergrundreviews an das Hilfsmodell. Absichtlich ohne Hauptmodell-Lease und Profilwechsel: Der Review darf den aktiven Qwen-Chat weder anhalten noch dessen Prompt-Cache ersetzen. """ try: data = _normalize_llamacpp_reasoning(data) data = _cap_chat_generation(data) if _request_has_image(data): raise ValueError("qwen-review unterstützt keine Bilder") data["model"] = REVIEW_MODEL_NAME self._proxy_to( json.dumps(data).encode(), REVIEW_UPSTREAM_HOST, REVIEW_UPSTREAM_PORT, "Review-Modell", ) except ValueError as exc: self._send_error(400, str(exc), "invalid_request_error", "invalid_review_request") 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: self._proxy_to(body, UPSTREAM_HOST, UPSTREAM_PORT, "llama.cpp") def _proxy_to(self, body: bytes | None, host: str, port: int, upstream_name: str) -> None: # An llama.cpp weiterleiten (Streaming bleibt erhalten). try: conn = http.client.HTTPConnection(host, port, timeout=CONNECT_TIMEOUT) conn.connect() conn.sock.settimeout(REQUEST_TIMEOUT) headers = {k: v for k, v in self.headers.items() if k.lower() not in HOP_BY_HOP} conn.request(self.command, self.path, body=body, headers=headers) resp = conn.getresponse() except (OSError, http.client.HTTPException) as e: self._send_error(502, f"{upstream_name} nicht erreichbar: {e}", "server_error", "upstream_unavailable") return self._last_code = resp.status self.send_response(resp.status) for k, v in resp.getheaders(): if k.lower() not in HOP_BY_HOP: self.send_header(k, v) self.send_header("Connection", "close") self.end_headers() try: while True: # ``read(n)`` may wait until the complete buffer is filled. # That defeats SSE: a client sees no token for a long time and # may time out while llama.cpp is already generating. read1() # returns the next currently available wire chunk instead. chunk = resp.read1(16384) if not chunk: break # A closed downstream socket may still accept one small write # into the kernel buffer. Detect the peer FIN before writing so # slow upstream streams are cancelled promptly and do not keep # an inference lease occupied until generation completes. try: readable, _, _ = select.select([self.connection], [], [], 0) if readable: flags = socket.MSG_PEEK | getattr(socket, "MSG_DONTWAIT", 0) if self.connection.recv(1, flags) == b"": log.info("Client trennte Streaming-Verbindung; Upstream wird abgebrochen") break except (BlockingIOError, InterruptedError): pass except OSError: break self.wfile.write(chunk) self.wfile.flush() except (OSError, http.client.HTTPException) as e: log.warning("Upstream-Stream abgebrochen: %s", e) finally: conn.close() # ---------- Antworten ---------- def _send_json(self, code: int, payload: dict) -> None: body = json.dumps(payload).encode() self._last_code = code self.send_response(code) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.send_header("Connection", "close") self.end_headers() self.wfile.write(body) def _send_error(self, code: int, message: str, etype: str, ecode: str) -> None: # OpenAI-kompatibles Fehlerformat self._send_json(code, {"error": {"message": message, "type": etype, "code": ecode}}) def _safe_error(self, code: int, message: str) -> None: try: self._send_error(code, message, "server_error", "internal_error") except Exception: pass # --------------------------------------------------------------------------- # Main # --------------------------------------------------------------------------- class _FlushHandler(logging.StreamHandler): """StreamHandler, der nach jedem Record flusht (journald).""" def emit(self, record): super().emit(record) self.flush() class RouterHTTPServer(ThreadingHTTPServer): allow_reuse_address = True def _startup_reconcile() -> None: """Reconcile persisted worker/model state before accepting requests.""" previous = RUNTIME.load() worker = previous.get("worker") if worker == "image": markers = [IMAGE_WORKER] terminated = terminate_recorded_worker( previous, markers, lambda msg: log.warning("Recovery: %s", msg)) if terminated: time.sleep(1) RUNTIME.clear_worker(worker) removed = enforce_artifact_retention( IMAGE_DIR, IMAGE_RETENTION_FILES, IMAGE_RETENTION_BYTES, IMAGE_RETENTION_DAYS) if removed: log.info("Startup-Retention: %d alte Bilder entfernt", len(removed)) special_mode = previous.get("mode") if ENABLE_MUSIC_MODE and special_mode in {"music", "yue2", "separation", "voice", "voicechange", "applio", "trellis", "video"}: STATE.mode = special_mode STATE.mode_phase = f"starting-{special_mode}" _set_qwen_unavailable(True) try: path, _worker_state, wait_ready = _special_worker(special_mode) _profile_controller_request( "POST", path, timeout=120) wait_ready() STATE.mode_phase = "ready" RUNTIME.save(mode=special_mode, phase=special_mode) log.info("Recovery: Spezialmodus %s wiederhergestellt", special_mode) except Exception as exc: STATE.mode_error = str(exc) STATE.mode_phase = "error" log.error("Recovery: Spezialmodus %s konnte nicht gestartet werden: %s", special_mode, exc) return profile = current_profile() if profile is None: saved = previous.get("last_profile") if saved in PROFILES: try: switch_profile(saved) log.info("Recovery: gespeichertes Profil %s neu angewendet", saved) return except Exception as exc: log.error("Recovery: gespeichertes Profil %s konnte nicht " "angewendet werden: %s", saved, exc) profile = None if profile is None: log.error("Recovery: kein gültiges Profil gefunden; Router startet degraded") _set_qwen_unavailable(True) return up = upstream_status() if _profile_is_ready(profile, up): _set_qwen_unavailable(False) STATE.mode = "llm" STATE.mode_phase = "ready" RUNTIME.save(last_profile=profile, mode="llm", phase="idle") log.info("Recovery: Profil %s ist bereits bereit", profile) return _set_qwen_unavailable(True) try: _restore_qwen(profile) except Exception as exc: log.error("Recovery: Profil %s konnte nicht gestartet werden: %s", profile, exc) return _set_qwen_unavailable(False) log.info("Recovery: Profil %s wurde wiederhergestellt", profile) def main() -> None: global AUTH handler = _FlushHandler(sys.stdout) handler.setFormatter(logging.Formatter( "%(asctime)s %(levelname)s %(message)s")) logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"), handlers=[handler]) try: AUTH = AuthPolicy.from_environment() if not AUTH.enabled and HOST not in {"127.0.0.1", "::1", "localhost"}: raise ConfigurationError( "ROUTER_AUTH_MODE=off ist nur an einer Loopback-Adresse erlaubt") except ConfigurationError as exc: log.critical("Unsichere Router-Konfiguration: %s", exc) raise SystemExit(2) log.info("AI Profile Router startet: %s:%s -> %s (Profile: %s)", HOST, PORT, UPSTREAM_URL, ", ".join(PROFILES)) log.info("Authentifizierung: %s", "aktiv" if AUTH.enabled else "deaktiviert") _startup_reconcile() server = RouterHTTPServer((HOST, PORT), Handler) server.daemon_threads = True try: server.serve_forever() except KeyboardInterrupt: pass finally: server.server_close() if __name__ == "__main__": main()