Fall back to disk when in-memory store is empty on client connect
storage_helpers: load_today_records() scans logs/ for all files matching today's UTC date and VIN, merges and sorts their records. network: _accept_loop receives an optional load_today callable; if store_ref["current"] has no records it calls load_today() to replay the day's history from disk before adding the client. collect: passes a load_today lambda (None when -o is used). Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+8
-3
@@ -57,7 +57,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_store, save_store
|
from server.storage_helpers import auto_out_path, load_store, load_today_records, save_store
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -93,10 +93,14 @@ def run(args: argparse.Namespace) -> None:
|
|||||||
logs_dir.mkdir(exist_ok=True)
|
logs_dir.mkdir(exist_ok=True)
|
||||||
out = auto_out_path(logs_dir, vin)
|
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()
|
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
|
||||||
log.info(
|
log.info(
|
||||||
"Collecting [%s] every %ds → %s (Ctrl-C to stop)",
|
"Collecting [%s] every %ds → %s (Ctrl-C to stop)",
|
||||||
", ".join(domains),
|
", ".join(domains),
|
||||||
@@ -111,6 +115,7 @@ def run(args: argparse.Namespace) -> None:
|
|||||||
if today != current_day:
|
if today != current_day:
|
||||||
out = auto_out_path(logs_dir, vin)
|
out = auto_out_path(logs_dir, vin)
|
||||||
store = load_store(out, vin, args.interval)
|
store = load_store(out, vin, args.interval)
|
||||||
|
push_store_ref["current"] = store
|
||||||
current_day = today
|
current_day = today
|
||||||
log.info("Midnight UTC rotation → %s", out)
|
log.info("Midnight UTC rotation → %s", out)
|
||||||
|
|
||||||
@@ -119,7 +124,7 @@ 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, push_last)
|
network.broadcast(snapshot, push_clients, push_lock)
|
||||||
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
|
||||||
|
|||||||
+26
-15
@@ -2,6 +2,7 @@ import json
|
|||||||
import logging
|
import logging
|
||||||
import socket
|
import socket
|
||||||
import threading
|
import threading
|
||||||
|
from typing import Callable
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
@@ -10,30 +11,36 @@ def _accept_loop(
|
|||||||
server_sock: socket.socket,
|
server_sock: socket.socket,
|
||||||
clients: list,
|
clients: list,
|
||||||
lock: threading.Lock,
|
lock: threading.Lock,
|
||||||
last_snapshot: dict,
|
store_ref: dict,
|
||||||
|
load_today: Callable | None,
|
||||||
) -> None:
|
) -> None:
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
conn, addr = server_sock.accept()
|
conn, addr = server_sock.accept()
|
||||||
log.info("Push client connected from %s", addr)
|
log.info("Push client connected from %s", addr)
|
||||||
with lock:
|
with lock:
|
||||||
if last_snapshot["data"] is not None:
|
store = store_ref["current"]
|
||||||
try:
|
records = store["records"] if store else []
|
||||||
line = (json.dumps(last_snapshot["data"], default=str) + "\n").encode()
|
if not records and load_today is not None:
|
||||||
conn.sendall(line)
|
log.info("In-memory store empty — loading today's records from disk")
|
||||||
except OSError:
|
records = load_today()
|
||||||
log.debug("Failed to send initial snapshot to %s", addr)
|
try:
|
||||||
conn.close()
|
for record in records:
|
||||||
continue
|
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)
|
clients.append(conn)
|
||||||
except OSError:
|
except OSError:
|
||||||
break
|
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()
|
line = (json.dumps(snapshot, default=str) + "\n").encode()
|
||||||
with lock:
|
with lock:
|
||||||
last_snapshot["data"] = snapshot
|
|
||||||
dead = []
|
dead = []
|
||||||
for conn in clients:
|
for conn in clients:
|
||||||
try:
|
try:
|
||||||
@@ -46,19 +53,23 @@ def broadcast(snapshot: dict, clients: list, lock: threading.Lock, last_snapshot
|
|||||||
log.info("Push client disconnected")
|
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 = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||||
server_sock.bind((host, port))
|
server_sock.bind((host, port))
|
||||||
server_sock.listen()
|
server_sock.listen()
|
||||||
clients: list = []
|
clients: list = []
|
||||||
lock = threading.Lock()
|
lock = threading.Lock()
|
||||||
last_snapshot: dict = {"data": None}
|
store_ref: dict = {"current": None}
|
||||||
t = threading.Thread(
|
t = threading.Thread(
|
||||||
target=_accept_loop,
|
target=_accept_loop,
|
||||||
args=(server_sock, clients, lock, last_snapshot),
|
args=(server_sock, clients, lock, store_ref, load_today),
|
||||||
daemon=True,
|
daemon=True,
|
||||||
)
|
)
|
||||||
t.start()
|
t.start()
|
||||||
log.info("Push server listening on %s:%d", host, port)
|
log.info("Push server listening on %s:%d", host, port)
|
||||||
return clients, lock, last_snapshot
|
return clients, lock, store_ref
|
||||||
@@ -33,3 +33,19 @@ def save_store(path: Path, store: dict, max_records: int | None) -> None:
|
|||||||
if max_records and len(store["records"]) > max_records:
|
if max_records and len(store["records"]) > max_records:
|
||||||
store["records"] = store["records"][-max_records:]
|
store["records"] = store["records"][-max_records:]
|
||||||
path.write_text(json.dumps(store, indent=2, default=str))
|
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
|
||||||
Reference in New Issue
Block a user