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/secret" ) // Secrets decrypts values the server must hand back out - today, each site's // broker password. Nil until configured, and every path that needs it says so // rather than silently returning an empty credential. func (s *Store) UseSecrets(b *secret.Box) { s.secrets = b } // likePattern escapes the wildcards so a customer searching for "50%" finds // the person called "50%" instead of matching everybody. func likePattern(q string) string { r := strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`) return "%" + r.Replace(q) + "%" } // searchNumber is the customer number behind a query, or 0 for a query that is // not one. // // The reference is what staff now READ on screen - "V-13" - so it is what they // paste into the search box, and matching only `label ILIKE '%V-13%'` finds // nothing at all, because the stored label says "Visitor 13". A search that // comes back empty for the identifier the product just showed you is worse // than no search at all. func searchNumber(query string) int64 { n, ok := api.ParseVisitorRef(query) if !ok { return 0 } return n } func (s *Store) SearchVisitors(ctx context.Context, clientID, query string, limit int) ( []api.Customer, error) { rows, err := s.pool.Query(ctx, ` SELECT v.id::text, v.number, v.label, COALESCE(p.full_name, ''), COALESCE(p.phone, ''), COALESCE(p.email, ''), v.visit_count, v.first_seen_at, v.last_seen_at, (p.id IS NOT NULL), EXISTS (SELECT 1 FROM consents c WHERE c.visitor_id = v.id AND c.revoked_at IS NULL) FROM visitors v LEFT JOIN visitor_profiles p ON p.visitor_id = v.id AND p.client_id = v.client_id WHERE v.client_id = $1 AND v.deleted_at IS NULL AND ($2 = '' OR v.number = $5 OR v.label ILIKE $3 ESCAPE '\' OR p.full_name ILIKE $3 ESCAPE '\' OR p.phone ILIKE $3 ESCAPE '\' OR p.email ILIKE $3 ESCAPE '\' -- Notes are searched for ONE reason: a merge records the -- phone and name it had to discard there, and a customer -- reached by their old number is exactly who somebody is -- looking for when they type it. Retained-but-unfindable -- answers the letter of "nothing is lost" and not the point. OR p.notes ILIKE $3 ESCAPE '\') ORDER BY v.last_seen_at DESC NULLS LAST, v.first_seen_at DESC LIMIT $4`, clientID, strings.TrimSpace(query), likePattern(strings.TrimSpace(query)), limit, searchNumber(query)) if err != nil { return nil, err } defer rows.Close() var out []api.Customer for rows.Next() { var c api.Customer var first time.Time var last *time.Time var number int64 if err := rows.Scan(&c.ID, &number, &c.Label, &c.FullName, &c.Phone, &c.Email, &c.VisitCount, &first, &last, &c.HasProfile, &c.HasConsent); err != nil { return nil, err } c.Ref = api.VisitorRef(number) c.FirstSeenAt = first.UTC().Format(time.RFC3339) if last != nil { c.LastSeenAt = last.UTC().Format(time.RFC3339) } out = append(out, c) } return out, rows.Err() } func (s *Store) VisitorHistory(ctx context.Context, clientID, visitorID string, limit int) ( []api.VisitRow, error) { rows, err := s.pool.Query(ctx, ` SELECT vi.id::text, vi.occurred_at, si.name, vi.camera_id, vi.is_new_visitor, vi.similarity, vi.quality, vi.attributes, pu.n, pu.total, pu.cur, pu.currencies FROM visits vi JOIN sites si ON si.id = vi.site_id -- LATERAL rather than a join onto purchases directly: two sales on one -- visit would otherwise return that visit twice and the customer would -- appear to have been in more often than they were. LEFT JOIN LATERAL ( SELECT count(*) AS n, sum(p.amount) AS total, max(p.currency) AS cur, count(DISTINCT p.currency) AS currencies FROM purchases p WHERE p.visit_id = vi.id AND p.client_id = vi.client_id ) pu ON true WHERE vi.client_id = $1 AND vi.visitor_id = $2::uuid ORDER BY vi.occurred_at DESC LIMIT $3`, clientID, visitorID, limit) if err != nil { return nil, err } defer rows.Close() var out []api.VisitRow for rows.Next() { var v api.VisitRow var at time.Time var sim, qual, total *float64 var nPurchases, nCurrencies int var cur *string if err := rows.Scan(&v.ID, &at, &v.Site, &v.CameraID, &v.IsNew, &sim, &qual, &v.Attributes, &nPurchases, &total, &cur, &nCurrencies); err != nil { return nil, err } v.Purchases = nPurchases // One currency or none. Mixed is left as a count with no figure // rather than a sum that means nothing. if nCurrencies == 1 && total != nil && cur != nil { v.Spend, v.Currency = *total, *cur } v.OccurredAt = at.UTC().Format(time.RFC3339) if sim != nil { v.Similarity = *sim } if qual != nil { v.Quality = *qual } out = append(out, v) } return out, rows.Err() } // SaveProfile writes the in-store form, and the consent record with it. // // One transaction: a name saved without its consent row is a customer whose // personal data we hold with no record of being allowed to, which is the exact // state the consents table exists to make impossible. func (s *Store) SaveProfile(ctx context.Context, clientID string, p api.Profile, actor string) error { tx, err := s.pool.Begin(ctx) if err != nil { return err } defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed // Scoped to the client, so an id from another tenant is simply not found - // the same answer as a typo, which is what it should look like. var exists bool err = tx.QueryRow(ctx, ` SELECT true FROM visitors WHERE id = $1::uuid AND client_id = $2 AND deleted_at IS NULL`, p.VisitorID, clientID).Scan(&exists) if errors.Is(err, pgx.ErrNoRows) { return errors.New("no such visitor") } if err != nil { return err } var dob any if p.DateOfBirth != "" { dob = p.DateOfBirth } if _, err := tx.Exec(ctx, ` INSERT INTO visitor_profiles (visitor_id, client_id, full_name, phone, email, gender, date_of_birth, notes, collected_by) VALUES ($1::uuid, $2, $3, $4, $5, $6, $7::date, $8, NULLIF($9, '')::uuid) ON CONFLICT (visitor_id) DO UPDATE SET full_name = EXCLUDED.full_name, phone = EXCLUDED.phone, email = EXCLUDED.email, gender = EXCLUDED.gender, date_of_birth = EXCLUDED.date_of_birth, notes = EXCLUDED.notes, collected_by = EXCLUDED.collected_by, updated_at = now()`, p.VisitorID, clientID, p.FullName, p.Phone, p.Email, p.Gender, dob, p.Notes, actor); err != nil { return fmt.Errorf("save profile: %w", err) } if p.Consent { // Only if there is not already a live one. Re-saving the form must not // stack up consent records, or the audit trail stops being readable. if _, err := tx.Exec(ctx, ` INSERT INTO consents (visitor_id, client_id, scope, method, collected_by, evidence) SELECT $1::uuid, $2, 'biometric', 'in_store_form', NULLIF($3, '')::uuid, '{}'::jsonb WHERE NOT EXISTS ( SELECT 1 FROM consents WHERE visitor_id = $1::uuid AND scope = 'biometric' AND revoked_at IS NULL)`, p.VisitorID, clientID, actor); err != nil { return fmt.Errorf("record consent: %w", err) } } else { // Unticking the box is a withdrawal, and a withdrawal is a timestamp, // never a delete: the fact that they withdrew is itself the thing an // auditor asks to see. if _, err := tx.Exec(ctx, ` UPDATE consents SET revoked_at = now() WHERE visitor_id = $1::uuid AND client_id = $2 AND scope = 'biometric' AND revoked_at IS NULL`, p.VisitorID, clientID); err != nil { return fmt.Errorf("revoke consent: %w", err) } } return tx.Commit(ctx) } // RecordPurchase books a sale against a customer. // // When no site is given it uses the one where this customer was most recently // seen, which is what "the assistant on the floor just sold them something" // means. If they have never been seen anywhere the caller is told to pass a // site rather than being handed a foreign key error. func (s *Store) RecordPurchase(ctx context.Context, clientID string, p api.PurchaseInput, actor string) error { tx, err := s.pool.Begin(ctx) if err != nil { return err } defer tx.Rollback(ctx) //nolint:errcheck var exists bool err = tx.QueryRow(ctx, ` SELECT true FROM visitors WHERE id = $1::uuid AND client_id = $2 AND deleted_at IS NULL`, p.VisitorID, clientID).Scan(&exists) if errors.Is(err, pgx.ErrNoRows) { return errors.New("no such visitor") } if err != nil { return err } siteID := p.SiteID var visitID any if siteID == "" { var sid, vid *string err = tx.QueryRow(ctx, ` SELECT site_id::text, id::text FROM visits WHERE client_id = $1 AND visitor_id = $2::uuid ORDER BY occurred_at DESC LIMIT 1`, clientID, p.VisitorID).Scan(&sid, &vid) if errors.Is(err, pgx.ErrNoRows) || sid == nil { return errors.New("no site for this visitor") } if err != nil { return err } siteID = *sid // Attaching the sale to the visit it belongs to is what makes // "did this visit convert" answerable at all, rather than only // "did this person ever buy". visitID = vid } else { // A site passed in must still belong to the caller's client. var ok bool err = tx.QueryRow(ctx, `SELECT true FROM sites WHERE id = $1::uuid AND client_id = $2`, siteID, clientID).Scan(&ok) if errors.Is(err, pgx.ErrNoRows) { return errors.New("no site for this visitor") } if err != nil { return err } } items := p.Items if items == nil { items = []string{} } if _, err := tx.Exec(ctx, ` INSERT INTO purchases (client_id, site_id, visitor_id, visit_id, amount, currency, items, source, external_ref, recorded_by) VALUES ($1, $2::uuid, $3::uuid, $4::uuid, $5, $6, $7, $8, $9, NULLIF($10, '')::uuid)`, clientID, siteID, p.VisitorID, visitID, p.Amount, p.Currency, items, p.Source, p.Notes, actor); err != nil { return fmt.Errorf("insert purchase: %w", err) } return tx.Commit(ctx) } // ---------------------------------------------------------------- enrolment // RedeemEnrolment spends an installation code and returns the broker login. // // Single-use is enforced by the UPDATE itself: the `used_at IS NULL` predicate // and the write are one statement, so two PCs racing on the same code cannot // both win. Checking first and updating after would be exactly that race. func (s *Store) RedeemEnrolment(ctx context.Context, hash []byte) (api.Enrolment, error) { var en api.Enrolment err := s.pool.QueryRow(ctx, ` UPDATE site_enrolment_tokens SET used_at = now() WHERE token_hash = $1 AND used_at IS NULL AND expires_at > now() RETURNING client_id::text, site_id::text`, hash). Scan(&en.ClientID, &en.SiteID) if errors.Is(err, pgx.ErrNoRows) { return en, errors.New("enrolment token is unknown, expired or already used") } if err != nil { return en, err } var sealed []byte if err := s.pool.QueryRow(ctx, ` SELECT si.name, si.slug, a.id::text, a.mqtt_username, a.mqtt_password_enc FROM sites si JOIN agents a ON a.site_id = si.id WHERE si.id = $1::uuid AND si.client_id = $2`, en.SiteID, en.ClientID). Scan(&en.SiteName, &en.SiteSlug, &en.AgentID, &en.MQTTUser, &sealed); err != nil { if errors.Is(err, pgx.ErrNoRows) { return en, errors.New("site has no agent provisioned - " + "create the broker user before issuing an enrolment token") } return en, err } if len(sealed) == 0 { return en, errors.New("site has no broker password stored") } if s.secrets == nil { return en, errors.New("BEHAVISION_SECRET_KEY is not configured, " + "so stored broker passwords cannot be read") } pass, err := s.secrets.OpenString(sealed, en.AgentID) if err != nil { return en, fmt.Errorf("broker password for %s: %w", en.SiteSlug, err) } en.MQTTPass = pass return en, nil } // SetAgentSecret stores a site's broker password, sealed to that agent's id. // Used by provisioning, never by a request handler. func (s *Store) SetAgentSecret(ctx context.Context, agentID, password string) error { if s.secrets == nil { return errors.New("BEHAVISION_SECRET_KEY is not configured") } sealed, err := s.secrets.SealString(password, agentID) if err != nil { return err } _, err = s.pool.Exec(ctx, `UPDATE agents SET mqtt_password_enc = $2 WHERE id = $1::uuid`, agentID, sealed) return err } // Audit never fails a request. // // A refused audit write is worth knowing about, but refusing the action it was // recording is worse: it would mean an outage in the logging table stops staff // serving customers. func (s *Store) Audit(ctx context.Context, e api.AuditEntry) { detail := e.Detail if detail == nil { detail = map[string]any{} } if _, err := s.pool.Exec(ctx, ` INSERT INTO audit_log (client_id, actor_id, actor_kind, action, entity, entity_id, detail) VALUES (NULLIF($1, '')::uuid, NULLIF($2, '')::uuid, $3, $4, $5, $6, $7)`, e.ClientID, e.ActorID, e.ActorKind, e.Action, e.Entity, e.EntityID, detail); err != nil { s.auditFailed(e.Action, err) } }