140 lines
4.7 KiB
Python
140 lines
4.7 KiB
Python
import json
|
|
import socket
|
|
import threading
|
|
import time
|
|
import lz4.frame
|
|
import zstandard as zstd
|
|
from numcodecs import Blosc
|
|
import numpy as np
|
|
|
|
|
|
class StreamReceiver:
|
|
def __init__(self, host="0.0.0.0", port=6001):
|
|
self.host = host
|
|
self.port = port
|
|
|
|
self._server_sock = None
|
|
self._client_sock = None
|
|
self._thread = None
|
|
self._running = False
|
|
|
|
self.last_frame = None
|
|
self.last_meta = None
|
|
self.last_receive_ts = None
|
|
|
|
self._zstd_d = zstd.ZstdDecompressor()
|
|
self._codec = Blosc(cname="lz4", clevel=1, shuffle=Blosc.SHUFFLE)
|
|
|
|
@property
|
|
def is_running(self):
|
|
return self._running
|
|
|
|
def start(self):
|
|
if self._running:
|
|
return
|
|
|
|
self._running = True
|
|
self._thread = threading.Thread(target=self._worker, daemon=True)
|
|
self._thread.start()
|
|
|
|
def stop(self):
|
|
self._running = False
|
|
|
|
try:
|
|
if self._client_sock:
|
|
self._client_sock.close()
|
|
except:
|
|
pass
|
|
|
|
try:
|
|
if self._server_sock:
|
|
self._server_sock.close()
|
|
except:
|
|
pass
|
|
|
|
self._client_sock = None
|
|
self._server_sock = None
|
|
|
|
def _recv_exact(self, sock: socket.socket, n: int) -> bytes:
|
|
chunks = []
|
|
remaining = n
|
|
|
|
while remaining > 0:
|
|
chunk = sock.recv(remaining)
|
|
if not chunk:
|
|
raise ConnectionError("Conexão encerrada durante recv")
|
|
chunks.append(chunk)
|
|
remaining -= len(chunk)
|
|
|
|
return b"".join(chunks)
|
|
|
|
def _worker(self):
|
|
try:
|
|
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as server:
|
|
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
server.bind((self.host, self.port))
|
|
server.listen(1)
|
|
server.settimeout(1.0)
|
|
|
|
self._server_sock = server
|
|
print(f"[INFO] StreamReceiver ouvindo em {self.host}:{self.port}")
|
|
|
|
while self._running:
|
|
try:
|
|
client, addr = server.accept()
|
|
except socket.timeout:
|
|
continue
|
|
|
|
print(f"[INFO] StreamReceiver conectado por {addr}")
|
|
self._client_sock = client
|
|
|
|
with client:
|
|
while self._running:
|
|
header_len = int.from_bytes(self._recv_exact(client, 4), "big")
|
|
header_bytes = self._recv_exact(client, header_len)
|
|
header = json.loads(header_bytes.decode("utf-8"))
|
|
|
|
payload_len = int.from_bytes(self._recv_exact(client, 4), "big")
|
|
payload_comp = self._recv_exact(client, payload_len)
|
|
codec = header["codec"]
|
|
|
|
if codec == "lz4":
|
|
payload = lz4.frame.decompress(payload_comp)
|
|
elif codec == "zstd":
|
|
payload = self._zstd_d.decompress(payload_comp)
|
|
elif codec == "numcodecs":
|
|
payload = self._codec.decode(payload_comp)
|
|
else:
|
|
raise ValueError(f"Codec não suportado: {codec}")
|
|
|
|
expected = header["payload_size_raw"]
|
|
if len(payload) != expected:
|
|
raise ValueError(
|
|
f"Tamanho descomprimido inválido: {len(payload)} != {expected}"
|
|
)
|
|
|
|
height = header["height"]
|
|
width = header["width"]
|
|
channels = header["channels"]
|
|
dtype = np.uint8 if header["dtype"] == "uint8" else None
|
|
|
|
if dtype is None:
|
|
raise RuntimeError(f"dtype não suportado: {header['dtype']}")
|
|
|
|
frame = np.frombuffer(payload, dtype=dtype).reshape(height, width, channels)
|
|
|
|
self.last_frame = frame
|
|
self.last_meta = header
|
|
self.last_receive_ts = time.perf_counter()
|
|
|
|
print("[INFO] StreamReceiver cliente desconectado")
|
|
self._client_sock = None
|
|
|
|
except Exception as e:
|
|
print(f"[WARN] StreamReceiver encerrado com erro: {e}")
|
|
|
|
finally:
|
|
self._running = False
|
|
self._client_sock = None
|
|
self._server_sock = None
|