From 4f48c2298ac925484c30198f5891652645519368 Mon Sep 17 00:00:00 2001 From: Ardit Marku Date: Sun, 23 Aug 2026 17:16:34 +0300 Subject: [PATCH 1/3] Remove entirely blocking on gracefulDone channel readiness --- module/grpcserver/server.go | 1 - module/grpcserver/server_test.go | 23 ++++++++++++++++++++--- 2 files changed, 20 insertions(+), 4 deletions(-) diff --git a/module/grpcserver/server.go b/module/grpcserver/server.go index 1af83735b0d..14c756950f2 100644 --- a/module/grpcserver/server.go +++ b/module/grpcserver/server.go @@ -137,5 +137,4 @@ func (g *GrpcServer) shutdownWorker(ctx irrecoverable.SignalerContext, ready com Msg("graceful stop timed out; force-stopping gRPC server") g.server.Stop() } - <-gracefulDone } diff --git a/module/grpcserver/server_test.go b/module/grpcserver/server_test.go index 9ae439f343a..0bac6f4e6db 100644 --- a/module/grpcserver/server_test.go +++ b/module/grpcserver/server_test.go @@ -39,12 +39,17 @@ var blockingStreamServiceDesc = grpc.ServiceDesc{ type blockingStreamServer struct { // started is closed when the stream handler has been entered. started chan struct{} + // blockDuration will cause `Stream` to block for the set duration + blockDuration time.Duration } var _ blockingStreamService = (*blockingStreamServer)(nil) func (s *blockingStreamServer) Stream(stream grpc.ServerStream) error { close(s.started) + if s.blockDuration > 0 { + time.Sleep(s.blockDuration) + } <-stream.Context().Done() return nil } @@ -60,7 +65,10 @@ func TestGrpcServerShutdown_WithActiveStream(t *testing.T) { gracefulStopTimeout := 200 * time.Millisecond rawServer := grpc.NewServer() - handler := &blockingStreamServer{started: make(chan struct{})} + handler := &blockingStreamServer{ + started: make(chan struct{}), + blockDuration: gracefulStopTimeout * 10, + } rawServer.RegisterService(&blockingStreamServiceDesc, handler) signalerCtx := atomic.NewPointer[irrecoverable.SignalerContext](nil) @@ -114,7 +122,10 @@ func TestGrpcServerShutdown_ShutdownStreamInterceptor(t *testing.T) { rawServer := grpc.NewServer( grpc.ChainStreamInterceptor(grpcserver.ShutdownStreamInterceptor(signalerCtx)), ) - handler := &blockingStreamServer{started: make(chan struct{})} + handler := &blockingStreamServer{ + started: make(chan struct{}), + blockDuration: time.Second, + } rawServer.RegisterService(&blockingStreamServiceDesc, handler) server := grpcserver.NewGrpcServer( @@ -159,7 +170,13 @@ func TestGrpcServerShutdown_NoActiveStreams(t *testing.T) { gracefulStopTimeout := 5 * time.Second rawServer := grpc.NewServer() - rawServer.RegisterService(&blockingStreamServiceDesc, &blockingStreamServer{started: make(chan struct{})}) + rawServer.RegisterService( + &blockingStreamServiceDesc, + &blockingStreamServer{ + started: make(chan struct{}), + blockDuration: gracefulStopTimeout * 10, + }, + ) signalerCtx := atomic.NewPointer[irrecoverable.SignalerContext](nil) server := grpcserver.NewGrpcServer( From 54cdafe1a1e8f7797822f9d333c5162d0b9ba51b Mon Sep 17 00:00:00 2001 From: Ardit Marku Date: Wed, 26 Aug 2026 09:17:38 +0300 Subject: [PATCH 2/3] Update comment regarding time.Sleep on Stream() Co-authored-by: Leo Zhang --- module/grpcserver/server_test.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/module/grpcserver/server_test.go b/module/grpcserver/server_test.go index 0bac6f4e6db..46002f3c8ee 100644 --- a/module/grpcserver/server_test.go +++ b/module/grpcserver/server_test.go @@ -48,6 +48,10 @@ var _ blockingStreamService = (*blockingStreamServer)(nil) func (s *blockingStreamServer) Stream(stream grpc.ServerStream) error { close(s.started) if s.blockDuration > 0 { + // this is to simulate the case that after `grpcServer.Stop()` is called, + // `<-gracefulDone` channel is still blocking, so that we can verify + // the caller is not waiting for `<-gracefulDone` return before shutdown, + // otherwise, the waiting might be still blocking for longer or indefinitely. time.Sleep(s.blockDuration) } <-stream.Context().Done() From 9dec93491f52254c41af21dc762532d87083e8b1 Mon Sep 17 00:00:00 2001 From: Ardit Marku Date: Wed, 26 Aug 2026 09:27:10 +0300 Subject: [PATCH 3/3] Fix linting issue around comment --- module/grpcserver/server_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/module/grpcserver/server_test.go b/module/grpcserver/server_test.go index 46002f3c8ee..a2083a1edc6 100644 --- a/module/grpcserver/server_test.go +++ b/module/grpcserver/server_test.go @@ -48,7 +48,7 @@ var _ blockingStreamService = (*blockingStreamServer)(nil) func (s *blockingStreamServer) Stream(stream grpc.ServerStream) error { close(s.started) if s.blockDuration > 0 { - // this is to simulate the case that after `grpcServer.Stop()` is called, + // this is to simulate the case that after `grpcServer.Stop()` is called, // `<-gracefulDone` channel is still blocking, so that we can verify // the caller is not waiting for `<-gracefulDone` return before shutdown, // otherwise, the waiting might be still blocking for longer or indefinitely.