Audited the engine for what it does when something goes wrong rather than when it goes right. Each of these left the process healthy, the dashboard green and the product not working. A gallery the running encoder cannot read. Embeddings are model-tagged, so when the fallback chain fires every vector the previous encoder wrote goes invisible: the shop keeps its customer list and recognises nobody on it, enrolling each regular a second time. Footfall stays correct, which is why nothing looks wrong. The only evidence was an INFO line reading 'gallery ready: 0 embeddings (model w600k_mbf) across 21 identities' - a sentence that states the disaster and calls it ready. Gallery.health now warns with the count of PEOPLE lost, not vectors, and carries the same numbers to /api/stats and /api/health, because a log line on a shop PC is read by nobody. Proved against the real 87-embedding gallery. Connected, and sending nothing. 'connected' meant the socket opened, so a stream that went quiet kept it true while last_frame_age_s climbed and the heartbeat told head office the camera was up. OpenCV breaks a blocked read at 30s, but a camera trickling a frame every 20s never trips that and never recovers. streaming/stalled are reported beside connected and the dashboard says live/stalled/offline - three states because offline sends you to the network and stalled says the camera is answering and sending nothing. The 5-second RTSP timeout that never existed. stimeout;5000000 carried a comment claiming it bounded a dead camera. Measured on OpenCV 4.11 / FFmpeg 7.1 against a socket that accepts and then says nothing: 30.0s with stimeout, 30.0s with timeout, 30.3s with no option at all - identical, so it was never honoured. stimeout became timeout in FFmpeg 5.0 and neither reaches the RTSP protocol through this path; the real bound is OpenCV's own interrupt constant. Replaced by the _tcp_reachable pre-flight probe_source already used, in code we own: 30.3s -> 0.00-2.02s, each naming its cause. That matters beyond speed - the VideoCapture constructor is not interruptible, so stop() could not cut it short and a camera removed from head office left a daemon thread holding a socket for half a minute. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
712 lines
32 KiB
Python
712 lines
32 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 . import capture
|
|
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._was_stalled = 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))
|
|
|
|
# A stall is not a disconnect and must not be reported as one: the
|
|
# socket is fine, the camera is answering, and nothing is arriving.
|
|
# Logged on the transition only — a per-frame warning would bury the
|
|
# one line that matters under thousands of copies of itself.
|
|
stalled = self.source.stalled()
|
|
if stalled != self._was_stalled:
|
|
self._was_stalled = stalled
|
|
if stalled:
|
|
log.warning("[%s] connected but no frame for over %.0fs - the "
|
|
"camera is answering and sending nothing",
|
|
self.cam_cfg.id, capture.STALL_AFTER_S)
|
|
else:
|
|
log.info("[%s] frames resumed", 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"),
|
|
},
|
|
# Counts, plus whether the running encoder can actually SEARCH
|
|
# them. A gallery of 21 identities that the loaded model cannot
|
|
# read is the silent version of an empty one.
|
|
"gallery": {**self.store.stats(), **self.gallery.health},
|
|
"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())
|