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
641 lines
29 KiB
Python
641 lines
29 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
|
|
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,
|
|
"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},
|
|
}
|
|
|
|
# -- 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
|
|
|
|
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())
|