Files
NmeaClient/network_backend.py
T
2026-03-23 15:42:24 +01:00

223 lines
5.9 KiB
Python

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_recv(nmea_packetizer)
transceiver[grm].register_recv(ubx_packetizer)
# Register Receiver to NMEA-Packetizer
for rx in nmea_rx_list:
nmea_packetizer.register_recv(rx)
# Register Receiver to UBX-Packetizer
for rx in ubx_rx_list:
ubx_packetizer.register_recv(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()