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>
88 lines
3.1 KiB
Python
88 lines
3.1 KiB
Python
import json
|
|
import logging
|
|
import socket
|
|
import threading
|
|
|
|
from utils.jay_diff import jay_diff_full
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
|
|
def _accept_loop(
|
|
server_sock: socket.socket,
|
|
clients: list,
|
|
lock: threading.Lock,
|
|
store_ref: dict,
|
|
) -> None:
|
|
while True:
|
|
try:
|
|
conn, addr = server_sock.accept()
|
|
log.info("Push client connected from %s", addr)
|
|
with lock:
|
|
store = store_ref["current"]
|
|
in_mem = store["records"] if store else []
|
|
preloaded = store_ref.get("preloaded", [])
|
|
|
|
if preloaded:
|
|
seen_ts = {r.get("ts") for r in preloaded}
|
|
extra = [r for r in in_mem if r.get("ts") not in seen_ts]
|
|
records = sorted(preloaded + extra, key=lambda r: r.get("ts", ""))
|
|
else:
|
|
records = in_mem
|
|
|
|
try:
|
|
if store_ref.get("diff_mode") and len(records) > 1:
|
|
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:
|
|
log.info("Sent %d record(s) to %s", len(records), addr)
|
|
except OSError:
|
|
log.debug("Failed to send history to %s", addr)
|
|
conn.close()
|
|
continue
|
|
clients.append(conn)
|
|
except OSError:
|
|
break
|
|
|
|
|
|
def broadcast(snapshot: dict, clients: list, lock: threading.Lock) -> None:
|
|
line = (json.dumps(snapshot, default=str) + "\n").encode()
|
|
with lock:
|
|
dead = []
|
|
for conn in clients:
|
|
try:
|
|
conn.sendall(line)
|
|
except OSError:
|
|
dead.append(conn)
|
|
for conn in dead:
|
|
clients.remove(conn)
|
|
conn.close()
|
|
log.info("Push client disconnected")
|
|
|
|
|
|
def start_push_server(host: str, port: int) -> tuple:
|
|
server_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
server_sock.bind((host, port))
|
|
server_sock.listen()
|
|
clients: list = []
|
|
lock = threading.Lock()
|
|
store_ref: dict = {"current": None, "preloaded": [], "diff_mode": False}
|
|
t = threading.Thread(
|
|
target=_accept_loop,
|
|
args=(server_sock, clients, lock, store_ref),
|
|
daemon=True,
|
|
)
|
|
t.start()
|
|
log.info("Push server listening on %s:%d", host, port)
|
|
return clients, lock, store_ref |