Every id in the schema is a uuid and stays one. What was wrong was putting one in front of a person: RecordVisit named every new customer 'Visitor ' || left(id::text, 8), so the arrivals feed, the shop PC and the mobile app all read "Visitor 3446ec35" - the string a shop assistant reads to a colleague and types into a search box. label is a stored column staff can overwrite and SearchVisitors matches on, so formatting around it in a front end would have left the data wrong on three surfaces. Migration 012 adds a per-client visitors.number, taken from a counter on clients with UPDATE ... RETURNING inside the visit transaction. Per client rather than global: a global sequence would tell any customer who signs up how many people the whole platform has ever seen, from their own first visitor number. The backfill numbers existing rows by first_seen_at and relabels only the eight-hex pattern the old statement produced, so a human-typed name is never overwritten. Three of the four things anyone addresses by URL already had a human name and the API simply refused it - a site has a slug, a camera has the id the engine knows it by. refs.go accepts either form anywhere an id is taken; a uuid resolves with no lookup, so every URL a client already stored keeps working. - An ambiguous camera name resolves to nothing, never to a guess: two shops may each have an "Office1" and acting on the first row would edit the wrong shop's camera. - 404 on a path, 400 on a query filter. /api/visits answered fine and it was the filter that was wrong. - site and site_id are both accepted everywhere now. They differed per endpoint, and an unknown query parameter is silently ignored, so getting it the wrong way round returned the whole estate. - The search matches V-13, which is what the product now shows. Two bugs found by running it rather than testing it: - 'Visitor ' || $2::text beside number = $2 makes Postgres deduce two types for one parameter and refuse the insert. It compiled and passed every in-memory test; the first real database rejected it, along with the existing face tests that share the path. - The fallback avatar said "V1" for Visitor 13, Visitor 10 and Visitor 15 alike, and read as the V-1 reference for a fourth person. It shows the number now. The prop is customerRef, not ref - React reserves that name and it would never have arrived. Verified on the live database and through the running API: 13 hex labels became Visitor 1-13 in first-seen order, two typed names left alone, and the same customer reachable by uuid, V-13 and 13. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
367 lines
13 KiB
Go
367 lines
13 KiB
Go
// Package store is the Postgres implementation of the ingest Store.
|
|
//
|
|
// Every statement filters or writes client_id explicitly, even where a join
|
|
// could derive it. That redundancy is the point: a cross-tenant leak then
|
|
// requires a deliberately wrong WHERE clause rather than one forgotten join
|
|
// condition.
|
|
package store
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
|
|
"github.com/loyaly/behavision-server/internal/contract"
|
|
"github.com/loyaly/behavision-server/internal/ingest"
|
|
"github.com/loyaly/behavision-server/internal/secret"
|
|
)
|
|
|
|
type Store struct {
|
|
pool *pgxpool.Pool
|
|
// secrets decrypts the few values the server must hand back out - today,
|
|
// each site's broker password. Nil until UseSecrets is called.
|
|
secrets *secret.Box
|
|
log *log.Logger
|
|
}
|
|
|
|
// UseLogger gives the store somewhere to report failures it deliberately does
|
|
// not surface to the caller, such as a refused audit write.
|
|
func (s *Store) UseLogger(l *log.Logger) { s.log = l }
|
|
|
|
func (s *Store) auditFailed(action string, err error) {
|
|
if s.log != nil {
|
|
s.log.Printf("WARN audit write failed for %s: %v", action, err)
|
|
}
|
|
}
|
|
|
|
func Open(ctx context.Context, dsn string) (*Store, error) {
|
|
cfg, err := pgxpool.ParseConfig(dsn)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("bad database url: %w", err)
|
|
}
|
|
// Small pool on purpose. This box has 2 vCPU and runs someone else's
|
|
// services; a large idle pool costs memory to no benefit at this volume.
|
|
cfg.MaxConns = 8
|
|
cfg.MinConns = 1
|
|
cfg.MaxConnLifetime = time.Hour
|
|
pool, err := pgxpool.NewWithConfig(ctx, cfg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := pool.Ping(ctx); err != nil {
|
|
pool.Close()
|
|
return nil, fmt.Errorf("database unreachable: %w", err)
|
|
}
|
|
return &Store{pool: pool}, nil
|
|
}
|
|
|
|
func (s *Store) Close() { s.pool.Close() }
|
|
|
|
// Pool exposes the connection pool for the schema migrator.
|
|
//
|
|
// Deliberately narrow in intent: the migrator has to run arbitrary DDL and
|
|
// take an advisory lock, neither of which belongs behind a typed store method.
|
|
// Nothing else should reach for this - a query that lives out here is a query
|
|
// nothing tenant-scopes.
|
|
func (s *Store) Pool() *pgxpool.Pool { return s.pool }
|
|
|
|
func (s *Store) Ping(ctx context.Context) error { return s.pool.Ping(ctx) }
|
|
|
|
// ResolveSite maps an authenticated MQTT username to a provisioned tenant.
|
|
// It only ever reads: see the comment in ingest.Consumer.Handle.
|
|
func (s *Store) ResolveSite(ctx context.Context, mqttUsername string) (ingest.Site, error) {
|
|
var site ingest.Site
|
|
err := s.pool.QueryRow(ctx, `
|
|
SELECT a.client_id::text, a.site_id::text, a.id::text, a.mqtt_username
|
|
FROM agents a
|
|
JOIN sites si ON si.id = a.site_id AND si.active
|
|
JOIN clients c ON c.id = a.client_id AND c.active
|
|
WHERE a.mqtt_username = $1`, mqttUsername).
|
|
Scan(&site.ClientID, &site.SiteID, &site.AgentID, &site.Slug)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return ingest.Site{}, ingest.ErrUnknownSite
|
|
}
|
|
if err != nil {
|
|
return ingest.Site{}, err
|
|
}
|
|
return site, nil
|
|
}
|
|
|
|
// RecordVisit writes one visit, resolving or creating the visitor.
|
|
//
|
|
// Returns inserted=false when the event was already stored. That is not an
|
|
// error: MQTT delivery is at-least-once, so a redelivery after a reconnect is
|
|
// expected, and treating it as a failure would make every reconnect look like
|
|
// an outage.
|
|
func (s *Store) RecordVisit(ctx context.Context, site ingest.Site,
|
|
v *contract.Visit) (bool, error) {
|
|
|
|
tx, err := s.pool.Begin(ctx)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed
|
|
|
|
// Claim the event id first. If it is already there, nothing else in this
|
|
// transaction should run - in particular we must not create a second
|
|
// visitor for a visit we already recorded.
|
|
var visitID string
|
|
err = tx.QueryRow(ctx, `
|
|
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at,
|
|
camera_id, is_new_visitor, similarity, quality,
|
|
attributes, image_key)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, COALESCE($9, '{}'::jsonb), $10)
|
|
ON CONFLICT (client_id, source_event_id) DO NOTHING
|
|
RETURNING id::text`,
|
|
site.ClientID, site.SiteID, v.EventID, v.OccurredAt, v.CameraID,
|
|
v.IsNew, nullFloat(v.Similarity), nullFloat(v.Quality),
|
|
v.Attributes, v.ImageKey).Scan(&visitID)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return false, tx.Commit(ctx) // already recorded
|
|
}
|
|
if err != nil {
|
|
return false, fmt.Errorf("insert visit: %w", err)
|
|
}
|
|
|
|
// Only now decide who this was. Matching is scoped to the client, never
|
|
// global: linking a face across unrelated clients would build a
|
|
// cross-company biometric tracking network.
|
|
if len(v.Embedding) == contract.EmbeddingDim {
|
|
visitorID, err := s.matchOrCreateVisitor(ctx, tx, site, v)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if _, err := tx.Exec(ctx,
|
|
`UPDATE visits SET visitor_id = $1 WHERE id = $2`,
|
|
visitorID, visitID); err != nil {
|
|
return false, err
|
|
}
|
|
if _, err := tx.Exec(ctx, `
|
|
UPDATE visitors
|
|
SET last_seen_at = GREATEST(COALESCE(last_seen_at, $2), $2),
|
|
visit_count = visit_count + 1
|
|
WHERE id = $1 AND client_id = $3`,
|
|
visitorID, v.OccurredAt, site.ClientID); err != nil {
|
|
return false, err
|
|
}
|
|
// Now that we know who this was, drop any face this server was holding
|
|
// for them from an earlier visit. Only ever one survives per person,
|
|
// which is what bounds visit_faces to the customer base rather than to
|
|
// footfall - see migration 011. A bucket deployment writes no such keys
|
|
// and this does nothing.
|
|
if err := pruneVisitorFaces(ctx, tx, site.ClientID, visitorID, visitID); err != nil {
|
|
return false, err
|
|
}
|
|
}
|
|
|
|
if _, err := tx.Exec(ctx,
|
|
`UPDATE agents SET last_event_at = now() WHERE id = $1`,
|
|
site.AgentID); err != nil {
|
|
return false, err
|
|
}
|
|
return true, tx.Commit(ctx)
|
|
}
|
|
|
|
// Thresholds mirror the edge defaults. Server-side matching answers a
|
|
// different question than the agent's - "has this person been to ANY of this
|
|
// client's sites" - but the vectors and the geometry are identical, so the
|
|
// numbers must be too. Diverging would mean two components disagreeing about
|
|
// who someone is.
|
|
const (
|
|
matchThreshold = 0.42
|
|
enrollThreshold = 0.32
|
|
reinforceThreshold = 0.55
|
|
maxEmbeddings = 5
|
|
// The server cannot know each camera's own quality gate, and every camera
|
|
// writes into ONE client-wide gallery, so it applies its own floor. Without
|
|
// it a loosely-gated camera could weld a poor view onto an identity that a
|
|
// strict camera then trusts.
|
|
minReinforceQuality = 0.45
|
|
)
|
|
|
|
func (s *Store) matchOrCreateVisitor(ctx context.Context, tx pgx.Tx,
|
|
site ingest.Site, v *contract.Visit) (string, error) {
|
|
|
|
vec := pgVector(v.Embedding)
|
|
|
|
// Exact nearest neighbour, scoped to this client and this encoder.
|
|
// `<=>` is cosine distance, so similarity is 1 - distance.
|
|
var visitorID string
|
|
var similarity float64
|
|
err := tx.QueryRow(ctx, `
|
|
SELECT e.visitor_id::text, 1 - (e.embedding <=> $1::vector) AS sim
|
|
FROM visitor_embeddings e
|
|
JOIN visitors vi ON vi.id = e.visitor_id AND vi.deleted_at IS NULL
|
|
WHERE e.client_id = $2 AND e.model = $3
|
|
ORDER BY e.embedding <=> $1::vector
|
|
LIMIT 1`, vec, site.ClientID, v.Model).Scan(&visitorID, &similarity)
|
|
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
|
|
return "", fmt.Errorf("match visitor: %w", err)
|
|
}
|
|
|
|
if err == nil && similarity >= matchThreshold {
|
|
// Known person. Consider keeping this view too.
|
|
//
|
|
// Without this an identity is born holding the single embedding from
|
|
// the first second it was ever seen, and the next encounter at a
|
|
// different angle has one reference vector to beat. That is not
|
|
// hypothetical: measured live on the Office1 camera, exactly this
|
|
// produced one person as two identities at similarity 0.304. The edge
|
|
// fixes it with reinforce_identity; the server has the same problem
|
|
// with the same cause and needs the same fix.
|
|
if err := reinforce(ctx, tx, site, v, visitorID, similarity, vec); err != nil {
|
|
return "", err
|
|
}
|
|
return visitorID, nil
|
|
}
|
|
|
|
// New person for this client.
|
|
//
|
|
// The number comes off the tenant's own counter rather than being derived
|
|
// from the uuid, because it is what a human will read, say and search for:
|
|
// "Visitor 42", not "Visitor 3446ec35". UPDATE ... RETURNING yields the
|
|
// value AFTER the update, which is what is wanted here, and it row-locks
|
|
// the client for the length of the insert so two shops cannot take the
|
|
// same number. That lock is free - this runs only for a face nobody in the
|
|
// estate has ever seen, not once per visit.
|
|
var number int64
|
|
if err := tx.QueryRow(ctx, `
|
|
UPDATE clients SET visitor_seq = visitor_seq + 1
|
|
WHERE id = $1 RETURNING visitor_seq`,
|
|
site.ClientID).Scan(&number); err != nil {
|
|
return "", fmt.Errorf("next visitor number: %w", err)
|
|
}
|
|
// The label is formatted here rather than as `'Visitor ' || $2::text` in
|
|
// the statement: reusing one parameter as a bigint and as a string operand
|
|
// makes Postgres deduce two types for it and refuse the whole insert
|
|
// ("inconsistent types deduced for parameter $2"). It compiled, it passed
|
|
// every in-memory test, and it failed on the first real database.
|
|
var newID string
|
|
if err := tx.QueryRow(ctx, `
|
|
INSERT INTO visitors (client_id, number, label, first_seen_at)
|
|
VALUES ($1, $2, $3, $4) RETURNING id::text`,
|
|
site.ClientID, number, fmt.Sprintf("Visitor %d", number),
|
|
v.OccurredAt).Scan(&newID); err != nil {
|
|
return "", fmt.Errorf("create visitor: %w", err)
|
|
}
|
|
if _, err := tx.Exec(ctx, `
|
|
INSERT INTO visitor_embeddings
|
|
(visitor_id, client_id, model, embedding, quality, source_site_id)
|
|
VALUES ($1, $2, $3, $4::vector, $5, $6)`,
|
|
newID, site.ClientID, v.Model, vec, v.Quality, site.SiteID); err != nil {
|
|
return "", fmt.Errorf("store embedding: %w", err)
|
|
}
|
|
return newID, nil
|
|
}
|
|
|
|
func (s *Store) RecordHeartbeat(ctx context.Context, site ingest.Site,
|
|
h *contract.Heartbeat) error {
|
|
up, total := 0, len(h.Cameras)
|
|
for _, ok := range h.Cameras {
|
|
if ok {
|
|
up++
|
|
}
|
|
}
|
|
// fraction_below_gate is NULL until a site has actually measured one.
|
|
// Storing 0.0 for "not reported" would read as a perfectly placed camera,
|
|
// which is the opposite of what an unmeasured site means.
|
|
var gate any
|
|
if h.FractionBelowGate > 0 {
|
|
gate = h.FractionBelowGate
|
|
}
|
|
_, err := s.pool.Exec(ctx, `
|
|
UPDATE agents
|
|
SET last_heartbeat_at = now(),
|
|
agent_version = COALESCE(NULLIF($2, ''), agent_version),
|
|
engine_version = COALESCE(NULLIF($3, ''), engine_version),
|
|
recognition_model = COALESCE(NULLIF($4, ''), recognition_model),
|
|
cameras_up = $5,
|
|
cameras_total = $6,
|
|
spool_queued = $7,
|
|
-- Never decreases. Dropped events are footfall a site permanently
|
|
-- lost; a restart that reset the agent's own counter must not make
|
|
-- that loss disappear from the report.
|
|
spool_dropped = GREATEST(spool_dropped, $8),
|
|
fraction_below_gate = COALESCE($9::real, fraction_below_gate)
|
|
WHERE id = $1`,
|
|
site.AgentID, h.AgentVersion, h.EngineVersion, h.RecognitionModel,
|
|
up, total, h.Queued, int64(h.Dropped), gate)
|
|
return err
|
|
}
|
|
|
|
// reinforce adds another view of an already-identified person.
|
|
//
|
|
// Guarded the same three ways as the edge, and for the same reasons:
|
|
//
|
|
// - Similar enough to believe it is them. Below enrollThreshold the matcher
|
|
// would call this vector a DIFFERENT person, so attaching it here would
|
|
// contradict the number driving every other decision.
|
|
// - Different enough to be worth storing. Above reinforceThreshold it is a
|
|
// near-duplicate of what we already hold and teaches the gallery nothing.
|
|
// - Good enough to keep. A blurred view welded onto an identity is
|
|
// unrecoverable; a missed hard angle is not. The risk is asymmetric, so
|
|
// the gate leans towards refusing.
|
|
//
|
|
// Capped, because an identity holding fifty vectors starts matching everyone.
|
|
func reinforce(ctx context.Context, tx pgx.Tx, site ingest.Site,
|
|
v *contract.Visit, visitorID string, similarity float64, vec string) error {
|
|
|
|
if similarity < enrollThreshold || similarity >= reinforceThreshold {
|
|
return nil
|
|
}
|
|
if v.Quality < minReinforceQuality {
|
|
return nil
|
|
}
|
|
var n int
|
|
if err := tx.QueryRow(ctx, `
|
|
SELECT count(*) FROM visitor_embeddings
|
|
WHERE visitor_id = $1 AND client_id = $2 AND model = $3`,
|
|
visitorID, site.ClientID, v.Model).Scan(&n); err != nil {
|
|
return fmt.Errorf("count embeddings: %w", err)
|
|
}
|
|
if n >= maxEmbeddings {
|
|
return nil
|
|
}
|
|
if _, err := tx.Exec(ctx, `
|
|
INSERT INTO visitor_embeddings
|
|
(visitor_id, client_id, model, embedding, quality, source_site_id)
|
|
VALUES ($1, $2, $3, $4::vector, $5, $6)`,
|
|
visitorID, site.ClientID, v.Model, vec, v.Quality, site.SiteID); err != nil {
|
|
return fmt.Errorf("reinforce: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// pgVector renders a float slice in pgvector's literal form. Built by hand
|
|
// rather than with a driver type so the store has no dependency on a pgvector
|
|
// Go package - the format is a bracketed comma list and nothing more.
|
|
func pgVector(v []float32) string {
|
|
var b strings.Builder
|
|
b.Grow(len(v) * 12)
|
|
b.WriteByte('[')
|
|
for i, f := range v {
|
|
if i > 0 {
|
|
b.WriteByte(',')
|
|
}
|
|
fmt.Fprintf(&b, "%g", f)
|
|
}
|
|
b.WriteByte(']')
|
|
return b.String()
|
|
}
|
|
|
|
// nullFloat keeps "not measured" distinct from "measured as zero". A track
|
|
// that never reached a decision has no similarity, and storing 0 would drag
|
|
// every percentile down.
|
|
func nullFloat(f float32) any {
|
|
if f == 0 {
|
|
return nil
|
|
}
|
|
return f
|
|
}
|