main: full CLI config with network/serial/file sources and JSON project export
- --network NAME HOST PORT (repeatable) - --serial NAME DEVICE BAUD (repeatable) - --file PATH (repeatable, replayed via FileBackend) - --log-dir DIR (frame logging, default ./log) - --delay SECS (file replay inter-frame delay, default 1.0) - --project FILE load config from JSON project file - --save FILE export current config as JSON project file Multiple source types can be combined. Network sources share one NetworkBackend; serial sources each get their own SerialBackend; file sources share one FileBackend. The poll loop runs only for network sources; file-only mode joins the reader threads; serial-only mode blocks on signal.pause(). file_backend: rename _source_id -> source_id_from_path (public), add join() to wait for reader threads. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,77 @@
|
|||||||
|
import re
|
||||||
|
import time
|
||||||
|
from threading import Thread
|
||||||
|
|
||||||
|
from backend.a_backend import ABackend
|
||||||
|
from msg.container import MsgContainer
|
||||||
|
|
||||||
|
_FILENAME_RE = re.compile(r'\d{4}_\d{2}_\d{2}_\d{2}_\d{2}_\d{2}_(.+)\.log$')
|
||||||
|
_CHUNK = 64
|
||||||
|
|
||||||
|
|
||||||
|
def source_id_from_path(path: str) -> str | None:
|
||||||
|
m = _FILENAME_RE.search(path)
|
||||||
|
return m.group(1) if m else None
|
||||||
|
|
||||||
|
|
||||||
|
class FileBackend(ABackend):
|
||||||
|
def __init__(self, files: list[str], delay: float = 1.0):
|
||||||
|
self._delay = max(0.0, min(10.0, delay))
|
||||||
|
self._entries = [{'path': p, 'source_id': source_id_from_path(p), 'xcvr': None, 'file': None, 'thread': None}
|
||||||
|
for p in files]
|
||||||
|
self._running = False
|
||||||
|
|
||||||
|
@property
|
||||||
|
def delay(self) -> float:
|
||||||
|
return self._delay
|
||||||
|
|
||||||
|
@delay.setter
|
||||||
|
def delay(self, value: float):
|
||||||
|
self._delay = max(0.0, min(10.0, value))
|
||||||
|
|
||||||
|
def register_xcvr(self, xcvr):
|
||||||
|
for e in self._entries:
|
||||||
|
if e['source_id'] == xcvr.name.strip():
|
||||||
|
e['xcvr'] = xcvr
|
||||||
|
xcvr.on_register(self)
|
||||||
|
return
|
||||||
|
raise Exception(f"No log file found for source \"{xcvr.name}\"")
|
||||||
|
|
||||||
|
def connect(self):
|
||||||
|
for e in self._entries:
|
||||||
|
e['file'] = open(e['path'], 'rb')
|
||||||
|
|
||||||
|
def start(self):
|
||||||
|
self._running = True
|
||||||
|
for e in self._entries:
|
||||||
|
e['thread'] = Thread(target=self._run, args=(e,), daemon=True)
|
||||||
|
e['thread'].start()
|
||||||
|
|
||||||
|
def stop(self):
|
||||||
|
self._running = False
|
||||||
|
for e in self._entries:
|
||||||
|
if e['thread'] and e['thread'].is_alive():
|
||||||
|
e['thread'].join(timeout=self._delay + 1.0)
|
||||||
|
|
||||||
|
def join(self):
|
||||||
|
for e in self._entries:
|
||||||
|
if e['thread']:
|
||||||
|
e['thread'].join()
|
||||||
|
|
||||||
|
def disconnect(self):
|
||||||
|
for e in self._entries:
|
||||||
|
if e['file']:
|
||||||
|
e['file'].close()
|
||||||
|
e['file'] = None
|
||||||
|
|
||||||
|
def _run(self, entry: dict):
|
||||||
|
xcvr = entry['xcvr']
|
||||||
|
f = entry['file']
|
||||||
|
while self._running:
|
||||||
|
data = f.read(_CHUNK)
|
||||||
|
if not data:
|
||||||
|
break
|
||||||
|
if xcvr:
|
||||||
|
xcvr.on_recv(MsgContainer(data, xcvr))
|
||||||
|
if self._delay > 0:
|
||||||
|
time.sleep(self._delay)
|
||||||
@@ -1,14 +1,17 @@
|
|||||||
|
import sys
|
||||||
import time
|
import time
|
||||||
|
import json
|
||||||
import signal
|
import signal
|
||||||
import argparse
|
import argparse
|
||||||
|
|
||||||
from backend.network_backend import NetworkBackend
|
from backend.network_backend import NetworkBackend
|
||||||
|
from backend.serial_backend import SerialBackend
|
||||||
|
from backend.file_backend import FileBackend, source_id_from_path
|
||||||
from msg.frame_logger import FrameLogger
|
from msg.frame_logger import FrameLogger
|
||||||
from transceiver import Transceiver
|
from transceiver import Transceiver
|
||||||
|
|
||||||
from nmea.packetizer import NmeaPacketizer
|
from nmea.packetizer import NmeaPacketizer
|
||||||
from nmea.messages import Gll, Gsa, Gsv, Gga, Rmc
|
from nmea.messages import Gll, Gsa, Gsv, Gga, Rmc
|
||||||
|
|
||||||
from ubx.ubx import frame_create
|
from ubx.ubx import frame_create
|
||||||
from ubx.msg_types import *
|
from ubx.msg_types import *
|
||||||
from ubx.packetizer import UbxPacketizer
|
from ubx.packetizer import UbxPacketizer
|
||||||
@@ -16,120 +19,190 @@ from ubx.messages import Sig, Sat, RawX, UbxMsg
|
|||||||
from ubx.listener import UbxListener
|
from ubx.listener import UbxListener
|
||||||
|
|
||||||
|
|
||||||
def handler(signum, frame, be: NetworkBackend, loggers: list):
|
# ---------------------------------------------------------------------------
|
||||||
sig_name = signal.Signals(signum).name
|
# Argument parsing
|
||||||
print(f'Signal handler called with signal {sig_name} ({signum})')
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
def parse_args():
|
||||||
|
ap = argparse.ArgumentParser(description='GNSS receiver client')
|
||||||
|
ap.add_argument('--project', metavar='FILE',
|
||||||
|
help='Load config from JSON project file')
|
||||||
|
ap.add_argument('--save', metavar='FILE',
|
||||||
|
help='Save current config as JSON project file')
|
||||||
|
ap.add_argument('--log-dir', default='./log',
|
||||||
|
help='Frame log directory (default: ./log)')
|
||||||
|
ap.add_argument('--delay', type=float, default=1.0,
|
||||||
|
help='File replay inter-frame delay in seconds, 0–10 (default: 1.0)')
|
||||||
|
ap.add_argument('--network', nargs=3, action='append',
|
||||||
|
metavar=('NAME', 'HOST', 'PORT'),
|
||||||
|
help='Add a network source')
|
||||||
|
ap.add_argument('--serial', nargs=3, action='append',
|
||||||
|
metavar=('NAME', 'DEVICE', 'BAUD'),
|
||||||
|
help='Add a serial source')
|
||||||
|
ap.add_argument('--file', action='append', metavar='PATH',
|
||||||
|
help='Add a log file source for replay')
|
||||||
|
return ap.parse_args()
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Config (de)serialisation
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
def args_to_config(args) -> dict:
|
||||||
|
sources = []
|
||||||
|
for name, host, port in (args.network or []):
|
||||||
|
sources.append({'type': 'network', 'name': name,
|
||||||
|
'host': host, 'port': int(port)})
|
||||||
|
for name, device, baud in (args.serial or []):
|
||||||
|
sources.append({'type': 'serial', 'name': name,
|
||||||
|
'device': device, 'baudrate': int(baud)})
|
||||||
|
for path in (args.file or []):
|
||||||
|
sources.append({'type': 'file', 'path': path})
|
||||||
|
return {'log_dir': args.log_dir, 'delay': args.delay, 'sources': sources}
|
||||||
|
|
||||||
|
|
||||||
|
def load_config(path: str) -> dict:
|
||||||
|
with open(path) as f:
|
||||||
|
return json.load(f)
|
||||||
|
|
||||||
|
|
||||||
|
def save_config(config: dict, path: str):
|
||||||
|
with open(path, 'w') as f:
|
||||||
|
json.dump(config, f, indent=2)
|
||||||
|
print(f'Config saved to {path}')
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Wiring helpers
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
def _make_listeners(name: str):
|
||||||
|
nmea = [Gll(name, 'GN'), Gsa(name, 'GN'), Gsv(name, 'GP'),
|
||||||
|
Gga(name, 'GP'), Rmc(name, 'GN'), Gga(name, 'GN')]
|
||||||
|
ubx = [UbxListener(Sig()), UbxListener(Sat()), UbxListener(RawX()),
|
||||||
|
UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_ACK)),
|
||||||
|
UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_NACK))]
|
||||||
|
return nmea, ubx
|
||||||
|
|
||||||
|
|
||||||
|
def wire_source(name: str, backend, log_dir: str | None,
|
||||||
|
loggers: list, transceivers: dict):
|
||||||
|
xcvr = Transceiver(name)
|
||||||
|
backend.register_xcvr(xcvr)
|
||||||
|
transceivers[name] = xcvr
|
||||||
|
|
||||||
|
if log_dir:
|
||||||
|
lg = FrameLogger(name, log_dir)
|
||||||
|
loggers.append(lg)
|
||||||
|
xcvr.register_listener(lg)
|
||||||
|
|
||||||
|
nmea_pkt = NmeaPacketizer()
|
||||||
|
ubx_pkt = UbxPacketizer()
|
||||||
|
xcvr.register_listener(nmea_pkt)
|
||||||
|
xcvr.register_listener(ubx_pkt)
|
||||||
|
|
||||||
|
nmea_rx, ubx_rx = _make_listeners(name)
|
||||||
|
for rx in nmea_rx:
|
||||||
|
nmea_pkt.register_listener(rx)
|
||||||
|
for rx in ubx_rx:
|
||||||
|
ubx_pkt.register_listener(rx)
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Signal handler
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
def handler(signum, _frame, backends: list, loggers: list):
|
||||||
|
print(f'Signal {signal.Signals(signum).name} – shutting down')
|
||||||
|
for be in backends:
|
||||||
be.stop()
|
be.stop()
|
||||||
be.disconnect()
|
be.disconnect()
|
||||||
for lg in loggers:
|
for lg in loggers:
|
||||||
lg.close()
|
lg.close()
|
||||||
|
sys.exit(0)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
# ---------------------------------------------------------------------------
|
||||||
ap = argparse.ArgumentParser()
|
# Entry point
|
||||||
ap.add_argument('--log-dir', default='./log', help='Directory for frame log files (default: ./log)')
|
# ---------------------------------------------------------------------------
|
||||||
args = ap.parse_args()
|
|
||||||
|
|
||||||
server_list = [
|
if __name__ == '__main__':
|
||||||
{'name': " neo-f9p", 'host': "192.168.22.93", 'port': 8721},
|
args = parse_args()
|
||||||
{'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731}
|
|
||||||
]
|
|
||||||
|
|
||||||
transceiver = {}
|
config = load_config(args.project) if args.project else args_to_config(args)
|
||||||
|
|
||||||
|
if args.save:
|
||||||
|
save_config(config, args.save)
|
||||||
|
|
||||||
|
sources = config.get('sources', [])
|
||||||
|
if not sources:
|
||||||
|
print('No sources specified. Use --network, --serial, --file or --project.')
|
||||||
|
sys.exit(1)
|
||||||
|
|
||||||
|
log_dir = config.get('log_dir', './log')
|
||||||
|
delay = config.get('delay', 1.0)
|
||||||
|
|
||||||
|
backends = []
|
||||||
loggers = []
|
loggers = []
|
||||||
|
transceivers = {}
|
||||||
|
|
||||||
# TCPIP multi source receiver
|
# Network sources → one shared NetworkBackend
|
||||||
mc = NetworkBackend(server_list)
|
net_sources = [s for s in sources if s['type'] == 'network']
|
||||||
|
if net_sources:
|
||||||
|
net_be = NetworkBackend([{'name': s['name'], 'host': s['host'], 'port': s['port']}
|
||||||
|
for s in net_sources])
|
||||||
|
backends.append(net_be)
|
||||||
|
for s in net_sources:
|
||||||
|
wire_source(s['name'], net_be, log_dir, loggers, transceivers)
|
||||||
|
|
||||||
# Set the signal handler
|
# Serial sources → one SerialBackend each
|
||||||
signal.signal(signal.SIGINT, lambda signum, frame: handler(signum, frame, mc, loggers))
|
for s in [s for s in sources if s['type'] == 'serial']:
|
||||||
|
ser_be = SerialBackend(s['device'], s.get('baudrate', 115200))
|
||||||
|
backends.append(ser_be)
|
||||||
|
wire_source(s['name'], ser_be, log_dir, loggers, transceivers)
|
||||||
|
|
||||||
for grm in [" neo-f9p", "zed-x20p"]:
|
# File sources → one shared FileBackend (no secondary logging)
|
||||||
nmea_rx_list = [
|
file_sources = [s for s in sources if s['type'] == 'file']
|
||||||
Gll(grm, 'GN'),
|
if file_sources:
|
||||||
Gsa(grm, 'GN'),
|
file_be = FileBackend([s['path'] for s in file_sources], delay=delay)
|
||||||
Gsv(grm, 'GP'),
|
backends.append(file_be)
|
||||||
Gga(grm, 'GP'),
|
for s in file_sources:
|
||||||
Rmc(grm, 'GN'),
|
name = source_id_from_path(s['path'])
|
||||||
Gga(grm, 'GN'),
|
if name:
|
||||||
]
|
wire_source(name, file_be, log_dir=None, loggers=loggers,
|
||||||
ubx_rx_list = [
|
transceivers=transceivers)
|
||||||
UbxListener(Sig()),
|
|
||||||
UbxListener(Sat()),
|
|
||||||
UbxListener(RawX()),
|
|
||||||
|
|
||||||
UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_ACK)),
|
signal.signal(signal.SIGINT,
|
||||||
UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_NACK)),
|
lambda sig, frm: handler(sig, frm, backends, loggers))
|
||||||
]
|
|
||||||
|
|
||||||
# Topology
|
for be in backends:
|
||||||
# Receiver-0.0 \
|
be.connect()
|
||||||
# Receiver-0.1 <- Packetizer-0 <- Transceiver-0 \
|
be.start()
|
||||||
# / \
|
|
||||||
# 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 and register frame logger
|
|
||||||
lg = FrameLogger(grm, args.log_dir)
|
|
||||||
loggers.append(lg)
|
|
||||||
transceiver[grm].register_listener(lg)
|
|
||||||
|
|
||||||
# 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()
|
|
||||||
|
|
||||||
|
# Poll loop – only meaningful for network/serial sources
|
||||||
|
net_names = [s['name'] for s in net_sources]
|
||||||
|
if net_names:
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
for grm in [" neo-f9p", "zed-x20p"]:
|
for name in net_names:
|
||||||
# Poll UBX-NAV-CLOCK (0x01 0x22)
|
transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_CLOCK))
|
||||||
frame = frame_create(UBX_CLASS_NAV, UBX_ID_CLOCK)
|
transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SIG))
|
||||||
transceiver[grm].send(frame)
|
transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SAT))
|
||||||
# Poll UBX-NAV-SIG (0x01 0x43)
|
transceivers[name].send(frame_create(UBX_CLASS_RXM, UBX_ID_RXM_RAW))
|
||||||
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:
|
except ValueError:
|
||||||
break
|
break
|
||||||
|
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
|
elif file_sources:
|
||||||
|
# Wait until all file threads finish
|
||||||
|
for be in backends:
|
||||||
|
if isinstance(be, FileBackend):
|
||||||
|
be.join()
|
||||||
|
else:
|
||||||
|
# Serial-only: block until interrupted
|
||||||
|
signal.pause()
|
||||||
|
|
||||||
# Stop thread
|
for be in backends:
|
||||||
mc.stop()
|
be.stop()
|
||||||
|
be.disconnect()
|
||||||
# code never reached
|
|
||||||
mc.disconnect()
|
|
||||||
|
|
||||||
for lg in loggers:
|
for lg in loggers:
|
||||||
lg.close()
|
lg.close()
|
||||||
Reference in New Issue
Block a user