import socket import selectors import time import signal from threading import Thread from queue import Queue from msg.container import MsgContainer from transceiver import Transceiver, ATransceiver from a_backend import ABackend from nmea.packetizer import NmeaPacketizer from nmea.messages import Gll, Gsa, Gsv, Gga, Rmc from ubx.ubx import frame_create from ubx.packetizer import UbxPacketizer from ubx.messages import Sig, Sat, RawX, UbxMsg from ubx.listener import UbxListener class NetworkBackend(ABackend): def __init__(self, servers: list[dict]): self.sel = None self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port']), 'xcvr': None, 'queue': Queue()} for s in servers] self.thread = Thread(target=self.event_loop) self.loop_enable = True 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, None) 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_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: 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=None) 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) 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: timestamp = time.time() 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) except KeyboardInterrupt: print("Caught keyboard interrupt, exiting") finally: self.sel.close() def handler(signum, frame, be: NetworkBackend): sig_name = signal.Signals(signum).name print(f'Signal handler called with signal {sig_name} ({signum})') be.stop() be.disconnect() if __name__ == "__main__": server_list = [ {'name': " neo-f9p", 'host': "192.168.22.93", 'port': 8721}, {'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731} ] transceiver = {} # TCPIP multi source receiver mc = NetworkBackend(server_list) # Set the signal handler signal.signal(signal.SIGINT, lambda signum, frame: handler(signum, frame, mc)) for grm in [" neo-f9p", "zed-x20p"]: nmea_rx_list = [ Gll(grm, 'GN'), Gsa(grm, 'GN'), Gsv(grm, 'GP'), Gga(grm, 'GP'), Rmc(grm, 'GN'), Gga(grm, 'GN'), ] ubx_rx_list = [ UbxListener(Sig()), UbxListener(Sat()), UbxListener(RawX()), UbxListener(UbxMsg(0x05, 0x01)), UbxListener(UbxMsg(0x05, 0x00)), ] # Topology # Receiver-0.0 \ # Receiver-0.1 <- Packetizer-0 <- Transceiver-0 \ # / \ # Send-0->/ \ # + Backend # Send-1->\ / # \ / # Receiver-1.0 <- Packetizer-1 <- Transceiver-1 / # Receiver-1.1 / # Create transceiver transceiver[grm] = Transceiver(grm) mc.register_xcvr(transceiver[grm]) # Create NMEA-Packetizer nmea_packetizer = NmeaPacketizer() ubx_packetizer = UbxPacketizer() # Register NMEA-Packetizer to transceiver transceiver[grm].register_listener(nmea_packetizer) transceiver[grm].register_listener(ubx_packetizer) # Register Receiver to NMEA-Packetizer for rx in nmea_rx_list: nmea_packetizer.register_listener(rx) # Register Receiver to UBX-Packetizer for rx in ubx_rx_list: ubx_packetizer.register_listener(rx) # Connect all sources mc.connect() # Start thread mc.start() while True: try: for grm in [" neo-f9p", "zed-x20p"]: # Poll UBX-NAV-CLOCK (0x01 0x22) frame = frame_create(0x01, 0x22, b'') transceiver[grm].send(frame) # Poll UBX-NAV-SIG (0x01 0x43) frame = frame_create(0x01, 0x43, b'') transceiver[grm].send(frame) # Poll UBX-NAV-SAT (0x01 0x35) frame = frame_create(0x01, 0x35, b'') transceiver[grm].send(frame) # Poll UBX-RXM-RAWX (0x02 0x15) frame = frame_create(0x02, 0x15, b'') transceiver[grm].send(frame) except ValueError: break time.sleep(1) # Stop thread mc.stop() # code never reached mc.disconnect()