Files
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

391 lines
13 KiB
Go

// Package blob puts face images in S3-compatible object storage.
//
// Written against the S3 REST API directly rather than pulling in the AWS SDK:
// this needs four operations, the SDK is tens of megabytes of dependency, and
// SigV4 is a hash chain and a sorted header list. The whole protocol surface
// used here is in this file.
//
// Two rules the rest of the system depends on:
//
// 1. Objects are written PRIVATE, always. The bucket this ships against is
// world-readable at the bucket level - anonymous listing and anonymous GET
// both work on it today - so a face image written with a public ACL would
// be downloadable by anyone who guessed or listed the key. Private objects
// stay private even in that bucket; measured, not assumed.
//
// 2. Nothing outside this package ever sees a permanent URL. Reads go through
// a presigned GET that expires. A stored public URL is unrevocable, and
// these are pictures of customers' faces: "delete my data" has to mean the
// link stops working, not that we stop linking to it.
package blob
import (
"bytes"
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"strings"
"time"
)
type Config struct {
Region string // e.g. sgp1
Endpoint string // e.g. sgp1.digitaloceanspaces.com
Bucket string
AccessKey string
SecretKey string
// Prefix namespaces this deployment inside a bucket that may be shared
// with other applications. The bucket in use already holds an unrelated
// app's uploads, so writing to the root would mingle two systems'
// retention and deletion rules.
Prefix string
}
type Store struct {
cfg Config
host string
http *http.Client
}
var ErrNotConfigured = errors.New("object storage is not configured")
func New(cfg Config) (*Store, error) {
for name, v := range map[string]string{
"region": cfg.Region, "endpoint": cfg.Endpoint, "bucket": cfg.Bucket,
"access key": cfg.AccessKey, "secret key": cfg.SecretKey,
} {
if strings.TrimSpace(v) == "" {
return nil, fmt.Errorf("%w: %s is missing", ErrNotConfigured, name)
}
}
if cfg.Prefix == "" {
cfg.Prefix = "behavision"
}
return &Store{
cfg: cfg,
host: cfg.Bucket + "." + cfg.Endpoint,
// Generous: this is a cross-region upload from a 2 vCPU box, and a
// deletion that times out is a deletion that did not happen.
http: &http.Client{Timeout: 30 * time.Second},
}, nil
}
// Key builds the object key for one visit's face image.
//
// The SERVER builds this, never the agent. The site slug is taken from the
// credential the request authenticated with, so a site physically cannot write
// into another site's prefix - which is the whole reason upload goes through a
// presigned URL rather than by handing shop PCs a bucket key.
//
// Dated path segments are not decoration: retention runs by prefix, so
// "delete everything older than ninety days" is a listing over date prefixes
// rather than a scan of the whole bucket.
func (s *Store) Key(client, site, visitID string, at time.Time) string {
at = at.UTC()
return fmt.Sprintf("%s/%s/%s/%04d/%02d/%02d/%s.jpg",
s.cfg.Prefix, safeSegment(client), safeSegment(site),
at.Year(), int(at.Month()), at.Day(), safeSegment(visitID))
}
// safeSegment keeps a caller-influenced string from escaping its prefix.
//
// Slugs are already constrained by a CHECK in the schema and visit ids are
// uuids, so this should never change anything - which is exactly why it is
// cheap to keep. A single "../" reaching a key builder is a cross-tenant
// overwrite.
func safeSegment(s string) string {
var b strings.Builder
for _, c := range s {
switch {
case c >= 'a' && c <= 'z', c >= 'A' && c <= 'Z',
c >= '0' && c <= '9', c == '-', c == '_':
b.WriteRune(c)
default:
b.WriteByte('-')
}
}
out := b.String()
if out == "" {
return "unknown"
}
return out
}
// OwnsKey reports whether a key belongs to this deployment's prefix.
//
// The agent sends back the key it was given, and a compromised or buggy one
// could send any string. Without this check the server would hand out
// presigned reads for arbitrary objects in a bucket it shares with another
// application.
func (s *Store) OwnsKey(key string) bool {
return key != "" &&
strings.HasPrefix(key, s.cfg.Prefix+"/") &&
!strings.Contains(key, "..") &&
!strings.Contains(key, "//")
}
// PresignPut returns a URL the agent can PUT one image to, and nothing else.
//
// The shop PC never holds bucket credentials. A stolen or resold machine
// therefore gives up at most a few minutes of write access to one key, instead
// of read and write over a bucket that also holds another application's data.
//
// The ACL is signed into the URL, so the agent cannot choose to make the object
// public: it must send the matching header or the signature fails.
func (s *Store) PresignPut(key string, ttl time.Duration) (string, http.Header, error) {
if !s.OwnsKey(key) {
return "", nil, fmt.Errorf("refusing to presign a key outside %s/", s.cfg.Prefix)
}
hdr := http.Header{}
hdr.Set("x-amz-acl", "private")
hdr.Set("Content-Type", "image/jpeg")
u, err := s.presign(http.MethodPut, key, ttl, map[string]string{
"x-amz-acl": "private",
"content-type": "image/jpeg",
})
return u, hdr, err
}
// PresignGet returns a short-lived read URL.
//
// Short because it ends up in a browser's history, a screenshot and a support
// ticket. Fifteen minutes is long enough to look at a page and too short to be
// worth passing on.
func (s *Store) PresignGet(key string, ttl time.Duration) (string, error) {
if !s.OwnsKey(key) {
return "", errors.New("not an object this deployment owns")
}
return s.presign(http.MethodGet, key, ttl, nil)
}
// Delete removes an object. Used by the erasure path, where it is the whole
// point: a GDPR/DPDP deletion that leaves the face image in a bucket has not
// deleted anything.
//
// A missing object is success. Erasure must be idempotent - a retry after a
// half-finished deletion has to be able to finish, not fail forever on the
// object that already went.
func (s *Store) Delete(ctx context.Context, key string) error {
if !s.OwnsKey(key) {
return errors.New("not an object this deployment owns")
}
code, body, err := s.do(ctx, http.MethodDelete, key, nil, nil)
if err != nil {
return err
}
if code == http.StatusNoContent || code == http.StatusOK || code == http.StatusNotFound {
return nil
}
return fmt.Errorf("delete %s: HTTP %d: %s", key, code, snippet(body))
}
// Put uploads directly. The agent uses a presigned URL instead; this exists for
// the server's own writes and for the connectivity self-test, so a
// misconfiguration is found at deploy time rather than by a shop PC at 9am.
func (s *Store) Put(ctx context.Context, key string, body []byte, contentType string) error {
if !s.OwnsKey(key) {
return errors.New("not an object this deployment owns")
}
if contentType == "" {
contentType = "application/octet-stream"
}
code, resp, err := s.do(ctx, http.MethodPut, key, body, map[string]string{
"x-amz-acl": "private",
"content-type": contentType,
})
if err != nil {
return err
}
if code != http.StatusOK && code != http.StatusCreated {
return fmt.Errorf("put %s: HTTP %d: %s", key, code, snippet(resp))
}
return nil
}
// Check proves the credentials work and, more importantly, that an object
// written here is NOT publicly readable.
//
// The second half is the one worth running: the bucket this ships against is
// world-readable at the bucket level, so "the upload worked" and "the face
// image is safe" are completely different questions.
func (s *Store) Check(ctx context.Context) error {
key := s.cfg.Prefix + "/_selftest/reachability"
if err := s.Put(ctx, key, []byte("behavision self-test"), "text/plain"); err != nil {
return fmt.Errorf("cannot write to the bucket: %w", err)
}
defer s.Delete(context.WithoutCancel(ctx), key) //nolint:errcheck
req, err := http.NewRequestWithContext(ctx, http.MethodGet,
"https://"+s.host+"/"+key, nil)
if err != nil {
return err
}
resp, err := s.http.Do(req)
if err != nil {
return fmt.Errorf("cannot reach the bucket: %w", err)
}
defer resp.Body.Close()
io.Copy(io.Discard, io.LimitReader(resp.Body, 1<<10)) //nolint:errcheck
if resp.StatusCode == http.StatusOK {
return errors.New("objects written here are readable with NO credentials - " +
"face images must not go in this bucket until that is fixed")
}
return nil
}
// -- signing ---------------------------------------------------------------
func (s *Store) do(ctx context.Context, method, key string, body []byte,
extra map[string]string) (int, []byte, error) {
path := "/" + key
now := time.Now().UTC()
amzDate := now.Format("20060102T150405Z")
dateStamp := now.Format("20060102")
payloadHash := hex.EncodeToString(sha256sum(body))
headers := map[string]string{
"host": s.host,
"x-amz-content-sha256": payloadHash,
"x-amz-date": amzDate,
}
for k, v := range extra {
headers[strings.ToLower(k)] = v
}
signed := sortedKeys(headers)
var canonHeaders strings.Builder
for _, k := range signed {
canonHeaders.WriteString(k + ":" + headers[k] + "\n")
}
signedList := strings.Join(signed, ";")
canon := method + "\n" + escapePath(path) + "\n\n" +
canonHeaders.String() + "\n" + signedList + "\n" + payloadHash
scope := dateStamp + "/" + s.cfg.Region + "/s3/aws4_request"
toSign := "AWS4-HMAC-SHA256\n" + amzDate + "\n" + scope + "\n" +
hex.EncodeToString(sha256sum([]byte(canon)))
sig := hex.EncodeToString(hmacSHA256(s.signingKey(dateStamp), toSign))
var rdr io.Reader
if body != nil {
rdr = bytes.NewReader(body)
}
req, err := http.NewRequestWithContext(ctx, method, "https://"+s.host+escapePath(path), rdr)
if err != nil {
return 0, nil, err
}
for k, v := range headers {
if k == "host" {
continue
}
req.Header.Set(k, v)
}
req.Header.Set("Authorization", fmt.Sprintf(
"AWS4-HMAC-SHA256 Credential=%s/%s, SignedHeaders=%s, Signature=%s",
s.cfg.AccessKey, scope, signedList, sig))
resp, err := s.http.Do(req)
if err != nil {
return 0, nil, err
}
defer resp.Body.Close()
out, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10))
return resp.StatusCode, out, nil
}
func (s *Store) presign(method, key string, ttl time.Duration,
signedHeaders map[string]string) (string, error) {
if ttl <= 0 || ttl > 12*time.Hour {
// A presigned URL is a bearer token for one object. A day-long one
// forwarded in an email outlives every reason it was issued for.
ttl = 15 * time.Minute
}
path := escapePath("/" + key)
now := time.Now().UTC()
amzDate := now.Format("20060102T150405Z")
dateStamp := now.Format("20060102")
scope := dateStamp + "/" + s.cfg.Region + "/s3/aws4_request"
headers := map[string]string{"host": s.host}
for k, v := range signedHeaders {
headers[strings.ToLower(k)] = v
}
signed := sortedKeys(headers)
var canonHeaders strings.Builder
for _, k := range signed {
canonHeaders.WriteString(k + ":" + headers[k] + "\n")
}
signedList := strings.Join(signed, ";")
q := url.Values{}
q.Set("X-Amz-Algorithm", "AWS4-HMAC-SHA256")
q.Set("X-Amz-Credential", s.cfg.AccessKey+"/"+scope)
q.Set("X-Amz-Date", amzDate)
q.Set("X-Amz-Expires", fmt.Sprintf("%d", int(ttl.Seconds())))
q.Set("X-Amz-SignedHeaders", signedList)
canonQuery := q.Encode()
canon := method + "\n" + path + "\n" + canonQuery + "\n" +
canonHeaders.String() + "\n" + signedList + "\nUNSIGNED-PAYLOAD"
toSign := "AWS4-HMAC-SHA256\n" + amzDate + "\n" + scope + "\n" +
hex.EncodeToString(sha256sum([]byte(canon)))
sig := hex.EncodeToString(hmacSHA256(s.signingKey(dateStamp), toSign))
return "https://" + s.host + path + "?" + canonQuery + "&X-Amz-Signature=" + sig, nil
}
func (s *Store) signingKey(dateStamp string) []byte {
k := hmacSHA256([]byte("AWS4"+s.cfg.SecretKey), dateStamp)
k = hmacSHA256(k, s.cfg.Region)
k = hmacSHA256(k, "s3")
return hmacSHA256(k, "aws4_request")
}
func hmacSHA256(key []byte, msg string) []byte {
m := hmac.New(sha256.New, key)
m.Write([]byte(msg))
return m.Sum(nil)
}
func sha256sum(b []byte) []byte {
h := sha256.New()
h.Write(b)
return h.Sum(nil)
}
func sortedKeys(m map[string]string) []string {
out := make([]string, 0, len(m))
for k := range m {
out = append(out, k)
}
sort.Strings(out)
return out
}
// escapePath encodes each segment but keeps the separators. url.PathEscape on
// the whole path would turn every "/" into %2F and address one object with a
// very strange name.
func escapePath(p string) string {
parts := strings.Split(p, "/")
for i, s := range parts {
parts[i] = url.PathEscape(s)
}
return strings.Join(parts, "/")
}
func snippet(b []byte) string {
s := strings.TrimSpace(string(b))
if len(s) > 300 {
return s[:300]
}
return s
}