// 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/.". 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...) } }