"""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())