Files
Behavision/server/internal/store/api_store.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

366 lines
13 KiB
Go

package store
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/loyaly/behavision-server/internal/api"
"github.com/loyaly/behavision-server/internal/auth"
)
// ---------------------------------------------------------------- identity
func (s *Store) UserByEmail(ctx context.Context, email string) (api.UserRecord, error) {
var u api.UserRecord
err := s.pool.QueryRow(ctx, `
SELECT u.id::text, COALESCE(u.client_id::text, ''), COALESCE(c.name, ''),
u.email, u.full_name, u.role, u.active, u.password_hash
FROM app_users u
LEFT JOIN clients c ON c.id = u.client_id
WHERE lower(u.email) = $1`, email).
Scan(&u.ID, &u.ClientID, &u.ClientName, &u.Email, &u.FullName,
&u.Role, &u.Active, &u.PasswordHash)
if errors.Is(err, pgx.ErrNoRows) {
// Not an error. The handler must still spend the same time verifying a
// password, so "no such user" has to come back as data rather than as a
// short-circuit.
return api.UserRecord{Found: false}, nil
}
if err != nil {
return api.UserRecord{}, err
}
// A user whose client has been deactivated must not be able to sign in and
// read that client's customers.
if u.ClientID != "" {
var active bool
if err := s.pool.QueryRow(ctx,
`SELECT active FROM clients WHERE id = $1`, u.ClientID).
Scan(&active); err != nil {
return api.UserRecord{}, err
}
u.Active = u.Active && active
}
u.Found = true
return u, nil
}
func (s *Store) TouchUserLogin(ctx context.Context, userID string) error {
_, err := s.pool.Exec(ctx,
`UPDATE app_users SET last_login_at = now() WHERE id = $1`, userID)
return err
}
func (s *Store) CreateSession(ctx context.Context, n api.NewSession) error {
_, err := s.pool.Exec(ctx, `
INSERT INTO sessions (user_id, client_id, access_hash, refresh_hash,
access_expires_at, refresh_expires_at, device)
VALUES ($1, NULLIF($2, '')::uuid, $3, $4, $5, $6, $7)`,
n.UserID, n.ClientID, n.AccessHash, n.RefreshHash,
n.AccessExpiry, n.RefreshExp, n.Device)
return err
}
// sessionQuery is shared by the access and refresh lookups so the two can
// never disagree about what makes a session valid.
const sessionQuery = `
SELECT s.id::text, s.user_id::text, COALESCE(s.client_id::text, ''),
COALESCE(c.name, ''), u.email, u.full_name, u.role, %s
FROM sessions s
JOIN app_users u ON u.id = s.user_id AND u.active
LEFT JOIN clients c ON c.id = s.client_id
WHERE s.%s = $1 AND s.revoked_at IS NULL`
func (s *Store) SessionByAccess(ctx context.Context, hash []byte) (auth.Principal, time.Time, error) {
// last_used_at is refreshed at most every five minutes. Writing it on every
// authenticated request would turn a read-only API call into a row update
// and a WAL record, for a column nothing needs to the second.
if _, err := s.pool.Exec(ctx, `
UPDATE sessions SET last_used_at = now()
WHERE access_hash = $1 AND revoked_at IS NULL
AND (last_used_at IS NULL OR last_used_at < now() - interval '5 minutes')`,
hash); err != nil {
// Bookkeeping. Refusing the request because a timestamp would not
// update would log everybody out over nothing.
_ = err
}
return s.session(ctx, hash, "access_expires_at", "access_hash")
}
func (s *Store) SessionByRefresh(ctx context.Context, hash []byte) (auth.Principal, time.Time, error) {
return s.session(ctx, hash, "refresh_expires_at", "refresh_hash")
}
func (s *Store) session(ctx context.Context, hash []byte, expiryCol, hashCol string) (
auth.Principal, time.Time, error) {
var p auth.Principal
var expires time.Time
err := s.pool.QueryRow(ctx,
fmt.Sprintf(sessionQuery, expiryCol, hashCol), hash).
Scan(&p.SessionID, &p.UserID, &p.ClientID, &p.ClientName,
&p.Email, &p.FullName, &p.Role, &expires)
if errors.Is(err, pgx.ErrNoRows) {
return auth.Principal{}, time.Time{}, auth.ErrNoSession
}
return p, expires, err
}
// RotateSession replaces the tokens on an existing row rather than inserting a
// new one. The old refresh token stops working the moment this commits, which
// is the point: a token copied off a resold shop PC must not keep working
// alongside the real one.
func (s *Store) RotateSession(ctx context.Context, sessionID string, n api.NewSession) error {
tag, err := s.pool.Exec(ctx, `
UPDATE sessions
SET access_hash = $2, refresh_hash = $3,
access_expires_at = $4, refresh_expires_at = $5,
last_used_at = now(),
device = COALESCE(NULLIF($6, ''), device)
WHERE id = $1 AND revoked_at IS NULL`,
sessionID, n.AccessHash, n.RefreshHash, n.AccessExpiry, n.RefreshExp, n.Device)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return auth.ErrNoSession
}
return nil
}
func (s *Store) RevokeSession(ctx context.Context, sessionID string) error {
_, err := s.pool.Exec(ctx,
`UPDATE sessions SET revoked_at = now()
WHERE id = $1 AND revoked_at IS NULL`, sessionID)
return err
}
// ---------------------------------------------------------------- reports
func nullUUID(s string) any {
if strings.TrimSpace(s) == "" {
return nil
}
return s
}
// Footfall buckets visits in the requested timezone.
//
// Two things here are easy to get wrong and expensive to notice:
//
// - "New" means first-ever, computed over all of time, not first-in-window.
// Otherwise every report re-labels your regulars as new customers the
// moment the window starts after their last visit.
// - A visit with no visitor_id (a site sending counts without templates) is
// real footfall but an unknown person. It is counted in `visitors` and in
// neither `new` nor `returning`, so those two may sum to less than the
// total. Guessing either way would put a number in a marketing report that
// nothing supports.
func (s *Store) Footfall(ctx context.Context, q api.ReportQuery) (
[]api.FootfallPoint, api.Totals, error) {
site := nullUUID(q.SiteID)
rows, err := s.pool.Query(ctx, `
WITH scoped AS (
SELECT v.visitor_id, v.occurred_at
FROM visits v
WHERE v.client_id = $1
AND v.occurred_at >= $2 AND v.occurred_at < $3
AND ($4::uuid IS NULL OR v.site_id = $4::uuid)
),
firsts AS (
SELECT v.visitor_id, min(v.occurred_at) AS first_at
FROM visits v
WHERE v.client_id = $1
AND v.visitor_id IS NOT NULL
AND ($4::uuid IS NULL OR v.site_id = $4::uuid)
GROUP BY v.visitor_id
)
SELECT date_trunc($5, s.occurred_at AT TIME ZONE $6) AS bucket,
count(DISTINCT s.visitor_id) AS identified,
count(*) FILTER (WHERE s.visitor_id IS NULL) AS anonymous,
count(DISTINCT s.visitor_id) FILTER (
WHERE date_trunc($5, f.first_at AT TIME ZONE $6)
= date_trunc($5, s.occurred_at AT TIME ZONE $6)) AS newcomers
FROM scoped s
LEFT JOIN firsts f ON f.visitor_id = s.visitor_id
GROUP BY 1
ORDER BY 1`,
q.ClientID, q.From, q.To, site, q.Bucket, q.Timezone)
if err != nil {
return nil, api.Totals{}, fmt.Errorf("footfall buckets: %w", err)
}
defer rows.Close()
var points []api.FootfallPoint
for rows.Next() {
var t time.Time
var identified, anonymous, newcomers int
if err := rows.Scan(&t, &identified, &anonymous, &newcomers); err != nil {
return nil, api.Totals{}, err
}
points = append(points, api.FootfallPoint{
// Local wall time, with no offset, because the label belongs to the
// timezone named alongside it in the report. Stamping it with Z
// would say 09:00 UTC when the shop means 09:00 in Chennai.
Bucket: t.Format("2006-01-02T15:04:05"),
Visitors: identified + anonymous,
New: newcomers,
Returning: identified - newcomers,
})
}
if err := rows.Err(); err != nil {
return nil, api.Totals{}, err
}
if points == nil {
points = []api.FootfallPoint{}
}
var totals api.Totals
var identified, anonymous int
if err := s.pool.QueryRow(ctx, `
SELECT count(DISTINCT v.visitor_id),
count(*) FILTER (WHERE v.visitor_id IS NULL),
count(*)
FROM visits v
WHERE v.client_id = $1
AND v.occurred_at >= $2 AND v.occurred_at < $3
AND ($4::uuid IS NULL OR v.site_id = $4::uuid)`,
q.ClientID, q.From, q.To, site).
Scan(&identified, &anonymous, &totals.Visits); err != nil {
return nil, api.Totals{}, fmt.Errorf("footfall totals: %w", err)
}
totals.UniqueVisitors = identified + anonymous
// Worst site, not the average. One badly placed camera is a hole in this
// report, and averaging it against three good ones hides the only site
// anyone needs to do something about.
err = s.pool.QueryRow(ctx, `
SELECT a.fraction_below_gate, si.name
FROM agents a
JOIN sites si ON si.id = a.site_id
WHERE a.client_id = $1
AND a.fraction_below_gate IS NOT NULL
AND ($2::uuid IS NULL OR a.site_id = $2::uuid)
ORDER BY a.fraction_below_gate DESC
LIMIT 1`, q.ClientID, site).
Scan(&totals.FractionBelowGate, &totals.WorstSite)
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
return nil, api.Totals{}, fmt.Errorf("gate fraction: %w", err)
}
return points, totals, nil
}
// Conversion answers "how many of the people who walked in bought something".
//
// Revenue is summed for ONE currency - whichever accounts for the most of it.
// Adding rupees to dollars produces a number that looks like money and is not,
// and this figure is the one a customer judges the product by.
func (s *Store) Conversion(ctx context.Context, q api.ReportQuery) (api.SalesReport, error) {
site := nullUUID(q.SiteID)
var rep api.SalesReport
var identified, anonymous int
if err := s.pool.QueryRow(ctx, `
SELECT count(DISTINCT v.visitor_id),
count(*) FILTER (WHERE v.visitor_id IS NULL)
FROM visits v
WHERE v.client_id = $1
AND v.occurred_at >= $2 AND v.occurred_at < $3
AND ($4::uuid IS NULL OR v.site_id = $4::uuid)`,
q.ClientID, q.From, q.To, site).Scan(&identified, &anonymous); err != nil {
return rep, fmt.Errorf("conversion visitors: %w", err)
}
rep.Visitors = identified + anonymous
var baskets int
if err := s.pool.QueryRow(ctx, `
WITH scoped AS (
SELECT p.visitor_id, p.amount, p.currency
FROM purchases p
WHERE p.client_id = $1
AND p.occurred_at >= $2 AND p.occurred_at < $3
AND ($4::uuid IS NULL OR p.site_id = $4::uuid)
),
dominant AS (
SELECT currency FROM scoped
GROUP BY currency ORDER BY sum(amount) DESC LIMIT 1
)
SELECT COALESCE((SELECT currency FROM dominant), 'INR'),
count(DISTINCT s.visitor_id),
count(*),
COALESCE(sum(s.amount), 0)::float8
FROM scoped s
WHERE s.currency = COALESCE((SELECT currency FROM dominant), 'INR')`,
q.ClientID, q.From, q.To, site).
Scan(&rep.Currency, &rep.Purchasers, &baskets, &rep.Revenue); err != nil {
return rep, fmt.Errorf("conversion purchases: %w", err)
}
if rep.Visitors > 0 {
rep.Conversion = float64(rep.Purchasers) / float64(rep.Visitors)
}
if baskets > 0 {
// Per basket, not per purchaser: a customer who bought twice in the
// window had two baskets, and averaging over people would overstate
// what a single transaction is worth.
rep.AvgBasket = rep.Revenue / float64(baskets)
}
return rep, nil
}
func (s *Store) SiteHealth(ctx context.Context, clientID string) ([]api.SiteHealth, error) {
rows, err := s.pool.Query(ctx, `
SELECT si.id::text, si.slug, si.name, si.timezone,
a.last_heartbeat_at, a.last_event_at,
COALESCE(a.recognition_model, ''), COALESCE(a.agent_version, ''),
COALESCE(a.cameras_up, 0), COALESCE(a.cameras_total, 0),
a.fraction_below_gate,
COALESCE(a.spool_queued, 0), COALESCE(a.spool_dropped, 0)
FROM sites si
LEFT JOIN agents a ON a.site_id = si.id
WHERE si.client_id = $1 AND si.active
ORDER BY si.name`, clientID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []api.SiteHealth
for rows.Next() {
var h api.SiteHealth
var beat, event *time.Time
var gate *float64
if err := rows.Scan(&h.SiteID, &h.Slug, &h.Name, &h.Timezone,
&beat, &event, &h.RecognitionModel, &h.AgentVersion,
&h.CamerasUp, &h.CamerasTotal, &gate,
&h.Queued, &h.Dropped); err != nil {
return nil, err
}
if beat != nil {
h.LastHeartbeatAt = beat.UTC().Format(time.RFC3339)
// Three missed beats. One missed beat is a dropped packet; three is
// a site that has actually gone away, and calling that out too
// eagerly trains people to ignore the indicator.
h.Online = time.Since(*beat) < 3*time.Minute
}
if event != nil {
h.LastEventAt = event.UTC().Format(time.RFC3339)
}
if gate != nil {
h.FractionBelowGate = *gate
}
out = append(out, h)
}
return out, rows.Err()
}
// Compile-time proof that the store satisfies what the API asks for. Without
// it a missing method is only discovered when main.go is wired up, which is the
// one file least covered by tests.
var _ api.Store = (*Store)(nil)