From 373f2f9f0e6e0b86d8e3f7ebeb7b9f6c7855c8d8 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 4 Aug 2021 18:03:54 +0300 Subject: [PATCH] added additionalInfo to the rpc --- application/src/main/data/upgrade/3.2.2/schema_update.sql | 1 + .../server/actors/device/DeviceActorMessageProcessor.java | 1 + .../server/controller/AbstractRpcController.java | 4 +++- .../server/service/rpc/DefaultTbCoreDeviceRpcService.java | 2 ++ .../server/service/rpc/DefaultTbRuleEngineRpcService.java | 2 +- .../org/thingsboard/server/common/data/DataConstants.java | 1 + .../java/org/thingsboard/server/common/data/rpc/Rpc.java | 2 ++ .../server/common/msg/rpc/ToDeviceRpcRequest.java | 3 +++ .../org/thingsboard/server/dao/model/ModelConstants.java | 1 + .../org/thingsboard/server/dao/model/sql/RpcEntity.java | 7 +++++++ dao/src/main/resources/sql/schema-entities-hsql.sql | 1 + dao/src/main/resources/sql/schema-entities.sql | 1 + .../rule/engine/api/RuleEngineDeviceRpcRequest.java | 1 + .../thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java | 3 +++ 14 files changed, 28 insertions(+), 2 deletions(-) diff --git a/application/src/main/data/upgrade/3.2.2/schema_update.sql b/application/src/main/data/upgrade/3.2.2/schema_update.sql index 50f538d5dc..f0077287e9 100644 --- a/application/src/main/data/upgrade/3.2.2/schema_update.sql +++ b/application/src/main/data/upgrade/3.2.2/schema_update.sql @@ -209,6 +209,7 @@ CREATE TABLE IF NOT EXISTS rpc ( expiration_time bigint NOT NULL, request varchar(10000000) NOT NULL, response varchar(10000000), + additional_info varchar(10000000), status varchar(255) NOT NULL ); 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 c3aa3cd9d6..ac0a59e0cf 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 @@ -234,6 +234,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { rpc.setExpirationTime(request.getExpirationTime()); rpc.setRequest(JacksonUtil.valueToTree(request)); rpc.setStatus(status); + rpc.setAdditionalInfo(JacksonUtil.valueToTree(request.getAdditionalInfo())); return systemContext.getTbRpcService().save(tenantId, rpc); } diff --git a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java index 7b2ec1f02f..f7aa23a37b 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java @@ -90,6 +90,7 @@ public abstract class AbstractRpcController extends BaseController { long expTime = System.currentTimeMillis() + Math.max(minTimeout, timeout); UUID rpcRequestUUID = rpcRequestBody.has("requestUUID") ? UUID.fromString(rpcRequestBody.get("requestUUID").asText()) : UUID.randomUUID(); boolean persisted = rpcRequestBody.has(DataConstants.PERSISTENT) && rpcRequestBody.get(DataConstants.PERSISTENT).asBoolean(); + String additionalInfo = JacksonUtil.toString(rpcRequestBody.get(DataConstants.ADDITIONAL_INFO)); accessValidator.validate(currentUser, Operation.RPC_CALL, deviceId, new HttpValidationCallback(response, new FutureCallback<>() { @Override public void onSuccess(@Nullable DeferredResult result) { @@ -99,7 +100,8 @@ public abstract class AbstractRpcController extends BaseController { oneWay, expTime, body, - persisted + persisted, + additionalInfo ); deviceRpcService.processRestApiRpcRequest(rpcRequest, fromDeviceRpcResponse -> reply(new LocalRequestMetaData(rpcRequest, currentUser, result), fromDeviceRpcResponse, timeoutStatus, noActiveConnectionStatus), currentUser); } diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java index 1ef8d9361e..9fbbb142cf 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java @@ -168,6 +168,8 @@ public class DefaultTbCoreDeviceRpcService implements TbCoreDeviceRpcService { entityNode.put("method", msg.getBody().getMethod()); entityNode.put("params", msg.getBody().getParams()); + entityNode.put(DataConstants.ADDITIONAL_INFO, msg.getAdditionalInfo()); + try { TbMsg tbMsg = TbMsg.newMsg(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE, msg.getDeviceId(), currentUser.getCustomerId(), metaData, TbMsgDataType.JSON, json.writeValueAsString(entityNode)); clusterService.pushMsgToRuleEngine(msg.getTenantId(), msg.getDeviceId(), tbMsg, null); diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index e16af3695b..505c8db7d2 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -100,7 +100,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi @Override public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer consumer) { ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(), - src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted()); + src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); forwardRpcRequestToDeviceActor(request, response -> { if (src.isRestApiCall()) { sendRpcResponseToTbCore(src.getOriginServiceId(), response); 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 329e74d702..d42efcd770 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 @@ -36,6 +36,7 @@ public class DataConstants { public static final String ALARM_CONDITION_REPEATS = "alarmConditionRepeats"; public static final String ALARM_CONDITION_DURATION = "alarmConditionDuration"; public static final String PERSISTENT = "persistent"; + public static final String ADDITIONAL_INFO = "additionalInfo"; 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/rpc/Rpc.java b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/Rpc.java index 504c89b569..27326087f7 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rpc/Rpc.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/Rpc.java @@ -33,6 +33,7 @@ public class Rpc extends BaseData implements HasTenantId { private JsonNode request; private JsonNode response; private RpcStatus status; + private JsonNode additionalInfo; public Rpc() { super(); @@ -50,5 +51,6 @@ public class Rpc extends BaseData implements HasTenantId { this.request = rpc.getRequest(); this.response = rpc.getResponse(); this.status = rpc.getStatus(); + this.additionalInfo = rpc.getAdditionalInfo(); } } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java index b3b33146d0..f9bb2b3810 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.msg.rpc; +import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; @@ -35,5 +36,7 @@ public class ToDeviceRpcRequest implements Serializable { private final long expirationTime; private final ToDeviceRpcRequestBody body; private final boolean persisted; + @JsonIgnore + private final String additionalInfo; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 525630fc27..e8e90f2d3c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -519,6 +519,7 @@ public class ModelConstants { public static final String RPC_REQUEST = "request"; public static final String RPC_RESPONSE = "response"; public static final String RPC_STATUS = "status"; + public static final String RPC_ADDITIONAL_INFO = ADDITIONAL_INFO_PROPERTY; /** * Edge constants. diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/RpcEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/RpcEntity.java index a8823cb8cd..d68aa1afb8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/RpcEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/RpcEntity.java @@ -36,6 +36,7 @@ import javax.persistence.Enumerated; import javax.persistence.Table; import java.util.UUID; +import static org.thingsboard.server.dao.model.ModelConstants.RPC_ADDITIONAL_INFO; import static org.thingsboard.server.dao.model.ModelConstants.RPC_DEVICE_ID; import static org.thingsboard.server.dao.model.ModelConstants.RPC_EXPIRATION_TIME; import static org.thingsboard.server.dao.model.ModelConstants.RPC_REQUEST; @@ -72,6 +73,10 @@ public class RpcEntity extends BaseSqlEntity implements BaseEntity { @Column(name = RPC_STATUS) private RpcStatus status; + @Type(type = "json") + @Column(name = RPC_ADDITIONAL_INFO) + private JsonNode additionalInfo; + public RpcEntity() { super(); } @@ -85,6 +90,7 @@ public class RpcEntity extends BaseSqlEntity implements BaseEntity { this.request = rpc.getRequest(); this.response = rpc.getResponse(); this.status = rpc.getStatus(); + this.additionalInfo = rpc.getAdditionalInfo(); } @Override @@ -97,6 +103,7 @@ public class RpcEntity extends BaseSqlEntity implements BaseEntity { rpc.setRequest(request); rpc.setResponse(response); rpc.setStatus(status); + rpc.setAdditionalInfo(additionalInfo); return rpc; } } diff --git a/dao/src/main/resources/sql/schema-entities-hsql.sql b/dao/src/main/resources/sql/schema-entities-hsql.sql index f1180748c8..e58ca10459 100644 --- a/dao/src/main/resources/sql/schema-entities-hsql.sql +++ b/dao/src/main/resources/sql/schema-entities-hsql.sql @@ -582,5 +582,6 @@ CREATE TABLE IF NOT EXISTS rpc ( expiration_time bigint NOT NULL, request varchar(10000000) NOT NULL, response varchar(10000000), + additional_info varchar(10000000), status varchar(255) NOT NULL ); diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 68d6521167..69760c377a 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -616,6 +616,7 @@ CREATE TABLE IF NOT EXISTS rpc ( expiration_time bigint NOT NULL, request varchar(10000000) NOT NULL, response varchar(10000000), + additional_info varchar(10000000), status varchar(255) NOT NULL ); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java index 1d74040b34..afdee27b74 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java @@ -40,5 +40,6 @@ public final class RuleEngineDeviceRpcRequest { private final String body; private final long expirationTime; private final boolean restApiCall; + private final String additionalInfo; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index 1018d70a2a..336dbeb6ce 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -100,6 +100,8 @@ public class TbSendRPCRequestNode implements TbNode { params = gson.toJson(paramsEl); } + String additionalInfo = gson.toJson(json.get(DataConstants.ADDITIONAL_INFO)); + RuleEngineDeviceRpcRequest request = RuleEngineDeviceRpcRequest.builder() .oneway(oneway) .method(json.get("method").getAsString()) @@ -112,6 +114,7 @@ public class TbSendRPCRequestNode implements TbNode { .expirationTime(expirationTime) .restApiCall(restApiCall) .persisted(persisted) + .additionalInfo(additionalInfo) .build(); ctx.getRpcService().sendRpcRequestToDevice(request, ruleEngineDeviceRpcResponse -> {