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; }