replay: delay per complete frame instead of per chunk

ReplayPacer (MsgListener+MsgTalker) sits between packetizer output and
downstream listeners; sleeps after each complete NMEA/UBX frame before
forwarding. FileBackend no longer owns the delay — it reads as fast as
possible and the pacer provides back-pressure on the same thread.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
2026-05-25 13:11:40 +02:00
co-authored by Claude Sonnet 4.6
parent ff686b64f7
commit dce570460e
4 changed files with 58 additions and 26 deletions
+2 -14
View File
@@ -1,5 +1,4 @@
import re import re
import time
from threading import Thread from threading import Thread
from backend.a_backend import ABackend from backend.a_backend import ABackend
@@ -15,20 +14,11 @@ def source_id_from_path(path: str) -> str | None:
class FileBackend(ABackend): class FileBackend(ABackend):
def __init__(self, files: list[str], delay: float = 1.0): def __init__(self, files: list[str]):
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} self._entries = [{'path': p, 'source_id': source_id_from_path(p), 'xcvr': None, 'file': None, 'thread': None}
for p in files] for p in files]
self._running = False 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): def register_xcvr(self, xcvr):
for e in self._entries: for e in self._entries:
if e['source_id'] == xcvr.name.strip(): if e['source_id'] == xcvr.name.strip():
@@ -56,7 +46,7 @@ class FileBackend(ABackend):
self._running = False self._running = False
for e in self._entries: for e in self._entries:
if e['thread'] and e['thread'].is_alive(): if e['thread'] and e['thread'].is_alive():
e['thread'].join(timeout=self._delay + 1.0) e['thread'].join(timeout=2.0)
def join(self): def join(self):
for e in self._entries: for e in self._entries:
@@ -78,5 +68,3 @@ class FileBackend(ABackend):
break break
if xcvr: if xcvr:
xcvr.on_recv(MsgContainer(data, xcvr)) xcvr.on_recv(MsgContainer(data, xcvr))
if self._delay > 0:
time.sleep(self._delay)
+13 -3
View File
@@ -153,6 +153,7 @@ class ReceiverManager(QObject):
from listeners import NavSatQtListener, GsaQtListener from listeners import NavSatQtListener, GsaQtListener
from backend.serial_backend import SerialBackend from backend.serial_backend import SerialBackend
from backend.file_backend import FileBackend from backend.file_backend import FileBackend
from msg.replay_pacer import ReplayPacer
config = self._configs.get(rid) config = self._configs.get(rid)
state = self._states.get(rid) state = self._states.get(rid)
@@ -168,7 +169,7 @@ class ReceiverManager(QObject):
elif config.conn_type == 'serial': elif config.conn_type == 'serial':
backend = SerialBackend(config.serial_port, config.baud_rate) backend = SerialBackend(config.serial_port, config.baud_rate)
else: else:
backend = FileBackend([config.file_path], delay=config.delay) backend = FileBackend([config.file_path])
backend.register_xcvr(xcvr) backend.register_xcvr(xcvr)
@@ -177,10 +178,19 @@ class ReceiverManager(QObject):
xcvr.register_listener(nmea_pkt) xcvr.register_listener(nmea_pkt)
xcvr.register_listener(ubx_pkt) xcvr.register_listener(ubx_pkt)
ubx_pkt.register_listener( if config.conn_type == 'file' and config.delay > 0:
nmea_sink = ReplayPacer(config.delay)
ubx_sink = ReplayPacer(config.delay)
nmea_pkt.register_listener(nmea_sink)
ubx_pkt.register_listener(ubx_sink)
else:
nmea_sink = nmea_pkt
ubx_sink = ubx_pkt
ubx_sink.register_listener(
NavSatQtListener(lambda sats, r=rid: self._on_sat(r, sats)) NavSatQtListener(lambda sats, r=rid: self._on_sat(r, sats))
) )
nmea_pkt.register_listener( nmea_sink.register_listener(
GsaQtListener(rid, lambda p, t, r=rid: self._on_pdop(r, p, t)) GsaQtListener(rid, lambda p, t, r=rid: self._on_pdop(r, p, t))
) )
+16 -4
View File
@@ -8,6 +8,7 @@ from backend.network_backend import NetworkBackend
from backend.serial_backend import SerialBackend from backend.serial_backend import SerialBackend
from backend.file_backend import FileBackend, source_id_from_path from backend.file_backend import FileBackend, source_id_from_path
from msg.frame_logger import FrameLogger from msg.frame_logger import FrameLogger
from msg.replay_pacer import ReplayPacer
from transceiver import Transceiver from transceiver import Transceiver
from nmea.packetizer import NmeaPacketizer from nmea.packetizer import NmeaPacketizer
@@ -31,7 +32,7 @@ def parse_args():
help='Save current config as JSON project file') help='Save current config as JSON project file')
ap.add_argument('--log-dir', default='./log', ap.add_argument('--log-dir', default='./log',
help='Frame log directory (default: ./log)') help='Frame log directory (default: ./log)')
ap.add_argument('--delay', type=float, default=1.0, ap.add_argument('--delay', type=float, default=0.1,
help='File replay inter-frame delay in seconds, 010 (default: 1.0)') help='File replay inter-frame delay in seconds, 010 (default: 1.0)')
ap.add_argument('--network', nargs=3, action='append', ap.add_argument('--network', nargs=3, action='append',
metavar=('NAME', 'HOST', 'PORT'), metavar=('NAME', 'HOST', 'PORT'),
@@ -86,7 +87,7 @@ def _make_listeners(name: str):
def wire_source(name: str, backend, log_dir: str | None, def wire_source(name: str, backend, log_dir: str | None,
loggers: list, transceivers: dict): loggers: list, transceivers: dict, pacer_delay: float = 0.0):
xcvr = Transceiver(name) xcvr = Transceiver(name)
backend.register_xcvr(xcvr) backend.register_xcvr(xcvr)
transceivers[name] = xcvr transceivers[name] = xcvr
@@ -102,6 +103,17 @@ def wire_source(name: str, backend, log_dir: str | None,
xcvr.register_listener(ubx_pkt) xcvr.register_listener(ubx_pkt)
nmea_rx, ubx_rx = _make_listeners(name) nmea_rx, ubx_rx = _make_listeners(name)
if pacer_delay > 0:
nmea_pacer = ReplayPacer(pacer_delay)
ubx_pacer = ReplayPacer(pacer_delay)
nmea_pkt.register_listener(nmea_pacer)
ubx_pkt.register_listener(ubx_pacer)
for rx in nmea_rx:
nmea_pacer.register_listener(rx)
for rx in ubx_rx:
ubx_pacer.register_listener(rx)
else:
for rx in nmea_rx: for rx in nmea_rx:
nmea_pkt.register_listener(rx) nmea_pkt.register_listener(rx)
for rx in ubx_rx: for rx in ubx_rx:
@@ -164,13 +176,13 @@ if __name__ == '__main__':
# File sources → one shared FileBackend (no secondary logging) # File sources → one shared FileBackend (no secondary logging)
file_sources = [s for s in sources if s['type'] == 'file'] file_sources = [s for s in sources if s['type'] == 'file']
if file_sources: if file_sources:
file_be = FileBackend([s['path'] for s in file_sources], delay=delay) file_be = FileBackend([s['path'] for s in file_sources])
backends.append(file_be) backends.append(file_be)
for s in file_sources: for s in file_sources:
name = source_id_from_path(s['path']) name = source_id_from_path(s['path'])
if name: if name:
wire_source(name, file_be, log_dir=None, loggers=loggers, wire_source(name, file_be, log_dir=None, loggers=loggers,
transceivers=transceivers) transceivers=transceivers, pacer_delay=delay)
signal.signal(signal.SIGINT, signal.signal(signal.SIGINT,
lambda sig, frm: handler(sig, frm, backends, loggers)) lambda sig, frm: handler(sig, frm, backends, loggers))
+22
View File
@@ -0,0 +1,22 @@
import time
from msg.listener import MsgListener
from msg.talker import MsgTalker
from msg.container import MsgContainer
class ReplayPacer(MsgListener, MsgTalker):
"""Sits between a packetizer and its listeners; sleeps after each
complete frame so file replay runs at a controlled rate."""
def __init__(self, delay: float = 1.0):
MsgListener.__init__(self)
MsgTalker.__init__(self)
self.delay = delay
def is_msg(self, msg: MsgContainer) -> bool:
return True
def on_recv(self, msg: MsgContainer):
if self.delay > 0:
time.sleep(self.delay)
self.call_listener(msg)