From 27ea28ac165bfa6ff0fb10ead8d8e01d2eeb5010 Mon Sep 17 00:00:00 2001 From: Meemaw Date: Thu, 1 Oct 2026 13:15:49 +0200 Subject: [PATCH 1/2] netty: drain before closing a gracefully shut down server connection After the second GOAWAY, Http2ConnectionHandler closed the socket as soon as the last stream's final frame was written to the kernel. If the client had not read the whole response yet, the socket was closed with response bytes still queued, and the client's next frame (WINDOW_UPDATE, BDP or keepalive PING) made the server's kernel answer with RST. The RST drops the untransmitted bytes, including the trailers, and the client fails an RPC the server completed with "Connection closed after GOAWAY ... max_age". Response headers were already delivered, so the call is not retried. Keep the connection open and reading for up to 1 second once no streams are active, mirroring golang/net cd69bc3. The client normally closes first, as soon as it has read the GOAWAY and its streams completed. The drain never extends past the configured grace time. Fixes #9566 Co-Authored-By: Claude Opus 5.5 --- .../io/grpc/netty/NettyServerHandler.java | 62 +++++++++++++++ .../io/grpc/netty/NettyServerHandlerTest.java | 79 +++++++++++++++++-- 2 files changed, 133 insertions(+), 8 deletions(-) diff --git a/netty/src/main/java/io/grpc/netty/NettyServerHandler.java b/netty/src/main/java/io/grpc/netty/NettyServerHandler.java index 846c45cd459..c98f313a32f 100644 --- a/netty/src/main/java/io/grpc/netty/NettyServerHandler.java +++ b/netty/src/main/java/io/grpc/netty/NettyServerHandler.java @@ -122,6 +122,14 @@ class NettyServerHandler extends AbstractNettyHandler { @VisibleForTesting static final long GRACEFUL_SHUTDOWN_PING = 0x97ACEF001L; private static final long GRACEFUL_SHUTDOWN_PING_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(10); + /** + * How long to keep the connection open and reading after the second GOAWAY, once no streams are + * active. Closing a socket while the peer is still sending (WINDOW_UPDATE, PING) makes the kernel + * answer with RST, which discards response bytes that have not been transmitted yet (#9566). + * The peer normally closes first, as soon as it has read the GOAWAY and its streams completed. + */ + @VisibleForTesting + static final long GRACEFUL_SHUTDOWN_DRAIN_NANOS = TimeUnit.SECONDS.toNanos(1); /** Temporary workaround for #8674. Fine to delete after v1.45 release, and maybe earlier. */ private static final boolean DISABLE_CONNECTION_HEADER_CHECK = Boolean.parseBoolean( System.getProperty("io.grpc.netty.disableConnectionHeaderCheck", "false")); @@ -369,6 +377,9 @@ public void onStreamClosed(Http2Stream stream) { if (maxConnectionIdleManager != null) { maxConnectionIdleManager.onTransportIdle(); } + if (gracefulShutdown != null) { + gracefulShutdown.drainIfIdle(); + } } } }); @@ -694,6 +705,9 @@ public void channelInactive(ChannelHandlerContext ctx) throws Exception { if (maxConnectionAgeMonitor != null) { maxConnectionAgeMonitor.cancel(false); } + if (gracefulShutdown != null) { + gracefulShutdown.cancelDrain(); + } final Status status = Status.UNAVAILABLE.withDescription("connection terminated for unknown reason"); // Any streams that are still active must be closed @@ -745,6 +759,12 @@ public void close(ChannelHandlerContext ctx, ChannelPromise promise) throws Exce ctx.flush(); } + @Override + protected boolean isGracefulShutdownComplete() { + return super.isGracefulShutdownComplete() + && (gracefulShutdown == null || gracefulShutdown.drainComplete()); + } + /** * Returns the given processed bytes back to inbound flow control. */ @@ -1090,6 +1110,13 @@ private final class GracefulShutdown { Future pingFuture; + ChannelHandlerContext ctx; + + /** Scheduled once the second GOAWAY has been sent and no streams are active. */ + Future drainFuture; + + boolean drained; + GracefulShutdown(String goAwayMessage, @Nullable Long graceTimeInNanos) { this.goAwayMessage = goAwayMessage; @@ -1100,6 +1127,7 @@ private final class GracefulShutdown { * Sends out first GOAWAY and ping, and schedules second GOAWAY and close. */ void start(final ChannelHandlerContext ctx) { + this.ctx = ctx; goAway( ctx, Integer.MAX_VALUE, @@ -1142,12 +1170,46 @@ void secondGoAwayAndClose(ChannelHandlerContext ctx) { long overriddenGraceTime = graceTimeOverrideMillis(savedGracefulShutdownTimeMillis); try { gracefulShutdownTimeoutMillis(overriddenGraceTime); + // Closes once isGracefulShutdownComplete(), i.e. after the drain, or when the grace time + // runs out. NettyServerHandler.super.close(ctx, ctx.newPromise()); } catch (Exception e) { onError(ctx, /* outbound= */ true, e); } finally { gracefulShutdownTimeoutMillis(savedGracefulShutdownTimeMillis); } + drainIfIdle(); + } + + boolean drainComplete() { + return !pingAckedOrTimeout || drained; + } + + void drainIfIdle() { + if (!pingAckedOrTimeout || drainFuture != null || connection().numActiveStreams() != 0) { + return; + } + drainFuture = ctx.executor().schedule( + new Runnable() { + @Override + public void run() { + drained = true; + try { + // No streams are active, so this closes as soon as pending writes are flushed. + NettyServerHandler.super.close(ctx, ctx.newPromise()); + } catch (Exception e) { + onError(ctx, /* outbound= */ true, e); + } + } + }, + GRACEFUL_SHUTDOWN_DRAIN_NANOS, + TimeUnit.NANOSECONDS); + } + + void cancelDrain() { + if (drainFuture != null) { + drainFuture.cancel(false); + } } private long graceTimeOverrideMillis(long originalMillis) { diff --git a/netty/src/test/java/io/grpc/netty/NettyServerHandlerTest.java b/netty/src/test/java/io/grpc/netty/NettyServerHandlerTest.java index 84a1a48b37f..f373384d0a0 100644 --- a/netty/src/test/java/io/grpc/netty/NettyServerHandlerTest.java +++ b/netty/src/test/java/io/grpc/netty/NettyServerHandlerTest.java @@ -366,7 +366,9 @@ public void closeShouldGracefullyCloseChannel() throws Exception { verifyWrite().writeGoAway(eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), isA(ByteBuf.class), any(ChannelPromise.class)); - // Verify that the channel was closed. + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -388,7 +390,9 @@ public void gracefulCloseShouldGracefullyCloseChannel() throws Exception { verifyWrite().writeGoAway(eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), isA(ByteBuf.class), any(ChannelPromise.class)); - // Verify that the channel was closed. + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -417,6 +421,8 @@ public void secondGracefulCloseIsSafe() throws Exception { channelRead(pingFrame(/*ack=*/ true , NettyServerHandler.GRACEFUL_SHUTDOWN_PING)); verifyWrite().writeGoAway(eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), isA(ByteBuf.class), any(ChannelPromise.class)); + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -992,7 +998,9 @@ public void maxConnectionIdle_goAwaySent_pingAck() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -1022,7 +1030,9 @@ public void maxConnectionIdle_goAwaySent_pingTimeout() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -1066,7 +1076,9 @@ public void maxConnectionIdle_activeThenRst_pingAck() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(STREAM_ID), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -1110,7 +1122,9 @@ public void maxConnectionIdle_activeThenRst_pingTimeoutk() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(STREAM_ID), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -1164,7 +1178,9 @@ public void maxConnectionAge_goAwaySent_pingAck() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); assertFalse(channel().isOpen()); } @@ -1195,7 +1211,54 @@ public void maxConnectionAge_goAwaySent_pingTimeout() throws Exception { verifyWrite().writeGoAway( eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), any(ChannelPromise.class)); - // channel closed + // channel stays open while draining, then closes + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); + assertFalse(channel().isOpen()); + } + + @Test + public void maxConnectionAge_drainStartsWhenLastStreamCloses() throws Exception { + maxConnectionAgeInNanos = TimeUnit.MILLISECONDS.toNanos(10L); + manualSetUp(); + createStream(); + + fakeClock().forwardNanos(maxConnectionAgeInNanos); + channelRead(pingFrame(true /* isAck */, NettyServerHandler.GRACEFUL_SHUTDOWN_PING)); + + // second GO_AWAY sent + verifyWrite().writeGoAway( + eq(ctx()), eq(STREAM_ID), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), + any(ChannelPromise.class)); + // stream still active, so the drain has not started + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS); + assertTrue(channel().isOpen()); + + channelRead(rstStreamFrame(STREAM_ID, (int) Http2Error.CANCEL.code())); + + // the peer may still be reading the response and sending WINDOW_UPDATE or PING + fakeClock().forwardNanos(NettyServerHandler.GRACEFUL_SHUTDOWN_DRAIN_NANOS - 1); + assertTrue(channel().isOpen()); + fakeClock().forwardNanos(1); + assertFalse(channel().isOpen()); + } + + @Test + public void maxConnectionAgeGrace_shorterThanDrain_closesAtGrace() throws Exception { + maxConnectionAgeInNanos = TimeUnit.MILLISECONDS.toNanos(10L); + maxConnectionAgeGraceInNanos = TimeUnit.MILLISECONDS.toNanos(100L); + manualSetUp(); + + fakeClock().forwardNanos(maxConnectionAgeInNanos); + channelRead(pingFrame(true /* isAck */, NettyServerHandler.GRACEFUL_SHUTDOWN_PING)); + + // second GO_AWAY sent + verifyWrite().writeGoAway( + eq(ctx()), eq(0), eq(Http2Error.NO_ERROR.code()), any(ByteBuf.class), + any(ChannelPromise.class)); + fakeClock().forwardTime(99, TimeUnit.MILLISECONDS); + assertTrue(channel().isOpen()); + fakeClock().forwardTime(1, TimeUnit.MILLISECONDS); assertFalse(channel().isOpen()); } From 26f66c762c39d741e660d45469b92029fedb3268 Mon Sep 17 00:00:00 2001 From: Meemaw Date: Thu, 1 Oct 2026 14:48:23 +0200 Subject: [PATCH 2/2] netty: add end-to-end regression test for max-age close racing a response Runs a real server with maxConnectionAge on loopback. A unary call outlives the max age, so both GOAWAYs are sent while it is in flight, and the client's event loop stalls while the 256 KiB response is sent. Without the drain, the server closes the socket with the response still queued, the client's next frame triggers an RST, and the call fails with "UNAVAILABLE: Connection closed after GOAWAY ... max_age" (3/3 runs on Linux; macOS loopback does not lose the data). With the drain it passes. Co-Authored-By: Claude Opus 5.5 --- .../netty/NettyServerGracefulCloseTest.java | 138 ++++++++++++++++++ 1 file changed, 138 insertions(+) create mode 100644 netty/src/test/java/io/grpc/netty/NettyServerGracefulCloseTest.java diff --git a/netty/src/test/java/io/grpc/netty/NettyServerGracefulCloseTest.java b/netty/src/test/java/io/grpc/netty/NettyServerGracefulCloseTest.java new file mode 100644 index 00000000000..85a538f1e7a --- /dev/null +++ b/netty/src/test/java/io/grpc/netty/NettyServerGracefulCloseTest.java @@ -0,0 +1,138 @@ +/* + * Copyright 2026 The gRPC Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.grpc.netty; + +import static com.google.common.truth.Truth.assertThat; + +import com.google.common.io.ByteStreams; +import io.grpc.CallOptions; +import io.grpc.ManagedChannel; +import io.grpc.MethodDescriptor; +import io.grpc.Server; +import io.grpc.ServerServiceDefinition; +import io.grpc.stub.ClientCalls; +import io.grpc.stub.ServerCalls; +import io.grpc.testing.GrpcCleanupRule; +import io.netty.channel.EventLoopGroup; +import io.netty.channel.MultiThreadIoEventLoopGroup; +import io.netty.channel.nio.NioIoHandler; +import io.netty.channel.socket.nio.NioSocketChannel; +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.Rule; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** + * Regression test for https://github.com/grpc/grpc-java/issues/9566: a unary RPC that the server + * completed must not fail because the server closed the connection after a max-age GOAWAY. + * + *

The client's event loop is stalled while the server sends the response, so the response is + * still unread in the server's kernel send buffer when the server would close the connection. If + * the server closes it then, the client's next frame (a BDP PING or WINDOW_UPDATE) makes the + * server's kernel answer with RST, which discards the unsent trailers. This depends on the kernel: + * without the fix it fails on Linux, while macOS loopback does not lose the data. + */ +@RunWith(JUnit4.class) +public class NettyServerGracefulCloseTest { + private static final int RESPONSE_BYTES = 256 * 1024; + private static final int ITERATIONS = 3; + private static final long CLIENT_STALL_MILLIS = 300; + + private static final MethodDescriptor.Marshaller BYTES = + new MethodDescriptor.Marshaller() { + @Override + public InputStream stream(byte[] value) { + return new ByteArrayInputStream(value); + } + + @Override + public byte[] parse(InputStream stream) { + try { + return ByteStreams.toByteArray(stream); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + }; + + private static final MethodDescriptor METHOD = + MethodDescriptor.newBuilder() + .setType(MethodDescriptor.MethodType.UNARY) + .setFullMethodName("test.Test/Get") + .setRequestMarshaller(BYTES) + .setResponseMarshaller(BYTES) + .build(); + + @Rule public final GrpcCleanupRule grpcCleanup = new GrpcCleanupRule(); + + private final AtomicReference clientGroup = new AtomicReference<>(); + + @Test + public void maxConnectionAge_unaryCallInFlight_succeeds() throws Exception { + ServerServiceDefinition service = ServerServiceDefinition.builder("test.Test") + .addMethod(METHOD, ServerCalls.asyncUnaryCall((request, responseObserver) -> { + // Outlive the max age (1 second plus up to 10% jitter), so both GOAWAYs are sent while + // the call is in flight. + sleep(1500); + clientGroup.get().execute(() -> sleep(CLIENT_STALL_MILLIS)); + responseObserver.onNext(new byte[RESPONSE_BYTES]); + responseObserver.onCompleted(); + })) + .build(); + Server server = grpcCleanup.register( + NettyServerBuilder.forAddress(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0)) + .maxConnectionAge(1, TimeUnit.SECONDS) + .maxConnectionAgeGrace(30, TimeUnit.SECONDS) + .addService(service) + .build() + .start()); + + for (int i = 0; i < ITERATIONS; i++) { + EventLoopGroup group = new MultiThreadIoEventLoopGroup(1, NioIoHandler.newFactory()); + clientGroup.set(group); + ManagedChannel channel = NettyChannelBuilder + .forAddress(new InetSocketAddress(InetAddress.getLoopbackAddress(), server.getPort())) + .channelType(NioSocketChannel.class) + .eventLoopGroup(group) + .usePlaintext() + .build(); + try { + byte[] response = + ClientCalls.blockingUnaryCall(channel, METHOD, CallOptions.DEFAULT, new byte[16]); + assertThat(response).hasLength(RESPONSE_BYTES); + } finally { + channel.shutdownNow().awaitTermination(5, TimeUnit.SECONDS); + group.shutdownGracefully(0, 1, TimeUnit.SECONDS).sync(); + } + } + } + + private static void sleep(long millis) { + try { + Thread.sleep(millis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +}