From e496ef783956f486c0247e53903acbedacfddc29 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 10 Apr 2023 10:52:51 +0300 Subject: [PATCH 1/4] Fix process of messages that exceeds max inbound message size --- .../service/edge/rpc/EdgeGrpcService.java | 2 +- .../service/edge/rpc/EdgeGrpcSession.java | 26 +++++++++++++++---- .../thingsboard/edge/rpc/EdgeGrpcClient.java | 9 +++++++ .../thingsboard/edge/rpc/EdgeRpcClient.java | 2 ++ common/edge-api/src/main/proto/edge.proto | 2 ++ 5 files changed, 35 insertions(+), 6 deletions(-) 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..677f743d01 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,8 @@ public final class EdgeGrpcSession implements Closeable { if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); } else { + log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); + clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); connected = true; } } @@ -408,9 +414,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 {}", + 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 +652,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/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..acbe2024a0 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 { @@ -102,6 +108,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { .setEdgeRoutingKey(edgeKey) .setEdgeSecret(edgeSecret) .setEdgeVersion(EdgeVersion.V_3_3_3) + .setMaxInboundMessageSize(maxInboundMessageSize) .build()) .build()); } @@ -117,6 +124,8 @@ public class EdgeGrpcClient implements EdgeRpcClient { if (responseMsg.hasConnectResponseMsg()) { ConnectResponseMsg connectResponseMsg = responseMsg.getConnectResponseMsg(); if (connectResponseMsg.getResponseCode().equals(ConnectResponseCode.ACCEPTED)) { + 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..afe98b11ae 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; + int32 maxInboundMessageSize = 4; } enum ConnectResponseCode { @@ -80,6 +81,7 @@ message ConnectResponseMsg { ConnectResponseCode responseCode = 1; string errorMsg = 2; EdgeConfiguration configuration = 3; + int32 maxInboundMessageSize = 4; } message SyncRequestMsg { From 63703bcfdb0f23abc883ad1bf6c7df8fc0804980 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 10 Apr 2023 11:03:49 +0300 Subject: [PATCH 2/4] Edge client - updated edge version --- .../src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 acbe2024a0..0dfd0f7cb2 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 @@ -107,7 +107,7 @@ 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()); From c95cff74f7a7ff6cba94a869a363b85359bcd6c5 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 10 Apr 2023 13:12:11 +0300 Subject: [PATCH 3/4] Fixed edge client tests by setting correctly max inbound message size --- .../server/service/edge/rpc/EdgeGrpcSession.java | 6 ++++-- application/src/main/resources/thingsboard.yml | 1 - .../server/edge/imitator/EdgeImitator.java | 11 ++++++----- .../java/org/thingsboard/edge/rpc/EdgeGrpcClient.java | 6 ++++-- common/edge-api/src/main/proto/edge.proto | 4 ++-- 5 files changed, 16 insertions(+), 12 deletions(-) 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 677f743d01..5931b413fd 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 @@ -134,8 +134,10 @@ public final class EdgeGrpcSession implements Closeable { if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); } else { - log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); - clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); + if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) { + log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); + clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); + } connected = true; } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f745e03ffc..54dda31576 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -960,7 +960,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 0dfd0f7cb2..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 @@ -124,8 +124,10 @@ public class EdgeGrpcClient implements EdgeRpcClient { if (responseMsg.hasConnectResponseMsg()) { ConnectResponseMsg connectResponseMsg = responseMsg.getConnectResponseMsg(); if (connectResponseMsg.getResponseCode().equals(ConnectResponseCode.ACCEPTED)) { - log.debug("[{}] Server max inbound message size: {}", edgeKey, connectResponseMsg.getMaxInboundMessageSize()); - serverMaxInboundMessageSize = connectResponseMsg.getMaxInboundMessageSize(); + 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/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index afe98b11ae..afcf4056c4 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -68,7 +68,7 @@ message ConnectRequestMsg { string edgeRoutingKey = 1; string edgeSecret = 2; EdgeVersion edgeVersion = 3; - int32 maxInboundMessageSize = 4; + optional int32 maxInboundMessageSize = 4; } enum ConnectResponseCode { @@ -81,7 +81,7 @@ message ConnectResponseMsg { ConnectResponseCode responseCode = 1; string errorMsg = 2; EdgeConfiguration configuration = 3; - int32 maxInboundMessageSize = 4; + optional int32 maxInboundMessageSize = 4; } message SyncRequestMsg { From 9bda3880eb0ea30c1cfbca9a57276e61e1f2a11d Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 10 Apr 2023 15:58:02 +0300 Subject: [PATCH 4/4] Edge session - added tenant and edge id logging in case message exceeds size limits --- .../server/service/edge/rpc/EdgeGrpcSession.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 5931b413fd..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 @@ -417,10 +417,10 @@ public final class EdgeGrpcSession implements Closeable { log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size()); for (DownlinkMsg downlinkMsg : copy) { if (this.clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > this.clientMaxInboundMessageSize) { - log.error("[{}] Downlink msg size [{}] exceeds client max inbound message size [{}]. Skipping this message. " + + 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 {}", - this.sessionId, downlinkMsg.getSerializedSize(), this.clientMaxInboundMessageSize, downlinkMsg); + "Message {}", edge.getTenantId(), edge.getId(), this.sessionId, downlinkMsg.getSerializedSize(), + this.clientMaxInboundMessageSize, downlinkMsg); sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId()); } else { sendDownlinkMsg(ResponseMsg.newBuilder()