|
|
|
@ -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<ResponseMsg> outputStream, BiConsumer<EdgeId, EdgeGrpcSession> sessionOpenListener, |
|
|
|
Consumer<EdgeId> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService) { |
|
|
|
Consumer<EdgeId> 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) |
|
|
|
|