From 77903ebebd6613950849d212b4fa5bbb82ea4869 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 24 Aug 2026 14:06:43 +0800 Subject: [PATCH] test(HH-604): stabilize notification close ordering (#164) Co-authored-by: Rogee --- .../notification_delivery_lifecycle_test.go | 39 ++++++++++++++++++- 1 file changed, 38 insertions(+), 1 deletion(-) diff --git a/backend/internal/service/notification_delivery_lifecycle_test.go b/backend/internal/service/notification_delivery_lifecycle_test.go index 8a8767fc..186ca128 100644 --- a/backend/internal/service/notification_delivery_lifecycle_test.go +++ b/backend/internal/service/notification_delivery_lifecycle_test.go @@ -39,6 +39,19 @@ type notificationAckGate struct { doneOnce sync.Once } +type notificationCloseSignalSubscriber struct { + message.Subscriber + started chan struct{} + release <-chan struct{} + once *sync.Once +} + +func (s *notificationCloseSignalSubscriber) Close() error { + s.once.Do(func() { close(s.started) }) + <-s.release + return s.Subscriber.Close() +} + func (g *notificationAckGate) DialHook(next redis.DialHook) redis.DialHook { return next } func (g *notificationAckGate) ProcessPipelineHook(next redis.ProcessPipelineHook) redis.ProcessPipelineHook { return next @@ -155,10 +168,23 @@ func TestNotificationShutdownWaitsForRedisStreamXAck(t *testing.T) { require.NoError(t, err) handlers := lifecycle.NewHandlerGroup() service.SetHandlerGroup(handlers) + subscriberCloseStarted := make(chan struct{}) + subscriberCloseRelease := make(chan struct{}) + subscriberCloseOnce := &sync.Once{} + service.router.AddSubscriberDecorators(func(subscriber message.Subscriber) (message.Subscriber, error) { + return ¬ificationCloseSignalSubscriber{ + Subscriber: subscriber, + started: subscriberCloseStarted, + release: subscriberCloseRelease, + once: subscriberCloseOnce, + }, nil + }) releaseAck := sync.OnceFunc(func() { close(gate.release) }) + releaseSubscriberClose := sync.OnceFunc(func() { close(subscriberCloseRelease) }) t.Cleanup(func() { releaseAck() + releaseSubscriberClose() _ = service.Close(context.Background()) }) runDone := make(chan error, 1) @@ -180,14 +206,25 @@ func TestNotificationShutdownWaitsForRedisStreamXAck(t *testing.T) { require.NoError(t, err) require.EqualValues(t, 1, pending.Count) - handlers.Stop() closeDone := make(chan error, 1) go func() { closeDone <- service.Close(context.Background()) }() + select { + case <-subscriberCloseStarted: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for subscriber close") + } + select { + case err := <-closeDone: + t.Fatalf("notification close returned before XAck release: %v", err) + default: + } require.NoError(t, publisherClient.Ping(context.Background()).Err()) releaseAck() + releaseSubscriberClose() require.NoError(t, <-closeDone) require.NoError(t, <-runDone) + handlers.Stop() select { case <-gate.done: default: