Files
we_monitor/network.py
T
jensandClaude Sonnet 4.6 bc016cd7cb Refactor collect.py into focused modules
- data_model.py  : _str, _phys, load_store, save_store, auto_out_path
- we_connect.py  : domain extractors, ALL_DOMAINS, connect,
                   select_vehicle, collect_snapshot
- network.py     : TCP push server (start_push_server, broadcast)
- log_config.py  : logging setup
- collect.py     : run() loop and CLI, wired to the above modules

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-25 19:41:54 +02:00

49 lines
1.4 KiB
Python

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