From 6071b1a2cd3f8c0f84d447e2a35856df1c5fe0cd Mon Sep 17 00:00:00 2001 From: Mikei386 <44135113+Mikei386@users.noreply.github.com> Date: Sun, 20 Sep 2026 19:10:59 +0200 Subject: [PATCH] Reduce controller polling and bound chat admission waits --- dev/check.sh | 12 ++ dev/test_profile_controller.py | 33 +++- dev/test_recovery.py | 9 +- dev/test_router_coordination.py | 102 ++++++++++++ manage.sh | 4 +- .../profile-controller/profile_controller.py | 149 ++++++++---------- platform/scripts/sync-profile-matrix.py | 1 + router/ai_profile_router.py | 104 +++++++----- router/router_profiles.json | 15 +- router/router_support.py | 8 +- 10 files changed, 303 insertions(+), 134 deletions(-) create mode 100755 dev/check.sh create mode 100644 dev/test_router_coordination.py diff --git a/dev/check.sh b/dev/check.sh new file mode 100755 index 0000000..35597ef --- /dev/null +++ b/dev/check.sh @@ -0,0 +1,12 @@ +#!/usr/bin/env bash +# Fast, GPU-free release checks. Integration mocks are opt-in. +set -Eeuo pipefail +cd "$(dirname "${BASH_SOURCE[0]}")/.." +python3 platform/scripts/sync-profile-matrix.py --check +python3 -m unittest discover -s dev -p 'test_*.py' +if [[ ${1:-} == --integration ]]; then + bash dev/test_local.sh +elif [[ $# -gt 0 ]]; then + echo 'Usage: dev/check.sh [--integration]' >&2 + exit 2 +fi diff --git a/dev/test_profile_controller.py b/dev/test_profile_controller.py index c84106d..743724a 100644 --- a/dev/test_profile_controller.py +++ b/dev/test_profile_controller.py @@ -1,4 +1,5 @@ import importlib.util +import json import os import pathlib import unittest @@ -42,6 +43,34 @@ def music_item(state="exited"): class ProfileControllerTests(unittest.TestCase): + def test_status_keeps_llm_when_optional_container_is_missing(self): + payload = json.dumps([item("medium", "running")]).encode() + with patch.object(controller, "MUSIC_WORKER", "missing-music"), \ + patch.object(controller, "docker_request", return_value=(200, payload)) as request: + result = controller.status_snapshot() + self.assertEqual(result["active_profile"], "medium") + self.assertEqual(result["music_worker"], "missing") + self.assertIn("music", result["worker_errors"]) + request.assert_called_once_with("GET", "/containers/json?all=1") + + def test_empty_snapshot_does_not_repeat_docker_query(self): + with patch.object(controller, "docker_request", return_value=(200, b"[]")) as request: + result = controller.status_snapshot() + self.assertIsNone(result["active_profile"]) + self.assertEqual(request.call_count, 1) + + def test_video_ui_failure_does_not_hide_llm(self): + video = {"Id": "video", "State": "running", + "Labels": {controller.VIDEO_LABEL_KEY: "ltx2"}} + payload = json.dumps([item("medium", "running"), video]).encode() + with patch.object(controller, "VIDEO_WORKER", "ltx2"), \ + patch.object(controller, "docker_request", return_value=(200, payload)), \ + patch.object(controller, "cached_video_ui_state", side_effect=RuntimeError("exec failed")): + result = controller.status_snapshot() + self.assertEqual(result["active_profile"], "medium") + self.assertEqual(result["video_worker"], "running") + self.assertIn("video_ui", result["worker_errors"]) + def test_music_start_exclusively_stops_gpu_workers(self): profiles = {name: item(name) for name in controller.ALLOWED} profiles["ultra"] = item("ultra", "running") @@ -99,9 +128,11 @@ class ProfileControllerTests(unittest.TestCase): profiles = {name: item(name) for name in controller.ALLOWED[:-1]} with patch.object(controller, "containers", return_value=profiles), \ patch.object(controller, "image_containers", return_value=[image_item()]), \ - patch.object(controller, "tts_container", return_value=tts_item()): + patch.object(controller, "tts_container", return_value=tts_item()), \ + patch.object(controller, "docker_request") as request: with self.assertRaisesRegex(RuntimeError, "missing"): controller.activate("fast") + request.assert_not_called() def test_image_start_stops_inference_first(self): profiles = {name: item(name) for name in controller.ALLOWED} diff --git a/dev/test_recovery.py b/dev/test_recovery.py index b377496..400eef1 100644 --- a/dev/test_recovery.py +++ b/dev/test_recovery.py @@ -1,4 +1,5 @@ import io +import re import subprocess import tarfile import tempfile @@ -94,7 +95,13 @@ class RecoveryScriptTests(unittest.TestCase): gateway = (ROOT / "platform/docker/wireguard-gateway/entrypoint.sh").read_text( encoding="utf-8" ) - self.assertNotIn('network_mode: "service:wireguard-gateway"', compose) + # WebRTC deliberately shares the gateway for private ICE candidates. + # Dashboard and Portainer must retain their independent namespaces. + for service in ("llama-dashboard", "portainer"): + block = re.search( + rf"(?ms)^ {service}:\n(.*?)(?=^ [a-zA-Z0-9_-]+:|\Z)", compose) + self.assertIsNotNone(block) + self.assertNotIn('network_mode: "service:wireguard-gateway"', block.group(1)) self.assertNotIn('"8099:8099"', compose) self.assertNotIn('"9443:9443"', compose) self.assertIn('start_proxy 8099 llama-dashboard:8099', gateway) diff --git a/dev/test_router_coordination.py b/dev/test_router_coordination.py new file mode 100644 index 0000000..a8b2f30 --- /dev/null +++ b/dev/test_router_coordination.py @@ -0,0 +1,102 @@ +"""Regression coverage for bounded admission and inexpensive status reads.""" +import json +import os +import sys +import threading +import unittest +from pathlib import Path +from unittest.mock import patch + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'router')) +os.environ.setdefault('ROUTER_PROFILES_FILE', '') +import ai_profile_router as router + + +class CoordinationTests(unittest.TestCase): + def test_mode_uses_one_consistent_controller_snapshot(self): + with patch.object(router, 'PROFILE_CONTROL_URL', 'http://controller'), \ + patch.object(router, '_profile_controller_request', return_value={ + 'music_worker': 'running', 'music_health': 'healthy'}) as request, \ + patch.object(router.RUNTIME, 'load', return_value={}): + result = router.Handler._mode_payload() + self.assertEqual(result['music_worker'], 'running') + self.assertEqual(result['music_health'], 'healthy') + request.assert_called_once_with('GET', '/status', timeout=3) + + def test_active_profile_does_not_request_optional_workers(self): + with patch.object(router, 'PROFILE_CONTROL_URL', 'http://controller'), \ + patch.object(router, '_profile_controller_request', + return_value={'active_profile': 'medium'}) as request: + self.assertEqual(router.current_profile(), 'medium') + request.assert_called_once_with('GET', '/profiles/status', timeout=3) + + def test_waiting_chat_times_out_without_upstream_work(self): + handler = object.__new__(router.Handler) + errors = [] + handler._send_error = lambda *args: errors.append(args) + router.STATE.lock.acquire() + thread = None + try: + with patch.object(router, 'CHAT_WAIT_TIMEOUT', 0.02), \ + patch.object(router, 'upstream_status') as upstream: + thread = threading.Thread(target=handler._chat_proxy, + args=(None, {'messages': []}, None)) + thread.start() + thread.join(0.5) + self.assertFalse(thread.is_alive(), 'lock wait ignored timeout') + upstream.assert_not_called() + self.assertEqual(errors[0][0], 503) + self.assertEqual(errors[0][3], 'model_wait_timeout') + finally: + router.STATE.lock.release() + if thread: + thread.join(1) + + def test_ultra_catalog_is_text_only(self): + with patch.object(router, 'current_profile', return_value='medium'): + catalog = router.Handler._llamacpp_models_payload()['data'] + ultra = next(x for x in catalog if x['id'] == 'qwen-ultra') + self.assertEqual(ultra['architecture']['input_modalities'], ['text']) + + def test_ultra_image_rejected_before_profile_switch(self): + handler = object.__new__(router.Handler) + errors = [] + handler._send_error = lambda *args: errors.append(args) + data = {'messages': [{'role': 'user', 'content': [{'type': 'image_url', + 'image_url': {'url': 'data:image/png;base64,iVBORw0KGgo='}}]}]} + with patch.object(router, 'switch_profile') as switch: + handler._chat_proxy(None, data, 'ultra') + switch.assert_not_called() + self.assertEqual(errors[0][0], 400) + + def test_ready_chat_reuses_profile_readiness_result_and_releases_lease(self): + handler = object.__new__(router.Handler) + before = router.STATE.active_chats + available = router.STATE.qwen_unavailable + router.STATE.qwen_unavailable = False + handler._proxy = lambda body: self.assertEqual( + json.loads(body)['model'], 'qwen-medium') + try: + with patch.object(router, 'switch_profile', return_value={ + 'reachable': True, 'model': 'qwen-medium', 'ctx': 160000}), \ + patch.object(router, 'upstream_status') as upstream: + handler._chat_proxy(None, {'messages': []}, 'medium') + upstream.assert_not_called() + self.assertEqual(router.STATE.active_chats, before) + finally: + router.STATE.qwen_unavailable = available + + def test_startup_accepts_tested_mtp_context_overhead_without_restart(self): + with patch.object(router.RUNTIME, 'load', return_value={}), \ + patch.object(router.RUNTIME, 'save'), \ + patch.object(router, 'enforce_artifact_retention', return_value=[]), \ + patch.object(router, 'current_profile', return_value='medium'), \ + patch.object(router, 'upstream_status', return_value={ + 'reachable': True, 'model': 'qwen-medium', 'ctx': 160128}), \ + patch.object(router, '_restore_qwen') as restore: + router._startup_reconcile() + restore.assert_not_called() + + +if __name__ == '__main__': + unittest.main() diff --git a/manage.sh b/manage.sh index c97c6de..ea66b50 100755 --- a/manage.sh +++ b/manage.sh @@ -38,12 +38,12 @@ shift || true case "$command" in validate) - run python3 "$ROOT_DIR/platform/scripts/sync-profile-matrix.py" --check + run bash "$ROOT_DIR/dev/check.sh" run "${compose[@]}" config --quiet ;; deploy) [[ $# -gt 0 ]] || { usage >&2; exit 2; } - run python3 "$ROOT_DIR/platform/scripts/sync-profile-matrix.py" --check + run bash "$ROOT_DIR/dev/check.sh" if ! docker volume inspect portainer_data >/dev/null 2>&1; then run docker volume create portainer_data fi diff --git a/platform/docker/profile-controller/profile_controller.py b/platform/docker/profile-controller/profile_controller.py index 0fa9c33..4c3686a 100644 --- a/platform/docker/profile-controller/profile_controller.py +++ b/platform/docker/profile-controller/profile_controller.py @@ -323,6 +323,8 @@ def stop_container(item: dict, timeout: int = 120) -> None: def start_container(item: dict) -> None: + if item.get("State") == "running": + return status, _ = docker_request("POST", f"/containers/{item['Id']}/start") if status not in (204, 304): raise RuntimeError(f"failed to start container: HTTP {status}") @@ -571,7 +573,7 @@ def set_video_worker(running: bool) -> dict: def active_profile(items: dict[str, dict] | None = None) -> str | None: - items = items or containers() + items = containers() if items is None else items active = [name for name, item in items.items() if item.get("State") == "running"] if len(active) > 1: raise RuntimeError(f"multiple llama profiles active: {', '.join(active)}") @@ -582,6 +584,11 @@ def activate(profile: str) -> dict: if profile not in ALLOWED: raise ValueError("profile is not allowlisted") with LOCK: + # Validate before stopping any working service. + items = containers() + missing = [name for name in ALLOWED if name not in items] + if missing: + raise RuntimeError("profile containers missing: " + ", ".join(missing)) # Defensive mutual exclusion even if a caller bypasses the router. for worker in image_containers(): stop_container(worker) @@ -592,10 +599,6 @@ def activate(profile: str) -> dict: stop_trellis_if_configured() stop_video_if_configured() start_container(tts_container()) - items = containers() - missing = [name for name in ALLOWED if name not in items] - if missing: - raise RuntimeError("profile containers missing: " + ", ".join(missing)) current = active_profile(items) if current == profile: return {"active_profile": current, "changed": False} @@ -613,6 +616,58 @@ def activate(profile: str) -> dict: return {"active_profile": profile, "changed": True} +def status_snapshot() -> dict: + """One Docker snapshot; missing optional workers cannot hide the LLM.""" + status, body = docker_request("GET", "/containers/json?all=1") + if status != 200: + raise RuntimeError(f"Docker list failed with HTTP {status}") + records = json.loads(body) + profiles = {} + for item in records: + name = item.get("Labels", {}).get(LABEL_KEY) + if name in ALLOWED: + if name in profiles: + raise RuntimeError(f"duplicate container for profile {name}") + profiles[name] = item + result = {"active_profile": active_profile(profiles), + "profiles": {name: profiles.get(name, {}).get("State", "missing") + for name in ALLOWED}, "worker_errors": {}} + workers = ( + ("music", MUSIC_LABEL_KEY, MUSIC_WORKER), + ("yue2", MUSIC_LABEL_KEY, YUE2_WORKER), + ("separator", SEPARATOR_LABEL_KEY, SEPARATOR_WORKER), + ("voice", VOICE_LABEL_KEY, VOICE_WORKER), + ("voice_change", VOICE_CHANGE_LABEL_KEY, VOICE_CHANGE_WORKER), + ("applio", APPLIO_LABEL_KEY, APPLIO_WORKER), + ("trellis", TRELLIS_LABEL_KEY, TRELLIS_WORKER), + ("video", VIDEO_LABEL_KEY, VIDEO_WORKER), + ) + result["video_ui"] = "unavailable" + for name, label, kind in workers: + matches = [item for item in records + if kind and item.get("Labels", {}).get(label) == kind] + if not kind: + state = health = "disabled" + elif len(matches) != 1: + state = health = "missing" if not matches else "error" + result["worker_errors"][name] = f"expected one container, found {len(matches)}" + else: + item = matches[0] + state = item.get("State", "unknown") + description = item.get("Status", "") + health = ("healthy" if "(healthy)" in description else + "unhealthy" if "(unhealthy)" in description else + "starting" if state == "running" else "stopped") + if name == "video": + try: + result["video_ui"] = cached_video_ui_state(item) + except Exception as exc: + result["worker_errors"]["video_ui"] = str(exc) + result[f"{name}_worker"] = state + result[f"{name}_health"] = health + return result + + def load_token() -> str: token = os.environ.get("CONTROLLER_TOKEN", "").strip() if not token: @@ -647,90 +702,18 @@ class Handler(BaseHTTPRequestHandler): if self.path == "/health": self.reply(200, {"status": "ok"}) return - if self.path != "/status": + if self.path not in {"/status", "/profiles/status"}: self.reply(404, {"error": "not found"}) return if not self.authenticated(): self.reply(401, {"error": "unauthorized"}) return try: - items = containers() - music = music_container() if MUSIC_WORKER else {} - music_status = music.get("Status", "") - music_health = ("disabled" if not MUSIC_WORKER else - "healthy" if "(healthy)" in music_status else - "unhealthy" if "(unhealthy)" in music_status else - "starting" if music.get("State") == "running" else - "stopped") - yue2 = yue2_container() if YUE2_WORKER else {} - yue2_status = yue2.get("Status", "") - yue2_health = ("disabled" if not YUE2_WORKER else - "healthy" if "(healthy)" in yue2_status else - "unhealthy" if "(unhealthy)" in yue2_status else - "starting" if yue2.get("State") == "running" else - "stopped") - separator = separator_container() if SEPARATOR_WORKER else {} - separator_status = separator.get("Status", "") - separator_health = ("disabled" if not SEPARATOR_WORKER else - "healthy" if "(healthy)" in separator_status else - "unhealthy" if "(unhealthy)" in separator_status else - "starting" if separator.get("State") == "running" else - "stopped") - voice = voice_container() if VOICE_WORKER else {} - voice_status = voice.get("Status", "") - voice_health = ("disabled" if not VOICE_WORKER else - "healthy" if "(healthy)" in voice_status else - "unhealthy" if "(unhealthy)" in voice_status else - "starting" if voice.get("State") == "running" else - "stopped") - voice_change = voice_change_container() if VOICE_CHANGE_WORKER else {} - voice_change_status = voice_change.get("Status", "") - voice_change_health = ("disabled" if not VOICE_CHANGE_WORKER else - "healthy" if "(healthy)" in voice_change_status else - "unhealthy" if "(unhealthy)" in voice_change_status else - "starting" if voice_change.get("State") == "running" else - "stopped") - applio = applio_container() if APPLIO_WORKER else {} - applio_status = applio.get("Status", "") - applio_health = ("disabled" if not APPLIO_WORKER else - "healthy" if "(healthy)" in applio_status else - "unhealthy" if "(unhealthy)" in applio_status else - "starting" if applio.get("State") == "running" else - "stopped") - trellis = trellis_container() if TRELLIS_WORKER else {} - trellis_status = trellis.get("Status", "") - trellis_health = ("disabled" if not TRELLIS_WORKER else - "healthy" if "(healthy)" in trellis_status else - "unhealthy" if "(unhealthy)" in trellis_status else - "starting" if trellis.get("State") == "running" else - "stopped") - video = video_container() if VIDEO_WORKER else {} - video_status = video.get("Status", "") - video_health = ("disabled" if not VIDEO_WORKER else - "healthy" if "(healthy)" in video_status else - "unhealthy" if "(unhealthy)" in video_status else - "starting" if video.get("State") == "running" else - "stopped") - self.reply(200, {"active_profile": active_profile(items), - "music_worker": music.get("State", "disabled"), - "music_health": music_health, - "yue2_worker": yue2.get("State", "disabled"), - "yue2_health": yue2_health, - "separator_worker": separator.get("State", "disabled"), - "separator_health": separator_health, - "voice_worker": voice.get("State", "disabled"), - "voice_health": voice_health, - "voice_change_worker": voice_change.get("State", "disabled"), - "voice_change_health": voice_change_health, - "applio_worker": applio.get("State", "disabled"), - "applio_health": applio_health, - "trellis_worker": trellis.get("State", "disabled"), - "trellis_health": trellis_health, - "video_worker": video.get("State", "disabled"), - "video_health": video_health, - "video_ui": cached_video_ui_state(video) if video else "unavailable", - "profiles": {name: items.get(name, {}).get( - "State", "missing") for name in ALLOWED}}) + if self.path == "/profiles/status": + items = containers() + self.reply(200, {"active_profile": active_profile(items)}) + else: + self.reply(200, status_snapshot()) except Exception as exc: log.exception("status failed") self.reply(503, {"error": str(exc)}) diff --git a/platform/scripts/sync-profile-matrix.py b/platform/scripts/sync-profile-matrix.py index ba7ac74..c33c7bb 100755 --- a/platform/scripts/sync-profile-matrix.py +++ b/platform/scripts/sync-profile-matrix.py @@ -44,6 +44,7 @@ def render_router(data: dict) -> str: item["id"]: { "context": item["context"], "model_alias": item["alias"], + "vision": bool(item.get("vision", False)), } for item in data["profiles"] } diff --git a/router/ai_profile_router.py b/router/ai_profile_router.py index e6ec128..7a945e5 100755 --- a/router/ai_profile_router.py +++ b/router/ai_profile_router.py @@ -73,6 +73,7 @@ import threading import time import uuid import http.client +from contextlib import contextmanager import urllib.error import urllib.parse import urllib.request @@ -299,6 +300,21 @@ class _State: 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. @@ -965,7 +981,7 @@ def current_profile() -> str | None: """ if PROFILE_CONTROL_URL: try: - profile = _profile_controller_request("GET", "/status").get( + profile = _profile_controller_request("GET", "/profiles/status", timeout=3).get( "active_profile") return profile if profile in PROFILES else None except Exception as exc: @@ -1029,7 +1045,7 @@ def _context_matches(expected: int, reported: object) -> bool: and expected <= reported <= expected + 1024) -def switch_profile(profile: str, implicit: bool = False) -> None: +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. @@ -1058,7 +1074,7 @@ def switch_profile(profile: str, implicit: bool = False) -> None: # so clear the stale flag before returning. _set_qwen_unavailable(False) log.info("Profil %s ist bereits aktiv", profile) - return + return up # Qwen wird neu geladen/gewechselt → für Chats nicht verfügbar. _set_qwen_unavailable(True) try: @@ -1067,7 +1083,7 @@ def switch_profile(profile: str, implicit: bool = False) -> None: # 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 + return upstream_status() if cur == profile and not up["reachable"] and implicit: raise RuntimeError( f"llama.cpp nicht erreichbar (Profil {profile} ist " @@ -1114,6 +1130,7 @@ def switch_profile(profile: str, implicit: bool = False) -> None: _set_qwen_unavailable(not available) finally: STATE.switching = None + return upstream_status() # --------------------------------------------------------------------------- @@ -2102,7 +2119,8 @@ class Handler(BaseHTTPRequestHandler): "value": "loaded" if name == active else "unloaded", }, "architecture": { - "input_modalities": ["text", "image"], + "input_modalities": (["text", "image"] + if PROFILE_REGISTRY[name]["vision"] else ["text"]), "output_modalities": ["text"], }, }) @@ -2155,29 +2173,27 @@ class Handler(BaseHTTPRequestHandler): @staticmethod def _mode_payload() -> dict: state = RUNTIME.load() - return { + 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, - "music_worker": _music_worker_state(), - "music_health": _music_worker_health(), - "yue2_worker": _worker_field("yue2_worker"), - "yue2_health": _worker_field("yue2_health"), - "separator_worker": _separator_worker_state(), - "separator_health": _separator_worker_health(), - "voice_worker": _voice_worker_state(), - "voice_health": _voice_worker_health(), - "voice_change_worker": _voice_change_worker_state(), - "voice_change_health": _voice_change_worker_health(), - "applio_worker": _worker_field("applio_worker"), - "applio_health": _worker_field("applio_health"), - "trellis_worker": _worker_field("trellis_worker"), - "trellis_health": _worker_field("trellis_health"), - "video_worker": _worker_field("video_worker"), - "video_health": _worker_field("video_health"), "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: @@ -2940,10 +2956,9 @@ class Handler(BaseHTTPRequestHandler): def _acquire_model_lease(self, profile: str | None = None) -> dict: """Atomar Profil sicherstellen und einen aktiven Request registrieren.""" - with STATE.lock: - if profile is not None: - switch_profile(profile, implicit=True) - up = upstream_status() + 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: @@ -2961,6 +2976,9 @@ class Handler(BaseHTTPRequestHandler): 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") @@ -2976,19 +2994,22 @@ class Handler(BaseHTTPRequestHandler): """Bildvalidierung, Profilwahl und Chat-Lease als eine Transaktion.""" lease_acquired = False try: - with STATE.lock: - if profile is not None: - switch_profile(profile, implicit=True) - - data = _normalize_llamacpp_reasoning(data) - data = _cap_chat_generation(data) - - if _request_has_image(data): - data = _normalize_chat_images(data) + # 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") - up = upstream_status() if not up["reachable"] or not up.get("model"): raise RuntimeError("llama.cpp nicht erreichbar") if profile is not None: @@ -3000,7 +3021,11 @@ class Handler(BaseHTTPRequestHandler): STATE.active_chats += 1 lease_acquired = True self._proxy(body) - except (ValueError, RuntimeError) as e: + 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: @@ -3218,10 +3243,7 @@ def _startup_reconcile() -> None: return up = upstream_status() - if (up["reachable"] and up.get("model") - and up.get("ctx") == PROFILES[profile] - and (not EXPECTED_MODELS.get(profile) - or up.get("model") == EXPECTED_MODELS[profile])): + if _profile_is_ready(profile, up): _set_qwen_unavailable(False) STATE.mode = "llm" STATE.mode_phase = "ready" diff --git a/router/router_profiles.json b/router/router_profiles.json index ceb40c0..b6bf0d4 100644 --- a/router/router_profiles.json +++ b/router/router_profiles.json @@ -2,23 +2,28 @@ "profiles": { "fast": { "context": 76800, - "model_alias": "qwen-fast" + "model_alias": "qwen-fast", + "vision": true }, "medium": { "context": 160000, - "model_alias": "qwen-medium" + "model_alias": "qwen-medium", + "vision": true }, "large": { "context": 192000, - "model_alias": "qwen-large" + "model_alias": "qwen-large", + "vision": true }, "ultra": { "context": 262144, - "model_alias": "qwen-ultra" + "model_alias": "qwen-ultra", + "vision": false }, "uncensored": { "context": 80000, - "model_alias": "qwen-uncensored" + "model_alias": "qwen-uncensored", + "vision": true } } } diff --git a/router/router_support.py b/router/router_support.py index 0790cba..84cbae4 100644 --- a/router/router_support.py +++ b/router/router_support.py @@ -39,6 +39,8 @@ def load_profile_registry(path: str | None) -> dict[str, dict]: "ultra": {"context": 262144, "model_alias": None}, "uncensored": {"context": 80000, "model_alias": None}, } + for name, definition in fallback.items(): + definition["vision"] = name != "ultra" if not path: return fallback try: @@ -56,12 +58,16 @@ def load_profile_registry(path: str | None) -> dict[str, dict]: raise ConfigurationError(f"Profil {name!r} ist kein Objekt") context = definition.get("context") alias = definition.get("model_alias") + vision = definition.get("vision", name != "ultra") + if not isinstance(vision, bool): + raise ConfigurationError(f"ungültige Vision-Fähigkeit für Profil {name!r}") if not isinstance(context, int) or context < 1024: raise ConfigurationError(f"ungültiger Kontext für Profil {name!r}") if alias is not None and (not isinstance(alias, str) or not alias.strip()): raise ConfigurationError(f"ungültiger Modellalias für Profil {name!r}") validated[name] = {"context": context, - "model_alias": alias.strip() if alias else None} + "model_alias": alias.strip() if alias else None, + "vision": vision} return validated