A struct size-mismatch or any other parse error in a listener could propagate uncaught out of the socket-reading thread and silently kill it. Catch broadly in the event loop and per-listener in call_listener, and make the NAV-SAT/NAV-SIG/RXM-RAWX group parsers tolerate a payload length that isn't an exact multiple of the group size instead of crashing on a short struct.unpack. Also wire NetworkBackend's socket-loss detection through to ReceiverManager (via a queued signal, since the callback fires from the backend thread) so the GUI reflects a peer-initiated disconnect instead of leaving the row stuck showing "connected". Verified against a real ZED-X20P: it drops the TCP session on its own after ~20-28s regardless of traffic; the GUI now marks the receiver disconnected and frees its ID for reconnecting. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01AzmCaNjDb3TAqPpTwunKTY
101 lines
2.9 KiB
Python
101 lines
2.9 KiB
Python
import socket
|
|
import selectors
|
|
from threading import Thread
|
|
from queue import Queue
|
|
|
|
from msg.container import MsgContainer
|
|
from transceiver import Transceiver, ATransceiver
|
|
|
|
from backend.a_backend import ABackend
|
|
|
|
|
|
class NetworkBackend(ABackend):
|
|
def __init__(self, servers: list[dict], on_disconnect=None):
|
|
self.sel = None
|
|
self.sock_list = [{'name': s['name'], 'addr': (socket.gethostbyname(s['host']), s['port']), 'xcvr': None, 'queue': Queue()} for s in servers]
|
|
self.thread = Thread(target=self.event_loop)
|
|
self.loop_enable = True
|
|
self.on_disconnect = on_disconnect
|
|
|
|
def send(self, xcvr: ATransceiver, data: bytes):
|
|
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, sock)
|
|
|
|
def start(self):
|
|
self.thread.start()
|
|
|
|
def stop(self):
|
|
if self.loop_enable:
|
|
self.loop_enable = False
|
|
self.thread.join()
|
|
|
|
def register_xcvr(self, xcvr: Transceiver):
|
|
found = False
|
|
for sock in self.sock_list:
|
|
if xcvr.name in sock['name']:
|
|
found = True
|
|
sock['xcvr'] = xcvr
|
|
xcvr.on_register(self)
|
|
break
|
|
|
|
if not found:
|
|
raise Exception(f"No socket found for \"{xcvr.name}\"")
|
|
|
|
def find_by_name(self, name: str):
|
|
for sock in self.sock_list:
|
|
if sock['name'] == name:
|
|
return sock
|
|
return None
|
|
|
|
def _prepare(self):
|
|
for server in self.sock_list:
|
|
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
sock.setblocking(False)
|
|
server['sock'] = sock
|
|
|
|
def connect(self):
|
|
self._prepare()
|
|
self.sel = selectors.DefaultSelector()
|
|
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=sock)
|
|
sock['sock'].connect_ex(sock['addr'])
|
|
|
|
def disconnect(self):
|
|
for sock in self.sock_list:
|
|
sock['sock'].close()
|
|
|
|
def event_loop(self):
|
|
try:
|
|
while self.loop_enable:
|
|
events = self.sel.select(timeout=0.1)
|
|
for key, mask in events:
|
|
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()
|
|
if self.on_disconnect:
|
|
self.on_disconnect(sock['name'])
|
|
except Exception as exc:
|
|
print(f"Connection {sock['name']}: error processing data, ignoring: {exc}")
|
|
|
|
except KeyboardInterrupt:
|
|
print("Caught keyboard interrupt, exiting")
|
|
finally:
|
|
self.sel.close() |