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(), 'display_name': s['name']} for s in servers] self.thread = Thread(target=self.event_loop) self.loop_enable = True self.on_disconnect = on_disconnect def set_display_name(self, name: str, display_name: str): sock = self.find_by_name(name) if sock is not None: sock['display_name'] = display_name 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['display_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['display_name']}: {data}") if sock['queue'].empty(): self.sel.modify(sock['sock'], selectors.EVENT_READ, sock) except OSError as exc: print(f"Connection {sock['display_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['display_name']}: error processing data, ignoring: {exc}") except KeyboardInterrupt: print("Caught keyboard interrupt, exiting") finally: self.sel.close()