474 lines
15 KiB
Python
474 lines
15 KiB
Python
import json
|
|
import socket
|
|
import threading
|
|
import time
|
|
import queue
|
|
from numcodecs import Blosc
|
|
|
|
|
|
class StreamSender:
|
|
_QUEUE_SENTINEL = object()
|
|
|
|
def __init__(self, state, frame_service):
|
|
self.state = state
|
|
self.frame_service = frame_service
|
|
|
|
self._lock = threading.RLock()
|
|
self._stop_event = threading.Event()
|
|
|
|
self._sock = None
|
|
self._thread_capture = None
|
|
self._thread_send = None
|
|
|
|
self._queue = queue.Queue(maxsize=2)
|
|
|
|
self._codec = None
|
|
self._codec_signature = None
|
|
|
|
self._t_json_pack = 0.0
|
|
self._t_send_header = 0.0
|
|
self._t_send_payload = 0.0
|
|
|
|
self._frames_dropped = 0
|
|
self._capture_errors = 0
|
|
self._send_errors = 0
|
|
self._last_queued_capture_signature = None
|
|
self._last_sent_capture_signature = None
|
|
|
|
@property
|
|
def is_running(self):
|
|
return not self._stop_event.is_set() and (
|
|
self._thread_capture is not None or self._thread_send is not None
|
|
)
|
|
|
|
# =========================================================
|
|
# Lifecycle
|
|
# =========================================================
|
|
|
|
def start(self, host: str, port: int, fps: float):
|
|
with self._lock:
|
|
if self.is_running:
|
|
raise RuntimeError("Stream já está em execução")
|
|
|
|
if not self.state.initialized:
|
|
raise RuntimeError("Módulo não inicializado")
|
|
|
|
self._stop_event.clear()
|
|
self._clear_queue()
|
|
|
|
self._t_json_pack = 0.0
|
|
self._t_send_header = 0.0
|
|
self._t_send_payload = 0.0
|
|
self._frames_dropped = 0
|
|
self._capture_errors = 0
|
|
self._send_errors = 0
|
|
|
|
self.state.streaming = True
|
|
self.state.stream_host = host
|
|
self.state.stream_port = port
|
|
self.state.stream_fps = fps
|
|
self._last_queued_capture_signature = None
|
|
self._last_sent_capture_signature = None
|
|
|
|
self._thread_capture = threading.Thread(
|
|
target=self._worker_capture,
|
|
args=(fps,),
|
|
daemon=True
|
|
)
|
|
self._thread_send = threading.Thread(
|
|
target=self._worker_send,
|
|
args=(host, port),
|
|
daemon=True
|
|
)
|
|
|
|
self._thread_capture.start()
|
|
self._thread_send.start()
|
|
|
|
def stop(self):
|
|
with self._lock:
|
|
self._stop_event.set()
|
|
|
|
try:
|
|
self._queue.put_nowait(self._QUEUE_SENTINEL)
|
|
except queue.Full:
|
|
try:
|
|
self._queue.get_nowait()
|
|
except queue.Empty:
|
|
pass
|
|
try:
|
|
self._queue.put_nowait(self._QUEUE_SENTINEL)
|
|
except queue.Full:
|
|
pass
|
|
|
|
sock = self._sock
|
|
self._sock = None
|
|
|
|
if sock is not None:
|
|
try:
|
|
sock.shutdown(socket.SHUT_RDWR)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
sock.close()
|
|
except Exception:
|
|
pass
|
|
|
|
if self._thread_capture is not None:
|
|
self._thread_capture.join(timeout=2.0)
|
|
self._thread_capture = None
|
|
|
|
if self._thread_send is not None:
|
|
self._thread_send.join(timeout=2.0)
|
|
self._thread_send = None
|
|
|
|
self._clear_queue()
|
|
|
|
self.state.streaming = False
|
|
self.state.stream_host = None
|
|
self.state.stream_port = None
|
|
self.state.stream_fps = None
|
|
self.state.stream_frame_id_sent = 0
|
|
|
|
def _clear_queue(self):
|
|
while not self._queue.empty():
|
|
try:
|
|
self._queue.get_nowait()
|
|
except queue.Empty:
|
|
break
|
|
|
|
# =========================================================
|
|
# Codec
|
|
# =========================================================
|
|
|
|
def _normalize_shuffle(self, shuffle_value):
|
|
if isinstance(shuffle_value, int):
|
|
return shuffle_value
|
|
|
|
mapping = {
|
|
"NOSHUFFLE": Blosc.NOSHUFFLE,
|
|
"SHUFFLE": Blosc.SHUFFLE,
|
|
"BITSHUFFLE": Blosc.BITSHUFFLE,
|
|
}
|
|
|
|
key = str(shuffle_value).upper()
|
|
if key not in mapping:
|
|
raise ValueError(f"shuffle inválido: {shuffle_value}")
|
|
|
|
return mapping[key]
|
|
|
|
def _build_codec_from_state(self):
|
|
family = self.state.codec_family
|
|
name = self.state.codec_name
|
|
params = dict(self.state.codec_params)
|
|
|
|
if family == "none":
|
|
return None
|
|
|
|
if family != "numcodecs":
|
|
raise ValueError(f"Família de codec não suportada: {family}")
|
|
|
|
if name == "blosc":
|
|
params["shuffle"] = self._normalize_shuffle(params.get("shuffle", "SHUFFLE"))
|
|
return Blosc(**params)
|
|
|
|
raise ValueError(f"Codec numcodecs não suportado: {name}")
|
|
|
|
def _get_codec_signature_from_state(self):
|
|
return (
|
|
self.state.codec_family,
|
|
self.state.codec_name,
|
|
tuple(sorted(self.state.codec_params.items()))
|
|
)
|
|
|
|
def _ensure_codec(self):
|
|
sig = self._get_codec_signature_from_state()
|
|
if self._codec is None or self._codec_signature != sig:
|
|
self._codec = self._build_codec_from_state()
|
|
self._codec_signature = sig
|
|
|
|
def _compress(self, frame_bytes: bytes) -> bytes:
|
|
self._ensure_codec()
|
|
|
|
if self.state.codec_family == "none":
|
|
return frame_bytes
|
|
|
|
return self._codec.encode(frame_bytes)
|
|
|
|
# =========================================================
|
|
# Packet helpers
|
|
# =========================================================
|
|
|
|
def _send_packet(self, sock: socket.socket, header: dict, payload: bytes):
|
|
t0 = time.perf_counter()
|
|
header_bytes = json.dumps(header, separators=(",", ":")).encode("utf-8")
|
|
t1 = time.perf_counter()
|
|
|
|
sock.sendall(len(header_bytes).to_bytes(4, "big"))
|
|
sock.sendall(header_bytes)
|
|
sock.sendall(len(payload).to_bytes(4, "big"))
|
|
t2 = time.perf_counter()
|
|
|
|
sock.sendall(payload)
|
|
t3 = time.perf_counter()
|
|
|
|
self._t_json_pack = t1 - t0
|
|
self._t_send_header = t2 - t1
|
|
self._t_send_payload = t3 - t2
|
|
|
|
def _build_capture_signature(self, meta):
|
|
camera_frames = meta.get("camera_frames", {})
|
|
if not camera_frames:
|
|
return None
|
|
|
|
items = []
|
|
for cam_id in sorted(camera_frames.keys()):
|
|
items.append((cam_id, camera_frames[cam_id].get("camera_frame_id")))
|
|
return tuple(items)
|
|
|
|
def _describe_single_array(self, frame, meta):
|
|
return {
|
|
"multi_payload": False,
|
|
"payload_kind": "single_array",
|
|
"payload_parts": [
|
|
{
|
|
"kind": "single_array",
|
|
"size_raw": int(frame.nbytes),
|
|
"dtype": str(frame.dtype),
|
|
"shape": list(frame.shape),
|
|
"layout": meta.get("output_layout"),
|
|
"channel_names": meta.get("output_channel_names"),
|
|
}
|
|
],
|
|
}
|
|
|
|
def _describe_multi_array_part(self, cam_id, arr, cam_meta):
|
|
return {
|
|
"camera_id": cam_id,
|
|
"size_raw": int(arr.nbytes),
|
|
"dtype": str(arr.dtype),
|
|
"shape": list(arr.shape),
|
|
"width": cam_meta.get("width"),
|
|
"height": cam_meta.get("height"),
|
|
"channels": cam_meta.get("channels"),
|
|
"role": cam_meta.get("role"),
|
|
"interface": cam_meta.get("interface"),
|
|
}
|
|
|
|
def _serialize_frame_payload(self, frame, meta):
|
|
"""
|
|
Para array único:
|
|
payload = bytes diretos do array
|
|
Para dict de arrays:
|
|
payload = repetição de:
|
|
[4 bytes tamanho][bytes da parte]
|
|
"""
|
|
if not isinstance(frame, dict):
|
|
frame_bytes = frame.tobytes()
|
|
payload_meta = self._describe_single_array(frame, meta)
|
|
return frame_bytes, payload_meta
|
|
|
|
payload = bytearray()
|
|
payload_parts = []
|
|
|
|
camera_frames_meta = meta.get("camera_frames", {})
|
|
|
|
for cam_id in meta.get("payload_sources", list(frame.keys())):
|
|
arr = frame[cam_id]
|
|
part_bytes = arr.tobytes()
|
|
|
|
payload.extend(len(part_bytes).to_bytes(4, "big"))
|
|
payload.extend(part_bytes)
|
|
|
|
payload_parts.append(
|
|
self._describe_multi_array_part(
|
|
cam_id,
|
|
arr,
|
|
camera_frames_meta.get(cam_id, {})
|
|
)
|
|
)
|
|
|
|
payload_meta = {
|
|
"multi_payload": True,
|
|
"payload_kind": "multi_array",
|
|
"payload_parts": payload_parts,
|
|
}
|
|
return bytes(payload), payload_meta
|
|
|
|
def _build_stream_header(self, meta, raw_payload, comp_bytes, payload_meta, dt_bytes, dt_comp, dt_frame_period):
|
|
codec_name = None if self.state.codec_family == "none" else self.state.codec_name
|
|
codec_params = {} if self.state.codec_family == "none" else dict(self.state.codec_params)
|
|
|
|
header = {
|
|
"frame_id": meta.get("frame_id"),
|
|
|
|
"module": self.state.module_name,
|
|
"module_version": self.state.version,
|
|
|
|
"frame_type": meta.get("frame_type"),
|
|
"payload_format_version": meta.get("payload_format_version"),
|
|
"capture_mode_resolved": meta.get("capture_mode_resolved"),
|
|
|
|
"output_dtype": meta.get("output_dtype"),
|
|
"output_layout": meta.get("output_layout"),
|
|
"output_channels": meta.get("output_channels"),
|
|
"output_channel_names": meta.get("output_channel_names"),
|
|
"output_width": meta.get("output_width"),
|
|
"output_height": meta.get("output_height"),
|
|
|
|
"payload_sources": meta.get("payload_sources"),
|
|
"camera_frames": meta.get("camera_frames"),
|
|
|
|
"codec_family": self.state.codec_family,
|
|
"codec_name": codec_name,
|
|
"codec_params": codec_params,
|
|
|
|
"payload_size_raw": len(raw_payload),
|
|
"payload_size_comp": len(comp_bytes),
|
|
|
|
"dt_frame_period": dt_frame_period,
|
|
"dt_bytes": dt_bytes,
|
|
"dt_comp": dt_comp,
|
|
|
|
"dt_trigger": meta.get("dt_trigger"),
|
|
"dt_settle": meta.get("dt_settle"),
|
|
"dt_capture": meta.get("dt_capture"),
|
|
"dt_process": meta.get("dt_process"),
|
|
"dt_total_pi": meta.get("dt_total_pi"),
|
|
|
|
"ts_pi": meta.get("ts_pi"),
|
|
"ts_pi_monotonic": meta.get("ts_pi_monotonic"),
|
|
|
|
**payload_meta,
|
|
}
|
|
|
|
# carrega o restante do meta sem sobrescrever campos já consolidados
|
|
reserved = set(header.keys())
|
|
for k, v in meta.items():
|
|
if k not in reserved:
|
|
header[k] = v
|
|
|
|
return header
|
|
|
|
# =========================================================
|
|
# Workers
|
|
# =========================================================
|
|
|
|
def _worker_capture(self, fps: float):
|
|
frame_interval = 1.0 / fps if fps > 0 else 0.0
|
|
next_deadline = time.perf_counter()
|
|
last_frame_ts = None
|
|
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
frame, meta = self.frame_service.capture_frame_raw()
|
|
except Exception as e:
|
|
self._capture_errors += 1
|
|
print(f"[WARN] Capture falhou: {e}")
|
|
time.sleep(0.05)
|
|
continue
|
|
|
|
if frame is None:
|
|
time.sleep(0.01)
|
|
continue
|
|
|
|
capture_signature = self._build_capture_signature(meta)
|
|
if capture_signature is not None and capture_signature == self._last_queued_capture_signature:
|
|
time.sleep(0.001)
|
|
continue
|
|
|
|
t_frame_ready = time.perf_counter()
|
|
dt_frame_period = 0.0 if last_frame_ts is None else (t_frame_ready - last_frame_ts)
|
|
last_frame_ts = t_frame_ready
|
|
|
|
t_bytes0 = time.perf_counter()
|
|
raw_payload, payload_meta = self._serialize_frame_payload(frame, meta)
|
|
t_bytes1 = time.perf_counter()
|
|
|
|
t_comp0 = time.perf_counter()
|
|
comp_bytes = self._compress(raw_payload)
|
|
t_comp1 = time.perf_counter()
|
|
|
|
header = self._build_stream_header(
|
|
meta=meta,
|
|
raw_payload=raw_payload,
|
|
comp_bytes=comp_bytes,
|
|
payload_meta=payload_meta,
|
|
dt_bytes=(t_bytes1 - t_bytes0),
|
|
dt_comp=(t_comp1 - t_comp0),
|
|
dt_frame_period=dt_frame_period,
|
|
)
|
|
|
|
queued = False
|
|
try:
|
|
self._queue.put_nowait((header, comp_bytes))
|
|
queued = True
|
|
except queue.Full:
|
|
self._frames_dropped += 1
|
|
try:
|
|
self._queue.get_nowait()
|
|
except queue.Empty:
|
|
pass
|
|
try:
|
|
self._queue.put_nowait((header, comp_bytes))
|
|
queued = True
|
|
except queue.Full:
|
|
self._frames_dropped += 1
|
|
|
|
if queued:
|
|
self._last_queued_capture_signature = capture_signature
|
|
|
|
if frame_interval > 0:
|
|
next_deadline += frame_interval
|
|
now = time.perf_counter()
|
|
|
|
if next_deadline < now - frame_interval:
|
|
next_deadline = now
|
|
|
|
sleep_time = next_deadline - now
|
|
if sleep_time > 0:
|
|
time.sleep(sleep_time)
|
|
|
|
def _worker_send(self, host: str, port: int):
|
|
try:
|
|
with socket.create_connection((host, port), timeout=5) as sock:
|
|
sock.settimeout(10)
|
|
with self._lock:
|
|
self._sock = sock
|
|
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
item = self._queue.get(timeout=0.2)
|
|
except queue.Empty:
|
|
continue
|
|
|
|
if item is self._QUEUE_SENTINEL:
|
|
break
|
|
|
|
header, payload = item
|
|
|
|
header["dt_pack_prev"] = self._t_json_pack
|
|
header["dt_send_header_prev"] = self._t_send_header
|
|
header["dt_send_payload_prev"] = self._t_send_payload
|
|
header["stream_drops"] = self._frames_dropped
|
|
header["capture_errors"] = self._capture_errors
|
|
header["send_errors"] = self._send_errors
|
|
|
|
self._send_packet(sock, header, payload)
|
|
self.state.stream_frame_id_sent = header["frame_id"]
|
|
capture_signature = self._build_capture_signature(header)
|
|
self._last_sent_capture_signature = capture_signature
|
|
|
|
except Exception as e:
|
|
self._send_errors += 1
|
|
if not self._stop_event.is_set():
|
|
print(f"[WARN] StreamSender encerrado com erro: {e}")
|
|
finally:
|
|
with self._lock:
|
|
self._sock = None
|
|
|
|
self.state.streaming = False
|
|
self.state.stream_host = None
|
|
self.state.stream_port = None
|
|
self.state.stream_fps = None
|
|
|
|
self._stop_event.set() |