diff --git a/backend/network_backend.py b/backend/network_backend.py index 47f9369..c5e608e 100644 --- a/backend/network_backend.py +++ b/backend/network_backend.py @@ -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")