From 45e63b6def5b52ff616b6d46aebea3bb893b21b1 Mon Sep 17 00:00:00 2001 From: duffyduck Date: Sat, 11 Jul 2026 11:11:49 +0200 Subject: [PATCH] =?UTF-8?q?feat(local-llm):=20B0=20Consumer=20=E2=80=94=20?= =?UTF-8?q?Bridge-Relay=20+=20Brain-Client=20(spiegelt=20FLUX)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- aria-brain/local_llm.py | 61 ++++++++++++++++++++++++++ bridge/aria_bridge.py | 94 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 155 insertions(+) create mode 100644 aria-brain/local_llm.py diff --git a/aria-brain/local_llm.py b/aria-brain/local_llm.py new file mode 100644 index 0000000..a8b9353 --- /dev/null +++ b/aria-brain/local_llm.py @@ -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 diff --git a/bridge/aria_bridge.py b/bridge/aria_bridge.py index 3b3f167..ef2cfd3 100644 --- a/bridge/aria_bridge.py +++ b/bridge/aria_bridge.py @@ -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"))