From 22118af4db7cc2fba246dc8f82f313983a46121a Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 8 Nov 2022 13:13:58 +0200 Subject: [PATCH 01/16] Added Edge RPC response/request support --- .../device/DeviceActorMessageProcessor.java | 5 +- .../service/edge/rpc/EdgeGrpcSession.java | 5 +- .../rpc/constructor/DeviceMsgConstructor.java | 46 +++++-- .../rpc/processor/DeviceEdgeProcessor.java | 113 ++++++++++++++++-- .../rpc/processor/EdgeRpcRequestMetadata.java | 28 +++++ .../server/edge/BaseDeviceEdgeTest.java | 3 +- .../server/common/data/DataConstants.java | 2 + .../common/data/edge/EdgeEventActionType.java | 3 +- common/edge-api/src/main/proto/edge.proto | 3 + .../rule/engine/api/RuleEngineRpcService.java | 2 - .../rule/engine/rpc/TbSendRPCReplyNode.java | 62 +++++++++- 11 files changed, 244 insertions(+), 28 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 2ffe2002f9..9b95b188d2 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -823,8 +823,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { body.put("expirationTime", msg.getExpirationTime()); body.put("method", msg.getBody().getMethod()); body.put("params", msg.getBody().getParams()); + body.put("persisted", msg.isPersisted()); + body.put("retries", msg.getRetries()); + body.put("additionalInfo", msg.getAdditionalInfo()); - EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body); + EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL_REQUEST, deviceId, body); return Futures.transform(systemContext.getEdgeEventService().saveAsync(edgeEvent), unused -> { systemContext.getClusterService().onEdgeEventUpdate(tenantId, edgeId); 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 2342997c65..aca2f8b562 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 @@ -463,7 +463,8 @@ public final class EdgeGrpcSession implements Closeable { case UNASSIGNED_FROM_CUSTOMER: case CREDENTIALS_REQUEST: case ENTITY_MERGE_REQUEST: - case RPC_CALL: + case RPC_CALL_REQUEST: + case RPC_CALL_RESPONSE: downlinkMsg = convertEntityEventToDownlink(edgeEvent); log.trace("[{}][{}] entity message processed [{}]", edgeEvent.getTenantId(), this.sessionId, downlinkMsg); break; @@ -611,7 +612,7 @@ public final class EdgeGrpcSession implements Closeable { } if (uplinkMsg.getDeviceRpcCallMsgCount() > 0) { for (DeviceRpcCallMsg deviceRpcCallMsg : uplinkMsg.getDeviceRpcCallMsgList()) { - result.add(ctx.getDeviceProcessor().processDeviceRpcCallResponseFromEdge(edge.getTenantId(), deviceRpcCallMsg)); + result.add(ctx.getDeviceProcessor().processDeviceRpcCallFromEdge(edge.getTenantId(), edge, deviceRpcCallMsg)); } } if (uplinkMsg.getWidgetBundleTypesRequestMsgCount() > 0) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java index 511910dd8a..49438792fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java @@ -28,6 +28,7 @@ import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.v1.RpcRequestMsg; +import org.thingsboard.server.gen.edge.v1.RpcResponseMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -96,26 +97,51 @@ public class DeviceMsgConstructor { .setIdLSB(deviceId.getId().getLeastSignificantBits()).build(); } - public DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode body) { - int requestId = body.get("requestId").asInt(); - boolean oneway = body.get("oneway").asBoolean(); - UUID requestUUID = UUID.fromString(body.get("requestUUID").asText()); - long expirationTime = body.get("expirationTime").asLong(); + public DeviceRpcCallMsg constructDeviceRpcRequestMsg(UUID deviceId, JsonNode body) { + DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body); + String method = body.get("method").asText(); String params = body.get("params").asText(); - RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder(); requestBuilder.setMethod(method); requestBuilder.setParams(params); - DeviceRpcCallMsg.Builder builder = DeviceRpcCallMsg.newBuilder() + builder.setRequestMsg(requestBuilder.build()); + + return builder.build(); + } + + public DeviceRpcCallMsg constructDeviceRpcResponseMsg(UUID deviceId, JsonNode body) { + DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body); + + RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder(); + if (body.has("error")) { + responseBuilder.setError(body.get("error").asText()); + } else { + responseBuilder.setResponse(body.get("response").asText()); + } + builder.setResponseMsg(responseBuilder.build()); + + return builder.build(); + } + + private DeviceRpcCallMsg.Builder constructDeviceRpcMsg(UUID deviceId, JsonNode body) { + int requestId = body.get("requestId").asInt(); + boolean oneway = body.get("oneway").asBoolean(); + UUID requestUUID = UUID.fromString(body.get("requestUUID").asText()); + long expirationTime = body.get("expirationTime").asLong(); + boolean persisted = body.get("persisted").asBoolean(); + int retries = body.get("retries").asInt(); + String additionalInfo = body.get("additionalInfo").asText(); + return DeviceRpcCallMsg.newBuilder() .setDeviceIdMSB(deviceId.getMostSignificantBits()) .setDeviceIdLSB(deviceId.getLeastSignificantBits()) .setRequestUuidMSB(requestUUID.getMostSignificantBits()) .setRequestUuidLSB(requestUUID.getLeastSignificantBits()) - .setRequestId(requestId) .setExpirationTime(expirationTime) + .setRequestId(requestId) .setOneway(oneway) - .setRequestMsg(requestBuilder.build()); - return builder.build(); + .setPersisted(persisted) + .setRetries(retries) + .setAdditionalInfo(additionalInfo); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index e09e036953..235fef7431 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -25,6 +25,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -52,6 +53,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; @@ -66,8 +68,14 @@ import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; +import javax.annotation.PostConstruct; +import java.util.Map; import java.util.Optional; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantLock; @Component @@ -75,11 +83,19 @@ import java.util.concurrent.locks.ReentrantLock; @TbCoreComponent public class DeviceEdgeProcessor extends BaseEdgeProcessor { + private final Map toServerRpcPendingMap = new ConcurrentHashMap<>(); + private ScheduledExecutorService scheduler; + @Autowired private DataDecodingEncodingService dataDecodingEncodingService; private static final ReentrantLock deviceCreationLock = new ReentrantLock(); + @PostConstruct + public void init(){ + this.scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-edge-processor-scheduler")); + } + public ListenableFuture processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { log.trace("[{}] onDeviceUpdate [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); switch (deviceUpdateMsg.getMsgType()) { @@ -325,8 +341,17 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return metaData; } - public ListenableFuture processDeviceRpcCallResponseFromEdge(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { - log.trace("[{}] processDeviceRpcCallResponseMsg [{}]", tenantId, deviceRpcCallMsg); + public ListenableFuture processDeviceRpcCallFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) { + log.trace("[{}] processDeviceRpcCallFromEdge [{}]", tenantId, deviceRpcCallMsg); + if (deviceRpcCallMsg.hasResponseMsg()) { + return processDeviceRpcResponseFromEdge(tenantId, deviceRpcCallMsg); + } else if (deviceRpcCallMsg.hasRequestMsg()) { + return processDeviceRpcRequestFromEdge(tenantId, edge, deviceRpcCallMsg); + } + return Futures.immediateFuture(null); + } + + private ListenableFuture processDeviceRpcResponseFromEdge(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { SettableFuture futureToSet = SettableFuture.create(); UUID requestUuid = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB()); DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB())); @@ -357,6 +382,68 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return futureToSet; } + private ListenableFuture processDeviceRpcRequestFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) { + DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB())); + UUID requestUUID = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB()); + try { + ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + TbMsgMetaData metaData = new TbMsgMetaData(); + String requestId = Integer.toString(deviceRpcCallMsg.getRequestId()); + metaData.putValue("requestId", requestId); + metaData.putValue("requestUUID", requestUUID.toString()); + // ?? metaData.putValue("originServiceId", deviceRpcRequestMsg.get); + metaData.putValue("expirationTime", Long.toString(deviceRpcCallMsg.getExpirationTime())); + metaData.putValue("oneway", Boolean.toString(deviceRpcCallMsg.getOneway())); + metaData.putValue(DataConstants.PERSISTENT, Boolean.toString(deviceRpcCallMsg.getPersisted())); + + if (deviceRpcCallMsg.getRetries() > 0) { + metaData.putValue(DataConstants.RETRIES, Integer.toString(deviceRpcCallMsg.getRetries())); + } + + metaData.putValue(DataConstants.EDGE_ID, edge.getId().toString()); + + Device device = deviceService.findDeviceById(tenantId, deviceId); + if (device != null) { + metaData.putValue("deviceName", device.getName()); + metaData.putValue("deviceType", device.getType()); + metaData.putValue(DataConstants.DEVICE_ID, deviceId.getId().toString()); + } + + entityNode.put("method", deviceRpcCallMsg.getRequestMsg().getMethod()); + entityNode.put("params", deviceRpcCallMsg.getRequestMsg().getParams()); + + entityNode.put(DataConstants.ADDITIONAL_INFO, deviceRpcCallMsg.getAdditionalInfo()); + TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, null, metaData, + TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode)); + tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", device); + } + + @Override + public void onFailure(Throwable t) { + log.debug("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", device, t); + } + }); + toServerRpcPendingMap.put(requestId, new EdgeRpcRequestMetadata(tenantId, edge.getId(), deviceId)); + scheduler.schedule(() -> processTimeout(requestId), 60000, TimeUnit.MILLISECONDS); + } catch (JsonProcessingException | IllegalArgumentException e) { + log.warn("[{}] Failed to push device action to rule engine: {}", deviceId, DataConstants.ENTITY_CREATED, e); + } + + return Futures.immediateFuture(null); + } + + private void processTimeout(String requestId) { + EdgeRpcRequestMetadata data = toServerRpcPendingMap.remove(requestId); + if (data != null) { + // TODO: add failure body + saveEdgeEvent(data.getTenantId(), data.getEdgeId(), EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL_RESPONSE, + data.getDeviceId(), JacksonUtil.OBJECT_MAPPER.valueToTree("{}")); + } + } + public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent) { DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; @@ -401,8 +488,10 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { .build(); } break; - case RPC_CALL: - return convertRpcCallEventToDownlink(edgeEvent); + case RPC_CALL_REQUEST: + return convertRpcCallRequestEventToDownlink(edgeEvent); + case RPC_CALL_RESPONSE: + return convertRpcCallResponseEventToDownlink(edgeEvent); case CREDENTIALS_REQUEST: return convertCredentialsRequestEventToDownlink(edgeEvent); case ENTITY_MERGE_REQUEST: @@ -411,10 +500,18 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } - private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) { - log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent); - DeviceRpcCallMsg deviceRpcCallMsg = - deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()); + private DownlinkMsg convertRpcCallRequestEventToDownlink(EdgeEvent edgeEvent) { + log.trace("Executing convertRpcCallRequestEventToDownlink, edgeEvent [{}]", edgeEvent); + return DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcRequestMsg(edgeEvent.getEntityId(), edgeEvent.getBody())) + .build(); + } + + private DownlinkMsg convertRpcCallResponseEventToDownlink(EdgeEvent edgeEvent) { + log.trace("Executing convertRpcCallResponseEventToDownlink, edgeEvent [{}]", edgeEvent); + DeviceRpcCallMsg deviceRpcCallMsg = deviceMsgConstructor.constructDeviceRpcResponseMsg(edgeEvent.getEntityId(), edgeEvent.getBody()); + toServerRpcPendingMap.remove(Integer.toString(deviceRpcCallMsg.getRequestId())); return DownlinkMsg.newBuilder() .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .addDeviceRpcCallMsg(deviceRpcCallMsg) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java new file mode 100644 index 0000000000..a5a71d0743 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java @@ -0,0 +1,28 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.edge.rpc.processor; + +import lombok.Data; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; + +@Data +public class EdgeRpcRequestMetadata { + private final TenantId tenantId; + private final EdgeId edgeId; + private final DeviceId deviceId; +} diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java index 392a6814d4..8e8c616d7a 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java @@ -512,7 +512,8 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { body.put("method", "test_method"); body.put("params", "{\"param1\":\"value1\"}"); - EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL_REQUEST, + device.getId().getId(), EdgeEventType.DEVICE, body); edgeImitator.expectMessageAmount(1); edgeEventService.saveAsync(edgeEvent).get(); clusterService.onEdgeEventUpdate(tenantId, edge.getId()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 85b9e681a1..69c36697d3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -40,6 +40,8 @@ public class DataConstants { public static final String EXPIRATION_TIME = "expirationTime"; public static final String ADDITIONAL_INFO = "additionalInfo"; public static final String RETRIES = "retries"; + public static final String EDGE_ID = "edgeId"; + public static final String DEVICE_ID = "deviceId"; public static final String COAP_TRANSPORT_NAME = "COAP"; public static final String LWM2M_TRANSPORT_NAME = "LWM2M"; public static final String MQTT_TRANSPORT_NAME = "MQTT"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java index bba8767fca..7fae6c56b1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java @@ -28,7 +28,8 @@ public enum EdgeEventActionType { UNASSIGNED_FROM_CUSTOMER, RELATION_ADD_OR_UPDATE, RELATION_DELETED, - RPC_CALL, + RPC_CALL_REQUEST, + RPC_CALL_RESPONSE, ALARM_ACK, ALARM_CLEAR, ASSIGNED_TO_EDGE, diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index a2e3b5989b..a5785d6870 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -430,6 +430,9 @@ message DeviceRpcCallMsg { bool oneway = 7; RpcRequestMsg requestMsg = 8; RpcResponseMsg responseMsg = 9; + bool persisted = 10; + int32 retries = 11; + string additionalInfo = 12; } message RpcRequestMsg { diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineRpcService.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineRpcService.java index abb5ace0f8..85596d7e7c 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineRpcService.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineRpcService.java @@ -15,8 +15,6 @@ */ package org.thingsboard.rule.engine.api; -import org.thingsboard.server.common.data.id.DeviceId; - import java.util.UUID; import java.util.function.Consumer; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index ce9b5050f0..22d13f196b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -15,15 +15,27 @@ */ package org.thingsboard.rule.engine.rpc; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.StringUtils; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -65,9 +77,53 @@ public class TbSendRPCReplyNode implements TbNode { } else if (StringUtils.isEmpty(msg.getData())) { ctx.tellFailure(msg, new RuntimeException("Request body is empty!")); } else { - ctx.getRpcService().sendRpcReplyToDevice(serviceIdStr, UUID.fromString(sessionIdStr), Integer.parseInt(requestIdStr), msg.getData()); - ctx.tellSuccess(msg); + if (StringUtils.isNotBlank(msg.getMetaData().getValue(DataConstants.EDGE_ID))) { + saveRpcResponseToEdgeQueue(ctx, msg); + } else { + ctx.getRpcService().sendRpcReplyToDevice(serviceIdStr, UUID.fromString(sessionIdStr), Integer.parseInt(requestIdStr), msg.getData()); + ctx.tellSuccess(msg); + } } } + private void saveRpcResponseToEdgeQueue(TbContext ctx, TbMsg msg) { +// EdgeEvent edgeEvent = new EdgeEvent(); +// edgeEvent.setTenantId(tenantId); +// edgeEvent.setAction(eventAction); +// edgeEvent.setEntityId(entityId); +// edgeEvent.setType(eventType); +// edgeEvent.setBody(entityBody); +// edgeEvent.setEdgeId(edgeId); +// +// ObjectNode body = mapper.createObjectNode(); +// body.put("requestId", requestId); +// body.put("requestUUID", msg.getId().toString()); +// body.put("oneway", msg.isOneway()); +// body.put("expirationTime", msg.getExpirationTime()); +// body.put("method", msg.getBody().getMethod()); +// body.put("params", msg.getBody().getParams()); +// body.put("persisted", msg.isPersisted()); +// body.put("retries", msg.getRetries()); +// body.put("additionalInfo", msg.getAdditionalInfo()); + + EdgeId edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID))); + DeviceId deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID))); + // TODO: add body + EdgeEvent edgeEvent = + EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE, + EdgeEventActionType.RPC_CALL_RESPONSE, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree("{}")); + + ListenableFuture future = ctx.getEdgeEventService().saveAsync(edgeEvent); + Futures.addCallback(future, new FutureCallback() { + @Override + public void onSuccess(@Nullable Void result) { + ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId); + ctx.tellSuccess(msg); + } + + @Override + public void onFailure(Throwable t) { + } + }, ctx.getDbCallbackExecutor()); + } } From 545790fc5b8430357f1f5a47291a92eda36b5f5d Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 10:18:20 +0200 Subject: [PATCH 02/16] Do not register timeout for edge processor on cloud --- .../rpc/processor/DeviceEdgeProcessor.java | 27 ------------------- 1 file changed, 27 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index 235fef7431..b340b252ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -25,7 +25,6 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -68,14 +67,8 @@ import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; -import javax.annotation.PostConstruct; -import java.util.Map; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantLock; @Component @@ -83,19 +76,11 @@ import java.util.concurrent.locks.ReentrantLock; @TbCoreComponent public class DeviceEdgeProcessor extends BaseEdgeProcessor { - private final Map toServerRpcPendingMap = new ConcurrentHashMap<>(); - private ScheduledExecutorService scheduler; - @Autowired private DataDecodingEncodingService dataDecodingEncodingService; private static final ReentrantLock deviceCreationLock = new ReentrantLock(); - @PostConstruct - public void init(){ - this.scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-edge-processor-scheduler")); - } - public ListenableFuture processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { log.trace("[{}] onDeviceUpdate [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); switch (deviceUpdateMsg.getMsgType()) { @@ -426,8 +411,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { log.debug("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", device, t); } }); - toServerRpcPendingMap.put(requestId, new EdgeRpcRequestMetadata(tenantId, edge.getId(), deviceId)); - scheduler.schedule(() -> processTimeout(requestId), 60000, TimeUnit.MILLISECONDS); } catch (JsonProcessingException | IllegalArgumentException e) { log.warn("[{}] Failed to push device action to rule engine: {}", deviceId, DataConstants.ENTITY_CREATED, e); } @@ -435,15 +418,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return Futures.immediateFuture(null); } - private void processTimeout(String requestId) { - EdgeRpcRequestMetadata data = toServerRpcPendingMap.remove(requestId); - if (data != null) { - // TODO: add failure body - saveEdgeEvent(data.getTenantId(), data.getEdgeId(), EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL_RESPONSE, - data.getDeviceId(), JacksonUtil.OBJECT_MAPPER.valueToTree("{}")); - } - } - public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent) { DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; @@ -511,7 +485,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { private DownlinkMsg convertRpcCallResponseEventToDownlink(EdgeEvent edgeEvent) { log.trace("Executing convertRpcCallResponseEventToDownlink, edgeEvent [{}]", edgeEvent); DeviceRpcCallMsg deviceRpcCallMsg = deviceMsgConstructor.constructDeviceRpcResponseMsg(edgeEvent.getEntityId(), edgeEvent.getBody()); - toServerRpcPendingMap.remove(Integer.toString(deviceRpcCallMsg.getRequestId())); return DownlinkMsg.newBuilder() .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .addDeviceRpcCallMsg(deviceRpcCallMsg) From d1312cb917c15eddfc685a8502a67ffcadda3db0 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 12:54:17 +0200 Subject: [PATCH 03/16] Added serviceId and sessionId to deviceRpcCallmsg --- .../rpc/constructor/DeviceMsgConstructor.java | 45 ++++++++++++------- .../rpc/processor/DeviceEdgeProcessor.java | 11 +++-- common/edge-api/src/main/proto/edge.proto | 2 + .../rule/engine/rpc/TbSendRPCReplyNode.java | 20 ++++++--- 4 files changed, 53 insertions(+), 25 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java index 49438792fa..ebea5d482f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java @@ -125,23 +125,36 @@ public class DeviceMsgConstructor { } private DeviceRpcCallMsg.Builder constructDeviceRpcMsg(UUID deviceId, JsonNode body) { - int requestId = body.get("requestId").asInt(); - boolean oneway = body.get("oneway").asBoolean(); - UUID requestUUID = UUID.fromString(body.get("requestUUID").asText()); - long expirationTime = body.get("expirationTime").asLong(); - boolean persisted = body.get("persisted").asBoolean(); - int retries = body.get("retries").asInt(); - String additionalInfo = body.get("additionalInfo").asText(); - return DeviceRpcCallMsg.newBuilder() + DeviceRpcCallMsg.Builder builder = DeviceRpcCallMsg.newBuilder() .setDeviceIdMSB(deviceId.getMostSignificantBits()) .setDeviceIdLSB(deviceId.getLeastSignificantBits()) - .setRequestUuidMSB(requestUUID.getMostSignificantBits()) - .setRequestUuidLSB(requestUUID.getLeastSignificantBits()) - .setExpirationTime(expirationTime) - .setRequestId(requestId) - .setOneway(oneway) - .setPersisted(persisted) - .setRetries(retries) - .setAdditionalInfo(additionalInfo); + .setRequestId(body.get("requestId").asInt()); + if (body.get("oneway") != null) { + builder.setOneway(body.get("oneway").asBoolean()); + } + if (body.get("requestUUID") != null) { + UUID requestUUID = UUID.fromString(body.get("requestUUID").asText()); + builder.setRequestUuidMSB(requestUUID.getMostSignificantBits()) + .setRequestUuidLSB(requestUUID.getLeastSignificantBits()); + } + if (body.get("expirationTime") != null) { + builder.setExpirationTime(body.get("expirationTime").asLong()); + } + if (body.get("persisted") != null) { + builder.setPersisted(body.get("persisted").asBoolean()); + } + if (body.get("retries") != null) { + builder.setRetries(body.get("retries").asInt()); + } + if (body.get("additionalInfo") != null) { + builder.setAdditionalInfo(JacksonUtil.toString(body.get("additionalInfo"))); + } + if (body.get("serviceId") != null) { + builder.setServiceId(body.get("serviceId").asText()); + } + if (body.get("sessionId") != null) { + builder.setSessionId(body.get("sessionId").asText()); + } + return builder; } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index b340b252ef..959f5c5fe8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -376,7 +376,8 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { String requestId = Integer.toString(deviceRpcCallMsg.getRequestId()); metaData.putValue("requestId", requestId); metaData.putValue("requestUUID", requestUUID.toString()); - // ?? metaData.putValue("originServiceId", deviceRpcRequestMsg.get); + metaData.putValue("serviceId", deviceRpcCallMsg.getServiceId()); + metaData.putValue("sessionId", deviceRpcCallMsg.getSessionId()); metaData.putValue("expirationTime", Long.toString(deviceRpcCallMsg.getExpirationTime())); metaData.putValue("oneway", Boolean.toString(deviceRpcCallMsg.getOneway())); metaData.putValue(DataConstants.PERSISTENT, Boolean.toString(deviceRpcCallMsg.getPersisted())); @@ -403,16 +404,18 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() { @Override public void onSuccess(TbQueueMsgMetadata metadata) { - log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", device); + log.debug("Successfully send TO_SERVER_RPC_REQUEST to rule engine [{}], deviceRpcCallMsg {}", + device, deviceRpcCallMsg); } @Override public void onFailure(Throwable t) { - log.debug("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", device, t); + log.debug("Failed to send TO_SERVER_RPC_REQUEST to rule engine [{}], deviceRpcCallMsg {}", + device, deviceRpcCallMsg, t); } }); } catch (JsonProcessingException | IllegalArgumentException e) { - log.warn("[{}] Failed to push device action to rule engine: {}", deviceId, DataConstants.ENTITY_CREATED, e); + log.warn("[{}] Failed to push TO_SERVER_RPC_REQUEST to rule engine. deviceRpcCallMsg {}", deviceId, deviceRpcCallMsg, e); } return Futures.immediateFuture(null); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index a5785d6870..d70179cfe0 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -433,6 +433,8 @@ message DeviceRpcCallMsg { bool persisted = 10; int32 retries = 11; string additionalInfo = 12; + string serviceId = 13; + string sessionId = 14; } message RpcRequestMsg { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index 22d13f196b..89bc55bb54 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -15,6 +15,8 @@ */ package org.thingsboard.rule.engine.rpc; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -78,7 +80,11 @@ public class TbSendRPCReplyNode implements TbNode { ctx.tellFailure(msg, new RuntimeException("Request body is empty!")); } else { if (StringUtils.isNotBlank(msg.getMetaData().getValue(DataConstants.EDGE_ID))) { - saveRpcResponseToEdgeQueue(ctx, msg); + try { + saveRpcResponseToEdgeQueue(ctx, msg, serviceIdStr, sessionIdStr, requestIdStr); + } catch (Exception e) { + ctx.tellFailure(msg, e); + } } else { ctx.getRpcService().sendRpcReplyToDevice(serviceIdStr, UUID.fromString(sessionIdStr), Integer.parseInt(requestIdStr), msg.getData()); ctx.tellSuccess(msg); @@ -86,7 +92,7 @@ public class TbSendRPCReplyNode implements TbNode { } } - private void saveRpcResponseToEdgeQueue(TbContext ctx, TbMsg msg) { + private void saveRpcResponseToEdgeQueue(TbContext ctx, TbMsg msg, String serviceIdStr, String sessionIdStr, String requestIdStr) throws JsonProcessingException { // EdgeEvent edgeEvent = new EdgeEvent(); // edgeEvent.setTenantId(tenantId); // edgeEvent.setAction(eventAction); @@ -95,8 +101,11 @@ public class TbSendRPCReplyNode implements TbNode { // edgeEvent.setBody(entityBody); // edgeEvent.setEdgeId(edgeId); // -// ObjectNode body = mapper.createObjectNode(); -// body.put("requestId", requestId); + ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + body.put("serviceId", serviceIdStr); + body.put("sessionId", sessionIdStr); + body.put("requestId", requestIdStr); + body.put("response", JacksonUtil.OBJECT_MAPPER.writeValueAsString(msg.getData())); // body.put("requestUUID", msg.getId().toString()); // body.put("oneway", msg.isOneway()); // body.put("expirationTime", msg.getExpirationTime()); @@ -111,7 +120,7 @@ public class TbSendRPCReplyNode implements TbNode { // TODO: add body EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE, - EdgeEventActionType.RPC_CALL_RESPONSE, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree("{}")); + EdgeEventActionType.RPC_CALL_RESPONSE, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body)); ListenableFuture future = ctx.getEdgeEventService().saveAsync(edgeEvent); Futures.addCallback(future, new FutureCallback() { @@ -123,6 +132,7 @@ public class TbSendRPCReplyNode implements TbNode { @Override public void onFailure(Throwable t) { + ctx.tellFailure(msg, t); } }, ctx.getDbCallbackExecutor()); } From 20609722f52dd201c287cb7b6ce7776b65c5bb8a Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 15:30:45 +0200 Subject: [PATCH 04/16] Added RPC_REQUEST and RPC_RESPONSE to the front-end --- ui-ngx/src/app/shared/models/edge.models.ts | 10 ++++++++-- ui-ngx/src/assets/locale/locale.constant-en_US.json | 2 ++ 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/ui-ngx/src/app/shared/models/edge.models.ts b/ui-ngx/src/app/shared/models/edge.models.ts index 59f9d54318..6a84b6d29c 100644 --- a/ui-ngx/src/app/shared/models/edge.models.ts +++ b/ui-ngx/src/app/shared/models/edge.models.ts @@ -77,7 +77,9 @@ export enum EdgeEventActionType { UNASSIGNED_FROM_CUSTOMER = 'UNASSIGNED_FROM_CUSTOMER', RELATION_ADD_OR_UPDATE = 'RELATION_ADD_OR_UPDATE', RELATION_DELETED = 'RELATION_DELETED', - RPC_CALL = 'RPC_CALL', + RPC_CALL = 'RPC_CALL', // deprecated - to be removed in 4.x + RPC_CALL_REQUEST = 'RPC_CALL_REQUEST', + RPC_CALL_RESPONSE = 'RPC_CALL_RESPONSE', ALARM_ACK = 'ALARM_ACK', ALARM_CLEAR = 'ALARM_CLEAR', ASSIGNED_TO_EDGE = 'ASSIGNED_TO_EDGE', @@ -128,6 +130,8 @@ export const edgeEventActionTypeTranslations = new Map( diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 2b1eb11de7..c4031b6d0c 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1791,6 +1791,8 @@ "action-type-relation-add-or-update": "Relation Add or Update", "action-type-relation-deleted": "Relation Deleted", "action-type-rpc-call": "RPC Call", + "action-type-rpc-call-request": "RPC Call Request", + "action-type-rpc-call-response": "RPC Call Response", "action-type-alarm-ack": "Alarm Ack", "action-type-alarm-clear": "Alarm Clear", "action-type-assigned-to-edge": "Assigned to Edge", From 5ff8144f8db2f5f2fe95e9bb6e0658faeab78367 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 18:31:08 +0200 Subject: [PATCH 05/16] Code cleanup --- .../rpc/processor/DeviceEdgeProcessor.java | 16 +-------- common/edge-api/src/main/proto/edge.proto | 10 +++--- .../engine/edge/AbstractTbMsgPushNode.java | 5 ++- .../rule/engine/rpc/TbSendRPCReplyNode.java | 33 +++---------------- 4 files changed, 14 insertions(+), 50 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index 959f5c5fe8..147bb454fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -369,36 +369,22 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { private ListenableFuture processDeviceRpcRequestFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) { DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB())); - UUID requestUUID = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB()); try { - ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.createObjectNode(); TbMsgMetaData metaData = new TbMsgMetaData(); String requestId = Integer.toString(deviceRpcCallMsg.getRequestId()); metaData.putValue("requestId", requestId); - metaData.putValue("requestUUID", requestUUID.toString()); metaData.putValue("serviceId", deviceRpcCallMsg.getServiceId()); metaData.putValue("sessionId", deviceRpcCallMsg.getSessionId()); - metaData.putValue("expirationTime", Long.toString(deviceRpcCallMsg.getExpirationTime())); - metaData.putValue("oneway", Boolean.toString(deviceRpcCallMsg.getOneway())); - metaData.putValue(DataConstants.PERSISTENT, Boolean.toString(deviceRpcCallMsg.getPersisted())); - - if (deviceRpcCallMsg.getRetries() > 0) { - metaData.putValue(DataConstants.RETRIES, Integer.toString(deviceRpcCallMsg.getRetries())); - } - metaData.putValue(DataConstants.EDGE_ID, edge.getId().toString()); - Device device = deviceService.findDeviceById(tenantId, deviceId); if (device != null) { metaData.putValue("deviceName", device.getName()); metaData.putValue("deviceType", device.getType()); metaData.putValue(DataConstants.DEVICE_ID, deviceId.getId().toString()); } - + ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.createObjectNode(); entityNode.put("method", deviceRpcCallMsg.getRequestMsg().getMethod()); entityNode.put("params", deviceRpcCallMsg.getRequestMsg().getParams()); - - entityNode.put(DataConstants.ADDITIONAL_INFO, deviceRpcCallMsg.getAdditionalInfo()); TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, null, metaData, TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode)); tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() { diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index d70179cfe0..249f251238 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -430,11 +430,11 @@ message DeviceRpcCallMsg { bool oneway = 7; RpcRequestMsg requestMsg = 8; RpcResponseMsg responseMsg = 9; - bool persisted = 10; - int32 retries = 11; - string additionalInfo = 12; - string serviceId = 13; - string sessionId = 14; + optional bool persisted = 10; + optional int32 retries = 11; + optional string additionalInfo = 12; + optional string serviceId = 13; + optional string sessionId = 14; } message RpcRequestMsg { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java index 28dc64c068..f120276df3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java @@ -142,8 +142,11 @@ public abstract class AbstractTbMsgPushNode future = ctx.getEdgeEventService().saveAsync(edgeEvent); Futures.addCallback(future, new FutureCallback() { @Override From 1635d883a77c41ab4714f80030bf92594bb56fc0 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 18:32:51 +0200 Subject: [PATCH 06/16] Rename variable --- .../service/edge/rpc/processor/DeviceEdgeProcessor.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index 147bb454fa..4a517e19e9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -382,11 +382,11 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { metaData.putValue("deviceType", device.getType()); metaData.putValue(DataConstants.DEVICE_ID, deviceId.getId().toString()); } - ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.createObjectNode(); - entityNode.put("method", deviceRpcCallMsg.getRequestMsg().getMethod()); - entityNode.put("params", deviceRpcCallMsg.getRequestMsg().getParams()); + ObjectNode data = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + data.put("method", deviceRpcCallMsg.getRequestMsg().getMethod()); + data.put("params", deviceRpcCallMsg.getRequestMsg().getParams()); TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, null, metaData, - TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode)); + TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(data)); tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() { @Override public void onSuccess(TbQueueMsgMetadata metadata) { From 54f967e944b20ad5ea4b2e77ac67d6d5beacc34d Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 9 Nov 2022 18:35:18 +0200 Subject: [PATCH 07/16] Remove unused file --- .../rpc/processor/EdgeRpcRequestMetadata.java | 28 ------------------- 1 file changed, 28 deletions(-) delete mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java deleted file mode 100644 index a5a71d0743..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeRpcRequestMetadata.java +++ /dev/null @@ -1,28 +0,0 @@ -/** - * Copyright © 2016-2022 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.service.edge.rpc.processor; - -import lombok.Data; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.EdgeId; -import org.thingsboard.server.common.data.id.TenantId; - -@Data -public class EdgeRpcRequestMetadata { - private final TenantId tenantId; - private final EdgeId edgeId; - private final DeviceId deviceId; -} From 0a56f065f7218472d23818ac83e66c0daba9c3d9 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 10 Nov 2022 09:35:06 +0200 Subject: [PATCH 08/16] Merge RPC request/response actions --- .../device/DeviceActorMessageProcessor.java | 2 +- .../service/edge/rpc/EdgeGrpcSession.java | 3 +- .../rpc/constructor/DeviceMsgConstructor.java | 35 +++++++------------ .../rpc/processor/DeviceEdgeProcessor.java | 21 +++-------- .../server/edge/BaseDeviceEdgeTest.java | 2 +- .../common/data/edge/EdgeEventActionType.java | 3 +- .../rule/engine/rpc/TbSendRPCReplyNode.java | 2 +- ui-ngx/src/app/shared/models/edge.models.ts | 10 ++---- .../assets/locale/locale.constant-en_US.json | 2 -- 9 files changed, 25 insertions(+), 55 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 9b95b188d2..8a365b4e56 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -827,7 +827,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { body.put("retries", msg.getRetries()); body.put("additionalInfo", msg.getAdditionalInfo()); - EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL_REQUEST, deviceId, body); + EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body); return Futures.transform(systemContext.getEdgeEventService().saveAsync(edgeEvent), unused -> { systemContext.getClusterService().onEdgeEventUpdate(tenantId, edgeId); 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 aca2f8b562..957c699e32 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 @@ -463,8 +463,7 @@ public final class EdgeGrpcSession implements Closeable { case UNASSIGNED_FROM_CUSTOMER: case CREDENTIALS_REQUEST: case ENTITY_MERGE_REQUEST: - case RPC_CALL_REQUEST: - case RPC_CALL_RESPONSE: + case RPC_CALL: downlinkMsg = convertEntityEventToDownlink(edgeEvent); log.trace("[{}][{}] entity message processed [{}]", edgeEvent.getTenantId(), this.sessionId, downlinkMsg); break; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java index ebea5d482f..522bcd3a4e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java @@ -21,7 +21,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; -import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; @@ -97,30 +96,22 @@ public class DeviceMsgConstructor { .setIdLSB(deviceId.getId().getLeastSignificantBits()).build(); } - public DeviceRpcCallMsg constructDeviceRpcRequestMsg(UUID deviceId, JsonNode body) { + public DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode body) { DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body); - - String method = body.get("method").asText(); - String params = body.get("params").asText(); - RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder(); - requestBuilder.setMethod(method); - requestBuilder.setParams(params); - builder.setRequestMsg(requestBuilder.build()); - - return builder.build(); - } - - public DeviceRpcCallMsg constructDeviceRpcResponseMsg(UUID deviceId, JsonNode body) { - DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body); - - RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder(); - if (body.has("error")) { - responseBuilder.setError(body.get("error").asText()); + if (body.has("error") || body.has("response")) { + RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder(); + if (body.has("error")) { + responseBuilder.setError(body.get("error").asText()); + } else { + responseBuilder.setResponse(body.get("response").asText()); + } + builder.setResponseMsg(responseBuilder.build()); } else { - responseBuilder.setResponse(body.get("response").asText()); + RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder(); + requestBuilder.setMethod(body.get("method").asText()); + requestBuilder.setParams(body.get("params").asText()); + builder.setRequestMsg(requestBuilder.build()); } - builder.setResponseMsg(responseBuilder.build()); - return builder.build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index 4a517e19e9..81e7970611 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -451,10 +451,8 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { .build(); } break; - case RPC_CALL_REQUEST: - return convertRpcCallRequestEventToDownlink(edgeEvent); - case RPC_CALL_RESPONSE: - return convertRpcCallResponseEventToDownlink(edgeEvent); + case RPC_CALL: + return convertRpcCallEventToDownlink(edgeEvent); case CREDENTIALS_REQUEST: return convertCredentialsRequestEventToDownlink(edgeEvent); case ENTITY_MERGE_REQUEST: @@ -463,20 +461,11 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } - private DownlinkMsg convertRpcCallRequestEventToDownlink(EdgeEvent edgeEvent) { - log.trace("Executing convertRpcCallRequestEventToDownlink, edgeEvent [{}]", edgeEvent); + private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) { + log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent); return DownlinkMsg.newBuilder() .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcRequestMsg(edgeEvent.getEntityId(), edgeEvent.getBody())) - .build(); - } - - private DownlinkMsg convertRpcCallResponseEventToDownlink(EdgeEvent edgeEvent) { - log.trace("Executing convertRpcCallResponseEventToDownlink, edgeEvent [{}]", edgeEvent); - DeviceRpcCallMsg deviceRpcCallMsg = deviceMsgConstructor.constructDeviceRpcResponseMsg(edgeEvent.getEntityId(), edgeEvent.getBody()); - return DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addDeviceRpcCallMsg(deviceRpcCallMsg) + .addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody())) .build(); } diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java index 8e8c616d7a..be8089e4c0 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java @@ -512,7 +512,7 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { body.put("method", "test_method"); body.put("params", "{\"param1\":\"value1\"}"); - EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL_REQUEST, + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); edgeImitator.expectMessageAmount(1); edgeEventService.saveAsync(edgeEvent).get(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java index 7fae6c56b1..bba8767fca 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java @@ -28,8 +28,7 @@ public enum EdgeEventActionType { UNASSIGNED_FROM_CUSTOMER, RELATION_ADD_OR_UPDATE, RELATION_DELETED, - RPC_CALL_REQUEST, - RPC_CALL_RESPONSE, + RPC_CALL, ALARM_ACK, ALARM_CLEAR, ASSIGNED_TO_EDGE, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index cf0ebefc01..db3953f6cb 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -96,7 +96,7 @@ public class TbSendRPCReplyNode implements TbNode { EdgeId edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID))); DeviceId deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID))); EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE, - EdgeEventActionType.RPC_CALL_RESPONSE, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body)); + EdgeEventActionType.RPC_CALL, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body)); ListenableFuture future = ctx.getEdgeEventService().saveAsync(edgeEvent); Futures.addCallback(future, new FutureCallback() { @Override diff --git a/ui-ngx/src/app/shared/models/edge.models.ts b/ui-ngx/src/app/shared/models/edge.models.ts index 6a84b6d29c..59f9d54318 100644 --- a/ui-ngx/src/app/shared/models/edge.models.ts +++ b/ui-ngx/src/app/shared/models/edge.models.ts @@ -77,9 +77,7 @@ export enum EdgeEventActionType { UNASSIGNED_FROM_CUSTOMER = 'UNASSIGNED_FROM_CUSTOMER', RELATION_ADD_OR_UPDATE = 'RELATION_ADD_OR_UPDATE', RELATION_DELETED = 'RELATION_DELETED', - RPC_CALL = 'RPC_CALL', // deprecated - to be removed in 4.x - RPC_CALL_REQUEST = 'RPC_CALL_REQUEST', - RPC_CALL_RESPONSE = 'RPC_CALL_RESPONSE', + RPC_CALL = 'RPC_CALL', ALARM_ACK = 'ALARM_ACK', ALARM_CLEAR = 'ALARM_CLEAR', ASSIGNED_TO_EDGE = 'ASSIGNED_TO_EDGE', @@ -130,8 +128,6 @@ export const edgeEventActionTypeTranslations = new Map( diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 17b4533b03..3f3f70860f 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1791,8 +1791,6 @@ "action-type-relation-add-or-update": "Relation Add or Update", "action-type-relation-deleted": "Relation Deleted", "action-type-rpc-call": "RPC Call", - "action-type-rpc-call-request": "RPC Call Request", - "action-type-rpc-call-response": "RPC Call Response", "action-type-alarm-ack": "Alarm Ack", "action-type-alarm-clear": "Alarm Clear", "action-type-assigned-to-edge": "Assigned to Edge", From 15b26f4317f8675b73bd2cfc0ff55336c633daf1 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 10 Nov 2022 10:37:14 +0200 Subject: [PATCH 09/16] Remove Nullable annotation --- .../org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index db3953f6cb..796aa5658c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -20,7 +20,6 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.checkerframework.checker.nullness.qual.Nullable; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -98,9 +97,9 @@ public class TbSendRPCReplyNode implements TbNode { EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body)); ListenableFuture future = ctx.getEdgeEventService().saveAsync(edgeEvent); - Futures.addCallback(future, new FutureCallback() { + Futures.addCallback(future, new FutureCallback<>() { @Override - public void onSuccess(@Nullable Void result) { + public void onSuccess(Void result) { ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId); ctx.tellSuccess(msg); } From d8a609344289f9ae7de77712cbbd1793b1a912a1 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 10 Nov 2022 11:36:46 +0200 Subject: [PATCH 10/16] Added TbSendRPCReplyNodeTest --- .../server/edge/BaseDeviceEdgeTest.java | 4 + .../rule/engine/rpc/TbSendRPCReplyNode.java | 13 +- .../engine/rpc/TbSendRPCReplyNodeTest.java | 119 ++++++++++++++++++ 3 files changed, 134 insertions(+), 2 deletions(-) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java index be8089e4c0..ffe21ae961 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java @@ -511,6 +511,8 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { body.put("expirationTime", System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)); body.put("method", "test_method"); body.put("params", "{\"param1\":\"value1\"}"); + body.put("persisted", true); + body.put("retries", 2); EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); @@ -523,6 +525,8 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); + Assert.assertTrue(latestDeviceRpcCallMsg.getPersisted()); + Assert.assertEquals(2, latestDeviceRpcCallMsg.getRetries()); } private void sendAttributesRequestAndVerify(Device device, String scope, String attributesDataStr, String expectedKey, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index 796aa5658c..405a73944e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -87,13 +87,22 @@ public class TbSendRPCReplyNode implements TbNode { } private void saveRpcResponseToEdgeQueue(TbContext ctx, TbMsg msg, String serviceIdStr, String sessionIdStr, String requestIdStr) { + EdgeId edgeId; + DeviceId deviceId; + try { + edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID))); + deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID))); + } catch (Exception e) { + String errMsg = String.format("[%s] Failed to parse edgeId or deviceId from metadata %s!", ctx.getTenantId(), msg.getMetaData()); + ctx.tellFailure(msg, new RuntimeException(errMsg)); + return; + } + ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode(); body.put("serviceId", serviceIdStr); body.put("sessionId", sessionIdStr); body.put("requestId", requestIdStr); body.put("response", msg.getData()); - EdgeId edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID))); - DeviceId deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID))); EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body)); ListenableFuture future = ctx.getEdgeEventService().saveAsync(edgeEvent); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java new file mode 100644 index 0000000000..095505aff3 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java @@ -0,0 +1,119 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.rule.engine.rpc; + +import com.google.common.util.concurrent.SettableFuture; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.common.util.ListeningExecutor; +import org.thingsboard.rule.engine.api.RuleEngineRpcService; +import org.thingsboard.rule.engine.api.TbContext; +import org.thingsboard.rule.engine.api.TbNodeConfiguration; +import org.thingsboard.rule.engine.api.TbNodeException; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgDataType; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.session.SessionMsgType; +import org.thingsboard.server.dao.edge.EdgeEventService; + +import java.util.UUID; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +@RunWith(MockitoJUnitRunner.class) +public class TbSendRPCReplyNodeTest { + + private static final String DUMMY_SERVICE_ID = "testServiceId"; + private static final int DUMMY_REQUEST_ID = 0; + private static final UUID DUMMY_SESSION_ID = UUID.randomUUID(); + private static final String DUMMY_DATA = "{\"key\":\"value\"}"; + + TbSendRPCReplyNode node; + + private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); + private final DeviceId deviceId = new DeviceId(UUID.randomUUID()); + + @Mock + private TbContext ctx; + + @Mock + private RuleEngineRpcService rpcService; + + @Mock + private EdgeEventService edgeEventService; + + @Mock + private ListeningExecutor listeningExecutor; + + @Before + public void setUp() throws TbNodeException { + node = new TbSendRPCReplyNode(); + TbSendRpcReplyNodeConfiguration config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration(); + node.init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + } + + @Test + public void sendReplyToTransport() { + Mockito.when(ctx.getRpcService()).thenReturn(rpcService); + + + TbMsg msg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, getDefaultMetadata(), + TbMsgDataType.JSON, DUMMY_DATA, null, null); + + node.onMsg(ctx, msg); + + verify(rpcService).sendRpcReplyToDevice(DUMMY_SERVICE_ID, DUMMY_SESSION_ID, DUMMY_REQUEST_ID, DUMMY_DATA); + verify(edgeEventService, never()).saveAsync(any()); + } + + @Test + public void sendReplyToEdgeQueue() { + Mockito.when(ctx.getTenantId()).thenReturn(tenantId); + Mockito.when(ctx.getEdgeEventService()).thenReturn(edgeEventService); + Mockito.when(edgeEventService.saveAsync(any())).thenReturn(SettableFuture.create()); + Mockito.when(ctx.getDbCallbackExecutor()).thenReturn(listeningExecutor); + + TbMsgMetaData defaultMetadata = getDefaultMetadata(); + defaultMetadata.putValue(DataConstants.EDGE_ID, UUID.randomUUID().toString()); + defaultMetadata.putValue(DataConstants.DEVICE_ID, UUID.randomUUID().toString()); + TbMsg msg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, defaultMetadata, + TbMsgDataType.JSON, DUMMY_DATA, null, null); + + node.onMsg(ctx, msg); + + verify(edgeEventService).saveAsync(any()); + verify(rpcService, never()).sendRpcReplyToDevice(DUMMY_SERVICE_ID, DUMMY_SESSION_ID, DUMMY_REQUEST_ID, DUMMY_DATA); + } + + private TbMsgMetaData getDefaultMetadata() { + TbSendRpcReplyNodeConfiguration config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration(); + TbMsgMetaData metadata = new TbMsgMetaData(); + metadata.putValue(config.getServiceIdMetaDataAttribute(), DUMMY_SERVICE_ID); + metadata.putValue(config.getSessionIdMetaDataAttribute(), DUMMY_SESSION_ID.toString()); + metadata.putValue(config.getRequestIdMetaDataAttribute(), Integer.toString(DUMMY_REQUEST_ID)); + return metadata; + } +} From 6223307cd8eb7e5ab7d5c022762a2926ef658f75 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 10 Nov 2022 17:04:48 +0200 Subject: [PATCH 11/16] Add validation for Asset Profile FK in Rule Chains --- .../org/thingsboard/server/dao/rule/BaseRuleChainService.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 420739e37c..1d11d34b2d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -731,6 +731,8 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_device_profile")) { throw new DataValidationException("The rule chain referenced by the device profiles cannot be deleted!"); + } else if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_asset_profile")) { + throw new DataValidationException("The rule chain referenced by the asset profiles cannot be deleted!"); } else { throw t; } From e32bd456b7875555f0b0b0654288f297e368fa0f Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 14 Nov 2022 11:07:09 +0100 Subject: [PATCH 12/16] added memory usage log to the js-executors --- msa/js-executor/api/jsInvokeMessageProcessor.ts | 5 +++++ msa/js-executor/config/custom-environment-variables.yml | 1 + msa/js-executor/config/default.yml | 1 + 3 files changed, 7 insertions(+) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.ts b/msa/js-executor/api/jsInvokeMessageProcessor.ts index 668cd61f50..e69209fa01 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.ts +++ b/msa/js-executor/api/jsInvokeMessageProcessor.ts @@ -39,6 +39,7 @@ const TIMEOUT_ERROR = 2; const NOT_FOUND_ERROR = 3; const statFrequency = Number(config.get('script.stat_print_frequency')); +const memoryUsageTraceFrequency = Number(config.get('script.memory_usage_trace_frequency')); const scriptBodyTraceFrequency = Number(config.get('script.script_body_trace_frequency')); const useSandbox = config.get('script.use_sandbox') === 'true'; const maxActiveScripts = Number(config.get('script.max_active_scripts')); @@ -167,6 +168,10 @@ export class JsInvokeMessageProcessor { if (this.executedScriptsCounter % scriptBodyTraceFrequency == 0) { this.logger.info('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); } + if (this.executedScriptsCounter % memoryUsageTraceFrequency == 0) { + this.logger.info('Current memory usage: [%s]', process.memoryUsage()); + } + this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then( (script) => { this.executor.executeScript(script, invokeRequest.args, invokeRequest.timeout).then( diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index b9c24c8d8d..2ebea4ccc1 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -75,6 +75,7 @@ logger: script: use_sandbox: "SCRIPT_USE_SANDBOX" + memory_usage_trace_frequency: "MEMORY_USAGE_TRACE_FREQUENCY" stat_print_frequency: "SCRIPT_STAT_PRINT_FREQUENCY" script_body_trace_frequency: "SCRIPT_BODY_TRACE_FREQUENCY" max_active_scripts: "MAX_ACTIVE_SCRIPTS" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 64829ef792..805c175dce 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -64,6 +64,7 @@ logger: script: use_sandbox: "true" + memory_usage_trace_frequency: "10000" script_body_trace_frequency: "10000" stat_print_frequency: "10000" max_active_scripts: "1000" From 11bd9b5257344b60c7e5eb1a2bbe6af175ffde5b Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 14 Nov 2022 11:28:30 +0100 Subject: [PATCH 13/16] added nodejs memory leak workaround --- docker/tb-js-executor.env | 3 ++- msa/js-executor/config/default.yml | 2 +- msa/js-executor/docker/start-js-executor.sh | 2 +- 3 files changed, 4 insertions(+), 3 deletions(-) diff --git a/docker/tb-js-executor.env b/docker/tb-js-executor.env index e080906549..1938449d53 100644 --- a/docker/tb-js-executor.env +++ b/docker/tb-js-executor.env @@ -3,4 +3,5 @@ LOGGER_LEVEL=info LOG_FOLDER=logs LOGGER_FILENAME=tb-js-executor-%DATE%.log DOCKER_MODE=true -SCRIPT_BODY_TRACE_FREQUENCY=1000 \ No newline at end of file +SCRIPT_BODY_TRACE_FREQUENCY=1000 +NODE_OPTIONS="--max-old-space-size=200" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 805c175dce..96f3401da5 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -64,7 +64,7 @@ logger: script: use_sandbox: "true" - memory_usage_trace_frequency: "10000" + memory_usage_trace_frequency: "1000" script_body_trace_frequency: "10000" stat_print_frequency: "10000" max_active_scripts: "1000" diff --git a/msa/js-executor/docker/start-js-executor.sh b/msa/js-executor/docker/start-js-executor.sh index 575f93c389..b27c1b7167 100755 --- a/msa/js-executor/docker/start-js-executor.sh +++ b/msa/js-executor/docker/start-js-executor.sh @@ -27,4 +27,4 @@ source "${CONF_FOLDER}/${configfile}" cd ${pkg.installFolder} # This will forward this PID 1 to the node.js and forward SIGTERM for graceful shutdown as well -exec node server.js +exec --no-compilation-cache node server.js From b02a2215db8f01b41a18049790d411ba5a3be647 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Mon, 14 Nov 2022 12:33:17 +0200 Subject: [PATCH 14/16] Update MVEL version to 2.4.24TB --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 60e8783bca..e7e75c32f8 100755 --- a/pom.xml +++ b/pom.xml @@ -77,7 +77,7 @@ 3.5.5 3.21.9 1.42.1 - 2.4.23TB + 2.4.24TB 1.18.18 1.2.4 4.1.75.Final From 5b1642246d722a67869ee4bb8ecca511cc138608 Mon Sep 17 00:00:00 2001 From: Yevhen Bondarenko <56396344+YevhenBondarenko@users.noreply.github.com> Date: Mon, 14 Nov 2022 13:30:05 +0100 Subject: [PATCH 15/16] [3.4.2] fix start js (#7614) * refactoring * fixed start js-executor typo --- msa/js-executor/docker/start-js-executor.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/msa/js-executor/docker/start-js-executor.sh b/msa/js-executor/docker/start-js-executor.sh index b27c1b7167..d30b62c145 100755 --- a/msa/js-executor/docker/start-js-executor.sh +++ b/msa/js-executor/docker/start-js-executor.sh @@ -27,4 +27,4 @@ source "${CONF_FOLDER}/${configfile}" cd ${pkg.installFolder} # This will forward this PID 1 to the node.js and forward SIGTERM for graceful shutdown as well -exec --no-compilation-cache node server.js +exec node --no-compilation-cache server.js From 51ec17d9d1782b755bd8aea10eeac1897c050166 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Mon, 14 Nov 2022 14:41:10 +0200 Subject: [PATCH 16/16] Update MVEL version to 2.4.25TB --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index e7e75c32f8..5f5aa3703b 100755 --- a/pom.xml +++ b/pom.xml @@ -77,7 +77,7 @@ 3.5.5 3.21.9 1.42.1 - 2.4.24TB + 2.4.25TB 1.18.18 1.2.4 4.1.75.Final