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
374 lines
18 KiB
Python
374 lines
18 KiB
Python
"""Identity resolution: match, reinforce, or auto-enroll — with hysteresis.
|
|
|
|
Three-zone decision instead of one threshold:
|
|
similarity >= match_threshold -> same person
|
|
similarity < enroll_threshold -> genuinely new person
|
|
in between -> ambiguous: do NOTHING
|
|
The ambiguous zone is what prevents both duplicate identities and wrong
|
|
merges — the two failure modes the previous system had simultaneously.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import threading
|
|
import time
|
|
from dataclasses import dataclass
|
|
from typing import Optional
|
|
|
|
import numpy as np
|
|
|
|
from ..config import RecognitionSection
|
|
from .index import VectorIndex
|
|
from .store import IdentityStore
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass
|
|
class Resolution:
|
|
kind: str # known | new | ambiguous | skipped
|
|
identity_id: Optional[int] = None
|
|
label: Optional[str] = None
|
|
similarity: float = 0.0
|
|
new_sighting: bool = False
|
|
|
|
|
|
class Gallery:
|
|
"""One gallery shared by every camera.
|
|
|
|
`cfg` here is the global recognition section — the default. Callers that
|
|
belong to a camera pass that camera's merged section as `rcfg`, because
|
|
the gates describe a view and cameras do not share one.
|
|
"""
|
|
|
|
def __init__(self, store: IdentityStore, index: VectorIndex,
|
|
cfg: RecognitionSection, model_name: str = "default"):
|
|
self.store = store
|
|
self.index = index
|
|
self.cfg = cfg
|
|
self.model_name = model_name
|
|
self._lock = threading.Lock()
|
|
self._last_sighting: dict[tuple[int, str], float] = {}
|
|
# Only embeddings produced by the active encoder enter the index;
|
|
# vectors from a different model are numerically incompatible.
|
|
ids, vecs = store.all_embeddings(index.dim, model=model_name)
|
|
index.add(ids, vecs)
|
|
self.health = self._assess(len(ids))
|
|
if self.health["stranded"]:
|
|
# Not an INFO line. The encoder fallback chain exists so a
|
|
# memory-starved box still runs, and when it fires every vector
|
|
# written by the previous encoder becomes invisible: the shop
|
|
# keeps its customer list and recognises nobody on it, greeting
|
|
# every regular as new and enrolling them a second time. Footfall
|
|
# stays right, which is exactly why nothing looks wrong. The old
|
|
# message for that state was "gallery ready: 0 embeddings".
|
|
log.warning(
|
|
"gallery: %d of %d stored embeddings were written by a "
|
|
"DIFFERENT encoder (%s) and cannot be searched - %d known "
|
|
"%s unrecognisable under the running model '%s'. Either "
|
|
"restore that model or accept that these identities start "
|
|
"over.",
|
|
self.health["stranded"], self.health["stored"],
|
|
", ".join(sorted(self.health["other_models"])),
|
|
self.health["identities_stranded"],
|
|
"person is" if self.health["identities_stranded"] == 1
|
|
else "people are",
|
|
model_name)
|
|
log.info("gallery ready: %d embeddings (model '%s') across %d "
|
|
"identities", len(ids), model_name,
|
|
store.stats()["identities"])
|
|
|
|
def _assess(self, usable: int) -> dict:
|
|
"""What share of the gallery the running encoder can actually reach.
|
|
|
|
Reported rather than merely logged, because a log line on a shop PC
|
|
is read by nobody: this travels to head office the same way
|
|
`fraction_below_gate` does, beside the number it qualifies.
|
|
"""
|
|
counts = self.store.model_counts()
|
|
stored = sum(counts.values())
|
|
others = {m: n for m, n in counts.items() if m != self.model_name}
|
|
identities = self.store.stats()["identities"]
|
|
return {
|
|
"model": self.model_name,
|
|
"stored": stored,
|
|
"usable": usable,
|
|
"stranded": sum(others.values()),
|
|
"other_models": sorted(others),
|
|
"identities": identities,
|
|
"identities_usable": self.store.identities_with_model(
|
|
self.model_name),
|
|
"identities_stranded": max(
|
|
0, identities - self.store.identities_with_model(
|
|
self.model_name)),
|
|
}
|
|
|
|
def resolve(self, embedding: np.ndarray, quality: float, camera_id: str,
|
|
ts: "float | None" = None,
|
|
attributes: "dict | None" = None,
|
|
rcfg: "RecognitionSection | None" = None) -> Resolution:
|
|
"""`rcfg` is the calling camera's merged thresholds; the gallery is
|
|
shared across cameras but the gates that decide a view are not."""
|
|
ts = ts or time.time()
|
|
cfg = rcfg or self.cfg
|
|
with self._lock:
|
|
matches = self.index.search(embedding, k=1)
|
|
top_id, top_sim = matches[0] if matches else (None, -1.0)
|
|
|
|
if top_id is not None and top_sim >= cfg.match_threshold:
|
|
ident = self.store.identity_for_embedding(top_id)
|
|
if ident is None: # index/store race — treat as ambiguous
|
|
return Resolution(kind="ambiguous", similarity=top_sim)
|
|
self._maybe_reinforce(ident["id"], embedding, quality,
|
|
top_sim, cfg)
|
|
fresh = self._record_sighting(
|
|
ident["id"], camera_id, ts, top_sim, quality, attributes)
|
|
return Resolution(kind="known", identity_id=ident["id"],
|
|
label=ident["label"], similarity=top_sim,
|
|
new_sighting=fresh)
|
|
|
|
if top_id is None or top_sim < cfg.enroll_threshold:
|
|
if not cfg.auto_enroll:
|
|
return Resolution(kind="skipped", similarity=top_sim)
|
|
if quality < cfg.min_enroll_quality:
|
|
# Not confident enough in this face to mint an identity.
|
|
return Resolution(kind="skipped", similarity=top_sim)
|
|
identity_id, label = self.store.create_auto_identity()
|
|
emb_id = self.store.add_embedding(identity_id, embedding,
|
|
quality, self.model_name)
|
|
self.index.add([emb_id], embedding.reshape(1, -1))
|
|
self._record_sighting(identity_id, camera_id, ts, 1.0, quality,
|
|
attributes)
|
|
log.info("auto-enrolled %s (quality %.2f)", label, quality)
|
|
return Resolution(kind="new", identity_id=identity_id,
|
|
label=label, similarity=top_sim,
|
|
new_sighting=True)
|
|
|
|
return Resolution(kind="ambiguous", similarity=top_sim)
|
|
|
|
def enroll(self, label: str, embeddings: "list[np.ndarray]",
|
|
quality: float = 1.0) -> int:
|
|
"""Explicit enrollment (CLI / API) with a known name."""
|
|
with self._lock:
|
|
identity_id = self.store.create_identity(label, kind="enrolled")
|
|
for emb in embeddings[: self.cfg.max_embeddings_per_identity]:
|
|
emb_id = self.store.add_embedding(identity_id, emb, quality,
|
|
self.model_name)
|
|
self.index.add([emb_id], emb.reshape(1, -1))
|
|
return identity_id
|
|
|
|
def reinforce_identity(self, identity_id: int, embedding: np.ndarray,
|
|
quality: float,
|
|
rcfg: "RecognitionSection | None" = None) -> bool:
|
|
"""Add another view of an ALREADY-identified person.
|
|
|
|
A track is resolved once and then stops contributing, so an identity
|
|
was born holding a single embedding from the first second of a visit —
|
|
and the next encounter at a different angle had one reference vector to
|
|
beat. This lets the rest of the visit fill the gallery out.
|
|
|
|
Guarded three ways: the view must still map to *this* identity (a
|
|
track that drifted onto another face must not poison the gallery), it
|
|
must be similar enough that we actually believe it is this person
|
|
(>= enroll_threshold), and different enough to be worth storing
|
|
(< reinforce_threshold).
|
|
"""
|
|
cfg = rcfg or self.cfg
|
|
with self._lock:
|
|
if quality < cfg.min_enroll_quality:
|
|
return False
|
|
if (self.store.embedding_count(identity_id)
|
|
>= cfg.max_embeddings_per_identity):
|
|
return False
|
|
matches = self.index.search(embedding, k=1)
|
|
if not matches:
|
|
return False
|
|
top_id, top_sim = matches[0]
|
|
ident = self.store.identity_for_embedding(top_id)
|
|
if ident is None or ident["id"] != identity_id:
|
|
return False # looks more like someone else - do not store
|
|
if top_sim < cfg.enroll_threshold:
|
|
# Nearest neighbour is this identity, but only barely. Below
|
|
# enroll_threshold resolve() would call this a DIFFERENT
|
|
# person, so gluing it on here would contradict the decision
|
|
# the same numbers drive everywhere else. Measured on the
|
|
# overhead camera, unfloored reinforcement gave one identity
|
|
# two vectors 0.195 apart. The risk is asymmetric: a wrong
|
|
# face welded into an identity is unrecoverable, a missed
|
|
# hard angle is not.
|
|
return False
|
|
if top_sim >= cfg.reinforce_threshold:
|
|
return False # near-duplicate of what we already have
|
|
emb_id = self.store.add_embedding(identity_id, embedding, quality,
|
|
self.model_name)
|
|
self.index.add([emb_id], embedding.reshape(1, -1))
|
|
log.debug("reinforced identity %d (sim %.3f, quality %.2f)",
|
|
identity_id, top_sim, quality)
|
|
return True
|
|
|
|
def merge_identities(self, source_id: int, target_id: int,
|
|
force: bool = False) -> "dict":
|
|
"""Fold one identity into another — the repair for a person who was
|
|
enrolled twice.
|
|
|
|
Duplicates are not a hypothetical: two views of one face can score
|
|
below `match_threshold`, and when they do the system mints a second
|
|
identity and there is no way back. Deleting one loses that person's
|
|
history; leaving both means the same customer is greeted as new.
|
|
|
|
Merging is destructive and, unlike a duplicate, *unrecoverable* — two
|
|
different people welded together cannot be separated afterwards,
|
|
because nothing records which embedding came from whom. So the two
|
|
identities must look at least plausibly alike: below
|
|
`enroll_threshold` resolve() positively asserts they are different
|
|
people, and overriding that assertion requires `force`.
|
|
|
|
Returns a dict with `ok`; on refusal `reason` says why, so the UI can
|
|
offer the override instead of failing silently.
|
|
"""
|
|
cfg = self.cfg
|
|
with self._lock:
|
|
if source_id == target_id:
|
|
return {"ok": False, "reason": "cannot merge an identity "
|
|
"into itself"}
|
|
if self.store.get_identity(source_id) is None:
|
|
return {"ok": False, "reason": f"identity {source_id} not found"}
|
|
if self.store.get_identity(target_id) is None:
|
|
return {"ok": False, "reason": f"identity {target_id} not found"}
|
|
|
|
sim, checkable = self._identity_similarity(source_id, target_id)
|
|
if not force:
|
|
if not checkable:
|
|
return {"ok": False, "similarity": None,
|
|
"reason": "no comparable embeddings (different "
|
|
"encoder model) - cannot verify these "
|
|
"are the same person"}
|
|
if sim < cfg.enroll_threshold:
|
|
return {"ok": False, "similarity": round(sim, 3),
|
|
"threshold": cfg.enroll_threshold,
|
|
"reason": "these look like different people "
|
|
f"(best similarity {sim:.3f} < "
|
|
f"{cfg.enroll_threshold})"}
|
|
|
|
result = self.store.merge_identities(
|
|
source_id, target_id, cfg.max_embeddings_per_identity)
|
|
if result is None:
|
|
return {"ok": False, "reason": "identity not found"}
|
|
# Trimmed vectors must leave the index or it keeps answering with
|
|
# embedding ids that no longer exist in SQLite.
|
|
self.index.remove(result["dropped_embeddings"])
|
|
# The per-camera sighting cooldown is keyed by identity; the
|
|
# source's keys now point at an identity that is gone.
|
|
for key in [k for k in self._last_sighting if k[0] == source_id]:
|
|
self._last_sighting.pop(key, None)
|
|
log.warning("merged identity %d into %d (%s): %d embeddings, "
|
|
"%d sightings, similarity %s%s", source_id, target_id,
|
|
result["label"], result["embeddings_moved"],
|
|
result["sightings_moved"],
|
|
f"{sim:.3f}" if checkable else "n/a",
|
|
" [FORCED]" if force else "")
|
|
result.update(ok=True, forced=force,
|
|
similarity=round(sim, 3) if checkable else None)
|
|
return result
|
|
|
|
def duplicate_candidates(self, limit: int = 20, k: int = 6
|
|
) -> "list[dict]":
|
|
"""Identity pairs that look like the same person.
|
|
|
|
Found through the index rather than an all-pairs comparison: every
|
|
stored vector asks for its `k` nearest neighbours and any that belong
|
|
to a *different* identity is evidence those two are one person. That
|
|
is O(n*k) and needs no big matrix — an all-pairs float32 matrix over
|
|
10k embeddings is 400 MB, and this runs on a box that already OOMs on
|
|
a 250 MB model.
|
|
|
|
Only pairs at or above `enroll_threshold` are reported: below it the
|
|
gallery's own numbers say these are different people, and offering
|
|
that as a suggestion would invite exactly the merge that cannot be
|
|
undone.
|
|
"""
|
|
with self._lock:
|
|
owners = self.store.embedding_owners(self.model_name)
|
|
if not owners:
|
|
return []
|
|
ids, vecs = self.store.all_embeddings(self.index.dim,
|
|
model=self.model_name)
|
|
best: dict[tuple[int, int], float] = {}
|
|
for emb_id, vec in zip(ids, vecs):
|
|
mine = owners.get(emb_id)
|
|
if mine is None:
|
|
continue
|
|
for other_id, sim in self.index.search(vec, k=k):
|
|
theirs = owners.get(other_id)
|
|
if theirs is None or theirs == mine:
|
|
continue
|
|
if sim < self.cfg.enroll_threshold:
|
|
continue
|
|
pair = (min(mine, theirs), max(mine, theirs))
|
|
if sim > best.get(pair, -1.0):
|
|
best[pair] = float(sim)
|
|
out = []
|
|
for (a, b), sim in sorted(best.items(), key=lambda kv: -kv[1])[:limit]:
|
|
ia, ib = self.store.get_identity(a), self.store.get_identity(b)
|
|
if ia is None or ib is None:
|
|
continue
|
|
out.append({
|
|
"a": {"id": a, "label": ia["label"], "kind": ia["kind"],
|
|
"sighting_count": ia["sighting_count"]},
|
|
"b": {"id": b, "label": ib["label"], "kind": ib["kind"],
|
|
"sighting_count": ib["sighting_count"]},
|
|
"similarity": round(sim, 3),
|
|
"confident": sim >= self.cfg.match_threshold})
|
|
return out
|
|
|
|
def _identity_similarity(self, a: int, b: int) -> "tuple[float, bool]":
|
|
"""Best cosine similarity between any view of `a` and any view of `b`.
|
|
|
|
Best, not mean: two identities of one person exist precisely because
|
|
their *typical* views disagree. If any pair of views agrees, that is
|
|
the evidence they are the same person.
|
|
"""
|
|
_, va = self.store.identity_embeddings(a, self.index.dim,
|
|
self.model_name)
|
|
_, vb = self.store.identity_embeddings(b, self.index.dim,
|
|
self.model_name)
|
|
if len(va) == 0 or len(vb) == 0:
|
|
return 0.0, False
|
|
return float((va @ vb.T).max()), True
|
|
|
|
def delete_identity(self, identity_id: int) -> bool:
|
|
with self._lock:
|
|
removed = self.store.delete_identity(identity_id)
|
|
self.index.remove(removed)
|
|
return bool(removed)
|
|
|
|
# -- internals ------------------------------------------------------
|
|
def _maybe_reinforce(self, identity_id: int, embedding: np.ndarray,
|
|
quality: float, similarity: float,
|
|
cfg: RecognitionSection) -> None:
|
|
"""Add an extra embedding for a known person when this view is
|
|
confidently theirs but usefully different (pose/lighting), improving
|
|
recall over time without letting the identity drift."""
|
|
if similarity >= cfg.reinforce_threshold:
|
|
return # too similar to what we already have — adds nothing
|
|
if quality < cfg.min_enroll_quality:
|
|
return
|
|
if (self.store.embedding_count(identity_id)
|
|
>= cfg.max_embeddings_per_identity):
|
|
return
|
|
emb_id = self.store.add_embedding(identity_id, embedding, quality,
|
|
self.model_name)
|
|
self.index.add([emb_id], embedding.reshape(1, -1))
|
|
|
|
def _record_sighting(self, identity_id: int, camera_id: str, ts: float,
|
|
similarity: float, quality: float,
|
|
attributes: "dict | None" = None) -> bool:
|
|
key = (identity_id, camera_id)
|
|
last = self._last_sighting.get(key, 0.0)
|
|
if ts - last < self.cfg.sighting_cooldown_seconds:
|
|
return False
|
|
self._last_sighting[key] = ts
|
|
self.store.record_sighting(identity_id, camera_id, ts, similarity,
|
|
quality, attributes)
|
|
return True
|