"""Resilient video capture: RTSP (or webcam) reader thread with reconnect. Design: one daemon thread per source holds the newest frame in a single slot. Consumers always get the latest frame (never a backlog), and a lost camera reconnects with exponential backoff instead of killing the pipeline. """ from __future__ import annotations import logging import os import threading import time from typing import Optional import cv2 import numpy as np log = logging.getLogger(__name__) # Force TCP transport and a 5s socket timeout for RTSP before OpenCV loads # ffmpeg. UDP is the default and silently drops frames on lossy Wi-Fi. os.environ.setdefault( "OPENCV_FFMPEG_CAPTURE_OPTIONS", "rtsp_transport;tcp|stimeout;5000000" ) def _tcp_reachable(source: "str | int", timeout: float ) -> "tuple[bool, str]": """Cheap pre-flight for an rtsp:// URL. Non-URL sources pass through.""" import socket from urllib.parse import urlparse if isinstance(source, int): return True, "" parsed = urlparse(source) if not parsed.hostname: return True, "" # not a form we can pre-check; let OpenCV try port = parsed.port or (554 if parsed.scheme == "rtsp" else 80) try: with socket.create_connection((parsed.hostname, port), timeout): return True, "" except socket.timeout: return False, (f"no response from {parsed.hostname}:{port} within " f"{timeout:.0f}s - check the IP address and that the " f"camera is on the same network") except OSError as exc: return False, f"cannot reach {parsed.hostname}:{port} - {exc.strerror or exc}" def _fourcc(cap) -> str: """The stream's codec as a four-character code, or "" if unknown. FFmpeg reports H.265 as "hevc" and H.264 as "h264"/"avc1" depending on the container. Returned as-is rather than mapped to a friendly name: the raw value is what somebody searching their camera's manual will match. """ try: raw = int(cap.get(cv2.CAP_PROP_FOURCC)) except Exception: return "" if raw <= 0: return "" code = "".join(chr((raw >> (8 * i)) & 0xFF) for i in range(4)) return code.strip().strip("\x00") def probe_source(source: "str | int", max_width: int = 1280, timeout: float = 12.0, connect_timeout: float = 3.0) -> dict: """Open a candidate camera, grab one frame, and let go. Backs the UI's Test button, so it must answer for a *wrong* URL as reliably as a right one: no retries, no reconnect loop, and a hard deadline because a bad host makes cv2.VideoCapture block until FFmpeg gives up. Returns a JPEG snapshot so the user can confirm the camera is pointing where they think it is. """ import base64 # cv2.VideoCapture blocks inside the constructor while FFmpeg completes a # TCP connect, and against an unroutable host that is the OS connect # timeout (~75s), not our deadline. A wrong IP or port is the single most # likely thing a user types, so check reachability first — it turns the # common failure into a sub-second answer instead of a frozen UI. reachable, why = _tcp_reachable(source, connect_timeout) if not reachable: return {"ok": False, "error": why} cap = None try: cap = (cv2.VideoCapture(source) if isinstance(source, int) else cv2.VideoCapture(source, cv2.CAP_FFMPEG)) cap.set(cv2.CAP_PROP_BUFFERSIZE, 1) if not cap.isOpened(): return {"ok": False, "error": "could not open stream - check the " "host, port, path and credentials"} deadline = time.time() + timeout frame = None while time.time() < deadline: ok, candidate = cap.read() if ok and candidate is not None and candidate.size: frame = candidate break if frame is None: return {"ok": False, "error": "connected but no frame arrived " f"within {timeout:.0f}s"} height, width = frame.shape[:2] preview = frame if max_width and width > max_width: scale = max_width / width preview = cv2.resize(frame, (max_width, int(height * scale)), interpolation=cv2.INTER_AREA) ok, buf = cv2.imencode(".jpg", preview, [int(cv2.IMWRITE_JPEG_QUALITY), 70]) return { "ok": True, "width": int(width), "height": int(height), "downscaled_to": int(preview.shape[1]) if preview is not frame else None, "fps": round(cap.get(cv2.CAP_PROP_FPS) or 0, 1), # The codec decides whether head office can ever show TRUE live # video from this camera. A browser plays H.264 everywhere; H.265 # only on some platforms, so a passthrough relay cannot rely on it # and the picture has to be re-encoded frame by frame instead. # Reported here because it is a property of the camera's settings # that an installer can usually change, and because otherwise the # only way to learn it is to read RTSP by hand — which is how this # was found: a camera whose paths end in ".264" was emitting H.265 # on both streams. "codec": _fourcc(cap), "snapshot": (base64.b64encode(buf.tobytes()).decode("ascii") if ok else None), } except (cv2.error, MemoryError, OSError) as exc: return {"ok": False, "error": f"{type(exc).__name__}: {exc}"} finally: if cap is not None: cap.release() class VideoSource(threading.Thread): def __init__(self, camera_id: str, source: "str | int", display_url: str = "", max_width: int = 1280): super().__init__(daemon=True, name=f"capture-{camera_id}") self.camera_id = camera_id self._source = source self._display_url = display_url or str(source) # Downscale at ingest: 3MP+ streams waste memory and detector time, # and on tight machines a full-res frame copy alone can OOM. self.max_width = max_width self._lock = threading.Lock() self._frame: Optional[np.ndarray] = None self._frame_ts: float = 0.0 # _stopping, NOT _stop. threading.Thread has its own private _stop(), # and join() calls it: shadowing the name with an Event made every # join() on a started worker raise "'Event' object is not callable". # It only surfaces when a camera is removed or edited at runtime, so # the engine answered 500 to every camera edit from head office while # every test using a stubbed worker passed. self._stopping = threading.Event() self.connected = False self.frames_total = 0 self.reconnects = 0 self._ever_connected = False # -- public --------------------------------------------------------- def latest(self) -> "tuple[Optional[np.ndarray], float]": with self._lock: if self._frame is None: return None, 0.0 try: return self._frame.copy(), self._frame_ts except MemoryError: return None, 0.0 def latest_since(self, known_ts: float) -> "tuple[Optional[np.ndarray], float]": """Latest frame, but only if it is newer than `known_ts`. The staleness check happens under the lock so no frame is copied just to be discarded — the worker polls far faster than the stream delivers, and a discarded full-frame copy per poll is exactly the allocation pattern that used to exhaust memory on small machines. """ with self._lock: if self._frame is None or self._frame_ts == known_ts: return None, self._frame_ts try: return self._frame.copy(), self._frame_ts except MemoryError: return None, 0.0 def stop(self) -> None: self._stopping.set() def stats(self) -> dict: return { "camera_id": self.camera_id, "url": self._display_url, "connected": self.connected, "frames_total": self.frames_total, "reconnects": self.reconnects, "last_frame_age_s": round(time.time() - self._frame_ts, 1) if self._frame_ts else None, } # -- thread --------------------------------------------------------- def run(self) -> None: backoff = 1.0 while not self._stopping.is_set(): cap = self._open() if cap is None: self.connected = False log.warning("[%s] connect failed, retrying in %.0fs (%s)", self.camera_id, backoff, self._display_url) if self._stopping.wait(backoff): break backoff = min(backoff * 2, 30.0) continue self.connected = True if self._ever_connected: # the first connect is not a reconnect self.reconnects += 1 self._ever_connected = True backoff = 1.0 log.info("[%s] connected (%s)", self.camera_id, self._display_url) while not self._stopping.is_set(): try: ok, frame = cap.read() except (cv2.error, SystemError, MemoryError): log.warning("[%s] read failed (low memory?), reconnecting", self.camera_id) break if not ok or frame is None: log.warning("[%s] stream dropped, reconnecting", self.camera_id) break try: if self.max_width and frame.shape[1] > self.max_width: scale = self.max_width / frame.shape[1] frame = cv2.resize( frame, (self.max_width, int(frame.shape[0] * scale)), interpolation=cv2.INTER_AREA) except (cv2.error, MemoryError): time.sleep(0.1) # transient allocation failure: drop frame continue with self._lock: self._frame = frame self._frame_ts = time.time() self.frames_total += 1 cap.release() self.connected = False log.info("[%s] capture stopped", self.camera_id) def _open(self) -> Optional[cv2.VideoCapture]: try: if isinstance(self._source, int): cap = cv2.VideoCapture(self._source) else: cap = cv2.VideoCapture(self._source, cv2.CAP_FFMPEG) cap.set(cv2.CAP_PROP_BUFFERSIZE, 1) if not cap.isOpened(): cap.release() return None return cap except cv2.error: log.exception("[%s] VideoCapture error", self.camera_id) return None