package mqtt import ( "context" "errors" "sync" "testing" "time" "github.com/loyaly/behavision-agent/pkg/spool" ) type fakeBroker struct { mu sync.Mutex connected bool sent []string failAfter int // fail every publish once this many have succeeded err error } func (f *fakeBroker) Publish(ctx context.Context, topic string, payload []byte) error { f.mu.Lock() defer f.mu.Unlock() if f.failAfter > 0 && len(f.sent) >= f.failAfter { if f.err != nil { return f.err } return errors.New("broker refused") } f.sent = append(f.sent, string(payload)) return nil } func (f *fakeBroker) Connected() bool { f.mu.Lock() defer f.mu.Unlock() return f.connected } func (f *fakeBroker) delivered() []string { f.mu.Lock() defer f.mu.Unlock() return append([]string(nil), f.sent...) } func queue(t *testing.T, payloads ...string) *spool.Spool { t.Helper() s, err := spool.Open(t.TempDir(), 100) if err != nil { t.Fatal(err) } for _, p := range payloads { if err := s.Append("visit", p); err != nil { t.Fatal(err) } } return s } func TestItDrainsInOrderAndAcks(t *testing.T) { q := queue(t, "a", "b", "c") b := &fakeBroker{connected: true} p := &Pump{Queue: q, Publisher: b} sent, err := p.drainOnce(context.Background()) if err != nil { t.Fatal(err) } if sent != 3 || q.Len() != 0 { t.Fatalf("sent %d, %d left in queue", sent, q.Len()) } got := b.delivered() if len(got) != 3 || got[0] != `"a"` || got[2] != `"c"` { t.Fatalf("wrong order: %v", got) } } func TestNothingIsAckedWhileTheBrokerIsDown(t *testing.T) { // Acking an event the broker never took is how footfall disappears. q := queue(t, "a", "b") b := &fakeBroker{connected: false} p := &Pump{Queue: q, Publisher: b} if _, err := p.drainOnce(context.Background()); err == nil { t.Fatal("a disconnected broker was treated as success") } if q.Len() != 2 { t.Fatalf("events were dropped while offline: %d left", q.Len()) } } func TestAFailureStopsTheBatchInsteadOfSkippingPast(t *testing.T) { // Events are a per-visitor timeline read in order; publishing around a // stuck one would reorder a customer's visits on the server. q := queue(t, "a", "b", "c") b := &fakeBroker{connected: true, failAfter: 1} p := &Pump{Queue: q, Publisher: b} sent, err := p.drainOnce(context.Background()) if err == nil { t.Fatal("failure not reported") } if sent != 1 { t.Fatalf("sent %d, want 1 before stopping", sent) } if q.Len() != 2 { t.Fatalf("%d left in queue, want the 2 unsent", q.Len()) } // And the survivors are the RIGHT two, still in order. rest, _ := q.Peek(10) if string(rest[0].Payload) != `"b"` { t.Fatalf("queue head is %s, want b", rest[0].Payload) } } func TestConfirmedEventsSurviveAMidBatchFailure(t *testing.T) { // Acking per-event rather than per-batch: a batch ack would re-send // everything before the failure after a restart, duplicating footfall. q := queue(t, "a", "b", "c") b := &fakeBroker{connected: true, failAfter: 2} p := &Pump{Queue: q, Publisher: b} p.drainOnce(context.Background()) if q.Len() != 1 { t.Fatalf("%d left, want only the unsent one", q.Len()) } b.failAfter = 0 sent, err := p.drainOnce(context.Background()) if err != nil || sent != 1 { t.Fatalf("recovery sent %d (%v)", sent, err) } got := b.delivered() if len(got) != 3 { t.Fatalf("delivered %v - duplicates or losses", got) } } func TestRunRecoversWhenTheBrokerComesBack(t *testing.T) { q := queue(t, "a") b := &fakeBroker{connected: false} p := &Pump{Queue: q, Publisher: b} ctx, cancel := context.WithCancel(context.Background()) defer cancel() go p.Run(ctx) time.Sleep(100 * time.Millisecond) if len(b.delivered()) != 0 { t.Fatal("published while disconnected") } b.mu.Lock() b.connected = true b.mu.Unlock() deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { if len(b.delivered()) == 1 { return } time.Sleep(10 * time.Millisecond) } t.Fatal("queue never drained after the broker returned") } func TestHeartbeatIsSentSeparatelyFromTheQueue(t *testing.T) { // "Site offline" and "nobody visited" must be distinguishable on the // server. And a heartbeat is only meaningful now, so it is never spooled - // otherwise a reconnecting site replays a week of "I am alive". q := queue(t) b := &fakeBroker{connected: true} p := &Pump{Queue: q, Publisher: b, HeartbeatTopic: "site/alive", HeartbeatPayload: func() []byte { return []byte(`{"up":true}`) }, HeartbeatInterval: 20 * time.Millisecond} ctx, cancel := context.WithCancel(context.Background()) defer cancel() go p.Run(ctx) time.Sleep(200 * time.Millisecond) if len(b.delivered()) == 0 { t.Fatal("no heartbeat was sent") } if q.Len() != 0 { t.Fatal("heartbeats were written to the durable queue") } } func TestRunStopsPromptlyOnCancel(t *testing.T) { q := queue(t) b := &fakeBroker{connected: true} p := &Pump{Queue: q, Publisher: b} ctx, cancel := context.WithCancel(context.Background()) done := make(chan struct{}) go func() { p.Run(ctx); close(done) }() cancel() select { case <-done: case <-time.After(3 * time.Second): t.Fatal("Run ignored cancellation") } } // ---------------------------------------------------------------- waking // The delay this removes is on the path between a person walking in and their // face reaching a screen, so the test asserts a real wall-clock bound rather // than that a channel was read. func TestAWakeDrainsWithoutWaitingOutTheIdleInterval(t *testing.T) { q := queue(t) pub := &fakeBroker{connected: true} waker := NewWaker() p := &Pump{Queue: q, Publisher: pub, Wake: waker.C(), HeartbeatInterval: time.Hour} ctx, cancel := context.WithCancel(context.Background()) defer cancel() go p.Run(ctx) // Let it reach the idle wait with an empty queue first, so what follows is // genuinely the wake path and not the drain it does on startup. waitUntil(t, func() bool { return len(pub.delivered()) == 0 }, time.Second) time.Sleep(50 * time.Millisecond) start := time.Now() if err := q.Append("visit", "e1"); err != nil { t.Fatal(err) } waker.Wake() waitUntil(t, func() bool { return len(pub.delivered()) == 1 }, 2*time.Second) if took := time.Since(start); took >= idleInterval { t.Fatalf("took %s - the wake did not beat the %s idle tick", took, idleInterval) } } // A pump with no waker must behave exactly as it did before: a nil channel // blocks forever in a select, which is the correct fallback, not a hang. func TestAPumpWithNoWakerStillDrainsOnItsTimer(t *testing.T) { q := queue(t) pub := &fakeBroker{connected: true} p := &Pump{Queue: q, Publisher: pub, HeartbeatInterval: time.Hour} if err := q.Append("visit", "e1"); err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() go p.Run(ctx) waitUntil(t, func() bool { return len(pub.delivered()) == 1 }, 3*time.Second) } // The waker runs on the engine's webhook request. If it could ever block, a // burst of arrivals would apply backpressure into the recognition loop. func TestWakingNeverBlocksEvenWithNobodyListening(t *testing.T) { waker := NewWaker() done := make(chan struct{}) go func() { defer close(done) for i := 0; i < 10000; i++ { waker.Wake() } }() select { case <-done: case <-time.After(2 * time.Second): t.Fatal("Wake blocked with no pump reading - this would stall recognition") } } func TestANilWakerIsSafe(t *testing.T) { var w *Waker w.Wake() // an agent assembled without one must still run if w.C() != nil { t.Fatal("a nil waker handed out a channel") } } func waitUntil(t *testing.T, cond func() bool, within time.Duration) { t.Helper() deadline := time.Now().Add(within) for time.Now().Before(deadline) { if cond() { return } time.Sleep(2 * time.Millisecond) } t.Fatalf("condition not met within %s", within) }