Skip to content

Commit d2c52f0

Browse files
transport: fix server Recv hang on initial HEADERS with EndStream (#9391)
Fix server-side Recv hang when client sends HEADERS with EndStream set. When a client initiates a stream with a HEADERS frame that has EndStream=true, the server-side stream state is set to streamReadDone, but the reader was not notified. This caused server-side Recv() calls to hang indefinitely waiting for data. This change writes an io.EOF error to the stream's receive buffer when the stream is created with EndStream=true, ensuring Recv() returns immediately with io.EOF. A new transport test has been added to verify this behavior, and existing tests have been updated. RELEASE NOTES: * transport: fix an issue where server-side `Recv()` hung indefinitely when a client initiated a stream with `HEADERS` having `EndStream` set.
1 parent a9cbf17 commit d2c52f0

3 files changed

Lines changed: 64 additions & 6 deletions

File tree

internal/transport/http2_server.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -531,6 +531,7 @@ func (t *http2Server) operateHeaders(ctx context.Context, frame *http2.MetaHeade
531531
if frame.StreamEnded() {
532532
// s is just created by the caller. No lock needed.
533533
s.state = streamReadDone
534+
s.write(recvMsg{err: io.EOF})
534535
}
535536
if timeoutSet {
536537
s.ctx, s.cancel = context.WithTimeout(ctx, timeout)

test/end2end_test.go

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4765,14 +4765,20 @@ func testClientInitialHeaderEndStream(t *testing.T, e env) {
47654765
te := newTest(t, e)
47664766
ts := &funcServer{streamingInputCall: func(stream testgrpc.TestService_StreamingInputCallServer) error {
47674767
defer close(handlerDone)
4768-
// Block on serverTester receiving RST_STREAM. This ensures server has closed
4769-
// stream before stream.Recv().
4768+
// Block on serverTester receiving RST_STREAM. This ensures server has
4769+
// closed stream before stream.Recv().
47704770
<-frameCheckingDone
4771-
data, err := stream.Recv()
4772-
if err == nil {
4773-
t.Errorf("unexpected data received in func server method: '%v'", data)
4771+
// Depending on whether the context cancellation (due to the illegal data
4772+
// RST_STREAM) or the buffered EOF (from the initial HEADERS END_STREAM) is
4773+
// selected first in recvBufferReader, stream.Recv() can return either
4774+
// io.EOF or Canceled.
4775+
if _, err := stream.Recv(); err != io.EOF && status.Code(err) != codes.Canceled {
4776+
t.Errorf("stream.Recv() returned error = %v, expected EOF or canceled error", err)
4777+
}
4778+
if err := stream.SendMsg(nil); err == nil {
4779+
t.Error("stream.SendMsg() returned nil, expected cancel error")
47744780
} else if status.Code(err) != codes.Canceled {
4775-
t.Errorf("expected canceled error, instead received '%v'", err)
4781+
t.Errorf("stream.SendMsg() returned error = %v, expected cancel error", err)
47764782
}
47774783
return nil
47784784
}}

test/transport_test.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"net"
2525
"sync"
2626
"testing"
27+
"time"
2728

2829
"golang.org/x/net/http2"
2930
"google.golang.org/grpc"
@@ -302,3 +303,53 @@ func (s) TestCancelWhileServerWaitingForFlowControl(t *testing.T) {
302303
t.Fatalf("Failed to read from the stream: %v", err)
303304
}
304305
}
306+
307+
// Tests that when a client sends a HEADERS frame with EndStream=true, the
308+
// server-side stream receives an io.EOF on Recv() and does not hang waiting
309+
// for data frames.
310+
func (s) TestHeadersEndStreamNoHang(t *testing.T) {
311+
receivedErr := make(chan error, 1)
312+
ss := &stubserver.StubServer{
313+
FullDuplexCallF: func(stream testgrpc.TestService_FullDuplexCallServer) error {
314+
_, err := stream.Recv()
315+
receivedErr <- err
316+
return nil
317+
},
318+
}
319+
if err := ss.Start(nil); err != nil {
320+
t.Fatalf("Error starting endpoint server: %v", err)
321+
}
322+
defer ss.Stop()
323+
324+
conn, err := net.DialTimeout("tcp", ss.Address, defaultTestTimeout)
325+
if err != nil {
326+
t.Fatalf("Failed to dial: %v", err)
327+
}
328+
defer conn.Close()
329+
330+
st := newServerTesterFromConn(t, conn)
331+
st.greet()
332+
333+
// Send HEADERS with EndStream = true and no grpc-timeout header.
334+
st.writeHeaders(http2.HeadersFrameParam{
335+
StreamID: 1,
336+
BlockFragment: st.encodeHeader(
337+
":method", "POST",
338+
":path", "/grpc.testing.TestService/FullDuplexCall",
339+
":authority", "localhost",
340+
"content-type", "application/grpc",
341+
"te", "trailers",
342+
),
343+
EndStream: true,
344+
EndHeaders: true,
345+
})
346+
347+
select {
348+
case err := <-receivedErr:
349+
if err != io.EOF {
350+
t.Errorf("Streaming handler returned error = %v, expected io.EOF", err)
351+
}
352+
case <-time.After(defaultTestTimeout):
353+
t.Fatalf("Timed out waiting for Recv() on the server to complete")
354+
}
355+
}

0 commit comments

Comments
 (0)