network_backend: fix crash on lost/refused connection
getpeername() raised ENOTCONN when a connect failed or the peer closed the socket, crashing the event loop thread. Identify sockets via the selector key's data (set at register time) instead, and catch OSError on recv/send to cleanly close and unregister a dead socket rather than crashing. Also fix select(timeout=100) which waited 100 seconds instead of 100ms, causing stop() to hang after a dead socket was unregistered. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+20
-26
@@ -20,7 +20,7 @@ class NetworkBackend(ABackend):
|
||||
sock = self.find_by_name(xcvr.name)
|
||||
sock['queue'].put(data)
|
||||
if sock is not None:
|
||||
self.sel.modify(sock['sock'], selectors.EVENT_READ + selectors.EVENT_WRITE, None)
|
||||
self.sel.modify(sock['sock'], selectors.EVENT_READ + selectors.EVENT_WRITE, sock)
|
||||
|
||||
def start(self):
|
||||
self.thread.start()
|
||||
@@ -42,12 +42,6 @@ class NetworkBackend(ABackend):
|
||||
if not found:
|
||||
raise Exception(f"No socket found for \"{xcvr.name}\"")
|
||||
|
||||
def find_by_addr(self, addr):
|
||||
for sock in self.sock_list:
|
||||
if sock['addr'] == addr:
|
||||
return sock
|
||||
return None
|
||||
|
||||
def find_by_name(self, name: str):
|
||||
for sock in self.sock_list:
|
||||
if sock['name'] == name:
|
||||
@@ -66,35 +60,35 @@ class NetworkBackend(ABackend):
|
||||
events = selectors.EVENT_READ
|
||||
for sock in self.sock_list:
|
||||
print(f"Starting connection {sock['name']} to {sock['addr']}")
|
||||
self.sel.register(sock['sock'], events, data=None)
|
||||
self.sel.register(sock['sock'], events, data=sock)
|
||||
sock['sock'].connect_ex(sock['addr'])
|
||||
|
||||
def disconnect(self):
|
||||
for sock in self.sock_list:
|
||||
sock['sock'].close()
|
||||
|
||||
def on_send(self, name: str, data: bytes):
|
||||
sock = self.find_by_name(name)
|
||||
if sock is not None:
|
||||
self.sel.modify(sock['sock'], selectors.EVENT_READ + selectors.EVENT_WRITE, data)
|
||||
|
||||
def event_loop(self):
|
||||
try:
|
||||
while self.loop_enable:
|
||||
events = self.sel.select(timeout=100)
|
||||
events = self.sel.select(timeout=0.1)
|
||||
for key, mask in events:
|
||||
descr = self.find_by_addr(key.fileobj.getpeername())
|
||||
name = descr['name']
|
||||
sock = self.find_by_name(name)
|
||||
if mask & selectors.EVENT_READ:
|
||||
data = key.fileobj.recv(64)
|
||||
sock['xcvr'].on_recv(MsgContainer(data, sock['xcvr']))
|
||||
if mask & selectors.EVENT_WRITE:
|
||||
data = sock['queue'].get()
|
||||
key.fileobj.send(data)
|
||||
print(f"Send: {sock['name']}: {data}")
|
||||
if sock['queue'].empty():
|
||||
self.sel.modify(sock['sock'], selectors.EVENT_READ, None)
|
||||
sock = key.data
|
||||
try:
|
||||
if mask & selectors.EVENT_READ:
|
||||
data = key.fileobj.recv(64)
|
||||
if not data:
|
||||
raise OSError("connection closed by peer")
|
||||
sock['xcvr'].on_recv(MsgContainer(data, sock['xcvr']))
|
||||
if mask & selectors.EVENT_WRITE:
|
||||
data = sock['queue'].get()
|
||||
key.fileobj.send(data)
|
||||
print(f"Send: {sock['name']}: {data}")
|
||||
if sock['queue'].empty():
|
||||
self.sel.modify(sock['sock'], selectors.EVENT_READ, sock)
|
||||
except OSError as exc:
|
||||
print(f"Connection {sock['name']} to {sock['addr']} lost: {exc}")
|
||||
self.sel.unregister(key.fileobj)
|
||||
key.fileobj.close()
|
||||
|
||||
except KeyboardInterrupt:
|
||||
print("Caught keyboard interrupt, exiting")
|
||||
|
||||
Reference in New Issue
Block a user