From 348e9783aaf6b9670dac00d8d49c431315b6034c Mon Sep 17 00:00:00 2001 From: duffyduck Date: Sat, 19 Sep 2026 11:10:49 +0200 Subject: [PATCH] =?UTF-8?q?fix(fleet):=20Empfangs-Watchdog=20in=20Workern?= =?UTF-8?q?=20=E2=80=94=20halb-tote=20RVS-Verbindung=20erkennen=20+=20reco?= =?UTF-8?q?nnect?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root Cause der unsichtbaren Box: die WS-Verbindung wird ueber die Zeit halb-tot — Caddy (TLS-Terminator vor dem RVS) haelt sie offen und pongt die WS-Pings selbst, waehrend der rvs-Backend sie laengst fallen liess (rooms:1 = Box nicht mehr Raum-Mitglied). Die Box sendet worker_hello ins Leere; der WS-Ping erkennt den toten Backend nie. (Token/URL/Port/IP/RVS-Instanz alle bewiesen identisch.) Fix: recv() mit Timeout (RX_STALE_S=60s) statt `async for`. In einem echten Raum kommt staendig Broadcast-Traffic rein (sat_hello alle 25s, Brain-Polling) → Timer wird laufend resettet. Bleibt der Traffic aus, ist die Verbindung tot → ConnectionError → Reconnect ueber die bestehende run_loop-Backoff-Logik. Betrifft alle vier Worker (f5tts/whisper/voxtral/llm-adapter). Co-Authored-By: Claude Opus 4.8 --- xtts/f5tts/bridge.py | 11 ++++++++++- xtts/llm-adapter/adapter.py | 11 ++++++++++- xtts/voxtral/bridge.py | 11 ++++++++++- xtts/whisper/bridge.py | 11 ++++++++++- 4 files changed, 40 insertions(+), 4 deletions(-) diff --git a/xtts/f5tts/bridge.py b/xtts/f5tts/bridge.py index 43907e1..137e7aa 100644 --- a/xtts/f5tts/bridge.py +++ b/xtts/f5tts/bridge.py @@ -67,6 +67,10 @@ 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")) +# Empfangs-Watchdog: kommt in RX_STALE_S kein Broadcast rein (ein echter Raum hat +# staendig Traffic, z.B. sat_hello alle 25s / Brain-Polling), gilt die Verbindung +# als halb-tot (Caddy pongt die WS-Pings selbst) -> Zwangs-Reconnect. +RX_STALE_S = int(os.getenv("RX_STALE_S", "60")) _tts_busy = False # True waehrend eine Synthese laeuft (busy-Report im ping) # ── Auslastungs-Monitor (Stage E) ────────────────────────── @@ -917,7 +921,12 @@ async def run_loop(runner: F5Runner) -> None: busy_fn=lambda: _tts_busy or not _tts_queue.empty())) try: - async for raw in ws: + while True: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=RX_STALE_S) + except asyncio.TimeoutError: + logger.warning("Kein RVS-Traffic seit %ds — Verbindung halb-tot, reconnect", RX_STALE_S) + raise ConnectionError("rvs-stale") try: msg = json.loads(raw) except Exception: diff --git a/xtts/llm-adapter/adapter.py b/xtts/llm-adapter/adapter.py index aebb333..208955c 100644 --- a/xtts/llm-adapter/adapter.py +++ b/xtts/llm-adapter/adapter.py @@ -55,6 +55,10 @@ 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")) +# Empfangs-Watchdog: kommt in RX_STALE_S kein Broadcast rein (ein echter Raum hat +# staendig Traffic, z.B. sat_hello alle 25s / Brain-Polling), gilt die Verbindung +# als halb-tot (Caddy pongt die WS-Pings selbst) -> Zwangs-Reconnect. +RX_STALE_S = int(os.getenv("RX_STALE_S", "60")) _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 @@ -424,7 +428,12 @@ async def _run() -> None: retry_s = 2 tls_fallback_tried = False ping_task = asyncio.create_task(_worker_register(ws)) - async for raw in ws: + while True: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=RX_STALE_S) + except asyncio.TimeoutError: + logger.warning("Kein RVS-Traffic seit %ds — Verbindung halb-tot, reconnect", RX_STALE_S) + raise ConnectionError("rvs-stale") try: msg = json.loads(raw) except Exception: diff --git a/xtts/voxtral/bridge.py b/xtts/voxtral/bridge.py index f3e2051..99bc92e 100644 --- a/xtts/voxtral/bridge.py +++ b/xtts/voxtral/bridge.py @@ -67,6 +67,10 @@ 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")) +# Empfangs-Watchdog: kommt in RX_STALE_S kein Broadcast rein (ein echter Raum hat +# staendig Traffic, z.B. sat_hello alle 25s / Brain-Polling), gilt die Verbindung +# als halb-tot (Caddy pongt die WS-Pings selbst) -> Zwangs-Reconnect. +RX_STALE_S = int(os.getenv("RX_STALE_S", "60")) # ── Auslastungs-Monitor (Stage E) ────────────────────────── import node_stats @@ -753,7 +757,12 @@ async def run_loop(sessions: SessionManager) -> None: ping_task = asyncio.create_task(_worker_register( ws, model=VOXTRAL_MODEL, busy_fn=lambda: bool(sessions._sessions))) - async for raw in ws: + while True: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=RX_STALE_S) + except asyncio.TimeoutError: + logger.warning("Kein RVS-Traffic seit %ds — Verbindung halb-tot, reconnect", RX_STALE_S) + raise ConnectionError("rvs-stale") try: msg = json.loads(raw) except Exception: diff --git a/xtts/whisper/bridge.py b/xtts/whisper/bridge.py index 4247152..42e254a 100644 --- a/xtts/whisper/bridge.py +++ b/xtts/whisper/bridge.py @@ -66,6 +66,10 @@ 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")) +# Empfangs-Watchdog: kommt in RX_STALE_S kein Broadcast rein (ein echter Raum hat +# staendig Traffic, z.B. sat_hello alle 25s / Brain-Polling), gilt die Verbindung +# als halb-tot (Caddy pongt die WS-Pings selbst) -> Zwangs-Reconnect. +RX_STALE_S = int(os.getenv("RX_STALE_S", "60")) # ── Auslastungs-Monitor (Stage E) ────────────────────────── import node_stats @@ -903,7 +907,12 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None: ws, model=(runner.model_size or WHISPER_MODEL), busy_fn=lambda: bool(sessions._sessions))) - async for raw in ws: + while True: + try: + raw = await asyncio.wait_for(ws.recv(), timeout=RX_STALE_S) + except asyncio.TimeoutError: + logger.warning("Kein RVS-Traffic seit %ds — Verbindung halb-tot, reconnect", RX_STALE_S) + raise ConnectionError("rvs-stale") try: msg = json.loads(raw) except Exception: