diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index f695fa1d0e..1e30980df0 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -174,7 +174,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService).getInputStream(); + return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, this.maxInboundMessageSize).getInputStream(); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index efa0fb4de7..9e3a22f64f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -105,16 +105,20 @@ public final class EdgeGrpcSession implements Closeable { private EdgeVersion edgeVersion; + private int maxInboundMessageSize; + private int clientMaxInboundMessageSize; + private ScheduledExecutorService sendDownlinkExecutorService; EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, BiConsumer sessionOpenListener, - Consumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService) { + Consumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize) { this.sessionId = UUID.randomUUID(); this.ctx = ctx; this.outputStream = outputStream; this.sessionOpenListener = sessionOpenListener; this.sessionCloseListener = sessionCloseListener; this.sendDownlinkExecutorService = sendDownlinkExecutorService; + this.maxInboundMessageSize = maxInboundMessageSize; initInputStream(); } @@ -130,6 +134,10 @@ public final class EdgeGrpcSession implements Closeable { if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); } else { + if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) { + log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); + clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); + } connected = true; } } @@ -408,9 +416,17 @@ public final class EdgeGrpcSession implements Closeable { } log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size()); for (DownlinkMsg downlinkMsg : copy) { - sendDownlinkMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .build()); + if (this.clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > this.clientMaxInboundMessageSize) { + log.error("[{}][{}][{}] Downlink msg size [{}] exceeds client max inbound message size [{}]. Skipping this message. " + + "Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE env variable on the edge and restart it." + + "Message {}", edge.getTenantId(), edge.getId(), this.sessionId, downlinkMsg.getSerializedSize(), + this.clientMaxInboundMessageSize, downlinkMsg); + sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId()); + } else { + sendDownlinkMsg(ResponseMsg.newBuilder() + .setDownlinkMsg(downlinkMsg) + .build()); + } } if (attempt < MAX_DOWNLINK_ATTEMPTS) { scheduleDownlinkMsgsPackSend(attempt + 1); @@ -638,7 +654,9 @@ public final class EdgeGrpcSession implements Closeable { return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.ACCEPTED) .setErrorMsg("") - .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build(); + .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)) + .setMaxInboundMessageSize(maxInboundMessageSize) + .build(); } return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 5eb0a70cf4..94751fc23b 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -961,7 +961,6 @@ edges: scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}" send_scheduler_pool_size: "${EDGES_SEND_SCHEDULER_POOL_SIZE:1}" grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}" - edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}" state: persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}" diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 26dd49b90e..b618ac7718 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -104,13 +104,14 @@ public class EdgeImitator { ignoredTypes = new ArrayList<>(); this.routingKey = routingKey; this.routingSecret = routingSecret; - setEdgeCredentials("rpcHost", host); - setEdgeCredentials("rpcPort", port); - setEdgeCredentials("timeoutSecs", 3); - setEdgeCredentials("keepAliveTimeSec", 300); + updateEdgeClientFields("rpcHost", host); + updateEdgeClientFields("rpcPort", port); + updateEdgeClientFields("timeoutSecs", 3); + updateEdgeClientFields("keepAliveTimeSec", 300); + updateEdgeClientFields("maxInboundMessageSize", 4194304); } - private void setEdgeCredentials(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException { + private void updateEdgeClientFields(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException { Field fieldToSet = edgeRpcClient.getClass().getDeclaredField(fieldName); fieldToSet.setAccessible(true); fieldToSet.set(edgeRpcClient, value); diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java index 750c285317..2acbe49f98 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java @@ -20,6 +20,7 @@ import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts; import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder; import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder; import io.grpc.stub.StreamObserver; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; @@ -62,6 +63,10 @@ public class EdgeGrpcClient implements EdgeRpcClient { private boolean sslEnabled; @Value("${cloud.rpc.ssl.cert:}") private String certResource; + @Value("${cloud.rpc.max_inbound_message_size:4194304}") + private int maxInboundMessageSize; + @Getter + private int serverMaxInboundMessageSize; private ManagedChannel channel; @@ -77,6 +82,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { Consumer onDownlink, Consumer onError) { NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort) + .maxInboundMessageSize(maxInboundMessageSize) .keepAliveTime(keepAliveTimeSec, TimeUnit.SECONDS); if (sslEnabled) { try { @@ -101,7 +107,8 @@ public class EdgeGrpcClient implements EdgeRpcClient { .setConnectRequestMsg(ConnectRequestMsg.newBuilder() .setEdgeRoutingKey(edgeKey) .setEdgeSecret(edgeSecret) - .setEdgeVersion(EdgeVersion.V_3_3_3) + .setEdgeVersion(EdgeVersion.V_3_4_0) + .setMaxInboundMessageSize(maxInboundMessageSize) .build()) .build()); } @@ -117,6 +124,10 @@ public class EdgeGrpcClient implements EdgeRpcClient { if (responseMsg.hasConnectResponseMsg()) { ConnectResponseMsg connectResponseMsg = responseMsg.getConnectResponseMsg(); if (connectResponseMsg.getResponseCode().equals(ConnectResponseCode.ACCEPTED)) { + if (connectResponseMsg.hasMaxInboundMessageSize()) { + log.debug("[{}] Server max inbound message size: {}", edgeKey, connectResponseMsg.getMaxInboundMessageSize()); + serverMaxInboundMessageSize = connectResponseMsg.getMaxInboundMessageSize(); + } log.info("[{}] Configuration received: {}", edgeKey, connectResponseMsg.getConfiguration()); onEdgeUpdate.accept(connectResponseMsg.getConfiguration()); } else { diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java index 5fc4e86c60..44d00e22a8 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java @@ -41,4 +41,6 @@ public interface EdgeRpcClient { void sendUplinkMsg(UplinkMsg uplinkMsg); void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); + + int getServerMaxInboundMessageSize(); } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 76fc1be7d1..afcf4056c4 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -68,6 +68,7 @@ message ConnectRequestMsg { string edgeRoutingKey = 1; string edgeSecret = 2; EdgeVersion edgeVersion = 3; + optional int32 maxInboundMessageSize = 4; } enum ConnectResponseCode { @@ -80,6 +81,7 @@ message ConnectResponseMsg { ConnectResponseCode responseCode = 1; string errorMsg = 2; EdgeConfiguration configuration = 3; + optional int32 maxInboundMessageSize = 4; } message SyncRequestMsg {