Skip to content

Commit af55d61

Browse files
author
Michal Witkowski
authored
Merge pull request mwitkow#10 from mwitkow/bugfix/streaming-fix
Streaming Fix: handle bidirectional case
2 parents c2f7c98 + de4d3db commit af55d61

File tree

4 files changed

+162
-45
lines changed

4 files changed

+162
-45
lines changed

proxy/handler.go

Lines changed: 35 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ package proxy
66
import (
77
"io"
88

9+
"golang.org/x/net/context"
910
"google.golang.org/grpc"
1011
"google.golang.org/grpc/codes"
1112
"google.golang.org/grpc/transport"
@@ -64,27 +65,49 @@ func (s *handler) handler(srv interface{}, serverStream grpc.ServerStream) error
6465
return grpc.Errorf(codes.Internal, "lowLevelServerStream not exists in context")
6566
}
6667
fullMethodName := lowLevelServerStream.Method()
68+
clientCtx, clientCancel := context.WithCancel(serverStream.Context())
6769
backendConn, err := s.director(serverStream.Context(), fullMethodName)
6870
if err != nil {
6971
return err
7072
}
7173
// TODO(mwitkow): Add a `forwarded` header to metadata, https://en.wikipedia.org/wiki/X-Forwarded-For.
72-
clientStream, err := grpc.NewClientStream(serverStream.Context(), clientStreamDescForProxying, backendConn, fullMethodName)
74+
clientStream, err := grpc.NewClientStream(clientCtx, clientStreamDescForProxying, backendConn, fullMethodName)
7375
if err != nil {
7476
return err
7577
}
76-
defer clientStream.CloseSend() // always close this!
77-
s2cErr := <-s.forwardServerToClient(serverStream, clientStream)
78-
c2sErr := <-s.forwardClientToServer(clientStream, serverStream)
79-
if s2cErr != io.EOF {
80-
return grpc.Errorf(codes.Internal, "failed proxying s2c: %v", s2cErr)
81-
}
82-
serverStream.SetTrailer(clientStream.Trailer())
83-
// c2sErr will contain RPC error from client code. If not io.EOF return the RPC error as server stream error.
84-
if c2sErr != io.EOF {
85-
return c2sErr
78+
s2cErrChan := s.forwardServerToClient(serverStream, clientStream)
79+
defer close(s2cErrChan)
80+
c2sErrChan := s.forwardClientToServer(clientStream, serverStream)
81+
defer close(c2sErrChan)
82+
// We don't know which side is going to stop sending first, so we need a select between the two.
83+
for i := 0; i < 2; i++ {
84+
select {
85+
case s2cErr := <-s2cErrChan:
86+
if s2cErr == io.EOF {
87+
// this is the happy case where the sender has encountered io.EOF, and won't be sending anymore./
88+
// the clientStream>serverStream may continue pumping though.
89+
clientStream.CloseSend()
90+
break
91+
} else {
92+
// however, we may have gotten a receive error (stream disconnected, a read error etc) in which case we need
93+
// to cancel the clientStream to the backend, let all of its goroutines be freed up by the CancelFunc and
94+
// exit with an error to the stack
95+
clientCancel()
96+
return grpc.Errorf(codes.Internal, "failed proxying s2c: %v", s2cErr)
97+
}
98+
case c2sErr := <-c2sErrChan:
99+
// This happens when the clientStream has nothing else to offer (io.EOF), returned a gRPC error. In those two
100+
// cases we may have received Trailers as part of the call. In case of other errors (stream closed) the trailers
101+
// will be nil.
102+
serverStream.SetTrailer(clientStream.Trailer())
103+
// c2sErr will contain RPC error from client code. If not io.EOF return the RPC error as server stream error.
104+
if c2sErr != io.EOF {
105+
return c2sErr
106+
}
107+
return nil
108+
}
86109
}
87-
return nil
110+
return grpc.Errorf(codes.Internal, "gRPC proxying should never reach this stage.")
88111
}
89112

90113
func (s *handler) forwardClientToServer(src grpc.ClientStream, dst grpc.ServerStream) chan error {
@@ -115,7 +138,6 @@ func (s *handler) forwardClientToServer(src grpc.ClientStream, dst grpc.ServerSt
115138
break
116139
}
117140
}
118-
close(ret)
119141
}()
120142
return ret
121143
}
@@ -134,7 +156,6 @@ func (s *handler) forwardServerToClient(src grpc.ServerStream, dst grpc.ClientSt
134156
break
135157
}
136158
}
137-
close(ret)
138159
}()
139160
return ret
140161
}

proxy/handler_test.go

Lines changed: 42 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ import (
2222
"google.golang.org/grpc/grpclog"
2323
"google.golang.org/grpc/metadata"
2424

25+
"fmt"
26+
2527
pb "github.com/mwitkow/grpc-proxy/testservice"
2628
)
2729

@@ -72,6 +74,27 @@ func (s *assertingService) PingList(ping *pb.PingRequest, stream pb.TestService_
7274
return nil
7375
}
7476

77+
func (s *assertingService) PingStream(stream pb.TestService_PingStreamServer) error {
78+
stream.SendHeader(metadata.Pairs(serverHeaderMdKey, "I like turtles."))
79+
counter := int32(0)
80+
for {
81+
ping, err := stream.Recv()
82+
if err == io.EOF {
83+
break
84+
} else if err != nil {
85+
require.NoError(s.t, err, "can't fail reading stream")
86+
return err
87+
}
88+
pong := &pb.PingResponse{Value: ping.Value, Counter: counter}
89+
if err := stream.Send(pong); err != nil {
90+
require.NoError(s.t, err, "can't fail sending back a pong")
91+
}
92+
counter += 1
93+
}
94+
stream.SetTrailer(metadata.Pairs(serverTrailerMdKey, "I like ending turtles."))
95+
return nil
96+
}
97+
7598
// ProxyHappySuite tests the "happy" path of handling: that everything works in absence of connection issues.
7699
type ProxyHappySuite struct {
77100
suite.Suite
@@ -125,24 +148,28 @@ func (s *ProxyHappySuite) TestDirectorErrorIsPropagated() {
125148
assert.Equal(s.T(), "testing rejection", grpc.ErrorDesc(err))
126149
}
127150

128-
func (s *ProxyHappySuite) TestPingListStreamsAll() {
129-
stream, err := s.testClient.PingList(s.ctx(), &pb.PingRequest{Value: "foo"})
130-
require.NoError(s.T(), err, "PingList request should be successful.")
131-
// Check that the header arrives before all entries.
132-
headerMd, err := stream.Header()
133-
require.NoError(s.T(), err, "PingList headers should not error.")
134-
assert.Len(s.T(), headerMd, 1, "PingList response headers user contain metadata")
135-
count := 0
136-
for {
151+
func (s *ProxyHappySuite) TestPingStream_FullDuplexWorks() {
152+
stream, err := s.testClient.PingStream(s.ctx())
153+
require.NoError(s.T(), err, "PingStream request should be successful.")
154+
155+
for i := 0; i < countListResponses; i++ {
156+
ping := &pb.PingRequest{Value: fmt.Sprintf("foo:%d", i)}
157+
require.NoError(s.T(), stream.Send(ping), "sending to PingStream must not fail")
137158
resp, err := stream.Recv()
138159
if err == io.EOF {
139160
break
140161
}
141-
require.NoError(s.T(), err, "PingList stream should not be interrupted.")
142-
require.Equal(s.T(), "foo", resp.Value)
143-
count = count + 1
162+
if i == 0 {
163+
// Check that the header arrives before all entries.
164+
headerMd, err := stream.Header()
165+
require.NoError(s.T(), err, "PingStream headers should not error.")
166+
assert.Len(s.T(), headerMd, 1, "PingStream response headers user contain metadata")
167+
}
168+
assert.EqualValues(s.T(), i, resp.Counter, "ping roundtrip must succeed with the correct id")
144169
}
145-
assert.Equal(s.T(), countListResponses, count, "PingList must successfully return all outputs")
170+
require.NoError(s.T(), stream.CloseSend(), "no error on close send")
171+
_, err = stream.Recv()
172+
require.Equal(s.T(), io.EOF, err, "stream should close with io.EOF, meaining OK")
146173
// Check that the trailer headers are here.
147174
trailerMd := stream.Trailer()
148175
assert.Len(s.T(), trailerMd, 1, "PingList trailer headers user contain metadata")
@@ -183,12 +210,12 @@ func (s *ProxyHappySuite) SetupSuite() {
183210
"Ping")
184211

185212
// Start the serving loops.
213+
s.T().Logf("starting grpc.Server at: %v", s.serverListener.Addr().String())
186214
go func() {
187-
s.T().Logf("starting grpc.Server at: %v", s.serverListener.Addr().String())
188215
s.server.Serve(s.serverListener)
189216
}()
217+
s.T().Logf("starting grpc.Proxy at: %v", s.proxyListener.Addr().String())
190218
go func() {
191-
s.T().Logf("starting grpc.Proxy at: %v", s.proxyListener.Addr().String())
192219
s.proxy.Serve(s.proxyListener)
193220
}()
194221

testservice/test.pb.go

Lines changed: 82 additions & 16 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

testservice/test.proto

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,5 +22,8 @@ service TestService {
2222
rpc PingError(PingRequest) returns (Empty) {}
2323

2424
rpc PingList(PingRequest) returns (stream PingResponse) {}
25+
26+
rpc PingStream(stream PingRequest) returns (stream PingResponse) {}
27+
2528
}
2629

0 commit comments

Comments
 (0)