Skip to content

Commit

Permalink
fix(pubsub): close grpc streams on retry (#10624)
Browse files Browse the repository at this point in the history
follow-up as 66581c4 is not enough to
plug all stream leaks

Co-authored-by: Alex Hong <9397363+hongalex@users.noreply.github.com>
  • Loading branch information
itizir and hongalex authored Sep 5, 2024
1 parent 02b2d12 commit 79a0e11
Showing 1 changed file with 18 additions and 11 deletions.
29 changes: 18 additions & 11 deletions pubsub/pullstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,9 @@ import (
// the stream on a retryable error.
type pullStream struct {
ctx context.Context
open func() (pb.Subscriber_StreamingPullClient, error)
cancel context.CancelFunc
cancel context.CancelFunc // cancel function of the context above
open func() (pb.Subscriber_StreamingPullClient, context.CancelFunc, error)
close context.CancelFunc // cancel function to close down the currently open stream

mu sync.Mutex
spc *pb.Subscriber_StreamingPullClient
Expand All @@ -50,8 +51,9 @@ func newPullStream(ctx context.Context, streamingPull streamingPullFunc, subName
return &pullStream{
ctx: ctx,
cancel: cancel,
open: func() (pb.Subscriber_StreamingPullClient, error) {
spc, err := streamingPull(ctx, gax.WithGRPCOptions(grpc.MaxCallRecvMsgSize(maxSendRecvBytes)))
open: func() (pb.Subscriber_StreamingPullClient, context.CancelFunc, error) {
sctx, close := context.WithCancel(ctx)
spc, err := streamingPull(sctx, gax.WithGRPCOptions(grpc.MaxCallRecvMsgSize(maxSendRecvBytes)))
if err == nil {
recordStat(ctx, StreamRequestCount, 1)
streamAckDeadline := int32(maxDurationPerLeaseExtension / time.Second)
Expand All @@ -69,9 +71,10 @@ func newPullStream(ctx context.Context, streamingPull streamingPullFunc, subName
})
}
if err != nil {
return nil, err
close()
return nil, nil, err
}
return spc, nil
return spc, close, nil
},
}
}
Expand Down Expand Up @@ -100,29 +103,33 @@ func (s *pullStream) get(spc *pb.Subscriber_StreamingPullClient) (*pb.Subscriber
if spc != s.spc {
return s.spc, nil
}
// we are about to open a new stream: if necessary, make sure the previous one is closed
if s.close != nil {
s.close()
}
// Either this is the very first call on this stream (s.spc == nil), or we have a valid
// retry request. Either way, open a new stream.
// The lock is held here for a long time, but it doesn't matter because no callers could get
// anything done anyway.
s.spc = new(pb.Subscriber_StreamingPullClient)
*s.spc, s.err = s.openWithRetry() // Any error from openWithRetry is permanent.
*s.spc, s.close, s.err = s.openWithRetry() // Any error from openWithRetry is permanent.
return s.spc, s.err
}

func (s *pullStream) openWithRetry() (pb.Subscriber_StreamingPullClient, error) {
func (s *pullStream) openWithRetry() (pb.Subscriber_StreamingPullClient, context.CancelFunc, error) {
r := defaultRetryer{}
for {
recordStat(s.ctx, StreamOpenCount, 1)
spc, err := s.open()
spc, close, err := s.open()
bo, shouldRetry := r.Retry(err)
if err != nil && shouldRetry {
recordStat(s.ctx, StreamRetryCount, 1)
if err := gax.Sleep(s.ctx, bo); err != nil {
return nil, err
return nil, nil, err
}
continue
}
return spc, err
return spc, close, err
}
}

Expand Down

0 comments on commit 79a0e11

Please sign in to comment.