diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 283868a4d2..62d4829085 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -42,6 +42,7 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.HasRuleEngineProfile; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.alarm.Alarm; @@ -339,7 +340,7 @@ class DefaultTbContext implements TbContext { if (device.getDeviceProfileId() != null) { deviceProfile = mainCtx.getDeviceProfileCache().find(device.getDeviceProfileId()); } - return deviceActionMsg(device, device.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, deviceProfile); + return entityActionMsg(device, device.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, deviceProfile); } public TbMsg assetCreatedMsg(Asset asset, RuleNodeId ruleNodeId) { @@ -347,20 +348,19 @@ class DefaultTbContext implements TbContext { if (asset.getAssetProfileId() != null) { assetProfile = mainCtx.getAssetProfileCache().find(asset.getAssetProfileId()); } - return assetActionMsg(asset, asset.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, assetProfile); + return entityActionMsg(asset, asset.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, assetProfile); } public TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action) { + HasRuleEngineProfile profile = null; if (EntityType.DEVICE.equals(alarm.getOriginator().getEntityType())) { DeviceId deviceId = new DeviceId(alarm.getOriginator().getId()); - DeviceProfile deviceProfile = mainCtx.getDeviceProfileCache().get(getTenantId(), deviceId); - return deviceActionMsg(alarm, alarm.getOriginator(), ruleNodeId, action, deviceProfile); + profile = mainCtx.getDeviceProfileCache().get(getTenantId(), deviceId); } else if (EntityType.ASSET.equals(alarm.getOriginator().getEntityType())) { AssetId assetId = new AssetId(alarm.getOriginator().getId()); - AssetProfile assetProfile = mainCtx.getAssetProfileCache().get(getTenantId(), assetId); - return assetActionMsg(alarm, alarm.getOriginator(), ruleNodeId, action, assetProfile); + profile = mainCtx.getAssetProfileCache().get(getTenantId(), assetId); } - return entityActionMsg(alarm, alarm.getOriginator(), ruleNodeId, action); + return entityActionMsg(alarm, alarm.getOriginator(), ruleNodeId, action, profile); } public TbMsg attributesUpdatedActionMsg(EntityId originator, RuleNodeId ruleNodeId, String scope, List attributes) { @@ -383,16 +383,15 @@ class DefaultTbContext implements TbContext { private TbMsg attributesActionMsg(EntityId originator, RuleNodeId ruleNodeId, String scope, String action, String msgData) { TbMsgMetaData tbMsgMetaData = getActionMetaData(ruleNodeId); tbMsgMetaData.putValue("scope", scope); + HasRuleEngineProfile profile = null; if (EntityType.DEVICE.equals(originator.getEntityType())) { DeviceId deviceId = new DeviceId(originator.getId()); - DeviceProfile deviceProfile = mainCtx.getDeviceProfileCache().get(getTenantId(), deviceId); - return deviceActionMsg(originator, tbMsgMetaData, msgData, action, deviceProfile); + profile = mainCtx.getDeviceProfileCache().get(getTenantId(), deviceId); } else if (EntityType.ASSET.equals(originator.getEntityType())) { AssetId assetId = new AssetId(originator.getId()); - AssetProfile assetProfile = mainCtx.getAssetProfileCache().get(getTenantId(), assetId); - return assetActionMsg(originator, tbMsgMetaData, msgData, action, assetProfile); + profile = mainCtx.getAssetProfileCache().get(getTenantId(), assetId); } - return entityActionMsg(originator, tbMsgMetaData, msgData, action, null, null); + return entityActionMsg(originator, tbMsgMetaData, msgData, action, profile); } @Override @@ -400,52 +399,26 @@ class DefaultTbContext implements TbContext { mainCtx.getClusterService().onEdgeEventUpdate(tenantId, edgeId); } - public TbMsg deviceActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action, DeviceProfile deviceProfile) { - try { - return deviceActionMsg(id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)), action, deviceProfile); - } catch (JsonProcessingException | IllegalArgumentException e) { - throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " " + action + " msg: " + e); - } + public TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) { + return entityActionMsg(entity, id, ruleNodeId, action, null); } - public TbMsg assetActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action, AssetProfile assetProfile) { + public TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action, K profile) { try { - return assetActionMsg(id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)), action, assetProfile); + return entityActionMsg(id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)), action, profile); } catch (JsonProcessingException | IllegalArgumentException e) { throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " " + action + " msg: " + e); } } - public TbMsg assetActionMsg(I id, TbMsgMetaData msgMetaData, String msgData, String action, AssetProfile assetProfile) { - RuleChainId ruleChainId = null; - String queueName = null; - if (assetProfile != null) { - ruleChainId = assetProfile.getDefaultRuleChainId(); - queueName = assetProfile.getDefaultQueueName(); + private TbMsg entityActionMsg(I id, TbMsgMetaData msgMetaData, String msgData, String action, K profile) { + String defaultQueueName = null; + RuleChainId defaultRuleChainId = null; + if (profile != null) { + defaultQueueName = profile.getDefaultQueueName(); + defaultRuleChainId = profile.getDefaultRuleChainId(); } - return entityActionMsg(id, msgMetaData, msgData, action, ruleChainId, queueName); - } - - public TbMsg deviceActionMsg(I id, TbMsgMetaData msgMetaData, String msgData, String action, DeviceProfile deviceProfile) { - RuleChainId ruleChainId = null; - String queueName = null; - if (deviceProfile != null) { - ruleChainId = deviceProfile.getDefaultRuleChainId(); - queueName = deviceProfile.getDefaultQueueName(); - } - return entityActionMsg(id, msgMetaData, msgData, action, ruleChainId, queueName); - } - - public TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) { - try { - return entityActionMsg(id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)), action, null, null); - } catch (JsonProcessingException | IllegalArgumentException e) { - throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " " + action + " msg: " + e); - } - } - - private TbMsg entityActionMsg(I id, TbMsgMetaData msgMetaData, String msgData, String action, RuleChainId ruleChainId, String queueName) { - return TbMsg.newMsg(queueName, action, id, msgMetaData, msgData, ruleChainId, null); + return TbMsg.newMsg(defaultQueueName, action, id, msgMetaData, msgData, defaultRuleChainId, null); } @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java b/common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java index da8edc695e..a66f3a40d9 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java @@ -45,7 +45,7 @@ import static org.thingsboard.server.common.data.SearchTextBasedWithAdditionalIn @ToString(exclude = {"image", "profileDataBytes"}) @EqualsAndHashCode(callSuper = true) @Slf4j -public class DeviceProfile extends SearchTextBased implements HasName, HasTenantId, HasOtaPackage, ExportableEntity { +public class DeviceProfile extends SearchTextBased implements HasName, HasTenantId, HasOtaPackage, HasRuleEngineProfile, ExportableEntity { private static final long serialVersionUID = 6998485460273302018L; @@ -141,7 +141,7 @@ public class DeviceProfile extends SearchTextBased implements H } @ApiModelProperty(position = 5, value = "Used to mark the default profile. Default profile is used when the device profile is not specified during device creation.") - public boolean isDefault(){ + public boolean isDefault() { return isDefault; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/HasRuleEngineProfile.java b/common/data/src/main/java/org/thingsboard/server/common/data/HasRuleEngineProfile.java new file mode 100644 index 0000000000..450a2235a7 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/HasRuleEngineProfile.java @@ -0,0 +1,26 @@ +/** + * 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.common.data; + +import org.thingsboard.server.common.data.id.RuleChainId; + +public interface HasRuleEngineProfile { + + RuleChainId getDefaultRuleChainId(); + + String getDefaultQueueName(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetProfile.java b/common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetProfile.java index 1bd9d9c6a5..01657feed6 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetProfile.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetProfile.java @@ -23,6 +23,7 @@ import lombok.ToString; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.ExportableEntity; import org.thingsboard.server.common.data.HasName; +import org.thingsboard.server.common.data.HasRuleEngineProfile; import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.SearchTextBased; import org.thingsboard.server.common.data.id.AssetProfileId; @@ -37,7 +38,7 @@ import org.thingsboard.server.common.data.validation.NoXss; @ToString(exclude = {"image"}) @EqualsAndHashCode(callSuper = true) @Slf4j -public class AssetProfile extends SearchTextBased implements HasName, HasTenantId, ExportableEntity { +public class AssetProfile extends SearchTextBased implements HasName, HasTenantId, HasRuleEngineProfile, ExportableEntity { private static final long serialVersionUID = 6998485460273302018L; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index 44cdc82c36..b810774832 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -66,6 +66,10 @@ public class TbMsgAttributesNode implements TbNode { } String src = msg.getData(); List attributes = new ArrayList<>(JsonConverter.convertToAttributes(new JsonParser().parse(src))); + if (attributes.isEmpty()) { + ctx.tellSuccess(msg); + return; + } String notifyDeviceStr = msg.getMetaData().getValue("notifyDevice"); ctx.getTelemetryService().saveAndNotify( ctx.getTenantId(),