From 04578781af904b5f52090a496b99527e89f96e6d Mon Sep 17 00:00:00 2001 From: Jens Ahrensfeld Date: Mon, 25 May 2026 21:39:22 +0200 Subject: [PATCH] Eagerly pre-load today's disk records into push_store_ref at startup Previously load_today() was only called if in-memory records were empty, which was never true once the first WeConnect poll had completed. Now disk records are loaded unconditionally at startup and merged with live in-memory records (dedup by ts) when a client connects. Co-Authored-By: Claude Sonnet 4.6 --- collect.py | 9 +++++---- server/network.py | 32 ++++++++++++++++---------------- 2 files changed, 21 insertions(+), 20 deletions(-) diff --git a/collect.py b/collect.py index 19de6b1..1a621c1 100644 --- a/collect.py +++ b/collect.py @@ -93,10 +93,11 @@ def run(args: argparse.Namespace) -> None: logs_dir.mkdir(exist_ok=True) out = auto_out_path(logs_dir, vin) - load_today = (lambda: load_today_records(logs_dir, vin)) if logs_dir else None - push_clients, push_lock, push_store_ref = network.start_push_server( - args.host, args.port, load_today=load_today - ) + push_clients, push_lock, push_store_ref = network.start_push_server(args.host, args.port) + if logs_dir is not None: + push_store_ref["preloaded"] = load_today_records(logs_dir, vin) + if push_store_ref["preloaded"]: + log.info("Pre-loaded %d record(s) from today's log files", len(push_store_ref["preloaded"])) current_day = datetime.now(timezone.utc).date() store = load_store(out, vin, args.interval) diff --git a/server/network.py b/server/network.py index 90ae519..01520fe 100644 --- a/server/network.py +++ b/server/network.py @@ -2,7 +2,6 @@ import json import logging import socket import threading -from typing import Callable log = logging.getLogger(__name__) @@ -12,18 +11,23 @@ def _accept_loop( clients: list, lock: threading.Lock, store_ref: dict, - load_today: Callable | None, ) -> None: while True: try: conn, addr = server_sock.accept() log.info("Push client connected from %s", addr) with lock: - store = store_ref["current"] - records = store["records"] if store else [] - if not records and load_today is not None: - log.info("In-memory store empty — loading today's records from disk") - records = load_today() + 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: for record in records: conn.sendall((json.dumps(record, default=str) + "\n").encode()) @@ -53,21 +57,17 @@ def broadcast(snapshot: dict, clients: list, lock: threading.Lock) -> None: log.info("Push client disconnected") -def start_push_server( - host: str, - port: int, - load_today: Callable | None = None, -) -> tuple: +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} + clients: list = [] + lock = threading.Lock() + store_ref: dict = {"current": None, "preloaded": []} t = threading.Thread( target=_accept_loop, - args=(server_sock, clients, lock, store_ref, load_today), + args=(server_sock, clients, lock, store_ref), daemon=True, ) t.start()