219 lines
5.7 KiB
Python
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': (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()
|
|
|
|
|