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 <noreply@anthropic.com>
This commit is contained in:
+5
-4
@@ -93,10 +93,11 @@ 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)
|
||||||
|
|
||||||
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)
|
||||||
push_clients, push_lock, push_store_ref = network.start_push_server(
|
if logs_dir is not None:
|
||||||
args.host, args.port, load_today=load_today
|
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()
|
current_day = datetime.now(timezone.utc).date()
|
||||||
store = load_store(out, vin, args.interval)
|
store = load_store(out, vin, args.interval)
|
||||||
|
|||||||
+13
-13
@@ -2,7 +2,6 @@ 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__)
|
||||||
|
|
||||||
@@ -12,7 +11,6 @@ def _accept_loop(
|
|||||||
clients: list,
|
clients: list,
|
||||||
lock: threading.Lock,
|
lock: threading.Lock,
|
||||||
store_ref: dict,
|
store_ref: dict,
|
||||||
load_today: Callable | None,
|
|
||||||
) -> None:
|
) -> None:
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
@@ -20,10 +18,16 @@ def _accept_loop(
|
|||||||
log.info("Push client connected from %s", addr)
|
log.info("Push client connected from %s", addr)
|
||||||
with lock:
|
with lock:
|
||||||
store = store_ref["current"]
|
store = store_ref["current"]
|
||||||
records = store["records"] if store else []
|
in_mem = store["records"] if store else []
|
||||||
if not records and load_today is not None:
|
preloaded = store_ref.get("preloaded", [])
|
||||||
log.info("In-memory store empty — loading today's records from disk")
|
|
||||||
records = load_today()
|
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:
|
try:
|
||||||
for record in records:
|
for record in records:
|
||||||
conn.sendall((json.dumps(record, default=str) + "\n").encode())
|
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")
|
log.info("Push client disconnected")
|
||||||
|
|
||||||
|
|
||||||
def start_push_server(
|
def start_push_server(host: str, port: int) -> tuple:
|
||||||
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()
|
||||||
store_ref: dict = {"current": None}
|
store_ref: dict = {"current": None, "preloaded": []}
|
||||||
t = threading.Thread(
|
t = threading.Thread(
|
||||||
target=_accept_loop,
|
target=_accept_loop,
|
||||||
args=(server_sock, clients, lock, store_ref, load_today),
|
args=(server_sock, clients, lock, store_ref),
|
||||||
daemon=True,
|
daemon=True,
|
||||||
)
|
)
|
||||||
t.start()
|
t.start()
|
||||||
|
|||||||
Reference in New Issue
Block a user