From fef6845c037786f63a1f85048bab0c7e18ae8b96 Mon Sep 17 00:00:00 2001 From: MarekWo Date: Sun, 7 Jun 2026 21:10:03 +0200 Subject: [PATCH] fix(connection): self-heal degraded long-lived TCP via in-place reconnect MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Long-lived TCP against the meshcore-proxy can degrade in a way the socket can't see: some commands (set_flood_scope_key with all-zero key) start timing out while RX events and other commands keep working. The 5 s execute() timeout fires with concurrent.futures.TimeoutError() — whose str() is empty — so the UI showed "Could not set region scope (none):" with no error text, and only channels with a mapped region could send because their non-zero scope_key happened to keep working. Two recovery paths: - send_channel_message now detects the timeout case (set_flood_scope_key surfaces timed_out=True) and runs force_reconnect() + one retry before failing. The user sees a brief delay instead of a cryptic error and having to restart the container. - A new _liveness_watcher_loop task runs on the DM event loop and forces a reconnect when no RX event has arrived for HEALTH_STRICT_MAX_RX_STALE_SEC (5 min). /health/strict now also reports rx_stale for TCP (previously serial/USB only), so an external watchdog could act on it too. force_reconnect() runs on the DM loop via run_coroutine_threadsafe with a 20 s cap, a 30 s cooldown to avoid churn under fire, and a _reconnect_lock to prevent concurrent attempts. mc.disconnect() fires DISCONNECTED — _intentional_disconnect tells _on_disconnected to skip its own reconnect loop so the two don't race. Co-Authored-By: Claude Opus 4.7 --- app/device_manager.py | 135 ++++++++++++++++++++++++++++++++++++++++++ app/routes/views.py | 5 +- 2 files changed, 139 insertions(+), 1 deletion(-) diff --git a/app/device_manager.py b/app/device_manager.py index 3c9fe69..95f0b95 100644 --- a/app/device_manager.py +++ b/app/device_manager.py @@ -14,6 +14,7 @@ import logging import struct import threading import time +from concurrent.futures import TimeoutError as FuturesTimeoutError from typing import Optional, Any, Dict, List, Tuple from urllib.parse import urlparse, parse_qs @@ -168,11 +169,21 @@ class DeviceManager: self._retry_context = {} # {dm_id: {attempt, max_attempts, path}} — for _on_ack self._ble_keepalive_task = None # asyncio.Task for BLE keepalive self._ble_permanently_failed = False # True when all reconnect attempts exhausted + self._liveness_task = None # asyncio.Task for periodic rx-stale check # Liveness telemetry for /health/strict and the watchdog self._last_rx_at: float = 0.0 # unix ts of last RX_LOG_DATA / event from device self._consecutive_stats_failures: int = 0 # incremented on get_stats_* / get_bat failures + # In-place reconnect (heals degraded long-lived TCP without container restart) + self._reconnect_lock = threading.Lock() # prevents concurrent force_reconnect calls + self._last_force_reconnect_at: float = 0.0 + self._force_reconnect_cooldown_sec = 30.0 # cap reconnect attempts under fire + # mc.disconnect() fires DISCONNECTED — our handler runs its own + # reconnect loop, which would race with _force_reconnect_async(). + # Flip this before the manual disconnect to suppress the race. + self._intentional_disconnect: bool = False + @property def is_connected(self) -> bool: return self._connected and self.mc is not None @@ -387,6 +398,17 @@ class DeviceManager: self._ble_keepalive_task = asyncio.ensure_future(self._ble_keepalive_loop()) self._ble_permanently_failed = False + # Seed liveness so the rx-stale watcher doesn't reconnect us + # immediately on a quiet mesh. + self._last_rx_at = time.time() + + # Start periodic liveness watcher — heals degraded long-lived + # connections (TCP especially) without waiting for an external + # watchdog. Restart cleanly across reconnects. + if self._liveness_task and not self._liveness_task.done(): + self._liveness_task.cancel() + self._liveness_task = asyncio.ensure_future(self._liveness_watcher_loop()) + except Exception as e: logger.error(f"Device connection failed: {e}") self._connected = False @@ -596,6 +618,11 @@ class DeviceManager: self._ble_keepalive_task.cancel() self._ble_keepalive_task = None + # Cancel liveness watcher + if self._liveness_task and not self._liveness_task.done(): + self._liveness_task.cancel() + self._liveness_task = None + if self.mc and self._loop and self._loop.is_running(): try: future = asyncio.run_coroutine_threadsafe( @@ -1376,6 +1403,10 @@ class DeviceManager: async def _on_disconnected(self, event): """Handle device disconnection with auto-reconnect.""" + if self._intentional_disconnect: + logger.debug("Device disconnected (intentional — force_reconnect in progress)") + self._connected = False + return logger.warning("Device disconnected") self._connected = False @@ -1419,6 +1450,94 @@ class DeviceManager: logger.error("Failed to reconnect after 3 attempts") + async def _force_reconnect_async(self) -> bool: + """Close the current mc cleanly and run _connect() again. Returns True on success.""" + self._intentional_disconnect = True + try: + if self.mc: + try: + await self.mc.disconnect() + except Exception as e: + logger.debug(f"disconnect() during force-reconnect raised: {e}") + self.mc = None + self._connected = False + if self.socketio: + self.socketio.emit('device_status', {'connected': False}, namespace='/chat') + await self._connect() + if self._connected and self.socketio: + self.socketio.emit('device_status', {'connected': True}, namespace='/chat') + return self._connected + except Exception as e: + logger.error(f"Force reconnect failed: {e}") + return False + finally: + self._intentional_disconnect = False + + def force_reconnect(self, timeout: float = 20.0) -> bool: + """Force an in-place reconnect from a sync (Flask) context. + + Used when commands time out while the socket still looks alive (a + degraded long-lived TCP we've seen against the meshcore-proxy). A + cooldown caps how often we'll churn the connection if the device + is truly wedged. Returns True only if the device is connected + again when the call returns. + """ + if not self._loop or not self._loop.is_running(): + return False + if not self._reconnect_lock.acquire(blocking=False): + logger.debug("force_reconnect skipped — another reconnect is in progress") + return False + try: + now = time.time() + since_last = now - self._last_force_reconnect_at + if since_last < self._force_reconnect_cooldown_sec: + logger.info( + f"force_reconnect skipped — cooldown active " + f"({since_last:.1f}s < {self._force_reconnect_cooldown_sec}s)" + ) + return self._connected + self._last_force_reconnect_at = now + logger.warning("Forcing in-place reconnect (degraded connection detected)") + future = asyncio.run_coroutine_threadsafe(self._force_reconnect_async(), self._loop) + try: + return bool(future.result(timeout=timeout)) + except FuturesTimeoutError: + logger.error(f"force_reconnect did not complete within {timeout}s") + return False + finally: + self._reconnect_lock.release() + + async def _liveness_watcher_loop(self): + """Watch for stale RX and self-heal via force_reconnect(). + + The /health/strict endpoint surfaces the same signal for an external + watchdog, but nothing in the host setup actually polls it — so the DM + also acts on the signal itself. BLE is excluded (handled by its own + keepalive + permanent-fail path). + """ + STALE_THRESHOLD_SEC = 300 # mirrors HEALTH_STRICT_MAX_RX_STALE_SEC + CHECK_INTERVAL_SEC = 60 + try: + while True: + await asyncio.sleep(CHECK_INTERVAL_SEC) + if not self._connected or self.config.use_ble: + continue + stale = time.time() - (self._last_rx_at or 0.0) + if self._last_rx_at and stale > STALE_THRESHOLD_SEC: + logger.warning( + f"Liveness watcher: no RX for {int(stale)}s " + f"(> {STALE_THRESHOLD_SEC}s) — forcing reconnect" + ) + # force_reconnect() blocks on the same loop we're on, so + # delegate to a thread to avoid deadlocking ourselves. + await asyncio.get_event_loop().run_in_executor( + None, self.force_reconnect + ) + except asyncio.CancelledError: + raise + except Exception as e: + logger.error(f"Liveness watcher crashed: {e}") + # ================================================================ # Command Methods (sync — called from Flask routes) # ================================================================ @@ -1446,6 +1565,15 @@ class DeviceManager: try: with self._send_lock: scope_res = self.set_flood_scope_key(scope['key_hex'] if scope else None) + # Degraded long-lived TCP: set_flood_scope_key can time out while + # other commands still work. Force one in-place reconnect and + # retry the scope+send pair before surfacing an error. + if not scope_res.get('success') and scope_res.get('timed_out'): + logger.warning( + "set_flood_scope_key timed out — forcing reconnect and retrying send" + ) + if self.force_reconnect(): + scope_res = self.set_flood_scope_key(scope['key_hex'] if scope else None) if not scope_res.get('success'): scope_name = scope['name'] if scope else 'none' return { @@ -3070,6 +3198,10 @@ class DeviceManager: Passing None or empty hex clears the scope (firmware falls back to its default). Used on the channel-send hot path in PR #4. + + Returns ``{'timed_out': True}`` when ``execute()`` raises + ``FuturesTimeoutError`` so callers can recover via force_reconnect() + instead of just surfacing an empty error string to the UI. """ if not self.is_connected: return {'success': False, 'error': 'Device not connected'} @@ -3082,6 +3214,9 @@ class DeviceManager: return {'success': False, 'error': 'Scope key must be 16 bytes (32 hex chars)'} self.execute(self.mc.commands.set_flood_scope(key_bytes), timeout=5) return {'success': True} + except FuturesTimeoutError: + logger.error("set_flood_scope_key timed out — device connection likely degraded") + return {'success': False, 'error': 'Device not responding (timeout)', 'timed_out': True} except Exception as e: logger.error(f"set_flood_scope_key failed: {e}") return {'success': False, 'error': str(e)} diff --git a/app/routes/views.py b/app/routes/views.py index 343217c..290e6cb 100644 --- a/app/routes/views.py +++ b/app/routes/views.py @@ -160,7 +160,10 @@ def health_strict(): transport = getattr(config, 'transport_type', 'serial') last_rx = getattr(dm, '_last_rx_at', 0.0) or 0.0 - if transport in ('serial', 'usb') and last_rx > 0: + # TCP included: long-lived TCP to meshcore-proxy can degrade in ways the + # socket can't detect (commands time out while events still trickle in or + # vice versa). rx_stale is the cheapest external symptom. + if transport in ('serial', 'usb', 'tcp') and last_rx > 0: stale = time.time() - last_rx if stale > HEALTH_STRICT_MAX_RX_STALE_SEC: return jsonify({