diff --git a/collect.py b/collect.py index 2c7b05c..19de6b1 100644 --- a/collect.py +++ b/collect.py @@ -57,7 +57,7 @@ from pathlib import Path from server import log_config, network, we_connect from server.data_model import ALL_DOMAINS, collect_snapshot -from server.storage_helpers import auto_out_path, load_store, save_store +from server.storage_helpers import auto_out_path, load_store, load_today_records, save_store log = logging.getLogger(__name__) @@ -93,10 +93,14 @@ def run(args: argparse.Namespace) -> None: logs_dir.mkdir(exist_ok=True) out = auto_out_path(logs_dir, vin) - push_clients, push_lock, push_last = network.start_push_server(args.host, args.port) + 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 + ) current_day = datetime.now(timezone.utc).date() store = load_store(out, vin, args.interval) + push_store_ref["current"] = store log.info( "Collecting [%s] every %ds → %s (Ctrl-C to stop)", ", ".join(domains), @@ -111,6 +115,7 @@ def run(args: argparse.Namespace) -> None: if today != current_day: out = auto_out_path(logs_dir, vin) store = load_store(out, vin, args.interval) + push_store_ref["current"] = store current_day = today log.info("Midnight UTC rotation → %s", out) @@ -119,7 +124,7 @@ def run(args: argparse.Namespace) -> None: snapshot = collect_snapshot(vehicle, domains) store["records"].append(snapshot) save_store(out, store, args.max_records) - network.broadcast(snapshot, push_clients, push_lock, push_last) + network.broadcast(snapshot, push_clients, push_lock) log.info("Record #%d saved at %s", len(store["records"]), snapshot["ts"]) except KeyboardInterrupt: raise diff --git a/server/network.py b/server/network.py index 8ccee12..90ae519 100644 --- a/server/network.py +++ b/server/network.py @@ -2,6 +2,7 @@ import json import logging import socket import threading +from typing import Callable log = logging.getLogger(__name__) @@ -10,30 +11,36 @@ def _accept_loop( server_sock: socket.socket, clients: list, lock: threading.Lock, - last_snapshot: dict, + 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: - if last_snapshot["data"] is not None: - try: - line = (json.dumps(last_snapshot["data"], default=str) + "\n").encode() - conn.sendall(line) - except OSError: - log.debug("Failed to send initial snapshot to %s", addr) - conn.close() - continue + 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() + try: + 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, last_snapshot: dict) -> None: +def broadcast(snapshot: dict, clients: list, lock: threading.Lock) -> None: line = (json.dumps(snapshot, default=str) + "\n").encode() with lock: - last_snapshot["data"] = snapshot dead = [] for conn in clients: try: @@ -46,19 +53,23 @@ def broadcast(snapshot: dict, clients: list, lock: threading.Lock, last_snapshot log.info("Push client disconnected") -def start_push_server(host: str, port: int) -> tuple: +def start_push_server( + host: str, + port: int, + load_today: Callable | None = None, +) -> 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() - last_snapshot: dict = {"data": None} + store_ref: dict = {"current": None} t = threading.Thread( target=_accept_loop, - args=(server_sock, clients, lock, last_snapshot), + args=(server_sock, clients, lock, store_ref, load_today), daemon=True, ) t.start() log.info("Push server listening on %s:%d", host, port) - return clients, lock, last_snapshot \ No newline at end of file + return clients, lock, store_ref \ No newline at end of file diff --git a/server/storage_helpers.py b/server/storage_helpers.py index 9445e3f..7ca5e46 100644 --- a/server/storage_helpers.py +++ b/server/storage_helpers.py @@ -32,4 +32,20 @@ def load_store(path: Path, vin: str, interval: int) -> dict: def save_store(path: Path, store: dict, max_records: int | None) -> None: if max_records and len(store["records"]) > max_records: store["records"] = store["records"][-max_records:] - path.write_text(json.dumps(store, indent=2, default=str)) \ No newline at end of file + path.write_text(json.dumps(store, indent=2, default=str)) + + +def load_today_records(logs_dir: Path, vin: str) -> list: + """Load and merge all records from today's log files for this VIN.""" + today_prefix = datetime.now(timezone.utc).strftime("%Y_%m_%d_") + records = [] + for path in sorted(logs_dir.glob(f"{today_prefix}*_{vin}.json")): + try: + data = json.loads(path.read_text()) + if isinstance(data, dict) and "records" in data: + records.extend(data["records"]) + log.debug("Loaded %d records from %s", len(data["records"]), path.name) + except (json.JSONDecodeError, OSError): + log.warning("Could not load %s", path) + records.sort(key=lambda r: r.get("ts", "")) + return records \ No newline at end of file