diff --git a/bridge/aria_bridge.py b/bridge/aria_bridge.py index a417aa8..7a935f1 100644 --- a/bridge/aria_bridge.py +++ b/bridge/aria_bridge.py @@ -766,6 +766,11 @@ class ARIABridge: # requestId → Future (sat_devices / sat_result), analog _pending_flux. self._satellites: dict[str, dict] = {} self._pending_sat: dict[str, asyncio.Future] = {} + # Compute-Fleet: GPU-Worker (voxtral/whisper/f5tts/llm) melden sich per + # worker_hello, halten sich per worker_ping frisch. instanceId → + # {service, node, gpus, model, busy, last_seen}. Genutzt fuer die + # Diagnostic-Flotten-Anzeige und (Stage 3) targetInstance-Routing. + self._workers: dict[str, dict] = {} # FLUX-Render-Requests die aktuell auf Antwort der flux-bridge (Gamebox) warten. # requestId → Future mit dem flux_response-Payload (oder None bei Fehler). self._pending_flux: dict[str, asyncio.Future] = {} @@ -3528,6 +3533,40 @@ class ARIABridge: self._satellites[sid]["caps"], self._satellites[sid]["control"]) return + elif msg_type == "worker_hello": + iid = (payload.get("instanceId") or "").strip() + if iid: + prev = self._workers.get(iid, {}) + self._workers[iid] = { + "instanceId": iid, + "service": payload.get("service") or "", + "node": payload.get("node") or "", + "gpus": payload.get("gpus") or "", + "model": payload.get("model") or "", + "busy": bool(prev.get("busy", False)), + "last_seen": time.time(), + } + logger.info("[worker] online: %s (service=%s node=%s gpus=%s model=%s)", + iid, self._workers[iid]["service"], self._workers[iid]["node"], + self._workers[iid]["gpus"] or "?", self._workers[iid]["model"] or "?") + return + + elif msg_type == "worker_ping": + iid = (payload.get("instanceId") or "").strip() + if iid: + w = self._workers.get(iid) + if w is None: + # Ping ohne vorheriges hello (Bridge-Neustart) → Minimal-Eintrag, + # service aus der instanceId ableiten (Form: "service@node"). + svc = iid.split("@", 1)[0] + w = self._workers[iid] = { + "instanceId": iid, "service": svc, "node": "", "gpus": "", + "model": "", "busy": False, "last_seen": 0.0, + } + w["busy"] = bool(payload.get("busy", False)) + w["last_seen"] = time.time() + return + elif msg_type in ("sat_devices", "sat_result"): req_id = payload.get("requestId", "") future = self._pending_sat.get(req_id) @@ -4435,6 +4474,9 @@ class ARIABridge: elif method == "POST" and path == "/internal/satellite-list": # Brain fragt: welche Satelliten/Netze sind online + Capabilities. await _send_response(writer, 200, {"ok": True, "satellites": self._satellite_list()}) + elif method in ("GET", "POST") and path == "/internal/worker-list": + # Diagnostic/Brain fragt: welche Compute-Worker sind online (Flotte). + await _send_response(writer, 200, {"ok": True, "workers": self._worker_list()}) elif method == "POST" and path == "/internal/satellite": # Brain-Tool: Discovery oder Command an einen Satelliten. # body: {op:'discover'|'command', satellite, device?, action?, params?} @@ -4684,6 +4726,44 @@ class ARIABridge: }) return out + # worker_ping kommt alle ~10s; nach 35s ohne Ping gilt ein Worker als offline. + WORKER_OFFLINE_S = 35 + + def _worker_list(self) -> list[dict]: + """Bekannte Compute-Worker (Flotte). online = kuerzlich per Ping gesehen.""" + now = time.time() + out = [] + for w in self._workers.values(): + out.append({ + "instanceId": w["instanceId"], "service": w.get("service") or "", + "node": w.get("node") or "", "gpus": w.get("gpus") or "", + "model": w.get("model") or "", "busy": bool(w.get("busy")), + "online": (now - w.get("last_seen", 0)) < self.WORKER_OFFLINE_S, + }) + return out + + def _pick_worker(self, service: str) -> Optional[str]: + """Waehlt eine online, moeglichst freie Instanz des Diensts (Round-Robin + ueber die freien). Gibt die instanceId oder None (keine online). Fuer + Stage-3-Routing (targetInstance).""" + now = time.time() + online = [w for w in self._workers.values() + if w.get("service") == service + and (now - w.get("last_seen", 0)) < self.WORKER_OFFLINE_S] + if not online: + return None + free = [w for w in online if not w.get("busy")] + pool = free or online # alle busy → trotzdem eine nehmen (least-bad) + # Round-Robin: rotierender Zeiger pro Dienst. + rr = getattr(self, "_worker_rr", None) + if rr is None: + rr = self._worker_rr = {} + idx = rr.get(service, 0) % len(pool) + rr[service] = idx + 1 + chosen = pool[idx] + chosen["busy"] = True # optimistisch, bis der naechste Ping korrigiert + return chosen["instanceId"] + async def _satellite_request(self, op: str, satellite: str = "", device: str = "", action: str = "", params: Optional[dict] = None, diff --git a/diagnostic/index.html b/diagnostic/index.html index 1222c73..bd9d1df 100644 --- a/diagnostic/index.html +++ b/diagnostic/index.html @@ -1307,6 +1307,22 @@
(Lade...)
+ +
+

Compute-Flotte 🖥️

+ +
+
+

+ GPU-Dienste (STT / TTS / lokales LLM), verteilt auf einen oder mehrere + Nodes. Jeder Node startet per COMPOSE_PROFILES nur die + Dienste, die er anbietet. Hier siehst Du, welche Instanz auf welchem + Node laeuft, mit welchem Modell/GPU und ob sie gerade beschaeftigt ist. +

+
+
+
(Lade...)
+
@@ -1958,6 +1974,7 @@ } if (msg.type === 'sat_update') { satellites = msg.satellites || []; renderSatellites(); return; } + if (msg.type === 'worker_update') { workers = msg.workers || []; renderWorkers(); return; } if (msg.type === 'sat_devices') { if (msg.satellite) { satDevices[msg.satellite] = { devices: msg.devices || [], location: msg.location, ts: Date.now() }; } satScanning = null; @@ -4179,9 +4196,50 @@ loadTriggers(); } else if (tab === 'satellites') { requestSatellites(); + requestWorkers(); } } + // ── Compute-Flotte ───────────────────────────────────── + let workers = []; // [{instanceId, service, node, gpus, model, busy, online}] + const WORKER_SVC_META = { + voxtral: { icon: '🎙️', label: 'STT (Voxtral)' }, + whisper: { icon: '🎙️', label: 'STT (Whisper)' }, + f5tts: { icon: '🔊', label: 'TTS (F5-TTS)' }, + llm: { icon: '🧠', label: 'LLM (lokal)' }, + }; + function requestWorkers() { send({ action: 'worker_list' }); } + function renderWorkers() { + const box = document.getElementById('worker-list'); + if (!box) return; + if (!workers.length) { + box.innerHTML = 'Kein Compute-Node verbunden. ' + + 'Starte auf einem GPU-Rechner cd xtts && docker compose up -d ' + + '(mit gesetztem COMPOSE_PROFILES).'; + return; + } + // Nach Node gruppieren. + const byNode = {}; + for (const w of workers) { (byNode[w.node || '?'] = byNode[w.node || '?'] || []).push(w); } + box.innerHTML = Object.keys(byNode).sort().map(node => { + const rows = byNode[node].map(w => { + const meta = WORKER_SVC_META[w.service] || { icon: '⚙️', label: w.service }; + const dot = !w.online ? '#666' : (w.busy ? '#FFB020' : '#3FFF3F'); + const stat = !w.online ? 'offline' : (w.busy ? 'beschaeftigt' : 'frei'); + return '
' + + '' + + '' + meta.icon + ' ' + escapeHtml(meta.label) + '' + + '' + escapeHtml(w.model || '') + '' + + (w.gpus ? 'GPU ' + escapeHtml(w.gpus) + '' : '') + + '' + stat + '' + + '
'; + }).join(''); + return '
' + + '
🖥️ ' + escapeHtml(node) + '
' + + rows + '
'; + }).join(''); + } + // ── Satelliten-Ansicht ───────────────────────────────── let satellites = []; // [{id, location, caps, control, online}] let satDevices = {}; // id → {devices, location, ts} diff --git a/diagnostic/server.js b/diagnostic/server.js index 8567e6b..dfe2b45 100644 --- a/diagnostic/server.js +++ b/diagnostic/server.js @@ -490,6 +490,24 @@ function broadcastSatellites() { broadcast({ type: "sat_update", satellites: satelliteList() }); } +// ── Compute-Fleet: GPU-Worker (voxtral/whisper/f5tts/llm) ────────── +// Worker melden sich per worker_hello + halten sich per worker_ping (busy) frisch. +const workers = new Map(); // instanceId → {instanceId, service, node, gpus, model, busy, last_seen} +const WORKER_OFFLINE_MS = 35000; // ping ~10s; nach 35s ohne Ping = offline + +function workerList() { + const now = Date.now(); + return Array.from(workers.values()).map(w => ({ + instanceId: w.instanceId, service: w.service, node: w.node, + gpus: w.gpus, model: w.model, busy: !!w.busy, + online: (now - (w.last_seen || 0)) < WORKER_OFFLINE_MS, + })); +} + +function broadcastWorkers() { + broadcast({ type: "worker_update", workers: workerList() }); +} + // ── OpenClaw Gateway Verbindung ───────────────────────── async function connectGateway() { @@ -935,6 +953,32 @@ function connectRVS(forcePlain) { if (p.satellite && satellites.has(p.satellite)) satellites.get(p.satellite).last_seen = Date.now(); broadcast({ type: "sat_devices", satellite: p.satellite || "", location: p.location || "", devices: p.devices || [] }); + } else if (msg.type === "worker_hello") { + // Ein Compute-Worker (GPU-Dienst) meldet sich mit seiner Identitaet. + const p = msg.payload || {}; + if (p.instanceId) { + const prev = workers.get(p.instanceId) || {}; + workers.set(p.instanceId, { + instanceId: p.instanceId, service: p.service || "", + node: p.node || "", gpus: p.gpus || "", model: p.model || "", + busy: !!prev.busy, last_seen: Date.now(), + }); + broadcastWorkers(); + } + } else if (msg.type === "worker_ping") { + // Heartbeat eines Workers (traegt busy-Status). + const p = msg.payload || {}; + if (p.instanceId) { + let w = workers.get(p.instanceId); + if (!w) { + const svc = String(p.instanceId).split("@")[0]; + w = { instanceId: p.instanceId, service: svc, node: "", gpus: "", model: "", busy: false, last_seen: 0 }; + workers.set(p.instanceId, w); + } + w.busy = !!p.busy; + w.last_seen = Date.now(); + broadcastWorkers(); + } } else if (msg.type === "agent_activity") { // Bridge meldet "ARIA denkt/schreibt/tool" oder "idle" — an Browser // weiterreichen, damit der Thinking-Indikator im Chat erscheint. @@ -2505,6 +2549,8 @@ wss.on("connection", (ws) => { if (currentDiskStatus) ws.send(JSON.stringify(currentDiskStatus)); // Aktuell bekannte Satelliten mitgeben (RVS replayt sat_hello nicht). ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() })); + // Aktuell bekannte Compute-Worker mitgeben (RVS replayt worker_hello nicht). + ws.send(JSON.stringify({ type: "worker_update", workers: workerList() })); ws.on("message", (raw) => { try { @@ -2536,6 +2582,9 @@ wss.on("connection", (ws) => { } else if (msg.action === "sat_list") { // Browser will die aktuelle Satelliten-Liste. ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() })); + } else if (msg.action === "worker_list") { + // Browser will die aktuelle Compute-Flotte. + ws.send(JSON.stringify({ type: "worker_update", workers: workerList() })); } else if (msg.action === "sat_discover") { // Browser triggert einen Geraete-Scan auf einem Satelliten. sendToRVS_raw({ type: "sat_discover", diff --git a/xtts/f5tts/bridge.py b/xtts/f5tts/bridge.py index 9146d63..b1bc087 100644 --- a/xtts/f5tts/bridge.py +++ b/xtts/f5tts/bridge.py @@ -60,6 +60,15 @@ RVS_TOKEN = os.getenv("RVS_TOKEN", "").strip() # f5ttsCkptFile, f5ttsVocabFile, f5ttsCfgStrength, f5ttsNfeStep). F5TTS_DEVICE = os.getenv("F5TTS_DEVICE", "cuda") # nur Bootstrap +# ── Compute-Fleet: Worker-Identitaet & Registrierung ────────────── +# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping. +NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node" +GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip() +WORKER_SERVICE = "f5tts" +INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}" +WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10")) +_tts_busy = False # True waehrend eine Synthese laeuft (busy-Report im ping) + DEFAULT_F5TTS_MODEL = "F5TTS_v1_Base" DEFAULT_F5TTS_CKPT_FILE = "" # leer = Default-Checkpoint von HF DEFAULT_F5TTS_VOCAB_FILE = "" # leer = Default-Vocab vom Modell @@ -460,13 +469,16 @@ _tts_queue: asyncio.Queue[tuple] = asyncio.Queue() async def _tts_worker(ws, runner: F5Runner) -> None: """Serialisiert Synthesen — GPU kann sonst OOM gehen.""" + global _tts_busy while True: text, voice, request_id, message_id, language, speed = await _tts_queue.get() + _tts_busy = True try: await _do_tts(ws, runner, text, voice, request_id, message_id, language, speed) except Exception: logger.exception("TTS-Worker Fehler") finally: + _tts_busy = False _tts_queue.task_done() @@ -808,6 +820,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None: await _send(ws, "service_status", payload) +async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None: + """Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die + Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch.""" + try: + await _send(ws, "worker_hello", { + "instanceId": INSTANCE_ID, "service": WORKER_SERVICE, + "node": NODE_NAME, "gpus": GPU_IDS, "model": model, + }) + while True: + await asyncio.sleep(WORKER_PING_INTERVAL_S) + busy = bool(busy_fn()) if busy_fn else False + await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy}) + except asyncio.CancelledError: + raise + except Exception: + return # Socket tot → still beenden; run_loop reconnectet + startet neu + + async def run_loop(runner: F5Runner) -> None: use_tls = RVS_TLS retry_s = 2 @@ -855,6 +885,9 @@ async def run_loop(runner: F5Runner) -> None: # TTS-Worker fuer diese Verbindung starten worker = asyncio.create_task(_tts_worker(ws, runner)) + ping_task = asyncio.create_task(_worker_register( + ws, model=runner.model_id, + busy_fn=lambda: _tts_busy or not _tts_queue.empty())) try: async for raw in ws: @@ -958,6 +991,7 @@ async def run_loop(runner: F5Runner) -> None: _last_diag_voice = "" finally: worker.cancel() + ping_task.cancel() try: await worker except asyncio.CancelledError: diff --git a/xtts/llm-adapter/adapter.py b/xtts/llm-adapter/adapter.py index db049bb..bdd157a 100644 --- a/xtts/llm-adapter/adapter.py +++ b/xtts/llm-adapter/adapter.py @@ -47,6 +47,15 @@ RVS_TOKEN = os.getenv("RVS_TOKEN", "").strip() LLAMA_URL = os.getenv("LLAMA_URL", "http://llama:8081").rstrip("/") LLM_MODEL = os.getenv("LLM_MODEL", "qwen3-8b") LLM_TIMEOUT_SEC = float(os.getenv("LLM_TIMEOUT_SEC", "60")) + +# ── Compute-Fleet: Worker-Identitaet & Registrierung ────────────── +# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping. +NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node" +GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip() +WORKER_SERVICE = "llm" +INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}" +WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10")) +_inflight = 0 # laufende llm_requests (busy-Report im ping) # Qwen3 hat Thinking-Mode default AN — dann verbraet es Tokens in einem # -Block und liefert (bei kleinem max_tokens) leeren/abgeschnittenen # content, ausserdem 3x langsamer. ARIAs schnelles Tier will KEIN Grübeln @@ -124,7 +133,34 @@ async def _emit_llm_status(ws, state: str, model: str, **extra) -> None: {"service": "llm", "state": state, "model": model, **extra}) +async def _worker_register(ws) -> None: + """Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die + Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch.""" + try: + await _send(ws, "worker_hello", { + "instanceId": INSTANCE_ID, "service": WORKER_SERVICE, + "node": NODE_NAME, "gpus": GPU_IDS, "model": LLM_MODEL, + }) + while True: + await asyncio.sleep(WORKER_PING_INTERVAL_S) + await _send(ws, "worker_ping", + {"instanceId": INSTANCE_ID, "busy": _inflight > 0}) + except asyncio.CancelledError: + raise + except Exception: + return # Socket tot → still beenden; _run reconnectet + startet neu + + async def _handle_llm_request(ws, payload: dict) -> None: + global _last_model, _inflight + _inflight += 1 + try: + await _do_llm_request(ws, payload) + finally: + _inflight -= 1 + + +async def _do_llm_request(ws, payload: dict) -> None: global _last_model req_id = payload.get("requestId", "") messages = payload.get("messages") or [] @@ -201,6 +237,7 @@ async def _run() -> None: logger.info("RVS verbunden — llm-adapter online") retry_s = 2 tls_fallback_tried = False + ping_task = asyncio.create_task(_worker_register(ws)) async for raw in ws: try: msg = json.loads(raw) @@ -214,6 +251,10 @@ async def _run() -> None: asyncio.create_task(_handle_llm_request(ws, payload)) except Exception as e: logger.warning("RVS-Verbindung verloren/fehlgeschlagen: %s", e) + try: + ping_task.cancel() + except NameError: + pass if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried: tls_fallback_tried = True use_tls = False diff --git a/xtts/voxtral/bridge.py b/xtts/voxtral/bridge.py index cb13048..c3ae726 100644 --- a/xtts/voxtral/bridge.py +++ b/xtts/voxtral/bridge.py @@ -58,6 +58,16 @@ VOXTRAL_MODEL = os.getenv("VOXTRAL_MODEL", "mistralai/Voxtral-Mini-3B-2507") VOXTRAL_LANGUAGE = os.getenv("VOXTRAL_LANGUAGE", "de") VOXTRAL_DEVICE = os.getenv("VOXTRAL_DEVICE", "cuda") +# ── Compute-Fleet: Worker-Identitaet & Registrierung ────────────── +# Jeder Node meldet sich bei der aria-bridge (worker_hello) und haelt die +# Registry per periodischem worker_ping frisch. INSTANCE_ID adressiert diesen +# Worker bei Redundanz (targetInstance-Routing, Stage 3). +NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node" +GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip() +WORKER_SERVICE = "voxtral" +INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}" +WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10")) + STREAM_TRANSCRIBE_INTERVAL_MS = int(os.getenv("STREAM_TRANSCRIBE_INTERVAL_MS", "1000")) STREAM_DEFAULT_ENDPOINT_MS = 2400 STREAM_DEFAULT_HARD_CAP_MS = 300000 @@ -695,6 +705,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None: await _send(ws, "service_status", payload) +async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None: + """Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die + Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch.""" + try: + await _send(ws, "worker_hello", { + "instanceId": INSTANCE_ID, "service": WORKER_SERVICE, + "node": NODE_NAME, "gpus": GPU_IDS, "model": model, + }) + while True: + await asyncio.sleep(WORKER_PING_INTERVAL_S) + busy = bool(busy_fn()) if busy_fn else False + await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy}) + except asyncio.CancelledError: + raise + except Exception: + return # Socket tot → still beenden; run_loop reconnectet + startet neu + + async def run_loop(sessions: SessionManager) -> None: use_tls = RVS_TLS retry_s = 2 @@ -713,6 +741,9 @@ async def run_loop(sessions: SessionManager) -> None: sessions.attach_ws(ws) await _broadcast_status(ws, "ready", model=VOXTRAL_MODEL) await _send(ws, "config_request", {"service": "voxtral"}) + ping_task = asyncio.create_task(_worker_register( + ws, model=VOXTRAL_MODEL, + busy_fn=lambda: bool(sessions._sessions))) async for raw in ws: try: msg = json.loads(raw) @@ -797,6 +828,10 @@ async def run_loop(sessions: SessionManager) -> None: "AN" if SPEAKER_ID_ENABLED else "AUS") except Exception as e: logger.warning("RVS-Verbindung verloren: %s — retry in %ds", e, retry_s) + try: + ping_task.cancel() + except NameError: + pass if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried: use_tls = False tls_fallback_tried = True diff --git a/xtts/whisper/bridge.py b/xtts/whisper/bridge.py index 0cf6e98..7a5d41e 100644 --- a/xtts/whisper/bridge.py +++ b/xtts/whisper/bridge.py @@ -59,6 +59,14 @@ WHISPER_DEVICE = os.getenv("WHISPER_DEVICE", "cuda") WHISPER_COMPUTE_TYPE = os.getenv("WHISPER_COMPUTE_TYPE", "float16") WHISPER_LANGUAGE = os.getenv("WHISPER_LANGUAGE", "de") +# ── Compute-Fleet: Worker-Identitaet & Registrierung ────────────── +# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping. +NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node" +GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip() +WORKER_SERVICE = "whisper" +INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}" +WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10")) + ALLOWED_MODELS = {"tiny", "base", "small", "medium", "large-v3"} # Streaming-Parameter (Defaults — koennen pro Session vom App-Payload ueberschrieben werden) @@ -822,6 +830,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None: await _send(ws, "service_status", payload) +async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None: + """Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die + Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch.""" + try: + await _send(ws, "worker_hello", { + "instanceId": INSTANCE_ID, "service": WORKER_SERVICE, + "node": NODE_NAME, "gpus": GPU_IDS, "model": model, + }) + while True: + await asyncio.sleep(WORKER_PING_INTERVAL_S) + busy = bool(busy_fn()) if busy_fn else False + await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy}) + except asyncio.CancelledError: + raise + except Exception: + return # Socket tot → still beenden; run_loop reconnectet + startet neu + + # ────────────────────────────────────────────────────────────── # WS-LOOP # ────────────────────────────────────────────────────────────── @@ -862,6 +888,9 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None: except Exception as e: logger.exception("Initial-Handshake crashed: %s", e) asyncio.create_task(_initial_handshake()) + ping_task = asyncio.create_task(_worker_register( + ws, model=(runner.model_size or WHISPER_MODEL), + busy_fn=lambda: bool(sessions._sessions))) async for raw in ws: try: @@ -1038,6 +1067,10 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None: except Exception as e: logger.warning("Verbindung verloren: %s", e) sessions.detach_ws() + try: + ping_task.cancel() + except NameError: + pass if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried: logger.info("TLS-Verbindung fehlgeschlagen — Fallback auf ws://") use_tls = False