Reduce controller polling and bound chat admission waits

This commit is contained in:
Mikei386
2026-09-20 19:10:59 +02:00
parent 53320fa6c3
commit 6071b1a2cd
10 changed files with 303 additions and 134 deletions
Executable
+12
View File
@@ -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
+32 -1
View File
@@ -1,4 +1,5 @@
import importlib.util import importlib.util
import json
import os import os
import pathlib import pathlib
import unittest import unittest
@@ -42,6 +43,34 @@ def music_item(state="exited"):
class ProfileControllerTests(unittest.TestCase): 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): def test_music_start_exclusively_stops_gpu_workers(self):
profiles = {name: item(name) for name in controller.ALLOWED} profiles = {name: item(name) for name in controller.ALLOWED}
profiles["ultra"] = item("ultra", "running") profiles["ultra"] = item("ultra", "running")
@@ -99,9 +128,11 @@ class ProfileControllerTests(unittest.TestCase):
profiles = {name: item(name) for name in controller.ALLOWED[:-1]} profiles = {name: item(name) for name in controller.ALLOWED[:-1]}
with patch.object(controller, "containers", return_value=profiles), \ with patch.object(controller, "containers", return_value=profiles), \
patch.object(controller, "image_containers", return_value=[image_item()]), \ 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"): with self.assertRaisesRegex(RuntimeError, "missing"):
controller.activate("fast") controller.activate("fast")
request.assert_not_called()
def test_image_start_stops_inference_first(self): def test_image_start_stops_inference_first(self):
profiles = {name: item(name) for name in controller.ALLOWED} profiles = {name: item(name) for name in controller.ALLOWED}
+8 -1
View File
@@ -1,4 +1,5 @@
import io import io
import re
import subprocess import subprocess
import tarfile import tarfile
import tempfile import tempfile
@@ -94,7 +95,13 @@ class RecoveryScriptTests(unittest.TestCase):
gateway = (ROOT / "platform/docker/wireguard-gateway/entrypoint.sh").read_text( gateway = (ROOT / "platform/docker/wireguard-gateway/entrypoint.sh").read_text(
encoding="utf-8" 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('"8099:8099"', compose)
self.assertNotIn('"9443:9443"', compose) self.assertNotIn('"9443:9443"', compose)
self.assertIn('start_proxy 8099 llama-dashboard:8099', gateway) self.assertIn('start_proxy 8099 llama-dashboard:8099', gateway)
+102
View File
@@ -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()
+2 -2
View File
@@ -38,12 +38,12 @@ shift || true
case "$command" in case "$command" in
validate) validate)
run python3 "$ROOT_DIR/platform/scripts/sync-profile-matrix.py" --check run bash "$ROOT_DIR/dev/check.sh"
run "${compose[@]}" config --quiet run "${compose[@]}" config --quiet
;; ;;
deploy) deploy)
[[ $# -gt 0 ]] || { usage >&2; exit 2; } [[ $# -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 if ! docker volume inspect portainer_data >/dev/null 2>&1; then
run docker volume create portainer_data run docker volume create portainer_data
fi fi
@@ -323,6 +323,8 @@ def stop_container(item: dict, timeout: int = 120) -> None:
def start_container(item: dict) -> None: def start_container(item: dict) -> None:
if item.get("State") == "running":
return
status, _ = docker_request("POST", f"/containers/{item['Id']}/start") status, _ = docker_request("POST", f"/containers/{item['Id']}/start")
if status not in (204, 304): if status not in (204, 304):
raise RuntimeError(f"failed to start container: HTTP {status}") 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: 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"] active = [name for name, item in items.items() if item.get("State") == "running"]
if len(active) > 1: if len(active) > 1:
raise RuntimeError(f"multiple llama profiles active: {', '.join(active)}") raise RuntimeError(f"multiple llama profiles active: {', '.join(active)}")
@@ -582,6 +584,11 @@ def activate(profile: str) -> dict:
if profile not in ALLOWED: if profile not in ALLOWED:
raise ValueError("profile is not allowlisted") raise ValueError("profile is not allowlisted")
with LOCK: 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. # Defensive mutual exclusion even if a caller bypasses the router.
for worker in image_containers(): for worker in image_containers():
stop_container(worker) stop_container(worker)
@@ -592,10 +599,6 @@ def activate(profile: str) -> dict:
stop_trellis_if_configured() stop_trellis_if_configured()
stop_video_if_configured() stop_video_if_configured()
start_container(tts_container()) 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) current = active_profile(items)
if current == profile: if current == profile:
return {"active_profile": current, "changed": False} return {"active_profile": current, "changed": False}
@@ -613,6 +616,58 @@ def activate(profile: str) -> dict:
return {"active_profile": profile, "changed": True} 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: def load_token() -> str:
token = os.environ.get("CONTROLLER_TOKEN", "").strip() token = os.environ.get("CONTROLLER_TOKEN", "").strip()
if not token: if not token:
@@ -647,90 +702,18 @@ class Handler(BaseHTTPRequestHandler):
if self.path == "/health": if self.path == "/health":
self.reply(200, {"status": "ok"}) self.reply(200, {"status": "ok"})
return return
if self.path != "/status": if self.path not in {"/status", "/profiles/status"}:
self.reply(404, {"error": "not found"}) self.reply(404, {"error": "not found"})
return return
if not self.authenticated(): if not self.authenticated():
self.reply(401, {"error": "unauthorized"}) self.reply(401, {"error": "unauthorized"})
return return
try: try:
if self.path == "/profiles/status":
items = containers() items = containers()
music = music_container() if MUSIC_WORKER else {} self.reply(200, {"active_profile": active_profile(items)})
music_status = music.get("Status", "") else:
music_health = ("disabled" if not MUSIC_WORKER else self.reply(200, status_snapshot())
"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}})
except Exception as exc: except Exception as exc:
log.exception("status failed") log.exception("status failed")
self.reply(503, {"error": str(exc)}) self.reply(503, {"error": str(exc)})
+1
View File
@@ -44,6 +44,7 @@ def render_router(data: dict) -> str:
item["id"]: { item["id"]: {
"context": item["context"], "context": item["context"],
"model_alias": item["alias"], "model_alias": item["alias"],
"vision": bool(item.get("vision", False)),
} }
for item in data["profiles"] for item in data["profiles"]
} }
+60 -38
View File
@@ -73,6 +73,7 @@ import threading
import time import time
import uuid import uuid
import http.client import http.client
from contextlib import contextmanager
import urllib.error import urllib.error
import urllib.parse import urllib.parse
import urllib.request import urllib.request
@@ -299,6 +300,21 @@ class _State:
STATE = _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: def _wait_chats_drained(timeout: float | None = None) -> None:
"""Wartet, bis keine aktiven Chat-Requests mehr laufen. """Wartet, bis keine aktiven Chat-Requests mehr laufen.
@@ -965,7 +981,7 @@ def current_profile() -> str | None:
""" """
if PROFILE_CONTROL_URL: if PROFILE_CONTROL_URL:
try: try:
profile = _profile_controller_request("GET", "/status").get( profile = _profile_controller_request("GET", "/profiles/status", timeout=3).get(
"active_profile") "active_profile")
return profile if profile in PROFILES else None return profile if profile in PROFILES else None
except Exception as exc: except Exception as exc:
@@ -1029,7 +1045,7 @@ def _context_matches(expected: int, reported: object) -> bool:
and expected <= reported <= expected + 1024) 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. """Stellt sicher, dass das Profil aktiv ist, und wartet bis es geladen ist.
Wirft RuntimeError, wenn das Profil nicht aktiviert werden konnte. 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. # so clear the stale flag before returning.
_set_qwen_unavailable(False) _set_qwen_unavailable(False)
log.info("Profil %s ist bereits aktiv", profile) log.info("Profil %s ist bereits aktiv", profile)
return return up
# Qwen wird neu geladen/gewechselt → für Chats nicht verfügbar. # Qwen wird neu geladen/gewechselt → für Chats nicht verfügbar.
_set_qwen_unavailable(True) _set_qwen_unavailable(True)
try: try:
@@ -1067,7 +1083,7 @@ def switch_profile(profile: str, implicit: bool = False) -> None:
# Modell wird gerade geladen (z.B. nach einem Wechsel) # Modell wird gerade geladen (z.B. nach einem Wechsel)
log.info("Warte, bis Profil %s geladen ist ...", profile) log.info("Warte, bis Profil %s geladen ist ...", profile)
_wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT) _wait_ready(profile, time.monotonic() + SWITCH_TIMEOUT)
return return upstream_status()
if cur == profile and not up["reachable"] and implicit: if cur == profile and not up["reachable"] and implicit:
raise RuntimeError( raise RuntimeError(
f"llama.cpp nicht erreichbar (Profil {profile} ist " 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) _set_qwen_unavailable(not available)
finally: finally:
STATE.switching = None STATE.switching = None
return upstream_status()
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -2102,7 +2119,8 @@ class Handler(BaseHTTPRequestHandler):
"value": "loaded" if name == active else "unloaded", "value": "loaded" if name == active else "unloaded",
}, },
"architecture": { "architecture": {
"input_modalities": ["text", "image"], "input_modalities": (["text", "image"]
if PROFILE_REGISTRY[name]["vision"] else ["text"]),
"output_modalities": ["text"], "output_modalities": ["text"],
}, },
}) })
@@ -2155,29 +2173,27 @@ class Handler(BaseHTTPRequestHandler):
@staticmethod @staticmethod
def _mode_payload() -> dict: def _mode_payload() -> dict:
state = RUNTIME.load() 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, "active": STATE.mode,
"phase": STATE.mode_phase, "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"), "return_profile": state.get("return_profile"),
"last_error": STATE.mode_error, "last_error": STATE.mode_error,
"enabled": ENABLE_MUSIC_MODE, "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: def _mode_change(self) -> None:
try: try:
@@ -2940,10 +2956,9 @@ class Handler(BaseHTTPRequestHandler):
def _acquire_model_lease(self, profile: str | None = None) -> dict: def _acquire_model_lease(self, profile: str | None = None) -> dict:
"""Atomar Profil sicherstellen und einen aktiven Request registrieren.""" """Atomar Profil sicherstellen und einen aktiven Request registrieren."""
with STATE.lock: with _model_lock():
if profile is not None: up = (switch_profile(profile, implicit=True)
switch_profile(profile, implicit=True) if profile is not None else upstream_status())
up = upstream_status()
if not up["reachable"] or not up.get("model"): if not up["reachable"] or not up.get("model"):
raise RuntimeError("llama.cpp nicht erreichbar") raise RuntimeError("llama.cpp nicht erreichbar")
with STATE.avail_lock: with STATE.avail_lock:
@@ -2961,6 +2976,9 @@ class Handler(BaseHTTPRequestHandler):
profile: str) -> None: profile: str) -> None:
try: try:
up = self._acquire_model_lease(profile) 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: except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error", self._send_error(502, str(e), "server_error",
"upstream_unavailable") "upstream_unavailable")
@@ -2976,19 +2994,22 @@ class Handler(BaseHTTPRequestHandler):
"""Bildvalidierung, Profilwahl und Chat-Lease als eine Transaktion.""" """Bildvalidierung, Profilwahl und Chat-Lease als eine Transaktion."""
lease_acquired = False lease_acquired = False
try: try:
with STATE.lock: # Validation and image decoding need no GPU ownership.
if profile is not None:
switch_profile(profile, implicit=True)
data = _normalize_llamacpp_reasoning(data) data = _normalize_llamacpp_reasoning(data)
data = _cap_chat_generation(data) data = _cap_chat_generation(data)
has_image = _request_has_image(data)
if _request_has_image(data): if has_image:
data = _normalize_chat_images(data) 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 " log.info("Vision: Bild wird direkt an das aktive "
"multimodale Qwen-Profil weitergeleitet") "multimodale Qwen-Profil weitergeleitet")
up = upstream_status()
if not up["reachable"] or not up.get("model"): if not up["reachable"] or not up.get("model"):
raise RuntimeError("llama.cpp nicht erreichbar") raise RuntimeError("llama.cpp nicht erreichbar")
if profile is not None: if profile is not None:
@@ -3000,7 +3021,11 @@ class Handler(BaseHTTPRequestHandler):
STATE.active_chats += 1 STATE.active_chats += 1
lease_acquired = True lease_acquired = True
self._proxy(body) 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") self._send_error(502, str(e), "server_error", "upstream_unavailable")
finally: finally:
if lease_acquired: if lease_acquired:
@@ -3218,10 +3243,7 @@ def _startup_reconcile() -> None:
return return
up = upstream_status() up = upstream_status()
if (up["reachable"] and up.get("model") if _profile_is_ready(profile, up):
and up.get("ctx") == PROFILES[profile]
and (not EXPECTED_MODELS.get(profile)
or up.get("model") == EXPECTED_MODELS[profile])):
_set_qwen_unavailable(False) _set_qwen_unavailable(False)
STATE.mode = "llm" STATE.mode = "llm"
STATE.mode_phase = "ready" STATE.mode_phase = "ready"
+10 -5
View File
@@ -2,23 +2,28 @@
"profiles": { "profiles": {
"fast": { "fast": {
"context": 76800, "context": 76800,
"model_alias": "qwen-fast" "model_alias": "qwen-fast",
"vision": true
}, },
"medium": { "medium": {
"context": 160000, "context": 160000,
"model_alias": "qwen-medium" "model_alias": "qwen-medium",
"vision": true
}, },
"large": { "large": {
"context": 192000, "context": 192000,
"model_alias": "qwen-large" "model_alias": "qwen-large",
"vision": true
}, },
"ultra": { "ultra": {
"context": 262144, "context": 262144,
"model_alias": "qwen-ultra" "model_alias": "qwen-ultra",
"vision": false
}, },
"uncensored": { "uncensored": {
"context": 80000, "context": 80000,
"model_alias": "qwen-uncensored" "model_alias": "qwen-uncensored",
"vision": true
} }
} }
} }
+7 -1
View File
@@ -39,6 +39,8 @@ def load_profile_registry(path: str | None) -> dict[str, dict]:
"ultra": {"context": 262144, "model_alias": None}, "ultra": {"context": 262144, "model_alias": None},
"uncensored": {"context": 80000, "model_alias": None}, "uncensored": {"context": 80000, "model_alias": None},
} }
for name, definition in fallback.items():
definition["vision"] = name != "ultra"
if not path: if not path:
return fallback return fallback
try: try:
@@ -56,12 +58,16 @@ def load_profile_registry(path: str | None) -> dict[str, dict]:
raise ConfigurationError(f"Profil {name!r} ist kein Objekt") raise ConfigurationError(f"Profil {name!r} ist kein Objekt")
context = definition.get("context") context = definition.get("context")
alias = definition.get("model_alias") 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: if not isinstance(context, int) or context < 1024:
raise ConfigurationError(f"ungültiger Kontext für Profil {name!r}") 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()): 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}") raise ConfigurationError(f"ungültiger Modellalias für Profil {name!r}")
validated[name] = {"context": context, validated[name] = {"context": context,
"model_alias": alias.strip() if alias else None} "model_alias": alias.strip() if alias else None,
"vision": vision}
return validated return validated