Files
NmeaClient/network_backend.py
T
2026-03-25 05:59:58 +01:00

220 lines
5.7 KiB
Python

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
from a_backend import ABackend
from nmea 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, 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 = [
UbxListener(Sig(), 0x01, 0x43),
UbxListener(Sat(), 0x01, 0x35),
UbxListener(RawX(), 0x02, 0x15),
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] = 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()