Browse Source

Merge pull request #8340 from volodymyr-babak/max-inbound-message-size

[3.5] Handle gRPC messages exceeding default max message size
pull/8350/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
d37283a292
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  2. 28
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 1
      application/src/main/resources/thingsboard.yml
  4. 11
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  5. 13
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  6. 2
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java
  7. 2
      common/edge-api/src/main/proto/edge.proto

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -174,7 +174,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override @Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) { public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> 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 @Override

28
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 EdgeVersion edgeVersion;
private int maxInboundMessageSize;
private int clientMaxInboundMessageSize;
private ScheduledExecutorService sendDownlinkExecutorService; private ScheduledExecutorService sendDownlinkExecutorService;
EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream, BiConsumer<EdgeId, EdgeGrpcSession> sessionOpenListener, EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream, BiConsumer<EdgeId, EdgeGrpcSession> sessionOpenListener,
Consumer<EdgeId> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService) { Consumer<EdgeId> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize) {
this.sessionId = UUID.randomUUID(); this.sessionId = UUID.randomUUID();
this.ctx = ctx; this.ctx = ctx;
this.outputStream = outputStream; this.outputStream = outputStream;
this.sessionOpenListener = sessionOpenListener; this.sessionOpenListener = sessionOpenListener;
this.sessionCloseListener = sessionCloseListener; this.sessionCloseListener = sessionCloseListener;
this.sendDownlinkExecutorService = sendDownlinkExecutorService; this.sendDownlinkExecutorService = sendDownlinkExecutorService;
this.maxInboundMessageSize = maxInboundMessageSize;
initInputStream(); initInputStream();
} }
@ -130,6 +134,10 @@ public final class EdgeGrpcSession implements Closeable {
if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) { if (ConnectResponseCode.ACCEPTED != responseMsg.getResponseCode()) {
outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); outputStream.onError(new RuntimeException(responseMsg.getErrorMsg()));
} else { } else {
if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) {
log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize());
clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize();
}
connected = true; 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()); log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size());
for (DownlinkMsg downlinkMsg : copy) { for (DownlinkMsg downlinkMsg : copy) {
sendDownlinkMsg(ResponseMsg.newBuilder() if (this.clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > this.clientMaxInboundMessageSize) {
.setDownlinkMsg(downlinkMsg) log.error("[{}][{}][{}] Downlink msg size [{}] exceeds client max inbound message size [{}]. Skipping this message. " +
.build()); "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) { if (attempt < MAX_DOWNLINK_ATTEMPTS) {
scheduleDownlinkMsgsPackSend(attempt + 1); scheduleDownlinkMsgsPackSend(attempt + 1);
@ -638,7 +654,9 @@ public final class EdgeGrpcSession implements Closeable {
return ConnectResponseMsg.newBuilder() return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.ACCEPTED) .setResponseCode(ConnectResponseCode.ACCEPTED)
.setErrorMsg("") .setErrorMsg("")
.setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build(); .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge))
.setMaxInboundMessageSize(maxInboundMessageSize)
.build();
} }
return ConnectResponseMsg.newBuilder() return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS)

1
application/src/main/resources/thingsboard.yml

@ -961,7 +961,6 @@ edges:
scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}" scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}"
send_scheduler_pool_size: "${EDGES_SEND_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}" grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}"
edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}"
state: state:
persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}" persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}"

11
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -104,13 +104,14 @@ public class EdgeImitator {
ignoredTypes = new ArrayList<>(); ignoredTypes = new ArrayList<>();
this.routingKey = routingKey; this.routingKey = routingKey;
this.routingSecret = routingSecret; this.routingSecret = routingSecret;
setEdgeCredentials("rpcHost", host); updateEdgeClientFields("rpcHost", host);
setEdgeCredentials("rpcPort", port); updateEdgeClientFields("rpcPort", port);
setEdgeCredentials("timeoutSecs", 3); updateEdgeClientFields("timeoutSecs", 3);
setEdgeCredentials("keepAliveTimeSec", 300); 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); Field fieldToSet = edgeRpcClient.getClass().getDeclaredField(fieldName);
fieldToSet.setAccessible(true); fieldToSet.setAccessible(true);
fieldToSet.set(edgeRpcClient, value); fieldToSet.set(edgeRpcClient, value);

13
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.grpc.netty.NettyChannelBuilder;
import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder; import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
@ -62,6 +63,10 @@ public class EdgeGrpcClient implements EdgeRpcClient {
private boolean sslEnabled; private boolean sslEnabled;
@Value("${cloud.rpc.ssl.cert:}") @Value("${cloud.rpc.ssl.cert:}")
private String certResource; private String certResource;
@Value("${cloud.rpc.max_inbound_message_size:4194304}")
private int maxInboundMessageSize;
@Getter
private int serverMaxInboundMessageSize;
private ManagedChannel channel; private ManagedChannel channel;
@ -77,6 +82,7 @@ public class EdgeGrpcClient implements EdgeRpcClient {
Consumer<DownlinkMsg> onDownlink, Consumer<DownlinkMsg> onDownlink,
Consumer<Exception> onError) { Consumer<Exception> onError) {
NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort) NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort)
.maxInboundMessageSize(maxInboundMessageSize)
.keepAliveTime(keepAliveTimeSec, TimeUnit.SECONDS); .keepAliveTime(keepAliveTimeSec, TimeUnit.SECONDS);
if (sslEnabled) { if (sslEnabled) {
try { try {
@ -101,7 +107,8 @@ public class EdgeGrpcClient implements EdgeRpcClient {
.setConnectRequestMsg(ConnectRequestMsg.newBuilder() .setConnectRequestMsg(ConnectRequestMsg.newBuilder()
.setEdgeRoutingKey(edgeKey) .setEdgeRoutingKey(edgeKey)
.setEdgeSecret(edgeSecret) .setEdgeSecret(edgeSecret)
.setEdgeVersion(EdgeVersion.V_3_3_3) .setEdgeVersion(EdgeVersion.V_3_4_0)
.setMaxInboundMessageSize(maxInboundMessageSize)
.build()) .build())
.build()); .build());
} }
@ -117,6 +124,10 @@ public class EdgeGrpcClient implements EdgeRpcClient {
if (responseMsg.hasConnectResponseMsg()) { if (responseMsg.hasConnectResponseMsg()) {
ConnectResponseMsg connectResponseMsg = responseMsg.getConnectResponseMsg(); ConnectResponseMsg connectResponseMsg = responseMsg.getConnectResponseMsg();
if (connectResponseMsg.getResponseCode().equals(ConnectResponseCode.ACCEPTED)) { 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()); log.info("[{}] Configuration received: {}", edgeKey, connectResponseMsg.getConfiguration());
onEdgeUpdate.accept(connectResponseMsg.getConfiguration()); onEdgeUpdate.accept(connectResponseMsg.getConfiguration());
} else { } else {

2
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java

@ -41,4 +41,6 @@ public interface EdgeRpcClient {
void sendUplinkMsg(UplinkMsg uplinkMsg); void sendUplinkMsg(UplinkMsg uplinkMsg);
void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg);
int getServerMaxInboundMessageSize();
} }

2
common/edge-api/src/main/proto/edge.proto

@ -68,6 +68,7 @@ message ConnectRequestMsg {
string edgeRoutingKey = 1; string edgeRoutingKey = 1;
string edgeSecret = 2; string edgeSecret = 2;
EdgeVersion edgeVersion = 3; EdgeVersion edgeVersion = 3;
optional int32 maxInboundMessageSize = 4;
} }
enum ConnectResponseCode { enum ConnectResponseCode {
@ -80,6 +81,7 @@ message ConnectResponseMsg {
ConnectResponseCode responseCode = 1; ConnectResponseCode responseCode = 1;
string errorMsg = 2; string errorMsg = 2;
EdgeConfiguration configuration = 3; EdgeConfiguration configuration = 3;
optional int32 maxInboundMessageSize = 4;
} }
message SyncRequestMsg { message SyncRequestMsg {

Loading…
Cancel
Save