Files
NmeaClient/network_backend.py
T

219 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, 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.msg_types import *
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': (socket.gethostbyname(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(UBX_CLASS_ACK, UBX_ID_ACK_ACK)),
UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_NACK)),
]
# 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(UBX_CLASS_NAV, UBX_ID_CLOCK)
transceiver[grm].send(frame)
# Poll UBX-NAV-SIG (0x01 0x43)
frame = frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SIG)
transceiver[grm].send(frame)
# Poll UBX-NAV-SAT (0x01 0x35)
frame = frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SAT)
transceiver[grm].send(frame)
# Poll UBX-RXM-RAWX (0x02 0x15)
frame = frame_create(UBX_CLASS_RXM, UBX_ID_RXM_RAW)
transceiver[grm].send(frame)
except ValueError:
break
time.sleep(1)
# Stop thread
mc.stop()
# code never reached
mc.disconnect()