Reported from the first Windows install: the camera feed lags. It did, and not because of the network, the proxy or the webview. The MJPEG stream served _annotated_jpeg - the frame the pipeline had most recently FINISHED with, encoded after detection, quality scoring, tracking and identification had all run on it. On a modest shop PC that is a few frames a second, and every picture was already as old as that processing. It looked like lag because it was lag. On the fast machine it was developed on the pipeline kept up with the stream's own 10 fps cap, which is why nobody here ever saw it. Two more things compounded it. Every processed frame was JPEG-encoded whether or not a viewer existed - CPU spent on precisely the machine short of it. And ffmpeg ran its RTSP demuxer with default buffering, which holds a comfortable queue of frames before handing over the first: half a second to two seconds a live view can never recover. Now the picture and the boxes are decoupled. latest_jpeg_since takes the capture thread's freshest frame at the camera's own rate and draws the boxes from the last processed frame over it - encoded on demand, per request, so a camera nobody watches costs no encode at all. The stream sends a frame only when the camera has a newer one, capped at 15 fps; nothing is sent twice. Boxes older than a second are not drawn, so a stalled pipeline cannot leave one floating over an empty spot. _publish_annotated becomes _remember_tracks: a handful of tuples under the lock, no copy, no encode. ffmpeg gets nobuffer / low_delay / max_delay. Measured on cam2's sub-stream, same machine, ten seconds each: before 99 frames sent, 98 distinct 9.8 new pictures/s after 141 frames sent, 141 distinct 14.0 new pictures/s against a 15 fps camera, with the pipeline still processing 166 of 181 captured frames alongside - and engine CPU DOWN from 90% with no viewer to 62% with one attached. Engine version 1.0.0 -> 1.1.0 so a re-run of setup reinstalls it rather than pip deciding the requirement is already satisfied. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
276 lines
12 KiB
Python
276 lines
12 KiB
Python
"""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. stimeout: a 5s socket timeout so a dead camera is noticed.
|
|
#
|
|
# 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|fflags;nobuffer|flags;low_delay|max_delay;200000",
|
|
)
|
|
|
|
|
|
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
|