Found by walking the scenario against production rather than by a test. Two records, each with a phone; the survivor kept its own, and the source's simply stopped existing. Searching for it returned nothing. The first version's rule was "fill the survivor's blanks, never overwrite what it has", which is right about which value WINS and said nothing about the one that loses. One person can have two numbers, two spellings of a name, a work address and a personal one - and a merge that quietly deletes one is exactly the data loss this file already refuses elsewhere: "silently turning Alice back into Visitor 3 is data loss the operator cannot see happen." The profile is now reconciled field by field in Go rather than in one clever upsert, because the interesting case was never the winner. Blanks are still filled and the survivor still keeps its own values, but every losing value is returned in `discarded` AND appended to the survivor's notes - the response is read once and the record is read forever. Notes themselves are additive rather than a winner: two people writing about one customer wrote two different true things. mergeProfiles is pure, so the rule is asserted directly - four cases including the ordinary one, a typed record joining a camera record with no profile at all, which must add no noise. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
314 lines
11 KiB
Go
314 lines
11 KiB
Go
// Creating a customer nobody has photographed, and joining two records that
|
|
// turn out to be one person.
|
|
//
|
|
// These ship together on purpose. A customer created by hand has no face, so
|
|
// when a camera later sees that person the matcher has nothing to compare
|
|
// against and records them as somebody new - by construction, not by failure.
|
|
// Shipping the create without the merge would mean manufacturing duplicates
|
|
// with no way back, which is the state CLAUDE.md already flags for the server:
|
|
// "there is no merge endpoint server-side, so its duplicates would be
|
|
// unrecoverable."
|
|
package store
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
|
|
"github.com/loyaly/behavision-server/internal/api"
|
|
)
|
|
|
|
// CreateCustomer registers a person before any camera has seen them.
|
|
//
|
|
// The number comes from the same counter, taken the same way, as a customer
|
|
// the engine enrols: `UPDATE ... RETURNING` inside the transaction. Two
|
|
// sources of visitor numbers that could disagree would be worse than none,
|
|
// and V-42 has to mean one person whichever way they arrived.
|
|
func (s *Store) CreateCustomer(ctx context.Context, clientID string,
|
|
in api.Profile, createdBy string) (api.Customer, error) {
|
|
|
|
tx, err := s.pool.Begin(ctx)
|
|
if err != nil {
|
|
return api.Customer{}, err
|
|
}
|
|
defer tx.Rollback(ctx)
|
|
|
|
var number int64
|
|
if err := tx.QueryRow(ctx, `
|
|
UPDATE clients SET visitor_seq = visitor_seq + 1
|
|
WHERE id = $1::uuid RETURNING visitor_seq`, clientID).Scan(&number); err != nil {
|
|
return api.Customer{}, fmt.Errorf("next visitor number: %w", err)
|
|
}
|
|
|
|
// The label is the person's name when they gave one, and "Visitor N"
|
|
// otherwise - the same string the engine would have written, so a record
|
|
// created by hand is indistinguishable from an enrolled one afterwards.
|
|
// Formatted in Go, never as `'Visitor ' || $2::text` beside `number = $2`:
|
|
// one parameter used as a bigint and as a string operand makes Postgres
|
|
// deduce two types for it and refuse the whole insert.
|
|
label := in.FullName
|
|
if label == "" {
|
|
label = fmt.Sprintf("Visitor %d", number)
|
|
}
|
|
|
|
now := time.Now().UTC()
|
|
var id string
|
|
if err := tx.QueryRow(ctx, `
|
|
INSERT INTO visitors (client_id, number, label, first_seen_at, visit_count)
|
|
VALUES ($1::uuid, $2, $3, $4, 0) RETURNING id::text`,
|
|
clientID, number, label, now).Scan(&id); err != nil {
|
|
return api.Customer{}, err
|
|
}
|
|
|
|
if _, err := tx.Exec(ctx, `
|
|
INSERT INTO visitor_profiles (visitor_id, client_id, full_name, phone,
|
|
email, gender, notes, collected_by)
|
|
VALUES ($1::uuid, $2::uuid, $3, $4, $5, $6, $7, NULLIF($8,'')::uuid)`,
|
|
id, clientID, in.FullName, in.Phone, in.Email, in.Gender, in.Notes,
|
|
createdBy); err != nil {
|
|
return api.Customer{}, err
|
|
}
|
|
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return api.Customer{}, err
|
|
}
|
|
return api.Customer{
|
|
ID: id, Ref: api.VisitorRef(number), Label: label,
|
|
FullName: in.FullName, Phone: in.Phone, Email: in.Email,
|
|
VisitCount: 0, HasProfile: true,
|
|
FirstSeenAt: now.Format(time.RFC3339),
|
|
}, nil
|
|
}
|
|
|
|
// ErrSameVisitor is the API package's sentinel, aliased rather than
|
|
// redeclared - two values would compare unequal and errors.Is would miss.
|
|
var ErrSameVisitor = api.ErrSameVisitor
|
|
|
|
// MergeVisitors folds `sourceID` into `targetID` and deletes the source.
|
|
//
|
|
// One transaction, because a half-merge - visits moved, profile not - leaves
|
|
// two records each holding part of one person, which is strictly worse than
|
|
// the duplicate it was called to fix.
|
|
//
|
|
// Five tables reference visitors and every one is re-pointed here. A merge
|
|
// that misses a table is the same half-merge arrived at by omission, and
|
|
// ON DELETE CASCADE means the miss is not an error: the rows are silently
|
|
// destroyed with the source row.
|
|
func (s *Store) MergeVisitors(ctx context.Context, clientID, sourceID, targetID string) (
|
|
api.MergeResult, error) {
|
|
|
|
var out api.MergeResult
|
|
if sourceID == targetID {
|
|
return out, ErrSameVisitor
|
|
}
|
|
|
|
tx, err := s.pool.Begin(ctx)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
defer tx.Rollback(ctx)
|
|
|
|
// Both must exist, belong to this tenant, and not already be erased.
|
|
// Locked in a stable order so two operators merging the same pair in
|
|
// opposite directions deadlock on nothing and one simply loses.
|
|
var srcNum, dstNum int64
|
|
var srcLabel, dstLabel string
|
|
var srcFirst, dstFirst time.Time
|
|
rows, err := tx.Query(ctx, `
|
|
SELECT id::text, number, label, first_seen_at FROM visitors
|
|
WHERE client_id = $1::uuid AND id::text IN ($2, $3)
|
|
AND deleted_at IS NULL
|
|
ORDER BY id FOR UPDATE`, clientID, sourceID, targetID)
|
|
if err != nil {
|
|
return out, err
|
|
}
|
|
found := 0
|
|
for rows.Next() {
|
|
var id, label string
|
|
var num int64
|
|
var first time.Time
|
|
if err := rows.Scan(&id, &num, &label, &first); err != nil {
|
|
rows.Close()
|
|
return out, err
|
|
}
|
|
found++
|
|
if id == sourceID {
|
|
srcNum, srcLabel, srcFirst = num, label, first
|
|
} else {
|
|
dstNum, dstLabel, dstFirst = num, label, first
|
|
}
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return out, err
|
|
}
|
|
if found != 2 {
|
|
return out, pgx.ErrNoRows
|
|
}
|
|
|
|
// Profile: visitor_profiles is UNIQUE on visitor_id, so the two cannot both
|
|
// move and something has to win. Merged field by field in Go rather than in
|
|
// one clever upsert, because the interesting case is not which value wins -
|
|
// it is what happens to the one that loses.
|
|
//
|
|
// Blanks on the survivor are filled from the source. Where BOTH hold a
|
|
// value the survivor keeps its own and the loser is recorded in `discarded`
|
|
// and appended to notes. Dropping it silently was the first version's
|
|
// behaviour and it lost a phone number on the very first live run: one
|
|
// person can have two numbers, and a merge that quietly deletes one is
|
|
// precisely the data loss an operator cannot see happen.
|
|
var src, dst profileFields
|
|
if err := readProfile(ctx, tx, sourceID, &src); err != nil {
|
|
return out, fmt.Errorf("read source profile: %w", err)
|
|
}
|
|
if err := readProfile(ctx, tx, targetID, &dst); err != nil {
|
|
return out, fmt.Errorf("read target profile: %w", err)
|
|
}
|
|
if src.exists {
|
|
merged, discarded := mergeProfiles(src, dst)
|
|
out.Discarded = discarded
|
|
if _, err := tx.Exec(ctx, `
|
|
INSERT INTO visitor_profiles (visitor_id, client_id, full_name, phone,
|
|
email, gender, notes)
|
|
VALUES ($1::uuid, $2::uuid, $3, $4, $5, $6, $7)
|
|
ON CONFLICT (visitor_id) DO UPDATE SET
|
|
full_name = EXCLUDED.full_name, phone = EXCLUDED.phone,
|
|
email = EXCLUDED.email, gender = EXCLUDED.gender,
|
|
notes = EXCLUDED.notes, updated_at = now()`,
|
|
targetID, clientID, merged.fullName, merged.phone, merged.email,
|
|
merged.gender, merged.notes); err != nil {
|
|
return out, fmt.Errorf("merge profile: %w", err)
|
|
}
|
|
}
|
|
|
|
for _, q := range []struct {
|
|
name, sql string
|
|
count *int
|
|
}{
|
|
{"visits", `UPDATE visits SET visitor_id = $2::uuid
|
|
WHERE visitor_id = $1::uuid AND client_id = $3::uuid`, &out.Visits},
|
|
{"purchases", `UPDATE purchases SET visitor_id = $2::uuid
|
|
WHERE visitor_id = $1::uuid AND client_id = $3::uuid`, &out.Purchases},
|
|
{"embeddings", `UPDATE visitor_embeddings SET visitor_id = $2::uuid
|
|
WHERE visitor_id = $1::uuid AND client_id = $3::uuid`, &out.Embeddings},
|
|
{"consents", `UPDATE consents SET visitor_id = $2::uuid
|
|
WHERE visitor_id = $1::uuid AND client_id = $3::uuid`, &out.Consents},
|
|
} {
|
|
tag, err := tx.Exec(ctx, q.sql, sourceID, targetID, clientID)
|
|
if err != nil {
|
|
return out, fmt.Errorf("merge %s: %w", q.name, err)
|
|
}
|
|
*q.count = int(tag.RowsAffected())
|
|
}
|
|
|
|
// A human-assigned name outranks an auto "Visitor N", whichever direction
|
|
// the operator merged in. Silently turning "Alice" back into "Visitor 3"
|
|
// is data loss they cannot see happen.
|
|
label := dstLabel
|
|
if isAutoLabel(dstLabel, dstNum) && !isAutoLabel(srcLabel, srcNum) {
|
|
label = srcLabel
|
|
}
|
|
// first_seen_at takes the earlier of the two: it is one person and always
|
|
// was. visit_count is recomputed with COUNT(*), never summed - the stored
|
|
// counters may themselves be stale, and the row count cannot be.
|
|
first := dstFirst
|
|
if srcFirst.Before(first) {
|
|
first = srcFirst
|
|
}
|
|
if _, err := tx.Exec(ctx, `
|
|
UPDATE visitors SET
|
|
label = $2, first_seen_at = $3,
|
|
last_seen_at = GREATEST(last_seen_at,
|
|
(SELECT max(occurred_at) FROM visits WHERE visitor_id = $1::uuid)),
|
|
visit_count = (SELECT count(*) FROM visits WHERE visitor_id = $1::uuid)
|
|
WHERE id = $1::uuid`, targetID, label, first); err != nil {
|
|
return out, fmt.Errorf("merge totals: %w", err)
|
|
}
|
|
|
|
// The source goes for real. A soft delete would leave its number resolving
|
|
// to a record with nothing in it, which reads as "this customer exists and
|
|
// has never been here" - a worse answer than "no such customer".
|
|
if _, err := tx.Exec(ctx, `DELETE FROM visitors WHERE id = $1::uuid`, sourceID); err != nil {
|
|
return out, fmt.Errorf("delete merged customer: %w", err)
|
|
}
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return out, err
|
|
}
|
|
|
|
out.VisitorID = targetID
|
|
out.Ref = api.VisitorRef(dstNum)
|
|
out.RetiredRef = api.VisitorRef(srcNum)
|
|
out.Label = label
|
|
return out, nil
|
|
}
|
|
|
|
// isAutoLabel reports whether a label is the one the system writes itself.
|
|
// Compared against the record's OWN number: "Visitor 7" on customer 42 was
|
|
// typed by a person and is a name, however unhelpful.
|
|
func isAutoLabel(label string, number int64) bool {
|
|
return label == fmt.Sprintf("Visitor %d", number)
|
|
}
|
|
|
|
// profileFields is the part of a profile a merge has to reconcile.
|
|
type profileFields struct {
|
|
exists bool
|
|
fullName, phone, email, gender, notes string
|
|
}
|
|
|
|
func readProfile(ctx context.Context, tx pgx.Tx, visitorID string, out *profileFields) error {
|
|
err := tx.QueryRow(ctx, `
|
|
SELECT full_name, phone, email, gender, notes
|
|
FROM visitor_profiles WHERE visitor_id = $1::uuid`, visitorID).
|
|
Scan(&out.fullName, &out.phone, &out.email, &out.gender, &out.notes)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil
|
|
}
|
|
out.exists = err == nil
|
|
return err
|
|
}
|
|
|
|
// mergeProfiles keeps the survivor's own values, fills its blanks from the
|
|
// source, and returns everything that lost so nothing disappears silently.
|
|
func mergeProfiles(src, dst profileFields) (profileFields, []string) {
|
|
out := dst
|
|
var discarded []string
|
|
keep := func(field string, mine *string, theirs string) {
|
|
switch {
|
|
case theirs == "":
|
|
case *mine == "":
|
|
*mine = theirs
|
|
case *mine != theirs:
|
|
discarded = append(discarded, field+": "+theirs)
|
|
}
|
|
}
|
|
keep("name", &out.fullName, src.fullName)
|
|
keep("phone", &out.phone, src.phone)
|
|
keep("email", &out.email, src.email)
|
|
keep("gender", &out.gender, src.gender)
|
|
|
|
// Notes are additive rather than a winner: two people writing about one
|
|
// customer wrote two different true things.
|
|
if src.notes != "" && src.notes != out.notes {
|
|
if out.notes == "" {
|
|
out.notes = src.notes
|
|
} else {
|
|
out.notes += "\n" + src.notes
|
|
}
|
|
}
|
|
// And the losers land in notes too, because the response is read once and
|
|
// the record is read forever.
|
|
if len(discarded) > 0 {
|
|
line := "merged, also known as - " + strings.Join(discarded, ", ")
|
|
if out.notes == "" {
|
|
out.notes = line
|
|
} else {
|
|
out.notes += "\n" + line
|
|
}
|
|
}
|
|
return out, discarded
|
|
}
|