import json import logging import socket import threading log = logging.getLogger(__name__) def _accept_loop( server_sock: socket.socket, clients: list, lock: threading.Lock, ) -> None: while True: try: conn, addr = server_sock.accept() log.info("Push client connected from %s", addr) with lock: 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() t = threading.Thread(target=_accept_loop, args=(server_sock, clients, lock), daemon=True) t.start() log.info("Push server listening on %s:%d", host, port) return clients, lock