diff --git a/const.py b/const.py index 0d78da4..c710385 100644 --- a/const.py +++ b/const.py @@ -55,7 +55,7 @@ REGISTER_START_ADDRESS_3 = 0x8C # Register 140 (decimal) REGISTER_MONITOR_COUNT_3 = 23 # Number of registers to monitor (140-162) # Availability timeout (if no data for 5 minutes, mark unavailable) -AVAILABILITY_TIMEOUT = timedelta(minutes=5) +AVAILABILITY_TIMEOUT = timedelta(minutes=1) # Update throttling (prevent overwhelming HA with high-frequency Modbus traffic) COORDINATOR_UPDATE_INTERVAL = timedelta(seconds=5) # Minimum time between coordinator updates diff --git a/hub.py b/hub.py index 76ddb73..2882935 100644 --- a/hub.py +++ b/hub.py @@ -212,96 +212,58 @@ class ModbusRTUMonitorHub: return True async def _connect(self) -> None: - """Establish TCP connection.""" - self._reader, self._writer = await asyncio.wait_for( - asyncio.open_connection(self.host, self.port), - timeout=10.0, - ) - _LOGGER.info( - "Connected to Modbus RTU gateway at %s:%s", self.host, self.port - ) - - async def _reconnect(self) -> None: - """Reconnect to Modbus gateway after connection loss.""" - # Close existing connection and null out references + """Connect or reconnect to Modbus gateway.""" + # Close existing connection if present if self._writer: try: self._writer.close() await self._writer.wait_closed() except Exception as ex: - _LOGGER.debug("Error closing writer during reconnect: %s", ex) + _LOGGER.debug("Error closing writer: %s", ex) - # Always null out reader/writer to ensure clean state + # Null out reader/writer self._reader = None self._writer = None - # Wait before reconnecting (30 seconds to prevent DDOS during server outages) - await asyncio.sleep(30) - - # Attempt reconnection + # Attempt connection try: - await self._connect() - _LOGGER.info("Reconnected to Modbus gateway at %s:%s", self.host, self.port) + self._reader, self._writer = await asyncio.wait_for( + asyncio.open_connection(self.host, self.port), + timeout=30.0, + ) + _LOGGER.info("Connected to Modbus gateway at %s:%s", self.host, self.port) except Exception as ex: - _LOGGER.error( - "Reconnection failed for %s:%s - %s (will retry in 30 seconds)", + _LOGGER.warning( + "Connection failed for %s:%s - %s (will retry)", self.host, self.port, ex, ) - # Reader/writer are already None, loop will detect this and retry + # Reader/writer are already None, will retry on next iteration async def _monitor_loop(self) -> None: """Main monitoring loop for passive frame reception.""" buffer = bytearray() - first_connection_attempt = True - consecutive_timeouts = 0 - TIMEOUT_THRESHOLD = 30 # 30 timeouts = 30 seconds of no data while True: try: # Check if reader is valid before attempting to read if self._reader is None: - consecutive_timeouts = 0 # Reset counter after reconnect attempt - if first_connection_attempt: - # First connection attempt - try immediately without delay - _LOGGER.info("Attempting initial connection to %s:%s", self.host, self.port) - first_connection_attempt = False - try: - await self._connect() - _LOGGER.info("Connected to Modbus gateway at %s:%s", self.host, self.port) - except Exception as ex: - _LOGGER.warning( - "Initial connection failed to %s:%s - %s (will retry)", - self.host, - self.port, - ex, - ) - # Fall through to reconnect logic with delay - await asyncio.sleep(30) - else: - # Subsequent reconnection attempts - use delay to prevent DDOS - _LOGGER.info("No active connection to %s:%s, attempting reconnect", self.host, self.port) - await self._reconnect() - # Clear buffer for fresh start after connection attempt + _LOGGER.info("Attempting connection to %s:%s", self.host, self.port) + await self._connect() buffer.clear() continue - # Read with timeout to allow graceful shutdown - data = await asyncio.wait_for(self._reader.read(1024), timeout=1.0) + # Read with 30-second timeout (includes reconnection delay) + data = await asyncio.wait_for(self._reader.read(1024), timeout=30.0) if not data: # Connection closed gracefully - _LOGGER.warning("Connection closed to %s:%s, attempting reconnect", self.host, self.port) - consecutive_timeouts = 0 # Reset counter - await self._reconnect() - # Clear buffer for fresh start after reconnection + _LOGGER.warning("Connection closed to %s:%s, reconnecting", self.host, self.port) + await self._connect() buffer.clear() continue - # Successfully received data - reset timeout counter - consecutive_timeouts = 0 - buffer.extend(data) # Process frames from buffer (limit iterations to prevent blocking) @@ -336,26 +298,14 @@ class ModbusRTUMonitorHub: await asyncio.sleep(FRAME_PROCESSING_DELAY) except asyncio.TimeoutError: - # Track consecutive timeouts to detect dead connections - consecutive_timeouts += 1 - - if consecutive_timeouts >= TIMEOUT_THRESHOLD: - # Connection appears dead (30+ seconds of no data) - _LOGGER.warning( - "Connection timeout to %s:%s (%d consecutive timeouts, %ds with no data) - assuming connection dead, will reconnect", - self.host, - self.port, - consecutive_timeouts, - consecutive_timeouts, - ) - # Null reader/writer to trigger reconnection on next iteration - self._reader = None - self._writer = None - consecutive_timeouts = 0 - # Clear buffer for fresh start - buffer.clear() - - # Continue to next iteration (either to reconnect or wait for next timeout) + # No data for 30 seconds - connection dead, reconnect directly + _LOGGER.warning( + "Read timeout after 30 seconds on %s:%s - connection dead, reconnecting", + self.host, + self.port, + ) + await self._connect() + buffer.clear() continue except asyncio.CancelledError: break @@ -370,14 +320,12 @@ class ModbusRTUMonitorHub: # Null out connection and trigger reconnect self._reader = None self._writer = None - consecutive_timeouts = 0 # Reset counter - await self._reconnect() + await self._connect() buffer.clear() continue except Exception as ex: # Other unexpected errors _LOGGER.error("Unexpected error in monitor loop for %s:%s - %s", self.host, self.port, ex) - consecutive_timeouts = 0 # Reset counter on unexpected errors await asyncio.sleep(5) async def _handle_frame(self, frame: ModbusRTUFrame) -> None: @@ -691,7 +639,7 @@ class ModbusRTUMonitorHub: """Periodically check slave availability.""" while True: try: - await asyncio.sleep(60) # Check every minute + await asyncio.sleep(15) # Check every 15 s self._check_slave_availability() except asyncio.CancelledError: break