from __future__ import annotations import logging import os import socket import threading import time import uuid import cv2 from .core import Frame, LatestQueue LOG = logging.getLogger("tail2.ndi") class NDIReceiver(threading.Thread): """Receive only. This class intentionally exposes no NDI PTZ methods.""" daemon = True def __init__(self, frames: LatestQueue, state) -> None: super().__init__(name="ndi-receiver") self.frames, self.state = frames, state self.source_name = os.environ["NDI_SOURCE_NAME"] self.source_ip = os.environ["NDI_SOURCE_IP"] self.stop_event = threading.Event() def stop(self) -> None: self.stop_event.set() def run(self) -> None: while not self.stop_event.is_set(): try: self._receive_session() except Exception as exc: self.state.ndi_connected = False self.state.error_code = "NDI_DISCONNECTED" LOG.warning("NDI receive session ended: %s: %s", type(exc).__name__, str(exc)[:300]) self.stop_event.wait(2.0) def _receive_session(self) -> None: from cyndilib.finder import Finder from cyndilib.receiver import Receiver from cyndilib.video_frame import VideoFrameSync from cyndilib.wrapper.ndi_recv import RecvBandwidth, RecvColorFormat # NDI SDK reads this file before finder initialization. It supplies the # camera as an additional unicast discovery server without changing it. ndi_dir = "/tmp/ndi" os.makedirs(ndi_dir, exist_ok=True) with open(f"{ndi_dir}/ndi-config.v1.json", "w", encoding="utf-8") as fh: fh.write('{"ndi":{"networks":{"ips":"' + self.source_ip + '"}}}') os.environ["NDI_CONFIG_DIR"] = ndi_dir finder = Finder() finder.open() receiver = None try: deadline = time.monotonic() + 30 source = None while not self.stop_event.is_set() and time.monotonic() < deadline: finder.wait_for_sources(1) names = finder.get_source_names() if self.source_name in names: source = finder.get_source(self.source_name) break if source is None: raise RuntimeError(f"configured NDI source not discovered; last_sources={names}") receiver = Receiver(color_format=RecvColorFormat.BGRX_BGRA, bandwidth=RecvBandwidth.highest) video = VideoFrameSync() receiver.frame_sync.set_video_frame(video) receiver.set_source(source) deadline = time.monotonic() + 15 while not receiver.is_connected() and time.monotonic() < deadline: time.sleep(0.2) if not receiver.is_connected(): raise RuntimeError("NDI source connection timeout") self.state.session_id = str(uuid.uuid4()) while not self.stop_event.is_set() and receiver.is_connected(): started = time.perf_counter_ns() receiver.frame_sync.capture_video() if min(video.xres, video.yres) <= 0: time.sleep(0.005) continue raw = video.get_array().reshape(video.yres, video.xres, 4) pixels = cv2.cvtColor(raw, cv2.COLOR_BGRA2BGR) now = time.time_ns() self.frames.put(Frame(pixels.copy(), now, (time.perf_counter_ns() - started) / 1e6)) self.state.note_frame(now, video.xres, video.yres) finally: if receiver is not None and receiver.is_connected(): receiver.disconnect() finder.close()