diff --git a/backend/file_backend.py b/backend/file_backend.py new file mode 100644 index 0000000..a751c0e --- /dev/null +++ b/backend/file_backend.py @@ -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) \ No newline at end of file diff --git a/main.py b/main.py index b27e363..f701912 100644 --- a/main.py +++ b/main.py @@ -1,135 +1,208 @@ +import sys import time +import json import signal import argparse from backend.network_backend import NetworkBackend -from msg.frame_logger import FrameLogger -from transceiver import Transceiver +from backend.serial_backend import SerialBackend +from backend.file_backend import FileBackend, source_id_from_path +from msg.frame_logger import FrameLogger +from transceiver import Transceiver -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 +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 -def handler(signum, frame, be: NetworkBackend, loggers: list): - sig_name = signal.Signals(signum).name - print(f'Signal handler called with signal {sig_name} ({signum})') - be.stop() - be.disconnect() +# --------------------------------------------------------------------------- +# Argument parsing +# --------------------------------------------------------------------------- + +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.disconnect() for lg in loggers: lg.close() + sys.exit(0) -if __name__ == "__main__": - ap = argparse.ArgumentParser() - ap.add_argument('--log-dir', default='./log', help='Directory for frame log files (default: ./log)') - args = ap.parse_args() +# --------------------------------------------------------------------------- +# Entry point +# --------------------------------------------------------------------------- - server_list = [ - {'name': " neo-f9p", 'host': "192.168.22.93", 'port': 8721}, - {'name': "zed-x20p", 'host': "192.168.22.93", 'port': 8731} - ] +if __name__ == '__main__': + args = parse_args() - transceiver = {} - loggers = [] + config = load_config(args.project) if args.project else args_to_config(args) - # TCPIP multi source receiver - mc = NetworkBackend(server_list) + if args.save: + save_config(config, args.save) - # Set the signal handler - signal.signal(signal.SIGINT, lambda signum, frame: handler(signum, frame, mc, loggers)) + sources = config.get('sources', []) + if not sources: + print('No sources specified. Use --network, --serial, --file or --project.') + sys.exit(1) - 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()), + log_dir = config.get('log_dir', './log') + delay = config.get('delay', 1.0) - UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_ACK)), - UbxListener(UbxMsg(UBX_CLASS_ACK, UBX_ID_ACK_NACK)), - ] + backends = [] + loggers = [] + transceivers = {} - # 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 / + # Network sources → one shared NetworkBackend + 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) - # Create transceiver - transceiver[grm] = Transceiver(grm) - mc.register_xcvr(transceiver[grm]) + # Serial sources → one SerialBackend each + 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) - # Create and register frame logger - lg = FrameLogger(grm, args.log_dir) - loggers.append(lg) - transceiver[grm].register_listener(lg) + # File sources → one shared FileBackend (no secondary logging) + file_sources = [s for s in sources if s['type'] == 'file'] + if file_sources: + file_be = FileBackend([s['path'] for s in file_sources], delay=delay) + backends.append(file_be) + for s in file_sources: + name = source_id_from_path(s['path']) + if name: + wire_source(name, file_be, log_dir=None, loggers=loggers, + transceivers=transceivers) - # Create NMEA-Packetizer - nmea_packetizer = NmeaPacketizer() - ubx_packetizer = UbxPacketizer() + signal.signal(signal.SIGINT, + lambda sig, frm: handler(sig, frm, backends, loggers)) - # Register NMEA-Packetizer to transceiver - transceiver[grm].register_listener(nmea_packetizer) - transceiver[grm].register_listener(ubx_packetizer) + for be in backends: + be.connect() + be.start() - # 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() + # Poll loop – only meaningful for network/serial sources + net_names = [s['name'] for s in net_sources] + if net_names: + while True: + try: + for name in net_names: + transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_CLOCK)) + transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SIG)) + transceivers[name].send(frame_create(UBX_CLASS_NAV, UBX_ID_NAV_SAT)) + transceivers[name].send(frame_create(UBX_CLASS_RXM, UBX_ID_RXM_RAW)) + except ValueError: + break + 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() + for be in backends: + be.stop() + be.disconnect() for lg in loggers: lg.close() \ No newline at end of file