Files
Behavision/server/internal/store/api_arrivals_live_test.go
Suriyakumarvijayanayagam 30e01765ae The live tests seeded a tenant per run and never took it back
Each live store test makes its own client - deliberately, so they can
run in any order and so the isolation assertions have a real neighbour
to be isolated from - and none of them removed it afterwards. The dev
database had reached 242 abandoned tenants against the one real
company.

That is not untidy, it is a broken screen. The platform admin's
Companies view lists every client, so the real company sat under pages
of `walk1788761685056287000`, which is the first thing anyone opening
tenant administration would see.

dropTenant registers the cleanup against the CLIENT rather than each
table: every foreign key onto clients is ON DELETE CASCADE, so one
delete takes the sites, visitors, visits, face images, embeddings,
cameras and agents with it. A per-table list would rot the first time a
migration adds a table, and it would rot silently - the same shape as
the leak it replaces.

A failed cleanup calls t.Errorf rather than being ignored. A tenant
left behind is precisely what this exists to prevent, and swallowing
the error would let the leak come back with nothing to show for it.

Verified against the live database: three consecutive runs of the store
suite leave clients, sites and visits unchanged.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Qiy5iKfz4L8S4vRaYPBdaU
2026-09-09 12:43:21 +05:30

448 lines
15 KiB
Go

package store
import (
"context"
"fmt"
"os"
"testing"
"time"
"github.com/loyaly/behavision-server/internal/api"
)
// Live database tests for the arrivals feed.
//
// Skipped unless TEST_DATABASE_URL is set, following the same rule as the
// bucket tests: the suite must stay runnable with no network and no services.
// They exist because the rest of the arrivals suite runs against an in-memory
// fake, and a fake cannot catch what actually goes wrong in this file - a
// keyset comparison Postgres plans differently than expected, a LEFT JOIN
// silently promoted to an inner one by a WHERE clause, a column list that
// drifts between the two directions of the query. Those only fail against a
// real planner.
//
// docker run -d -p 55432:5432 -e POSTGRES_PASSWORD=test \
// -e POSTGRES_DB=behavision pgvector/pgvector:pg16
// psql < server/migrations/*.sql
// TEST_DATABASE_URL='postgres://postgres:test@127.0.0.1:55432/behavision' go test ./internal/store/
func liveStore(t *testing.T) *Store {
t.Helper()
dsn := os.Getenv("TEST_DATABASE_URL")
if dsn == "" {
t.Skip("set TEST_DATABASE_URL to run the live store tests")
}
st, err := Open(context.Background(), dsn)
if err != nil {
t.Fatalf("open: %v", err)
}
t.Cleanup(st.Close)
return st
}
// dropTenant removes a seeded tenant when the test that made it finishes.
//
// Without this these tests are a slow leak. Every one of them seeds its own
// tenant - deliberately, so they can run in any order and so the isolation
// assertions have a real neighbour - and none of them ever removed it. A dev
// database reached 242 abandoned tenants against the single real one, which is
// not merely untidy: the platform admin's Companies screen lists every client,
// so the one real company was buried under pages of `walk1788761685056287000`.
//
// Registered against the CLIENT rather than each table because every foreign
// key onto clients is ON DELETE CASCADE, so one delete takes the sites,
// visitors, visits, face images, embeddings, cameras and agents with it. A
// per-table list would rot the first time a migration adds a table, and it
// would rot silently - which is the shape of the bug it is cleaning up after.
//
// t.Cleanup runs LIFO and liveStore registers st.Close before any seeding, so
// the delete still has a live pool when it runs.
func dropTenant(t *testing.T, st *Store, clientID string) {
t.Helper()
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if _, err := st.pool.Exec(ctx,
`DELETE FROM clients WHERE id = $1::uuid`, clientID); err != nil {
// Reported rather than ignored: a tenant left behind is the very
// thing this exists to prevent, and swallowing the error would let
// the leak return with nothing to show for it.
t.Errorf("cleanup tenant %s: %v", clientID, err)
}
})
}
// seedTenant builds a client, a site and n visits, and returns the client id.
// Every test gets its own tenant so they can run in any order without a
// truncate between them - and so the isolation assertions below have a real
// neighbour to be isolated from.
func seedTenant(t *testing.T, st *Store, name string, n int, withImages bool) (clientID, siteID string) {
t.Helper()
ctx := context.Background()
err := st.pool.QueryRow(ctx, `
INSERT INTO clients (name, slug) VALUES ($1, $1) RETURNING id::text`, name).Scan(&clientID)
if err != nil {
t.Fatalf("seed client: %v", err)
}
dropTenant(t, st, clientID)
err = st.pool.QueryRow(ctx, `
INSERT INTO sites (client_id, name, slug) VALUES ($1::uuid, $2, $3)
RETURNING id::text`, clientID, name+" Main", name+"-main").Scan(&siteID)
if err != nil {
t.Fatalf("seed site: %v", err)
}
start := time.Date(2026, 9, 2, 10, 0, 0, 0, time.UTC)
for i := 0; i < n; i++ {
var visitorID string
if err := st.pool.QueryRow(ctx, `
INSERT INTO visitors (client_id, number, label, first_seen_at)
VALUES ($1::uuid, $2, $3, $4) RETURNING id::text`,
clientID, i+1, fmt.Sprintf("Visitor %d", i+1), start).Scan(&visitorID); err != nil {
t.Fatalf("seed visitor: %v", err)
}
key := ""
if withImages {
key = fmt.Sprintf("behavision/v2/%s/main/2026/09/02/%d.jpg", name, i)
}
if _, err := st.pool.Exec(ctx, `
INSERT INTO visits (client_id, site_id, visitor_id, source_event_id,
occurred_at, camera_id, is_new_visitor,
similarity, quality, image_key)
VALUES ($1::uuid, $2::uuid, $3::uuid, $4, $5, 'door', $6, 0.71, 0.66, $7)`,
clientID, siteID, visitorID, fmt.Sprintf("%s-e%d", name, i),
start.Add(time.Duration(i)*time.Second), i == 0, key); err != nil {
t.Fatalf("seed visit: %v", err)
}
}
return clientID, siteID
}
func TestLiveArrivalsWalkTheFeedWithoutLosingAnyone(t *testing.T) {
st := liveStore(t)
clientID, _ := seedTenant(t, st, "walk"+stamp(), 25, true)
ctx := context.Background()
seen := map[string]int{}
// Position zero: replay from the very beginning. Positions start at 1, so
// nothing is excluded.
from := int64(0)
after := &from
for poll := 0; poll < 6; poll++ {
rows, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: clientID, AfterSeq: after, Limit: 10})
if err != nil {
t.Fatalf("poll %d: %v", poll, err)
}
for _, a := range rows {
seen[a.VisitID]++
}
if len(rows) == 0 {
break
}
last := rows[len(rows)-1].Seq
after = &last
}
if len(seen) != 25 {
t.Fatalf("saw %d of 25 visits", len(seen))
}
for id, n := range seen {
if n != 1 {
t.Errorf("visit %s delivered %d times", id, n)
}
}
}
// The case the tuple comparison exists for. Four people through a door at once
// share a timestamp to the microsecond; ordering on time alone either repeats
// them forever or skips three of them.
func TestLiveArrivalsPageThroughASimultaneousBurst(t *testing.T) {
st := liveStore(t)
name := "burst" + stamp()
clientID, siteID := seedTenant(t, st, name, 0, false)
ctx := context.Background()
at := time.Date(2026, 9, 2, 11, 0, 0, 0, time.UTC)
for i := 0; i < 4; i++ {
if _, err := st.pool.Exec(ctx, `
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at, camera_id)
VALUES ($1::uuid, $2::uuid, $3, $4, 'door')`,
clientID, siteID, fmt.Sprintf("%s-b%d", name, i), at); err != nil {
t.Fatal(err)
}
}
seen := map[string]bool{}
from := int64(0)
after := &from
for poll := 0; poll < 5; poll++ {
rows, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: clientID, AfterSeq: after, Limit: 2})
if err != nil {
t.Fatal(err)
}
if len(rows) == 0 {
break
}
for _, a := range rows {
if seen[a.VisitID] {
t.Fatalf("visit %s came back twice - the cursor is stuck", a.VisitID)
}
seen[a.VisitID] = true
}
last := rows[len(rows)-1].Seq
after = &last
}
if len(seen) != 4 {
t.Fatalf("paged a 4-person burst two at a time and saw %d", len(seen))
}
}
// One tenant's feed must never contain another's customers. The site filter is
// caller-supplied, so this asks for a site id that exists - and belongs to
// somebody else.
func TestLiveArrivalsCannotReadAnotherTenant(t *testing.T) {
st := liveStore(t)
mine, _ := seedTenant(t, st, "mine"+stamp(), 3, false)
_, theirSite := seedTenant(t, st, "theirs"+stamp(), 3, false)
ctx := context.Background()
rows, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: mine, SiteID: theirSite, Limit: 50})
if err != nil {
t.Fatal(err)
}
if len(rows) != 0 {
t.Fatalf("read %d visits from another tenant's site", len(rows))
}
}
// A visit with no visitor_id is a site sending counts without templates. It is
// real footfall by an unknown person and an inner join would delete it from the
// feed while the footfall report still counted it.
func TestLiveArrivalsKeepVisitsWithNoVisitor(t *testing.T) {
st := liveStore(t)
name := "anon" + stamp()
clientID, siteID := seedTenant(t, st, name, 0, false)
ctx := context.Background()
if _, err := st.pool.Exec(ctx, `
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at, camera_id)
VALUES ($1::uuid, $2::uuid, $3, now(), 'door')`,
clientID, siteID, name+"-anon"); err != nil {
t.Fatal(err)
}
rows, err := st.Arrivals(ctx, api.ArrivalQuery{ClientID: clientID, Limit: 10})
if err != nil {
t.Fatal(err)
}
if len(rows) != 1 {
t.Fatalf("an anonymous visit vanished from the feed (%d rows)", len(rows))
}
if rows[0].VisitorID != "" {
t.Errorf("visitor id should be empty, got %q", rows[0].VisitorID)
}
}
// An erased customer's visits stay, unlinked - that is the documented erasure
// contract. The feed must still show them, or a shop's live count silently
// drops every time someone exercises their rights.
func TestLiveArrivalsKeepVisitsOfAnErasedCustomer(t *testing.T) {
st := liveStore(t)
clientID, _ := seedTenant(t, st, "erased"+stamp(), 2, false)
ctx := context.Background()
if _, err := st.pool.Exec(ctx,
`UPDATE visitors SET deleted_at = now() WHERE client_id = $1::uuid`,
clientID); err != nil {
t.Fatal(err)
}
rows, err := st.Arrivals(ctx, api.ArrivalQuery{ClientID: clientID, Limit: 10})
if err != nil {
t.Fatal(err)
}
if len(rows) != 2 {
t.Fatalf("erasing a customer removed %d visits from the feed", 2-len(rows))
}
if rows[0].Label != "" {
t.Errorf("an erased customer's label leaked into the feed: %q", rows[0].Label)
}
}
// A profile name must reach the feed, or a shop screen shows "Visitor 12" for
// a regular whose name staff typed in last week.
func TestLiveArrivalsCarryTheProfileName(t *testing.T) {
st := liveStore(t)
clientID, _ := seedTenant(t, st, "named"+stamp(), 1, false)
ctx := context.Background()
var visitorID string
if err := st.pool.QueryRow(ctx,
`SELECT id::text FROM visitors WHERE client_id = $1::uuid`, clientID).
Scan(&visitorID); err != nil {
t.Fatal(err)
}
if _, err := st.pool.Exec(ctx, `
INSERT INTO visitor_profiles (client_id, visitor_id, full_name)
VALUES ($1::uuid, $2::uuid, 'Asha Menon')`, clientID, visitorID); err != nil {
t.Fatal(err)
}
rows, err := st.Arrivals(ctx, api.ArrivalQuery{ClientID: clientID, Limit: 10})
if err != nil {
t.Fatal(err)
}
if rows[0].Name != "Asha Menon" {
t.Fatalf("profile name did not reach the feed: %q", rows[0].Name)
}
if rows[0].Label == "" {
t.Error("the system label should travel alongside the typed name")
}
}
// No cursor means "an app that has just opened": it wants the LAST few
// arrivals, not the first few ever recorded - but still ascending, so the
// caller's cursor handling is identical on every poll.
func TestLiveArrivalsFirstPollIsTheNewestWindowAscending(t *testing.T) {
st := liveStore(t)
clientID, _ := seedTenant(t, st, "newest"+stamp(), 12, false)
ctx := context.Background()
rows, err := st.Arrivals(ctx, api.ArrivalQuery{ClientID: clientID, Limit: 4})
if err != nil {
t.Fatal(err)
}
if len(rows) != 4 {
t.Fatalf("got %d rows", len(rows))
}
for i := 1; i < len(rows); i++ {
if rows[i-1].OccurredAt >= rows[i].OccurredAt {
t.Fatalf("not ascending at %d", i)
}
}
// Seeded one second apart from 10:00:00, so the newest four start at :08.
if want := "2026-09-02T10:00:08Z"; rows[0].OccurredAt != want {
t.Errorf("first poll started at %s, want the newest window at %s",
rows[0].OccurredAt, want)
}
}
func TestLiveArrivalsCarryTheImageKeyForPresigning(t *testing.T) {
st := liveStore(t)
clientID, _ := seedTenant(t, st, "img"+stamp(), 1, true)
rows, err := st.Arrivals(context.Background(),
api.ArrivalQuery{ClientID: clientID, Limit: 10})
if err != nil {
t.Fatal(err)
}
if rows[0].ImageKey == "" {
t.Fatal("no image key reached the handler, so no photo can be signed")
}
if rows[0].Image.URL != "" {
t.Error("the store must not presign - it has no bucket and checks nobody")
}
}
func stamp() string { return fmt.Sprintf("%d", time.Now().UnixNano()) }
// The regression test for the bug that shipped, and was caught only by running
// the real thing against a real broker.
//
// The feed used to be ordered by (occurred_at, id). Four people through one
// door share occurred_at to the microsecond, so the tie-break fell to `id` - a
// RANDOM uuid. A visit that committed AFTER the reader had moved its cursor but
// carried a lower uuid sorted behind that cursor and was never delivered.
// Measured live: four simultaneous visits published, two delivered, and no
// counter anywhere that would show the other two had been dropped.
//
// This reproduces the exact shape: read, move the cursor, THEN insert more rows
// carrying the same occurred_at. Every one of them must still arrive.
func TestLiveArrivalsDeliverLateInsertsThatShareATimestamp(t *testing.T) {
st := liveStore(t)
name := "late" + stamp()
clientID, siteID := seedTenant(t, st, name, 0, false)
ctx := context.Background()
// One instant for everybody - this is a single frame of one camera.
at := time.Date(2026, 9, 2, 12, 0, 0, 0, time.UTC)
insert := func(tag string) {
t.Helper()
if _, err := st.pool.Exec(ctx, `
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at, camera_id)
VALUES ($1::uuid, $2::uuid, $3, $4, 'door')`,
clientID, siteID, name+"-"+tag, at); err != nil {
t.Fatal(err)
}
}
insert("a")
insert("b")
from := int64(0)
rows, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: clientID, AfterSeq: &from, Limit: 50})
if err != nil {
t.Fatal(err)
}
if len(rows) != 2 {
t.Fatalf("first read got %d rows, want 2", len(rows))
}
cursor := rows[len(rows)-1].Seq
// Now two more arrive at the SAME instant, after the cursor has moved.
// Under the old ordering roughly half of these vanished, depending on how
// their random uuids happened to sort.
insert("c")
insert("d")
rest, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: clientID, AfterSeq: &cursor, Limit: 50})
if err != nil {
t.Fatal(err)
}
if len(rest) != 2 {
t.Fatalf("late inserts sharing a timestamp: got %d of 2 - people are being dropped from the feed",
len(rest))
}
}
// Run the same shape many times over. The old bug was probabilistic - it
// depended on how random uuids happened to sort - so a single pass could pass
// by luck. This one cannot.
func TestLiveArrivalsNeverDropAnyoneAcrossManySimultaneousBursts(t *testing.T) {
st := liveStore(t)
name := "many" + stamp()
clientID, siteID := seedTenant(t, st, name, 0, false)
ctx := context.Background()
at := time.Date(2026, 9, 2, 13, 0, 0, 0, time.UTC)
cursor := int64(0)
delivered := 0
for round := 0; round < 30; round++ {
for i := 0; i < 4; i++ {
if _, err := st.pool.Exec(ctx, `
INSERT INTO visits (client_id, site_id, source_event_id, occurred_at, camera_id)
VALUES ($1::uuid, $2::uuid, $3, $4, 'door')`,
clientID, siteID, fmt.Sprintf("%s-r%d-%d", name, round, i), at); err != nil {
t.Fatal(err)
}
}
rows, err := st.Arrivals(ctx, api.ArrivalQuery{
ClientID: clientID, AfterSeq: &cursor, Limit: 50})
if err != nil {
t.Fatal(err)
}
delivered += len(rows)
if len(rows) > 0 {
cursor = rows[len(rows)-1].Seq
}
}
if delivered != 120 {
t.Fatalf("30 bursts of 4 people delivered %d of 120", delivered)
}
}