Compare commits
1 Commits
v0.4.0-dem
...
v0.4.1-dem
| Author | SHA1 | Date | |
|---|---|---|---|
| e262fc8482 |
@@ -6,6 +6,7 @@ models finish loading without a single unguarded None dereference.
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import time
|
||||
import logging
|
||||
import secrets
|
||||
from pathlib import Path
|
||||
@@ -370,11 +371,23 @@ def create_app(engine: Engine) -> FastAPI:
|
||||
# Stop when the camera is deleted or its worker dies - otherwise a
|
||||
# removed camera leaves this generator running for the life of the
|
||||
# process, holding a reference to a worker nothing else can see.
|
||||
# Driven by the camera, not a timer: a frame goes out when the
|
||||
# capture thread has one newer than the last one sent, so nothing
|
||||
# is sent twice and nothing waits on the recognition pipeline.
|
||||
# Capped at 15 fps - the office cameras' own rate - so a viewer
|
||||
# never costs more encodes than the camera produces pictures.
|
||||
last_ts, min_gap, sent_at = 0.0, 1.0 / 15, 0.0
|
||||
while engine.workers.get(camera_id) is worker and worker.is_alive():
|
||||
jpeg = worker.latest_jpeg()
|
||||
if jpeg is not None:
|
||||
yield boundary + jpeg + b"\r\n"
|
||||
await asyncio.sleep(0.1) # ~10 fps to the browser
|
||||
now = time.time()
|
||||
if now - sent_at < min_gap:
|
||||
await asyncio.sleep(min_gap - (now - sent_at))
|
||||
continue
|
||||
jpeg, ts = worker.latest_jpeg_since(last_ts)
|
||||
if jpeg is None:
|
||||
await asyncio.sleep(0.02)
|
||||
continue
|
||||
last_ts, sent_at = ts, time.time()
|
||||
yield boundary + jpeg + b"\r\n"
|
||||
|
||||
return StreamingResponse(
|
||||
generate(),
|
||||
|
||||
@@ -17,10 +17,20 @@ 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.
|
||||
# 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"
|
||||
"OPENCV_FFMPEG_CAPTURE_OPTIONS",
|
||||
"rtsp_transport;tcp|stimeout;5000000|fflags;nobuffer|flags;low_delay|max_delay;200000",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -162,7 +162,11 @@ class CameraWorker(threading.Thread):
|
||||
# every test using a stubbed worker passed.
|
||||
self._stopping = threading.Event()
|
||||
self._lock = threading.Lock()
|
||||
self._annotated_jpeg: Optional[bytes] = None
|
||||
# What the live view draws over the freshest frame: the boxes from
|
||||
# the most recent processed frame, and when they were computed. NOT a
|
||||
# pre-rendered JPEG - see latest_jpeg for why.
|
||||
self._overlay: "list[tuple[tuple[int, int, int, int], tuple[int, int, int], str]]" = []
|
||||
self._overlay_ts = 0.0
|
||||
self._last_frame_ts = 0.0
|
||||
self._was_connected = False
|
||||
self.frames_processed = 0
|
||||
@@ -187,8 +191,46 @@ class CameraWorker(threading.Thread):
|
||||
self.source.stop()
|
||||
|
||||
def latest_jpeg(self) -> Optional[bytes]:
|
||||
jpeg, _ = self.latest_jpeg_since(0.0)
|
||||
return jpeg
|
||||
|
||||
def latest_jpeg_since(self, known_ts: float) -> "tuple[Optional[bytes], float]":
|
||||
"""The freshest captured frame with the latest boxes drawn on it, or
|
||||
(None, known_ts) if the camera has produced nothing newer.
|
||||
|
||||
The live picture is deliberately NOT the frame the pipeline last
|
||||
finished with. That version advanced only when detection, tracking and
|
||||
identification had all completed on a frame - a few times a second on a
|
||||
modest shop PC - and every picture it showed was already as old as that
|
||||
processing. It looked like lag because it was lag. Here the picture runs
|
||||
at the camera's rate off the capture thread's latest frame, and the
|
||||
boxes - which genuinely can only update at pipeline rate - are drawn
|
||||
over it from the last processed frame. Boxes may trail a fast walker by
|
||||
one pipeline period; the picture never does.
|
||||
|
||||
Encoded on demand, per request, so a camera nobody is watching pays for
|
||||
no JPEG at all. The old path encoded every processed frame whether or
|
||||
not a viewer existed - CPU spent on precisely the machine short of it.
|
||||
"""
|
||||
frame, ts = self.source.latest_since(known_ts)
|
||||
if frame is None:
|
||||
return None, known_ts
|
||||
with self._lock:
|
||||
return self._annotated_jpeg
|
||||
overlay, overlay_ts = list(self._overlay), self._overlay_ts
|
||||
# A stalled pipeline must not leave a box floating over an empty spot.
|
||||
# Older than a second and the person has walked out from under it.
|
||||
draw = overlay if (time.time() - overlay_ts) < 1.0 else []
|
||||
if draw:
|
||||
frame = frame.copy()
|
||||
for (x1, y1, x2, y2), color, text in draw:
|
||||
cv2.rectangle(frame, (x1, y1), (x2, y2), color, 2)
|
||||
if text:
|
||||
cv2.putText(frame, text, (x1, max(20, y1 - 8)),
|
||||
cv2.FONT_HERSHEY_SIMPLEX, 0.55, color, 2)
|
||||
ok, buf = cv2.imencode(".jpg", frame, [int(cv2.IMWRITE_JPEG_QUALITY), 80])
|
||||
if not ok:
|
||||
return None, known_ts
|
||||
return buf.tobytes(), ts
|
||||
|
||||
def stats(self) -> dict:
|
||||
return {
|
||||
@@ -236,7 +278,7 @@ class CameraWorker(threading.Thread):
|
||||
for track in ended:
|
||||
self._finish_track(track, ts)
|
||||
|
||||
self._publish_annotated(frame, active)
|
||||
self._remember_tracks(active)
|
||||
self.frames_processed += 1
|
||||
except Exception:
|
||||
log.exception("[%s] frame processing failed", self.cam_cfg.id)
|
||||
@@ -426,12 +468,14 @@ class CameraWorker(threading.Thread):
|
||||
track.quality, rcfg=self.rcfg):
|
||||
track.reinforcements += 1
|
||||
|
||||
def _publish_annotated(self, frame: np.ndarray, tracks: "list[Track]") -> None:
|
||||
canvas = frame.copy()
|
||||
def _remember_tracks(self, tracks: "list[Track]") -> None:
|
||||
"""Record what to draw. Cheap: a handful of tuples under the lock,
|
||||
no frame copy and no encode. The encode happens in latest_jpeg_since,
|
||||
only when somebody is looking."""
|
||||
overlay = []
|
||||
for t in tracks:
|
||||
if t.misses > 0:
|
||||
continue # only draw tracks matched in this frame
|
||||
x1, y1, x2, y2 = t.box
|
||||
if t.state == "resolved":
|
||||
color = _COLORS["known"] if t.label and not str(t.label).startswith(
|
||||
"Visitor") else _COLORS["new"]
|
||||
@@ -440,15 +484,10 @@ class CameraWorker(threading.Thread):
|
||||
color, text = _COLORS["ambiguous"], "?"
|
||||
else:
|
||||
color, text = _COLORS["pending"], ""
|
||||
cv2.rectangle(canvas, (x1, y1), (x2, y2), color, 2)
|
||||
if text:
|
||||
cv2.putText(canvas, text, (x1, max(20, y1 - 8)),
|
||||
cv2.FONT_HERSHEY_SIMPLEX, 0.55, color, 2)
|
||||
ok, buf = cv2.imencode(".jpg", canvas,
|
||||
[int(cv2.IMWRITE_JPEG_QUALITY), 80])
|
||||
if ok:
|
||||
with self._lock:
|
||||
self._annotated_jpeg = buf.tobytes()
|
||||
overlay.append((tuple(t.box), color, text))
|
||||
with self._lock:
|
||||
self._overlay = overlay
|
||||
self._overlay_ts = time.time()
|
||||
|
||||
|
||||
class Engine:
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "behavision"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
description = "Production face recognition over RTSP"
|
||||
requires-python = ">=3.10"
|
||||
dependencies = [
|
||||
|
||||
@@ -31,6 +31,7 @@ class FakeWorker:
|
||||
def stats(self): return {"camera_id": self.cam_cfg.id, "connected": True,
|
||||
"url": self.cam_cfg.safe_url()}
|
||||
def latest_jpeg(self): return None
|
||||
def latest_jpeg_since(self, known_ts): return None, known_ts
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
|
||||
144
tests/test_live_picture.py
Normal file
144
tests/test_live_picture.py
Normal file
@@ -0,0 +1,144 @@
|
||||
"""The live picture is the camera's latest frame, not the pipeline's.
|
||||
|
||||
Until this, the MJPEG stream served the frame the recognition pipeline had
|
||||
most recently *finished* - so on a shop PC where detection plus identification
|
||||
ran a few times a second, the live view ran a few times a second too, and every
|
||||
picture it showed was already as old as that processing. It read as lag
|
||||
because it was. These tests pin the decoupling: the picture comes from the
|
||||
capture thread at its own rate, the boxes come from the pipeline at theirs,
|
||||
and nothing is encoded for a camera nobody is watching.
|
||||
"""
|
||||
import time
|
||||
|
||||
import numpy as np
|
||||
|
||||
from behavision.config import CameraConfig, Config
|
||||
from behavision.engine import CameraWorker
|
||||
from behavision.tracking import Track
|
||||
|
||||
|
||||
class StillSource:
|
||||
"""A capture thread stand-in that hands out whatever frame it is given."""
|
||||
|
||||
def __init__(self):
|
||||
self.frame = None
|
||||
self.ts = 0.0
|
||||
|
||||
def set(self, frame):
|
||||
self.frame, self.ts = frame, time.time()
|
||||
|
||||
def latest(self):
|
||||
return self.frame, self.ts
|
||||
|
||||
def latest_since(self, known_ts):
|
||||
if self.frame is None or self.ts <= known_ts:
|
||||
return None, known_ts
|
||||
return self.frame, self.ts
|
||||
|
||||
def stop(self):
|
||||
pass
|
||||
|
||||
|
||||
def make_worker(tmp_path):
|
||||
cfg = Config()
|
||||
cfg.app.data_dir = tmp_path
|
||||
cam = CameraConfig(id="cam1", host="127.0.0.1", port=1, path="/none")
|
||||
w = CameraWorker(cam, cfg, detector=None, encoder=None, gallery=None,
|
||||
bus=None, attrs=None)
|
||||
w.source = StillSource()
|
||||
return w
|
||||
|
||||
|
||||
def grey(v=90):
|
||||
return np.full((120, 160, 3), v, dtype=np.uint8)
|
||||
|
||||
|
||||
def resolved_track(box=(30, 30, 90, 100), label="Priya"):
|
||||
t = Track(id=1, box=box, kps=np.zeros((5, 2), dtype=np.float32), score=0.9)
|
||||
t.state, t.label, t.similarity = "resolved", label, 0.71
|
||||
return t
|
||||
|
||||
|
||||
def test_the_picture_advances_with_the_camera_not_the_pipeline(tmp_path):
|
||||
w = make_worker(tmp_path)
|
||||
|
||||
# Nothing captured yet: nothing to show, and no encode happened.
|
||||
assert w.latest_jpeg_since(0.0) == (None, 0.0)
|
||||
|
||||
# A frame arrives from the camera. The pipeline has not touched it - and
|
||||
# the live view must not wait for it to.
|
||||
w.source.set(grey(80))
|
||||
jpeg1, ts1 = w.latest_jpeg_since(0.0)
|
||||
assert jpeg1 is not None and ts1 > 0
|
||||
|
||||
# Same frame again: the stream asks "anything newer than ts1?" and the
|
||||
# answer is no. This is what stops duplicates going down the wire.
|
||||
assert w.latest_jpeg_since(ts1) == (None, ts1)
|
||||
|
||||
# The camera produces a new frame; the pipeline still has not run.
|
||||
time.sleep(0.002)
|
||||
w.source.set(grey(160))
|
||||
jpeg2, ts2 = w.latest_jpeg_since(ts1)
|
||||
assert jpeg2 is not None and ts2 > ts1 and jpeg2 != jpeg1
|
||||
|
||||
|
||||
def test_boxes_from_the_last_processed_frame_are_drawn_on_the_fresh_one(tmp_path):
|
||||
w = make_worker(tmp_path)
|
||||
w.source.set(grey())
|
||||
plain, _ = w.latest_jpeg_since(0.0)
|
||||
|
||||
# The pipeline finishes a frame with one recognised person in it.
|
||||
w._remember_tracks([resolved_track()])
|
||||
|
||||
# The NEXT camera frame - which the pipeline has not seen - still carries
|
||||
# the box, because a person does not vanish between two frames.
|
||||
time.sleep(0.002)
|
||||
w.source.set(grey())
|
||||
boxed, _ = w.latest_jpeg_since(0.0)
|
||||
assert boxed != plain, "a resolved track should be drawn on the live picture"
|
||||
|
||||
|
||||
def test_a_stale_overlay_is_not_drawn(tmp_path):
|
||||
"""A stalled pipeline must not leave a box floating over an empty spot."""
|
||||
w = make_worker(tmp_path)
|
||||
w.source.set(grey())
|
||||
plain, _ = w.latest_jpeg_since(0.0)
|
||||
|
||||
w._remember_tracks([resolved_track()])
|
||||
# Pretend the pipeline last ran a while ago.
|
||||
with w._lock:
|
||||
w._overlay_ts = time.time() - 2.0
|
||||
|
||||
time.sleep(0.002)
|
||||
w.source.set(grey())
|
||||
fresh, _ = w.latest_jpeg_since(0.0)
|
||||
assert fresh == plain, "boxes older than a second should not be drawn"
|
||||
|
||||
|
||||
def test_only_tracks_matched_in_the_frame_are_drawn(tmp_path):
|
||||
"""A track being coasted on misses has no face under it right now."""
|
||||
w = make_worker(tmp_path)
|
||||
missed = resolved_track()
|
||||
missed.misses = 3
|
||||
w._remember_tracks([missed])
|
||||
assert w._overlay == []
|
||||
|
||||
seen = resolved_track()
|
||||
w._remember_tracks([seen])
|
||||
assert len(w._overlay) == 1
|
||||
(box, _color, text) = w._overlay[0]
|
||||
assert box == (30, 30, 90, 100) and text.startswith("Priya (")
|
||||
|
||||
|
||||
def test_remembering_tracks_does_not_encode_or_copy(tmp_path):
|
||||
"""The whole CPU argument: recording what to draw is a few tuples, and the
|
||||
frame is never touched. A shop PC with no viewer pays nothing."""
|
||||
w = make_worker(tmp_path)
|
||||
tracks = [resolved_track() for _ in range(5)]
|
||||
t0 = time.perf_counter()
|
||||
for _ in range(1000):
|
||||
w._remember_tracks(tracks)
|
||||
per_call_us = (time.perf_counter() - t0) / 1000 * 1e6
|
||||
# A JPEG encode of even a small frame is hundreds of microseconds; this
|
||||
# should be an order of magnitude under that.
|
||||
assert per_call_us < 100, f"remembering tracks took {per_call_us:.0f}us"
|
||||
Reference in New Issue
Block a user