Files
Suriyakumarvijayanayagam 62c2cc8a7b A visit is #1042, not 4cc216ca-dad3-4958-bb96-5f5a82022cf8
Every other thing in this product a person refers to already had a
readable reference: a shop is chennai, a camera cam1, a customer V-42, a
person their email. An audit of every list response found exactly one
gap, and it was the row people look at most - the arrivals feed showed a
visit as 36 hex characters.

012 argued no route takes a visit id so none was needed. That is true of
routing and false of everything else: it is what the feed shows, what a
support conversation quotes, and what somebody reading an API response
judges the product by.

Migration 014 mirrors the visitor scheme exactly - per client, so it
discloses no platform-wide volume, and beside the uuid rather than
instead of it. A stored counter is affordable on the busiest table
because visits from one tenant are already serialised by the consumer's
SetOrderMatters(true), so it adds no contention that was not already
there. A derived reference was the alternative and does not work:
several people through one door share occurred_at to the microsecond,
which is the collision 004 exists to handle.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-24 13:44:51 +05:30

372 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, `
WITH n AS (
UPDATE clients SET visit_seq = visit_seq + 1
WHERE id = $1 RETURNING visit_seq
)
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at,
camera_id, is_new_visitor, similarity, quality,
attributes, image_key, number)
SELECT $1, $2, $3, $4, $5, $6, $7, $8, COALESCE($9, '{}'::jsonb), $10, n.visit_seq
FROM n
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
}