← back to Exo
health check udp discovered peers before adding them
428bb6063ab06426d3ca9be5b4329bfc2af6a546 · 2024-09-24 19:11:45 +0100 · Alex Cheema
Files touched
M exo/networking/udp_discovery.py
Diff
commit 428bb6063ab06426d3ca9be5b4329bfc2af6a546
Author: Alex Cheema <alexcheema123@gmail.com>
Date: Tue Sep 24 19:11:45 2024 +0100
health check udp discovered peers before adding them
---
exo/networking/udp_discovery.py | 12 ++++++++++--
1 file changed, 10 insertions(+), 2 deletions(-)
diff --git a/exo/networking/udp_discovery.py b/exo/networking/udp_discovery.py
index bba84868..9c454f4b 100644
--- a/exo/networking/udp_discovery.py
+++ b/exo/networking/udp_discovery.py
@@ -53,7 +53,7 @@ class UDPDiscovery(Discovery):
self.broadcast_interval = broadcast_interval
self.discovery_timeout = discovery_timeout
self.device_capabilities = device_capabilities
- self.known_peers: Dict[str, Tuple[PeerHandle, float, float, bool]] = {}
+ self.known_peers: Dict[str, Tuple[PeerHandle, float, float]] = {}
self.broadcast_task = None
self.listen_task = None
self.cleanup_task = None
@@ -141,9 +141,17 @@ class UDPDiscovery(Discovery):
device_capabilities = DeviceCapabilities(**message["device_capabilities"])
if peer_id not in self.known_peers or self.known_peers[peer_id][0].addr() != f"{peer_host}:{peer_port}":
+ new_peer_handle = self.create_peer_handle(peer_id, f"{peer_host}:{peer_port}", device_capabilities)
+ if not await new_peer_handle.health_check():
+ if DEBUG >= 1: print(f"Peer {peer_id} at {peer_host}:{peer_port} is not healthy. Skipping.")
+ return
if DEBUG >= 1: print(f"Adding {peer_id=} at {peer_host}:{peer_port}. Replace existing peer_id: {peer_id in self.known_peers}")
- self.known_peers[peer_id] = (self.create_peer_handle(peer_id, f"{peer_host}:{peer_port}", device_capabilities), time.time(), time.time())
+ self.known_peers[peer_id] = (new_peer_handle, time.time(), time.time())
else:
+ if not await self.known_peers[peer_id][0].health_check():
+ if DEBUG >= 1: print(f"Peer {peer_id} at {peer_host}:{peer_port} is not healthy. Removing.")
+ if peer_id in self.known_peers: del self.known_peers[peer_id]
+ return
self.known_peers[peer_id] = (self.known_peers[peer_id][0], self.known_peers[peer_id][1], time.time())
async def task_listen_for_peers(self):
← da06fb3c check before removing
·
back to Exo
·
change the message we search for in ci 115f0eac →