feat(local-llm): B0 Consumer — Bridge-Relay + Brain-Client (spiegelt FLUX)
Brain → HTTP /internal/local-llm → Bridge → RVS → llm-adapter → llama.cpp,
1:1 nach dem FLUX-Roundtrip-Muster gebaut:
- Bridge: _pending_llm (requestId→Future), llm_response-Handler (setzt Future),
_local_llm() (sendet llm_request, wartet mit 30s-Timeout), HTTP-Route
POST /internal/local-llm ({messages, max_tokens?, temperature?, stop?}).
- Brain: local_llm.py mit local_llm_chat() — POSTet an die Bridge, gibt
{ok, content, model?, elapsedMs?} zurueck, wirft nie (Aufrufer eskaliert
bei ok=false auf Claude).
Provider-Kette bereits verifiziert (718ms). Als naechstes: Test brain→bridge→
gamebox end-to-end, dann B1 (Router-Heuristik + Escalation) und der
Diagnostic-Testchat/Status (B0.5).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,61 @@
|
||||
"""
|
||||
Local-LLM-Client (Plan B) — Brain-Seite.
|
||||
|
||||
Ruft das schnelle lokale LLM (Qwen3 auf der Gamebox) ueber die Bridge:
|
||||
Brain → HTTP /internal/local-llm → Bridge → RVS → llm-adapter → llama.cpp
|
||||
|
||||
Analog zum Claude-`proxy_client`, nur ueber die Bridge (die ist der RVS-Client;
|
||||
das Brain bleibt HTTP-only). Der Router im Brain (B1) entscheidet, welche Turns
|
||||
hierher gehen (einfach) und welche an Claude (schwer / Tool-Bedarf).
|
||||
|
||||
Rueckgabe von local_llm_chat: {ok, content, model?, elapsedMs?} oder {ok:False, error}.
|
||||
Nie werfen — der Aufrufer entscheidet bei ok=False, ob er auf Claude eskaliert.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
BRIDGE_URL = os.environ.get("BRIDGE_URL", "http://aria-bridge:8090")
|
||||
# Etwas ueber dem Bridge-seitigen _LLM_TIMEOUT_S (30s), damit der HTTP-Call nicht
|
||||
# vor dem eigentlichen LLM-Timeout abbricht.
|
||||
LOCAL_LLM_HTTP_TIMEOUT_SEC = float(os.environ.get("LOCAL_LLM_HTTP_TIMEOUT_SEC", "35"))
|
||||
|
||||
|
||||
def local_llm_chat(messages: list, *, max_tokens: int = 512,
|
||||
temperature: float = 0.7, stop=None) -> dict:
|
||||
"""Ein Chat-Call ans lokale LLM. messages = [{role, content}, ...].
|
||||
Blockierend (urllib) — im Brain laeuft chat() ohnehin im Executor-Thread."""
|
||||
if not isinstance(messages, list) or not messages:
|
||||
return {"ok": False, "error": "messages leer/ungueltig"}
|
||||
req = {"messages": messages, "max_tokens": max_tokens, "temperature": temperature}
|
||||
if stop:
|
||||
req["stop"] = stop
|
||||
try:
|
||||
body = json.dumps(req).encode("utf-8")
|
||||
http_req = urllib.request.Request(
|
||||
f"{BRIDGE_URL}/internal/local-llm", data=body, method="POST",
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
with urllib.request.urlopen(http_req, timeout=LOCAL_LLM_HTTP_TIMEOUT_SEC) as resp:
|
||||
result = json.loads(resp.read().decode("utf-8", "ignore"))
|
||||
except urllib.error.HTTPError as exc:
|
||||
try:
|
||||
err_data = json.loads(exc.read().decode("utf-8", "ignore"))
|
||||
err = err_data.get("error") or str(exc)
|
||||
except Exception:
|
||||
err = str(exc)
|
||||
return {"ok": False, "error": f"local-llm: {err}"}
|
||||
except Exception as exc:
|
||||
logger.warning("local_llm_chat HTTP-Call fehlgeschlagen: %s", exc)
|
||||
return {"ok": False, "error": f"local-llm nicht erreichbar ({exc})"}
|
||||
|
||||
if not isinstance(result, dict) or not result.get("ok"):
|
||||
return {"ok": False, "error": (result or {}).get("error", "unbekannt")}
|
||||
return result
|
||||
@@ -643,6 +643,11 @@ class ARIABridge:
|
||||
# flux-bridge service_status: True wenn ready. Render-Timeouts werden
|
||||
# bei 'loading' deutlich grosszuegiger gesetzt (Modell-Download ~24 GB).
|
||||
self._remote_flux_ready: bool = False
|
||||
# Lokales LLM (Plan B): requestId → Future mit dem llm_response-Payload.
|
||||
# Analog zu _pending_flux — Brain ruft /internal/local-llm, wir relayen
|
||||
# llm_request via RVS an den llm-adapter (Gamebox) und warten auf
|
||||
# llm_response.
|
||||
self._pending_llm: dict[str, asyncio.Future] = {}
|
||||
# User-Message-Counter fuer Auto-Compact. Bei zu langer Konversation
|
||||
# sprengt die argv-Liste beim Claude-Subprocess-Spawn (E2BIG). Bei
|
||||
# COMPACT_AFTER erreicht → Sessions reset + Container restart.
|
||||
@@ -2963,6 +2968,15 @@ class ARIABridge:
|
||||
future.set_result(payload)
|
||||
return
|
||||
|
||||
elif msg_type == "llm_response":
|
||||
# Antwort des llm-adapter (Gamebox) auf unseren llm_request.
|
||||
request_id = payload.get("requestId", "")
|
||||
future = self._pending_llm.get(request_id)
|
||||
if future is None or future.done():
|
||||
return
|
||||
future.set_result(payload)
|
||||
return
|
||||
|
||||
elif msg_type == "service_status":
|
||||
# Gamebox-Bridges (whisper / f5tts / flux) melden ihren Lade-Status.
|
||||
# Wir nutzen das fuer den dynamischen STT-Timeout: solange whisper
|
||||
@@ -3346,6 +3360,59 @@ class ARIABridge:
|
||||
_FLUX_TIMEOUT_READY_S = 240.0 # 4 min nach erstem Render
|
||||
_FLUX_TIMEOUT_LOADING_S = 900.0 # 15 min beim allerersten Mal (Modell-Download)
|
||||
|
||||
# ── Local-LLM-Roundtrip: Brain → Bridge → RVS → llm-adapter → zurueck ──
|
||||
# Qwen3 auf der Gamebox antwortet auf kurze Turns in <1 s. Grosszuegiger
|
||||
# Timeout deckt Kaltstart / laengere Antworten / Netz-Jitter (Gamebox@home)
|
||||
# ab. Bei Timeout faellt der Router im Brain per Escalation auf Claude.
|
||||
_LLM_TIMEOUT_S = 30.0
|
||||
|
||||
async def _local_llm(self, messages: list, max_tokens: int = 512,
|
||||
temperature: float = 0.7, stop=None) -> dict:
|
||||
"""Schickt einen llm_request an den llm-adapter (Gamebox), wartet auf
|
||||
llm_response. Rueckgabe: {ok, content, model, elapsedMs} oder {ok:False, error}."""
|
||||
if self.ws_rvs is None:
|
||||
return {"ok": False, "error": "RVS-Verbindung nicht aktiv"}
|
||||
if not isinstance(messages, list) or not messages:
|
||||
return {"ok": False, "error": "messages leer/ungueltig"}
|
||||
|
||||
request_id = str(uuid.uuid4())
|
||||
loop = asyncio.get_event_loop()
|
||||
future: asyncio.Future = loop.create_future()
|
||||
self._pending_llm[request_id] = future
|
||||
try:
|
||||
req_payload = {
|
||||
"requestId": request_id,
|
||||
"messages": messages,
|
||||
"max_tokens": max_tokens,
|
||||
"temperature": temperature,
|
||||
}
|
||||
if stop:
|
||||
req_payload["stop"] = stop
|
||||
logger.info("[rvs] llm_request → llm-adapter (id=%s, msgs=%d, max_tokens=%d)",
|
||||
request_id[:8], len(messages), max_tokens)
|
||||
ok = await self._send_to_rvs({
|
||||
"type": "llm_request",
|
||||
"payload": req_payload,
|
||||
"timestamp": int(time.time() * 1000),
|
||||
})
|
||||
if not ok:
|
||||
return {"ok": False, "error": "llm_request konnte nicht gesendet werden"}
|
||||
try:
|
||||
result = await asyncio.wait_for(future, timeout=self._LLM_TIMEOUT_S)
|
||||
except asyncio.TimeoutError:
|
||||
return {"ok": False, "error": f"Timeout ({self._LLM_TIMEOUT_S:.0f}s) — Gamebox nicht erreichbar?"}
|
||||
if not isinstance(result, dict) or not result.get("ok"):
|
||||
err = (result or {}).get("error") if isinstance(result, dict) else "leeres Resultat"
|
||||
return {"ok": False, "error": err or "llm-adapter Fehler"}
|
||||
return {
|
||||
"ok": True,
|
||||
"content": result.get("content", ""),
|
||||
"model": result.get("model"),
|
||||
"elapsedMs": result.get("elapsedMs"),
|
||||
}
|
||||
finally:
|
||||
self._pending_llm.pop(request_id, None)
|
||||
|
||||
async def _flux_generate(self, prompt: str, width: int, height: int,
|
||||
steps: Optional[int], guidance: Optional[float],
|
||||
seed: Optional[int], model: Optional[str] = None) -> dict:
|
||||
@@ -3797,6 +3864,33 @@ class ARIABridge:
|
||||
)
|
||||
status = 200 if result.get("ok") else 502
|
||||
await _send_response(writer, status, result)
|
||||
elif method == "POST" and path == "/internal/local-llm":
|
||||
# Vom Brain (Router / Testchat) gefeuert. Wir relayen den
|
||||
# Chat-Request via RVS an den llm-adapter (Gamebox Qwen3),
|
||||
# warten synchron auf llm_response und geben content zurueck.
|
||||
try:
|
||||
data = json.loads(body.decode("utf-8", "ignore"))
|
||||
except Exception as exc:
|
||||
await _send_response(writer, 400, {"error": f"bad json: {exc}"})
|
||||
return
|
||||
messages = data.get("messages")
|
||||
if not isinstance(messages, list) or not messages:
|
||||
await _send_response(writer, 400, {"error": "messages (nicht-leere Liste) erforderlich"})
|
||||
return
|
||||
try:
|
||||
max_tokens = int(data.get("max_tokens") or 512)
|
||||
except (TypeError, ValueError):
|
||||
max_tokens = 512
|
||||
try:
|
||||
temperature = float(data.get("temperature"))
|
||||
except (TypeError, ValueError):
|
||||
temperature = 0.7
|
||||
result = await self._local_llm(
|
||||
messages=messages, max_tokens=max_tokens,
|
||||
temperature=temperature, stop=data.get("stop"),
|
||||
)
|
||||
status = 200 if result.get("ok") else 502
|
||||
await _send_response(writer, status, result)
|
||||
elif method == "POST" and path == "/internal/delete-chat-message":
|
||||
try:
|
||||
data = json.loads(body.decode("utf-8", "ignore"))
|
||||
|
||||
Reference in New Issue
Block a user