package bridge import ( "context" "encoding/json" "errors" "net/http" "net/http/httptest" "strings" "testing" ) type fakeQueue struct { topics []string payloads []map[string]any err error } func (q *fakeQueue) Append(topic string, payload any) error { if q.err != nil { return q.err } b, _ := json.Marshal(payload) var m map[string]any json.Unmarshal(b, &m) q.topics = append(q.topics, topic) q.payloads = append(q.payloads, m) return nil } type fakeEmb struct { calls int err error } func (f *fakeEmb) Embedding(_ context.Context, id int64) ([]float32, string, error) { f.calls++ if f.err != nil { return nil, "", f.err } return make([]float32, 512), "w600k_r50.onnx", nil } func newBridge() (*Bridge, *fakeQueue, *fakeEmb) { q, e := &fakeQueue{}, &fakeEmb{} return &Bridge{Queue: q, Embeddings: e, TopicPrefix: "bv/acme.store1"}, q, e } func seen(id int64, ts float64) Event { return Event{Type: "person.seen", CameraID: "entrance", TS: ts, Data: map[string]any{"identity_id": float64(id), "similarity": 0.58, "quality": 0.74, "gender": "Male", "age": float64(34)}} } func TestADetectionIsQueuedForTheRightTopic(t *testing.T) { b, q, _ := newBridge() if err := b.Handle(context.Background(), seen(7, 1787996491)); err != nil { t.Fatal(err) } if len(q.topics) != 1 || q.topics[0] != "bv/acme.store1/visit" { t.Fatalf("topics=%v", q.topics) } p := q.payloads[0] if p["camera_id"] != "entrance" || p["is_new"] != false { t.Fatalf("%+v", p) } if p["quality"] != 0.74 || p["similarity"] != 0.58 { t.Fatalf("measurements lost: %+v", p) } } func TestTheEventIDIsStableForTheSameSighting(t *testing.T) { // This is what makes at-least-once delivery safe. A random id here would // defeat the server's idempotency check and double the store's footfall // after every reconnect. b, q, _ := newBridge() ctx := context.Background() b.Handle(ctx, seen(7, 1787996491.20)) b.Handle(ctx, seen(7, 1787996491.86)) // same second, redelivered if q.payloads[0]["event_id"] != q.payloads[1]["event_id"] { t.Fatalf("ids differ: %v vs %v", q.payloads[0]["event_id"], q.payloads[1]["event_id"]) } } func TestDifferentPeopleAndCamerasGetDifferentIDs(t *testing.T) { b, q, _ := newBridge() ctx := context.Background() b.Handle(ctx, seen(7, 1787996491)) b.Handle(ctx, seen(8, 1787996491)) // different person, same instant other := seen(7, 1787996491) other.CameraID = "till" b.Handle(ctx, other) // same person, different camera ids := map[any]bool{} for _, p := range q.payloads { ids[p["event_id"]] = true } if len(ids) != 3 { t.Fatalf("collided: %d distinct ids from 3 sightings", len(ids)) } } func TestLocalDiagnosticsAreNotCountedAsVisits(t *testing.T) { // person.missed and camera.up are real events, but sending them down the // footfall stream would inflate the headcount with things that are not // people. b, q, _ := newBridge() for _, typ := range []string{"person.missed", "camera.up", "camera.down", "identity.merged"} { b.Handle(context.Background(), Event{Type: typ, CameraID: "entrance"}) } if len(q.payloads) != 0 { t.Fatalf("queued %d non-visits", len(q.payloads)) } if b.Skipped != 4 { t.Fatalf("skipped=%d", b.Skipped) } } func TestAnUnclaimedPCQueuesNothing(t *testing.T) { // Without a tenant prefix an event would be addressed to nowhere, and the // broker would refuse it anyway. b, q, _ := newBridge() b.TopicPrefix = "" b.Handle(context.Background(), seen(7, 1)) if len(q.payloads) != 0 { t.Fatal("queued an event with no tenant") } } func TestTheTemplateIsFetchedOncePerIdentity(t *testing.T) { // A returning customer seen forty times a day would otherwise pull the // same 512 floats out of SQLite forty times. b, _, e := newBridge() ctx := context.Background() for i := 0; i < 5; i++ { b.Handle(ctx, seen(7, float64(1787996491+i*60))) } if e.calls != 1 { t.Fatalf("fetched the template %d times", e.calls) } } func TestAMissingTemplateStillQueuesTheVisit(t *testing.T) { // A footfall count without a template is still a real visit. Dropping it // would lose the number the customer is paying for over an optional field. b, q, e := newBridge() e.err = errors.New("no embedding") if err := b.Handle(context.Background(), seen(7, 1)); err != nil { t.Fatal(err) } if len(q.payloads) != 1 { t.Fatal("visit dropped because the template was missing") } if _, has := q.payloads[0]["embedding"]; has { t.Fatal("queued an embedding key with no embedding") } } func TestAttributesSurviveButBookkeepingIsDropped(t *testing.T) { b, q, _ := newBridge() b.Handle(context.Background(), seen(7, 1)) attrs := q.payloads[0]["attributes"].(map[string]any) if attrs["gender"] != "Male" || attrs["age"] != float64(34) { t.Fatalf("attributes lost: %+v", attrs) } for _, k := range []string{"identity_id", "similarity", "quality"} { if _, dup := attrs[k]; dup { t.Errorf("%q duplicated into attributes; it is already a column", k) } } } func TestMalformedJSONIsRejectedNotRetried(t *testing.T) { b, _, _ := newBridge() srv := httptest.NewServer(b.Handler()) defer srv.Close() resp, err := http.Post(srv.URL+"/events", "application/json", strings.NewReader("{ truncated")) if err != nil { t.Fatal(err) } defer resp.Body.Close() if resp.StatusCode != http.StatusBadRequest { t.Fatalf("status %d — the engine would retry this forever", resp.StatusCode) } } func TestTheWebhookQueuesARealPost(t *testing.T) { b, q, _ := newBridge() srv := httptest.NewServer(b.Handler()) defer srv.Close() body, _ := json.Marshal(seen(7, 1787996491)) resp, err := http.Post(srv.URL+"/events", "application/json", strings.NewReader(string(body))) if err != nil { t.Fatal(err) } defer resp.Body.Close() if resp.StatusCode != http.StatusNoContent { t.Fatalf("status %d", resp.StatusCode) } if len(q.payloads) != 1 { t.Fatal("nothing queued") } } func TestListenBindsLoopbackOnly(t *testing.T) { // The webhook accepts unauthenticated posts that become footfall rows; // it must not be reachable from the network. b, _, _ := newBridge() url, stop, err := b.Listen(context.Background()) if err != nil { t.Fatal(err) } defer stop() if !strings.HasPrefix(url, "http://127.0.0.1:") { t.Fatalf("bridge listening on %s", url) } } // ---------------------------------------------------------------- waking // Waking BEFORE the append would send the pump to look at a queue the event // has not reached yet: it finds nothing, goes back to sleep, and the visit then // waits out the full idle interval anyway - the exact delay the wake exists to // remove, with an extra wasted drain on top. func TestTheQueueIsWokenOnlyAfterTheVisitIsOnDisk(t *testing.T) { b, q, _ := newBridge() var depthWhenWoken = -1 b.Wake = func() { depthWhenWoken = len(q.payloads) } if err := b.Handle(context.Background(), seen(7, 1_700_000_000)); err != nil { t.Fatal(err) } if depthWhenWoken != 1 { t.Fatalf("woken with %d events queued, want 1 - the event must be durable first", depthWhenWoken) } } // A visit that never reached the queue must not wake anything: there is nothing // to drain, and the pump would spin on an empty spool. func TestAFailedAppendDoesNotWakeThePump(t *testing.T) { b, q, _ := newBridge() q.err = errors.New("disk full") woken := 0 b.Wake = func() { woken++ } if err := b.Handle(context.Background(), seen(7, 1_700_000_000)); err == nil { t.Fatal("a failed append should surface as an error") } if woken != 0 { t.Fatalf("woke the pump %d times for an event that was never queued", woken) } } // Diagnostics are not visits. They are skipped before the queue, so they must // not wake a pump either. func TestASkippedEventDoesNotWakeThePump(t *testing.T) { b, _, _ := newBridge() woken := 0 b.Wake = func() { woken++ } ev := seen(7, 1_700_000_000) ev.Type = "person.missed" if err := b.Handle(context.Background(), ev); err != nil { t.Fatal(err) } if woken != 0 { t.Fatalf("a diagnostic event woke the pump %d times", woken) } } // The bridge must run unchanged with no waker wired, which is what an agent // built before this existed looks like. func TestNoWakerIsFine(t *testing.T) { b, q, _ := newBridge() b.Wake = nil if err := b.Handle(context.Background(), seen(7, 1_700_000_000)); err != nil { t.Fatal(err) } if len(q.payloads) != 1 { t.Fatalf("queued %d visits, want 1", len(q.payloads)) } }