Measured rather than guessed, and the first guess was wrong. Wall clock said H.265 decode cost 58 ms a frame; cap.read() blocks until the next frame arrives, so that was the frame interval, not work. As CPU time: decode 3.7 ms, detection 31.0 ms - and detection ran on every frame whether or not anything was in front of the camera, 6,649 of 8,634 frames with faces_seen 0 and active_tracks 0 throughout. detect_threads: OpenCV spreads a small repeated job over eight threads, costing 31.0 ms of CPU for 8.9 ms of wall. One thread costs 15.3 ms for 15.3 ms, against a 66 ms budget at 15 fps. Half the CPU for latency nothing can notice. motion_gate: a 160x90 greyscale absdiff, 0.1 ms against detection's 15. Consulted only while no track is open; forced to look every motion_max_skip frames; compared against the last frame SEARCHED so a slow drift cannot creep under the threshold; and a threshold above this camera's measured noise and far below a person, so anything ambiguous detects. tests/test_motion_gate.py pins each of those rather than the saving, including asserting the longest run of skips rather than the total - counting the total would pass a gate that slept forty frames and then looked forty times. Together 80% -> 16% of a core, detection skipped on 92% of frames. faces_seen is still 0 and the gate is not why: run directly over the same frames the detector finds nothing at threshold 0.50 either. The placement is the limit, as recorded; the CPU was being spent to rediscover that fifteen times a second. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
693 lines
31 KiB
Python
693 lines
31 KiB
Python
"""Pipeline engine: one worker thread per camera, shared models and gallery.
|
|
|
|
Per frame: detect -> score quality -> update tracker. Identity is resolved
|
|
per TRACK (once a track has enough hits and a good-enough frame), never per
|
|
frame. Ambiguous matches retry on later, better frames up to a bounded
|
|
number of attempts.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import collections
|
|
import logging
|
|
import threading
|
|
import time
|
|
from typing import Optional
|
|
|
|
import cv2
|
|
import numpy as np
|
|
|
|
from .attributes import AttributeEstimator, aggregate as aggregate_attrs
|
|
from .cameras import CameraStore
|
|
from .capture import VideoSource
|
|
from .faces import FaceOutbox
|
|
from .commission import CommissionRun
|
|
from .config import CameraConfig, Config
|
|
from .detection import FaceDetector
|
|
from .events import EmailSink, Event, EventBus, LogSink, WebhookSink
|
|
from .gallery import Gallery, IdentityStore, VectorIndex
|
|
from .geometry import align_face
|
|
from .recognition import EMBEDDING_DIM, ArcFaceEncoder, face_quality
|
|
from .tracking import IouTracker, Track
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
_COLORS = {"known": (80, 200, 80), "new": (60, 160, 255),
|
|
"pending": (160, 160, 160), "ambiguous": (60, 120, 200)}
|
|
|
|
|
|
# Terminal outcomes that mean "a person was on camera and we failed to place
|
|
# them", as opposed to "a person walked through too fast to try".
|
|
_LOST_OUTCOMES = ("rejected_quality", "gave_up_ambiguous", "ended_ambiguous")
|
|
|
|
|
|
def _track_outcome(track: Track) -> str:
|
|
"""Classify a finished track. Exactly one label per track, decided once.
|
|
|
|
Order matters: a track that exhausted its attempts because every one was
|
|
refused for quality is a *quality* failure, and reporting it as "ambiguous"
|
|
would send anyone tuning the site to the match threshold instead of to the
|
|
camera mount.
|
|
"""
|
|
if track.state == "resolved":
|
|
return "enrolled" if track.is_new else "recognized"
|
|
if track.emb_count == 0:
|
|
return "no_embedding" # never held a frame worth encoding
|
|
if track.id_attempts == 0:
|
|
return "too_brief" # left before enough evidence accumulated
|
|
if track.quality_skips:
|
|
return "rejected_quality" # face seen, too poor to mint an identity
|
|
if track.state == "gave_up":
|
|
return "gave_up_ambiguous"
|
|
return "ended_ambiguous"
|
|
|
|
|
|
def _quantile(ordered: "list[float]", frac: float) -> float:
|
|
if not ordered:
|
|
return 0.0
|
|
idx = min(len(ordered) - 1, int(round(frac * (len(ordered) - 1))))
|
|
return ordered[idx]
|
|
|
|
|
|
def _spread(values: "list[float]", gate: "float | None" = None) -> dict:
|
|
ordered = sorted(values)
|
|
out = {"n": len(ordered)}
|
|
if not ordered:
|
|
return out
|
|
out["p05"] = round(_quantile(ordered, 0.05), 3)
|
|
out["p50"] = round(_quantile(ordered, 0.50), 3)
|
|
out["p95"] = round(_quantile(ordered, 0.95), 3)
|
|
if gate is not None:
|
|
out["fraction_below_gate"] = round(
|
|
sum(1 for v in ordered if v < gate) / len(ordered), 3)
|
|
return out
|
|
|
|
|
|
class PipelineStats:
|
|
"""Per-camera tally of what became of each track.
|
|
|
|
The pipeline used to be unfalsifiable from outside: the only numbers were
|
|
frames and faces, so "nobody visited" and "every visitor was refused by the
|
|
quality gate" produced identical output, and every diagnosis meant reading
|
|
SQLite by hand. Deciding whether a site's camera placement works needs the
|
|
rejection reasons and the quality spread, not the frame count.
|
|
|
|
Written by the worker thread, read by API threads, hence the lock. The
|
|
distribution windows are bounded so a camera running for weeks cannot grow
|
|
this without limit.
|
|
"""
|
|
|
|
WINDOW = 500
|
|
|
|
def __init__(self) -> None:
|
|
self._lock = threading.Lock()
|
|
self._outcomes: "collections.Counter[str]" = collections.Counter()
|
|
self._qualities: "collections.deque[float]" = collections.deque(
|
|
maxlen=self.WINDOW)
|
|
self._similarities: "collections.deque[float]" = collections.deque(
|
|
maxlen=self.WINDOW)
|
|
self.tracks_ended = 0
|
|
|
|
def record(self, track: Track, outcome: str) -> None:
|
|
with self._lock:
|
|
self.tracks_ended += 1
|
|
self._outcomes[outcome] += 1
|
|
if track.best_quality > 0:
|
|
self._qualities.append(track.best_quality)
|
|
# Only tracks that actually reached resolve() have a similarity;
|
|
# zero from the others would drag every percentile down.
|
|
if track.id_attempts:
|
|
self._similarities.append(track.similarity)
|
|
|
|
def snapshot(self, enroll_gate: "float | None" = None) -> dict:
|
|
with self._lock:
|
|
outcomes = dict(self._outcomes)
|
|
qualities = list(self._qualities)
|
|
similarities = list(self._similarities)
|
|
ended = self.tracks_ended
|
|
return {
|
|
"tracks_ended": ended,
|
|
"outcomes": outcomes,
|
|
# fraction_below_gate is the number that says whether the
|
|
# enrollment gate is set wrong for this camera.
|
|
"best_quality": _spread(qualities, enroll_gate),
|
|
"similarity": _spread(similarities),
|
|
}
|
|
|
|
|
|
class CameraWorker(threading.Thread):
|
|
def __init__(self, cam_cfg: CameraConfig, cfg: Config, detector: FaceDetector,
|
|
encoder: ArcFaceEncoder, gallery: Gallery, bus: EventBus,
|
|
attrs: Optional[AttributeEstimator]):
|
|
super().__init__(daemon=True, name=f"worker-{cam_cfg.id}")
|
|
self.cam_cfg = cam_cfg
|
|
self.cfg = cfg
|
|
self.detector = detector
|
|
self.encoder = encoder
|
|
self.gallery = gallery
|
|
self.bus = bus
|
|
self.attrs = attrs
|
|
# Gates describe a view, so they are resolved per camera: an overhead
|
|
# corridor and an entrance camera at head height cannot share a
|
|
# quality gate, and a site has both.
|
|
self.rcfg = cfg.recognition.merged(cam_cfg.tuning)
|
|
self.source = VideoSource(cam_cfg.id, cam_cfg.source(),
|
|
cam_cfg.safe_url(), cam_cfg.max_width)
|
|
self.tracker = IouTracker(cfg.tracking.iou_threshold,
|
|
cfg.tracking.max_misses)
|
|
# _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._lock = threading.Lock()
|
|
# 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
|
|
# Motion gate state: a 160x90 greyscale thumbnail of the last frame we
|
|
# actually searched, and how many frames we have skipped since.
|
|
self._motion_prev = None
|
|
self._motion_skipped = 0
|
|
self.frames_skipped = 0
|
|
self.faces_seen = 0
|
|
self.pipeline = PipelineStats()
|
|
# One outbox per worker, all writing into the same directory. Files are
|
|
# uuid-named so two cameras resolving a visit in the same millisecond
|
|
# cannot collide.
|
|
self.faces = FaceOutbox(cfg.app.data_dir, cfg.app.store_faces)
|
|
# Set while a placement check is running. The check reads the live
|
|
# pipeline rather than a probe of its own, so what it measures is
|
|
# exactly what production will see.
|
|
self.commission: Optional[CommissionRun] = None
|
|
|
|
# -- public ---------------------------------------------------------
|
|
def start(self) -> None:
|
|
self.source.start()
|
|
super().start()
|
|
|
|
def stop(self) -> None:
|
|
self._stopping.set()
|
|
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:
|
|
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 {
|
|
**self.source.stats(),
|
|
"frames_processed": self.frames_processed,
|
|
"frames_skipped": self.frames_skipped,
|
|
"faces_seen": self.faces_seen,
|
|
"active_tracks": len(self.tracker.tracks),
|
|
"pipeline": self.pipeline.snapshot(self.rcfg.min_enroll_quality),
|
|
"gates": {"min_enroll_quality": self.rcfg.min_enroll_quality,
|
|
"match_threshold": self.rcfg.match_threshold,
|
|
"enroll_threshold": self.rcfg.enroll_threshold},
|
|
}
|
|
|
|
def _nothing_moved(self, frame) -> bool:
|
|
"""True when this frame is close enough to the last searched one that
|
|
searching it again would find the same nothing.
|
|
|
|
It cannot lose a face, and that property is what makes it acceptable
|
|
rather than merely cheap. Three guards, in order:
|
|
|
|
* the caller only asks while NO track is open, so a person already
|
|
being followed is never affected by it;
|
|
* `motion_max_skip` forces a real detection about once a second
|
|
whatever the thumbnail says, which covers a change too small or too
|
|
gradual for it - someone easing into frame at the far edge;
|
|
* the threshold sits well above measured sensor noise and well below
|
|
a person, and anything ambiguous falls through to detection. When
|
|
in doubt it looks.
|
|
|
|
Cost is 0.1 ms against detection's 15 ms, so an empty shop stops paying
|
|
for a search of an empty room ~90 times a second.
|
|
"""
|
|
import cv2 as _cv2
|
|
small = _cv2.resize(_cv2.cvtColor(frame, _cv2.COLOR_BGR2GRAY), (160, 90),
|
|
interpolation=_cv2.INTER_AREA)
|
|
prev, self._motion_prev = self._motion_prev, small
|
|
if prev is None:
|
|
return False
|
|
if self._motion_skipped >= self.cfg.app.motion_max_skip:
|
|
self._motion_skipped = 0
|
|
return False
|
|
if float(_cv2.absdiff(small, prev).mean()) >= self.cfg.app.motion_threshold:
|
|
self._motion_skipped = 0
|
|
# Keep the thumbnail we just searched against, not this one, so a
|
|
# slow drift cannot creep past the threshold one frame at a time.
|
|
return False
|
|
self._motion_prev = prev
|
|
self._motion_skipped += 1
|
|
return True
|
|
|
|
# -- thread ---------------------------------------------------------
|
|
def run(self) -> None:
|
|
tcfg = self.cfg.tracking
|
|
while not self._stopping.is_set():
|
|
try:
|
|
self._emit_connection_events()
|
|
frame, ts = self.source.latest_since(self._last_frame_ts)
|
|
if frame is None:
|
|
time.sleep(0.02)
|
|
continue
|
|
self._last_frame_ts = ts
|
|
|
|
# An empty room costs exactly as much to search as a busy one,
|
|
# and a shop is empty most of the day. Only ever while nothing
|
|
# is being tracked - see _nothing_moved.
|
|
if (self.cfg.app.motion_gate and not self.tracker.tracks
|
|
and self._nothing_moved(frame)):
|
|
self.frames_skipped += 1
|
|
self._remember_tracks([])
|
|
continue
|
|
|
|
detections = self.detector.detect(frame)
|
|
for det in detections:
|
|
det.quality = face_quality(frame, det.box, det.kps)
|
|
active, ended = self.tracker.update(detections, ts)
|
|
self.faces_seen += sum(1 for t in active if t.hits == 1)
|
|
|
|
# A placement check needs to know a face is in view even while
|
|
# its track is still open: someone standing in front of their
|
|
# own camera to test it produces no finished tracks at all.
|
|
run = self.commission
|
|
if run is not None:
|
|
run.observe(len(active), ts)
|
|
|
|
for track in active:
|
|
if self._should_identify(track, tcfg, ts):
|
|
self._identify(track, frame, ts)
|
|
|
|
# Every track ends exactly once, so this is the one place a
|
|
# per-visit outcome can be tallied without double counting.
|
|
for track in ended:
|
|
self._finish_track(track, ts)
|
|
|
|
self._remember_tracks(active)
|
|
self.frames_processed += 1
|
|
except Exception:
|
|
log.exception("[%s] frame processing failed", self.cam_cfg.id)
|
|
time.sleep(0.5)
|
|
log.info("[%s] worker stopped", self.cam_cfg.id)
|
|
|
|
# -- internals ------------------------------------------------------
|
|
def _emit_connection_events(self) -> None:
|
|
connected = self.source.connected
|
|
if connected != self._was_connected:
|
|
self._was_connected = connected
|
|
self.bus.publish(Event(
|
|
type="camera.up" if connected else "camera.down",
|
|
camera_id=self.cam_cfg.id))
|
|
|
|
def _finish_track(self, track: Track, ts: float) -> None:
|
|
"""Record what became of a track, once, as it ends.
|
|
|
|
Tracks that never reached an identity previously vanished without a
|
|
trace. For a footfall product that is a headcount which is wrong in a
|
|
way nobody can detect, and it is why a mis-set quality gate was
|
|
indistinguishable from an empty corridor.
|
|
"""
|
|
outcome = _track_outcome(track)
|
|
self.pipeline.record(track, outcome)
|
|
run = self.commission
|
|
if run is not None:
|
|
run.record(track.best_quality, ts)
|
|
if outcome not in _LOST_OUTCOMES:
|
|
return
|
|
# Only worth an event once the track held enough evidence to have been
|
|
# a real decision; a face glimpsed for two frames is noise, not a loss.
|
|
if track.emb_count < self.cfg.tracking.min_embeddings_for_id:
|
|
return
|
|
self.bus.publish(Event(
|
|
type="person.missed", camera_id=self.cam_cfg.id, ts=ts,
|
|
data={"reason": outcome,
|
|
"quality": round(track.best_quality, 3),
|
|
"similarity": round(track.similarity, 3),
|
|
"attempts": track.id_attempts,
|
|
"embeddings": track.emb_count}))
|
|
|
|
def _should_identify(self, track: Track, tcfg, ts: float) -> bool:
|
|
if track.state == "resolved":
|
|
# Known person, still on screen: keep sampling their other angles
|
|
# so the identity does not stay frozen on the one embedding it was
|
|
# created with. Gated hard on quality — a blurred frame teaches
|
|
# the gallery nothing useful.
|
|
if not tcfg.reinforce_during_track:
|
|
return False
|
|
if track.identity_id is None:
|
|
return False
|
|
if track.quality < tcfg.min_quality_to_encode:
|
|
return False
|
|
return ts - track.last_reinforce_ts >= tcfg.reinforce_interval_seconds
|
|
if track.state not in ("pending", "ambiguous"):
|
|
return False
|
|
if track.id_attempts >= tcfg.max_id_attempts:
|
|
track.state = "gave_up"
|
|
return False
|
|
# Wait for a frame worth encoding, but don't wait forever: after
|
|
# twice the warmup period, take whatever the track has.
|
|
if (track.quality < tcfg.min_quality_to_encode
|
|
and track.hits < tcfg.min_hits_for_id * 2):
|
|
return False
|
|
return True
|
|
|
|
def _identify(self, track: Track, frame: np.ndarray, ts: float) -> None:
|
|
"""Accumulate an embedding for this frame; decide identity only from
|
|
the mean of several frames. Single-frame ArcFace embeddings under
|
|
steep camera angles / motion blur differ so much that one walk-by
|
|
can look like several people — the average is stable."""
|
|
if track.state == "resolved":
|
|
self._reinforce(track, frame, ts)
|
|
return
|
|
chip = align_face(frame, track.kps, size=self.encoder.size)
|
|
# Keep the best-looking view for the customer record. Quality is
|
|
# already computed for the enrolment gate, so choosing on it costs
|
|
# nothing and picks the frame a person would have picked.
|
|
if self.faces.enabled and track.quality > track.best_face_quality:
|
|
crop = self.faces.crop(frame, track.box)
|
|
if crop is not None:
|
|
track.best_face, track.best_face_quality = crop, track.quality
|
|
if self.cfg.app.debug_faces:
|
|
debug_dir = self.cfg.app.data_dir / "debug"
|
|
debug_dir.mkdir(parents=True, exist_ok=True)
|
|
cv2.imwrite(str(debug_dir / (
|
|
f"{ts:.1f}_track{track.id}_q{track.quality:.2f}.jpg")), chip)
|
|
embedding = self.encoder.encode_chip(chip)
|
|
if embedding is None:
|
|
return
|
|
if track.emb_sum is None:
|
|
track.emb_sum = embedding.copy()
|
|
else:
|
|
track.emb_sum += embedding
|
|
track.emb_count += 1
|
|
|
|
# One more attribute sample per accumulated frame, capped. A single
|
|
# frame's age estimate swings by a decade; a few frames median out.
|
|
tcfg = self.cfg.tracking
|
|
if (self.attrs is not None
|
|
and len(track.attr_samples) < tcfg.min_embeddings_for_id
|
|
and track.quality >= self.cfg.attributes.min_quality):
|
|
track.attr_samples.append(
|
|
self.attrs.estimate(frame, track.box, chip))
|
|
|
|
if (track.emb_count < tcfg.min_embeddings_for_id
|
|
or track.hits < tcfg.min_hits_for_id):
|
|
return # keep collecting evidence
|
|
|
|
# An already-ambiguous track keeps accumulating above (that is what
|
|
# improves the mean) but only re-decides after a real interval —
|
|
# otherwise max_id_attempts is spent on consecutive frames of the
|
|
# same instant instead of on the "later, better frame" it promises.
|
|
if (track.state == "ambiguous"
|
|
and ts - track.last_attempt_ts < tcfg.id_retry_interval_seconds):
|
|
return
|
|
|
|
mean = track.emb_sum / track.emb_count
|
|
norm = float(np.linalg.norm(mean))
|
|
if norm < 1e-6:
|
|
return
|
|
mean = (mean / norm).astype(np.float32)
|
|
|
|
# Aggregated before resolve() so the sighting row carries the settled
|
|
# verdict, not whichever frame happened to be first.
|
|
if self.attrs is not None and not track.attributes:
|
|
if not track.attr_samples:
|
|
track.attr_samples.append(
|
|
self.attrs.estimate(frame, track.box, chip))
|
|
track.attributes = aggregate_attrs(track.attr_samples)
|
|
|
|
track.id_attempts += 1
|
|
track.last_attempt_ts = ts
|
|
res = self.gallery.resolve(mean, track.best_quality,
|
|
self.cam_cfg.id, ts,
|
|
attributes=track.attributes or None,
|
|
rcfg=self.rcfg)
|
|
track.similarity = res.similarity
|
|
if res.kind in ("known", "new"):
|
|
track.state = "resolved"
|
|
track.is_new = res.kind == "new"
|
|
track.identity_id = res.identity_id
|
|
track.label = res.label
|
|
if res.new_sighting:
|
|
# Written once, at the moment the visit becomes real. Writing
|
|
# per frame would fill the outbox with images of visits that
|
|
# never resolved into anything.
|
|
image_path = self.faces.save(track.best_face)
|
|
track.best_face = None # let the array go; the file has it now
|
|
self.bus.publish(Event(
|
|
type="person.new" if res.kind == "new" else "person.seen",
|
|
camera_id=self.cam_cfg.id, ts=ts,
|
|
data={"identity_id": res.identity_id, "label": res.label,
|
|
"similarity": round(res.similarity, 3),
|
|
# A local file for the agent to upload and delete.
|
|
# The engine does not upload: a shop PC must never
|
|
# hold object-storage credentials.
|
|
**({"image_path": image_path} if image_path else {}),
|
|
# The gate ran on best_quality; reporting this
|
|
# frame's quality made events look like they had
|
|
# passed a threshold they were below.
|
|
"quality": round(track.best_quality, 3),
|
|
"frame_quality": round(track.quality, 3),
|
|
**track.attributes}))
|
|
elif res.kind == "ambiguous":
|
|
track.state = "ambiguous" # retried on a later, better frame
|
|
elif res.kind == "skipped":
|
|
# resolve() declined to mint an identity — in practice always
|
|
# because best_quality is under min_enroll_quality. This branch
|
|
# did not exist: the verdict fell through, the track stayed
|
|
# "pending", and the visitor was dropped with no event, no counter
|
|
# and no log line. Marking it ambiguous also buys the retry
|
|
# throttle, so the remaining attempts are spent on genuinely later
|
|
# frames instead of being burnt in one burst on the same instant.
|
|
track.quality_skips += 1
|
|
track.state = "ambiguous"
|
|
|
|
def _reinforce(self, track: Track, frame: np.ndarray, ts: float) -> None:
|
|
"""Feed one more view of an already-identified person to the gallery."""
|
|
track.last_reinforce_ts = ts
|
|
chip = align_face(frame, track.kps, size=self.encoder.size)
|
|
embedding = self.encoder.encode_chip(chip)
|
|
if embedding is None:
|
|
return
|
|
if self.gallery.reinforce_identity(track.identity_id, embedding,
|
|
track.quality, rcfg=self.rcfg):
|
|
track.reinforcements += 1
|
|
|
|
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
|
|
if t.state == "resolved":
|
|
color = _COLORS["known"] if t.label and not str(t.label).startswith(
|
|
"Visitor") else _COLORS["new"]
|
|
text = f"{t.label} ({t.similarity:.2f})"
|
|
elif t.state == "ambiguous":
|
|
color, text = _COLORS["ambiguous"], "?"
|
|
else:
|
|
color, text = _COLORS["pending"], ""
|
|
overlay.append((tuple(t.box), color, text))
|
|
with self._lock:
|
|
self._overlay = overlay
|
|
self._overlay_ts = time.time()
|
|
|
|
|
|
class Engine:
|
|
"""Owns all shared components and one CameraWorker per camera."""
|
|
|
|
def __init__(self, cfg: Config):
|
|
self.cfg = cfg
|
|
self.bus = EventBus()
|
|
self.bus.add_sink(LogSink())
|
|
if cfg.events.webhook_url:
|
|
self.bus.add_sink(WebhookSink(cfg.events.webhook_url))
|
|
email = cfg.events.email
|
|
if email.enabled:
|
|
self.bus.add_sink(EmailSink(
|
|
email.smtp_host, email.smtp_port, email.username,
|
|
email.password, email.to, email.min_interval_seconds))
|
|
|
|
# One detector PER CAMERA. cv2.FaceDetectorYN carries mutable state
|
|
# (setInputSize + the cached input size) and is not thread-safe, so a
|
|
# shared instance races as soon as a second camera worker runs — and
|
|
# corrupts inference outright if the two streams differ in resolution.
|
|
# The YuNet model is ~230 KB, so per-worker copies are essentially free
|
|
# and avoid serialising the hottest per-frame call behind a lock.
|
|
self.detectors: "dict[str, FaceDetector]" = {}
|
|
# Shared deliberately: onnxruntime InferenceSession.run is thread-safe
|
|
# and the encoder weights are worth sharing (13-260 MB).
|
|
self.encoder = ArcFaceEncoder(cfg.app.models_dir,
|
|
cfg.recognition.model_file,
|
|
cfg.recognition.color_order)
|
|
self.store = IdentityStore(cfg.app.data_dir / "behavision.db")
|
|
self.gallery = Gallery(self.store, VectorIndex(EMBEDDING_DIM),
|
|
cfg.recognition, self.encoder.model_name)
|
|
self.attributes = None
|
|
if cfg.attributes.enabled:
|
|
est = AttributeEstimator(cfg.app.models_dir)
|
|
self.attributes = est if est.any_loaded else None
|
|
|
|
# Cameras are added and removed at runtime from the API, so this dict
|
|
# is mutated by request threads while the worker loop and stats() read
|
|
# it. RLock because add_camera/remove_camera call each other via
|
|
# restart_camera.
|
|
self._lock = threading.RLock()
|
|
self.workers: "dict[str, CameraWorker]" = {}
|
|
self.started_at: Optional[float] = None
|
|
self._running = False
|
|
|
|
# YAML seeds the store on first run; after that the store is
|
|
# authoritative, or a camera deleted in the UI would come back on the
|
|
# next restart.
|
|
self.camera_store = CameraStore(cfg.app.data_dir / "cameras.json")
|
|
self.camera_store.seed(cfg.cameras)
|
|
for cam in self.camera_store.list():
|
|
self._build_worker(cam)
|
|
|
|
# -- camera lifecycle -----------------------------------------------
|
|
def _build_worker(self, cam_cfg: CameraConfig) -> "CameraWorker":
|
|
"""Construct (but do not start) a worker and its own detector."""
|
|
det = self.cfg.detection
|
|
detector = FaceDetector(self.cfg.app.models_dir, det.score_threshold,
|
|
det.nms_threshold, det.max_faces,
|
|
det.min_face_px)
|
|
worker = CameraWorker(cam_cfg, self.cfg, detector, self.encoder,
|
|
self.gallery, self.bus, self.attributes)
|
|
self.detectors[cam_cfg.id] = detector
|
|
self.workers[cam_cfg.id] = worker
|
|
return worker
|
|
|
|
def add_camera(self, cam_cfg: CameraConfig) -> "CameraWorker":
|
|
"""Attach a camera to a live engine. Raises if the id is taken."""
|
|
with self._lock:
|
|
if cam_cfg.id in self.workers:
|
|
raise ValueError(f"camera '{cam_cfg.id}' is already running")
|
|
worker = self._build_worker(cam_cfg)
|
|
if self._running:
|
|
worker.start()
|
|
log.info("camera '%s' added (%s)", cam_cfg.id, cam_cfg.safe_url())
|
|
return worker
|
|
|
|
def remove_camera(self, camera_id: str) -> bool:
|
|
with self._lock:
|
|
worker = self.workers.pop(camera_id, None)
|
|
self.detectors.pop(camera_id, None)
|
|
if worker is None:
|
|
return False
|
|
# Outside the lock: join() can take seconds and must not block the
|
|
# frame loop's stats() calls or another camera being added.
|
|
worker.stop()
|
|
if worker.is_alive():
|
|
worker.join(timeout=5)
|
|
log.info("camera '%s' removed", camera_id)
|
|
return True
|
|
|
|
def restart_camera(self, cam_cfg: CameraConfig) -> "CameraWorker":
|
|
"""Apply an edited URL/credential. CameraWorker is a Thread, and a
|
|
stopped Thread cannot be restarted, so this must build a new one."""
|
|
with self._lock:
|
|
self.remove_camera(cam_cfg.id)
|
|
return self.add_camera(cam_cfg)
|
|
|
|
# -- lifecycle ------------------------------------------------------
|
|
def start(self) -> None:
|
|
self.bus.start()
|
|
with self._lock:
|
|
self._running = True
|
|
workers = list(self.workers.values())
|
|
for worker in workers:
|
|
worker.start()
|
|
self.started_at = time.time()
|
|
log.info("engine started with %d camera(s)", len(workers))
|
|
|
|
def stop(self) -> None:
|
|
with self._lock:
|
|
self._running = False
|
|
workers = list(self.workers.values())
|
|
for worker in workers:
|
|
worker.stop()
|
|
for worker in workers:
|
|
if worker.is_alive():
|
|
worker.join(timeout=5)
|
|
self.bus.stop()
|
|
self.store.close()
|
|
log.info("engine stopped")
|
|
|
|
def stats(self) -> dict:
|
|
return {
|
|
"uptime_s": round(time.time() - self.started_at, 1)
|
|
if self.started_at else 0,
|
|
# Which encoder actually won the fallback chain. On a
|
|
# memory-constrained box the big model can silently lose to the
|
|
# 13 MB one, and every stored embedding is tagged with whichever
|
|
# loaded — so this is the first thing to check after a deploy.
|
|
"recognition": {
|
|
"model": self.encoder.model_name,
|
|
"color_order": self.encoder.color_order,
|
|
"input_size": self.encoder.size,
|
|
},
|
|
"attributes": {
|
|
"enabled": self.attributes is not None,
|
|
"age_model": ("genderage" if self.attributes is not None
|
|
and self.attributes.has_genderage else "caffe/none"),
|
|
},
|
|
"gallery": self.store.stats(),
|
|
"cameras": [w.stats() for w in self.snapshot_workers()],
|
|
}
|
|
|
|
def snapshot_workers(self) -> "list[CameraWorker]":
|
|
"""Point-in-time copy — callers must never iterate self.workers
|
|
directly now that cameras come and go from request threads."""
|
|
with self._lock:
|
|
return list(self.workers.values())
|