From 48f82ab70492beb4537e8b157a5ae06fab773970 Mon Sep 17 00:00:00 2001 From: Benjamin Peterson Date: Wed, 30 Sep 2026 06:47:04 -0700 Subject: [PATCH] okhttp: reset stream when trailers arrive with END_STREAM unsent This is the OkHttp counterpart of netty's https://github.com/grpc/grpc-java/issues/13078 and ba0aab4e082ec60e4ff4e81dce29b0178f82387a. --- .../io/grpc/okhttp/OkHttpClientStream.java | 2 +- .../okhttp/OkHttpClientTransportTest.java | 51 +++++++++++++++++++ 2 files changed, 52 insertions(+), 1 deletion(-) diff --git a/okhttp/src/main/java/io/grpc/okhttp/OkHttpClientStream.java b/okhttp/src/main/java/io/grpc/okhttp/OkHttpClientStream.java index 8dd55d9f23e..b1e65ec57fa 100644 --- a/okhttp/src/main/java/io/grpc/okhttp/OkHttpClientStream.java +++ b/okhttp/src/main/java/io/grpc/okhttp/OkHttpClientStream.java @@ -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); diff --git a/okhttp/src/test/java/io/grpc/okhttp/OkHttpClientTransportTest.java b/okhttp/src/test/java/io/grpc/okhttp/OkHttpClientTransportTest.java index 0b571530db4..1958065e176 100644 --- a/okhttp/src/test/java/io/grpc/okhttp/OkHttpClientTransportTest.java +++ b/okhttp/src/test/java/io/grpc/okhttp/OkHttpClientTransportTest.java @@ -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. */