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

289 lines
7.7 KiB
Go

package mqtt
import (
"context"
"errors"
"sync"
"testing"
"time"
"github.com/loyaly/behavision-agent/pkg/spool"
)
type fakeBroker struct {
mu sync.Mutex
connected bool
sent []string
failAfter int // fail every publish once this many have succeeded
err error
}
func (f *fakeBroker) Publish(ctx context.Context, topic string, payload []byte) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.failAfter > 0 && len(f.sent) >= f.failAfter {
if f.err != nil {
return f.err
}
return errors.New("broker refused")
}
f.sent = append(f.sent, string(payload))
return nil
}
func (f *fakeBroker) Connected() bool {
f.mu.Lock()
defer f.mu.Unlock()
return f.connected
}
func (f *fakeBroker) delivered() []string {
f.mu.Lock()
defer f.mu.Unlock()
return append([]string(nil), f.sent...)
}
func queue(t *testing.T, payloads ...string) *spool.Spool {
t.Helper()
s, err := spool.Open(t.TempDir(), 100)
if err != nil {
t.Fatal(err)
}
for _, p := range payloads {
if err := s.Append("visit", p); err != nil {
t.Fatal(err)
}
}
return s
}
func TestItDrainsInOrderAndAcks(t *testing.T) {
q := queue(t, "a", "b", "c")
b := &fakeBroker{connected: true}
p := &Pump{Queue: q, Publisher: b}
sent, err := p.drainOnce(context.Background())
if err != nil {
t.Fatal(err)
}
if sent != 3 || q.Len() != 0 {
t.Fatalf("sent %d, %d left in queue", sent, q.Len())
}
got := b.delivered()
if len(got) != 3 || got[0] != `"a"` || got[2] != `"c"` {
t.Fatalf("wrong order: %v", got)
}
}
func TestNothingIsAckedWhileTheBrokerIsDown(t *testing.T) {
// Acking an event the broker never took is how footfall disappears.
q := queue(t, "a", "b")
b := &fakeBroker{connected: false}
p := &Pump{Queue: q, Publisher: b}
if _, err := p.drainOnce(context.Background()); err == nil {
t.Fatal("a disconnected broker was treated as success")
}
if q.Len() != 2 {
t.Fatalf("events were dropped while offline: %d left", q.Len())
}
}
func TestAFailureStopsTheBatchInsteadOfSkippingPast(t *testing.T) {
// Events are a per-visitor timeline read in order; publishing around a
// stuck one would reorder a customer's visits on the server.
q := queue(t, "a", "b", "c")
b := &fakeBroker{connected: true, failAfter: 1}
p := &Pump{Queue: q, Publisher: b}
sent, err := p.drainOnce(context.Background())
if err == nil {
t.Fatal("failure not reported")
}
if sent != 1 {
t.Fatalf("sent %d, want 1 before stopping", sent)
}
if q.Len() != 2 {
t.Fatalf("%d left in queue, want the 2 unsent", q.Len())
}
// And the survivors are the RIGHT two, still in order.
rest, _ := q.Peek(10)
if string(rest[0].Payload) != `"b"` {
t.Fatalf("queue head is %s, want b", rest[0].Payload)
}
}
func TestConfirmedEventsSurviveAMidBatchFailure(t *testing.T) {
// Acking per-event rather than per-batch: a batch ack would re-send
// everything before the failure after a restart, duplicating footfall.
q := queue(t, "a", "b", "c")
b := &fakeBroker{connected: true, failAfter: 2}
p := &Pump{Queue: q, Publisher: b}
p.drainOnce(context.Background())
if q.Len() != 1 {
t.Fatalf("%d left, want only the unsent one", q.Len())
}
b.failAfter = 0
sent, err := p.drainOnce(context.Background())
if err != nil || sent != 1 {
t.Fatalf("recovery sent %d (%v)", sent, err)
}
got := b.delivered()
if len(got) != 3 {
t.Fatalf("delivered %v - duplicates or losses", got)
}
}
func TestRunRecoversWhenTheBrokerComesBack(t *testing.T) {
q := queue(t, "a")
b := &fakeBroker{connected: false}
p := &Pump{Queue: q, Publisher: b}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go p.Run(ctx)
time.Sleep(100 * time.Millisecond)
if len(b.delivered()) != 0 {
t.Fatal("published while disconnected")
}
b.mu.Lock()
b.connected = true
b.mu.Unlock()
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
if len(b.delivered()) == 1 {
return
}
time.Sleep(10 * time.Millisecond)
}
t.Fatal("queue never drained after the broker returned")
}
func TestHeartbeatIsSentSeparatelyFromTheQueue(t *testing.T) {
// "Site offline" and "nobody visited" must be distinguishable on the
// server. And a heartbeat is only meaningful now, so it is never spooled -
// otherwise a reconnecting site replays a week of "I am alive".
q := queue(t)
b := &fakeBroker{connected: true}
p := &Pump{Queue: q, Publisher: b,
HeartbeatTopic: "site/alive",
HeartbeatPayload: func() []byte { return []byte(`{"up":true}`) },
HeartbeatInterval: 20 * time.Millisecond}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go p.Run(ctx)
time.Sleep(200 * time.Millisecond)
if len(b.delivered()) == 0 {
t.Fatal("no heartbeat was sent")
}
if q.Len() != 0 {
t.Fatal("heartbeats were written to the durable queue")
}
}
func TestRunStopsPromptlyOnCancel(t *testing.T) {
q := queue(t)
b := &fakeBroker{connected: true}
p := &Pump{Queue: q, Publisher: b}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() { p.Run(ctx); close(done) }()
cancel()
select {
case <-done:
case <-time.After(3 * time.Second):
t.Fatal("Run ignored cancellation")
}
}
// ---------------------------------------------------------------- waking
// The delay this removes is on the path between a person walking in and their
// face reaching a screen, so the test asserts a real wall-clock bound rather
// than that a channel was read.
func TestAWakeDrainsWithoutWaitingOutTheIdleInterval(t *testing.T) {
q := queue(t)
pub := &fakeBroker{connected: true}
waker := NewWaker()
p := &Pump{Queue: q, Publisher: pub, Wake: waker.C(),
HeartbeatInterval: time.Hour}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go p.Run(ctx)
// Let it reach the idle wait with an empty queue first, so what follows is
// genuinely the wake path and not the drain it does on startup.
waitUntil(t, func() bool { return len(pub.delivered()) == 0 }, time.Second)
time.Sleep(50 * time.Millisecond)
start := time.Now()
if err := q.Append("visit", "e1"); err != nil {
t.Fatal(err)
}
waker.Wake()
waitUntil(t, func() bool { return len(pub.delivered()) == 1 }, 2*time.Second)
if took := time.Since(start); took >= idleInterval {
t.Fatalf("took %s - the wake did not beat the %s idle tick", took, idleInterval)
}
}
// A pump with no waker must behave exactly as it did before: a nil channel
// blocks forever in a select, which is the correct fallback, not a hang.
func TestAPumpWithNoWakerStillDrainsOnItsTimer(t *testing.T) {
q := queue(t)
pub := &fakeBroker{connected: true}
p := &Pump{Queue: q, Publisher: pub, HeartbeatInterval: time.Hour}
if err := q.Append("visit", "e1"); err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go p.Run(ctx)
waitUntil(t, func() bool { return len(pub.delivered()) == 1 }, 3*time.Second)
}
// The waker runs on the engine's webhook request. If it could ever block, a
// burst of arrivals would apply backpressure into the recognition loop.
func TestWakingNeverBlocksEvenWithNobodyListening(t *testing.T) {
waker := NewWaker()
done := make(chan struct{})
go func() {
defer close(done)
for i := 0; i < 10000; i++ {
waker.Wake()
}
}()
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("Wake blocked with no pump reading - this would stall recognition")
}
}
func TestANilWakerIsSafe(t *testing.T) {
var w *Waker
w.Wake() // an agent assembled without one must still run
if w.C() != nil {
t.Fatal("a nil waker handed out a channel")
}
}
func waitUntil(t *testing.T, cond func() bool, within time.Duration) {
t.Helper()
deadline := time.Now().Add(within)
for time.Now().Before(deadline) {
if cond() {
return
}
time.Sleep(2 * time.Millisecond)
}
t.Fatalf("condition not met within %s", within)
}