Fix: Worker Reconnect Attempt
Check the connection status of reconnect and loop retry drop to start
This commit is contained in:
+23
-17
@@ -282,24 +282,22 @@ class WorkerClient:
|
||||
logger.warning(f"[RECONNECT] Starting reconnection process | sio.connected={self.sio.connected}")
|
||||
reconnect_delay = 1 # Start with 1 second delay
|
||||
max_reconnect_delay = 30 # Maximum delay of 30 seconds
|
||||
|
||||
while self.reconnection_in_progress:
|
||||
|
||||
while self.reconnection_in_progress and self.running:
|
||||
try:
|
||||
logger.info(f"[RECONNECT] Attempt #{int(reconnect_delay)} | delay: {reconnect_delay}s | sio.connected={self.sio.connected}")
|
||||
await asyncio.sleep(reconnect_delay)
|
||||
|
||||
# Try to connect
|
||||
if not self.sio.connected:
|
||||
logger.info(f"[RECONNECT] Calling sio.connect({self.server_url})")
|
||||
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
|
||||
|
||||
if self.sio.connected:
|
||||
logger.info(f"[RECONNECT] ✓ Verified connected | sio.connected={self.sio.connected}")
|
||||
self.reconnection_in_progress = False
|
||||
break
|
||||
else:
|
||||
logger.warning(f"[RECONNECT] Already connected! sio.connected={self.sio.connected}")
|
||||
|
||||
# If we get here, we've successfully reconnected
|
||||
logger.info(f"[RECONNECT] ✓ Success | sio.connected={self.sio.connected}")
|
||||
self.reconnection_in_progress = False
|
||||
# The on_connect handler will send the join request
|
||||
break
|
||||
raise ConnectionError("connect() returned without establishing a connection")
|
||||
except Exception as e:
|
||||
logger.error(f"[RECONNECT] ✗ Failed: {e}")
|
||||
# Increase delay exponentially, up to max_reconnect_delay
|
||||
@@ -514,17 +512,25 @@ class WorkerClient:
|
||||
try:
|
||||
logger.info(f"[START] Worker {self.worker_id} connecting to {self.server_url}...")
|
||||
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
|
||||
|
||||
logger.info(f"[START] Initial connection established | sio.connected={self.sio.connected}")
|
||||
|
||||
# Start the event queue processor if we have an event queue
|
||||
|
||||
|
||||
event_processor = None
|
||||
if self.event_queue:
|
||||
event_processor = asyncio.create_task(self.process_event_queue())
|
||||
await self.sio.wait()
|
||||
event_processor.cancel()
|
||||
|
||||
else:
|
||||
await self.sio.wait()
|
||||
try:
|
||||
while self.running:
|
||||
await self.sio.wait()
|
||||
if not self.running:
|
||||
break
|
||||
logger.info("[START] Disconnected — waiting for reconnect() to bring the session back")
|
||||
while self.reconnection_in_progress and self.running:
|
||||
await asyncio.sleep(1)
|
||||
logger.info(f"[START] Reconnection settled | sio.connected={self.sio.connected}")
|
||||
finally:
|
||||
if event_processor:
|
||||
event_processor.cancel()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[START] Error: {e}")
|
||||
|
||||
Reference in New Issue
Block a user