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
372 lines
13 KiB
Go
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
|
|
}
|