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()