From 5028364e2c67dc3c2385a5ff9e2de2740bcd04cd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?J=C3=B6rn=20Friedrich=20Dreyer?= Date: Thu, 13 Aug 2026 12:37:11 +0200 Subject: [PATCH] add tests to cover more ack scenarios MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Jörn Friedrich Dreyer --- .../pkg/service/activitylog/debouncer_test.go | 183 ++++++++++++++++++ 1 file changed, 183 insertions(+) diff --git a/services/activitylog/pkg/service/activitylog/debouncer_test.go b/services/activitylog/pkg/service/activitylog/debouncer_test.go index 671c39590a..6305d51f80 100644 --- a/services/activitylog/pkg/service/activitylog/debouncer_test.go +++ b/services/activitylog/pkg/service/activitylog/debouncer_test.go @@ -12,6 +12,37 @@ import ( "github.com/opencloud-eu/opencloud/services/activitylog/pkg/service/activitylog" ) +// ackTracker tracks per-name ack invocations so we can assert each event's ack fires. +type ackTracker struct { + mu sync.Mutex + calls map[string]int + errors map[string]error +} + +func newAckTracker() *ackTracker { + return &ackTracker{ + calls: make(map[string]int), + errors: make(map[string]error), + } +} + +func (t *ackTracker) track(name string, fn func() error) func() error { + return func() error { + t.mu.Lock() + t.calls[name]++ + err := fn() + t.errors[name] = err + t.mu.Unlock() + return err + } +} + +func (t *ackTracker) count(name string) int { + t.mu.Lock() + defer t.mu.Unlock() + return t.calls[name] +} + var _ = Describe("Debouncer", func() { var ( testlog = log.NopLogger() @@ -249,6 +280,158 @@ var _ = Describe("Debouncer", func() { Expect(callbacks[0].EventID).To(Equal("act1")) Expect(callbacks[1].EventID).To(Equal("act2")) }) + + It("calls ack even when callback returns an error", func() { + cbErr = fmt.Errorf("store failed") + d := activitylog.NewDebouncer(testlog, 50*time.Millisecond, newCallback) + + ra := data.RawActivity{EventID: "act1"} + d.Debounce("space1", ra, newAck) + + Eventually(func() int { + ackMu.Lock() + defer ackMu.Unlock() + return ackCalled + }).Should(Equal(1)) + }) + + It("calls each ack when multiple events target the same id within one batch", func() { + tracker := newAckTracker() + d := activitylog.NewDebouncer(testlog, 200*time.Millisecond, newCallback) + + ra1 := data.RawActivity{EventID: "act1"} + ra2 := data.RawActivity{EventID: "act2"} + ra3 := data.RawActivity{EventID: "act3"} + + d.Debounce("space1", ra1, tracker.track("ack1", func() error { return nil })) + d.Debounce("space1", ra2, tracker.track("ack2", func() error { return nil })) + d.Debounce("space1", ra3, tracker.track("ack3", func() error { return nil })) + + // All three activities should be stored ... + Eventually(func() int { + mu.Lock() + defer mu.Unlock() + return len(callbacks) + }).Should(Equal(3)) + + // ... and all three acks should have been called. + Eventually(func() int { return tracker.count("ack1") }).Should(Equal(1)) + Eventually(func() int { return tracker.count("ack2") }).Should(Equal(1)) + Eventually(func() int { return tracker.count("ack3") }).Should(Equal(1)) + }) + + It("calls ack on the re-schedule path when timer fires during in-progress callback", func() { + tracker := newAckTracker() + unblockCh := make(chan struct{}) + + blockingCallback := func(_ string, ras []data.RawActivity) error { + mu.Lock() + callbacks = append(callbacks, ras...) + mu.Unlock() + + <-unblockCh + return nil + } + + d := activitylog.NewDebouncer(testlog, 50*time.Millisecond, blockingCallback) + + // First activity starts the timer. + ra1 := data.RawActivity{EventID: "act1"} + d.Debounce("space1", ra1, tracker.track("ack1", func() error { return nil })) + + // Wait for timer to fire so the blocking callback begins. + Eventually(func() int { + mu.Lock() + defer mu.Unlock() + return len(callbacks) + }, 1*time.Second).Should(Equal(1)) + + // Queue a second activity on a *different* id so its own timer fires + // while the first callback is still blocked. When that timer fires, + // it will complete immediately (no blocking), then we'll queue onto + // space1 again to hit the re-schedule path. + ra2 := data.RawActivity{EventID: "act2"} + d.Debounce("space2", ra2, tracker.track("ack2", func() error { return nil })) + + // Wait for space2 batch to complete (non-blocking callback). + Eventually(func() int { + mu.Lock() + defer mu.Unlock() + return len(callbacks) + }, 1*time.Second).Should(Equal(2)) + + // Now queue onto space1 while its callback is still in-progress. + // The timer for space1 has already fired once; new Debounce calls + // will append to the existing pending item and reset the timer, + // which triggers the re-schedule path when it fires again. + ra3 := data.RawActivity{EventID: "act3"} + d.Debounce("space1", ra3, tracker.track("ack3", func() error { return nil })) + + // Unblock the first callback so space1 can process act3. + close(unblockCh) + + Eventually(func() int { + mu.Lock() + defer mu.Unlock() + return len(callbacks) + }, 2*time.Second).Should(Equal(3)) + + // All three acks should have been called. + Eventually(func() int { return tracker.count("ack1") }).Should(Equal(1)) + Eventually(func() int { return tracker.count("ack2") }).Should(Equal(1)) + Eventually(func() int { return tracker.count("ack3") }).Should(Equal(1)) + }) + + It("acks event even when store callback fails (data-loss risk)", func() { + // This test documents current behaviour: ack fires regardless of + // whether StoreActivity succeeded. If storage fails but the event + // is still acked, it will never be replayed by JetStream. + cbErr = fmt.Errorf("nats kv put failed") + d := activitylog.NewDebouncer(testlog, 50*time.Millisecond, newCallback) + + ra := data.RawActivity{EventID: "act1"} + d.Debounce("space1", ra, newAck) + + Eventually(func() int { + mu.Lock() + defer mu.Unlock() + return len(callbacks) + }).Should(Equal(1)) + + // ack was called despite the store error -- event is lost. + Eventually(func() int { + ackMu.Lock() + defer ackMu.Unlock() + return ackCalled + }).Should(Equal(1)) + }) + }) + + Context("zero-duration ack behaviour", func() { + It("calls ack when duration is zero (immediate path)", func() { + d := activitylog.NewDebouncer(testlog, 0, newCallback) + + ra := data.RawActivity{EventID: "immediate"} + d.Debounce("space1", ra, newAck) + + ackMu.Lock() + defer ackMu.Unlock() + Expect(ackCalled).To(Equal(1)) + }) + + It("calls each ack when multiple immediate events target the same id", func() { + tracker := newAckTracker() + d := activitylog.NewDebouncer(testlog, 0, newCallback) + + ra1 := data.RawActivity{EventID: "act1"} + ra2 := data.RawActivity{EventID: "act2"} + + d.Debounce("space1", ra1, tracker.track("ack1", func() error { return nil })) + d.Debounce("space1", ra2, tracker.track("ack2", func() error { return nil })) + + Expect(tracker.count("ack1")).To(Equal(1)) + Expect(tracker.count("ack2")).To(Equal(1)) + }) }) })