test(HH-604): stabilize notification close ordering (#164)
Co-authored-by: Rogee <rogee@ipao.vip>
This commit is contained in:
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user