package ingest import ( "context" "encoding/json" "strings" "testing" "time" "github.com/loyaly/behavision-server/internal/contract" ) type fakeStore struct { sites map[string]Site visits []contract.Visit seenEvents map[string]bool heartbeats []contract.Heartbeat failVisit error } func newFake() *fakeStore { return &fakeStore{ sites: map[string]Site{"acme.store1": {ClientID: "c1", SiteID: "s1", Slug: "acme.store1"}}, seenEvents: map[string]bool{}, } } func (f *fakeStore) ResolveSite(_ context.Context, u string) (Site, error) { s, ok := f.sites[u] if !ok { return Site{}, ErrUnknownSite } return s, nil } func (f *fakeStore) RecordVisit(_ context.Context, _ Site, v *contract.Visit) (bool, error) { if f.failVisit != nil { return false, f.failVisit } if f.seenEvents[v.EventID] { return false, nil } f.seenEvents[v.EventID] = true f.visits = append(f.visits, *v) return true, nil } func (f *fakeStore) RecordHeartbeat(_ context.Context, _ Site, h *contract.Heartbeat) error { f.heartbeats = append(f.heartbeats, *h) return nil } func visitJSON(id string) []byte { b, _ := json.Marshal(contract.Visit{ EventID: id, OccurredAt: time.Now(), CameraID: "entrance"}) return b } func consumer() (*Consumer, *fakeStore) { f := newFake() return &Consumer{Store: f}, f } func TestAVisitIsStored(t *testing.T) { c, f := consumer() if err := c.Handle(context.Background(), "bv/acme.store1/visit", visitJSON("e1")); err != nil { t.Fatal(err) } if len(f.visits) != 1 || c.Accepted != 1 { t.Fatalf("visits=%d accepted=%d", len(f.visits), c.Accepted) } } func TestRedeliveryDoesNotDoubleFootfall(t *testing.T) { // MQTT is at-least-once by design, so this happens after every reconnect. // Getting it wrong inflates the one number the customer pays for. c, f := consumer() ctx := context.Background() c.Handle(ctx, "bv/acme.store1/visit", visitJSON("same-id")) c.Handle(ctx, "bv/acme.store1/visit", visitJSON("same-id")) if len(f.visits) != 1 { t.Fatalf("stored %d rows for one event", len(f.visits)) } if c.Accepted != 1 || c.Duplicate != 1 { t.Fatalf("accepted=%d duplicate=%d - a redelivery is normal, not an error", c.Accepted, c.Duplicate) } } func TestAnUnprovisionedSiteIsDroppedNotCreated(t *testing.T) { // A site typo'd into existence would silently become a tenant with its own // visitors and its own footfall report, and nobody would notice until the // numbers stopped adding up. Provisioning is a deliberate act. c, f := consumer() err := c.Handle(context.Background(), "bv/ghost.store9/visit", visitJSON("e1")) if err != nil { t.Fatalf("an unknown site must not be a transient error: %v", err) } if len(f.visits) != 0 || len(f.sites) != 1 { t.Fatal("an unprovisioned site produced rows") } if c.Dropped != 1 { t.Fatal("the drop was not counted") } } func TestMalformedPayloadIsDroppedNotRetriedForever(t *testing.T) { // A message retried forever stops every good one behind it - the same // failure the agent's spool quarantine exists to prevent. c, _ := consumer() for _, bad := range [][]byte{ []byte("{ truncated"), []byte(`{"event_id":"","occurred_at":"2026-01-01T00:00:00Z"}`), []byte(`{"event_id":"e","occurred_at":"2026-01-01T00:00:00Z","embedding":[1,2,3],"model":"m"}`), } { if err := c.Handle(context.Background(), "bv/acme.store1/visit", bad); err != nil { t.Errorf("malformed payload asked for redelivery: %v", err) } } if c.Dropped != 3 { t.Fatalf("dropped=%d, want 3", c.Dropped) } } func TestADatabaseFailureAsksForRedelivery(t *testing.T) { // The one path that must NOT swallow the message: the event is valid and // the store is broken, so losing it would lose real footfall. c, f := consumer() f.failVisit = context.DeadlineExceeded err := c.Handle(context.Background(), "bv/acme.store1/visit", visitJSON("e1")) if err == nil { t.Fatal("a database failure was swallowed") } if c.Dropped != 0 { t.Fatal("a database failure was counted as a drop") } } func TestABadTopicIsDropped(t *testing.T) { c, _ := consumer() for _, topic := range []string{"bv/acme/visit", "nope/acme.store1/visit", "bv/acme.store1/nonsense"} { if err := c.Handle(context.Background(), topic, visitJSON("e1")); err != nil { t.Errorf("%s: %v", topic, err) } } if c.Dropped != 3 { t.Fatalf("dropped=%d, want 3", c.Dropped) } } func TestHeartbeatIsRecorded(t *testing.T) { c, f := consumer() h, _ := json.Marshal(contract.Heartbeat{ SentAt: time.Now(), RecognitionModel: "w600k_r50.onnx"}) if err := c.Handle(context.Background(), "bv/acme.store1/heartbeat", h); err != nil { t.Fatal(err) } if len(f.heartbeats) != 1 || f.heartbeats[0].RecognitionModel != "w600k_r50.onnx" { t.Fatalf("%+v", f.heartbeats) } } func TestASiteReportingDroppedEventsIsLoggedLoudly(t *testing.T) { // That site lost footfall the customer paid for and will never get back. // It must not be discoverable only by staring at a graph. var sb strings.Builder c, _ := consumer() c.Log = newTestLogger(&sb) h, _ := json.Marshal(contract.Heartbeat{SentAt: time.Now(), Dropped: 417}) c.Handle(context.Background(), "bv/acme.store1/heartbeat", h) out := sb.String() if !strings.Contains(out, "417") || !strings.Contains(out, "WARNING") { t.Fatalf("dropped events not surfaced: %q", out) } } // ---------------------------------------------------------------- the doorbell // A live arrivals stream should learn about a visitor within milliseconds, not // whenever its fallback tick next comes round. func TestANewVisitRingsTheDoorbell(t *testing.T) { c, _ := consumer() var rung []string c.Notify = func(clientID string) { rung = append(rung, clientID) } if err := c.Handle(context.Background(), "bv/acme.store1/visit", visitJSON("e1")); err != nil { t.Fatal(err) } if len(rung) != 1 || rung[0] != "c1" { t.Fatalf("want one ring carrying the tenant, got %v", rung) } } // At-least-once delivery makes redelivery normal after any reconnect. Ringing // for one would wake every live stream on the estate to re-query for rows they // already have - and a reconnect redelivers a whole batch at once. func TestARedeliveredVisitDoesNotRingTheDoorbell(t *testing.T) { c, _ := consumer() rings := 0 c.Notify = func(string) { rings++ } for i := 0; i < 3; i++ { if err := c.Handle(context.Background(), "bv/acme.store1/visit", visitJSON("e1")); err != nil { t.Fatal(err) } } if rings != 1 { t.Fatalf("three deliveries of one event rang %d times, want 1", rings) } } // The consumer must work unchanged with no listener wired up: a server // assembled without a hub still has to ingest. func TestIngestWorksWithNoDoorbellWired(t *testing.T) { c, f := consumer() c.Notify = nil if err := c.Handle(context.Background(), "bv/acme.store1/visit", visitJSON("e1")); err != nil { t.Fatal(err) } if len(f.visits) != 1 { t.Fatalf("visits=%d", len(f.visits)) } }