From cfdc57c933277b22e86e1ad3c6d2844788b2da30 Mon Sep 17 00:00:00 2001 From: waynercheung Date: Tue, 8 Sep 2026 20:16:26 +0800 Subject: [PATCH] fix(grpc): apply RST_STREAM frame limit by default The node.rpc.maxRstStream and secondsPerWindow options defaulted to 0, which left the HTTP/2 RST_STREAM frame limit unset, and the limit was configured only when both values were positive. Give them sensible defaults, validate them at load time, and apply the limit unconditionally so the configuration behaves consistently. - NodeConfig: default maxRstStream/secondsPerWindow to 1000/5. In postProcess, reject negative values for either option and reject Integer.MAX_VALUE specifically for maxRstStream, which grpc-java treats as disabling the limit. Independently replace each zero with its default and log a warning. - RpcService: read and validate the two values before allocating the executor, then always configure maxRstFramesPerWindow. - reference.conf, docs/configuration.md and config/README.md: document the new defaults and the 0 fallback for compatibility. - Tests: pin the Java and reference defaults, preserve positive overrides, and verify PING at the RST limit followed by GOAWAY after one additional reset using a prebuilt HTTP/2 frame burst. Note: an explicit 0 now falls back to the default instead of leaving the limit unset, and maxRstStream = 2147483647 is rejected at startup. --- .../main/java/org/tron/core/config/README.md | 15 ++ .../org/tron/core/config/args/NodeConfig.java | 29 ++- common/src/main/resources/reference.conf | 10 +- .../tron/core/config/args/NodeConfigTest.java | 60 +++++++ docs/configuration.md | 13 ++ .../tron/common/application/RpcService.java | 17 +- .../RpcServiceHttp2SecurityTest.java | 169 +++++++++++++++++- 7 files changed, 298 insertions(+), 15 deletions(-) diff --git a/common/src/main/java/org/tron/core/config/README.md b/common/src/main/java/org/tron/core/config/README.md index 6bcd2ade1aa..33057b4322e 100644 --- a/common/src/main/java/org/tron/core/config/README.md +++ b/common/src/main/java/org/tron/core/config/README.md @@ -58,6 +58,12 @@ node { # default 100. Setting 0 also uses the secure default. # maxConcurrentCallsPerConnection = 100 + # Maximum RST_STREAM frames per connection per window; 0 uses the default of 1000. + # maxRstStream = 1000 + + # RST_STREAM counting window in seconds; 0 uses the default of 5. + # secondsPerWindow = 5 + # The HTTP/2 flow control window, default 1MB # flowControlWindow = @@ -80,6 +86,15 @@ node { > It now selects the secure default of 100. Configure an explicit positive value if a node > requires more than 100 concurrent calls on one connection. +> **Upgrade note:** `maxRstStream = 0` and `secondsPerWindow = 0` no longer disable +> RST_STREAM flood protection. Each zero independently falls back to its secure default +> (1000 frames / 5 seconds), with a startup warning. Negative values and +> `maxRstStream = 2147483647` (`Integer.MAX_VALUE`, grpc-java's disable sentinel) are rejected. +> These are java-tron defaults; grpc-java itself defaults to no limit. +> Exceeding the limit closes that connection with `GOAWAY(ENHANCE_YOUR_CALM)`. +> Clients that frequently cancel calls, including deadline cancellations, may need explicit +> positive limits tuned to their workload; keep `maxRstStream` below `2147483647`. + ## backup You can customize backup options in the `node.backup` part of `config.conf`, which looks like: ``` diff --git a/common/src/main/java/org/tron/core/config/args/NodeConfig.java b/common/src/main/java/org/tron/core/config/args/NodeConfig.java index 91945b5a73b..8c1bf69536d 100644 --- a/common/src/main/java/org/tron/core/config/args/NodeConfig.java +++ b/common/src/main/java/org/tron/core/config/args/NodeConfig.java @@ -208,6 +208,8 @@ public static class HttpConfig { public static class RpcConfig { public static final int DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION = 100; + public static final int DEFAULT_MAX_RST_STREAM = 1000; + public static final int DEFAULT_SECONDS_PER_WINDOW = 5; private boolean enable = true; private int port = 50051; @@ -224,8 +226,8 @@ public static class RpcConfig { private long maxConnectionAgeInMillis = 0; private int maxMessageSize = 4194304; private int maxHeaderListSize = 8192; - private int maxRstStream = 0; - private int secondsPerWindow = 0; + private int maxRstStream = DEFAULT_MAX_RST_STREAM; + private int secondsPerWindow = DEFAULT_SECONDS_PER_WINDOW; private int minEffectiveConnection = 1; private boolean reflectionService = false; private boolean trxCacheEnable = false; @@ -373,6 +375,29 @@ private void postProcess() { rpc.maxConcurrentCallsPerConnection = RpcConfig.DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION; } + if (rpc.maxRstStream < 0) { + throw new TronError("node.rpc.maxRstStream must be non-negative, got: " + + rpc.maxRstStream, PARAMETER_INIT); + } + if (rpc.secondsPerWindow < 0) { + throw new TronError("node.rpc.secondsPerWindow must be non-negative, got: " + + rpc.secondsPerWindow, PARAMETER_INIT); + } + // Only the frame count has a grpc-java disable sentinel; the window does not. + if (rpc.maxRstStream == Integer.MAX_VALUE) { + throw new TronError("node.rpc.maxRstStream must not be Integer.MAX_VALUE because grpc-java " + + "treats it as disabling RST_STREAM flood protection", PARAMETER_INIT); + } + if (rpc.maxRstStream == 0) { + logger.warn("Configuring [node.rpc.maxRstStream] as 0 no longer disables RST_STREAM flood " + + "protection; using the secure default of {}.", RpcConfig.DEFAULT_MAX_RST_STREAM); + rpc.maxRstStream = RpcConfig.DEFAULT_MAX_RST_STREAM; + } + if (rpc.secondsPerWindow == 0) { + logger.warn("Configuring [node.rpc.secondsPerWindow] as 0 no longer disables RST_STREAM flood " + + "protection; using the secure default of {}.", RpcConfig.DEFAULT_SECONDS_PER_WINDOW); + rpc.secondsPerWindow = RpcConfig.DEFAULT_SECONDS_PER_WINDOW; + } if (rpc.maxConnectionIdleInMillis == 0) { rpc.maxConnectionIdleInMillis = Long.MAX_VALUE; } diff --git a/common/src/main/resources/reference.conf b/common/src/main/resources/reference.conf index d8c483d932a..bc122444e61 100644 --- a/common/src/main/resources/reference.conf +++ b/common/src/main/resources/reference.conf @@ -308,11 +308,13 @@ node { # Maximum header list size (bytes), default 8192 maxHeaderListSize = 8192 - # RST_STREAM frames allowed per connection per period, 0 = no limit - maxRstStream = 0 + # RST_STREAM frames allowed per connection per period. Integer.MAX_VALUE is rejected. + # 0 is accepted for backward compatibility and falls back to the secure default of 1000. + maxRstStream = 1000 - # Seconds per period for gRPC RST_STREAM limit - secondsPerWindow = 0 + # Seconds per period for the gRPC RST_STREAM limit. + # 0 is accepted for backward compatibility and falls back to the secure default of 5. + secondsPerWindow = 5 # Minimum effective connections required to broadcast transactions minEffectiveConnection = 1 diff --git a/common/src/test/java/org/tron/core/config/args/NodeConfigTest.java b/common/src/test/java/org/tron/core/config/args/NodeConfigTest.java index bcb8b09dd7a..dac9dfdbdd6 100644 --- a/common/src/test/java/org/tron/core/config/args/NodeConfigTest.java +++ b/common/src/test/java/org/tron/core/config/args/NodeConfigTest.java @@ -4,6 +4,7 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; +import static org.tron.core.exception.TronError.ErrCode.PARAMETER_INIT; import com.typesafe.config.Config; import com.typesafe.config.ConfigFactory; @@ -111,6 +112,8 @@ public void testRpcDefaultsFromReference() { assertEquals(9223372036854775807L, rpc.getMaxConnectionAgeInMillis()); assertEquals(4194304, rpc.getMaxMessageSize()); assertEquals(8192, rpc.getMaxHeaderListSize()); + assertEquals(NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM, rpc.getMaxRstStream()); + assertEquals(NodeConfig.RpcConfig.DEFAULT_SECONDS_PER_WINDOW, rpc.getSecondsPerWindow()); assertEquals(1, rpc.getMinEffectiveConnection()); // thread=0 in reference.conf triggers auto-detect in postProcess assertTrue(rpc.getThread() > 0); @@ -146,6 +149,63 @@ public void testRpcNegativeConcurrentCallsRejected() { "node.rpc.maxConcurrentCallsPerConnection must be non-negative, got: -1")); } + @Test + public void testRpcRstDefaultsMatchReference() { + NodeConfig.RpcConfig rpc = new NodeConfig.RpcConfig(); + assertEquals(1000, rpc.getMaxRstStream()); + assertEquals(5, rpc.getSecondsPerWindow()); + + Config reference = withRef(); + assertEquals(rpc.getMaxRstStream(), reference.getInt("node.rpc.maxRstStream")); + assertEquals(rpc.getSecondsPerWindow(), reference.getInt("node.rpc.secondsPerWindow")); + } + + @Test + public void testRpcZeroRstLimitsUseSecureDefaultsIndependently() { + int[][] cases = { + {0, 0, NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM, + NodeConfig.RpcConfig.DEFAULT_SECONDS_PER_WINDOW}, + {0, 10, NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM, 10}, + {5, 0, 5, NodeConfig.RpcConfig.DEFAULT_SECONDS_PER_WINDOW} + }; + for (int[] values : cases) { + NodeConfig.RpcConfig rpc = NodeConfig.fromConfig(withRef( + "node.rpc { maxRstStream = " + values[0] + + ", secondsPerWindow = " + values[1] + " }")).getRpc(); + assertEquals(values[2], rpc.getMaxRstStream()); + assertEquals(values[3], rpc.getSecondsPerWindow()); + } + } + + @Test + public void testRpcInvalidRstLimitsRejectedBeforeFallback() { + int[][] cases = { + {-1, 5}, {1000, -1}, {-1, 0}, {0, -1}, + {Integer.MIN_VALUE, 5}, {1000, Integer.MIN_VALUE}, + {Integer.MAX_VALUE, 5}, {Integer.MAX_VALUE, 0} + }; + for (int[] values : cases) { + Config config = withRef("node.rpc { maxRstStream = " + values[0] + + ", secondsPerWindow = " + values[1] + " }"); + TronError exception = assertThrows(TronError.class, () -> NodeConfig.fromConfig(config)); + assertEquals(PARAMETER_INIT, exception.getErrCode()); + assertTrue(exception.getMessage().contains(values[1] < 0 + ? "node.rpc.secondsPerWindow" : "node.rpc.maxRstStream")); + } + } + + @Test + public void testRpcExplicitPositiveRstLimitsPreserved() { + int[][] cases = {{5, 10}, {200, 30}, {1, 1}, {Integer.MAX_VALUE - 1, Integer.MAX_VALUE}}; + for (int[] values : cases) { + NodeConfig.RpcConfig rpc = NodeConfig.fromConfig(withRef( + "node.rpc { maxRstStream = " + values[0] + + ", secondsPerWindow = " + values[1] + " }")).getRpc(); + assertEquals(values[0], rpc.getMaxRstStream()); + assertEquals(values[1], rpc.getSecondsPerWindow()); + } + } + @Test public void testRpcUserOverrideExplicitValues() { Config config = withRef( diff --git a/docs/configuration.md b/docs/configuration.md index d021326a15e..736111452d9 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -106,6 +106,10 @@ node { solidityPort = 50061 # Maximum concurrent calls per connection. 0 uses the secure default of 100. maxConcurrentCallsPerConnection = 100 + # Maximum RST_STREAM frames per connection per window. 0 uses the default of 1000. + maxRstStream = 1000 + # RST_STREAM counting window in seconds. 0 uses the default of 5. + secondsPerWindow = 5 # Idle connection timeout (ms). 0 = no limit. maxConnectionIdleInMillis = 0 # Minimum active connections required before broadcasting transactions. @@ -118,6 +122,15 @@ node { > It now selects the secure default of 100. Configure an explicit positive value if a client > needs more than 100 concurrent calls on one connection. +> **Upgrade note:** `node.rpc.maxRstStream = 0` and `node.rpc.secondsPerWindow = 0` +> no longer disable RST_STREAM flood protection. Each zero independently falls back to its +> secure default (1000 frames / 5 seconds), with a startup warning. Negative values and +> `maxRstStream = 2147483647` (`Integer.MAX_VALUE`, grpc-java's disable sentinel) are rejected. +> These are java-tron defaults; grpc-java itself defaults to no limit. +> Exceeding the limit closes that connection with `GOAWAY(ENHANCE_YOUR_CALM)`. +> Clients that frequently cancel calls, including deadline cancellations, may need explicit +> positive limits tuned to their workload; keep `maxRstStream` below `2147483647`. + To disable an API endpoint that you do not want to expose publicly, set its `Enable` flag to `false` or add endpoints to `node.disabledApi`: ```hocon diff --git a/framework/src/main/java/org/tron/common/application/RpcService.java b/framework/src/main/java/org/tron/common/application/RpcService.java index c398b71ae41..bc494de2035 100644 --- a/framework/src/main/java/org/tron/common/application/RpcService.java +++ b/framework/src/main/java/org/tron/common/application/RpcService.java @@ -15,6 +15,8 @@ package org.tron.common.application; +import static org.tron.core.exception.TronError.ErrCode.API_SERVER_INIT; + import io.grpc.Server; import io.grpc.netty.NettyServerBuilder; import io.grpc.protobuf.services.ProtoReflectionService; @@ -26,6 +28,7 @@ import org.tron.common.es.ExecutorServiceManager; import org.tron.common.parameter.CommonParameter; import org.tron.core.config.args.Args; +import org.tron.core.exception.TronError; import org.tron.core.services.filter.LiteFnQueryGrpcInterceptor; import org.tron.core.services.ratelimiter.PrometheusInterceptor; import org.tron.core.services.ratelimiter.RateLimiterInterceptor; @@ -92,8 +95,15 @@ public CompletableFuture start() { } protected NettyServerBuilder initServerBuilder() { - NettyServerBuilder serverBuilder = NettyServerBuilder.forPort(this.port); CommonParameter parameter = Args.getInstance(); + int maxRstStream = parameter.getRpcMaxRstStream(); + int secondsPerWindow = parameter.getRpcSecondsPerWindow(); + // Validate before allocating the executor, including callers bypassing NodeConfig. + if (maxRstStream <= 0 || secondsPerWindow <= 0 || maxRstStream == Integer.MAX_VALUE) { + throw new TronError("Invalid gRPC RST_STREAM limit config: maxRstStream=" + + maxRstStream + ", secondsPerWindow=" + secondsPerWindow, API_SERVER_INIT); + } + NettyServerBuilder serverBuilder = NettyServerBuilder.forPort(this.port); if (parameter.getRpcThreadNum() > 0) { this.executorService = ExecutorServiceManager.newFixedThreadPool( this.executorName, parameter.getRpcThreadNum()); @@ -107,10 +117,7 @@ protected NettyServerBuilder initServerBuilder() { .maxConnectionAge(parameter.getMaxConnectionAgeInMillis(), TimeUnit.MILLISECONDS) .maxInboundMessageSize(parameter.getMaxMessageSize()) .maxHeaderListSize(parameter.getMaxHeaderListSize()); - if (parameter.getRpcMaxRstStream() > 0 && parameter.getRpcSecondsPerWindow() > 0) { - serverBuilder.maxRstFramesPerWindow( - parameter.getRpcMaxRstStream(), parameter.getRpcSecondsPerWindow()); - } + serverBuilder.maxRstFramesPerWindow(maxRstStream, secondsPerWindow); if (parameter.isRpcReflectionServiceEnable()) { serverBuilder.addService(ProtoReflectionService.newInstance()); diff --git a/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java b/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java index 9ad83cccec1..63aa8a9ea7e 100644 --- a/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java +++ b/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java @@ -15,12 +15,18 @@ package org.tron.common.application; +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.tron.common.math.StrictMathWrapper.min; +import static org.tron.core.exception.TronError.ErrCode.API_SERVER_INIT; import com.google.protobuf.Empty; +import com.typesafe.config.ConfigFactory; import io.grpc.Metadata; import io.grpc.MethodDescriptor; import io.grpc.Server; @@ -48,20 +54,25 @@ import org.junit.After; import org.junit.Before; import org.junit.Test; +import org.springframework.test.util.ReflectionTestUtils; import org.tron.common.parameter.CommonParameter; import org.tron.core.config.args.Args; +import org.tron.core.config.args.NodeConfig; +import org.tron.core.exception.TronError; public class RpcServiceHttp2SecurityTest { private static final int HEADERS_FRAME_TYPE = 0x1; private static final int RST_STREAM_FRAME_TYPE = 0x3; private static final int SETTINGS_FRAME_TYPE = 0x4; + private static final int PING_FRAME_TYPE = 0x6; private static final int GO_AWAY_FRAME_TYPE = 0x7; private static final int SETTINGS_MAX_CONCURRENT_STREAMS = 0x3; private static final byte[] CLIENT_PREFACE = "PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n".getBytes(StandardCharsets.US_ASCII); private static final byte[] EMPTY_SETTINGS_FRAME = new byte[]{0, 0, 0, 4, 0, 0, 0, 0, 0}; + private static final byte[] PING_PAYLOAD = new byte[]{1, 2, 3, 4, 5, 6, 7, 8}; private static final String SERVICE_NAME = "test.HoldService"; private static final String METHOD_NAME = "Hold"; private static final String METHOD_PATH = "/" + SERVICE_NAME + "/" + METHOD_NAME; @@ -99,8 +110,8 @@ public void setUp() { parameter.setMaxConnectionAgeInMillis(Long.MAX_VALUE); parameter.setMaxMessageSize(4 * 1024 * 1024); parameter.setMaxHeaderListSize(8 * 1024); - parameter.setRpcMaxRstStream(0); - parameter.setRpcSecondsPerWindow(0); + parameter.setRpcMaxRstStream(NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM); + parameter.setRpcSecondsPerWindow(NodeConfig.RpcConfig.DEFAULT_SECONDS_PER_WINDOW); parameter.setRpcReflectionServiceEnable(false); } @@ -143,6 +154,153 @@ public void shouldRejectExcessStreamsBeforeClientAcknowledgesSettings() throws E } } + @Test + public void shouldEnforceDefaultRstLimit() throws Exception { + useRstConfig(""); + assertRstLimitEnforced(NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM); + } + + @Test + public void shouldEnforceRstLimitForLegacyZeroConfig() throws Exception { + useRstConfig("node.rpc { maxRstStream = 0, secondsPerWindow = 0 }"); + assertRstLimitEnforced(NodeConfig.RpcConfig.DEFAULT_MAX_RST_STREAM); + } + + @Test + public void shouldEnforceExplicitRstLimit() throws Exception { + parameter.setRpcMaxRstStream(2); + parameter.setRpcSecondsPerWindow(30); + assertRstLimitEnforced(2); + } + + @Test + public void shouldRejectInvalidRstLimitsBeforeAllocatingExecutor() throws Exception { + parameter.setRpcThreadNum(1); + int[][] cases = { + {0, 0}, {0, 5}, {1000, 0}, {-1, 5}, {1000, -1}, + {Integer.MIN_VALUE, 5}, {1000, Integer.MIN_VALUE}, + {Integer.MAX_VALUE, 5} + }; + for (int[] values : cases) { + parameter.setRpcMaxRstStream(values[0]); + parameter.setRpcSecondsPerWindow(values[1]); + TestRpcService rpcService = new TestRpcService(); + try { + TronError exception = assertThrows(TronError.class, rpcService::newServerBuilder); + assertEquals(API_SERVER_INIT, exception.getErrCode()); + assertTrue(exception.getMessage().contains("maxRstStream=" + values[0])); + assertTrue(exception.getMessage().contains("secondsPerWindow=" + values[1])); + assertNull("Invalid configuration must not allocate an executor", + ReflectionTestUtils.getField(rpcService, "executorService")); + } finally { + rpcService.innerStop(); + } + } + } + + @Test + public void shouldAcceptLargestFiniteRstLimitAndWindow() { + parameter.setRpcMaxRstStream(Integer.MAX_VALUE - 1); + parameter.setRpcSecondsPerWindow(Integer.MAX_VALUE); + assertNotNull(new TestRpcService().newServerBuilder()); + } + + private void useRstConfig(String hocon) { + NodeConfig.RpcConfig rpc = NodeConfig.fromConfig(ConfigFactory.parseString(hocon) + .withFallback(ConfigFactory.defaultReference())).getRpc(); + parameter.setRpcMaxRstStream(rpc.getMaxRstStream()); + parameter.setRpcSecondsPerWindow(rpc.getSecondsPerWindow()); + } + + private void assertRstLimitEnforced(int limit) throws Exception { + // Prepare frames before connecting so encoding does not consume the counting window. + ByteArrayOutputStream burst = new ByteArrayOutputStream(); + burst.write(CLIENT_PREFACE); + burst.write(EMPTY_SETTINGS_FRAME); + for (int i = 0; i < limit; i++) { + int streamId = 2 * i + 1; + burst.write(newHeadersFrame(streamId)); + burst.write(newRstStreamFrame(streamId)); + } + burst.write(newPingFrame()); + int excessStreamId = 2 * limit + 1; + burst.write(newHeadersFrame(excessStreamId)); + burst.write(newRstStreamFrame(excessStreamId)); + byte[] frames = burst.toByteArray(); + + parameter.setMaxConcurrentCallsPerConnection(100); + TestRpcService rpcService = new TestRpcService(); + Server server = rpcService.newServerBuilder() + .addService(newHoldService()) + .build() + .start(); + + try (Socket socket = new Socket("127.0.0.1", server.getPort())) { + socket.setSoTimeout(5_000); + OutputStream output = socket.getOutputStream(); + // Send PING and the excess reset together to avoid a round-trip within the short window. + output.write(frames); + output.flush(); + + // Exactly the configured number of resets is allowed on the connection. + assertPingAcknowledged(socket.getInputStream()); + assertRstFloodGoAway(socket.getInputStream()); + } finally { + server.shutdownNow(); + assertTrue(server.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private static byte[] newRstStreamFrame(int streamId) { + return ByteBuffer.allocate(13) + .put(new byte[]{0, 0, 4, RST_STREAM_FRAME_TYPE, 0}) + .putInt(streamId) + .putInt((int) Http2Error.CANCEL.code()) + .array(); + } + + private static byte[] newPingFrame() { + return ByteBuffer.allocate(17) + .put(new byte[]{0, 0, 8, PING_FRAME_TYPE, 0}) + .putInt(0) + .put(PING_PAYLOAD) + .array(); + } + + private static void assertPingAcknowledged(InputStream input) throws IOException { + for (int i = 0; i < 512; i++) { + Http2Frame frame = readFrame(input); + if (frame.type == GO_AWAY_FRAME_TYPE) { + fail("Server closed the connection before the RST_STREAM limit was exceeded"); + } + if (frame.type == PING_FRAME_TYPE && (frame.flags & 0x1) != 0) { + assertEquals(0, frame.streamId); + assertArrayEquals(PING_PAYLOAD, frame.payload); + return; + } + } + fail("No PING acknowledgement at the RST_STREAM limit"); + } + + private static void assertRstFloodGoAway(InputStream input) throws IOException { + for (int i = 0; i < 512; i++) { + Http2Frame frame; + try { + frame = readFrame(input); + } catch (EOFException e) { + break; + } + if (frame.type == GO_AWAY_FRAME_TYPE) { + assertEquals(0, frame.streamId); + assertTrue("Invalid GOAWAY payload", frame.payload.length >= 8); + long errorCode = ByteBuffer.wrap(frame.payload).getInt(4) & 0xffff_ffffL; + assertEquals(Http2Error.ENHANCE_YOUR_CALM.code(), errorCode); + return; + } + } + fail("No ENHANCE_YOUR_CALM GOAWAY after exceeding the RST_STREAM limit"); + } + private static ServerServiceDefinition newHoldService() { MethodDescriptor method = MethodDescriptor.newBuilder() @@ -236,8 +394,9 @@ private static Http2Frame readFrame(InputStream input) throws IOException { int payloadLength = ((header[0] & 0xff) << 16) | ((header[1] & 0xff) << 8) | (header[2] & 0xff); int type = header[3] & 0xff; + int flags = header[4] & 0xff; int streamId = ByteBuffer.wrap(header, 5, 4).getInt() & 0x7fff_ffff; - return new Http2Frame(type, streamId, readFully(input, payloadLength)); + return new Http2Frame(type, flags, streamId, readFully(input, payloadLength)); } private static byte[] readFully(InputStream input, int length) throws IOException { @@ -256,11 +415,13 @@ private static byte[] readFully(InputStream input, int length) throws IOExceptio private static final class Http2Frame { private final int type; + private final int flags; private final int streamId; private final byte[] payload; - private Http2Frame(int type, int streamId, byte[] payload) { + private Http2Frame(int type, int flags, int streamId, byte[] payload) { this.type = type; + this.flags = flags; this.streamId = streamId; this.payload = payload; }