"""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__) # Set before OpenCV loads ffmpeg, which reads this once. # # rtsp_transport=tcp: UDP is the default and silently drops frames on lossy # Wi-Fi. # # The timeout here is NOT what bounds a dead camera, and the comment that # once said it did was wrong. Measured against OpenCV 4.11 / FFmpeg 7.1 on a # socket that accepts the connection and then says nothing: # # stimeout;5000000 -> 30.0s timeout;5000000 -> 30.0s # stimeout;2000000 -> 30.5s timeout;2000000 -> 30.4s # no timeout option at all -> 30.3s # # Identical with the option absent, so it is not being honoured under either # name through this path. `stimeout` was renamed `timeout` in FFmpeg 5.0, and # neither reaches the RTSP protocol here. What actually bounds it is # OpenCV's own interrupt callback (30s for open, 30s for read), which is a # compile-time constant we do not control. # # Both names are still set, because on a build where they DO take effect the # shorter bound is what we want and an unrecognised option is ignored. But # nothing may depend on it: a wrong address is caught by _tcp_reachable # below, in code we own, in under a second. # # fflags=nobuffer and flags=low_delay: without them ffmpeg's RTSP demuxer # holds a comfortable queue of frames before handing over the first, which # on a live feed is half a second to two seconds of latency that no amount of # work downstream can recover - the frame is already old when we get it. A # recorder wants that buffer; a live view does not. max_delay caps the # reorder wait for the same reason. os.environ.setdefault( "OPENCV_FFMPEG_CAPTURE_OPTIONS", "rtsp_transport;tcp|stimeout;5000000|timeout;5000000" "|fflags;nobuffer|flags;low_delay|max_delay;200000", ) # A stream can stay open and stop delivering. OpenCV breaks a blocked read # after 30s and we reconnect, but for those 30s `connected` is True and the # camera is dead — and a stream that trickles a frame every 20s never trips # that timeout at all, so it never reconnects and never recovers either. # # 10s is not a preference. The tracker gives up on a face after `max_misses` # (25 frames, ~1.7s at 15 fps), so by 10s every track is long gone and 150 # frames are missing: whatever this is, it is not something recognition can # work with. Reported separately from `connected` because the two need # opposite actions — one says check the network, the other says the camera # is answering but sending nothing. STALL_AFTER_S = 10.0 def _local_ipv4() -> str: """This machine's address on the interface holding the default route. A UDP socket is `connect`ed and nothing is sent - it only fixes a route so the kernel will name the source address. No packet leaves, and it needs no dependency, which matters in an engine that already ships 200 MB of models. """ import socket try: with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as s: s.settimeout(0.5) s.connect(("8.8.8.8", 80)) return s.getsockname()[0] except OSError: return "" def _wrong_network_hint(host: str) -> str: """Why a private camera address is unreachable, when that is the reason. A camera lives on the shop's LAN behind a router, and 192.168.x.x means "something on the network I am attached to" - nothing more. From mobile data, a hotel, or head office it either resolves to nobody or to a completely different device that happens to hold that number. There is no route in from the internet and there must not be: an RTSP camera reachable from outside is how a shop's cameras end up being watched by strangers. Without this the answer was "cannot reach 192.168.1.121:554 - Operation timed out", which reads as a broken camera and sends somebody to re-type an address and a password that were always correct. Asked directly by the owner, about his own cameras, from his phone's connection. Two states, two different actions, so they must not share a sentence: on the same network the camera or its address is the problem; on a different one the COMPUTER is in the wrong place and no setting will fix it. """ import ipaddress try: addr = ipaddress.ip_address(host) except ValueError: return "" # a DNS name; nothing can be concluded from the string # The RFC1918 blocks and link-local, spelled out rather than `is_private`. # That property is broader than "an address on somebody's LAN": it also # covers the carrier-grade NAT range and the documentation networks # (192.0.2, 198.51.100, 203.0.113), and telling somebody who typed one of # those that it is "on the shop's own network" would be a confident wrong # answer in the place people look first. Found by a test using 203.0.113.9 # as an example of a PUBLIC address, which `is_private` calls private. lan = any(addr in ipaddress.ip_network(n) for n in ("10.0.0.0/8", "172.16.0.0/12", "192.168.0.0/16", "169.254.0.0/16") if addr.version == 4) if not lan: return "" mine = _local_ipv4() if not mine: return (" - that is a private address, reachable only from inside " "the network the camera is on") try: same = ipaddress.ip_network(f"{mine}/24", strict=False).supernet_of( ipaddress.ip_network(f"{host}/24", strict=False)) except (ValueError, TypeError): same = False if same: return (f" - this computer is on that network ({mine}), so check the " f"camera is powered on and that {host} is its address") return (f" - this computer is on {mine}, not the camera's network. A " f"private address like {host} is only reachable from inside the " f"shop's own network, never over the internet or mobile data, so " f"recognition has to run on a computer in the shop") 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{_wrong_network_hint(parsed.hostname)}") except OSError as exc: return False, (f"cannot reach {parsed.hostname}:{port} - " f"{exc.strerror or exc}" f"{_wrong_network_hint(parsed.hostname)}") 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 # Why the last open failed, in the words an installer can act on. # Without it a camera that never connects reports only `connected: # false`, which cannot distinguish a wrong IP from a wrong password. self.last_error = "" # -- 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 stalled(self) -> bool: """Open, but not delivering. See STALL_AFTER_S.""" if not self.connected or not self._frame_ts: return False return (time.time() - self._frame_ts) > STALL_AFTER_S def stats(self) -> dict: return { "camera_id": self.camera_id, "url": self._display_url, "connected": self.connected, # Connected AND delivering. `connected` alone stays true through # a stall, so it is the wrong thing for a dashboard to colour a # camera green on. "streaming": self.connected and not self.stalled(), "stalled": self.stalled(), "last_error": self.last_error, "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]: # Pre-flight the socket, exactly as probe_source does. Without it a # camera that is off, moved or mistyped costs 30s per attempt inside # the VideoCapture constructor (measured; it is OpenCV's interrupt # timeout, not ours to shorten) — and the constructor is not # interruptible, so stop() cannot cut it short and a removed camera # leaves a daemon thread holding a socket for half a minute. A # refused or unroutable address answers in well under a second, which # is also what lets the backoff below mean what it says. reachable, why = _tcp_reachable(self._source, 2.0) if not reachable: log.debug("[%s] %s", self.camera_id, why) self.last_error = why return None 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() self.last_error = ("reachable, but the stream would not open " "- check the path and credentials") return None self.last_error = "" return cap except cv2.error: log.exception("[%s] VideoCapture error", self.camera_id) self.last_error = "VideoCapture error - see the engine log" return None