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