Add opt-in diff mode for network transmission (--diff / Diff mode checkbox)
Server (collect.py --diff): - History burst: first record sent full, subsequent as jay_diff_full diffs - Live broadcast: each snapshot diffed against the previous one - Diff payloads tagged with "_diff": true Client (client.py --diff): - Reconstructs full state via jay_merge_full before printing GUI (Connector tab "Diff mode" checkbox): - TcpReader.connect_to() accepts diff_mode; maintains _state dict - Merges incoming diffs before emitting to Dashboard/Plot - Checkbox state persisted to plot_settings.json Both sides default to off; must be enabled consistently on server and client. utils/__init__.py added so utils.jay_diff is importable as a package. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -9,6 +9,7 @@ Usage
|
|||||||
-----
|
-----
|
||||||
python client.py
|
python client.py
|
||||||
python client.py --host 192.168.1.10 --port 9999
|
python client.py --host 192.168.1.10 --port 9999
|
||||||
|
python client.py --diff # reconstruct from diff stream
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
@@ -16,6 +17,8 @@ import json
|
|||||||
import socket
|
import socket
|
||||||
import sys
|
import sys
|
||||||
|
|
||||||
|
from utils.jay_diff import jay_merge_full
|
||||||
|
|
||||||
|
|
||||||
def main() -> None:
|
def main() -> None:
|
||||||
p = argparse.ArgumentParser(
|
p = argparse.ArgumentParser(
|
||||||
@@ -24,6 +27,7 @@ def main() -> None:
|
|||||||
)
|
)
|
||||||
p.add_argument("--host", default="127.0.0.1", help="Server host (default: 127.0.0.1)")
|
p.add_argument("--host", default="127.0.0.1", help="Server host (default: 127.0.0.1)")
|
||||||
p.add_argument("--port", type=int, default=9999, help="Server port (default: 9999)")
|
p.add_argument("--port", type=int, default=9999, help="Server port (default: 9999)")
|
||||||
|
p.add_argument("--diff", action="store_true", help="Reconstruct full snapshots from diff stream")
|
||||||
args = p.parse_args()
|
args = p.parse_args()
|
||||||
|
|
||||||
print(f"Connecting to {args.host}:{args.port} …", flush=True)
|
print(f"Connecting to {args.host}:{args.port} …", flush=True)
|
||||||
@@ -36,6 +40,7 @@ def main() -> None:
|
|||||||
print("Connected. Waiting for data (Ctrl-C to quit).\n", flush=True)
|
print("Connected. Waiting for data (Ctrl-C to quit).\n", flush=True)
|
||||||
buf = ""
|
buf = ""
|
||||||
count = 0
|
count = 0
|
||||||
|
state: dict = {}
|
||||||
try:
|
try:
|
||||||
with sock:
|
with sock:
|
||||||
while True:
|
while True:
|
||||||
@@ -51,8 +56,17 @@ def main() -> None:
|
|||||||
continue
|
continue
|
||||||
try:
|
try:
|
||||||
data = json.loads(line)
|
data = json.loads(line)
|
||||||
|
if args.diff:
|
||||||
|
if data.get("_diff"):
|
||||||
|
diff = {k: v for k, v in data.items() if k != "_diff"}
|
||||||
|
state = jay_merge_full(state, diff)
|
||||||
|
else:
|
||||||
|
state = data
|
||||||
|
display = state
|
||||||
|
else:
|
||||||
|
display = data
|
||||||
count += 1
|
count += 1
|
||||||
print(json.dumps(data, indent=2))
|
print(json.dumps(display, indent=2))
|
||||||
print(f"--- #{count} ", "-" * 50, flush=True)
|
print(f"--- #{count} ", "-" * 50, flush=True)
|
||||||
except json.JSONDecodeError as e:
|
except json.JSONDecodeError as e:
|
||||||
print(f"[invalid JSON] {e}: {line!r}", file=sys.stderr)
|
print(f"[invalid JSON] {e}: {line!r}", file=sys.stderr)
|
||||||
|
|||||||
+15
-1
@@ -55,6 +55,7 @@ from pathlib import Path
|
|||||||
from server import log_config, network, we_connect
|
from server import log_config, network, we_connect
|
||||||
from server.data_model import ALL_DOMAINS, collect_snapshot
|
from server.data_model import ALL_DOMAINS, collect_snapshot
|
||||||
from server.storage_helpers import auto_out_path, load_last_24h_records, load_store, save_store
|
from server.storage_helpers import auto_out_path, load_last_24h_records, load_store, save_store
|
||||||
|
from utils.jay_diff import jay_diff_full
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -92,6 +93,7 @@ def run(args: argparse.Namespace) -> None:
|
|||||||
out = auto_out_path(logs_dir, vin)
|
out = auto_out_path(logs_dir, vin)
|
||||||
|
|
||||||
push_clients, push_lock, push_store_ref = network.start_push_server(args.host, args.port)
|
push_clients, push_lock, push_store_ref = network.start_push_server(args.host, args.port)
|
||||||
|
push_store_ref["diff_mode"] = args.diff
|
||||||
push_store_ref["preloaded"] = load_last_24h_records(logs_dir, vin)
|
push_store_ref["preloaded"] = load_last_24h_records(logs_dir, vin)
|
||||||
if push_store_ref["preloaded"]:
|
if push_store_ref["preloaded"]:
|
||||||
log.info("Pre-loaded %d record(s) from the last 24 h", len(push_store_ref["preloaded"]))
|
log.info("Pre-loaded %d record(s) from the last 24 h", len(push_store_ref["preloaded"]))
|
||||||
@@ -99,6 +101,7 @@ def run(args: argparse.Namespace) -> None:
|
|||||||
current_day = datetime.now(timezone.utc).date()
|
current_day = datetime.now(timezone.utc).date()
|
||||||
store = load_store(out, vin, args.interval)
|
store = load_store(out, vin, args.interval)
|
||||||
push_store_ref["current"] = store
|
push_store_ref["current"] = store
|
||||||
|
last_broadcast: dict = {}
|
||||||
|
|
||||||
if args.dry:
|
if args.dry:
|
||||||
log.info("Push server running on %s:%d — no live collection (Ctrl-C to stop)", args.host, args.port)
|
log.info("Push server running on %s:%d — no live collection (Ctrl-C to stop)", args.host, args.port)
|
||||||
@@ -126,7 +129,13 @@ def run(args: argparse.Namespace) -> None:
|
|||||||
snapshot = collect_snapshot(vehicle, domains)
|
snapshot = collect_snapshot(vehicle, domains)
|
||||||
store["records"].append(snapshot)
|
store["records"].append(snapshot)
|
||||||
save_store(out, store, args.max_records)
|
save_store(out, store, args.max_records)
|
||||||
network.broadcast(snapshot, push_clients, push_lock)
|
if args.diff and last_broadcast:
|
||||||
|
payload = jay_diff_full(last_broadcast, snapshot, combine_upd_add=True)
|
||||||
|
payload["_diff"] = True
|
||||||
|
else:
|
||||||
|
payload = snapshot
|
||||||
|
network.broadcast(payload, push_clients, push_lock)
|
||||||
|
last_broadcast = snapshot
|
||||||
log.info("Record #%d saved at %s", len(store["records"]), snapshot["ts"])
|
log.info("Record #%d saved at %s", len(store["records"]), snapshot["ts"])
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
raise
|
raise
|
||||||
@@ -184,6 +193,11 @@ def main() -> None:
|
|||||||
help="Keep only the last N records per file (default: unlimited)",
|
help="Keep only the last N records per file (default: unlimited)",
|
||||||
)
|
)
|
||||||
p.add_argument("-v", "--verbose", action="store_true", help="Enable debug logging")
|
p.add_argument("-v", "--verbose", action="store_true", help="Enable debug logging")
|
||||||
|
p.add_argument(
|
||||||
|
"--diff",
|
||||||
|
action="store_true",
|
||||||
|
help="Send diffs instead of full snapshots over the push server (default: off)",
|
||||||
|
)
|
||||||
p.add_argument(
|
p.add_argument(
|
||||||
"--dry",
|
"--dry",
|
||||||
action="store_true",
|
action="store_true",
|
||||||
|
|||||||
+34
-10
@@ -23,7 +23,7 @@ from pathlib import Path
|
|||||||
from PyQt5.QtCore import QObject, QTimer, Qt, pyqtSignal
|
from PyQt5.QtCore import QObject, QTimer, Qt, pyqtSignal
|
||||||
from PyQt5.QtGui import QColor, QFont
|
from PyQt5.QtGui import QColor, QFont
|
||||||
from PyQt5.QtWidgets import (
|
from PyQt5.QtWidgets import (
|
||||||
QApplication, QComboBox, QFormLayout, QGridLayout, QGroupBox,
|
QApplication, QCheckBox, QComboBox, QFormLayout, QGridLayout, QGroupBox,
|
||||||
QHBoxLayout, QLabel, QLineEdit, QMainWindow, QPushButton, QScrollArea,
|
QHBoxLayout, QLabel, QLineEdit, QMainWindow, QPushButton, QScrollArea,
|
||||||
QSizePolicy, QSpinBox, QStatusBar, QTabWidget, QTreeWidget, QTreeWidgetItem,
|
QSizePolicy, QSpinBox, QStatusBar, QTabWidget, QTreeWidget, QTreeWidgetItem,
|
||||||
QVBoxLayout, QWidget,
|
QVBoxLayout, QWidget,
|
||||||
@@ -31,6 +31,8 @@ from PyQt5.QtWidgets import (
|
|||||||
|
|
||||||
import pyqtgraph as pg
|
import pyqtgraph as pg
|
||||||
|
|
||||||
|
from utils.jay_diff import jay_merge_full
|
||||||
|
|
||||||
pg.setConfigOption("background", "#12121e")
|
pg.setConfigOption("background", "#12121e")
|
||||||
pg.setConfigOption("foreground", "#aaaacc")
|
pg.setConfigOption("foreground", "#aaaacc")
|
||||||
|
|
||||||
@@ -114,17 +116,21 @@ class TcpReader(QObject):
|
|||||||
|
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
super().__init__()
|
super().__init__()
|
||||||
self._sock = None
|
self._sock = None
|
||||||
self._running = False
|
self._running = False
|
||||||
|
self._diff_mode = False
|
||||||
|
self._state: dict = {}
|
||||||
|
|
||||||
def connect_to(self, host: str, port: int) -> str | None:
|
def connect_to(self, host: str, port: int, diff_mode: bool = False) -> str | None:
|
||||||
"""Open connection; returns error string or None on success."""
|
"""Open connection; returns error string or None on success."""
|
||||||
try:
|
try:
|
||||||
self._sock = socket.create_connection((host, port), timeout=5)
|
self._sock = socket.create_connection((host, port), timeout=5)
|
||||||
self._sock.settimeout(None) # blocking reads after connect
|
self._sock.settimeout(None)
|
||||||
except OSError as exc:
|
except OSError as exc:
|
||||||
return str(exc)
|
return str(exc)
|
||||||
self._running = True
|
self._diff_mode = diff_mode
|
||||||
|
self._state = {}
|
||||||
|
self._running = True
|
||||||
threading.Thread(target=self._read_loop, daemon=True).start()
|
threading.Thread(target=self._read_loop, daemon=True).start()
|
||||||
return None
|
return None
|
||||||
|
|
||||||
@@ -150,7 +156,16 @@ class TcpReader(QObject):
|
|||||||
line = line.strip()
|
line = line.strip()
|
||||||
if line:
|
if line:
|
||||||
try:
|
try:
|
||||||
self.message.emit(json.loads(line))
|
data = json.loads(line)
|
||||||
|
if self._diff_mode:
|
||||||
|
if data.get("_diff"):
|
||||||
|
diff = {k: v for k, v in data.items() if k != "_diff"}
|
||||||
|
self._state = jay_merge_full(self._state, diff)
|
||||||
|
else:
|
||||||
|
self._state = data
|
||||||
|
self.message.emit(dict(self._state))
|
||||||
|
else:
|
||||||
|
self.message.emit(data)
|
||||||
except json.JSONDecodeError:
|
except json.JSONDecodeError:
|
||||||
pass
|
pass
|
||||||
except OSError:
|
except OSError:
|
||||||
@@ -198,9 +213,11 @@ class ConnectorTab(QWidget):
|
|||||||
self._port = QSpinBox()
|
self._port = QSpinBox()
|
||||||
self._port.setRange(1, 65535)
|
self._port.setRange(1, 65535)
|
||||||
self._port.setValue(9999)
|
self._port.setValue(9999)
|
||||||
|
self._diff_cb = QCheckBox("Diff mode")
|
||||||
self._load_conn_settings()
|
self._load_conn_settings()
|
||||||
form.addRow("Host:", self._host)
|
form.addRow("Host:", self._host)
|
||||||
form.addRow("Port:", self._port)
|
form.addRow("Port:", self._port)
|
||||||
|
form.addRow(self._diff_cb)
|
||||||
|
|
||||||
btn_row = QHBoxLayout()
|
btn_row = QHBoxLayout()
|
||||||
self._btn = QPushButton("Connect")
|
self._btn = QPushButton("Connect")
|
||||||
@@ -223,6 +240,10 @@ class ConnectorTab(QWidget):
|
|||||||
else:
|
else:
|
||||||
self.disconnect_requested.emit()
|
self.disconnect_requested.emit()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def diff_mode(self) -> bool:
|
||||||
|
return self._diff_cb.isChecked()
|
||||||
|
|
||||||
def _save_conn_settings(self):
|
def _save_conn_settings(self):
|
||||||
try:
|
try:
|
||||||
_SETTINGS_PATH.parent.mkdir(parents=True, exist_ok=True)
|
_SETTINGS_PATH.parent.mkdir(parents=True, exist_ok=True)
|
||||||
@@ -231,8 +252,9 @@ class ConnectorTab(QWidget):
|
|||||||
data = json.loads(_SETTINGS_PATH.read_text())
|
data = json.loads(_SETTINGS_PATH.read_text())
|
||||||
except (OSError, json.JSONDecodeError):
|
except (OSError, json.JSONDecodeError):
|
||||||
pass
|
pass
|
||||||
data["host"] = self._host.text()
|
data["host"] = self._host.text()
|
||||||
data["port"] = self._port.value()
|
data["port"] = self._port.value()
|
||||||
|
data["diff_mode"] = self._diff_cb.isChecked()
|
||||||
_SETTINGS_PATH.write_text(json.dumps(data, indent=2))
|
_SETTINGS_PATH.write_text(json.dumps(data, indent=2))
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
@@ -244,6 +266,8 @@ class ConnectorTab(QWidget):
|
|||||||
self._host.setText(data["host"])
|
self._host.setText(data["host"])
|
||||||
if "port" in data:
|
if "port" in data:
|
||||||
self._port.setValue(int(data["port"]))
|
self._port.setValue(int(data["port"]))
|
||||||
|
if "diff_mode" in data:
|
||||||
|
self._diff_cb.setChecked(bool(data["diff_mode"]))
|
||||||
except (OSError, json.JSONDecodeError, KeyError, ValueError):
|
except (OSError, json.JSONDecodeError, KeyError, ValueError):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
@@ -692,7 +716,7 @@ class MainWindow(QMainWindow):
|
|||||||
|
|
||||||
def _connect(self, host: str, port: int):
|
def _connect(self, host: str, port: int):
|
||||||
self._plot.clear_data()
|
self._plot.clear_data()
|
||||||
err = self._reader.connect_to(host, port)
|
err = self._reader.connect_to(host, port, self._connector.diff_mode)
|
||||||
if err:
|
if err:
|
||||||
self._connector.set_disconnected(err)
|
self._connector.set_disconnected(err)
|
||||||
self._statusbar.showMessage(f"Connection failed: {err}")
|
self._statusbar.showMessage(f"Connection failed: {err}")
|
||||||
|
|||||||
+16
-3
@@ -3,6 +3,8 @@ import logging
|
|||||||
import socket
|
import socket
|
||||||
import threading
|
import threading
|
||||||
|
|
||||||
|
from utils.jay_diff import jay_diff_full
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
@@ -29,8 +31,19 @@ def _accept_loop(
|
|||||||
records = in_mem
|
records = in_mem
|
||||||
|
|
||||||
try:
|
try:
|
||||||
for record in records:
|
if store_ref.get("diff_mode") and len(records) > 1:
|
||||||
conn.sendall((json.dumps(record, default=str) + "\n").encode())
|
lines = [json.dumps(records[0], default=str) + "\n"]
|
||||||
|
prev = records[0]
|
||||||
|
for record in records[1:]:
|
||||||
|
diff = jay_diff_full(prev, record, combine_upd_add=True)
|
||||||
|
diff["_diff"] = True
|
||||||
|
lines.append(json.dumps(diff, default=str) + "\n")
|
||||||
|
prev = record
|
||||||
|
for line in lines:
|
||||||
|
conn.sendall(line.encode())
|
||||||
|
else:
|
||||||
|
for record in records:
|
||||||
|
conn.sendall((json.dumps(record, default=str) + "\n").encode())
|
||||||
if records:
|
if records:
|
||||||
log.info("Sent %d record(s) to %s", len(records), addr)
|
log.info("Sent %d record(s) to %s", len(records), addr)
|
||||||
except OSError:
|
except OSError:
|
||||||
@@ -64,7 +77,7 @@ def start_push_server(host: str, port: int) -> tuple:
|
|||||||
server_sock.listen()
|
server_sock.listen()
|
||||||
clients: list = []
|
clients: list = []
|
||||||
lock = threading.Lock()
|
lock = threading.Lock()
|
||||||
store_ref: dict = {"current": None, "preloaded": []}
|
store_ref: dict = {"current": None, "preloaded": [], "diff_mode": False}
|
||||||
t = threading.Thread(
|
t = threading.Thread(
|
||||||
target=_accept_loop,
|
target=_accept_loop,
|
||||||
args=(server_sock, clients, lock, store_ref),
|
args=(server_sock, clients, lock, store_ref),
|
||||||
|
|||||||
Reference in New Issue
Block a user