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"): prev: dict = {} for record in records: diff = jay_diff_full(prev, record, combine_upd_add=True) prev = record conn.sendall((json.dumps(diff, default=str) + "\n").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