Fixing reconnect issues after outage by simplifying
This commit is contained in:
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user