Files
Behavision/agent/pkg/bridge/bridge.go
Suriyakumarvijayanayagam dad04e8cda Behavision: face recognition for retail, edge to head office
Five components that ship as one product:

- behavision/  the recognition engine. RTSP ingest, YuNet detection, IoU
               tracking, ArcFace embeddings, a FAISS/SQLite gallery, and a
               FastAPI dashboard. Identity is decided once per TRACK from an
               average of at least three embeddings, never per frame.
- agent/       the Go edge agent: supervises the engine, holds a durable
               spool, and drains it to MQTT. Nothing is acked before the
               broker confirms.
- desktop/     the shop PC application (Wails + React + tray).
- server/      the cloud API, MQTT consumer, reports and assistant.
- web/         platform.loyaly.ai, the head-office app, embedded in the
               server binary.

The gallery stores 512-float embeddings and timestamps - no images unless
`app.store_faces` is switched on. Those embeddings are biometric personal
data under GDPR and India's DPDP: template inversion reconstructs a
recognisable face from an ArcFace vector, so data/behavision.db is treated
as a biometric database and DELETE /api/visitors/{id} is a real erasure.

CLAUDE.md carries the reasoning behind every non-obvious decision here,
including the ones that were measured and the ones that were wrong first.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
2026-09-04 11:14:18 +05:30

280 lines
8.4 KiB
Go

// Package bridge turns engine detections into queued MQTT messages.
//
// This is the link that was missing: the engine detects a person and fires an
// event onto its own bus; nothing turned that into something the server would
// ever see. The engine already has a webhook sink, so the agent listens on
// loopback and points `events.webhook_url` at itself.
//
// A webhook rather than the agent polling the engine: polling would either miss
// events between polls or need cursor state the engine does not keep, and the
// sink already exists and already runs off the hot path.
package bridge
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net"
"net/http"
"strings"
"sync"
"time"
)
// Queue is the durable spool, reduced to what the bridge needs.
type Queue interface {
Append(topic string, payload any) error
}
// Embeddings fetches an identity's template from the engine.
//
// The event bus deliberately does not carry embeddings — a 512-float template
// on the bus would reach the log sink and the email sink too — so the bridge
// asks for it separately, once per identity.
type Embeddings interface {
Embedding(ctx context.Context, identityID int64) (vector []float32, model string, err error)
}
// Event is the engine's wire shape (behavision/events.py).
type Event struct {
Type string `json:"type"`
CameraID string `json:"camera_id"`
TS float64 `json:"ts"`
Data map[string]any `json:"data"`
}
type Bridge struct {
Queue Queue
Embeddings Embeddings
// TopicPrefix is "bv/<client>.<site>". The broker enforces that a site can
// only publish under its own, so an empty one means this PC is not claimed
// yet and events stay on disk rather than being addressed to nowhere.
TopicPrefix string
Log *log.Logger
// Uploader sends face images to object storage. Nil when the engine is not
// writing them, which is the default.
Uploader Uploader
// Wake, when set, is rung after a visit reaches the queue so the pump
// drains it now instead of on its next idle tick. That tick is two seconds,
// and it sits squarely on the path between a person walking in and their
// face appearing on a screen - the one delay in this chain that costs
// nothing to remove.
//
// Must not block: it runs on the engine's webhook request, so a slow pump
// would apply backpressure all the way into the recognition loop.
Wake func()
mu sync.Mutex
// Templates are fetched once per identity, not once per sighting. A
// returning customer seen forty times a day would otherwise pull the same
// 512 floats out of SQLite forty times.
seen map[int64]cached
Accepted uint64
Skipped uint64
Failed uint64
}
type cached struct {
vector []float32
model string
at time.Time
}
const cacheTTL = 30 * time.Minute
// Handler is the HTTP endpoint the engine posts to.
func (b *Bridge) Handler() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "post only", http.StatusMethodNotAllowed)
return
}
var ev Event
if err := json.NewDecoder(io.LimitReader(r.Body, 1<<20)).Decode(&ev); err != nil {
// 400, not 500: the engine must not retry a payload that will
// never parse, and its sink logs failures without blocking.
http.Error(w, "bad json", http.StatusBadRequest)
return
}
if err := b.Handle(r.Context(), ev); err != nil {
b.logf("event %s: %v", ev.Type, err)
http.Error(w, "queue failed", http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
})
return mux
}
// Handle converts one engine event and queues it.
func (b *Bridge) Handle(ctx context.Context, ev Event) error {
switch ev.Type {
case "person.new", "person.seen":
default:
// camera.up/down and person.missed are local diagnostics. They belong
// in the heartbeat, not in the footfall stream, where they would be
// counted as visits.
b.bump(&b.Skipped)
return nil
}
if b.TopicPrefix == "" {
b.bump(&b.Skipped)
return nil
}
identityID := asInt(ev.Data["identity_id"])
visit := map[string]any{
// Deterministic from what identifies the sighting, so the SAME event
// redelivered after a crash carries the SAME id and the server's
// idempotency check catches it. A random uuid here would defeat the
// entire at-least-once design.
"event_id": eventID(b.TopicPrefix, ev.CameraID, identityID, ev.TS),
"occurred_at": time.Unix(0, int64(ev.TS*float64(time.Second))).UTC(),
"camera_id": ev.CameraID,
"is_new": ev.Type == "person.new",
"similarity": asFloat(ev.Data["similarity"]),
"quality": asFloat(ev.Data["quality"]),
"local_visitor_id": identityID,
"attributes": attributes(ev.Data),
}
if identityID > 0 && b.Embeddings != nil {
vec, model, err := b.embedding(ctx, identityID)
if err != nil {
// Queue the visit anyway. A footfall count without a template is
// still a real visit; dropping it would lose the one number the
// customer is paying for over an optional field.
b.logf("no embedding for identity %d, sending counts only: %v",
identityID, err)
} else {
visit["embedding"] = vec
visit["model"] = model
}
}
// After the embedding, before the queue: the key has to be on the event
// that gets queued, and the local file is removed either way so a failed
// upload cannot leave a picture of a customer on a shop PC forever.
if path, _ := ev.Data["image_path"].(string); path != "" {
b.attachImage(ctx, visit, path)
}
if err := b.Queue.Append(b.TopicPrefix+"/visit", visit); err != nil {
b.bump(&b.Failed)
return fmt.Errorf("queue visit: %w", err)
}
b.bump(&b.Accepted)
// After the append, never before: waking a pump for an event that is not
// on disk yet is a drain that finds nothing and an event that then waits
// out the full idle interval anyway.
if b.Wake != nil {
b.Wake()
}
return nil
}
func (b *Bridge) embedding(ctx context.Context, id int64) ([]float32, string, error) {
b.mu.Lock()
if c, ok := b.seen[id]; ok && time.Since(c.at) < cacheTTL {
b.mu.Unlock()
return c.vector, c.model, nil
}
b.mu.Unlock()
vec, model, err := b.Embeddings.Embedding(ctx, id)
if err != nil {
return nil, "", err
}
b.mu.Lock()
if b.seen == nil {
b.seen = map[int64]cached{}
}
b.seen[id] = cached{vector: vec, model: model, at: time.Now()}
b.mu.Unlock()
return vec, model, nil
}
// Listen serves the webhook on loopback and returns the URL to configure in
// the engine. Port 0 so two instances on one machine cannot collide.
func (b *Bridge) Listen(ctx context.Context) (url string, stop func(), err error) {
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
return "", nil, err
}
srv := &http.Server{
Handler: b.Handler(),
ReadHeaderTimeout: 5 * time.Second,
}
go func() {
if err := srv.Serve(ln); err != nil && !errors.Is(err, http.ErrServerClosed) {
b.logf("bridge server stopped: %v", err)
}
}()
return fmt.Sprintf("http://%s/events", ln.Addr().String()), func() {
c, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
_ = srv.Shutdown(c)
}, nil
}
// eventID is stable for one sighting: same camera, same identity, same second.
//
// The engine's sighting cooldown is 30 s, so two genuinely different visits by
// one person at one camera cannot share a second. Truncating to the second
// rather than using the raw float also survives the engine re-sending after a
// restart with a marginally different timestamp.
func eventID(prefix, camera string, identity int64, ts float64) string {
return fmt.Sprintf("%s|%s|%d|%d",
strings.TrimPrefix(prefix, "bv/"), camera, identity, int64(ts))
}
// attributes keeps the estimator output and drops the bookkeeping fields the
// server already has as columns.
func attributes(data map[string]any) map[string]any {
out := map[string]any{}
for k, v := range data {
switch k {
case "identity_id", "label", "similarity", "quality", "frame_quality":
continue
}
out[k] = v
}
return out
}
func asInt(v any) int64 {
switch n := v.(type) {
case float64:
return int64(n)
case int64:
return n
case int:
return int64(n)
}
return 0
}
func asFloat(v any) float64 {
if f, ok := v.(float64); ok {
return f
}
return 0
}
func (b *Bridge) bump(p *uint64) {
b.mu.Lock()
*p++
b.mu.Unlock()
}
func (b *Bridge) logf(format string, args ...any) {
if b.Log != nil {
b.Log.Printf(format, args...)
}
}