import socket import selectors import time import signal, os from msg_container import MsgContainer from receiver.nmea_packetizer import NmeaPacketizer from receiver.ubx_packetizer import UbxPacketizer from receiver.ubx_receiver import UbxReceiver from transceiver import Transceiver from a_backend import ABackend from nmea.GLL import Gll from nmea.GSA import Gsa from nmea.GSV import Gsv from nmea.GGA import Gga from nmea.RMC import Rmc from threading import Thread from queue import Queue from ubx.ubx import frame_create from ubx.msg import UbxMsg from ubx.messages.nav import Sig, Sat from ubx.messages.rxm import RawX 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, name: str, data: bytes): sock = self.find_by_name(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 create_xcvr(self, name: str) -> Transceiver: obj = Transceiver(self, name) if not self.register_xcvr(obj): raise Exception(f"No socket found for \"{name}\"") return obj self.sel.register(name, selectors.EVENT_READ, data=None) 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 break return found 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(name, timestamp, data)) 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 = [ UbxReceiver(Sig(), 0x01, 0x43), UbxReceiver(Sat(), 0x01, 0x35), UbxReceiver(RawX(), 0x02, 0x15), UbxReceiver(UbxMsg(), 0x05, 0x01), UbxReceiver(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] = mc.create_xcvr(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()