Replace vision hotswap with native multimodal profiles

This commit is contained in:
Mikei386 committed 2026-08-20 14:24:56 +02:00
1 parent 8a45a0d851
commit cb07779f5a
19 files changed
+169 -719

No files matched your search

+83 -605
View File
@@ -7,7 +7,7 @@ Profilen um:
Profil Kontext
------ --------
fast 73728
fast 76800
medium 94208
long 131072
@@ -39,14 +39,11 @@ 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-Orchestrierung (Q3 "Augen", temporär):
POST /v1/chat/completions mit Bild im letzten
User-Message → ein temporäres Q3-Vision-Modell
(llama-server + mmproj) wird geladen, analysiert
das Bild und wird wieder entladen; das Hauptprofil
wird immer wiederhergestellt (try/finally) und
erzeugt die Endantwort. Die Vision-Analyse bleibt
intern (kein sichtbarer Assistant-Turn).
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).
"""
@@ -56,7 +53,6 @@ from __future__ import annotations
import base64
import binascii
import email
import hashlib
import ipaddress
import json
import logging
@@ -73,7 +69,6 @@ import http.client
import urllib.error
import urllib.parse
import urllib.request
from collections import OrderedDict
from email.parser import BytesParser
from email.policy import compat32
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
@@ -122,30 +117,12 @@ 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"))
# --- Vision-Orchestrierung (Q3 "Augen", temporär) ---
# Chat-Requests mit Bild im letzten User-Message lösen einen
# temporären Hotswap aus: Hauptprofil raus → Q3+mmproj rein →
# Bild analysieren → Q3 raus → Hauptprofil rein (try/finally).
LLAMA_SERVER_BIN = os.environ.get(
"LLAMA_SERVER_BIN", "/opt/mike-ai/llama.cpp/build/bin/llama-server")
VISION_MODEL = os.environ.get(
"VISION_MODEL", "/opt/mike-ai/models/qwen3.8-27b/Qwen3.8-27B-Q3_K_M.gguf")
VISION_MMPROJ = os.environ.get(
"VISION_MMPROJ", "/opt/mike-ai/models/qwen3.8-27b-nvfp4/mmproj-BF16.gguf")
VISION_CTX = int(os.environ.get("VISION_CTX", "32768"))
VISION_PORT = int(os.environ.get("VISION_PORT", "8086"))
VISION_ALIAS = "qwen38-27b-q3-vision"
VISION_LOAD_TIMEOUT = float(os.environ.get("VISION_LOAD_TIMEOUT", "300")) # s
VISION_INFER_TIMEOUT = float(os.environ.get("VISION_INFER_TIMEOUT", "300")) # s
VISION_UNLOAD_TIMEOUT = float(os.environ.get("VISION_UNLOAD_TIMEOUT", "120")) # s
VISION_MAX_TOKENS = int(os.environ.get("VISION_MAX_TOKENS", "4096"))
VISION_CACHE_MAX = int(os.environ.get("VISION_CACHE_MAX", "64"))
VISION_MAX_IMAGE_BYTES = int(os.environ.get(
"VISION_MAX_IMAGE_BYTES", str(20 * 1024 * 1024)))
VISION_ALLOW_REMOTE_URLS = os.environ.get(
"VISION_ALLOW_REMOTE_URLS", "false").lower() in {"1", "true", "yes"}
VISION_LOG = os.environ.get(
"VISION_LOG", "/opt/mike-ai/ai-profile-router/vision_server.log")
# --- 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). FLUX.2 klein ist für 1 MP
# ausgelegt; 1920x1088 (≈2 MP) wird zusätzlich unterstützt.
@@ -247,43 +224,20 @@ IMAGE_PHASES = (
)
class _VisionState:
"""Zustand der Vision-Orchestrierung (Status-Reporting + Analyse-Cache)."""
def __init__(self) -> None:
self.phase = "idle" # siehe VISION_PHASES unten
self.last_error: str | None = None
self.last_turnaround: float | None = None # s, letzter kompletter Swap
self.last_profile: str | None = None # Profil, das gesichert wurde
# Cache: stabiler Bild-Hash → Vision-Analyse-Text. Folgefragen im
# selben Chat (Open WebUI schickt den multimodalen Verlauf erneut)
# senden das Bild NICHT erneut durch Q3, sondern verwenden die
# gecachte Analyse. LRU-begrenzt auf VISION_CACHE_MAX Einträge.
self.analysis_cache: "OrderedDict[str, str]" = OrderedDict()
self.cache_lock = threading.Lock()
VISION_PHASES = (
"idle", "stopping-main", "loading-vision", "analyzing",
"unloading-vision", "restoring-main",
)
class _State:
"""Gemeinsamer, thread-sicherer Zustand.
lock : zentraler GPU-/Model-Lock. Wird von Profilwechsel,
Image-Generation UND Vision-Swap gehalten →
gegenseitiger Ausschluss, kein Race zwischen beiden.
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 + Vision + Chat-Lease,
# 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.vision = _VisionState()
self.qwen_unavailable = True
self.active_chats = 0
self.avail_lock = threading.Lock()
@@ -900,199 +854,43 @@ def _image_filename_ok(name: str) -> bool:
# ---------------------------------------------------------------------------
# Vision-Orchestrierung (Q3 "Augen", temporär)
# Multimodale Chat-Eingaben
# ---------------------------------------------------------------------------
VISION_ANALYST_PROMPT = (
"Du bist ein reiner Bild- und Screenshot-Analyst. Du beantwortest die "
"Benutzerfrage NICHT selbst. Erstelle eine umfassende, von der aktuellen "
"Frage unabhängige Bestandsaufnahme. Lasse keine sichtbaren Details aus, "
"damit auch spätere Folgefragen allein mit dieser Analyse beantwortet "
"werden können.\n\n"
"Antworte NUR mit einer strukturierten Analyse in dieser Form:\n"
"1. SIEHTBARER TEXT: alle Texte wörtlich und vollständig, mit Anordnung\n"
"2. UI-ELEMENTE: Felder, Buttons, Menüs, Tabs, Dropdowns, Checkboxen – "
"mit Namen, Werten und Zustand (aktiv/inaktiv, gefüllt/leer, ausgewählt)\n"
"3. FEHLER- UND WARNMELDUNGEN: wörtlich, mit Farbe und Position\n"
"4. POSITIONEN UND BEZIEHUNGEN: räumliche Anordnung (oben/unten, "
"links/rechts, Reihenfolge)\n"
"5. ZUSTÄNDE: Statusanzeigen, Farben (rot/grün/gelb), Ladezustände\n"
"6. OBJEKTE: relevante Objekte und Beziehungen zwischen Elementen\n"
"7. WEITERES: alle übrigen sichtbaren Details\n\n"
"Regeln: Bei Screenshots hat Text- und UI-Genauigkeit Vorrang vor "
"schöner Beschreibung. Keine Interpretation, keine Vermutungen – nur "
"was sichtbar ist. Unleserliches als [unleserlich] markieren."
)
IMAGE_PLACEHOLDER = "[Bild angehänggt – siehe Vision-Analyse]"
class _VisionServer:
"""Temporärer llama-server (Q3 + mmproj) für die Bildanalyse."""
def __init__(self) -> None:
self.proc: subprocess.Popen | 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
cmd = [
LLAMA_SERVER_BIN,
"--model", VISION_MODEL,
"--mmproj", VISION_MMPROJ,
"--alias", VISION_ALIAS,
"--ctx-size", str(VISION_CTX),
"--flash-attn", "on",
"--cache-type-k", "q4_0",
"--cache-type-v", "q4_0",
"--threads", "6",
"--threads-batch", "6",
"--batch-size", "64",
"--ubatch-size", "32",
"--parallel", "1",
"--jinja",
"--host", "127.0.0.1",
"--port", str(VISION_PORT),
"--metrics",
"--fit", "off",
"--n-gpu-layers", "all",
"--mmproj-offload",
"--no-mmap",
"--temperature", "0.2",
"--top-p", "0.8",
"--top-k", "20",
"--device", "CUDA0",
"--split-mode", "none",
]
self._logf = open(VISION_LOG, "ab")
self.proc = subprocess.Popen(
cmd, stdin=subprocess.DEVNULL,
stdout=self._logf, stderr=subprocess.STDOUT,
start_new_session=True)
RUNTIME.save(worker="vision", worker_pid=self.proc.pid,
phase="loading-vision")
log.info("Vision-Server gestartet (PID %d, Port %d, ctx %d)",
self.proc.pid, VISION_PORT, VISION_CTX)
def wait_ready(self, deadline: float) -> None:
"""Wartet, bis der Vision-Server das Modell mit erwartetem ctx meldet."""
while True:
try:
conn = http.client.HTTPConnection("127.0.0.1", VISION_PORT,
timeout=3)
conn.request("GET", "/v1/models")
resp = conn.getresponse()
data = json.loads(resp.read())
conn.close()
models = data.get("data") or []
if models and (models[0].get("meta") or {}).get("n_ctx") == VISION_CTX:
return
except (OSError, ValueError):
pass
if not self.alive():
raise RuntimeError(
"Vision-Server-Prozess beendet sich während des Ladens "
f"(Details: {VISION_LOG})")
if time.monotonic() > deadline:
raise RuntimeError(
f"Vision-Server nach {VISION_LOAD_TIMEOUT:.0f} s nicht "
f"bereit (Details: {VISION_LOG})")
time.sleep(2)
def stop(self) -> None:
if self.proc is not None and self.proc.poll() is None:
self.proc.terminate()
try:
self.proc.wait(timeout=30)
except subprocess.TimeoutExpired:
log.warning("Vision-Server reagiert nicht auf SIGTERM – SIGKILL")
self.proc.kill()
try:
self.proc.wait(timeout=10)
except subprocess.TimeoutExpired:
pass
if self._logf is not None:
try:
self._logf.close()
except OSError:
pass
self._logf = None
self.proc = None
RUNTIME.clear_worker("vision")
def _extract_last_user_image(data: dict) -> tuple[str | None, str]:
"""Liefert (image_url, question) aus der letzten User-Message.
Nur Bilder in der LETZTEN User-Message lösen eine Vision-Analyse aus.
Bilder in früheren Nachrichten sind bereits durch die vorherige
Antwort abgedeckt (Chat-Historie) und werden nur durch einen
Platzhalter ersetzt.
"""
messages = data.get("messages")
if not isinstance(messages, list):
return None, ""
for msg in reversed(messages):
if not isinstance(msg, dict) or msg.get("role") != "user":
continue
content = msg.get("content")
if isinstance(content, str):
return None, content
if isinstance(content, list):
image_url: str | None = None
texts: list[str] = []
for part in content:
if not isinstance(part, dict):
continue
if part.get("type") == "image_url":
iu = part.get("image_url")
url = iu.get("url") if isinstance(iu, dict) else iu
if isinstance(url, str) and url:
image_url = url
elif (part.get("type") == "text"
and isinstance(part.get("text"), str)):
texts.append(part["text"])
question = " ".join(t.strip() for t in texts if t.strip())
return image_url, question
return None, ""
_VISION_DATA_TYPES = {
_CHAT_IMAGE_DATA_TYPES = {
"image/jpeg", "image/png", "image/webp", "image/gif",
}
def _normalize_vision_image(image_url: str) -> str:
"""Validate an image and return a bounded data URL for llama.cpp.
def _normalize_chat_image(image_url: str) -> str:
"""Validiert ein Bild und liefert eine begrenzte data-URL.
Remote fetches are off by default. If explicitly enabled, the router
downloads the file itself after rejecting local/private destinations and
hands llama.cpp a data URL. The model server therefore never receives an
arbitrary user-controlled URL (SSRF protection).
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 _VISION_DATA_TYPES:
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) > VISION_MAX_IMAGE_BYTES:
if not raw or len(raw) > CHAT_IMAGE_MAX_BYTES:
raise ValueError(
f"Bildgröße außerhalb des Limits (max {VISION_MAX_IMAGE_BYTES} Bytes)")
"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 VISION_ALLOW_REMOTE_URLS:
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:
@@ -1100,7 +898,8 @@ def _normalize_vision_image(image_url: str) -> str:
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
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])
@@ -1115,49 +914,50 @@ def _normalize_vision_image(image_url: str) -> str:
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")
raise ValueError(
"Weiterleitungen zu einem anderen Bild-Host sind nicht erlaubt")
content_type = response.headers.get_content_type().lower()
if content_type not in _VISION_DATA_TYPES:
raise ValueError(f"Remote-Inhalt ist kein unterstütztes Bild ({content_type})")
raw = response.read(VISION_MAX_IMAGE_BYTES + 1)
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) > VISION_MAX_IMAGE_BYTES:
raise ValueError(
f"Bildgröße außerhalb des Limits (max {VISION_MAX_IMAGE_BYTES} Bytes)")
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 _image_hash(image_url: str) -> str:
"""Stabiler Hash für ein Bild (Base64-Payload oder URL).
Dasselbe Bild → derselbe Hash (unabhängig von Chat/Turn). Dient als
Schlüssel für den Vision-Analyse-Cache.
"""
if image_url.startswith("data:"):
payload = image_url.split(",", 1)[1] if "," in image_url else ""
try:
raw = base64.b64decode(payload)
return "img:" + hashlib.sha256(raw).hexdigest()[:32]
except (ValueError, TypeError):
return "b64:" + hashlib.sha256(payload.encode()).hexdigest()[:32]
return "url:" + hashlib.sha256(image_url.encode()).hexdigest()[:32]
def _vision_cached_analysis(img_hash: str) -> str | None:
"""Liefert die gecachte Vision-Analyse für einen Bild-Hash (oder None)."""
with STATE.vision.cache_lock:
return STATE.vision.analysis_cache.get(img_hash)
def _vision_store_analysis(img_hash: str, analysis: str) -> None:
"""Speichert eine Vision-Analyse im Cache (LRU-begrenzt)."""
with STATE.vision.cache_lock:
STATE.vision.analysis_cache[img_hash] = analysis
STATE.vision.analysis_cache.move_to_end(img_hash)
while len(STATE.vision.analysis_cache) > VISION_CACHE_MAX:
STATE.vision.analysis_cache.popitem(last=False)
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:
@@ -1165,260 +965,12 @@ def _request_has_image(data: dict) -> bool:
messages = data.get("messages")
if not isinstance(messages, list):
return False
for msg in messages:
if not isinstance(msg, dict):
continue
content = msg.get("content")
if isinstance(content, list):
for part in content:
if isinstance(part, dict) and part.get("type") == "image_url":
return True
return False
def _sanitize_for_main_model(data: dict) -> dict:
"""Erzeugt eine sanisierte Kopie des Requests für das Hauptmodell.
Das Hauptmodell (Fast/Medium/Long) hat KEIN Vision-Modell. Deshalb
werden alle Bild-Parts (image_url) durch den zu diesem Bild erzeugten
Vision-Analyse-Text (aus dem Cache) ersetzt:
- Bilddaten (Base64/URL) werden vollständig entfernt.
- Der ursprüngliche Text des Users bleibt erhalten.
- Die Analyse wird klar gekennzeichnet in denselben Turn eingesetzt.
Beispiel:
user: [image_url, "Was ist hier falsch?"]
→ user: "Was ist hier falsch?\n\n[Vision-Analyse des hochgeladenen
Bildes: ...]"
Wird bei JEDER Folgefrage erneut angewendet, weil Open WebUI den
ursprünglichen multimodalen Verlauf wieder mitsendet.
"""
out = json.loads(json.dumps(data))
messages = out.get("messages")
if not isinstance(messages, list):
return out
for msg in messages:
if not isinstance(msg, dict):
continue
content = msg.get("content")
if not isinstance(content, list):
continue
texts: list[str] = []
had_image = False
for part in content:
if isinstance(part, dict) and part.get("type") == "image_url":
had_image = True
iu = part.get("image_url")
url = iu.get("url") if isinstance(iu, dict) else iu
analysis = None
if isinstance(url, str) and url:
analysis = _vision_cached_analysis(_image_hash(url))
if analysis:
texts.append("[Vision-Analyse des hochgeladenen Bildes: "
+ analysis + "]")
else:
texts.append(IMAGE_PLACEHOLDER)
elif (isinstance(part, dict) and part.get("type") == "text"
and isinstance(part.get("text"), str)):
texts.append(part["text"])
# andere Part-Typen (z. B. audio) werden verworfen
if had_image:
msg["content"] = "\n\n".join(t for t in texts if t.strip())
return out
def _vision_analyze(image_url: str, question: str) -> str:
"""Sendet Bild + Frage an den Vision-Server, liefert den Analysen-Text."""
payload = {
"model": VISION_ALIAS,
"messages": [
{"role": "system", "content": VISION_ANALYST_PROMPT},
{"role": "user", "content": [
{"type": "image_url", "image_url": {"url": image_url}},
{"type": "text",
"text": ("Aktuelle Benutzerfrage nur als zusätzlicher "
"Kontext, ohne die Analyse darauf zu beschränken: "
+ (question or "(keine Frage)")
+ "\n\nErstelle jetzt eine vollständige, "
"fragenunabhängige Vision-Analyse.")},
]},
],
"max_tokens": VISION_MAX_TOKENS,
"temperature": 0.1,
"stream": False,
}
body = json.dumps(payload).encode()
# timeout=VISION_INFER_TIMEOUT gilt für Connect UND Read. (Nicht
# conn.sock.settimeout() – vor connect() ist conn.sock noch None.)
conn = http.client.HTTPConnection("127.0.0.1", VISION_PORT,
timeout=VISION_INFER_TIMEOUT)
conn.request("POST", "/v1/chat/completions", body=body,
headers={"Content-Type": "application/json"})
resp = conn.getresponse()
raw = resp.read()
conn.close()
try:
data = json.loads(raw)
except ValueError:
raise RuntimeError(
f"Vision-Inferenz: ungültige Antwort (HTTP {resp.status})")
if resp.status != 200:
msg = (data.get("error") or {}).get("message", str(data)) \
if isinstance(data, dict) else str(data)
raise RuntimeError(f"Vision-Inferenz fehlgeschlagen ({resp.status}): {msg}")
choices = data.get("choices") or []
if not choices:
raise RuntimeError("Vision-Inferenz: leere Antwort")
content = (choices[0].get("message") or {}).get("content")
if not isinstance(content, str) or not content.strip():
raise RuntimeError("Vision-Inferenz: leere Analyse")
return content.strip()
def _vision_swap(data: dict) -> str:
"""Kompletter Vision-Hotswap (Q3 "Augen").
Hält den zentralen GPU-Lock (gegenseitiger Ausschluss mit
Profilwechsel und FLUX). Normaler Ablauf (explizit, Schritt für
Schritt – die Wiederherstellung des Hauptprofils erfolgt erst NACH
abgeschlossener Vision-Inferenz):
1. Hauptprofil entladen (VRAM freigeben)
2. Q3-Vision-Server starten und auf Ready warten
3. Bild direkt an 127.0.0.1:VISION_PORT analysieren (Inferenz)
4. Q3-Vision-Server stoppen
5. Hauptprofil wiederherstellen
finally dient NUR als Fehler-/Cleanup-Sicherung (Q3 stoppen,
Hauptprofil retten, Phase zurücksetzen, Timing loggen) – nicht als
normaler Ablauf.
Liefert den Analysen-Text. Die Analyse wird zusätzlich im
Vision-Cache (keyed by Bild-Hash) gespeichert, damit Folgefragen das
Bild nicht erneut durch Q3 schicken (siehe _sanitize_for_main_model).
"""
vis = STATE.vision
profile = current_profile()
if profile is None:
raise RuntimeError("kein aktives Qwen-Profil (override.conf?)")
image_url, question = _extract_last_user_image(data)
if not image_url:
raise RuntimeError("kein Bild im Request")
img_hash = _image_hash(image_url)
safe_image_url = _normalize_vision_image(image_url)
vis.last_profile = profile
vis.last_error = None
t_total = time.monotonic()
timings: dict[str, float] = {}
with STATE.lock:
if vis.phase != "idle":
raise RuntimeError(f"Vision-Analyse läuft ({vis.phase})")
# Qwen wird gestoppt → für Chats nicht verfügbar (die warten).
_set_qwen_unavailable(True)
vision = _VisionServer()
analysis: str | None = None
main_restored = False
try:
_wait_chats_drained()
# 1) Hauptprofil entladen (VRAM freigeben).
vis.phase = "stopping-main"
RUNTIME.save(last_profile=profile, phase=vis.phase)
t0 = time.monotonic()
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)
timings["main_unload"] = time.monotonic() - t0
log.info("Vision: Hauptprofil %s entladen (%.1f s)",
profile, timings["main_unload"])
# 2) Q3-Vision-Server starten und auf Ready warten.
vis.phase = "loading-vision"
t0 = time.monotonic()
vision.start()
vision.wait_ready(time.monotonic() + VISION_LOAD_TIMEOUT)
timings["vision_load"] = time.monotonic() - t0
log.info("Vision: Q3-Vision-Modell geladen (%.1f s)",
timings["vision_load"])
# 3) Bild direkt an den Vision-Server analysieren.
vis.phase = "analyzing"
t0 = time.monotonic()
log.info("Vision: Inferenz gestartet (127.0.0.1:%d)", VISION_PORT)
analysis = _vision_analyze(safe_image_url, question)
timings["vision_infer"] = time.monotonic() - t0
log.info("Vision: Inferenz abgeschlossen (%.1f s, %d Zeichen)",
timings["vision_infer"], len(analysis))
_vision_store_analysis(img_hash, analysis)
log.info("Vision: Analyse im Cache gespeichert (Hash %s…)",
img_hash[:16])
# 4) Q3-Vision-Server stoppen.
vis.phase = "unloading-vision"
t0 = time.monotonic()
vision.stop()
try:
_wait_vram_free(timeout=VISION_UNLOAD_TIMEOUT)
except RuntimeError as e:
log.warning("Vision VRAM-Check: %s (fahre mit Restore fort)", e)
timings["vision_unload"] = time.monotonic() - t0
log.info("Vision: Q3 gestoppt (%.1f s)", timings["vision_unload"])
# 5) Hauptprofil wiederherstellen (erst NACH der Inferenz).
vis.phase = "restoring-main"
t0 = time.monotonic()
log.info("Vision: Hauptprofil %s wird wiederhergestellt", profile)
_restore_qwen(profile)
timings["main_restore"] = time.monotonic() - t0
log.info("Vision: Hauptprofil %s wiederhergestellt (%.1f s)",
profile, timings["main_restore"])
main_restored = True
_set_qwen_unavailable(False)
except Exception as e:
# Fehler-/Cleanup-Pfad: Q3 stoppen, Hauptprofil retten.
log.error("Vision-Fehler in Phase %s: %s", vis.phase, e)
vis.last_error = str(e)
if vision.alive():
try:
vision.stop()
log.info("Vision: Q3 gestoppt (Cleanup nach Fehler)")
except Exception:
log.exception("Vision: Q3-Cleanup fehlgeschlagen")
if not main_restored:
try:
_restore_qwen(profile)
_set_qwen_unavailable(False)
log.info("Vision: Hauptprofil %s wiederhergestellt "
"(Cleanup nach Fehler)", profile)
except Exception as e2:
vis.last_error = (f"Qwen-Wiederherstellung "
f"fehlgeschlagen: {e2}")
log.error("Vision: %s", vis.last_error)
# qwen_unavailable bleibt True (Qwen ist down).
raise
finally:
# NUR Cleanup: Phase zurücksetzen, Timing loggen.
vis.phase = "idle"
vis.last_turnaround = round(time.monotonic() - t_total, 1)
log.info("Vision-Timing: %s | Gesamt %.1f s",
" ".join(f"{k}={v:.1f}s" for k, v in timings.items()),
vis.last_turnaround)
if analysis is None:
# Defensive Absicherung (der except-Pfad wirft immer weiter).
raise RuntimeError(
f"Vision-Analyse fehlgeschlagen: "
f"{vis.last_error or 'unbekannter Fehler'}")
return analysis
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 [])
)
# ---------------------------------------------------------------------------
@@ -1494,8 +1046,6 @@ class Handler(BaseHTTPRequestHandler):
self._speech()
elif path == "/v1/audio/transcriptions" and self.command == "POST":
self._transcribe()
elif path == "/vision/test" and self.command == "POST":
self._vision_test()
elif path == "/images" and self.command == "GET":
self._images_list()
elif path.startswith("/images/") and self.command == "GET":
@@ -1694,13 +1244,6 @@ class Handler(BaseHTTPRequestHandler):
"last_seconds": img.last_seconds,
"last_error": img.last_error,
},
"vision": {
"phase": STATE.vision.phase,
"last_error": STATE.vision.last_error,
"last_profile": STATE.vision.last_profile,
"last_turnaround_seconds": STATE.vision.last_turnaround,
"analysis_cache_size": len(STATE.vision.analysis_cache),
},
"tts": tts_status(),
"stt": stt_status(),
}
@@ -2128,63 +1671,6 @@ class Handler(BaseHTTPRequestHandler):
"model": up.get("model"),
})
# ---------- Interner Vision-Test ----------
def _vision_test(self) -> None:
"""Interner Vision-Test: führt den kompletten Q3-Hotswap durch und
liefert die Analyse + Timing (ohne finale Hauptmodell-Inferenz).
Body: {"image_url": "data:image/png;base64,..." | "http://...",
"question": "optional"}
"""
try:
body = self._read_body()
except ValueError as e:
self._send_error(400, str(e),
"invalid_request_error", "invalid_body")
return
try:
req = json.loads(body) if body else {}
except ValueError:
self._send_error(400, "ungültiges JSON",
"invalid_request_error", "invalid_body")
return
if not isinstance(req, dict):
self._send_error(400, "Body muss ein JSON-Objekt sein",
"invalid_request_error", "invalid_body")
return
image_url = req.get("image_url")
if not isinstance(image_url, str) or not image_url:
self._send_error(400, "image_url fehlt (data-URL oder http-URL)",
"invalid_request_error", "missing_image")
return
question = req.get("question") or ""
# Request bauen, der den Vision-Pfad triggert.
data = {
"model": "qwen-medium",
"messages": [
{"role": "user", "content": [
{"type": "image_url", "image_url": {"url": image_url}},
{"type": "text",
"text": question or "Analysiere das Bild."},
]},
],
}
self.timeout = None # Vision-Swap kann Minuten dauern
try:
analysis = _vision_swap(data)
except (ValueError, RuntimeError) as e:
self._send_error(502, str(e), "server_error", "vision_failed")
return
vis = STATE.vision
self._send_json(200, {
"status": "ok",
"profile": vis.last_profile,
"turnaround_seconds": vis.last_turnaround,
"analysis_chars": len(analysis),
"analysis": analysis,
})
# ---------- Transparentes Forwarding ----------
def _forward(self) -> None:
@@ -2197,8 +1683,8 @@ class Handler(BaseHTTPRequestHandler):
return
# /v1/streams/lookup (Open-WebUI-Stream-Recovery): darf NIEMALS auf
# die Qwen-Wiederherstellung warten (während Vision-Hotswap,
# Profilwechsel oder Image-Job ist Qwen down). Wenn Qwen down ist,
# 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":
@@ -2276,24 +1762,17 @@ class Handler(BaseHTTPRequestHandler):
def _chat_proxy(self, body: bytes | None, data: dict,
profile: str | None) -> None:
"""Vision-Vorbereitung, Profilwahl und Chat-Lease als eine Transaktion."""
"""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)
image_url, _ = _extract_last_user_image(data)
if image_url:
if _vision_cached_analysis(_image_hash(image_url)) is None:
self.timeout = None
_vision_swap(data)
else:
log.info("Vision: Bild bereits analysiert "
"(Cache-Treffer) – kein Hotswap")
if _request_has_image(data):
data = _sanitize_for_main_model(data)
log.info("Vision: finale Hauptmodell-Inferenz gestartet")
data = _normalize_chat_images(data)
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"):
@@ -2425,9 +1904,8 @@ def _startup_reconcile() -> None:
"""Reconcile persisted worker/model state before accepting requests."""
previous = RUNTIME.load()
worker = previous.get("worker")
if worker in {"image", "vision"}:
markers = ([IMAGE_WORKER] if worker == "image" else
[VISION_ALIAS, VISION_MODEL, str(VISION_PORT)])
if worker == "image":
markers = [IMAGE_WORKER]
terminated = terminate_recorded_worker(
previous, markers, lambda msg: log.warning("Recovery: %s", msg))
if terminated: