"""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() self._annotated_jpeg: Optional[bytes] = None 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]: with self._lock: return self._annotated_jpeg 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._publish_annotated(frame, 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 _publish_annotated(self, frame: np.ndarray, tracks: "list[Track]") -> None: canvas = frame.copy() 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"] text = f"{t.label} ({t.similarity:.2f})" elif t.state == "ambiguous": 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() 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())