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

218 lines
7.0 KiB
Go

// Broker client: the thin adapter behind the Publisher interface.
//
// Everything that decides *what to send and when* is in pump.go and is tested
// without a broker. This file only knows how to put bytes on a topic, which is
// why it is the one part that needs a real connection to exercise.
//
// Targets Mosquitto. No clustering, no shared subscriptions, no broker-side
// rules — a store publishes its own events under its own prefix and that is
// the whole interaction.
package mqtt
import (
"context"
"crypto/tls"
"crypto/x509"
"errors"
"fmt"
"log"
neturl "net/url"
"os"
"strings"
"time"
paho "github.com/eclipse/paho.mqtt.golang"
)
// ClientOptions configures a broker connection.
type ClientOptions struct {
// BrokerURL is tls://host:8883 in production, tcp://host:1883 for local
// testing only. Credentials and footfall must never cross the internet in
// the clear, so Connect refuses tcp:// to a non-loopback host.
BrokerURL string
ClientID string
Username string
Password string
// CAFile pins a private CA. Empty uses the system roots, which is what a
// Let's Encrypt certificate on the broker needs.
CAFile string
// InsecureSkipVerify disables certificate checking. Only ever for a
// self-signed staging box, and it is logged loudly when set, because a
// forgotten one silently removes the protection TLS was added for.
InsecureSkipVerify bool
// PublishTimeout bounds a single publish. Without it a half-open
// connection blocks the pump indefinitely and the queue grows behind it.
PublishTimeout time.Duration
Log *log.Logger
}
// Client implements Publisher.
type Client struct {
opts ClientOptions
client paho.Client
}
// NewClient dials the broker. It returns as soon as the connection is
// established; reconnection afterwards is automatic and the pump reads
// Connected() to decide whether to try.
func NewClient(opts ClientOptions) (*Client, error) {
if opts.BrokerURL == "" {
return nil, errors.New("mqtt: no broker url")
}
if opts.PublishTimeout <= 0 {
opts.PublishTimeout = 10 * time.Second
}
if err := checkTransport(opts.BrokerURL); err != nil {
return nil, err
}
po := paho.NewClientOptions().
AddBroker(opts.BrokerURL).
SetClientID(opts.ClientID).
SetUsername(opts.Username).
SetPassword(opts.Password).
// The broker holds no state for us: every event is already durable on
// our own disk, so a clean session avoids the broker queueing a
// second copy we would then have to de-duplicate.
SetCleanSession(true).
SetAutoReconnect(true).
SetConnectRetry(true).
SetConnectRetryInterval(5 * time.Second).
SetMaxReconnectInterval(2 * time.Minute).
SetKeepAlive(30 * time.Second).
SetConnectTimeout(15 * time.Second).
// Publishes must fail fast rather than pile up in memory while the
// link is down; the spool is what holds them, not the client.
SetMessageChannelDepth(1).
SetOrderMatters(true)
if strings.HasPrefix(opts.BrokerURL, "tls://") ||
strings.HasPrefix(opts.BrokerURL, "ssl://") {
cfg, err := tlsConfig(opts)
if err != nil {
return nil, err
}
po.SetTLSConfig(cfg)
}
c := &Client{opts: opts}
po.OnConnect = func(paho.Client) { c.logf("broker connected: %s", opts.BrokerURL) }
po.OnConnectionLost = func(_ paho.Client, err error) {
c.logf("broker connection lost: %v", err)
}
c.client = paho.NewClient(po)
tok := c.client.Connect()
if !tok.WaitTimeout(20 * time.Second) {
return c, fmt.Errorf("mqtt: connect to %s timed out", opts.BrokerURL)
}
if err := tok.Error(); err != nil {
return c, fmt.Errorf("mqtt: connect to %s: %w", opts.BrokerURL, err)
}
return c, nil
}
// Publish sends one message at QoS 1 and waits for the broker's PUBACK.
//
// QoS 1, not 0 or 2. At QoS 0 the broker never confirms, so the pump would ack
// and delete an event that was dropped on the wire. QoS 2 costs two extra
// round trips to remove a duplicate the server can drop itself from the event
// id — at-least-once with idempotent consumers is the cheaper contract.
func (c *Client) Publish(ctx context.Context, topic string, payload []byte) error {
if c.client == nil {
return errors.New("mqtt: no client")
}
if !c.client.IsConnected() {
return errors.New("mqtt: not connected")
}
tok := c.client.Publish(topic, 1, false, payload)
// Honour both the caller's context and a hard timeout: a half-open TCP
// connection can leave a token that never completes, which would stall the
// pump forever with the queue growing behind it.
done := make(chan struct{})
go func() { tok.Wait(); close(done) }()
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(c.opts.PublishTimeout):
return fmt.Errorf("mqtt: publish to %s timed out", topic)
case <-done:
return tok.Error()
}
}
// Connected reports whether the broker link is up.
func (c *Client) Connected() bool {
return c.client != nil && c.client.IsConnected()
}
// Close disconnects cleanly, giving in-flight publishes a moment to land.
func (c *Client) Close() {
if c.client != nil && c.client.IsConnected() {
c.client.Disconnect(1000)
}
}
func (c *Client) logf(format string, args ...any) {
if c.opts.Log != nil {
c.opts.Log.Printf(format, args...)
}
}
// checkTransport refuses plaintext MQTT to anywhere but the local machine.
//
// The payloads carry customer visit records and the connection carries the
// tenant's broker password. A tcp:// URL to a public host is not a
// configuration choice, it is a mistake, and it is one that works — which is
// exactly why it has to be rejected here rather than noticed later.
func checkTransport(raw string) error {
if !strings.HasPrefix(raw, "tcp://") && !strings.HasPrefix(raw, "mqtt://") {
return nil
}
if os.Getenv("BEHAVISION_ALLOW_PLAINTEXT_MQTT") == "1" {
return nil
}
// url.Parse, not hand-rolled splitting: an IPv6 literal is bracketed and
// full of colons, so scanning for the first ":" turns "[::1]:1883" into
// "[" and refuses a perfectly good loopback address.
u, err := neturl.Parse(raw)
if err != nil {
return fmt.Errorf("mqtt: cannot parse broker url %q: %w", raw, err)
}
switch u.Hostname() {
case "localhost", "127.0.0.1", "::1", "":
return nil
}
return fmt.Errorf("mqtt: refusing plaintext connection to %q - use tls:// "+
"(set BEHAVISION_ALLOW_PLAINTEXT_MQTT=1 only for local testing)",
u.Hostname())
}
func tlsConfig(opts ClientOptions) (*tls.Config, error) {
cfg := &tls.Config{MinVersion: tls.VersionTLS12}
if opts.InsecureSkipVerify {
cfg.InsecureSkipVerify = true
if opts.Log != nil {
opts.Log.Print("WARNING: MQTT certificate verification is DISABLED")
}
return cfg, nil
}
if opts.CAFile == "" {
return cfg, nil // system roots
}
pem, err := os.ReadFile(opts.CAFile)
if err != nil {
return nil, fmt.Errorf("mqtt: ca file: %w", err)
}
pool := x509.NewCertPool()
if !pool.AppendCertsFromPEM(pem) {
return nil, fmt.Errorf("mqtt: no certificates found in %s", opts.CAFile)
}
cfg.RootCAs = pool
return cfg, nil
}
// osWriteFile is indirected so tests can build without importing os twice.
var osWriteFile = os.WriteFile