// 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 }