Files
Behavision/agent/pkg/mqtt/pump.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

219 lines
6.0 KiB
Go

// Package mqtt moves queued events to the broker.
//
// Split from the broker client on purpose: everything that decides *what to
// send and when* lives here and is testable without a broker, while the paho
// binding is a thin adapter that only knows how to put bytes on a topic.
package mqtt
import (
"context"
"errors"
"log"
"time"
"github.com/loyaly/behavision-agent/pkg/spool"
)
// Publisher is the broker, reduced to what the pump needs.
type Publisher interface {
// Publish must return nil only once the broker has confirmed receipt.
// Returning early would let the pump ack an event that never arrived.
Publish(ctx context.Context, topic string, payload []byte) error
Connected() bool
}
// Queue is the durable side, reduced likewise.
type Queue interface {
Peek(n int) ([]spool.Entry, error)
Ack(seqs ...uint64) error
Len() int
Dropped() uint64
}
const (
batchSize = 32
idleInterval = 2 * time.Second
minRetry = 1 * time.Second
maxRetry = 30 * time.Second
defaultHeartbe = 60 * time.Second
)
// Pump drains the queue into the broker and emits a heartbeat.
type Pump struct {
Queue Queue
Publisher Publisher
// Heartbeat topic. Without it "the site is offline" and "nobody visited"
// are indistinguishable on the server, which for a footfall product is a
// silent hole in the customer's report.
HeartbeatTopic string
HeartbeatPayload func() []byte
HeartbeatInterval time.Duration
Log *log.Logger
// Wake, when set, makes the pump drain immediately instead of waiting out
// idleInterval. Without it a visit that lands one millisecond after a drain
// sits on disk for two seconds before anyone is told - and that delay is on
// the path a shop screen or a mobile app sees as "how long after someone
// walks in does their face appear".
//
// A doorbell, not a queue: it carries nothing, because the pump re-reads
// the spool either way. Buffered by one and written non-blockingly, so a
// burst of arrivals cannot stall the recognition pipeline behind a pump
// that is mid-publish.
Wake <-chan struct{}
}
// Waker is the writing end of the Wake channel, held by whatever appends to the
// queue. NewWaker returns both halves so a caller cannot accidentally build one
// that blocks its own producer.
type Waker struct{ ch chan struct{} }
func NewWaker() *Waker { return &Waker{ch: make(chan struct{}, 1)} }
// Wake rings the pump. Never blocks: a full slot already means "there is work",
// which is the entire message, so a second ring adds nothing.
func (w *Waker) Wake() {
if w == nil {
return
}
select {
case w.ch <- struct{}{}:
default:
}
}
// C is the channel to hand the pump.
func (w *Waker) C() <-chan struct{} {
if w == nil {
return nil
}
return w.ch
}
// Run drains until ctx is cancelled.
func (p *Pump) Run(ctx context.Context) {
interval := p.HeartbeatInterval
if interval <= 0 {
interval = defaultHeartbe
}
beat := time.NewTicker(interval)
defer beat.Stop()
retry := minRetry
for {
if ctx.Err() != nil {
return
}
// Non-blocking, for the case where there is a backlog and the loop
// never reaches the waiting select below.
select {
case <-beat.C:
p.heartbeat(ctx)
default:
}
sent, err := p.drainOnce(ctx)
if ctx.Err() != nil {
return
}
var wait time.Duration
switch {
case err != nil:
// The broker is down or refusing. Back off rather than spinning:
// a store with no internet would otherwise burn a core all night.
p.logf("publish failed, retrying in %s: %v", retry, err)
wait = retry
if retry < maxRetry {
retry *= 2
if retry > maxRetry {
retry = maxRetry
}
}
case sent == 0:
retry = minRetry
wait = idleInterval
default:
// Something went through; there may be more waiting, so loop
// immediately rather than sleeping through a backlog.
retry = minRetry
}
if wait == 0 {
continue
}
// The heartbeat must be able to interrupt this wait. Sleeping through
// it would delay every beat by the idle interval, and on a quiet site
// the pump is idle essentially always.
// A nil Wake channel blocks forever in a select, which is exactly the
// right behaviour: an agent with no waker falls back to the timer.
select {
case <-ctx.Done():
return
case <-beat.C:
p.heartbeat(ctx)
case <-p.Wake:
// Something was queued. Loop straight round and drain it rather
// than sleeping out the rest of the idle interval.
case <-time.After(wait):
}
}
}
// drainOnce sends at most one batch and returns how many were acked.
func (p *Pump) drainOnce(ctx context.Context) (int, error) {
if !p.Publisher.Connected() {
return 0, errors.New("broker not connected")
}
entries, err := p.Queue.Peek(batchSize)
if err != nil || len(entries) == 0 {
return 0, err
}
sent := 0
for _, e := range entries {
if err := p.Publisher.Publish(ctx, e.Topic, e.Payload); err != nil {
// Stop at the first failure instead of skipping past it. Events
// are a per-visitor timeline and the server reads them in order;
// publishing around a stuck one would reorder a customer's visits.
return sent, err
}
// Acked one at a time, immediately after its own confirmation. A batch
// ack would re-send everything before a mid-batch failure on restart.
if err := p.Queue.Ack(e.Seq); err != nil {
return sent, err
}
sent++
}
return sent, nil
}
func (p *Pump) heartbeat(ctx context.Context) {
if p.HeartbeatTopic == "" || p.HeartbeatPayload == nil {
return
}
if !p.Publisher.Connected() {
return
}
// Not queued: a heartbeat is only meaningful now. Spooling them would
// replay a week of "I am alive" the moment a site reconnects.
if err := p.Publisher.Publish(ctx, p.HeartbeatTopic, p.HeartbeatPayload()); err != nil {
p.logf("heartbeat failed: %v", err)
}
}
func (p *Pump) logf(format string, args ...any) {
if p.Log != nil {
p.Log.Printf(format, args...)
}
}
func sleep(ctx context.Context, d time.Duration) bool {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return false
case <-t.C:
return true
}
}