diff --git a/build.gradle b/build.gradle index e143ab3a947..65e72c0fb73 100644 --- a/build.gradle +++ b/build.gradle @@ -5,6 +5,10 @@ plugins { id 'net.ltgt.errorprone' version '5.0.0' apply false } +ext { + grpcVersion = "1.83.0" +} + allprojects { version = "1.0.0" apply plugin: "java-library" @@ -41,7 +45,7 @@ ext.archInfo = [ // https://github.com/grpc/grpc-java/issues/7690 // https://github.com/grpc/grpc-java/pull/12319, Add support for macOS aarch64 with universal binary // https://github.com/grpc/grpc-java/pull/11371 , 1.64.x is not supported CentOS 7. - ProtocGenVersion: isArm64 || isMac ? '1.81.0' : '1.60.0' + ProtocGenVersion: isArm64 || isMac ? rootProject.grpcVersion : '1.60.0' ], VMOptions: isArm64 ? "${rootDir}/gradle/jdk17/java-tron.vmoptions" : "${rootDir}/gradle/java-tron.vmoptions" ] 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 1380c98984e..5618cb9fed9 100644 --- a/common/src/main/java/org/tron/core/config/README.md +++ b/common/src/main/java/org/tron/core/config/README.md @@ -54,8 +54,9 @@ node { # Number of gRPC thread, default availableProcessors / 2 # thread = 16 - # The maximum number of concurrent calls permitted for each incoming connection - # maxConcurrentCallsPerConnection = + # The maximum number of concurrent calls permitted for each incoming connection, + # default 100. Setting 0 also uses the secure default. + # maxConcurrentCallsPerConnection = 100 # The HTTP/2 flow control window, default 1MB # flowControlWindow = @@ -75,6 +76,10 @@ node { } ``` +> **Upgrade note:** `maxConcurrentCallsPerConnection = 0` previously disabled the limit. +> 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. + ## 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 2158f56d0ba..91945b5a73b 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 @@ -207,6 +207,8 @@ public static class HttpConfig { @Setter public static class RpcConfig { + public static final int DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION = 100; + private boolean enable = true; private int port = 50051; private boolean solidityEnable = true; @@ -215,7 +217,8 @@ public static class RpcConfig { private int pBFTPort = 50071; private int thread = 0; - private int maxConcurrentCallsPerConnection = 0; + private int maxConcurrentCallsPerConnection = + DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION; private int flowControlWindow = 1048576; private long maxConnectionIdleInMillis = 0; private long maxConnectionAgeInMillis = 0; @@ -358,8 +361,17 @@ private void postProcess() { rpc.thread = (Runtime.getRuntime().availableProcessors() + 1) / 2; } + if (rpc.maxConcurrentCallsPerConnection < 0) { + throw new TronError("node.rpc.maxConcurrentCallsPerConnection must be non-negative, got: " + + rpc.maxConcurrentCallsPerConnection, PARAMETER_INIT); + } if (rpc.maxConcurrentCallsPerConnection == 0) { - rpc.maxConcurrentCallsPerConnection = Integer.MAX_VALUE; + logger.warn("Configuring [node.rpc.maxConcurrentCallsPerConnection] as 0 no longer " + + "disables the limit; using the secure default of {}. Configure an explicit positive " + + "value if more concurrency is required.", + RpcConfig.DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION); + rpc.maxConcurrentCallsPerConnection = + RpcConfig.DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION; } 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 25fc4832e55..d8c483d932a 100644 --- a/common/src/main/resources/reference.conf +++ b/common/src/main/resources/reference.conf @@ -289,9 +289,8 @@ node { # Number of gRPC threads, 0 = auto (availableProcessors / 2) thread = 0 - # Maximum concurrent calls per incoming connection - # 0 means No limit on concurrent calls per connection - maxConcurrentCallsPerConnection = 0 + # Maximum concurrent calls per incoming connection. 0 falls back to the secure default of 100. + maxConcurrentCallsPerConnection = 100 # HTTP/2 flow control window (bytes), default 1MB flowControlWindow = 1048576 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 bbc2d2475ee..bcb8b09dd7a 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 @@ -2,11 +2,13 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import com.typesafe.config.Config; import com.typesafe.config.ConfigFactory; import org.junit.Test; +import org.tron.core.exception.TronError; public class NodeConfigTest { @@ -102,7 +104,8 @@ public void testRpcDefaultsFromReference() { NodeConfig.RpcConfig rpc = nc.getRpc(); // reference.conf provides actual final defaults, no sentinel conversion needed - assertEquals(2147483647, rpc.getMaxConcurrentCallsPerConnection()); + assertEquals(NodeConfig.RpcConfig.DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION, + rpc.getMaxConcurrentCallsPerConnection()); assertEquals(1048576, rpc.getFlowControlWindow()); assertEquals(9223372036854775807L, rpc.getMaxConnectionIdleInMillis()); assertEquals(9223372036854775807L, rpc.getMaxConnectionAgeInMillis()); @@ -122,6 +125,27 @@ public void testRpcUserOverrideZeroNotConverted() { assertEquals(0, nc.getRpc().getMinEffectiveConnection()); } + @Test + public void testRpcZeroConcurrentCallsUsesSecureDefault() { + Config config = withRef( + "node { rpc { maxConcurrentCallsPerConnection = 0 } }"); + NodeConfig nc = NodeConfig.fromConfig(config); + assertEquals(NodeConfig.RpcConfig.DEFAULT_MAX_CONCURRENT_CALLS_PER_CONNECTION, + nc.getRpc().getMaxConcurrentCallsPerConnection()); + } + + @Test + public void testRpcNegativeConcurrentCallsRejected() { + Config config = withRef( + "node { rpc { maxConcurrentCallsPerConnection = -1 } }"); + + TronError exception = assertThrows(TronError.class, + () -> NodeConfig.fromConfig(config)); + + assertTrue(exception.getMessage().contains( + "node.rpc.maxConcurrentCallsPerConnection must be non-negative, got: -1")); + } + @Test public void testRpcUserOverrideExplicitValues() { Config config = withRef( diff --git a/docs/configuration.md b/docs/configuration.md index 28b53b1970c..d021326a15e 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -104,8 +104,8 @@ node { port = 50051 solidityEnable = true solidityPort = 50061 - # Maximum concurrent calls per connection. 0 = no limit. - maxConcurrentCallsPerConnection = 0 + # Maximum concurrent calls per connection. 0 uses the secure default of 100. + maxConcurrentCallsPerConnection = 100 # Idle connection timeout (ms). 0 = no limit. maxConnectionIdleInMillis = 0 # Minimum active connections required before broadcasting transactions. @@ -114,6 +114,10 @@ node { } ``` +> **Upgrade note:** `node.rpc.maxConcurrentCallsPerConnection = 0` previously meant no limit. +> 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. + 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/build.gradle b/framework/build.gradle index 0ce33f253cf..8255fc30d18 100644 --- a/framework/build.gradle +++ b/framework/build.gradle @@ -40,6 +40,10 @@ dependencies { // end local libraries implementation group: 'com.beust', name: 'jcommander', version: '1.78' implementation group: 'io.dropwizard.metrics', name: 'metrics-core', version: '3.1.2' + implementation('io.netty:netty-codec-protobuf:4.2.15.Final') { + exclude group: 'com.google.protobuf' + exclude group: 'com.google.protobuf.nano' + } // http implementation 'org.eclipse.jetty:jetty-server:9.4.58.v20250814' implementation 'org.eclipse.jetty:jetty-servlet:9.4.58.v20250814' diff --git a/framework/src/main/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiter.java b/framework/src/main/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiter.java new file mode 100644 index 00000000000..cdd71ffee3c --- /dev/null +++ b/framework/src/main/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiter.java @@ -0,0 +1,79 @@ +/* + * java-tron is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * java-tron is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with java-tron. If not, see . + */ + +package org.tron.common.application; + +import static com.google.common.base.Preconditions.checkArgument; +import static com.google.common.base.Preconditions.checkNotNull; + +import io.grpc.netty.GrpcHttp2ConnectionHandler; +import io.grpc.netty.InternalProtocolNegotiator; +import io.grpc.netty.InternalProtocolNegotiators; +import io.grpc.netty.NettyServerBuilder; +import io.netty.channel.ChannelHandler; +import io.netty.util.AsciiString; + +/** Enforces the advertised HTTP/2 concurrent stream limit for grpc-netty servers. */ +final class GrpcNettyMaxConcurrentStreamsLimiter { + + private GrpcNettyMaxConcurrentStreamsLimiter() { + } + + static NettyServerBuilder configurePlaintext( + NettyServerBuilder builder, int maxConcurrentStreams) { + checkNotNull(builder, "builder"); + checkArgument(maxConcurrentStreams > 0, "maxConcurrentStreams must be positive"); + builder.maxConcurrentCallsPerConnection(maxConcurrentStreams); + // TODO: Remove this shim after https://github.com/grpc/grpc-java/issues/12930 is fixed. + return builder.protocolNegotiator(newPlaintextNegotiator(maxConcurrentStreams)); + } + + static InternalProtocolNegotiator.ProtocolNegotiator newPlaintextNegotiator( + int maxConcurrentStreams) { + checkArgument(maxConcurrentStreams > 0, "maxConcurrentStreams must be positive"); + return new EnforcingProtocolNegotiator( + InternalProtocolNegotiators.serverPlaintext(), maxConcurrentStreams); + } + + private static final class EnforcingProtocolNegotiator + implements InternalProtocolNegotiator.ProtocolNegotiator { + + private final InternalProtocolNegotiator.ProtocolNegotiator delegate; + private final int maxConcurrentStreams; + + private EnforcingProtocolNegotiator( + InternalProtocolNegotiator.ProtocolNegotiator delegate, int maxConcurrentStreams) { + this.delegate = checkNotNull(delegate, "delegate"); + this.maxConcurrentStreams = maxConcurrentStreams; + } + + @Override + public AsciiString scheme() { + return delegate.scheme(); + } + + @Override + public ChannelHandler newHandler(GrpcHttp2ConnectionHandler grpcHandler) { + // grpc-java builds the connection directly, bypassing Netty's builder-side enforcement. + grpcHandler.connection().remote().maxActiveStreams(maxConcurrentStreams); + return delegate.newHandler(grpcHandler); + } + + @Override + public void close() { + delegate.close(); + } + } +} 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..27fcc479f4e 100644 --- a/framework/src/main/java/org/tron/common/application/RpcService.java +++ b/framework/src/main/java/org/tron/common/application/RpcService.java @@ -100,8 +100,9 @@ protected NettyServerBuilder initServerBuilder() { serverBuilder = serverBuilder.executor(this.executorService); } // Set configs from config.conf or default value + serverBuilder = GrpcNettyMaxConcurrentStreamsLimiter.configurePlaintext( + serverBuilder, parameter.getMaxConcurrentCallsPerConnection()); serverBuilder - .maxConcurrentCallsPerConnection(parameter.getMaxConcurrentCallsPerConnection()) .flowControlWindow(parameter.getFlowControlWindow()) .maxConnectionIdle(parameter.getMaxConnectionIdleInMillis(), TimeUnit.MILLISECONDS) .maxConnectionAge(parameter.getMaxConnectionAgeInMillis(), TimeUnit.MILLISECONDS) diff --git a/framework/src/test/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiterTest.java b/framework/src/test/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiterTest.java new file mode 100644 index 00000000000..fc578ca7947 --- /dev/null +++ b/framework/src/test/java/org/tron/common/application/GrpcNettyMaxConcurrentStreamsLimiterTest.java @@ -0,0 +1,108 @@ +/* + * java-tron is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * java-tron is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with java-tron. If not, see . + */ + +package org.tron.common.application; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThrows; + +import io.grpc.ChannelLogger; +import io.grpc.ChannelLogger.ChannelLogLevel; +import io.grpc.netty.GrpcHttp2ConnectionHandler; +import io.grpc.netty.InternalProtocolNegotiator; +import io.netty.channel.ChannelHandler; +import io.netty.handler.codec.http2.DefaultHttp2Connection; +import io.netty.handler.codec.http2.DefaultHttp2ConnectionDecoder; +import io.netty.handler.codec.http2.DefaultHttp2ConnectionEncoder; +import io.netty.handler.codec.http2.DefaultHttp2FrameReader; +import io.netty.handler.codec.http2.DefaultHttp2FrameWriter; +import io.netty.handler.codec.http2.Http2Connection; +import io.netty.handler.codec.http2.Http2ConnectionDecoder; +import io.netty.handler.codec.http2.Http2ConnectionEncoder; +import io.netty.handler.codec.http2.Http2Error; +import io.netty.handler.codec.http2.Http2Exception; +import io.netty.handler.codec.http2.Http2FrameWriter; +import io.netty.handler.codec.http2.Http2Settings; +import org.junit.Test; + +public class GrpcNettyMaxConcurrentStreamsLimiterTest { + + private static final ChannelLogger NOOP_LOGGER = new ChannelLogger() { + @Override + public void log(ChannelLogLevel level, String message) { + } + + @Override + public void log(ChannelLogLevel level, String messageFormat, Object... args) { + } + }; + + @Test + public void shouldEnforceMaxStreamsBeforeSettingsAck() throws Exception { + Http2Connection connection = new DefaultHttp2Connection(true); + GrpcHttp2ConnectionHandler grpcHandler = newGrpcHandler(connection); + InternalProtocolNegotiator.ProtocolNegotiator negotiator = + GrpcNettyMaxConcurrentStreamsLimiter.newPlaintextNegotiator(2); + + ChannelHandler negotiationHandler = negotiator.newHandler(grpcHandler); + + assertNotNull(negotiationHandler); + assertEquals(2, connection.remote().maxActiveStreams()); + connection.remote().createStream(1, true); + connection.remote().createStream(3, true); + Http2Exception exception = assertThrows( + Http2Exception.class, () -> connection.remote().createStream(5, true)); + assertEquals(Http2Error.REFUSED_STREAM, exception.error()); + negotiator.close(); + } + + @Test + public void shouldIgnoreClientMaxHeaderListSizeOnServer() throws Exception { + Http2Connection connection = new DefaultHttp2Connection(true); + Http2FrameWriter frameWriter = new DefaultHttp2FrameWriter(); + Http2ConnectionEncoder encoder = + new DefaultHttp2ConnectionEncoder(connection, frameWriter); + long originalMaxHeaderListSize = + encoder.configuration().headersConfiguration().maxHeaderListSize(); + + encoder.remoteSettings(new Http2Settings().maxHeaderListSize(1)); + + assertEquals(originalMaxHeaderListSize, + encoder.configuration().headersConfiguration().maxHeaderListSize()); + encoder.close(); + } + + @Test + public void shouldRejectNonPositiveStreamLimit() { + IllegalArgumentException zeroLimitException = assertThrows(IllegalArgumentException.class, + () -> GrpcNettyMaxConcurrentStreamsLimiter.newPlaintextNegotiator(0)); + assertEquals("maxConcurrentStreams must be positive", zeroLimitException.getMessage()); + IllegalArgumentException negativeLimitException = assertThrows(IllegalArgumentException.class, + () -> GrpcNettyMaxConcurrentStreamsLimiter.newPlaintextNegotiator(-1)); + assertEquals("maxConcurrentStreams must be positive", negativeLimitException.getMessage()); + } + + private static GrpcHttp2ConnectionHandler newGrpcHandler(Http2Connection connection) { + Http2FrameWriter frameWriter = new DefaultHttp2FrameWriter(); + Http2ConnectionEncoder encoder = + new DefaultHttp2ConnectionEncoder(connection, frameWriter); + Http2ConnectionDecoder decoder = new DefaultHttp2ConnectionDecoder( + connection, encoder, new DefaultHttp2FrameReader()); + return new GrpcHttp2ConnectionHandler( + null, decoder, encoder, new Http2Settings(), NOOP_LOGGER) { + }; + } +} diff --git a/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java b/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java new file mode 100644 index 00000000000..9ad83cccec1 --- /dev/null +++ b/framework/src/test/java/org/tron/common/application/RpcServiceHttp2SecurityTest.java @@ -0,0 +1,283 @@ +/* + * java-tron is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * java-tron is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with java-tron. If not, see . + */ + +package org.tron.common.application; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; +import static org.tron.common.math.StrictMathWrapper.min; + +import com.google.protobuf.Empty; +import io.grpc.Metadata; +import io.grpc.MethodDescriptor; +import io.grpc.Server; +import io.grpc.ServerCall; +import io.grpc.ServerCallHandler; +import io.grpc.ServerServiceDefinition; +import io.grpc.netty.NettyServerBuilder; +import io.grpc.protobuf.ProtoUtils; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufUtil; +import io.netty.buffer.Unpooled; +import io.netty.handler.codec.http2.DefaultHttp2Headers; +import io.netty.handler.codec.http2.DefaultHttp2HeadersEncoder; +import io.netty.handler.codec.http2.Http2Error; +import io.netty.handler.codec.http2.Http2Headers; +import java.io.ByteArrayOutputStream; +import java.io.EOFException; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.net.Socket; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.util.concurrent.TimeUnit; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.tron.common.parameter.CommonParameter; +import org.tron.core.config.args.Args; + +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 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 String SERVICE_NAME = "test.HoldService"; + private static final String METHOD_NAME = "Hold"; + private static final String METHOD_PATH = "/" + SERVICE_NAME + "/" + METHOD_NAME; + + private CommonParameter parameter; + private int previousRpcThreadNum; + private int previousMaxConcurrentCalls; + private int previousFlowControlWindow; + private long previousMaxConnectionIdle; + private long previousMaxConnectionAge; + private int previousMaxMessageSize; + private int previousMaxHeaderListSize; + private int previousMaxRstStream; + private int previousSecondsPerWindow; + private boolean previousReflectionServiceEnable; + + @Before + public void setUp() { + parameter = Args.getInstance(); + previousRpcThreadNum = parameter.getRpcThreadNum(); + previousMaxConcurrentCalls = parameter.getMaxConcurrentCallsPerConnection(); + previousFlowControlWindow = parameter.getFlowControlWindow(); + previousMaxConnectionIdle = parameter.getMaxConnectionIdleInMillis(); + previousMaxConnectionAge = parameter.getMaxConnectionAgeInMillis(); + previousMaxMessageSize = parameter.getMaxMessageSize(); + previousMaxHeaderListSize = parameter.getMaxHeaderListSize(); + previousMaxRstStream = parameter.getRpcMaxRstStream(); + previousSecondsPerWindow = parameter.getRpcSecondsPerWindow(); + previousReflectionServiceEnable = parameter.isRpcReflectionServiceEnable(); + + parameter.setRpcThreadNum(0); + parameter.setMaxConcurrentCallsPerConnection(2); + parameter.setFlowControlWindow(NettyServerBuilder.DEFAULT_FLOW_CONTROL_WINDOW); + parameter.setMaxConnectionIdleInMillis(60_000); + parameter.setMaxConnectionAgeInMillis(Long.MAX_VALUE); + parameter.setMaxMessageSize(4 * 1024 * 1024); + parameter.setMaxHeaderListSize(8 * 1024); + parameter.setRpcMaxRstStream(0); + parameter.setRpcSecondsPerWindow(0); + parameter.setRpcReflectionServiceEnable(false); + } + + @After + public void tearDown() { + parameter.setRpcThreadNum(previousRpcThreadNum); + parameter.setMaxConcurrentCallsPerConnection(previousMaxConcurrentCalls); + parameter.setFlowControlWindow(previousFlowControlWindow); + parameter.setMaxConnectionIdleInMillis(previousMaxConnectionIdle); + parameter.setMaxConnectionAgeInMillis(previousMaxConnectionAge); + parameter.setMaxMessageSize(previousMaxMessageSize); + parameter.setMaxHeaderListSize(previousMaxHeaderListSize); + parameter.setRpcMaxRstStream(previousMaxRstStream); + parameter.setRpcSecondsPerWindow(previousSecondsPerWindow); + parameter.setRpcReflectionServiceEnable(previousReflectionServiceEnable); + } + + @Test + public void shouldRejectExcessStreamsBeforeClientAcknowledgesSettings() throws Exception { + 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(); + output.write(CLIENT_PREFACE); + output.write(EMPTY_SETTINGS_FRAME); + output.write(newHeadersFrame(1)); + output.write(newHeadersFrame(3)); + output.write(newHeadersFrame(5)); + output.flush(); + + assertSettingsAndRefusedStream(socket.getInputStream(), 2, 5); + } finally { + server.shutdownNow(); + assertTrue(server.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private static ServerServiceDefinition newHoldService() { + MethodDescriptor method = + MethodDescriptor.newBuilder() + .setType(MethodDescriptor.MethodType.BIDI_STREAMING) + .setFullMethodName(MethodDescriptor.generateFullMethodName( + SERVICE_NAME, METHOD_NAME)) + .setRequestMarshaller(ProtoUtils.marshaller(Empty.getDefaultInstance())) + .setResponseMarshaller(ProtoUtils.marshaller(Empty.getDefaultInstance())) + .build(); + + return ServerServiceDefinition.builder(SERVICE_NAME) + .addMethod(method, new ServerCallHandler() { + @Override + public ServerCall.Listener startCall( + ServerCall call, Metadata headers) { + return new ServerCall.Listener() { + }; + } + }) + .build(); + } + + private static byte[] newHeadersFrame(int streamId) throws Exception { + DefaultHttp2HeadersEncoder encoder = new DefaultHttp2HeadersEncoder(); + ByteBuf headerBlock = Unpooled.buffer(); + ByteBuf frame = Unpooled.buffer(); + try { + Http2Headers headers = new DefaultHttp2Headers() + .method("POST") + .scheme("http") + .authority("localhost") + .path(METHOD_PATH) + .set("content-type", "application/grpc") + .set("te", "trailers"); + encoder.encodeHeaders(streamId, headers, headerBlock); + + frame.writeMedium(headerBlock.readableBytes()); + frame.writeByte(HEADERS_FRAME_TYPE); + frame.writeByte(0x4); + frame.writeInt(streamId); + frame.writeBytes(headerBlock); + return ByteBufUtil.getBytes(frame); + } finally { + frame.release(); + headerBlock.release(); + encoder.close(); + } + } + + private static void assertSettingsAndRefusedStream( + InputStream input, long expectedMaxConcurrentStreams, int expectedRefusedStreamId) + throws IOException { + boolean advertisedLimitFound = false; + for (int i = 0; i < 20; i++) { + Http2Frame frame = readFrame(input); + if (frame.type == SETTINGS_FRAME_TYPE && frame.streamId == 0) { + advertisedLimitFound |= hasSetting( + frame.payload, SETTINGS_MAX_CONCURRENT_STREAMS, expectedMaxConcurrentStreams); + } + if (frame.type == RST_STREAM_FRAME_TYPE && frame.streamId == expectedRefusedStreamId) { + assertTrue("Server did not advertise the enforced concurrent-stream limit", + advertisedLimitFound); + assertEquals(4, frame.payload.length); + long errorCode = ByteBuffer.wrap(frame.payload).getInt() & 0xffff_ffffL; + assertEquals(Http2Error.REFUSED_STREAM.code(), errorCode); + return; + } + if (frame.type == GO_AWAY_FRAME_TYPE) { + fail("Server closed the connection instead of refusing only the excess stream"); + } + } + fail("No REFUSED_STREAM response for stream " + expectedRefusedStreamId); + } + + private static boolean hasSetting(byte[] payload, int expectedId, long expectedValue) { + assertEquals("Invalid HTTP/2 SETTINGS payload length", 0, payload.length % 6); + ByteBuffer settings = ByteBuffer.wrap(payload); + while (settings.remaining() >= 6) { + int id = settings.getShort() & 0xffff; + long value = settings.getInt() & 0xffff_ffffL; + if (id == expectedId) { + assertEquals(expectedValue, value); + return true; + } + } + return false; + } + + private static Http2Frame readFrame(InputStream input) throws IOException { + byte[] header = readFully(input, 9); + int payloadLength = + ((header[0] & 0xff) << 16) | ((header[1] & 0xff) << 8) | (header[2] & 0xff); + int type = header[3] & 0xff; + int streamId = ByteBuffer.wrap(header, 5, 4).getInt() & 0x7fff_ffff; + return new Http2Frame(type, streamId, readFully(input, payloadLength)); + } + + private static byte[] readFully(InputStream input, int length) throws IOException { + ByteArrayOutputStream output = new ByteArrayOutputStream(length); + byte[] buffer = new byte[min(length, 1024)]; + while (output.size() < length) { + int read = input.read(buffer, 0, min(buffer.length, length - output.size())); + if (read < 0) { + throw new EOFException("Unexpected end of HTTP/2 frame"); + } + output.write(buffer, 0, read); + } + return output.toByteArray(); + } + + private static final class Http2Frame { + + private final int type; + private final int streamId; + private final byte[] payload; + + private Http2Frame(int type, int streamId, byte[] payload) { + this.type = type; + this.streamId = streamId; + this.payload = payload; + } + } + + private static final class TestRpcService extends RpcService { + + private TestRpcService() { + port = 0; + } + + private NettyServerBuilder newServerBuilder() { + return initServerBuilder(); + } + + @Override + protected void addService(NettyServerBuilder serverBuilder) { + } + } +} diff --git a/framework/src/test/java/org/tron/core/config/args/ArgsTest.java b/framework/src/test/java/org/tron/core/config/args/ArgsTest.java index 076a8ab5387..36b8a3269c1 100644 --- a/framework/src/test/java/org/tron/core/config/args/ArgsTest.java +++ b/framework/src/test/java/org/tron/core/config/args/ArgsTest.java @@ -103,7 +103,7 @@ public void get() { // gRPC network configs checking Assert.assertEquals(50051, parameter.getRpcPort()); - Assert.assertEquals(Integer.MAX_VALUE, parameter.getMaxConcurrentCallsPerConnection()); + Assert.assertEquals(100, parameter.getMaxConcurrentCallsPerConnection()); Assert .assertEquals(NettyServerBuilder .DEFAULT_FLOW_CONTROL_WINDOW, parameter.getFlowControlWindow()); diff --git a/gradle/java-tron.vmoptions b/gradle/java-tron.vmoptions index e994a332740..cf34689ebdd 100644 --- a/gradle/java-tron.vmoptions +++ b/gradle/java-tron.vmoptions @@ -4,4 +4,5 @@ -XX:+PrintGCDateStamps -XX:+CMSParallelRemarkEnabled -XX:ReservedCodeCacheSize=256m --XX:+CMSScavengeBeforeRemark \ No newline at end of file +-XX:+CMSScavengeBeforeRemark +-Dio.netty.allocator.type=pooled diff --git a/gradle/jdk17/java-tron.vmoptions b/gradle/jdk17/java-tron.vmoptions index 7af3123d268..180ff1aa7c7 100644 --- a/gradle/jdk17/java-tron.vmoptions +++ b/gradle/jdk17/java-tron.vmoptions @@ -5,4 +5,5 @@ -XX:MetaspaceSize=256m -XX:MaxMetaspaceSize=512m -XX:MaxDirectMemorySize=1g --XX:+HeapDumpOnOutOfMemoryError \ No newline at end of file +-XX:+HeapDumpOnOutOfMemoryError +-Dio.netty.allocator.type=pooled diff --git a/gradle/verification-metadata.xml b/gradle/verification-metadata.xml index 832d2728f0b..61cdda33edd 100644 --- a/gradle/verification-metadata.xml +++ b/gradle/verification-metadata.xml @@ -432,6 +432,14 @@ + + + + + + + + @@ -445,6 +453,11 @@ + + + + + @@ -490,6 +503,14 @@ + + + + + + + + @@ -526,6 +547,11 @@ + + + + + @@ -635,6 +661,17 @@ + + + + + + + + + + + @@ -680,6 +717,11 @@ + + + + + @@ -730,6 +772,11 @@ + + + + + @@ -754,6 +801,14 @@ + + + + + + + + @@ -762,6 +817,14 @@ + + + + + + + + @@ -772,6 +835,11 @@ + + + + + @@ -1103,76 +1171,76 @@ - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + @@ -1183,111 +1251,127 @@ - - - + + + - - + + - - + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + - - + + - - - + + + + + + - - - + + + - - + + - - - + + + - - + + + + + + + + + + + + + + + - - - + + + - - + + @@ -1877,6 +1961,14 @@ + + + + + + + + diff --git a/protocol/build.gradle b/protocol/build.gradle index 0ce01a9bfb8..ed8914343b8 100644 --- a/protocol/build.gradle +++ b/protocol/build.gradle @@ -2,8 +2,6 @@ apply plugin: 'com.google.protobuf' apply from: 'protoLint.gradle' def protobufVersion = '3.25.8' -// keep same version as protoc-gen-grpc-java for arm64 or macOS, see rootProject.archInfo.requires.ProtocGenVersion -def grpcVersion = '1.81.0' dependencies { api group: 'com.google.protobuf', name: 'protobuf-java', version: protobufVersion @@ -12,11 +10,11 @@ dependencies { // checkstyleConfig "com.puppycrawl.tools:checkstyle:${versions.checkstyle}" // google grpc - api group: 'io.grpc', name: 'grpc-netty', version: grpcVersion - api group: 'io.grpc', name: 'grpc-protobuf', version: grpcVersion - api group: 'io.grpc', name: 'grpc-stub', version: grpcVersion - api group: 'io.grpc', name: 'grpc-core', version: grpcVersion - api group: 'io.grpc', name: 'grpc-services', version: grpcVersion + api group: 'io.grpc', name: 'grpc-netty', version: rootProject.grpcVersion + api group: 'io.grpc', name: 'grpc-protobuf', version: rootProject.grpcVersion + api group: 'io.grpc', name: 'grpc-stub', version: rootProject.grpcVersion + api group: 'io.grpc', name: 'grpc-core', version: rootProject.grpcVersion + api group: 'io.grpc', name: 'grpc-services', version: rootProject.grpcVersion // end google grpc diff --git a/start.sh b/start.sh index 89f13cf25a7..1472a94dc62 100644 --- a/start.sh +++ b/start.sh @@ -358,7 +358,8 @@ startService() { nohup $JAVACMD -Xms$JVM_MS -Xmx$JVM_MX -XX:+UseConcMarkSweepGC -XX:+PrintGCDetails -Xloggc:./gc.log \ -XX:+PrintGCDateStamps -XX:+CMSParallelRemarkEnabled -XX:ReservedCodeCacheSize=256m -XX:+UseCodeCacheFlushing \ -XX:MetaspaceSize=256m -XX:MaxMetaspaceSize=512m \ - -XX:MaxDirectMemorySize=$MAX_DIRECT_MEMORY -XX:+HeapDumpOnOutOfMemoryError \ + -XX:MaxDirectMemorySize=$MAX_DIRECT_MEMORY -Dio.netty.allocator.type=pooled \ + -XX:+HeapDumpOnOutOfMemoryError \ -XX:NewRatio=2 -jar \ $JAR_NAME $FULL_START_OPT >>start.log 2>&1 & checkPid diff --git a/start.sh.simple b/start.sh.simple index 52548dea62b..109f0dc85a3 100644 --- a/start.sh.simple +++ b/start.sh.simple @@ -137,6 +137,7 @@ startService() { -XX:MetaspaceSize=256m \ -XX:MaxMetaspaceSize=512m \ -XX:MaxDirectMemorySize=1g \ + -Dio.netty.allocator.type=pooled \ -XX:+HeapDumpOnOutOfMemoryError \ -jar "$FULL_NODE_JAR" "${FULL_START_OPT[@]}" \ >> start.log 2>&1 &