Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -343,7 +343,7 @@ public void transportDataReceived(okio.Buffer frame, boolean endOfStream, int pa

@GuardedBy("lock")
private void onEndOfStream() {
if (!isOutboundClosed()) {
if (!isOutboundClosed() || outboundFlowState.hasPendingData()) {
// If server's end-of-stream is received before client sends end-of-stream, we just send a
// reset to server to fully close the server side stream.
transport.finishStream(id(),null, PROCESSED, false, ErrorCode.CANCEL, null);
Expand Down
51 changes: 51 additions & 0 deletions okhttp/src/test/java/io/grpc/okhttp/OkHttpClientTransportTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -1007,6 +1007,57 @@ public void outboundFlowControl() throws Exception {
shutdownAndVerify();
}

/**
* The server closes the call while the client's END_STREAM is still queued in the outbound flow
* controller. That END_STREAM will never be written, so the client must reset the stream;
* otherwise the server never sees the stream close.
*/
@Test
public void serverClosesWhileEndOfStreamBlockedByFlowControl_sendsReset() throws Exception {
initTransport();
MockStreamListener listener = new MockStreamListener();
ClientStream stream =
clientTransport.newStream(method, new Metadata(), CallOptions.DEFAULT, tracers);
stream.start(listener);

// Larger than the outbound window, so the tail of the message stays queued.
stream.writeMessage(new ByteArrayInputStream(new byte[INITIAL_WINDOW_SIZE]));
stream.flush();
verify(frameWriter, timeout(TIME_OUT_MS))
.data(eq(false), eq(3), any(Buffer.class), eq(INITIAL_WINDOW_SIZE));
// END_STREAM is queued behind the tail.
stream.halfClose();

frameHandler().headers(true, true, 3, 0, grpcResponseTrailers(), HeadersMode.HTTP_20_HEADERS);
listener.waitUntilStreamClosed();

assertEquals(Status.Code.OK, listener.status.getCode());
verify(frameWriter, timeout(TIME_OUT_MS)).rstStream(eq(3), eq(ErrorCode.CANCEL));
verify(frameWriter, never()).data(eq(true), eq(3), any(Buffer.class), anyInt());
shutdownAndVerify();
}

@Test
public void serverClosesAfterEndOfStreamSent_noReset() throws Exception {
initTransport();
MockStreamListener listener = new MockStreamListener();
ClientStream stream =
clientTransport.newStream(method, new Metadata(), CallOptions.DEFAULT, tracers);
stream.start(listener);

stream.writeMessage(new ByteArrayInputStream(new byte[10]));
stream.halfClose();
verify(frameWriter, timeout(TIME_OUT_MS))
.data(eq(true), eq(3), any(Buffer.class), eq(10 + HEADER_LENGTH));

frameHandler().headers(true, true, 3, 0, grpcResponseTrailers(), HeadersMode.HTTP_20_HEADERS);
listener.waitUntilStreamClosed();

assertEquals(Status.Code.OK, listener.status.getCode());
verify(frameWriter, never()).rstStream(eq(3), any(ErrorCode.class));
shutdownAndVerify();
}

/**
* Outbound flow control where the initial window size is reduced before a stream is started.
*/
Expand Down
Loading