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
219 lines
6.0 KiB
Go
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
|
|
}
|
|
}
|