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
289 lines
7.7 KiB
Go
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)
|
|
}
|